From 7ef3c6f2fe2b60dfcab8217f7a2b3d55b0da3c62 Mon Sep 17 00:00:00 2001 From: Andre Masella Date: Fri, 4 Sep 2026 14:36:13 -0400 Subject: [PATCH] JAVA-6179: Send afterClusterTime on writes in causally-consistent sessions Send `readConcern.afterClusterTime` on write commands when they are run in a causally-consistent session. --- .../src/main/com/mongodb/ReadConcern.java | 33 ++++++++++++------- .../client/internal/ClientSessionBinding.java | 4 +++ .../client/internal/ClientSessionBinding.java | 4 +++ .../unified/UnifiedTestModifications.java | 8 ----- 4 files changed, 29 insertions(+), 20 deletions(-) diff --git a/driver-core/src/main/com/mongodb/ReadConcern.java b/driver-core/src/main/com/mongodb/ReadConcern.java index 395e9255a9d..cf086add1b1 100644 --- a/driver-core/src/main/com/mongodb/ReadConcern.java +++ b/driver-core/src/main/com/mongodb/ReadConcern.java @@ -19,6 +19,9 @@ import com.mongodb.lang.Nullable; import org.bson.BsonDocument; import org.bson.BsonString; +import org.bson.BsonTimestamp; + +import java.util.Objects; import static com.mongodb.assertions.Assertions.notNull; @@ -31,6 +34,7 @@ */ public final class ReadConcern { private final ReadConcernLevel level; + private final BsonTimestamp afterClusterTime; /** * Construct a new read concern @@ -39,6 +43,15 @@ public final class ReadConcern { */ public ReadConcern(final ReadConcernLevel level) { this.level = notNull("level", level); + this.afterClusterTime = null; + } + + /** + * Use an after-cluster-time read concern for causally-consistent sessions + */ + public ReadConcern(final BsonTimestamp afterClusterTime) { + this.afterClusterTime = notNull("afterClusterTime", afterClusterTime); + this.level = null; } /** @@ -99,7 +112,7 @@ public ReadConcernLevel getLevel() { * @return true if this is the server default read concern */ public boolean isServerDefault() { - return level == null; + return level == null && afterClusterTime == null; } /** @@ -112,37 +125,33 @@ public BsonDocument asDocument() { if (level != null) { readConcern.put("level", new BsonString(level.getValue())); } + if (afterClusterTime != null) { + readConcern.put("afterClusterTime", afterClusterTime); + } return readConcern; } @Override public boolean equals(final Object o) { - if (this == o) { - return true; - } if (o == null || getClass() != o.getClass()) { return false; } - ReadConcern that = (ReadConcern) o; - - return level == that.level; + return level == that.level && Objects.equals(afterClusterTime, that.afterClusterTime); } @Override public int hashCode() { - return level != null ? level.hashCode() : 0; + return Objects.hash(level, afterClusterTime); } - @Override public String toString() { - return "ReadConcern{" - + "level=" + level - + '}'; + return "ReadConcern{" + "level=" + level + ", afterClusterTime=" + afterClusterTime + '}'; } private ReadConcern() { this.level = null; + this.afterClusterTime = null; } } diff --git a/driver-reactive-streams/src/main/com/mongodb/reactivestreams/client/internal/ClientSessionBinding.java b/driver-reactive-streams/src/main/com/mongodb/reactivestreams/client/internal/ClientSessionBinding.java index a838e7e03bd..9aed64b8e88 100644 --- a/driver-reactive-streams/src/main/com/mongodb/reactivestreams/client/internal/ClientSessionBinding.java +++ b/driver-reactive-streams/src/main/com/mongodb/reactivestreams/client/internal/ClientSessionBinding.java @@ -241,6 +241,10 @@ public ReadConcern getReadConcern() { return assertNotNull(clientSession.getTransactionOptions().getReadConcern()); } else if (isSnapshot()) { return ReadConcern.SNAPSHOT; + } else if (!clientSession.getServerSession().isClosed() + && clientSession.isCausallyConsistent() + && clientSession.getOperationTime() != null) { + return new ReadConcern(clientSession.getOperationTime()); } else { return inheritedReadConcern; } diff --git a/driver-sync/src/main/com/mongodb/client/internal/ClientSessionBinding.java b/driver-sync/src/main/com/mongodb/client/internal/ClientSessionBinding.java index 6c6f8b6ced9..f559916aa05 100644 --- a/driver-sync/src/main/com/mongodb/client/internal/ClientSessionBinding.java +++ b/driver-sync/src/main/com/mongodb/client/internal/ClientSessionBinding.java @@ -222,6 +222,10 @@ public ReadConcern getReadConcern() { return assertNotNull(clientSession.getTransactionOptions().getReadConcern()); } else if (isSnapshot()) { return ReadConcern.SNAPSHOT; + } else if (!clientSession.getServerSession().isClosed() + && clientSession.isCausallyConsistent() + && clientSession.getOperationTime() != null) { + return new ReadConcern(clientSession.getOperationTime()); } else { return inheritedReadConcern; } diff --git a/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java b/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java index 67cb82f3656..1861d4a0a0b 100644 --- a/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java +++ b/driver-sync/src/test/functional/com/mongodb/client/unified/UnifiedTestModifications.java @@ -583,14 +583,6 @@ public static void applyCustomizations(final TestDef def) { .file("transactions", "backpressure-retryable-commit"); def.skipJira("https://jira.mongodb.org/browse/JAVA-5956 TODO-JAVA-5956") .file("transactions", "backpressure-retryable-abort"); - def.skipJira("https://jira.mongodb.org/browse/JAVA-6179") - .test("transactions", "retryable-writes", "increment txnNumber") - .test("transactions", "commit", "reset session state commit") - .test("transactions", "commit", "reset session state abort") - .test("transactions-convenient-api", "callback-commits", - "withTransaction still succeeds if callback commits and runs extra op") - .test("transactions-convenient-api", "callback-aborts", - "withTransaction still succeeds if callback aborts and runs extra op"); // valid-pass