This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/main/pr-6729-f9c7b613ef601e38bf8a2134463a623fa94def9a
in repository https://gitbox.apache.org/repos/asf/texera.git

commit 43f26495f36002ed2fa3aa3e0fc0c0e478cf37b0
Author: Xiaozhen Liu <[email protected]>
AuthorDate: Mon Sep 21 20:46:54 2026 +0000

    feat(amber): carry cache-reuse status as a metrics flag (#6729)
    
    ### What changes were proposed in this PR?
    
    This PR was reworked after the discussion in #5880. The old version
    added a new CACHE_REUSED state; the review threads below refer to that
    version. The discussion concluded that a reused operator should just
    report COMPLETED, because every place that checks state treats completed
    and reused the same way. What still needs to travel is one bit: whether
    the operator's results came from the cache.
    
    So the PR now adds only that bit:
    
    - A `reused_from_cache` boolean on `OperatorMetrics`. A reused operator
    still reports COMPLETED.
    - A logical operator counts as reused only when all of its physical
    operators are (`aggregateMetrics`).
    - The statistics websocket event and the TS `OperatorStatistics` type
    carry the flag to the frontend. No UI changes here; that is #5886.
    - `ExecutionStatsService.computeStatsDiff` no longer rebuilds each
    `OperatorMetrics` field by field. The rebuild copied every existing
    field, so removing it changes nothing today, but it would have silently
    reset the new flag to false on the statistics persistence path.
    
    Nothing sets the flag yet. The producer comes with #5884, so with an
    empty cache the engine behaves exactly like main.
    
    ### Any related issues, documentation, discussions?
    
    Closes #5883. Part of #5881. Design discussion: #5880. Related: #5884.
    
    ### How was this PR tested?
    
    New unit tests in ExecutionUtilsSpec cover the all-physical-operators
    rule and the flag staying false when nothing sets it. Existing specs
    pass unchanged, scalafmt is clean, and the frontend production build
    passes.
    
    ### Was this PR authored or co-authored using generative AI tooling?
    
    Generated-by: Claude (Claude Code)
---
 .../engine/common/executionruntimestate.proto      |  4 ++++
 .../coordinator/execution/ExecutionUtils.scala     |  9 +++++++-
 .../event/OperatorStatisticsUpdateEvent.scala      |  6 +++++-
 .../texera/web/service/ExecutionStatsService.scala | 21 +++----------------
 .../coordinator/execution/ExecutionUtilsSpec.scala | 24 ++++++++++++++++++++--
 .../websocket/event/TexeraWebSocketEventSpec.scala |  5 ++++-
 .../workspace/types/execute-workflow.interface.ts  |  2 ++
 7 files changed, 48 insertions(+), 23 deletions(-)

diff --git 
a/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto
 
b/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto
index e712b3adc8..e384e80c80 100644
--- 
a/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto
+++ 
b/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto
@@ -82,6 +82,10 @@ message OperatorStatistics{
 message OperatorMetrics{
   architecture.rpc.WorkflowAggregatedState operator_state = 1 
[(scalapb.field).no_box = true];
   OperatorStatistics operator_statistics = 2 [(scalapb.field).no_box = true];
+  // True when the operator's results were reused from the operator port cache
+  // instead of being computed by workers. Provenance of a completed operator,
+  // not a distinct state; the operator still reports COMPLETED.
+  bool reused_from_cache = 3;
 }
 
 message ExecutionStatsStore {
diff --git 
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala
 
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala
index 666aeece42..303fba9a1f 100644
--- 
a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala
+++ 
b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala
@@ -77,7 +77,14 @@ object ExecutionUtils {
         dataProcessingTimeSum,
         controlProcessingTimeSum,
         idleTimeSum
-      )
+      ),
+      // Fully-reused semantics: partial reuse is possible (HashJoin's build 
and
+      // probe sit in different regions), and a partially reused operator 
reports
+      // false; per-port detail comes from the cache entries. The input holds 
the
+      // operators that currently have a region execution, so this shares the
+      // state field's transient window until all regions exist. Non-empty 
here,
+      // so the forall cannot hold vacuously.
+      reusedFromCache = metrics.forall(_.reusedFromCache)
     )
   }
 
diff --git 
a/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala
 
b/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala
index d4aa6117c9..af91f65784 100644
--- 
a/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala
+++ 
b/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala
@@ -30,7 +30,11 @@ case class OperatorAggregatedMetrics(
     numWorkers: Long,
     aggregatedDataProcessingTime: Long,
     aggregatedControlProcessingTime: Long,
-    aggregatedIdleTime: Long
+    aggregatedIdleTime: Long,
+    // Provenance: the operator completed by reusing cached results (no 
workers ran).
+    // Deliberately no default: a new construction site must decide the flag
+    // explicitly instead of silently sending false.
+    reusedFromCache: Boolean
 )
 
 case class OperatorStatisticsUpdateEvent(operatorStatistics: Map[String, 
OperatorAggregatedMetrics])
diff --git 
a/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala
 
b/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala
index ac33478451..2c173c25a8 100644
--- 
a/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala
+++ 
b/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala
@@ -125,7 +125,8 @@ class ExecutionStatsService(
                 metrics.operatorStatistics.numWorkers,
                 metrics.operatorStatistics.dataProcessingTime,
                 metrics.operatorStatistics.controlProcessingTime,
-                metrics.operatorStatistics.idleTime
+                metrics.operatorStatistics.idleTime,
+                reusedFromCache = metrics.reusedFromCache
               )
               (x._1, res)
           })
@@ -237,23 +238,7 @@ class ExecutionStatsService(
     val updatedLastMetrics = lastPersistedMetrics ++ newKeys.map(_ -> 
defaultMetrics)
 
     // Combine new metrics with old metrics for keys that are no longer present
-    val completeMetricsMap = newMetrics ++ oldKeys.map(key => key -> 
updatedLastMetrics(key))
-
-    // Transform the complete metrics map to ensure consistent structure
-    completeMetricsMap.map {
-      case (key, metrics) =>
-        key -> OperatorMetrics(
-          metrics.operatorState,
-          OperatorStatistics(
-            metrics.operatorStatistics.inputMetrics,
-            metrics.operatorStatistics.outputMetrics,
-            metrics.operatorStatistics.numWorkers,
-            metrics.operatorStatistics.dataProcessingTime,
-            metrics.operatorStatistics.controlProcessingTime,
-            metrics.operatorStatistics.idleTime
-          )
-        )
-    }
+    newMetrics ++ oldKeys.map(key => key -> updatedLastMetrics(key))
   }
 
   private def storeRuntimeStatistics(
diff --git 
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala
 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala
index 237936fa06..3bdb6c2e86 100644
--- 
a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala
+++ 
b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala
@@ -190,11 +190,13 @@ class ExecutionUtilsSpec extends AnyFlatSpec {
       numWorkers: Int = 0,
       dataTime: Long = 0,
       controlTime: Long = 0,
-      idleTime: Long = 0
+      idleTime: Long = 0,
+      reused: Boolean = false
   ): OperatorMetrics =
     OperatorMetrics(
       state,
-      OperatorStatistics(input, output, numWorkers, dataTime, controlTime, 
idleTime)
+      OperatorStatistics(input, output, numWorkers, dataTime, controlTime, 
idleTime),
+      reusedFromCache = reused
     )
 
   "ExecutionUtils.aggregateMetrics" should "return UNINITIALIZED defaults when 
given no metrics" in {
@@ -337,4 +339,22 @@ class ExecutionUtilsSpec extends AnyFlatSpec {
     assert(result.operatorStatistics.numWorkers == 3)
     assert(result.operatorStatistics.dataProcessingTime == 12)
   }
+
+  // -- aggregateMetrics: reused-from-cache provenance ----------------------
+
+  it should "report reusedFromCache only when every physical operator is 
reused" in {
+    val reusedA = metricsWith(WorkflowAggregatedState.COMPLETED, reused = true)
+    val reusedB = metricsWith(WorkflowAggregatedState.COMPLETED, reused = true)
+    val computed = metricsWith(WorkflowAggregatedState.COMPLETED)
+
+    assert(ExecutionUtils.aggregateMetrics(List(reusedA, 
reusedB)).reusedFromCache)
+    assert(!ExecutionUtils.aggregateMetrics(List(reusedA, 
computed)).reusedFromCache)
+    assert(!ExecutionUtils.aggregateMetrics(List(computed)).reusedFromCache)
+  }
+
+  it should "default reusedFromCache to false for empty input" in {
+    // Empty input takes the early-return path, whose default is false; metrics
+    // that no producer has marked keep that default too.
+    assert(!ExecutionUtils.aggregateMetrics(Iterable.empty).reusedFromCache)
+  }
 }
diff --git 
a/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala
 
b/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala
index b66d4d0c99..fc119b9f60 100644
--- 
a/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala
+++ 
b/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala
@@ -123,7 +123,10 @@ class TexeraWebSocketEventSpec extends AnyFlatSpec with 
Matchers {
     numWorkers = 17L,
     aggregatedDataProcessingTime = 18L,
     aggregatedControlProcessingTime = 19L,
-    aggregatedIdleTime = 20L
+    aggregatedIdleTime = 20L,
+    // Non-default on purpose: the symmetric round trip below only pins this
+    // field on the wire if a drop would change the value read back.
+    reusedFromCache = true
   )
 
   private val resultRow = objectMapper.createObjectNode().put("city", "Irvine")
diff --git a/frontend/src/app/workspace/types/execute-workflow.interface.ts 
b/frontend/src/app/workspace/types/execute-workflow.interface.ts
index 79f8aa8d2b..a4ce01ce24 100644
--- a/frontend/src/app/workspace/types/execute-workflow.interface.ts
+++ b/frontend/src/app/workspace/types/execute-workflow.interface.ts
@@ -80,6 +80,8 @@ export enum OperatorState {
 
 export interface OperatorStatistics
   extends Readonly<{
+    // Provenance: the operator completed by reusing cached results (no 
workers ran).
+    reusedFromCache?: boolean;
     aggregatedInputRowCount: number;
     aggregatedInputSize?: number;
     /** Absent when the snapshot has no per-port information at all; `{}` 
means every port was zero. */

Reply via email to