Skip to content

Commit 8fac08a

Browse files
committed
mason: stabilize marker advance wire replay
1 parent dbdcda7 commit 8fac08a

3 files changed

Lines changed: 206 additions & 1 deletion

File tree

packages/plugin/src/hooks/magic-context/compaction-marker-manager.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,8 @@ export type MarkerUpdateOutcome =
5555
markerOrdinal: number;
5656
summaryMessageId: string;
5757
boundaryMessageId: string;
58+
/** Summary source removed from the database during an advance, if any. */
59+
removedSummaryMessageId: string | null;
5860
}
5961
| { kind: "already-current" }
6062
| {
@@ -255,6 +257,9 @@ export function applyDeferredCompactionMarker(
255257

256258
// Remove old marker if present. `removeCompactionMarker` returns false
257259
// only when the DELETE transaction itself failed (e.g. SQLITE_BUSY).
260+
// Keep the summary id so the transform can remove the old source from
261+
// the in-memory drain representation before it serves the new marker.
262+
const removedSummaryMessageId = existing?.summaryMessageId ?? null;
258263
// No-op success on already-missing rows is fine — that's why retry is
259264
// safe. False here means we couldn't even attempt the delete cleanly;
260265
// bail to retryable WITHOUT calling inject (avoids leaving two marker
@@ -309,6 +314,7 @@ export function applyDeferredCompactionMarker(
309314
markerOrdinal: pending.ordinal,
310315
summaryMessageId: result.summaryMessageId,
311316
boundaryMessageId: result.boundaryMessageId,
317+
removedSummaryMessageId,
312318
};
313319
} catch (err) {
314320
// Thrown paths:

packages/plugin/src/hooks/magic-context/transform-postprocess-phase.test.ts

Lines changed: 180 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,10 @@ import {
2020
setPendingCompactionMarkerState,
2121
} from "../../features/magic-context/storage";
2222
import { initializeDatabase } from "../../features/magic-context/storage-db";
23-
import { getPersistedCompactionMarkerState } from "../../features/magic-context/storage-meta-persisted";
23+
import {
24+
getPersistedCompactionMarkerState,
25+
setPersistedCompactionMarkerState,
26+
} from "../../features/magic-context/storage-meta-persisted";
2427
import { createTagger } from "../../features/magic-context/tagger";
2528
import { Database } from "../../shared/sqlite";
2629
import { MARKER_SUMMARY_TEXT } from "./compaction-marker-manager";
@@ -484,6 +487,182 @@ describe("deferred compaction marker representation", () => {
484487
});
485488
});
486489

490+
describe("deferred compaction marker advance representation", () => {
491+
it("keeps the advance drain byte-identical with the next pass after removing the old marker", async () => {
492+
db = new Database(":memory:");
493+
initializeDatabase(db);
494+
const sessionId = "ses-marker-advance-wire-stability";
495+
const dataHome = mkdtempSync(join(tmpdir(), "postprocess-marker-advance-wire-"));
496+
tempDirs.push(dataHome);
497+
process.env.XDG_DATA_HOME = dataHome;
498+
mkdirSync(join(dataHome, "opencode"), { recursive: true });
499+
const opencodeDb = new Database(join(dataHome, "opencode", "opencode.db"));
500+
opencodeDb.exec(
501+
"CREATE TABLE message (id TEXT PRIMARY KEY, session_id TEXT, time_created INTEGER, time_updated INTEGER, data TEXT)",
502+
);
503+
opencodeDb.exec(
504+
"CREATE TABLE part (id TEXT PRIMARY KEY, message_id TEXT, session_id TEXT, time_created INTEGER, time_updated INTEGER, data TEXT)",
505+
);
506+
const insertMessage = opencodeDb.prepare(
507+
"INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES (?, ?, ?, ?, ?)",
508+
);
509+
insertMessage.run(
510+
"old-boundary",
511+
sessionId,
512+
1_000,
513+
1_000,
514+
JSON.stringify({ role: "user" }),
515+
);
516+
insertMessage.run(
517+
"new-boundary",
518+
sessionId,
519+
2_000,
520+
2_000,
521+
JSON.stringify({ role: "user" }),
522+
);
523+
insertMessage.run(
524+
"new-end",
525+
sessionId,
526+
3_000,
527+
3_000,
528+
JSON.stringify({ role: "assistant", finish: "stop" }),
529+
);
530+
insertMessage.run(
531+
"old-summary",
532+
sessionId,
533+
1_001,
534+
1_001,
535+
JSON.stringify({ role: "assistant", summary: true, finish: "stop" }),
536+
);
537+
opencodeDb.prepare(
538+
"INSERT INTO part (id, message_id, session_id, time_created, time_updated, data) VALUES (?, ?, ?, ?, ?, ?)",
539+
).run(
540+
"old-compaction",
541+
"old-boundary",
542+
sessionId,
543+
1_000,
544+
1_000,
545+
JSON.stringify({ type: "compaction", auto: true }),
546+
);
547+
opencodeDb.prepare(
548+
"INSERT INTO part (id, message_id, session_id, time_created, time_updated, data) VALUES (?, ?, ?, ?, ?, ?)",
549+
).run(
550+
"old-summary-part",
551+
"old-summary",
552+
sessionId,
553+
1_001,
554+
1_001,
555+
JSON.stringify({ type: "text", text: MARKER_SUMMARY_TEXT }),
556+
);
557+
opencodeDb.close();
558+
559+
appendCompartments(db, sessionId, [
560+
{
561+
sequence: 0,
562+
startMessage: 1,
563+
endMessage: 20,
564+
startMessageId: "old-boundary",
565+
endMessageId: "new-end",
566+
title: "marker advance",
567+
content: "test content",
568+
},
569+
]);
570+
setPersistedCompactionMarkerState(db, sessionId, {
571+
boundaryMessageId: "old-boundary",
572+
summaryMessageId: "old-summary",
573+
compactionPartId: "old-compaction",
574+
summaryPartId: "old-summary-part",
575+
boundaryOrdinal: 10,
576+
targetEndMessageId: "old-end",
577+
});
578+
setPendingCompactionMarkerState(db, sessionId, {
579+
ordinal: 20,
580+
endMessageId: "new-end",
581+
publishedAt: 2,
582+
});
583+
584+
const tagger = createTagger();
585+
const drainMessages = [
586+
{
587+
info: { role: "user", sessionID: sessionId },
588+
parts: [{ type: "text", text: "<session-history>\\n\\n</session-history>" }],
589+
},
590+
{
591+
info: {
592+
id: "old-summary",
593+
role: "assistant",
594+
sessionID: sessionId,
595+
summary: true,
596+
finish: "stop",
597+
},
598+
parts: [{ type: "text", text: MARKER_SUMMARY_TEXT }],
599+
},
600+
{
601+
info: { id: "new-end", role: "assistant", sessionID: sessionId, finish: "stop" },
602+
parts: [{ type: "text", text: "tail content" }],
603+
},
604+
] as unknown as MessageLike[];
605+
const taggedDrain = tagMessages(sessionId, drainMessages, tagger, db);
606+
const deferredHistoryRefreshSessions = new Set<string>([sessionId]);
607+
608+
await runPostTransformPhase(
609+
basePostTransformArgs(db, sessionId, drainMessages, {
610+
fullFeatureMode: false,
611+
tagger,
612+
targets: taggedDrain.targets,
613+
reasoningByMessage: taggedDrain.reasoningByMessage,
614+
messageTagNumbers: taggedDrain.messageTagNumbers,
615+
batch: taggedDrain.batch,
616+
deferredHistoryWasPendingAtPassStart: true,
617+
historyRebuiltThisPass: true,
618+
canConsumeDeferredLate: true,
619+
deferredHistoryRefreshSessions,
620+
pendingCompartmentInjection: {
621+
block: "",
622+
compartmentEndMessage: 20,
623+
compartmentEndMessageId: "new-end",
624+
compartmentCount: 1,
625+
skippedVisibleMessages: 0,
626+
factCount: 0,
627+
memoryCount: 0,
628+
rebuiltFromDb: true,
629+
},
630+
}),
631+
);
632+
633+
const marker = getPersistedCompactionMarkerState(db, sessionId);
634+
expect(marker?.summaryMessageId).toBeString();
635+
expect(drainMessages.some((message) => message.info.id === "old-summary")).toBe(false);
636+
637+
const rebuiltMessages = [
638+
{
639+
info: { role: "user", sessionID: sessionId },
640+
parts: [{ type: "text", text: "<session-history>\\n\\n</session-history>" }],
641+
},
642+
{
643+
info: {
644+
id: marker?.summaryMessageId,
645+
role: "assistant",
646+
sessionID: sessionId,
647+
summary: true,
648+
finish: "stop",
649+
},
650+
parts: [{ type: "text", text: MARKER_SUMMARY_TEXT }],
651+
},
652+
{
653+
info: { id: "new-end", role: "assistant", sessionID: sessionId, finish: "stop" },
654+
parts: [{ type: "text", text: "tail content" }],
655+
},
656+
] as unknown as MessageLike[];
657+
tagger.initFromDb(sessionId, db);
658+
tagMessages(sessionId, rebuiltMessages, tagger, db);
659+
660+
expect(serializeAnthropicWireWithAdjacentAssistantMerge(drainMessages)).toBe(
661+
serializeAnthropicWireWithAdjacentAssistantMerge(rebuiltMessages),
662+
);
663+
});
664+
});
665+
487666
describe("deferred compaction marker CAS drain", () => {
488667
it("preserves the deferred-history signal when a newer pending blob exists", () => {
489668
db = new Database(":memory:");

packages/plugin/src/hooks/magic-context/transform-postprocess-phase.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -160,6 +160,22 @@ export function injectDeferredCompactionSummaryRepresentation(args: {
160160
return true;
161161
}
162162

163+
/**
164+
* Remove the prior marker source from the drain array when a boundary advances.
165+
* The database mutation happens before the next transform pass, so serving the
166+
* old summary during this pass would make the provider wire differ from replay.
167+
*/
168+
function removeDeferredCompactionSummaryRepresentation(
169+
messages: MessageLike[],
170+
summaryMessageId: string | null,
171+
): boolean {
172+
if (summaryMessageId === null) return false;
173+
const index = messages.findIndex((message) => message.info.id === summaryMessageId);
174+
if (index < 0) return false;
175+
messages.splice(index, 1);
176+
return true;
177+
}
178+
163179
function pendingMarkerCoveredByConsumedBoundary(
164180
pending: PendingCompactionMarker,
165181
injection: PreparedCompartmentInjection | null,
@@ -1363,6 +1379,10 @@ export async function runPostTransformPhase(
13631379
args.sessionDirectory,
13641380
);
13651381
if (outcome.kind === "applied") {
1382+
removeDeferredCompactionSummaryRepresentation(
1383+
args.messages,
1384+
outcome.removedSummaryMessageId,
1385+
);
13661386
injectDeferredCompactionSummaryRepresentation({
13671387
db: args.db,
13681388
sessionId: args.sessionId,

0 commit comments

Comments
 (0)