Skip to content

Commit 102bfc9

Browse files
authored
Add server timestamps to version mark resolution APIs (microsoft#28037)
1 parent e4c4fea commit 102bfc9

10 files changed

Lines changed: 330 additions & 118 deletions

File tree

.changeset/lazy-teams-tan.md

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
---
2+
"@fluidframework/container-runtime": minor
3+
"__section": feature
4+
---
5+
6+
Add server timestamps to beta version mark resolution APIs
7+
8+
The `@beta` `IVersionMarkResolver` APIs now expose the server timestamp of a resolved mark while remaining compatible with callers using the previous result and listener shapes:
9+
10+
- The resolved variants returned by `sealAndCaptureVersionMark()` and `resolve()` populate an optional `timestamp` property.
11+
- The `onBatchSequenced()` listener can receive `timestamp` as its optional third argument when an incoming batch resolves a pending mark.

packages/runtime/container-runtime/api-report/container-runtime.legacy.alpha.api.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -340,7 +340,7 @@ export interface IUploadSummaryResult extends Omit<IGenerateSummaryTreeResult, "
340340

341341
// @beta @legacy
342342
export interface IVersionMarkResolver {
343-
onBatchSequenced(listener: (batchId: string, sequenceNumber: number) => void): () => void;
343+
onBatchSequenced(listener: (batchId: string, sequenceNumber: number, timestamp?: number) => void): () => void;
344344
resolve(batchId: string, sequenceNumberLowerBound: number): Promise<ResolveResult>;
345345
sealAndCaptureVersionMark(): VersionMarkCapture;
346346
}
@@ -384,6 +384,7 @@ export type ReadFluidDataStoreAttributes = IFluidDataStoreAttributes0 | IFluidDa
384384
export type ResolveResult = {
385385
readonly kind: "resolved";
386386
readonly sequenceNumber: number;
387+
readonly timestamp?: number;
387388
} | {
388389
readonly kind: "pending";
389390
} | {
@@ -448,6 +449,7 @@ export type VersionMarkCapture = {
448449
} | {
449450
readonly kind: "resolved";
450451
readonly sequenceNumber: number;
452+
readonly timestamp?: number;
451453
};
452454

453455
// (No @packageDocumentation comment for this package)

packages/runtime/container-runtime/api-report/container-runtime.legacy.beta.api.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -340,7 +340,7 @@ export interface IUploadSummaryResult extends Omit<IGenerateSummaryTreeResult, "
340340

341341
// @beta @legacy
342342
export interface IVersionMarkResolver {
343-
onBatchSequenced(listener: (batchId: string, sequenceNumber: number) => void): () => void;
343+
onBatchSequenced(listener: (batchId: string, sequenceNumber: number, timestamp?: number) => void): () => void;
344344
resolve(batchId: string, sequenceNumberLowerBound: number): Promise<ResolveResult>;
345345
sealAndCaptureVersionMark(): VersionMarkCapture;
346346
}
@@ -379,6 +379,7 @@ export type ReadFluidDataStoreAttributes = IFluidDataStoreAttributes0 | IFluidDa
379379
export type ResolveResult = {
380380
readonly kind: "resolved";
381381
readonly sequenceNumber: number;
382+
readonly timestamp?: number;
382383
} | {
383384
readonly kind: "pending";
384385
} | {
@@ -443,6 +444,7 @@ export type VersionMarkCapture = {
443444
} | {
444445
readonly kind: "resolved";
445446
readonly sequenceNumber: number;
447+
readonly timestamp?: number;
446448
};
447449

448450
// (No @packageDocumentation comment for this package)

packages/runtime/container-runtime/src/containerRuntime.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1943,6 +1943,7 @@ export class ContainerRuntime
19431943
const fetchOps = (context as IContainerContextInternal).fetchOps;
19441944
this.versionMarkResolverInternal = new VersionMarkResolver({
19451945
getCurrentSequenceNumber: () => this.deltaManager.lastSequenceNumber,
1946+
getCurrentTimestamp: () => this.getCurrentReferenceTimestampMs(),
19461947
getCurrentMinimumSequenceNumber: () => this.deltaManager.minimumSequenceNumber,
19471948
getCurrentPendingBatchId: () => this.pendingStateManager.getMostRecentPendingBatchId(),
19481949
logger: createChildLogger({
@@ -3422,6 +3423,7 @@ export class ContainerRuntime
34223423
this.versionMarkResolverInternal.processInboundBatch(
34233424
versionMarkUpdate.sequenced.batchId,
34243425
versionMarkUpdate.sequenced.sequenceNumber,
3426+
versionMarkUpdate.sequenced.timestamp,
34253427
);
34263428
}
34273429
this.versionMarkInboundBatchId = versionMarkUpdate.carriedBatchId;

packages/runtime/container-runtime/src/test/containerRuntime.spec.ts

Lines changed: 107 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -104,7 +104,9 @@ import { pkgVersion } from "../packageVersion.js";
104104
import type { IPendingMessage, PendingStateManager } from "../pendingStateManager.js";
105105
import {
106106
type ISummaryCancellationToken,
107+
type IContainerRuntimeMetadata,
107108
neverCancelledSummaryToken,
109+
metadataBlobName,
108110
recentBatchInfoBlobName,
109111
type IRefreshSummaryAckOptions,
110112
} from "../summary/index.js";
@@ -2886,16 +2888,117 @@ describe("Runtime", () => {
28862888
"Expected DuplicateBatch telemetry to still be logged",
28872889
);
28882890
});
2891+
2892+
it("Version mark tracking keeps the earlier resolution for a tolerated duplicate batch with no explicit batchId", async () => {
2893+
const logger = new MockLogger();
2894+
const { runtime: containerRuntime } = await ContainerRuntime.loadRuntime2({
2895+
context: getMockContext({ logger }) as IContainerContext,
2896+
registry: new FluidDataStoreRegistry([]),
2897+
existing: false,
2898+
runtimeOptions: {
2899+
enableRuntimeIdCompressor: "on",
2900+
},
2901+
provideEntryPoint: mockProvideEntryPoint,
2902+
});
2903+
2904+
// Subscribing activates version-mark tracking on the inbound batch path.
2905+
const promotions: { sequenceNumber: number; timestamp: number | undefined }[] = [];
2906+
containerRuntime.versionMarkResolver.onBatchSequenced(
2907+
(_batchId, sequenceNumber, timestamp) => {
2908+
promotions.push({ sequenceNumber, timestamp });
2909+
},
2910+
);
2911+
2912+
const makeMessage = (
2913+
sequenceNumber: number,
2914+
timestamp: number,
2915+
): ISequencedDocumentMessage =>
2916+
({
2917+
clientId: "clientId",
2918+
clientSequenceNumber: 1,
2919+
sequenceNumber,
2920+
timestamp,
2921+
type: MessageType.Operation,
2922+
contents: { type: ContainerMessageType.Rejoin, contents: undefined },
2923+
}) satisfies Partial<ISequencedDocumentMessage> as ISequencedDocumentMessage;
2924+
2925+
containerRuntime.process(makeMessage(123, 1000), false);
2926+
2927+
// The service re-broadcasts the same batch under a different sequence number/timestamp.
2928+
// DuplicateBatchDetector tolerates this (no explicit batchId), so it must not throw, and
2929+
// VersionMarkResolver must keep the earlier (first-landed) resolution rather than remapping
2930+
// to the duplicate's sequence number.
2931+
assert.doesNotThrow(
2932+
() => containerRuntime.process(makeMessage(234, 2000), false),
2933+
"Should not throw when the duplicate batch has no explicit batchId",
2934+
);
2935+
logger.assertMatchAny(
2936+
[{ eventName: "ContainerRuntime:DuplicateBatch" }],
2937+
"Expected DuplicateBatch telemetry to be logged",
2938+
);
2939+
assert.deepStrictEqual(
2940+
promotions,
2941+
[{ sequenceNumber: 123, timestamp: 1000 }],
2942+
"Only the first (original) resolution should have been promoted to listeners",
2943+
);
2944+
});
28892945
});
28902946

28912947
describe("Version mark inbound update", () => {
2948+
it("captures the summary timestamp after load before any new ops arrive", async () => {
2949+
const sequenceNumber = 42;
2950+
const timestamp = 123456;
2951+
const metadata: IContainerRuntimeMetadata = {
2952+
summaryFormatVersion: 1,
2953+
message: {
2954+
clientId: mockClientId,
2955+
clientSequenceNumber: 1,
2956+
minimumSequenceNumber: 0,
2957+
referenceSequenceNumber: 41,
2958+
sequenceNumber,
2959+
timestamp,
2960+
type: MessageType.Operation,
2961+
},
2962+
};
2963+
const context = getMockContext({
2964+
baseSnapshot: {
2965+
trees: { ".channels": { trees: {}, blobs: {} } },
2966+
blobs: { [metadataBlobName]: "metadata-id" },
2967+
},
2968+
mockStorage: {
2969+
...defaultMockStorage,
2970+
readBlob: async (id) => {
2971+
assert.equal(id, "metadata-id");
2972+
return stringToBuffer(JSON.stringify(metadata), "utf8");
2973+
},
2974+
},
2975+
});
2976+
const deltaManager = context.deltaManager as MockDeltaManager;
2977+
deltaManager.initialSequenceNumber = sequenceNumber;
2978+
deltaManager.lastSequenceNumber = sequenceNumber;
2979+
deltaManager.lastMessage = undefined;
2980+
2981+
const { runtime: containerRuntime } = await ContainerRuntime.loadRuntime2({
2982+
context: context as IContainerContext,
2983+
registry: new FluidDataStoreRegistry([]),
2984+
existing: true,
2985+
provideEntryPoint: mockProvideEntryPoint,
2986+
});
2987+
2988+
assert.deepEqual(containerRuntime.versionMarkResolver.sealAndCaptureVersionMark(), {
2989+
kind: "resolved",
2990+
sequenceNumber,
2991+
timestamp,
2992+
});
2993+
});
2994+
28922995
it("resolves through context fetchOps and the historical unpack pipeline", async () => {
28932996
let capturedSignal: AbortSignal | undefined;
28942997
let readCount = 0;
28952998
const context = {
28962999
...getMockContext(),
28973000
fetchOps: async (from, to, abortSignal) => {
2898-
assert.equal(from, 11, "the history read starts at the inclusive lower bound");
3001+
assert.equal(from, 10, "the history read starts at the inclusive lower bound");
28993002
assert.equal(to, undefined, "the history read has no fixed upper bound");
29003003
capturedSignal = abortSignal;
29013004
return {
@@ -2919,6 +3022,7 @@ describe("Runtime", () => {
29193022
{
29203023
type: MessageType.Operation,
29213024
sequenceNumber: 13,
3025+
timestamp: 13000,
29223026
clientId: "targetClient",
29233027
clientSequenceNumber: 7,
29243028
contents: {
@@ -2943,8 +3047,8 @@ describe("Runtime", () => {
29433047
});
29443048

29453049
assert.deepEqual(
2946-
await containerRuntime.versionMarkResolver.resolve("targetBatch", 11),
2947-
{ kind: "resolved", sequenceNumber: 13 },
3050+
await containerRuntime.versionMarkResolver.resolve("targetBatch", 10),
3051+
{ kind: "resolved", sequenceNumber: 13, timestamp: 13000 },
29483052
);
29493053
assert.equal(readCount, 1, "the stream stops reading after the target batch");
29503054
assert.equal(capturedSignal?.aborted, true, "the fetch is aborted after the match");

packages/runtime/container-runtime/src/test/versionMarks/inboundBatch.spec.ts

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -18,17 +18,23 @@ import {
1818
// eslint-disable-next-line import-x/no-internal-modules -- test targets the inbound-batch helper directly
1919
} from "../../versionMarks/inboundBatch.js";
2020

21-
/** The helper only reads `sequenceNumber` off messages. */
21+
/** The helper reads the sequence number and server timestamp off messages. */
2222
function msg(sequenceNumber: number): InboundSequencedContainerRuntimeMessage {
23-
return { sequenceNumber } as unknown as InboundSequencedContainerRuntimeMessage;
23+
return {
24+
sequenceNumber,
25+
timestamp: sequenceNumber * 1000,
26+
} as unknown as InboundSequencedContainerRuntimeMessage;
2427
}
2528

2629
function batchStart(batchId: string | undefined, keyMessageSeq: number): BatchStartInfo {
2730
return {
2831
batchId,
2932
clientId: "client",
3033
batchStartCsn: 3,
31-
keyMessage: { sequenceNumber: keyMessageSeq } as unknown as ISequencedDocumentMessage,
34+
keyMessage: {
35+
sequenceNumber: keyMessageSeq,
36+
timestamp: keyMessageSeq * 1000,
37+
} as unknown as ISequencedDocumentMessage,
3238
};
3339
}
3440

@@ -58,7 +64,7 @@ describe("inboundVersionMarkUpdate", () => {
5864
"stale_[9]",
5965
);
6066
assert.deepEqual(result, {
61-
sequenced: { batchId: "b_[3]", sequenceNumber: 12 },
67+
sequenced: { batchId: "b_[3]", sequenceNumber: 12, timestamp: 12000 },
6268
carriedBatchId: undefined,
6369
});
6470
});
@@ -69,7 +75,7 @@ describe("inboundVersionMarkUpdate", () => {
6975
undefined,
7076
);
7177
assert.deepEqual(result, {
72-
sequenced: { batchId: "b_[3]", sequenceNumber: 42 },
78+
sequenced: { batchId: "b_[3]", sequenceNumber: 42, timestamp: 42000 },
7379
carriedBatchId: undefined,
7480
});
7581
});
@@ -80,7 +86,11 @@ describe("inboundVersionMarkUpdate", () => {
8086
undefined,
8187
);
8288
assert.deepEqual(result, {
83-
sequenced: { batchId: generateBatchId("client", 3), sequenceNumber: 7 },
89+
sequenced: {
90+
batchId: generateBatchId("client", 3),
91+
sequenceNumber: 7,
92+
timestamp: 7000,
93+
},
8494
carriedBatchId: undefined,
8595
});
8696
});
@@ -103,7 +113,7 @@ describe("inboundVersionMarkUpdate", () => {
103113
it("records the carried id at the batch's last op and clears the carry", () => {
104114
const result = inboundVersionMarkUpdate(nextBatchMessage(22, true), "b_[3]");
105115
assert.deepEqual(result, {
106-
sequenced: { batchId: "b_[3]", sequenceNumber: 22 },
116+
sequenced: { batchId: "b_[3]", sequenceNumber: 22, timestamp: 22000 },
107117
carriedBatchId: undefined,
108118
});
109119
});
@@ -123,7 +133,7 @@ describe("inboundVersionMarkUpdate", () => {
123133

124134
const end = inboundVersionMarkUpdate(nextBatchMessage(32, true), mid.carriedBatchId);
125135
assert.deepEqual(end, {
126-
sequenced: { batchId: "b_[3]", sequenceNumber: 32 },
136+
sequenced: { batchId: "b_[3]", sequenceNumber: 32, timestamp: 32000 },
127137
carriedBatchId: undefined,
128138
});
129139
});

0 commit comments

Comments
 (0)