Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -650,6 +650,57 @@ test('delivers a mid-session tail append even while a history window is resident
assert.equal(replica.durableThrough, 5);
});

test('advances a projected transcript across hidden durable records', async () => {
const visible = (sequence: number) => ({
identity: sequence,
message: userMessage(`Visible ${sequence}`, `user-${sequence}`),
});
const bootstrapPage = transcriptPage('older', null, 1);
const visibleAdvancePage = transcriptPage('newer', null, 5);
const hiddenAdvancePage = transcriptPage('newer', null, 6);
const changes: { durableUpserts: readonly { sequence: number }[] }[] = [];
const handle = runtimeHostSessionFixture({
snapshot: continuitySnapshot(),
transcript: Promise.resolve([]),
events: { async *[Symbol.asyncIterator]() {} },
transcriptBootstrap: {
throughSequence: 1,
durableCoverage: 'projected',
overlayMessageCount: 0,
durable: bootstrapPage,
overlay: { ...transcriptPage('older', null, 1), source: 'overlay' },
},
loadTranscriptOverlay: async () => [],
decodeTranscriptPage: async (page) => ({
messages:
page === bootstrapPage
? [visible(0)]
: page === visibleAdvancePage
? [visible(3), visible(4)]
: [],
nextCursor: null,
}),
loadTranscriptPage: async ({ throughSequence }) =>
throughSequence === 5 ? visibleAdvancePage : hiddenAdvancePage,
async close() {},
});
const replica = await DesktopTranscriptReplica.prepare(handle, {
onChange: (_replica, change) => changes.push(change),
});

// Sequences 1, 2, 5, and 6 are valid Host-private records omitted from the
// Guest projection. The physical watermark still advances across them.
await replica.advance(5);
await replica.advance(6);

assert.equal(replica.durableThrough, 6);
assert.deepEqual(replica.snapshot().durable.map(({ sequence }) => sequence), [0, 3, 4]);
assert.deepEqual(
changes.flatMap((change) => change.durableUpserts.map(({ sequence }) => sequence)),
[3, 4],
);
});

test('keeps an oversized streaming Turn visible when its overlay settles', async () => {
const older = {
identity: 0,
Expand Down
24 changes: 16 additions & 8 deletions apps/desktop/src/main/desktop-transcript-replica.ts
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ export class DesktopTranscriptReplica {
readonly generation: string;
readonly hostEpoch: string;
readonly #handle: DesktopRuntimeHostSession;
readonly #durableCoverage: DesktopRuntimeHostSession['transcriptBootstrap']['durableCoverage'];
readonly #maxResidentBytes: number;
readonly #maxResidentTurns: number;
readonly #maxOverlayBytes: number;
Expand Down Expand Up @@ -110,6 +111,7 @@ export class DesktopTranscriptReplica {
options: DesktopTranscriptReplicaOptions,
) {
this.#handle = handle;
this.#durableCoverage = handle.transcriptBootstrap.durableCoverage;
this.sessionId = handle.snapshot.session.sessionId;
this.generation = options.generation ?? randomUUID();
this.hostEpoch = handle.hostEpoch;
Expand Down Expand Up @@ -255,7 +257,7 @@ export class DesktopTranscriptReplica {
if (
anchor !== null &&
decoded.messages.length > 0 &&
decoded.messages.at(-1)!.identity !== anchor - 1
!this.#matchesCoverageStep(anchor, decoded.messages.at(-1)!.identity + 1)
) {
throw correlationError('Desktop transcript older page did not meet its anchor');
}
Expand Down Expand Up @@ -308,7 +310,7 @@ export class DesktopTranscriptReplica {
if (
decoded.messages.length > 0 &&
(loadTail
? decoded.messages.at(-1)!.identity !== sequence
? !this.#matchesCoverageStep(sequence, decoded.messages.at(-1)!.identity)
: decoded.messages[0]!.identity !== sequence)
) {
throw correlationError('Desktop transcript range did not meet its anchor');
Expand Down Expand Up @@ -406,7 +408,7 @@ export class DesktopTranscriptReplica {
}
let cursor: string | null = null;
const anchorSequence = this.#durableThrough;
let expectedSequence = (anchorSequence ?? -1) + 1;
let nextSequence = (anchorSequence ?? -1) + 1;
do {
if (!this.#resident) return;
const page: SessionTranscriptPage = await this.#handle.loadTranscriptPage({
Expand All @@ -426,12 +428,12 @@ export class DesktopTranscriptReplica {
this.#acceptRange(decoded.messages);
if (
decoded.messages.length > 0 &&
decoded.messages[0]!.identity !== expectedSequence
!this.#matchesCoverageStep(decoded.messages[0]!.identity, nextSequence)
) {
throw correlationError('Desktop transcript catch-up has a sequence gap');
}
if (decoded.messages.length > 0) {
expectedSequence = decoded.messages.at(-1)!.identity + 1;
nextSequence = decoded.messages.at(-1)!.identity + 1;
}
const completedOverlayMessageIds = this.#installDurable(decoded.messages);
const evictedDurableSequences = this.#evictToBudget(
Expand All @@ -445,13 +447,13 @@ export class DesktopTranscriptReplica {
} while (cursor !== null);
// A concurrent `discard()` (LRU reclaim for another observed session) can
// flip `#resident` to false across any page `await` above. The per-page
// callback already returns early in that case, so `expectedSequence` is
// callback already returns early in that case, so `nextSequence` is
// left short of the watermark. Without this guard the check below would
// turn a benign memory reclaim into a fatal `correlation_changed` that
// drives the session terminal. A discarded replica has no watermark to
// meet, so return cleanly and let a later resume re-catch-up.
if (!this.#resident) return;
if (expectedSequence !== target + 1) {
if (this.#durableCoverage === 'complete' && nextSequence !== target + 1) {
throw correlationError('Desktop transcript catch-up ended before its watermark');
}
this.#durableThrough = target;
Expand Down Expand Up @@ -505,12 +507,18 @@ export class DesktopTranscriptReplica {
for (let index = 1; index < messages.length; index += 1) {
const previous = messages[index - 1]!.identity;
const current = messages[index]!.identity;
if (current !== previous + 1) {
if (!this.#matchesCoverageStep(current, previous + 1)) {
throw correlationError('Desktop transcript page has a sequence gap');
}
}
}

#matchesCoverageStep(sequence: number, firstPossibleSequence: number): boolean {
return this.#durableCoverage === 'complete'
? sequence === firstPossibleSequence
: sequence >= firstPossibleSequence;
}

#publish(
messages: readonly {
readonly identity: number;
Expand Down