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
48 changes: 42 additions & 6 deletions packages/pi-plugin/src/commands/ctx-wrapup.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import {
clearProducerModelObservations,
observeProducerModelsForTest,
} from "@magic-context/core/hooks/magic-context/producer-window-test-support";
import { hasRawMessageProvider } from "@magic-context/core/hooks/magic-context/read-session-chunk";
import * as logger from "@magic-context/core/shared/logger";
import { Database } from "@magic-context/core/shared/sqlite";
import { closeQuietly } from "@magic-context/core/shared/sqlite-helpers";
Expand Down Expand Up @@ -485,10 +486,13 @@ describe("Pi /ctx-wrapup", () => {
}
});

it("persists tokens when wrapup uses the real Pi historian", async () => {
it.each([
[8, 100_000],
[24, 100],
])("persists tokens and covers the full target with the real Pi historian (%i messages, %i tokens)", async (messageCount, historianChunkTokens) => {
const db = createDb();
try {
const sessionId = "pi-wrapup-persisted-tokens";
const sessionId = `pi-wrapup-persisted-tokens-${messageCount}`;
const runner = {
harness: "pi",
run: mock(async (options: SubagentRunOptions) => {
Expand All @@ -510,9 +514,16 @@ describe("Pi /ctx-wrapup", () => {
const range = ranges.at(-1);
if (!range)
throw new Error("historian prompt did not include a message range");
const start = Number(range[1]);
const end = Number(range[2]);
// Non-final chunks need lookahead so the runner can discard only the last compartment.
const head =
start < end
? `<compartment start="${start}" end="${end - 1}" title="Pi wrapup"><p1>Summarized the eligible Pi history.</p1></compartment>`
: "";
return {
ok: true as const,
assistantText: `<compartment start="${range[1]}" end="${range[2]}" title="Pi wrapup"><p1>Summarized the eligible Pi history.</p1></compartment>`,
assistantText: `${head}<compartment start="${end}" end="${end}" title="Pi wrapup lookahead"><p1>Summarized the last message.</p1></compartment>`,
durationMs: 1,
};
}),
Expand All @@ -523,17 +534,42 @@ describe("Pi /ctx-wrapup", () => {
deps(db, {
runner,
runPiHistorianForWrapup: undefined,
historianChunkTokens: 100_000,
historianChunkTokens,
}),
ctx(sessionId, 8),
ctx(
sessionId,
branch(messageCount).map((entry, index) => ({
...entry,
message: {
role: index % 2 === 0 ? "user" : "assistant",
content: [{ type: "text", text: entry.message.content }],
},
})),
),
sessionId,
2,
);

expect(result).toContain("## Magic Wrapup");
expect(result).not.toContain("## Magic Wrapup — Partial");
expect(getLastCompartmentEndMessage(db, sessionId)).toBe(
messageCount - 2,
);
expect(getPendingPiCompactionMarkerState(db, sessionId)?.ordinal).toBe(
messageCount - 2,
);
expect(getWrapupInProgressState(db, sessionId)).toBeNull();
expect(hasRawMessageProvider(sessionId)).toBe(false);
const compartments = getCompartments(db, sessionId);
expect(compartments[0].startMessage).toBe(1);
for (let i = 1; i < compartments.length; i++) {
expect(compartments[i].startMessage).toBe(
compartments[i - 1].endMessage + 1,
);
}
const rows = getSubagentInvocations(db, sessionId);
expect(rows).toHaveLength(1);
if (messageCount === 8) expect(rows).toHaveLength(1);
else expect(rows.length).toBeGreaterThan(1);
expect(rows[0]).toMatchObject({
harness: "pi",
subagent: "historian",
Expand Down
117 changes: 117 additions & 0 deletions packages/plugin/src/hooks/magic-context/read-session-chunk.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,11 +11,14 @@ import { v2NonNarrativeStoredGapRanges } from "./compartment-runner-incremental"
import { validateHistorianOutput } from "./compartment-runner-validation";
import {
getProtectedTailStartOrdinal,
getRawSessionMessageCount,
getRawSessionMessageIdsThrough,
hasRawMessageProvider,
primeTailRawMessageCache,
readRawSessionMessageRange,
readRawSessionMessages,
readSessionChunk,
setRawMessageProvider,
withRawMessageProvider,
withRawSessionMessageCache,
} from "./read-session-chunk";
Expand Down Expand Up @@ -181,6 +184,120 @@ function appendOpenCodeMessage(
}
}

describe("raw message provider lifecycle", () => {
const provider = { readMessages: () => [], getMessageCount: () => 7 };

it("keeps the outer provider after nested synchronous return and throw", () => {
const sessionId = "provider-nested-sync";
const cleanup = setRawMessageProvider(sessionId, provider);
try {
expect(withRawMessageProvider(sessionId, provider, () => 42)).toBe(42);
expect(getRawSessionMessageCount(sessionId)).toBe(7);
expect(() =>
withRawMessageProvider(sessionId, provider, () => {
throw new Error("nested failure");
}),
).toThrow("nested failure");
expect(getRawSessionMessageCount(sessionId)).toBe(7);
} finally {
cleanup();
}
expect(hasRawMessageProvider(sessionId)).toBe(false);
});

it("keeps the outer provider after nested async settlement and rejection", async () => {
const sessionId = "provider-nested-async";
await withRawMessageProvider(sessionId, provider, async () => {
await withRawMessageProvider(sessionId, provider, async () => {
await Promise.resolve();
expect(getRawSessionMessageCount(sessionId)).toBe(7);
});
expect(getRawSessionMessageCount(sessionId)).toBe(7);
await expect(
withRawMessageProvider(sessionId, provider, async () => {
await Promise.resolve();
throw new Error("nested rejection");
}),
).rejects.toThrow("nested rejection");
expect(getRawSessionMessageCount(sessionId)).toBe(7);
});
expect(hasRawMessageProvider(sessionId)).toBe(false);
});

it("keeps a shared registration until every scope cleans up, exactly once", () => {
const sessionId = "provider-shared-cleanup";
const outerCleanup = setRawMessageProvider(sessionId, provider);
const innerCleanup = setRawMessageProvider(sessionId, provider);
try {
outerCleanup();
outerCleanup();
expect(getRawSessionMessageCount(sessionId)).toBe(7);
} finally {
innerCleanup();
outerCleanup();
}
expect(hasRawMessageProvider(sessionId)).toBe(false);
});

it.each(["old-first", "new-first"])(
"never overwrites or restores a replaced provider (%s cleanup)",
(order) => {
const sessionId = `provider-replaced-${order}`;
const oldCleanup = setRawMessageProvider(sessionId, provider);
const newCleanup = setRawMessageProvider(sessionId, {
readMessages: () => [],
getMessageCount: () => 11,
});
try {
if (order === "old-first") {
oldCleanup();
expect(getRawSessionMessageCount(sessionId)).toBe(11);
newCleanup();
} else {
newCleanup();
expect(hasRawMessageProvider(sessionId)).toBe(false);
oldCleanup();
}
expect(hasRawMessageProvider(sessionId)).toBe(false);
} finally {
newCleanup();
oldCleanup();
}
},
);

it("does not let an old scope remove a later registration of the same object", () => {
const sessionId = "provider-reregistered";
const oldCleanup = setRawMessageProvider(sessionId, provider);
const replacementCleanup = setRawMessageProvider(sessionId, { readMessages: () => [] });
const latestCleanup = setRawMessageProvider(sessionId, provider);
try {
oldCleanup();
replacementCleanup();
expect(getRawSessionMessageCount(sessionId)).toBe(7);
} finally {
latestCleanup();
replacementCleanup();
oldCleanup();
}
expect(hasRawMessageProvider(sessionId)).toBe(false);
});

it("owns registrations separately for each session", () => {
const firstCleanup = setRawMessageProvider("provider-session-one", provider);
const secondCleanup = setRawMessageProvider("provider-session-two", provider);
try {
firstCleanup();
expect(hasRawMessageProvider("provider-session-one")).toBe(false);
expect(getRawSessionMessageCount("provider-session-two")).toBe(7);
} finally {
secondCleanup();
firstCleanup();
}
expect(hasRawMessageProvider("provider-session-two")).toBe(false);
});
});

describe("readSessionChunk", () => {
it("reads raw OpenCode messages with stable ordinals and ids", () => {
useTempDataHome("read-session-chunk-");
Expand Down
Loading
Loading