summaryrefslogtreecommitdiff
path: root/src/mongo/db/pipeline/document_source_change_stream_check_resumability.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/pipeline/document_source_change_stream_check_resumability.cpp')
-rw-r--r--src/mongo/db/pipeline/document_source_change_stream_check_resumability.cpp66
1 files changed, 26 insertions, 40 deletions
diff --git a/src/mongo/db/pipeline/document_source_change_stream_check_resumability.cpp b/src/mongo/db/pipeline/document_source_change_stream_check_resumability.cpp
index 01159d3dc2b..3861b21693d 100644
--- a/src/mongo/db/pipeline/document_source_change_stream_check_resumability.cpp
+++ b/src/mongo/db/pipeline/document_source_change_stream_check_resumability.cpp
@@ -30,7 +30,6 @@
#include "mongo/platform/basic.h"
#include "mongo/db/curop.h"
-#include "mongo/db/pipeline/change_stream_helpers.h"
#include "mongo/db/pipeline/document_source_change_stream_check_resumability.h"
#include "mongo/db/query/query_feature_flags_gen.h"
#include "mongo/db/repl/oplog_entry.h"
@@ -48,15 +47,16 @@ REGISTER_INTERNAL_DOCUMENT_SOURCE(_internalChangeStreamCheckResumability,
// Returns ResumeStatus::kFoundToken if the document retrieved from the resumed pipeline satisfies
// the client's resume token, ResumeStatus::kCheckNextDoc if it is older than the client's token,
-// and ResumeToken::kSurpassedToken if it is more recent than the client's resume token, indicating
-// that we will never see the token. Return ResumeStatus::kNeedsSplit if we have found the event
-// that produced the resume token, but it was split in the original stream.
+// and ResumeToken::kSurpassedToken if it is more recent than the client's resume token (indicating
+// that we will never see the token).
DocumentSourceChangeStreamCheckResumability::ResumeStatus
DocumentSourceChangeStreamCheckResumability::compareAgainstClientResumeToken(
- const Document& eventFromResumedStream, const ResumeTokenData& tokenDataFromClient) {
+ const intrusive_ptr<ExpressionContext>& expCtx,
+ const Document& documentFromResumedStream,
+ const ResumeTokenData& tokenDataFromClient) {
// Parse the stream doc into comprehensible ResumeTokenData.
auto tokenDataFromResumedStream =
- ResumeToken::parse(eventFromResumedStream.metadata().getSortKey().getDocument()).getData();
+ ResumeToken::parse(documentFromResumedStream["_id"].getDocument()).getData();
// We start the resume with a $gte query on the timestamp, so we never expect it to be lower
// than our resume token's timestamp.
@@ -97,25 +97,21 @@ DocumentSourceChangeStreamCheckResumability::compareAgainstClientResumeToken(
// clusterTime. If the stream UUID sorts after the client's, however, then the stream is not
// resumable; we are past the point in the stream where the token should have appeared.
if (tokenDataFromResumedStream.uuid != tokenDataFromClient.uuid) {
+ // If we are running on a replica set deployment, we don't ever expect to see identical time
+ // stamps and txnOpIndex but differing UUIDs, and we reject the resume attempt at once.
+ if (!expCtx->inMongos && !expCtx->needsMerge) {
+ return ResumeStatus::kSurpassedToken;
+ }
+ // Otherwise, return a ResumeStatus based on the sort-order of the client and stream UUIDs.
return tokenDataFromResumedStream.uuid > tokenDataFromClient.uuid
? ResumeStatus::kSurpassedToken
: ResumeStatus::kCheckNextDoc;
}
- // If the eventIdentifier matches exactly, then we have found the resume point. However, this
- // event may have been split by the original stream; we must check the value of the resume
- // token's fragmentNum field to determine the correct return status.
+ // If all the fields match exactly, then we have found the token.
if (ValueComparator::kInstance.evaluate(tokenDataFromResumedStream.eventIdentifier ==
tokenDataFromClient.eventIdentifier)) {
- if (tokenDataFromClient.fragmentNum && !tokenDataFromResumedStream.fragmentNum) {
- return ResumeStatus::kNeedsSplit;
- }
- if (tokenDataFromResumedStream.fragmentNum == tokenDataFromClient.fragmentNum) {
- return ResumeStatus::kFoundToken;
- }
- return tokenDataFromResumedStream.fragmentNum > tokenDataFromClient.fragmentNum
- ? ResumeStatus::kSurpassedToken
- : ResumeStatus::kCheckNextDoc;
+ return ResumeStatus::kFoundToken;
}
// At this point, we know that the tokens differ only by eventIdentifier. The status we return
@@ -134,7 +130,7 @@ DocumentSourceChangeStreamCheckResumability::DocumentSourceChangeStreamCheckResu
intrusive_ptr<DocumentSourceChangeStreamCheckResumability>
DocumentSourceChangeStreamCheckResumability::create(const intrusive_ptr<ExpressionContext>& expCtx,
const DocumentSourceChangeStreamSpec& spec) {
- auto resumeToken = change_stream::resolveResumeTokenFromSpec(expCtx, spec);
+ auto resumeToken = DocumentSourceChangeStream::resolveResumeTokenFromSpec(expCtx, spec);
return new DocumentSourceChangeStreamCheckResumability(expCtx, std::move(resumeToken));
}
@@ -182,21 +178,15 @@ DocumentSource::GetNextResult DocumentSourceChangeStreamCheckResumability::doGet
// Determine whether the current event sorts before, equal to or after the resume token.
_resumeStatus =
DocumentSourceChangeStreamCheckResumability::compareAgainstClientResumeToken(
- nextInput.getDocument(), _tokenFromClient);
+ pExpCtx, nextInput.getDocument(), _tokenFromClient);
switch (_resumeStatus) {
case ResumeStatus::kCheckNextDoc:
// If the result was kCheckNextDoc, we are resumable but must swallow this event.
continue;
- case ResumeStatus::kNeedsSplit:
- // If the result was kNeedsSplit, we found a resume token which matches the client's
- // except for the splitNum attribute. Allow this document to pass through so that
- // the split stage can regenerate the original fragments and their resume tokens.
- return nextInput;
case ResumeStatus::kSurpassedToken:
// In this case the resume token wasn't found; it may be on another shard. However,
// since the oplog scan did not throw, we know that we are resumable. Fall through
// into the following case and return the document.
- return nextInput;
case ResumeStatus::kFoundToken:
// We found the actual token! Return the doc so DSEnsureResumeTokenPresent sees it.
return nextInput;
@@ -206,20 +196,16 @@ DocumentSource::GetNextResult DocumentSourceChangeStreamCheckResumability::doGet
}
Value DocumentSourceChangeStreamCheckResumability::serialize(
- const SerializationOptions& opts) const {
- BSONObjBuilder builder;
- if (opts.verbosity) {
- BSONObjBuilder sub(builder.subobjStart(DocumentSourceChangeStream::kStageName));
- sub.append("stage"_sd, kStageName);
- sub << "resumeToken"_sd << Value(ResumeToken(_tokenFromClient).toDocument(opts));
- sub.done();
- } else {
- builder.append(
- kStageName,
- DocumentSourceChangeStreamCheckResumabilitySpec(ResumeToken(_tokenFromClient))
- .toBSON(opts));
- }
- return Value(builder.obj());
+ boost::optional<ExplainOptions::Verbosity> explain) const {
+ return explain
+ ? Value(DOC(DocumentSourceChangeStream::kStageName
+ << DOC("stage"
+ << "internalCheckResumability"_sd
+ << "resumeToken" << ResumeToken(_tokenFromClient).toDocument())))
+ : Value(Document{
+ {DocumentSourceChangeStreamCheckResumability::kStageName,
+ DocumentSourceChangeStreamCheckResumabilitySpec(ResumeToken(_tokenFromClient))
+ .toBSON()}});
}
} // namespace mongo