test(agent): strengthened consistency and persistence logic in integration tests
- Added regression test in `remote-compaction` to verify that concurrent v2 compaction preparation correctly reuses preserved history and avoids redundant re-expansion. - Added mock-backend verification in `session-storage` to ensure that failed atomic title updates do not rollback newer optimistic state. - Updated `sql-session-storage` expectations to account for the preserved fixed-width title slot header in session files. - Expanded `remote-compaction` fetch header validation to include `x-client-request-id` assertion.
This commit is contained in:
@@ -4,6 +4,8 @@ import {
|
||||
compact,
|
||||
createFileOps,
|
||||
DEFAULT_COMPACTION_SETTINGS,
|
||||
prepareCompaction,
|
||||
type SessionEntry,
|
||||
} from "@oh-my-pi/pi-agent-core/compaction";
|
||||
import {
|
||||
buildCompactionV2Request,
|
||||
@@ -297,13 +299,19 @@ describe("requestCompactionV2Streaming", () => {
|
||||
);
|
||||
let requestBody: { model: string; input: Array<Record<string, unknown>>; prompt_cache_key?: string } | undefined;
|
||||
let sessionHeader: string | undefined;
|
||||
let clientRequestHeader: string | undefined;
|
||||
let legacySessionHeader: string | undefined;
|
||||
const fetchMock: FetchImpl = async (input, init) => {
|
||||
expect(String(input)).toBe("https://compact.example/v1/responses");
|
||||
if (!init?.headers || init.headers instanceof Headers || Array.isArray(init.headers)) {
|
||||
throw new Error("Expected V2 compaction to send headers as a plain object");
|
||||
}
|
||||
const rawSessionHeader = init.headers["session-id"];
|
||||
const rawSessionHeader = init.headers.session_id;
|
||||
const rawClientRequestHeader = init.headers["x-client-request-id"];
|
||||
const rawLegacySessionHeader = init.headers["session-id"];
|
||||
sessionHeader = typeof rawSessionHeader === "string" ? rawSessionHeader : undefined;
|
||||
clientRequestHeader = typeof rawClientRequestHeader === "string" ? rawClientRequestHeader : undefined;
|
||||
legacySessionHeader = typeof rawLegacySessionHeader === "string" ? rawLegacySessionHeader : undefined;
|
||||
requestBody = JSON.parse(String(init.body)) as {
|
||||
model: string;
|
||||
input: Array<Record<string, unknown>>;
|
||||
@@ -335,6 +343,8 @@ describe("requestCompactionV2Streaming", () => {
|
||||
const result = await requestCompactionV2Streaming(model, "test-key", request, undefined, { fetch: fetchMock });
|
||||
|
||||
expect(sessionHeader).toBe("session-1");
|
||||
expect(clientRequestHeader).toBe("session-1");
|
||||
expect(legacySessionHeader).toBeUndefined();
|
||||
expect(requestBody?.model).toBe("gpt-5-compact");
|
||||
expect(requestBody?.prompt_cache_key).toBe("cache-1");
|
||||
expect(requestBody?.input[requestBody.input.length - 1]).toEqual({ type: "compaction_trigger" });
|
||||
@@ -653,9 +663,88 @@ describe("compact() remote compaction failure handling", () => {
|
||||
const remote = getCompactionV2PreserveData(result.preserveData);
|
||||
expect(remote?.usedTokens).toBe(55);
|
||||
expect(remote?.replacementHistory.at(-1)).toEqual(compactionItem);
|
||||
expect(result.summary).toContain("Remote compaction preserved provider-native history");
|
||||
expect(completeSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
test("re-expands a prior V2 compaction's originals when no candidate can reuse the replay", async () => {
|
||||
vi.spyOn(ai, "completeSimple").mockResolvedValue(localSummaryMessage("re-expanded local summary"));
|
||||
const compactionItem = { type: "compaction", encrypted_content: "enc_v2" };
|
||||
const v2Model = makeOpenAiModel({
|
||||
remoteCompaction: {
|
||||
enabled: true,
|
||||
v2StreamingEnabled: true,
|
||||
v2Endpoint: "https://compact.example/v1/responses",
|
||||
},
|
||||
});
|
||||
// Produce a real V2 preserve payload (opaque placeholder summary, provider "openai").
|
||||
const v2Preparation = makePreparation();
|
||||
v2Preparation.messagesToSummarize = [{ role: "user", content: "ORIGINAL ALPHA port 4242", timestamp: 1 }];
|
||||
v2Preparation.recentMessages = [{ role: "user", content: "turn after", timestamp: 2 }];
|
||||
v2Preparation.settings = { ...v2Preparation.settings, remoteStreamingV2Enabled: true };
|
||||
const v2Result = await compact(v2Preparation, v2Model, "k", undefined, undefined, {
|
||||
fetch: async () =>
|
||||
sseResponse([
|
||||
{ type: "response.output_item.done", output_index: 0, item: compactionItem },
|
||||
{
|
||||
type: "response.completed",
|
||||
response: { usage: { input_tokens: 9, output_tokens: 1, total_tokens: 10 } },
|
||||
},
|
||||
]),
|
||||
});
|
||||
// V2 success persists only the opaque placeholder — no second local summarization round.
|
||||
expect(v2Result.summary).toContain("Remote compaction preserved provider-native history");
|
||||
|
||||
// Session branch after that V2 compaction: originals + compaction boundary + new turns.
|
||||
const ts = (n: number) => new Date(n).toISOString();
|
||||
const entries: SessionEntry[] = [
|
||||
{
|
||||
type: "message",
|
||||
id: "m1",
|
||||
parentId: null,
|
||||
timestamp: ts(1),
|
||||
message: { role: "user", content: "ORIGINAL ALPHA port 4242", timestamp: 1 },
|
||||
},
|
||||
{
|
||||
type: "compaction",
|
||||
id: "c1",
|
||||
parentId: "m1",
|
||||
timestamp: ts(2),
|
||||
summary: v2Result.summary,
|
||||
firstKeptEntryId: "m1",
|
||||
tokensBefore: 100_000,
|
||||
preserveData: v2Result.preserveData,
|
||||
},
|
||||
{
|
||||
type: "message",
|
||||
id: "m2",
|
||||
parentId: "c1",
|
||||
timestamp: ts(3),
|
||||
message: { role: "user", content: "second turn", timestamp: 3 },
|
||||
},
|
||||
{
|
||||
type: "message",
|
||||
id: "m3",
|
||||
parentId: "m2",
|
||||
timestamp: ts(4),
|
||||
message: { role: "user", content: "third turn", timestamp: 4 },
|
||||
},
|
||||
];
|
||||
const baseSettings = { ...DEFAULT_COMPACTION_SETTINGS, keepRecentTokens: 1 };
|
||||
|
||||
// Remote disabled → the V2 replay is unusable → re-expand the pre-V2 original.
|
||||
const reexpanded = prepareCompaction(entries, { ...baseSettings, remoteEnabled: false }, [v2Model]);
|
||||
expect(reexpanded).toBeDefined();
|
||||
const reexpandedText = JSON.stringify(reexpanded?.messagesToSummarize ?? []);
|
||||
expect(reexpandedText).toContain("ORIGINAL ALPHA port 4242");
|
||||
|
||||
// Remote + V2 still enabled, same provider → reuse the replay, don't re-summarize originals.
|
||||
const reused = prepareCompaction(entries, { ...baseSettings, remoteStreamingV2Enabled: true }, [v2Model]);
|
||||
expect(reused).toBeDefined();
|
||||
const reusedText = JSON.stringify(reused?.messagesToSummarize ?? []);
|
||||
expect(reusedText).not.toContain("ORIGINAL ALPHA port 4242");
|
||||
});
|
||||
|
||||
test("user abort during the remote compact request rejects without falling back to local summarization", async () => {
|
||||
// Contract: Esc is a cancellation, not a remote failure. Before the fix
|
||||
// the AbortError was swallowed by the fallback catch and compaction kept
|
||||
|
||||
@@ -3,9 +3,90 @@ import * as fs from "node:fs";
|
||||
import * as fsp from "node:fs/promises";
|
||||
import * as os from "node:os";
|
||||
import * as path from "node:path";
|
||||
import {
|
||||
IndexedSessionStorage,
|
||||
type SessionStorageBackend,
|
||||
type SessionStorageIndexEntry,
|
||||
} from "@oh-my-pi/pi-coding-agent/session/indexed-session-storage";
|
||||
import { FileSessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
|
||||
import { serializeTitleSlot } from "@oh-my-pi/pi-coding-agent/session/session-title-slot";
|
||||
import { type SessionTitleUpdate, serializeTitleSlot } from "@oh-my-pi/pi-coding-agent/session/session-title-slot";
|
||||
|
||||
class ControlledTitleUpdateBackend implements SessionStorageBackend {
|
||||
readonly #sessionPath: string;
|
||||
readonly #initialEntry: SessionStorageIndexEntry;
|
||||
#content: string;
|
||||
#firstUpdate: PromiseWithResolvers<void> | undefined;
|
||||
#updateCount = 0;
|
||||
|
||||
constructor(sessionPath: string, content: string) {
|
||||
this.#sessionPath = sessionPath;
|
||||
this.#content = content;
|
||||
this.#initialEntry = {
|
||||
path: sessionPath,
|
||||
size: content.length,
|
||||
mtimeMs: 1,
|
||||
title: "Old",
|
||||
titleSource: "auto",
|
||||
titleUpdatedAt: "t0",
|
||||
};
|
||||
}
|
||||
|
||||
init(): Promise<void> {
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
loadIndex(): Promise<Iterable<SessionStorageIndexEntry>> {
|
||||
return Promise.resolve([this.#initialEntry]);
|
||||
}
|
||||
|
||||
readFull(path: string): Promise<string | null> {
|
||||
return Promise.resolve(path === this.#sessionPath ? this.#content : null);
|
||||
}
|
||||
|
||||
readSlices(path: string, prefixBytes: number, suffixBytes: number): Promise<[string, string]> {
|
||||
if (path !== this.#sessionPath) return Promise.resolve(["", ""]);
|
||||
const suffix = suffixBytes > 0 ? this.#content.slice(-suffixBytes) : "";
|
||||
return Promise.resolve([this.#content.slice(0, prefixBytes), suffix]);
|
||||
}
|
||||
|
||||
writeFull(_path: string, content: string, _mtimeMs: number, _title?: SessionTitleUpdate): Promise<void> {
|
||||
this.#content = content;
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
append(_path: string, line: string, _mtimeMs: number): Promise<void> {
|
||||
this.#content += line;
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
updateSessionTitle(_path: string, _title: SessionTitleUpdate, _mtimeMs: number): Promise<void> {
|
||||
this.#updateCount++;
|
||||
if (this.#updateCount === 1) {
|
||||
this.#firstUpdate = Promise.withResolvers<void>();
|
||||
return this.#firstUpdate.promise;
|
||||
}
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
truncate(_path: string, _mtimeMs: number): Promise<void> {
|
||||
this.#content = "";
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
remove(_paths: string[]): Promise<void> {
|
||||
this.#content = "";
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
move(_src: string, _dst: string, _mtimeMs: number): Promise<void> {
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
rejectFirstUpdate(error: Error): void {
|
||||
if (!this.#firstUpdate) throw new Error("First title update has not started");
|
||||
this.#firstUpdate.reject(error);
|
||||
}
|
||||
}
|
||||
describe("FileSessionStorage.deleteSessionWithArtifacts", () => {
|
||||
let tempDir: string;
|
||||
let storage: FileSessionStorage;
|
||||
@@ -124,3 +205,29 @@ describe("FileSessionStorage.updateSessionTitle", () => {
|
||||
).rejects.toThrow(/ENOENT|no such file/i);
|
||||
});
|
||||
});
|
||||
|
||||
describe("IndexedSessionStorage.updateSessionTitle", () => {
|
||||
it("does not roll a newer optimistic title back when an older backend write fails", async () => {
|
||||
const sessionPath = "/sessions/session.jsonl";
|
||||
const content = `${serializeTitleSlot({ title: "Old", source: "auto", updatedAt: "t0" })}${JSON.stringify({
|
||||
type: "session",
|
||||
id: "session-id",
|
||||
timestamp: "t0",
|
||||
cwd: "/cwd",
|
||||
})}\n`;
|
||||
const backend = new ControlledTitleUpdateBackend(sessionPath, content);
|
||||
const storage = new IndexedSessionStorage(backend);
|
||||
await storage.initialize();
|
||||
|
||||
const first = storage.updateSessionTitle(sessionPath, { title: "First", source: "auto", updatedAt: "t1" });
|
||||
const second = storage.updateSessionTitle(sessionPath, { title: "Second", source: "user", updatedAt: "t2" });
|
||||
for (let i = 0; i < 10; i++) await Promise.resolve();
|
||||
|
||||
backend.rejectFirstUpdate(new Error("first title write failed"));
|
||||
await expect(first).rejects.toThrow("first title write failed");
|
||||
await expect(second).resolves.toBeUndefined();
|
||||
|
||||
const [slotLine] = (await storage.readText(sessionPath)).split("\n");
|
||||
expect(JSON.parse(slotLine)).toMatchObject({ type: "title", title: "Second", source: "user", updatedAt: "t2" });
|
||||
});
|
||||
});
|
||||
|
||||
@@ -58,8 +58,11 @@ describe("SessionManager + SqlSessionStorage (SQLite)", () => {
|
||||
])) as Array<{ content: string }>;
|
||||
expect(rows).toHaveLength(1);
|
||||
const lines = rows[0].content.trim().split("\n");
|
||||
expect(lines.length).toBeGreaterThanOrEqual(2);
|
||||
const header = JSON.parse(lines[0]);
|
||||
expect(lines.length).toBeGreaterThanOrEqual(3);
|
||||
// The fixed-width title slot is always the first physical line; the session header follows.
|
||||
const slot = JSON.parse(lines[0]);
|
||||
expect(slot.type).toBe("title");
|
||||
const header = JSON.parse(lines[1]);
|
||||
expect(header.type).toBe("session");
|
||||
const msg = JSON.parse(lines[lines.length - 1]);
|
||||
expect(msg.type).toBe("message");
|
||||
|
||||
Reference in New Issue
Block a user