This is an automated email from the ASF dual-hosted git repository.
xuang7 pushed a commit to branch release/v1.2
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/release/v1.2 by this push:
new 5c53cf73f0 fix(amber, v1.2): store execution sizes as BIGINT to stop
silent >2GiB truncation (#7050)
5c53cf73f0 is described below
commit 5c53cf73f074a1b1985555f26e12a0b2e979e138
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Wed Jul 29 23:24:41 2026 -0700
fix(amber, v1.2): store execution sizes as BIGINT to stop silent >2GiB
truncation (#7050)
### What changes were proposed in this PR?
Automated backport of #6980 to `release/v1.2`.
Source: 791d7df8afcc6b1b9af590960098128d40828436 · [automation
run](https://github.com/apache/texera/actions/runs/30488298522)
### Any related issues, documentation, discussions?
Backport of #6980. Originally linked #6978.
### How was this PR tested?
Release-branch CI runs on this branch once the conflicts are resolved
and this PR is marked ready for review.
### Was this PR authored or co-authored using generative AI tooling?
No.
---------
Co-authored-by: Meng Wang <[email protected]>
Co-authored-by: Yicong Huang
<[email protected]>
Co-authored-by: Xinyuan Lin <[email protected]>
Co-authored-by: mengw15 <[email protected]>
---
.../dashboard/user/quota/UserQuotaResource.scala | 10 +-
.../user/workflow/WorkflowExecutionsResource.scala | 54 ++++++---
.../workflow/WorkflowExecutionsResourceSpec.scala | 130 +++++++++++++++++++++
sql/changelog.xml | 10 ++
sql/texera_ddl.sql | 6 +-
sql/updates/29.sql | 40 +++++++
6 files changed, 225 insertions(+), 25 deletions(-)
diff --git
a/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/quota/UserQuotaResource.scala
b/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/quota/UserQuotaResource.scala
index 6f1bb79b15..7f89d3fe1c 100644
---
a/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/quota/UserQuotaResource.scala
+++
b/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/quota/UserQuotaResource.scala
@@ -149,10 +149,7 @@ object UserQuotaResource {
.where(OPERATOR_PORT_EXECUTIONS.WORKFLOW_EXECUTION_ID.eq(eid))
.fetch()
.asScala
- .map(r =>
-
Option(r.get(OPERATOR_PORT_EXECUTIONS.RESULT_SIZE)).getOrElse(0).asInstanceOf[Integer]
- )
- .map(_.toLong)
+ .map(r =>
Option(r.get(OPERATOR_PORT_EXECUTIONS.RESULT_SIZE)).map(_.toLong).getOrElse(0L))
.sum
val logSize = context
@@ -162,11 +159,8 @@ object UserQuotaResource {
.fetch()
.asScala
.map(r =>
- Option(r.get(OPERATOR_EXECUTIONS.CONSOLE_MESSAGES_SIZE))
- .getOrElse(0)
- .asInstanceOf[Integer]
+
Option(r.get(OPERATOR_EXECUTIONS.CONSOLE_MESSAGES_SIZE)).map(_.toLong).getOrElse(0L)
)
- .map(_.toLong)
.sum
QuotaStorage(
diff --git
a/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResource.scala
b/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResource.scala
index 72fb1c364e..8bd5eb504d 100644
---
a/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResource.scala
+++
b/amber/src/main/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResource.scala
@@ -420,7 +420,7 @@ object WorkflowExecutionsResource {
): Unit = {
context
.update(OPERATOR_PORT_EXECUTIONS)
- .set(OPERATOR_PORT_EXECUTIONS.RESULT_SIZE, Integer.valueOf(size.toInt))
+ .set(OPERATOR_PORT_EXECUTIONS.RESULT_SIZE, java.lang.Long.valueOf(size))
.where(OPERATOR_PORT_EXECUTIONS.WORKFLOW_EXECUTION_ID.eq(eid.id.toInt))
.and(OPERATOR_PORT_EXECUTIONS.GLOBAL_PORT_ID.eq(globalPortId.serializeAsString))
.execute()
@@ -440,15 +440,27 @@ object WorkflowExecutionsResource {
.map(URI.create)
if (statsUriOpt.isPresent) {
- val size =
DocumentFactory.openDocument(statsUriOpt.get)._1.getTotalFileSize
- context
- .update(WORKFLOW_EXECUTIONS)
- .set(WORKFLOW_EXECUTIONS.RUNTIME_STATS_SIZE,
Integer.valueOf(size.toInt))
- .where(WORKFLOW_EXECUTIONS.EID.eq(eid.id.toInt))
- .execute()
+ updateRuntimeStatsSize(
+ eid,
+ DocumentFactory.openDocument(statsUriOpt.get)._1.getTotalFileSize
+ )
}
}
+ /**
+ * Stores an already-measured runtime statistics size, mirroring
updateResultSize.
+ *
+ * @param eid Execution ID associated with the runtime statistics document.
+ * @param size Size of the runtime statistics in bytes.
+ */
+ def updateRuntimeStatsSize(eid: ExecutionIdentity, size: Long): Unit = {
+ context
+ .update(WORKFLOW_EXECUTIONS)
+ .set(WORKFLOW_EXECUTIONS.RUNTIME_STATS_SIZE,
java.lang.Long.valueOf(size))
+ .where(WORKFLOW_EXECUTIONS.EID.eq(eid.id.toInt))
+ .execute()
+ }
+
/**
* Updates the size of the console message stored via Iceberg document.
*
@@ -465,16 +477,30 @@ object WorkflowExecutionsResource {
.map(URI.create)
if (uriOpt.isPresent) {
- val size = DocumentFactory.openDocument(uriOpt.get)._1.getTotalFileSize
- context
- .update(OPERATOR_EXECUTIONS)
- .set(OPERATOR_EXECUTIONS.CONSOLE_MESSAGES_SIZE,
Integer.valueOf(size.toInt))
- .where(OPERATOR_EXECUTIONS.WORKFLOW_EXECUTION_ID.eq(eid.id.toInt))
- .and(OPERATOR_EXECUTIONS.OPERATOR_ID.eq(opId.id))
- .execute()
+ updateConsoleMessageSize(
+ eid,
+ opId,
+ DocumentFactory.openDocument(uriOpt.get)._1.getTotalFileSize
+ )
}
}
+ /**
+ * Stores an already-measured console message size, mirroring
updateResultSize.
+ *
+ * @param eid Execution ID associated with the console message.
+ * @param opId Operator ID of the corresponding operator.
+ * @param size Size of the console messages in bytes.
+ */
+ def updateConsoleMessageSize(eid: ExecutionIdentity, opId: OperatorIdentity,
size: Long): Unit = {
+ context
+ .update(OPERATOR_EXECUTIONS)
+ .set(OPERATOR_EXECUTIONS.CONSOLE_MESSAGES_SIZE,
java.lang.Long.valueOf(size))
+ .where(OPERATOR_EXECUTIONS.WORKFLOW_EXECUTION_ID.eq(eid.id.toInt))
+ .and(OPERATOR_EXECUTIONS.OPERATOR_ID.eq(opId.id))
+ .execute()
+ }
+
/**
* This method is mainly used for frontend requests. Given a logicalOpId
and an outputPortId of an execution,
* this method finds the URI for a globalPortId that both: 1. matches the
logicalOpId and outputPortId, and
diff --git
a/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala
b/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala
index 03adc82c96..75efa4da6a 100644
---
a/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala
+++
b/amber/src/test/scala/org/apache/texera/web/resource/dashboard/user/workflow/WorkflowExecutionsResourceSpec.scala
@@ -592,6 +592,136 @@ class WorkflowExecutionsResourceSpec
assert(row.getResultSize == 4096)
}
+ it should "store a >2GiB size without truncation (#6978)" in {
+ val execution = insertExecution()
+ val eid = ExecutionIdentity(execution.getEid.longValue())
+ val globalPortId = GlobalPortIdentity(
+ PhysicalOpIdentity(OperatorIdentity("op-big-size"), "main"),
+ PortIdentity(),
+ input = false
+ )
+ WorkflowExecutionsResource.insertOperatorPortResultUri(
+ eid,
+ globalPortId,
+ URI.create("vfs:///big")
+ )
+
+ // 3 GiB exceeds Int.MaxValue; a Long->Int narrowing would wrap it
negative.
+ val threeGiB = 3L * 1024 * 1024 * 1024
+ WorkflowExecutionsResource.updateResultSize(eid, globalPortId, threeGiB)
+
+ val row = getDSLContext
+ .selectFrom(OPERATOR_PORT_EXECUTIONS)
+
.where(OPERATOR_PORT_EXECUTIONS.WORKFLOW_EXECUTION_ID.eq(execution.getEid))
+
.and(OPERATOR_PORT_EXECUTIONS.GLOBAL_PORT_ID.eq(globalPortId.serializeAsString))
+ .fetchOne()
+ assert(row.getResultSize.longValue() == threeGiB)
+ }
+
+ // ─── new: updateRuntimeStatsSize / updateConsoleMessageSize ───────────────
+
+ "updateRuntimeStatsSize" should "store a >2GiB size on the matching
execution" in {
+ val execution = insertExecution()
+ val eid = ExecutionIdentity(execution.getEid.longValue())
+ val threeGiB = 3L * 1024 * 1024 * 1024
+
+ WorkflowExecutionsResource.updateRuntimeStatsSize(eid, threeGiB)
+
+ val row = getDSLContext
+ .selectFrom(WORKFLOW_EXECUTIONS)
+ .where(WORKFLOW_EXECUTIONS.EID.eq(execution.getEid))
+ .fetchOne()
+ assert(row.getRuntimeStatsSize.longValue() == threeGiB)
+ }
+
+ it should "leave the size untouched when the execution has no runtime stats
URI" in {
+ val execution = insertExecution(runtimeStatsUri = null)
+
+ WorkflowExecutionsResource.updateRuntimeStatsSize(
+ ExecutionIdentity(execution.getEid.longValue())
+ )
+
+ val row = getDSLContext
+ .selectFrom(WORKFLOW_EXECUTIONS)
+ .where(WORKFLOW_EXECUTIONS.EID.eq(execution.getEid))
+ .fetchOne()
+ // The fixture never set a size, so a no-op leaves the column as inserted.
+ assert(row.getRuntimeStatsSize == null)
+ }
+
+ it should "open the stored document for measuring when a runtime stats URI
is present" in {
+ // A URI is present, so the method must reach the document-open call. No
+ // document backend exists in this unit environment, so the open fails on
+ // the unsupported scheme — proving the branch executed and that the
+ // failure propagates instead of degrading into a silent no-op.
+ val execution = insertExecution(runtimeStatsUri = "mock:///runtime-stats")
+
+ val ex = intercept[UnsupportedOperationException] {
+ WorkflowExecutionsResource.updateRuntimeStatsSize(
+ ExecutionIdentity(execution.getEid.longValue())
+ )
+ }
+ assert(ex.getMessage.contains("mock"))
+ }
+
+ "updateConsoleMessageSize" should "store a >2GiB size on the matching (eid,
opId) row" in {
+ val execution = insertExecution()
+ val eid = ExecutionIdentity(execution.getEid.longValue())
+ val opId = OperatorIdentity("op-console-size")
+ WorkflowExecutionsResource.insertOperatorExecutions(
+ execution.getEid.longValue(),
+ opId.id,
+ URI.create("vfs:///console-big")
+ )
+
+ val threeGiB = 3L * 1024 * 1024 * 1024
+ WorkflowExecutionsResource.updateConsoleMessageSize(eid, opId, threeGiB)
+
+ val row = getDSLContext
+ .selectFrom(OPERATOR_EXECUTIONS)
+ .where(OPERATOR_EXECUTIONS.WORKFLOW_EXECUTION_ID.eq(execution.getEid))
+ .and(OPERATOR_EXECUTIONS.OPERATOR_ID.eq(opId.id))
+ .fetchOne()
+ assert(row.getConsoleMessagesSize.longValue() == threeGiB)
+ }
+
+ it should "leave the size untouched when the operator has no console
messages URI" in {
+ val execution = insertExecution()
+ val opId = OperatorIdentity("op-no-console-uri")
+
+ WorkflowExecutionsResource.updateConsoleMessageSize(
+ ExecutionIdentity(execution.getEid.longValue()),
+ opId
+ )
+
+ val row = getDSLContext
+ .selectFrom(OPERATOR_EXECUTIONS)
+ .where(OPERATOR_EXECUTIONS.WORKFLOW_EXECUTION_ID.eq(execution.getEid))
+ .and(OPERATOR_EXECUTIONS.OPERATOR_ID.eq(opId.id))
+ .fetchOne()
+ assert(row == null)
+ }
+
+ it should "open the stored document for measuring when a console messages
URI is present" in {
+ // Same shape as the runtime-stats case above: the stored URI forces the
+ // document-open call, whose unsupported-scheme failure propagates.
+ val execution = insertExecution()
+ val opId = OperatorIdentity("op-console-uri")
+ WorkflowExecutionsResource.insertOperatorExecutions(
+ execution.getEid.longValue(),
+ opId.id,
+ URI.create("mock:///console")
+ )
+
+ val ex = intercept[UnsupportedOperationException] {
+ WorkflowExecutionsResource.updateConsoleMessageSize(
+ ExecutionIdentity(execution.getEid.longValue()),
+ opId
+ )
+ }
+ assert(ex.getMessage.contains("mock"))
+ }
+
// ─── new: getResultUriByLogicalPortId ─────────────────────────────────────
"getResultUriByLogicalPortId" should "match by logical operator id, port id,
and resource type" in {
diff --git a/sql/changelog.xml b/sql/changelog.xml
index 7e6bd1606a..8c181882a2 100644
--- a/sql/changelog.xml
+++ b/sql/changelog.xml
@@ -24,6 +24,16 @@
xsi:schemaLocation="http://www.liquibase.org/xml/ns/dbchangelog
http://www.liquibase.org/xml/ns/dbchangelog/dbchangelog-4.9.xsd">
+ <!-- Widen execution size columns to BIGINT (#6978).
+ Kept identical to main's entry: Liquibase identifies an applied change
+ by its (id, author, changelog path) tuple, and all three match main's
+ changeSet 29 here, so a deployment upgrading from v1.2 to main sees it
+ as already run instead of re-executing it. Ids 23-28 are main-only
+ changes that do not exist on this branch. -->
+ <changeSet id="29" author="mengw15">
+ <sqlFile path="sql/updates/29.sql"/>
+ </changeSet>
+
<!-- example changeSet
<changeSet id="1" author="author">
<sqlFile path="sql/updates/1.sql"/>
diff --git a/sql/texera_ddl.sql b/sql/texera_ddl.sql
index d6b488e582..cf76ea1286 100644
--- a/sql/texera_ddl.sql
+++ b/sql/texera_ddl.sql
@@ -227,7 +227,7 @@ CREATE TABLE IF NOT EXISTS workflow_executions
environment_version VARCHAR(128) NOT NULL,
log_location TEXT,
runtime_stats_uri TEXT,
- runtime_stats_size INT DEFAULT 0,
+ runtime_stats_size BIGINT DEFAULT 0,
FOREIGN KEY (vid) REFERENCES workflow_version(vid) ON DELETE CASCADE,
FOREIGN KEY (uid) REFERENCES "user"(uid) ON DELETE CASCADE,
FOREIGN KEY (cuid) REFERENCES workflow_computing_unit(cuid) ON DELETE
CASCADE
@@ -331,7 +331,7 @@ CREATE TABLE IF NOT EXISTS operator_executions
workflow_execution_id INT NOT NULL,
operator_id VARCHAR(100) NOT NULL,
console_messages_uri TEXT,
- console_messages_size INT DEFAULT 0,
+ console_messages_size BIGINT DEFAULT 0,
PRIMARY KEY (workflow_execution_id, operator_id),
FOREIGN KEY (workflow_execution_id) REFERENCES workflow_executions(eid) ON
DELETE CASCADE
);
@@ -342,7 +342,7 @@ CREATE TABLE operator_port_executions
workflow_execution_id INT NOT NULL,
global_port_id VARCHAR(200) NOT NULL,
result_uri TEXT,
- result_size INT DEFAULT 0,
+ result_size BIGINT DEFAULT 0,
PRIMARY KEY (workflow_execution_id, global_port_id),
FOREIGN KEY (workflow_execution_id) REFERENCES workflow_executions(eid) ON
DELETE CASCADE
);
diff --git a/sql/updates/29.sql b/sql/updates/29.sql
new file mode 100644
index 0000000000..8870179f00
--- /dev/null
+++ b/sql/updates/29.sql
@@ -0,0 +1,40 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+\c texera_db
+
+SET search_path TO texera_db;
+
+BEGIN;
+
+-- Execution size columns hold byte counts that can exceed 2 GiB (see #6978):
+-- they were INT, so Long sizes were narrowed via Scala's silent `.toInt`,
+-- wrapping sizes in [2 GiB, 4 GiB) to negative values and corrupting the
+-- user-quota sums computed over these columns. Widen them to BIGINT to match
+-- the Long byte counts the application writes.
+ALTER TABLE workflow_executions
+ ALTER COLUMN runtime_stats_size TYPE BIGINT;
+
+ALTER TABLE operator_executions
+ ALTER COLUMN console_messages_size TYPE BIGINT;
+
+ALTER TABLE operator_port_executions
+ ALTER COLUMN result_size TYPE BIGINT;
+
+COMMIT;