diff --git a/.memory/perf-stream-accounting.md b/.memory/perf-stream-accounting.md new file mode 100644 index 000000000..85990d14e --- /dev/null +++ b/.memory/perf-stream-accounting.md @@ -0,0 +1,11 @@ +# Remote stream budget accounting — 2026-09-27 + +Implemented SDR-1 from the performance audit against `a9baa4aa3027893e5455043083465c34b4c8b4ac` on `feature/perf-stream-accounting`. + +`AidenRemoteStreamService` now accounts for retained events through per-stream byte totals and weakly cached exact UTF-8 serialized event sizes. A shared bounded metadata envelope includes stream fields, empty arrays, JSON separators and retained turn-index identities; adding event totals/commas yields the exact compact snapshot budget. Deletion/pruning/revocation/deferred terminal cleanup need no global delta ledger. Payloads are journal-owned clones. Restore counts the normalized events. State/timestamp are updated before enforcing the budget. + +Preserved existing trimming/eviction ordering, token/control events, sequence recovery, deferred terminal delivery and single-flight persistence. No persistence debounce, wire/schema, runtime admission, native-client or onboarding changes. iOS and Android SSE DTO/parser/recovery consumers were inspected; no native suites required for internal accounting. Existing array shifts and bounded metadata recomputation during pressure remain deliberate tradeoffs. + +Validation: final focused stream suite 62 passed; full `test:aiden-remote` 533 passed (19 peer-host, 507 Remote, 7 spike), one existing legacy-port-65535 occupancy skip. Exact independent serialized-byte oracles cover state/timestamp/Unicode/escaping, normalized restore, mutation ownership, turn-index eviction, pruning/revocation, per-stream count/byte trimming, ordinary terminal eviction and deferred-terminal pressure/delivery cleanup. TypeScript, scoped ESLint and diff check pass. Initial added test used Array.at against the project's ES2020 lib; replaced with ordinary indexing before final validation. + +Reproducible benchmark source and before/after details: `docs/performance/remote-stream-accounting.md` and companion `.mts`. Identical 64-appends fixture drops accounting-triggered full snapshot calls 64 → 0 with identical snapshot byte lengths at 0 and ~4.2 MiB retained history. Local timing is exploratory during concurrent builds, not an energy/whole-app claim. Persistence still serializes whole retained snapshots independently (SDR-3/4 deferred). Raw local evidence: `/tmp/aiden-performance-batch-1/stream-{before,after}.json`, `stream-tests-after.log`, `remote-tests-after.log`, `stream-typecheck.log`, `stream-lint.log`. diff --git a/docs/performance/remote-stream-accounting-benchmark.mts b/docs/performance/remote-stream-accounting-benchmark.mts new file mode 100644 index 000000000..443d12dac --- /dev/null +++ b/docs/performance/remote-stream-accounting-benchmark.mts @@ -0,0 +1,35 @@ +import { performance } from 'node:perf_hooks'; +import { writeFileSync } from 'node:fs'; +import { platform, release, arch } from 'node:os'; +import { pathToFileURL } from 'node:url'; +const root = process.argv[2]; +const output = process.argv[3]; +const { AidenRemoteStreamService } = await import(pathToFileURL(`${root}/main/services/aiden-remote-streams.ts`).href); +const rows = []; +for (const chunks of [0,32]) { + const timings=[]; let counters; + for(let trial=0;trial<6;trial++) { + const service = new AidenRemoteStreamService({now:()=>1000,cancel:()=>true,approve:()=>true}); + if(chunks) { + const retained = service.create('device','retained','old-chat','old-turn'); + for(let i=0;i{calls++;return original();}; + const start=performance.now(); + for(let i=0;i<64;i++) current.owner.send('chat:delta',{delta:'small synthetic delta'}); + const ms=performance.now()-start; + const count=calls; + service.snapshot=original; + const afterBytes=Buffer.byteLength(JSON.stringify(service.snapshot())); + if(trial>0)timings.push(ms); + counters={appends:64,snapshotCallsDuringAppends:count,beforeSnapshotBytes:beforeBytes,afterSnapshotBytes:afterBytes}; + } + rows.push({retainedChunks:chunks,warmups:1,samplesMs:timings,...counters}); +} +const result={scenario:'No persistence/subscribers; 64 synthetic appends with 0 or 32 retained 128-KiB chunks',runtime:process.version,platform:platform(),release:release(),arch:arch(),note:'Local microbenchmark during concurrent builds; wall time exploratory only. Snapshot-call counts are deterministic. No GPU/energy claim.',rows}; +writeFileSync(output,JSON.stringify(result,null,2)+'\n'); +console.log(JSON.stringify(result,null,2)); diff --git a/docs/performance/remote-stream-accounting.md b/docs/performance/remote-stream-accounting.md new file mode 100644 index 000000000..6e6cf27af --- /dev/null +++ b/docs/performance/remote-stream-accounting.md @@ -0,0 +1,28 @@ +# Remote stream budget accounting + +SDR-1 implementation against `a9baa4aa3027893e5455043083465c34b4c8b4ac`. + +The budget check no longer deep-clones and serializes retained event history on each append. Events are serialized once for their UTF-8 size; a weak cache supplies that same size when trimming. Newly appended payloads are cloned into journal ownership, and restore sizes are measured after legacy normalization. Existing per-stream event totals remain the accounting source. + +A shared snapshot envelope serializes only bounded metadata: at most 256 stream records and 1,024 retained turn identities. Empty event-array brackets are already in the envelope; cached event totals and the exact number of array commas complete the compact snapshot size. This deliberately avoids a global increment/decrement ledger spread across deletion, pruning, revocation and deferred-delivery cleanup. State and timestamp updates happen before enforcing the budget, so the check includes their final serialized sizes. + +Aggregate pressure keeps the existing oldest-terminal eviction and largest-journal trimming order, deferred terminal delivery protection, replay sequences and snapshot recovery. Each removal rechecks bounded metadata, without revisiting retained payloads. Array shifts and candidate sorting remain; this is not a deque rewrite. Persistence still clones snapshots and runs with the same single-flight coalescing, durability and error behavior. No schema, endpoint, native-client or onboarding change. + +## Reproducible measurement + +Run `node_modules/.bin/tsx docs/performance/remote-stream-accounting-benchmark.mts "$PWD" /tmp/stream-accounting.json` in either checkout. The same script ran before and after, with no persistence callback or subscribers, one warmup plus five measured trials, and 64 small appends. Runtime: Node v26.10.0, Darwin 27.0.0 arm64. Inputs are synthetic; no application state or provider traffic is used. + +| Prior retained chunks | Snapshot bytes before → after appends (identical both revisions) | Full snapshot calls before → after | Median elapsed before → after | +| --- | --- | --- | --- | +| 0 | 368 → 11,561 | 64 → 0 | 2.767 → 0.212 ms | +| 32 × 128 KiB | 4,200,168 → 4,211,361 | 64 → 0 | 93.925 → 0.223 ms | + +Before samples (ms): `[2.769, 2.763875, 2.637083, 2.766583, 2.864084]`; retained-history samples `[93.924875, 90.031208, 93.034, 104.148833, 96.301417]`. After samples: `[0.248875, 0.248666, 0.212458, 0.210042, 0.1905]`; retained-history samples `[0.260792, 0.222958, 0.221417, 0.216959, 0.262041]`. + +Snapshot-call counts and matching serialized lengths are deterministic evidence. Timing is exploratory: other builds were running on the host, this is a short synchronous fixture, and subscriber latency, allocation/RSS, event-loop delay, energy and GPU usage were not measured. No battery or whole-app speed claim follows. Persistence-related full-history work (SDR-3/4) remains outside this change. + +## Validation + +The existing registered stream suite gains independent `Buffer.byteLength(JSON.stringify(service.snapshot()), "utf8")` oracles for Unicode/escaping, timestamp widths, state changes, source/snapshot mutation isolation, normalization on restore, expiry, revocation, retained turn identities/index eviction, per-stream count/byte trimming, repeated aggregate pressure, terminal eviction and deferred terminal replay cleanup. Existing delivery, approval/question, cancellation and persistence tests remain in the suite. Under the budget, all 120 Unicode deltas are retained exactly; pressure retains contiguous newest sequences and the terminal event under the existing policy. + +Inspected iOS `AidenRemoteContract.swift` / `AidenRemoteClient.swift` and Android `AidenSSEParser.kt` / `AidenRemoteClient.kt`: DTOs, sequence recovery and terminal delivery rules are unchanged. Native builds/tests are not required for this internal accounting change. Focused stream, full Remote tests, TypeScript and scoped ESLint results are recorded in `.memory/perf-stream-accounting.md`. diff --git a/main/services/aiden-remote-streams.test.ts b/main/services/aiden-remote-streams.test.ts index 545e544e1..1a3c3a38d 100644 --- a/main/services/aiden-remote-streams.test.ts +++ b/main/services/aiden-remote-streams.test.ts @@ -9,6 +9,15 @@ import { } from "./aiden-remote-streams.js"; import { AidenRemoteServiceError } from "./aiden-remote-errors.js"; +function assertExactSnapshotBytes(service: AidenRemoteStreamService): void { + // Independent oracle includes UTF-8 encoding, all JSON punctuation, metadata, + // and identities retained after a journal disappears. + const actual = Buffer.byteLength(JSON.stringify(service.snapshot()), "utf8"); + const accounting = service as unknown as { snapshotBytes(): number }; + assert.equal(accounting.snapshotBytes(), actual); + assert.ok(actual <= 16 * 1_024 * 1_024); +} + function fixture() { let now = 1_000; const cancelled: string[] = []; @@ -73,6 +82,9 @@ test("remote stream forwards the renderer-safe chronological timeline without ra const event = app.service.snapshot().streams[0]?.events[1]; assert.equal(event?.type, "timeline"); assert.deepEqual(event?.payload, { timeline }); + timeline.steps[0]!.detail = "mutated caller value"; + assert.deepEqual(app.service.snapshot().streams[0]!.events[1]!.payload, event!.payload); + assertExactSnapshotBytes(app.service); assert.doesNotMatch(JSON.stringify(event), /cat |\.ssh/u); }); @@ -601,6 +613,10 @@ test("legacy label-only timeline journals load into the current safe timeline sh const legacy = app.service.snapshot(); legacy.streams[0]!.events[1]!.payload = { label: "Run command" }; + const restored = new AidenRemoteStreamService({ + now: () => 2_000, cancel: () => true, approve: () => true, snapshot: legacy, + }); + assertExactSnapshotBytes(restored); const normalized = normalizeAidenRemoteStreamSnapshot(legacy); const payload = normalized.streams[0]!.events[1]!.payload; assert.equal("label" in payload, false); @@ -668,16 +684,21 @@ test("aggregate stream journals stay within the durable snapshot budget", async approve: () => true, persist: async (snapshot) => { normalizeAidenRemoteStreamSnapshot(snapshot); }, }); + const old = service.create("device-1", "old", "old-chat", "old-turn"); + old.owner.send("chat:done", { chat: { messages: [{ id: "assistant", role: "assistant" }] } }); for (let streamIndex = 0; streamIndex < 3; streamIndex += 1) { const owner = service.create("device-1", `stream-${streamIndex}`, `chat-${streamIndex}`, `turn-${streamIndex}`); for (let index = 0; index < 35; index += 1) { owner.owner.send("chat:delta", { delta: `${streamIndex}:${index}:` + "x".repeat(199_990) }); + assertExactSnapshotBytes(service); } } await service.settlePersistence(); const snapshot = service.snapshot(); assert.doesNotThrow(() => normalizeAidenRemoteStreamSnapshot(snapshot)); assert.equal(Buffer.byteLength(JSON.stringify(snapshot), "utf8") <= 16 * 1_024 * 1_024, true); + assert.equal(snapshot.streams.some((stream) => stream.streamId === "old"), false); + assert.equal(service.turnIdFor("old-chat", "old"), "old-turn"); }); test("restart restores terminal journals and marks active work interrupted", () => { @@ -694,6 +715,7 @@ test("restart restores terminal journals and marks active work interrupted", () const status = restarted.status("device-1", "stream-1"); assert.equal(status.state, "interrupted"); assert.equal(status.lastSequence, 3); + assertExactSnapshotBytes(restarted); }); test("revocation closes only the selected device streams and approval expiry denies safely", async () => { @@ -889,6 +911,7 @@ test("disconnect and response errors release blocked SSE delivery", () => { assert.equal(client.output.length, 1); assert.equal(client.response.listenerCount("drain"), 0); assert.deepEqual(app.cancelled, []); + assertExactSnapshotBytes(app.service); } }); @@ -1024,6 +1047,7 @@ for (const settlement of ["drain", "timeout", "abort", "finish-abort", "finish-t (error: unknown) => error instanceof AidenRemoteServiceError && error.code === "not_found", ); assert.deepEqual(app.cancelled, []); + assertExactSnapshotBytes(app.service); }); } @@ -1691,3 +1715,87 @@ test("cancellation settles a pending question with a cancelled response", async (error: unknown) => (error as { code?: string }).code === "question_expired", ); }); + + +test("snapshot accounting follows Unicode payloads, state changes, pruning, restore and revocation", async () => { + const app = fixture(); + assertExactSnapshotBytes(app.service); + const device = 'device-雪😀"\\\n'; + const owner = app.service.create(device, "stream-1", "chat-1", "turn-1"); + const deltas = Array.from({ length: 120 }, (_, index) => `${index}:雪😀\n\t"\\`); + for (const delta of deltas) { + app.setNow(10 ** (1 + delta.length % 9)); + owner.owner.send("chat:delta", { delta }); + assertExactSnapshotBytes(app.service); + } + const retained = app.service.snapshot().streams[0]!; + assert.equal(retained.events.filter((event) => event.type === "text_delta") + .map((event) => event.payload.text).join(""), deltas.join("")); + // A consumer may mutate an exported snapshot without changing the journal. + retained.events[1]!.payload.text = "external mutation"; + assertExactSnapshotBytes(app.service); + assert.equal(app.service.snapshot().streams[0]!.events[1]!.payload.text, deltas[0]); + owner.owner.send("chat:approval", { approvalId: "approval-1", summary: "雪😀\nRun?" }); + assert.equal(app.service.status(device, "stream-1").state, "waiting_for_approval"); + assertExactSnapshotBytes(app.service); + owner.owner.send("chat:error", { cancelled: true }); + assertExactSnapshotBytes(app.service); + const saved = app.service.snapshot(); + const restored = new AidenRemoteStreamService({ + now: () => 1_000, cancel: () => true, approve: () => true, snapshot: saved, + }); + saved.streams[0]!.events[1]!.payload.text = "mutated restore source"; + assertExactSnapshotBytes(restored); + assert.equal(restored.snapshot().streams[0]!.events[1]!.payload.text, deltas[0]); + await restored.revokeDevice(device); + assertExactSnapshotBytes(restored); + assert.equal(restored.snapshot().streams.length, 0); + assert.equal(restored.turnIdFor("chat-1", "stream-1"), "turn-1"); + app.setNow(10 ** 12); + app.service.create("other", "stream-2", "chat-2", "turn-2"); + assertExactSnapshotBytes(app.service); + assert.equal(app.service.snapshot().streams.length, 1); + assert.equal(app.service.snapshot().turnIndex!.length, 2); +}); + +test("snapshot accounting covers escaped identities and bounded turn index eviction", async () => { + const app = fixture(); + const unusual = '雪😀"\\\n'; + app.service.create(unusual, unusual, unusual, unusual); + assertExactSnapshotBytes(app.service); + await app.service.revokeDevice(unusual); + for (let index = 0; index < 1_030; index++) { + const owner = app.service.create("device", `s-${index}`, `c-${index}`, `t-${index}`); + owner.owner.send("chat:error", {}); + app.setNow(1_000 + (index + 1) * 24 * 60 * 60 * 1_000); + } + assertExactSnapshotBytes(app.service); + assert.equal(app.service.snapshot().turnIndex!.length, 1_024); + assert.equal(app.service.turnIdFor("c-0", "s-0"), undefined); + assert.equal(app.service.turnIdFor("c-1029", "s-1029"), "t-1029"); +}); + +test("event count and byte trimming keep exact accounting and contiguous newest replay", () => { + const app = fixture(); + const owner = app.service.create("device", "stream", "chat", "turn"); + for (let index = 0; index < 4_100; index++) { + owner.owner.send("chat:delta", { delta: `chunk-${index}` }); + } + assertExactSnapshotBytes(app.service); + let events = app.service.snapshot().streams[0]!.events; + assert.equal(events.length, 4_096); + assert.equal(events[0]!.sequence, 6); + assert.equal(events[events.length - 1]!.payload.text, "chunk-4099"); + for (let index = 0; index < 50; index++) { + owner.owner.send("chat:delta", { delta: "雪".repeat(100_000) }); + } + owner.owner.send("chat:done", { chat: { messages: [{ id: "assistant", role: "assistant" }] } }); + assertExactSnapshotBytes(app.service); + events = app.service.snapshot().streams[0]!.events; + assert.ok(events.length < 50, "the per-stream byte limit prunes old events"); + assert.equal(events[events.length - 1]!.type, "done"); + assert.equal(events[events.length - 1]!.terminal, true); + for (let index = 1; index < events.length; index++) { + assert.equal(events[index]!.sequence, events[index - 1]!.sequence + 1); + } +}); diff --git a/main/services/aiden-remote-streams.ts b/main/services/aiden-remote-streams.ts index 7ce6f1a35..f1a50f04b 100644 --- a/main/services/aiden-remote-streams.ts +++ b/main/services/aiden-remote-streams.ts @@ -430,6 +430,7 @@ function sseFrame(event: AidenRemoteStreamEvent): string { export class AidenRemoteStreamService { private readonly streams = new Map(); + private readonly eventSizes = new WeakMap(); private readonly turnIndex = new Map(); private readonly approvals = new Map(); private readonly questions = new Map(); @@ -561,7 +562,7 @@ export class AidenRemoteStreamService { const base: Omit = { ...saved, eventBytes: saved.events.reduce( - (total, event) => total + Buffer.byteLength(JSON.stringify(event), "utf8"), + (total, event) => total + this.eventSize(event), 0, ), subscribers: new Set(), @@ -594,7 +595,7 @@ export class AidenRemoteStreamService { this.turnIndex.set(streamId, { chatId, turnId }); } - snapshot(): AidenRemoteStreamSnapshot { + private snapshotEnvelope(): AidenRemoteStreamSnapshot { return { version: 1, streams: [...this.streams.values()].map((stream) => ({ @@ -604,7 +605,7 @@ export class AidenRemoteStreamService { deviceId: stream.deviceId, state: stream.state, updatedAt: stream.updatedAt, - events: structuredClone(stream.events), + events: [], })), turnIndex: [...this.turnIndex.entries()].map(([streamId, entry]) => ({ streamId, @@ -614,6 +615,35 @@ export class AidenRemoteStreamService { }; } + snapshot(): AidenRemoteStreamSnapshot { + const snapshot = this.snapshotEnvelope(); + for (const stream of snapshot.streams) { + stream.events = structuredClone(this.streams.get(stream.streamId)!.events); + } + return snapshot; + } + + private eventSize(event: AidenRemoteStreamEvent): number { + let bytes = this.eventSizes.get(event); + if (bytes === undefined) { + bytes = Buffer.byteLength(JSON.stringify(event), "utf8"); + this.eventSizes.set(event, bytes); + } + return bytes; + } + + private snapshotBytes(): number { + // Serialize only bounded metadata (256 streams / 1,024 turn identities). + // The empty arrays already include brackets; add event bytes and commas. + // Rebuilding this envelope keeps deletion, expiry and delivery cleanup from + // having to maintain a second, fragile aggregate mutation ledger. + let bytes = Buffer.byteLength(JSON.stringify(this.snapshotEnvelope()), "utf8"); + for (const stream of this.streams.values()) { + bytes += stream.eventBytes + Math.max(0, stream.events.length - 1); + } + return bytes; + } + private persist(): void { if (!this.options.persist) return; this.persistDirty = true; @@ -853,9 +883,11 @@ export class AidenRemoteStreamService { timestamp: new Date(this.options.now()).toISOString(), type, terminal: isTerminal, - payload, + // Journal events own their payload: callers must not invalidate cached + // sizes (or change already-published replay) by mutating nested values. + payload: structuredClone(payload), }; - const bytes = Buffer.byteLength(JSON.stringify(event), "utf8"); + const bytes = this.eventSize(event); stream.events.push(event); stream.eventBytes += bytes; while ( @@ -863,11 +895,11 @@ export class AidenRemoteStreamService { (stream.eventBytes > MAX_STREAM_EVENT_BYTES && stream.events.length > 1) ) { const removed = stream.events.shift(); - if (removed) stream.eventBytes -= Buffer.byteLength(JSON.stringify(removed), "utf8"); + if (removed) stream.eventBytes -= this.eventSize(removed); } - this.enforceAggregateBudget(stream.streamId); stream.state = state ?? stream.state; stream.updatedAt = this.options.now(); + this.enforceAggregateBudget(stream.streamId); for (const subscriber of [...stream.subscribers]) subscriber.flush(); if (isTerminal) { stream.owner.invalidate(); @@ -892,8 +924,7 @@ export class AidenRemoteStreamService { } private enforceAggregateBudget(currentStreamId: string): void { - const snapshotBytes = () => Buffer.byteLength(JSON.stringify(this.snapshot()), "utf8"); - if (snapshotBytes() <= MAX_AIDEN_REMOTE_STREAM_SNAPSHOT_BYTES) return; + if (this.snapshotBytes() <= MAX_AIDEN_REMOTE_STREAM_SNAPSHOT_BYTES) return; const terminalStreams = [...this.streams.values()] .filter((entry) => entry.streamId !== currentStreamId && terminal(entry.state)) .sort((left, right) => left.updatedAt - right.updatedAt); @@ -909,16 +940,16 @@ export class AidenRemoteStreamService { } entry.owner.invalidate(); this.streams.delete(entry.streamId); - if (snapshotBytes() <= MAX_AIDEN_REMOTE_STREAM_SNAPSHOT_BYTES) return; + if (this.snapshotBytes() <= MAX_AIDEN_REMOTE_STREAM_SNAPSHOT_BYTES) return; } - while (snapshotBytes() > MAX_AIDEN_REMOTE_STREAM_SNAPSHOT_BYTES) { + while (this.snapshotBytes() > MAX_AIDEN_REMOTE_STREAM_SNAPSHOT_BYTES) { const candidate = [...this.streams.values()] .filter((entry) => entry.events.length > 0) .sort((left, right) => right.eventBytes - left.eventBytes)[0]; if (!candidate) break; if (candidate.events.length > 1) { const removed = candidate.events.shift(); - if (removed) candidate.eventBytes -= Buffer.byteLength(JSON.stringify(removed), "utf8"); + if (removed) candidate.eventBytes -= this.eventSize(removed); continue; } const retained = candidate.events[0]!; @@ -933,7 +964,7 @@ export class AidenRemoteStreamService { }, }; candidate.events[0] = replacement; - candidate.eventBytes = Buffer.byteLength(JSON.stringify(replacement), "utf8"); + candidate.eventBytes = this.eventSize(replacement); } }