diff --git a/.changeset/every-feet-post.md b/.changeset/every-feet-post.md new file mode 100644 index 000000000000..9953e737c118 --- /dev/null +++ b/.changeset/every-feet-post.md @@ -0,0 +1,9 @@ +--- +"@fluidframework/driver-definitions": minor +"@fluidframework/container-loader": minor +"@fluidframework/odsp-driver": minor +"__section": fix +--- +Preserve driver state in pending container state + +Pending container state now captures and restores opaque driver state before reconnecting. ODSP uses this to retain the cached epoch and validate document identity, allowing restored files to reject stale pending state after a server-side restore. diff --git a/packages/common/driver-definitions/api-report/driver-definitions.legacy.beta.api.md b/packages/common/driver-definitions/api-report/driver-definitions.legacy.beta.api.md index e20164ab4f7d..2335ba262475 100644 --- a/packages/common/driver-definitions/api-report/driver-definitions.legacy.beta.api.md +++ b/packages/common/driver-definitions/api-report/driver-definitions.legacy.beta.api.md @@ -267,6 +267,10 @@ export interface IDocumentService extends IEventProvider connectToDeltaStream(client: IClient): Promise; connectToStorage(): Promise; dispose(error?: any): void; + readonly driverStatePersistence?: { + get(): unknown; + set(state: unknown): void; + }; policies?: IDocumentServicePolicies | undefined; // (undocumented) resolvedUrl: IResolvedUrl; diff --git a/packages/common/driver-definitions/src/storage.ts b/packages/common/driver-definitions/src/storage.ts index 18809bcea913..590c54a6f09c 100644 --- a/packages/common/driver-definitions/src/storage.ts +++ b/packages/common/driver-definitions/src/storage.ts @@ -372,6 +372,29 @@ export interface IDocumentServicePolicies { export interface IDocumentService extends IEventProvider { resolvedUrl: IResolvedUrl; + /** + * Persists opaque driver state with pending container state. + * + * @remarks + * The value returned by `get` is serialized as part of the containing pending state. It must + * round-trip through `JSON.stringify` and `JSON.parse` without custom serialization or + * information loss, and it must not contain customer-identifying information. `get` returning + * `undefined` indicates that there is no driver state to preserve. `set` receives the + * JSON-deserialized value previously returned by `get` before the document service connects + * to storage or the delta stream. Persisted state must contain enough information for `set` to + * reject state captured for a different document. + * + * @privateRemarks + * Grouping `get` and `set` under one optional property makes support atomic: an implementation + * either omits persistence or provides the complete round-trip capability. Making the methods + * independently optional would permit state that can be captured but not restored, or restored + * but not captured. + */ + readonly driverStatePersistence?: { + get(): unknown; + set(state: unknown): void; + }; + /** * Policies implemented/instructed by driver. */ diff --git a/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.alpha.api.md b/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.alpha.api.md index 26506ee1658c..1b53c3a4a6c4 100644 --- a/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.alpha.api.md +++ b/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.alpha.api.md @@ -48,7 +48,8 @@ export class EpochTracker implements IPersistedFileCache { readonly rateLimiter: RateLimiter; // (undocumented) removeEntries(): Promise; - // (undocumented) + setEpoch(epoch: string, source: FetchTypeInternal | "pendingState"): void; + // @deprecated setEpoch(epoch: string, fromCache: boolean, fetchType: FetchTypeInternal): void; // (undocumented) validateEpoch(epoch: string | undefined, fetchType: FetchType): Promise; diff --git a/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.beta.api.md b/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.beta.api.md index 0e2c3a7039a4..95ef3cd002cf 100644 --- a/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.beta.api.md +++ b/packages/drivers/odsp-driver/api-report/odsp-driver.legacy.beta.api.md @@ -48,7 +48,8 @@ export class EpochTracker implements IPersistedFileCache { readonly rateLimiter: RateLimiter; // (undocumented) removeEntries(): Promise; - // (undocumented) + setEpoch(epoch: string, source: FetchTypeInternal | "pendingState"): void; + // @deprecated setEpoch(epoch: string, fromCache: boolean, fetchType: FetchTypeInternal): void; // (undocumented) validateEpoch(epoch: string | undefined, fetchType: FetchType): Promise; diff --git a/packages/drivers/odsp-driver/api-report/odsp-driver.point-in-time.legacy.beta.api.md b/packages/drivers/odsp-driver/api-report/odsp-driver.point-in-time.legacy.beta.api.md index 580a5c47b654..ba498ed70bbb 100644 --- a/packages/drivers/odsp-driver/api-report/odsp-driver.point-in-time.legacy.beta.api.md +++ b/packages/drivers/odsp-driver/api-report/odsp-driver.point-in-time.legacy.beta.api.md @@ -33,7 +33,8 @@ export class EpochTracker implements IPersistedFileCache { readonly rateLimiter: RateLimiter; // (undocumented) removeEntries(): Promise; - // (undocumented) + setEpoch(epoch: string, source: FetchTypeInternal | "pendingState"): void; + // @deprecated setEpoch(epoch: string, fromCache: boolean, fetchType: FetchTypeInternal): void; // (undocumented) validateEpoch(epoch: string | undefined, fetchType: FetchType): Promise; diff --git a/packages/drivers/odsp-driver/src/epochTracker.ts b/packages/drivers/odsp-driver/src/epochTracker.ts index 33d6ec6798d7..c6acbbb4babd 100644 --- a/packages/drivers/odsp-driver/src/epochTracker.ts +++ b/packages/drivers/odsp-driver/src/epochTracker.ts @@ -126,16 +126,38 @@ export class EpochTracker implements IPersistedFileCache { : maximumCacheDurationMs; } - // public for UT purposes only! - public setEpoch(epoch: string, fromCache: boolean, fetchType: FetchTypeInternal): void { + /** + * Sets the initial epoch and records where it was learned. + */ + public setEpoch(epoch: string, source: FetchTypeInternal | "pendingState"): void; + /** + * Sets the initial epoch. + * + * @deprecated Use the overload that accepts a source instead. + */ + public setEpoch(epoch: string, fromCache: boolean, fetchType: FetchTypeInternal): void; + public setEpoch( + epoch: string, + sourceOrFromCache: boolean | FetchTypeInternal | "pendingState", + fetchType?: FetchTypeInternal, + ): void { + const source = + typeof sourceOrFromCache === "boolean" + ? sourceOrFromCache + ? "cache" + : fetchType + : sourceOrFromCache; + assert( + source !== undefined, + "Fetch type is required when using the legacy setEpoch overload", + ); assert(this._fluidEpoch === undefined, 0x1db /* "epoch exists" */); this._fluidEpoch = epoch; this.loggerInternal.sendTelemetryEvent({ eventName: "EpochLearnedFirstTime", epoch, - fetchType, - fromCache, + source, }); } @@ -155,7 +177,7 @@ export class EpochTracker implements IPersistedFileCache { } assert(value.fluidEpoch !== undefined, 0x1dc /* "all entries have to have epoch" */); if (this._fluidEpoch === undefined) { - this.setEpoch(value.fluidEpoch, true, "cache"); + this.setEpoch(value.fluidEpoch, "cache"); // Epoch mismatch, the cached value is considerably different from what the current state of // the runtime and should not be used } else if (this._fluidEpoch !== value.fluidEpoch) { @@ -455,7 +477,7 @@ export class EpochTracker implements IPersistedFileCache { throw error; } if (epochFromResponse !== undefined && this._fluidEpoch === undefined) { - this.setEpoch(epochFromResponse, fromCache, fetchType); + this.setEpoch(epochFromResponse, fromCache ? "cache" : fetchType); } } diff --git a/packages/drivers/odsp-driver/src/odspDocumentService.ts b/packages/drivers/odsp-driver/src/odspDocumentService.ts index 884fc6a3ee26..cc16bd04f8d7 100644 --- a/packages/drivers/odsp-driver/src/odspDocumentService.ts +++ b/packages/drivers/odsp-driver/src/odspDocumentService.ts @@ -27,6 +27,7 @@ import { createChildMonitoringContext, type MonitoringContext, type TelemetryLoggerExt, + UsageError, } from "@fluidframework/telemetry-utils/internal"; import type { HostStoragePolicyInternal } from "./contracts.js"; @@ -166,6 +167,42 @@ export class OdspDocumentService return this._policies; } + public readonly driverStatePersistence = { + get: (): Record | undefined => { + const epoch = this.epochTracker.fluidEpoch; + return epoch === undefined + ? undefined + : { documentId: this.odspResolvedUrl.hashedDocumentId, epoch }; + }, + set: (state: unknown): void => { + if ( + typeof state !== "object" || + state === null || + Array.isArray(state) || + !("epoch" in state) || + typeof state.epoch !== "string" || + state.epoch.length === 0 || + !("documentId" in state) || + typeof state.documentId !== "string" + ) { + throw new UsageError( + "ODSP driver state must contain a non-empty epoch and document ID string", + ); + } + if (state.documentId !== this.odspResolvedUrl.hashedDocumentId) { + throw new UsageError("ODSP driver state belongs to a different document"); + } + const currentEpoch = this.epochTracker.fluidEpoch; + if (currentEpoch === state.epoch) { + return; + } + if (currentEpoch !== undefined) { + throw new UsageError("ODSP driver state epoch does not match the current epoch"); + } + this.epochTracker.setEpoch(state.epoch, "pendingState"); + }, + }; + /** * Connects to a storage endpoint for snapshot service. * diff --git a/packages/drivers/odsp-driver/src/test/joinSessionCacheTests.spec.ts b/packages/drivers/odsp-driver/src/test/joinSessionCacheTests.spec.ts index dd82bc12e8eb..7cc9a30dc3e9 100644 --- a/packages/drivers/odsp-driver/src/test/joinSessionCacheTests.spec.ts +++ b/packages/drivers/odsp-driver/src/test/joinSessionCacheTests.spec.ts @@ -56,6 +56,41 @@ describe("expose joinSessionInfo Tests", () => { async (_options) => "token", ); + it("preserves an epoch supplied in pending driver state", async () => { + const resolver = new OdspDriverUrlResolver(); + const odspResolvedUrl = await resolver.resolve({ + url: createOdspUrl({ driveId, itemId, siteUrl, dataStorePath: "/" }), + }); + const logger = new MockLogger(); + const service = await odspDocumentServiceFactory.createDocumentService( + odspResolvedUrl, + logger.toTelemetryLogger(), + ); + assert(service.driverStatePersistence !== undefined); + const driverState = { documentId: odspResolvedUrl.hashedDocumentId, epoch: "epoch1" }; + assert.throws( + () => service.driverStatePersistence?.set({ ...driverState, epoch: "" }), + /ODSP driver state must contain a non-empty epoch/, + ); + service.driverStatePersistence.set(driverState); + logger.assertMatch([ + { + eventName: "OdspDriver:EpochLearnedFirstTime", + source: "pendingState", + }, + ]); + service.driverStatePersistence.set(driverState); + assert.deepStrictEqual(service.driverStatePersistence.get(), driverState); + assert.throws( + () => service.driverStatePersistence?.set({ ...driverState, epoch: "epoch2" }), + /ODSP driver state epoch does not match the current epoch/, + ); + assert.throws( + () => service.driverStatePersistence?.set({ ...driverState, documentId: "other" }), + /ODSP driver state belongs to a different document/, + ); + }); + function addJoinSessionStub(): SinonStub { const joinSessionStub = stub(fetchJoinSession, mockify.key).callsFake( async () => joinSessionResponse, diff --git a/packages/loader/container-loader/src/container.ts b/packages/loader/container-loader/src/container.ts index 363e39459ec5..9916820eaff0 100644 --- a/packages/loader/container-loader/src/container.ts +++ b/packages/loader/container-loader/src/container.ts @@ -1198,6 +1198,7 @@ export class Container this.clientId, this.runtime, this.resolvedUrl, + this.service?.driverStatePersistence?.get(), ); return pendingState; } @@ -1567,7 +1568,9 @@ export class Container */ private async createDocumentService( resolvedUrl: IResolvedUrl, - props: { mode: "load" } | { mode: "attach"; summary: ISummaryTree | undefined }, + props: + | { mode: "load"; driverState?: unknown } + | { mode: "attach"; summary: ISummaryTree | undefined }, ): Promise { let service: IDocumentService; if (props.mode === "load") { @@ -1576,6 +1579,19 @@ export class Container this.subLogger, this.client.details.type === summarizerClientType, ); + if (props.driverState !== undefined) { + // Restoration errors deliberately reject the load: ignoring malformed or foreign + // state could reconnect without the driver's persisted consistency protection. + try { + if (service.driverStatePersistence === undefined) { + throw new UsageError("Document service cannot restore pending driver state"); + } + service.driverStatePersistence.set(props.driverState); + } catch (error) { + service.dispose(error); + throw error; + } + } if (service.on !== undefined) { // Back-compat for Old driver service.on("metadataUpdate", this.metadataUpdateHandler); @@ -1624,7 +1640,10 @@ export class Container numUnsummarizedOps: number; }> { const timings: Record = { phase1: performanceNow() }; - this.service = await this.createDocumentService(resolvedUrl, { mode: "load" }); + this.service = await this.createDocumentService(resolvedUrl, { + mode: "load", + driverState: pendingLocalState?.driverState, + }); // Except in cases where its requested by feature gate, the container will connect in "read" mode const mode = diff --git a/packages/loader/container-loader/src/createAndLoadContainerUtils.ts b/packages/loader/container-loader/src/createAndLoadContainerUtils.ts index f513e4a66c67..25d2c4432d23 100644 --- a/packages/loader/container-loader/src/createAndLoadContainerUtils.ts +++ b/packages/loader/container-loader/src/createAndLoadContainerUtils.ts @@ -419,8 +419,7 @@ export async function loadFrozenContainerFromPendingState( if (driverWiring === "none") { // Offline: synthesize the driver wiring from the URL captured in pending state. // The container's load pipeline is reused unchanged — the synthesized resolver - // returns a resolved URL whose `url` equals `pendingLocalState.url`, so the - // identity guard in `Loader.resolveCore` is trivially satisfied. + // returns a resolved URL whose `url` equals `pendingLocalState.url`. const pending = getAttachedContainerStateFromSerializedContainer(pendingLocalState); if (pending.blobContentsMode === "reference") { throw new UsageError( @@ -467,9 +466,7 @@ export async function loadFrozenContainerFromPendingState( return loadExistingContainer({ ...props, // `request.url` is unused: `synthesizedUrlResolver.resolve()` returns - // `synthesizedResolvedUrl` regardless of input, and the identity - // guard in `Loader.resolveCore` compares the resolver's output URL - // against `pendingLocalState.url` — both equal `pending.url` here. + // `synthesizedResolvedUrl` regardless of input. // Using a recognizable opaque placeholder instead of the resolved-form // URL avoids implying that any downstream stage interprets it as a // request-form URL. @@ -740,6 +737,7 @@ export async function captureFullContainerState({ pendingRuntimeState: undefined, savedOps, url: resolvedUrl.url, + driverState: documentService.driverStatePersistence?.get(), }; return JSON.stringify(pendingState); } finally { diff --git a/packages/loader/container-loader/src/frozenServices.ts b/packages/loader/container-loader/src/frozenServices.ts index d09408427d72..3a4c18a3cee1 100644 --- a/packages/loader/container-loader/src/frozenServices.ts +++ b/packages/loader/container-loader/src/frozenServices.ts @@ -96,6 +96,10 @@ class FrozenDocumentService // a single field because `IDocumentService.connectToStorage` is a public API that can be // called more than once — we cannot assume the Container holds a single instance. private readonly storageServices = new Set(); + private driverState: unknown; + public readonly driverStatePersistence?: NonNullable< + IDocumentService["driverStatePersistence"] + >; constructor( public readonly resolvedUrl: IResolvedUrl, @@ -116,6 +120,17 @@ class FrozenDocumentService // indistinguishable from a normal container at the policies layer; downstream behavior // flows through the live `WritableFrozenDeltaStream` instead. this.policies = readOnly ? { storageOnly: true } : {}; + const innerPersistence = this.documentService?.driverStatePersistence; + if (this.documentService === undefined || innerPersistence !== undefined) { + this.driverStatePersistence = { + get: () => + innerPersistence === undefined ? this.driverState : innerPersistence.get(), + set: (state) => { + this.driverState = state; + innerPersistence?.set(state); + }, + }; + } } public readonly policies: IDocumentServicePolicies; diff --git a/packages/loader/container-loader/src/loader.ts b/packages/loader/container-loader/src/loader.ts index f8609992ad94..ef9205ada0fa 100644 --- a/packages/loader/container-loader/src/loader.ts +++ b/packages/loader/container-loader/src/loader.ts @@ -333,7 +333,9 @@ export class Loader implements IHostLoader { throw new Error(`Invalid URL ${resolvedAsFluid.url}`); } - if (pendingLocalState !== undefined) { + // Driver state owns document identity when present. Older pending state falls back to the + // loader's URL-shape-dependent validation. + if (pendingLocalState !== undefined && pendingLocalState.driverState === undefined) { const parsedPendingUrl = tryParseCompatibleResolvedUrl(pendingLocalState.url); if ( parsedPendingUrl?.id !== parsed.id || diff --git a/packages/loader/container-loader/src/serializedStateManager.ts b/packages/loader/container-loader/src/serializedStateManager.ts index 45da251ad233..175a12dc8550 100644 --- a/packages/loader/container-loader/src/serializedStateManager.ts +++ b/packages/loader/container-loader/src/serializedStateManager.ts @@ -112,13 +112,22 @@ export interface IPendingContainerState extends SnapshotWithBlobs { */ savedOps: ISequencedDocumentMessage[]; /** - * The Container's URL in the service, needed to hook up the driver during rehydration + * The Container's URL in the service, needed to hook up the driver during rehydration and to + * validate the document identity of legacy pending state that has no {@link driverState}. */ url: string; /** * If the Container was connected when serialized, its clientId. Used as the initial clientId upon rehydration, until reconnected. */ clientId?: string; + /** + * Opaque state supplied by the document service for use when rehydrating. This value is + * serialized as part of the containing pending state and persisted by the host. It must + * round-trip through `JSON.stringify` and `JSON.parse` without custom serialization or + * information loss, must not contain customer-identifying information, and is responsible for + * validating that it belongs to the document being loaded. + */ + driverState?: unknown; } /** @@ -423,6 +432,7 @@ export class SerializedStateManager implements IDisposable { clientId: string | undefined, runtime: Pick, resolvedUrl: IResolvedUrl, + driverState?: unknown, ): Promise { this.verifyNotDisposed(); if (!this.offlineLoadEnabled) { @@ -476,6 +486,7 @@ export class SerializedStateManager implements IDisposable { savedOps: this.processedOps, url: resolvedUrl.url, clientId, + driverState, }; return JSON.stringify(pendingState); diff --git a/packages/loader/container-loader/src/test/frozenServices.spec.ts b/packages/loader/container-loader/src/test/frozenServices.spec.ts index 430ba510f229..67b647b50e69 100644 --- a/packages/loader/container-loader/src/test/frozenServices.spec.ts +++ b/packages/loader/container-loader/src/test/frozenServices.spec.ts @@ -169,6 +169,70 @@ describe("FrozenDocumentService.connectToDeltaStream", () => { }); }); +describe("FrozenDocumentService driver state", () => { + it("preserves state without an inner document service", async () => { + const service = await new FrozenDocumentServiceFactory(false).createDocumentService( + fakeUrl, + ); + const state = { epoch: "epoch1" }; + + service.driverStatePersistence?.set(state); + + assert.deepStrictEqual(service.driverStatePersistence?.get(), state); + }); + + it("uses the inner service state when its persistence capability is present", async () => { + let innerState: unknown; + const innerService = { + resolvedUrl: fakeUrl, + policies: {}, + driverStatePersistence: { + get: () => innerState, + set: () => {}, + }, + } as unknown as IDocumentService; + const innerFactory: IDocumentServiceFactory = { + createDocumentService: async () => innerService, + createContainer: async () => { + throw new Error("not used in this test"); + }, + }; + const service = await new FrozenDocumentServiceFactory( + false, + innerFactory, + ).createDocumentService(fakeUrl); + + service.driverStatePersistence?.set({ epoch: "restored" }); + assert.strictEqual(service.driverStatePersistence?.get(), undefined); + + const nullState: unknown = JSON.parse("null"); + innerState = nullState; + assert.strictEqual(service.driverStatePersistence?.get(), nullState); + + innerState = { epoch: "live" }; + assert.deepStrictEqual(service.driverStatePersistence?.get(), innerState); + }); + + it("does not expose persistence when the inner service cannot validate state", async () => { + const innerService = { + resolvedUrl: fakeUrl, + policies: {}, + } as unknown as IDocumentService; + const innerFactory: IDocumentServiceFactory = { + createDocumentService: async () => innerService, + createContainer: async () => { + throw new Error("not used in this test"); + }, + }; + const service = await new FrozenDocumentServiceFactory( + false, + innerFactory, + ).createDocumentService(fakeUrl); + + assert.strictEqual(service.driverStatePersistence, undefined); + }); +}); + describe("FrozenDocumentService disposal", () => { it("dispose() rejects in-flight createBlob promises on writable-frozen storage", async () => { // The writable-frozen `createBlob` returns a never-resolving promise so the diff --git a/packages/loader/container-loader/src/test/loader.spec.ts b/packages/loader/container-loader/src/test/loader.spec.ts index cdf7849a2ec2..3ce0a02c3680 100644 --- a/packages/loader/container-loader/src/test/loader.spec.ts +++ b/packages/loader/container-loader/src/test/loader.spec.ts @@ -237,6 +237,96 @@ describe("DisableLoadConnectionRetries", () => { return error; } + it("delegates identity validation and restores driver state before connecting", async () => { + const driverState = { epoch: "epoch1" }; + let restored = false; + let shouldThrowRestoreError = false; + let supportsDriverStatePersistence = true; + const restoreError = new Error("invalid driver state"); + let connectionAttempts = 0; + const disposalErrors: unknown[] = []; + const documentServiceFactory = failSometimeProxy< + IDocumentServiceFactory & + IProvideLayerCompatDetails & + IProvideLayerCompatSupportRequirements + >({ + createDocumentService: async () => + failSometimeProxy({ + policies: {}, + resolvedUrl, + driverStatePersistence: supportsDriverStatePersistence + ? { + get: () => undefined, + set: (state) => { + assert.deepStrictEqual(state, driverState); + if (shouldThrowRestoreError) { + throw restoreError; + } + restored = true; + }, + } + : AbsentProperty, + connectToStorage: async () => { + connectionAttempts++; + assert(restored, "Driver state must be restored before connecting to storage"); + throw new Error("stop after ordering check"); + }, + connectToDeltaStream: async () => new Promise(() => {}), + on: AbsentProperty, + off: AbsentProperty, + dispose: (error) => disposalErrors.push(error), + }), + ILayerCompatDetails: AbsentProperty, + ILayerCompatSupportRequirements: AbsentProperty, + }); + const loader = new Loader({ + codeLoader: createTestCodeLoaderProxy(), + documentServiceFactory, + urlResolver, + }); + const pendingStateData = { + attached: true, + baseSnapshot: { blobs: {}, trees: {} }, + snapshotBlobs: {}, + pendingRuntimeState: {}, + savedOps: [], + url: `https://localhost/tenant/${uuid()}`, + driverState, + }; + const pendingState = JSON.stringify(pendingStateData); + + await assert.rejects( + async () => loader.resolve({ url: "test" }, pendingState), + /stop after ordering check/, + ); + assert(restored); + assert.strictEqual(connectionAttempts, 1); + + shouldThrowRestoreError = true; + await assert.rejects( + async () => loader.resolve({ url: "test" }, pendingState), + restoreError, + ); + assert.strictEqual(connectionAttempts, 1); + assert.strictEqual(disposalErrors.at(-1), restoreError); + + await assert.rejects( + async () => + loader.resolve( + { url: "test" }, + JSON.stringify({ ...pendingStateData, driverState: undefined }), + ), + /does not match pending state URL/, + ); + + supportsDriverStatePersistence = false; + await assert.rejects( + async () => loader.resolve({ url: "test" }, pendingState), + /Document service cannot restore pending driver state/, + ); + assert.strictEqual(connectionAttempts, 1); + }); + it("load rejects when connectToStorage fails with retryable error and flag is enabled", async () => { const documentServiceFactory = failSometimeProxy< IDocumentServiceFactory & diff --git a/packages/loader/container-loader/src/test/serializedStateManager.spec.ts b/packages/loader/container-loader/src/test/serializedStateManager.spec.ts index 7f7bf2951786..d5b45402651a 100644 --- a/packages/loader/container-loader/src/test/serializedStateManager.spec.ts +++ b/packages/loader/container-loader/src/test/serializedStateManager.spec.ts @@ -243,11 +243,13 @@ describe("serializedStateManager", () => { ); // equivalent to attach serializedStateManager.setInitialSnapshot(initialSnapshot); - await serializedStateManager.getPendingLocalState( + const state = await serializedStateManager.getPendingLocalState( "clientId", new MockRuntime(), resolvedUrl, + "epoch", ); + assert.strictEqual((JSON.parse(state) as IPendingContainerState).driverState, "epoch"); }); it("can get pending local state from previous pending state", async () => { diff --git a/packages/loader/driver-utils/src/documentServiceProxy.ts b/packages/loader/driver-utils/src/documentServiceProxy.ts index b7a6133bd83a..1fe64ecbc066 100644 --- a/packages/loader/driver-utils/src/documentServiceProxy.ts +++ b/packages/loader/driver-utils/src/documentServiceProxy.ts @@ -23,8 +23,15 @@ export abstract class DocumentServiceProxy extends TypedEventEmitter implements IDocumentService { + public readonly driverStatePersistence?: NonNullable< + IDocumentService["driverStatePersistence"] + >; + constructor(private readonly _service: IDocumentService) { super(); + if (_service.driverStatePersistence !== undefined) { + this.driverStatePersistence = _service.driverStatePersistence; + } } public get service(): IDocumentService { diff --git a/packages/test/test-end-to-end-tests/src/test/pointInTime/epochMismatch.spec.ts b/packages/test/test-end-to-end-tests/src/test/pointInTime/epochMismatch.spec.ts index 320eac88c18b..8d758095eb9f 100644 --- a/packages/test/test-end-to-end-tests/src/test/pointInTime/epochMismatch.spec.ts +++ b/packages/test/test-end-to-end-tests/src/test/pointInTime/epochMismatch.spec.ts @@ -31,6 +31,11 @@ import { strict as assert } from "assert"; import { describeCompat, itExpects } from "@fluid-private/test-version-utils"; +import { + createLoader, + getRequiredPendingLocalState, + waitForContainerConnection, +} from "@fluidframework/test-utils/internal"; import { listFileVersions, restoreFileVersion } from "./odspVersionTestApi.js"; import { @@ -45,6 +50,57 @@ describeCompat( (getTestObjectProvider, apis) => { const suite = setupPointInTimeSuite(getTestObjectProvider, apis); + itExpects( + "rejects pending state captured before a file restore", + [ + { + eventName: "fluid:telemetry:Container:ContainerClose", + errorType: "fileOverwrittenInStorage", + }, + ], + async function (this: Mocha.Context) { + this.timeout(120_000); + + const ctx = await createPointInTimeTestContext(suite, apis, { + withSummarizer: true, + }); + await ctx.incrementAndSync(1); + const olderVersion = await ctx.snapVersion("pending-state-base"); + await ctx.incrementAndSync(1); + await ctx.snapVersion("pending-state-tip"); + + ctx.container.disconnect(); + ctx.dataObject.increment(); + const pendingState = await getRequiredPendingLocalState(ctx.container); + ctx.container.close(); + + const restored = await restoreFileVersion(ctx.versionApi, olderVersion.id); + assert.strictEqual(restored, true, "restore should succeed (HTTP 204)"); + + const provider = suite.provider(); + const loader = createLoader( + [[provider.defaultCodeDetails, suite.runtimeFactory()]], + provider.documentServiceFactory, + provider.urlResolver, + provider.logger, + ); + const url = await provider.driver.createContainerUrl(ctx.documentId); + + await assert.rejects( + async () => { + const resumed = await loader.resolve({ url }, pendingState); + await waitForContainerConnection(resumed, true, { + durationMs: 60_000, + errorMsg: "timed out waiting for the stale pending-state container to close", + }); + }, + (error: Error & { errorType?: string }) => + error.errorType === "fileOverwrittenInStorage", + "expected pending state from the previous epoch to be rejected", + ); + }, + ); + // The failed point-in-time load closes its container with the driver's non-retryable // fileOverwrittenInStorage (epoch-mismatch) error. That ContainerClose is the expected // outcome, so declare it via itExpects; otherwise describeCompat's afterEach hook would diff --git a/packages/test/test-service-load/src/faultInjectionDriver.ts b/packages/test/test-service-load/src/faultInjectionDriver.ts index 986ba343e99d..d1dd6df528fd 100644 --- a/packages/test/test-service-load/src/faultInjectionDriver.ts +++ b/packages/test/test-service-load/src/faultInjectionDriver.ts @@ -114,8 +114,15 @@ export class FaultInjectionDocumentService this._currentDeltaStorage?.goOnline(); } + public readonly driverStatePersistence?: NonNullable< + IDocumentService["driverStatePersistence"] + >; + constructor(private readonly internal: IDocumentService) { super(); + if (internal.driverStatePersistence !== undefined) { + this.driverStatePersistence = internal.driverStatePersistence; + } } public get resolvedUrl(): IResolvedUrl {