diff --git a/lib/sessionRepository.ts b/lib/sessionRepository.ts index 05308c44..10274ae1 100644 --- a/lib/sessionRepository.ts +++ b/lib/sessionRepository.ts @@ -827,10 +827,13 @@ function cloudBodyBytes(body: CloudPutBody): number { * #515 envelope/Blob path, where the transcript object is ferried **client→Blob** * and never through a Function body. */ +import { assertCompleteForCloudPut } from './viewportFirewall'; + export function trimForCloudPut( snapshot: SessionSnapshot, maxBytes: number = HARNESS_SESSION_MAX_FUNCTION_BODY_BYTES, ): CloudPutBody { + assertCompleteForCloudPut(snapshot); const messages = snapshot.messages.map((m) => ({ id: m.id, role: m.role, @@ -1306,6 +1309,8 @@ export function createHttpSessionRepository( // Canonical identity (parent #415): a snapshot never stores under a different // resource id than its own persisted id. if (snapshot.id !== id) return; + // Plan #960 firewall — never queue a partial view snapshot for cloud persist. + if (snapshot.historyComplete === false) return; const c = channel(id); c.pending = snapshot; void drain(id, c); diff --git a/lib/sessionStore.ts b/lib/sessionStore.ts index acdeb2c9..69e8bc2f 100644 --- a/lib/sessionStore.ts +++ b/lib/sessionStore.ts @@ -141,6 +141,14 @@ export type SessionSnapshot = { * host-observable without a protocol bump (documented residual on #815). */ queue?: string[]; + /** + * Plan #960 — explicit view marker: this snapshot is a partial (view-only) window + * rebuilt from a viewport snapshot (`historyComplete: false`). Cloud transcript + * persist (put/mint/upload/flatten) MUST skip it; metadata-only preferences + * (model/reasoning/queue mirror) ride the narrow PATCH channel. `undefined`/`true` + * = a fully restored session. Never inferred from length. + */ + historyComplete?: boolean; }; import { diff --git a/lib/sessions/sessionStore.ts b/lib/sessions/sessionStore.ts index dc19de54..1857d423 100644 --- a/lib/sessions/sessionStore.ts +++ b/lib/sessions/sessionStore.ts @@ -101,6 +101,8 @@ export const RESERVED_META_KEYS = [ 'workingNotes', 'freshnessReminderPointer', 'compactionPointer', + /** Plan #960 reserved carrier — JSON-array string (each entry a host-known prompt). `[]` means an explicit empty tombstone, not "unset". */ + 'queueMirror', ] as const; export type HarnessSessionMetaKey = (typeof RESERVED_META_KEYS)[number]; @@ -636,12 +638,15 @@ export function validateMeta(value: unknown): SessionStoreResult({ + start(controller) { + const enc = new TextEncoder(); + for (const e of records) controller.enqueue(enc.encode(e)); + controller.close(); + }, + }); + return new Response(body, { + status: 200, + headers: { 'Content-Type': 'text/event-stream; charset=utf-8', ...(viewHeader ? { 'x-viewport-version': '1' } : {}) }, + }); +} + +describe('viewportAttachTurnStream', () => { + afterEach(() => vi.unstubAllGlobals()); + + it('GETs cold hydrate=tail (no startIndex) and folds the decoded records', async () => { + const fetchMock = vi.fn(() => + Promise.resolve( + sseResponse([ + 'event: viewport_state\ndata: {"type":"viewport_state","version":1,"runId":"wr_1","phase":"recovering","status":"running"}\n\n', + ]), + ), + ); + vi.stubGlobal('fetch', fetchMock); + const events: unknown[] = []; + const result = await viewportAttachTurnStream('wr_1', { + sessionId: 's_1', + onEvent: (rec) => { events.push(rec); }, + }); + expect(result.ok).toBe(true); + if ('status' in result) expect(result.status).toBe(200); + const url = (fetchMock.mock.calls as unknown as Array<[RequestInfo | URL, RequestInit]>)[0]?.[0] as RequestInfo | URL; + expect(String(url as string).includes('viewportVersion=1&hydrate=tail')).toBe(true); + }); + + it('GETs indexed startIndex=N when explicitly supplied', async () => { + const fetchMock = vi.fn(() => + Promise.resolve( + sseResponse([ + 'event: viewport_state\ndata: {"type":"viewport_state","version":1,"runId":"wr_1","phase":"recovering","status":"running"}\n\n', + ]), + ), + ); + vi.stubGlobal('fetch', fetchMock); + const result = await viewportAttachTurnStream('wr_1', { sessionId: 's_1', startIndex: 7 }); + expect(result.ok).toBe(true); + const url = (fetchMock.mock.calls as unknown as Array<[RequestInfo | URL, RequestInit]>)[0]?.[0] as RequestInfo | URL; + expect(String(url as string).includes('startIndex=7')).toBe(true); + }); + + it('rejects an unknown/wrong viewportVersion (never legacy consume)', async () => { + const fetchMock = vi.fn(() => Promise.resolve(new Response('bad', { status: 400 }))); + vi.stubGlobal('fetch', fetchMock); + const result = await viewportAttachTurnStream('wr_1', { sessionId: 's_1' }); + expect(result.ok).toBe(false); + if ('status' in result) expect(result.status).toBe(400); + }); + + it('a legacy 200 SSE body without x-viewport-version is a client error', async () => { + const fetchMock = vi.fn(() => + Promise.resolve(sseResponse(['data: {"type":"done","text":"x"}'], false)), + ); + vi.stubGlobal('fetch', fetchMock); + const result = await viewportAttachTurnStream('wr_1', { sessionId: 's_1' }); + expect(result.ok).toBe(false); + if ('error' in result) expect(result.error).toMatch(/Viewport negotiation not accepted/); + }); +}); \ No newline at end of file diff --git a/lib/viewportAttach.ts b/lib/viewportAttach.ts new file mode 100644 index 00000000..d9716dff --- /dev/null +++ b/lib/viewportAttach.ts @@ -0,0 +1,167 @@ +import { normalizePrompt } from './chatApi'; +import { + AGENT_STREAM_ACCEPT, + type AgentStreamEvent, +} from './agent/agentStream'; +import type { AgentFailure } from './agentApi'; +import { sanitizeUsageSummary, type UsageSummary } from './agent/usageSummary'; +import { + isRedisSafeOpaqueId, + sanitizeResolvedProvider, + sanitizeTurnRunId, + sanitizeTurnStreamCursor, +} from './sessionCloudCaps'; +import { + ViewportStreamDecoder, + type ViewportRecord, + type ViewportSnapshot, +} from './viewportStreamProtocol'; + +function parseTurnRunId(res: Response): string | undefined { + const raw = res.headers.get('x-workflow-run-id'); + if (!raw) return undefined; + const trimmed = raw.trim(); + return trimmed || undefined; +} + +function isAbortError(err: unknown): boolean { + return err instanceof Error && err.name === 'AbortError'; +} + +export type ViewportNegotiation = 'cold' | 'indexed'; + +export type ViewportAttachInit = { + sessionId: string; + startIndex?: number; + negotiation?: ViewportNegotiation; + onTurnStarted?: (info: { turnRunId: string }) => Promise | void; + onEvent?: (event: ViewportRecord) => Promise | void; + signal?: AbortSignal; +}; + +export type ViewportAttachResult = + | AgentFailure + | { ok: true; text: string; turnRunId?: string; cursor?: number; snapshot?: ViewportSnapshot; usage?: UsageSummary; resolvedProvider?: string; } + | { ok: true; text: string; turnRunId?: string; cursor?: number; snapshot?: ViewportSnapshot; usage?: UsageSummary; resolvedProvider?: string; completedAt: 'done' }; + +/** Read a viewport body and dispatch each decoded record to opts.onEvent. */ +export async function readViewportBody( + body: ReadableStream, + runId: string, + opts: ViewportAttachInit, +): Promise<{ ok: boolean; error?: string; status?: number; turnRunId?: string; cursor?: number; snapshot?: ViewportSnapshot; text?: string; usage?: UsageSummary; resolvedProvider?: string }> { + const decoder = new ViewportStreamDecoder(runId); + let lastError: AgentFailure | undefined; + let lastTurnEventNextIndex: number | undefined; + let sawViewportEnd = false; + let doneEvent: Extract | undefined; + let sawDoneText: string | undefined; + let streamUsage: UsageSummary | undefined; + let streamProvider: string | undefined; + const reader = body.getReader(); + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + for (const rec of decoder.push(value)) { + if (opts.onEvent) { + await Promise.resolve(opts.onEvent(rec)); + } + if (rec.type === 'turn_event' && !rec.skipped && rec.event) { + const ev = rec.event; + if (ev.type === 'done') { + if (doneEvent === undefined) doneEvent = ev; + if (typeof ev.text === 'string') sawDoneText = ev.text; + } else if (ev.type === 'usage') { + streamUsage = sanitizeUsageSummary(ev.usage) ?? streamUsage; + } else if (ev.type === 'provider') { + streamProvider = sanitizeResolvedProvider(ev.provider) ?? streamProvider; + } + lastTurnEventNextIndex = rec.nextIndex; + } else if (rec.type === 'viewport_end') { + sawViewportEnd = true; + } else if (rec.type === 'viewport_error') { + lastError = { ok: false, error: `[viewport_error ${rec.code}] Viewport attach failed.` }; + } + if (sawViewportEnd || lastError) break; + } + } + if (!sawViewportEnd && !lastError) { + for (const rec of decoder.push(new Uint8Array(), true)) { + if (opts.onEvent) { + await Promise.resolve(opts.onEvent(rec)); + } + } + } + if (lastError) return lastError; + if (!sawViewportEnd) { + return { + ok: true, + text: sawDoneText, + cursor: lastTurnEventNextIndex, + ...(doneEvent !== undefined ? { doneEvent } : {}), + ...(streamUsage ? { usage: streamUsage } : {}), + ...(streamProvider ? { resolvedProvider: streamProvider } : {}), + ...(decoder.lastSnapshot !== undefined ? { snapshot: decoder.lastSnapshot } : {}), + }; + } + if (doneEvent !== undefined && sawDoneText !== undefined) { + return { + ok: true, + text: sawDoneText, + ...(streamUsage ? { usage: streamUsage } : {}), + ...(streamProvider ? { resolvedProvider: streamProvider } : {}), + }; + } + return { ok: true }; +} + +/** GET `/api/turns/:runId/stream` with negotiated viewport transport. */ +export async function viewportAttachTurnStream( + runId: string, + opts: ViewportAttachInit, +): Promise { + const cleanRunId = sanitizeTurnRunId(runId); + if (!cleanRunId) return { ok: false, status: 400, error: 'Invalid run id' }; + if (!normalizePrompt(opts.sessionId ?? '') || !isRedisSafeOpaqueId(opts.sessionId)) { + return { ok: false, status: 400, error: 'Invalid session id' }; + } + const params = new URLSearchParams(); + params.set('sessionId', opts.sessionId); + params.set('viewportVersion', '1'); + if (opts.negotiation === 'indexed' || opts.startIndex !== undefined) { + if (opts.startIndex === undefined) { + return { ok: false, status: 400, error: 'Indexed attach needs startIndex.' }; + } + const startIndex = sanitizeTurnStreamCursor(opts.startIndex); + if (startIndex === undefined) { + return { ok: false, status: 400, error: 'Invalid startIndex.' }; + } + params.set('startIndex', String(startIndex)); + } else { + params.set('hydrate', 'tail'); + } + const path = `/api/turns/${encodeURIComponent(cleanRunId)}/stream?${params.toString()}`; + let res: Response; + try { + res = await fetch(path, { + method: 'GET', + headers: { Accept: AGENT_STREAM_ACCEPT }, + signal: opts.signal, + }); + } catch (err) { + if (isAbortError(err)) return { ok: false, error: 'Aborted.' }; + return { ok: false, error: err instanceof Error ? err.message : 'Network request failed.' }; + } + const headerRunId = parseTurnRunId(res) ?? cleanRunId; + const contentType = res.headers.get('content-type') ?? ''; + if (!res.body || !contentType.includes('text/event-stream') || res.headers.get('x-viewport-version') !== '1') { + const status = res.status; + const error = !contentType.includes('text/event-stream') + ? 'View attach failed.' + : 'Viewport negotiation not accepted.'; + return { ok: false, status, error: !res.body ? 'Empty viewport stream body.' : error, turnRunId: headerRunId }; + } + await opts.onTurnStarted?.({ turnRunId: headerRunId }); + const bodyResult = await readViewportBody(res.body!, headerRunId, opts); + return bodyResult as AgentFailure | ViewportAttachResult; +} \ No newline at end of file diff --git a/lib/viewportFirewall.ts b/lib/viewportFirewall.ts new file mode 100644 index 00000000..d5360dd1 --- /dev/null +++ b/lib/viewportFirewall.ts @@ -0,0 +1,13 @@ +/** + * Plan #960 firewall — explicit view/partial guard on CLOUD transcript paths. + * Meta-only preferences (model / reasoning / queue mirror) ride the narrow + * PATCH channel instead of a transcript write. The marker is the snapshot's + * explicit `historyComplete` field, never length inference. + */ +export function assertCompleteForCloudPut(snapshot: { historyComplete?: boolean }): void { + if (snapshot.historyComplete === false) { + throw new Error( + 'snapshot has historyComplete:false — view only (blocked from transcript put by the partial-view firewall)', + ); + } +} \ No newline at end of file diff --git a/lib/viewportStreamProtocol.ts b/lib/viewportStreamProtocol.ts index 6b3b104c..58da0814 100644 --- a/lib/viewportStreamProtocol.ts +++ b/lib/viewportStreamProtocol.ts @@ -40,6 +40,7 @@ export class ViewportStreamDecoder { private decoder = new TextDecoder('utf-8', { fatal: true }); private pending = ''; private nextIndex: number | undefined; + private lastViewportSnapshot: ViewportSnapshot | undefined; constructor(private readonly runId: string, startIndex?: number) { this.nextIndex = startIndex; } push(bytes: Uint8Array, final = false): ViewportRecord[] { this.pending += this.decoder.decode(bytes, { stream: !final }); @@ -57,6 +58,9 @@ export class ViewportStreamDecoder { if (final && this.pending.trim()) throw new Error('Incomplete viewport record'); return out; } + public get lastSnapshot(): ViewportSnapshot | undefined { + return this.lastViewportSnapshot; + } private parse(block: string): ViewportRecord { const fields = (name: string) => block.split('\n').filter(l => l.startsWith(`${name}:`)).map(l => l.slice(name.length + 1).trimStart()); const types = fields('event'), ids = fields('id'); @@ -98,8 +102,7 @@ export class ViewportStreamDecoder { } const start = sanitizeTurnStreamCursor(range.start)!; const end = sanitizeTurnStreamCursor(range.end)!; - return { - type, + const recSnapshotData: ViewportSnapshot = { version: 1, runId: this.runId, sessionId: o.sessionId as string, @@ -114,6 +117,9 @@ export class ViewportStreamDecoder { sampledRange: { start, end }, ...(index !== undefined ? { resumeIndex: index } : {}), }; + const recSnapshot = { type, ...recSnapshotData }; + if (this.lastViewportSnapshot === undefined) this.lastViewportSnapshot = recSnapshot; + return recSnapshot as ({ type: 'viewport_snapshot' } & ViewportSnapshot); } if (type === 'viewport_state' && o.phase === 'recovering' && typeof o.status === 'string') return { type, version: 1, runId: this.runId, status: o.status, phase: 'recovering' };