diff --git a/driver-core/src/main/com/mongodb/ReadConcern.java b/driver-core/src/main/com/mongodb/ReadConcern.java index 395e9255a9..2cb4d1e5e8 100644 --- a/driver-core/src/main/com/mongodb/ReadConcern.java +++ b/driver-core/src/main/com/mongodb/ReadConcern.java @@ -19,6 +19,10 @@ import com.mongodb.lang.Nullable; import org.bson.BsonDocument; import org.bson.BsonString; +import org.bson.BsonTimestamp; +import org.bson.BsonValue; + +import java.util.Objects; import static com.mongodb.assertions.Assertions.notNull; @@ -31,6 +35,7 @@ */ public final class ReadConcern { private final ReadConcernLevel level; + private final BsonTimestamp afterClusterTime; /** * Construct a new read concern @@ -39,6 +44,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 +113,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 +126,37 @@ 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; + } + + public BsonTimestamp getAfterClusterTime() { + return afterClusterTime; } } diff --git a/driver-core/src/main/com/mongodb/internal/connection/ReadConcernHelper.java b/driver-core/src/main/com/mongodb/internal/connection/ReadConcernHelper.java index df1e1ed754..72b0796780 100644 --- a/driver-core/src/main/com/mongodb/internal/connection/ReadConcernHelper.java +++ b/driver-core/src/main/com/mongodb/internal/connection/ReadConcernHelper.java @@ -45,7 +45,10 @@ public static BsonDocument getReadConcernDocument(final SessionContext sessionCo if (sessionContext.isSnapshot() && maxWireVersion < FIVE_DOT_ZERO_WIRE_VERSION) { throw new MongoClientException("Snapshot reads require MongoDB 5.0 or later"); } - if (shouldAddAfterClusterTime(sessionContext)) { + if (sessionContext.getReadConcern().getAfterClusterTime() != null) { + readConcernDocument.append( + "afterClusterTime", sessionContext.getReadConcern().getAfterClusterTime()); + } else if (shouldAddAfterClusterTime(sessionContext)) { readConcernDocument.append("afterClusterTime", sessionContext.getOperationTime()); } else if (shouldAddAtClusterTime(sessionContext)) { readConcernDocument.append("atClusterTime", sessionContext.getSnapshotTimestamp()); @@ -61,6 +64,5 @@ private static boolean shouldAddAfterClusterTime(final SessionContext sessionCon return sessionContext.isCausallyConsistent() && sessionContext.getOperationTime() != null; } - private ReadConcernHelper() { - } + private ReadConcernHelper() {} } 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 a838e7e03b..fe1e24d647 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,11 @@ public ReadConcern getReadConcern() { return assertNotNull(clientSession.getTransactionOptions().getReadConcern()); } else if (isSnapshot()) { return ReadConcern.SNAPSHOT; + } else if (inheritedReadConcern == null + && !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 6c6f8b6ced..6004f99f9a 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,11 @@ public ReadConcern getReadConcern() { return assertNotNull(clientSession.getTransactionOptions().getReadConcern()); } else if (isSnapshot()) { return ReadConcern.SNAPSHOT; + } else if (inheritedReadConcern == null + && !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 67cb82f365..1861d4a0a0 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