diff --git a/AGENTS.md b/AGENTS.md index ba494a06..31770eca 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -337,6 +337,7 @@ invincible/ | Kind of change | Where | |----------------|--------| | UI page / layout | `app/` | +| Bounded partial viewport read (opt-in backend, current host still unnegotiated) | `lib/sessionViewport.ts` (pure recent reducer, allowlisted event/carrier parse, UTF-8 excerpts), `lib/viewportStreamProtocol.ts` (versioned parser/indices), `lib/sessions/viewportRead.ts` (one scoped Blob head + at most 2048 recent frames / 8 MiB / 5 seconds, one final 1-second probe), `lib/workflows/viewportRunReader.ts` (injected public SDK handle), `lib/agent/viewportStream.ts` (snapshot-first + indexed live), GET `/api/sessions/:id/viewport`, negotiated turns GET/POST. Best-effort history: no ancestry merge/certificates, no recovery writes/start/cancel. `historyComplete:false` is not a full transcript or model seed. `npm run test:int` includes real-Wasm recovery→parser→ring evidence; it does not by itself prove the host cold-policy flip. | | DOM site chrome nav (hamburger Account menu) | `app/components/AppNav.tsx` (brand wordmark; optional `busy?: boolean` from `HarnessHost` only — TEAL outline + neon bloom + sine pulse + motes while the active harness turn is Busy; never poll the bridge; settings/admin omit the prop), `app/components/AuthNavLinks.tsx` (server: `soleMembership`+`canAccessAdmin` → `showAdmin`), `app/components/NavMenu.tsx` (client dropdown: ARIA `menu`, Arrow/Tab/Home/End, Escape + click-outside close + focus return, ≥44px touch targets, palette-only TEAL), `lib/navMenu.ts` + `lib/navMenu.test.ts` (`buildSignedInNavItems` — pure ordering/gating rule, unit-tested), footer slot `app/logout/LogoutButton.tsx`; unauth keeps inline `Sign in` header control. Client holds **zero** role-gate logic — it renders pre-gated inert `items` only | | API / AI Gateway / agent | `app/api/*`, `lib/agent/*`, `lib/sandbox/*`, `lib/gateway/modelCatalog.ts` (Gateway catalog, models.dev `vercel.models` fills holes; rewrite catalog `max` → `xhigh` then drop remaining non-wire tokens; never fetch inside `'use step'` / `'use workflow'`) | | Agent SSE stream (tools + text + reasoning) | `lib/agent/agentStream.ts`, `lib/agent/runAgent.ts`, `lib/agent/reasoningConfig.ts` (`glm-5*` ids look reasoning-capable → **`low`** when env unset and the joined catalog is empty; catalog is Gateway, filled from models.dev overlay when Gateway omitted (`lib/gateway/modelCatalog.ts`) — never fetch inside `'use step'`; GLM-5.x always thinks; never auto `provider-default` / `max` / `xhigh`; catalog `max` → `xhigh`; request `max` sends `xhigh` when listed else `high`), `lib/agent/generateOneRound.ts`, `lib/workflows/modelGenerateStep.ts`, `lib/workflows/toolExecuteStep.ts` (one step per model round’s `toolCalls`; live `tool_result` via `withDefaultStreamWriter`; `maxRetries = 0` — 1-call infra throws retry in-process), `lib/workflows/turnSseWrite.ts` (`withDefaultStreamWriter` around `generateOneRound` for durable live `reasoning_delta` / `text_delta` / `tool_start` / `provider` — **one** `getWritable()` per model round, not per token; same held-writer for a tool batch; loop does not dump those events; `writeOnDefaultStream` is the sparse `writeTurnSse` path for `done` / `error`; stream PUT 5xx/429/timeout latch immediately — SDK already retried 429 inside the first `write()`; do not map writer I/O to `write_error`; AbortError is never latched), `app/api/agent/route.ts`, `lib/agentApi.ts`, `docs/agent-stream.md` | diff --git a/SECURITY.md b/SECURITY.md index 82694f6e..838f9a13 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -48,6 +48,29 @@ Never use `NEXT_PUBLIC_SANDBOX_*` (or any client-exposed sandbox secret). | Backfill | One-shot Postgres `harness_sessions` → Redis via GHA **`sessions-redis-backfill`** (per-`{tenant,user}` marker, idempotent); Postgres becomes a **read-only archive**; legacy `/api/session` write route removed | | Client bundle | Session repository is client-safe (`lib/sessionRepository.ts`); must **not** import server `db` / Drizzle modules | +### Read-only partial viewport + +`GET /api/sessions/:id/viewport` and the negotiated `viewportVersion=1` +durable stream mode require the same authenticated tenant/user session scope. +A run stream must match the owned envelope's **current** run id before SDK access; +cold handoff rechecks that binding before releasing recovered rows. A planted Blob +pointer is read only if `isObjectIdBoundTo` matches the same session, and the decoded +body's `id` must match. These routes never accept client URLs, read arbitrary old +Workflow inputs, start/cancel runs, execute tools, or emit secrets, raw run +inputs, or signed read URLs. The response is a **partial, disposable display +view** for the owning session — not a canonical transcript and not a public +dataset. + +Optional history is best effort, **authorization is not**. Missing/corrupt history +can produce `replace:false`; a missing/foreign session/run fails closed. Carriers +are allowlisted and exclude persona/notes/model/compaction bodies, raw metadata, +keys and signed read URLs. Negotiated stream failures use sanitized codes, not +backend connection details. Responses are private/no-store. Historical reasoning +is omitted from snapshots; post-handoff live events keep the existing redacted +producer contract. `historyComplete:false` must never be treated as a full +transcript replacement or inference seed. A transport resume index does not +certify that omitted history is visible. + Product behavior: [docs/session-model.md](docs/session-model.md). ## Builtin HTTPS fetch (Vercel Sandbox) diff --git a/app/api/sessions/[id]/viewport/route.test.ts b/app/api/sessions/[id]/viewport/route.test.ts new file mode 100644 index 00000000..691cad3e --- /dev/null +++ b/app/api/sessions/[id]/viewport/route.test.ts @@ -0,0 +1,37 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { newBlobObjectId } from '../../../../../lib/sessions/blobStore'; +const mocks=vi.hoisted(()=>({auth:vi.fn(),tenant:vi.fn(),store:vi.fn(),envelope:vi.fn(),read:vi.fn()})); +vi.mock('../../../../../lib/tenancy/session',()=>({requireSessionUser:mocks.auth})); +vi.mock('../../../../../lib/di',()=>({createProdServices:()=>({harnessSessionsRedis:{resolveTenantIdForUser:mocks.tenant},createBlobTranscriptStore:()=>({read:mocks.read})})})); +vi.mock('../../../../../lib/tenancy/harnessSessionsRedis',()=>({resolveSessionStore:mocks.store,sessionKeyFor:(tenantId:string,userId:string,sessionId:string)=>({tenantId,userId,sessionId})})); +vi.mock('../../../../../lib/sessions/sessionStore',()=>({isEnvelopeStore:()=>true})); +import {GET, maxDuration} from './route'; +const scope={tenantId:'t',userId:'u',sessionId:'s'}; +beforeEach(()=>{vi.clearAllMocks();mocks.auth.mockResolvedValue({ok:true,user:{id:'u'}});mocks.tenant.mockResolvedValue({ok:true,value:'t'}); + mocks.store.mockResolvedValue({ok:true,value:{readEnvelope:mocks.envelope}}); + mocks.envelope.mockResolvedValue({meta:{transcriptPointer:newBlobObjectId(scope),turnRunId:'run',workingNotes:'secret notes',personaSnapshot:'secret persona'}}); + mocks.read.mockResolvedValue(JSON.stringify({id:'s',prev:newBlobObjectId(scope),messages:[{role:'assistant',text:'recent'}],queue:['next']}));}); +const request=(query='')=>GET(new Request(`https://example.com/api/sessions/s/viewport${query}`),{params:Promise.resolve({id:'s'})}); +describe('head-only view route',()=>{ + it('is a short JSON function, not a 1800s SSE attach',()=>{ + expect(maxDuration).toBe(15); + }); + it('reads one head and safe carriers, no private meta or write surface',async()=>{ + const res=await request();expect(res.status).toBe(200);expect(res.headers.get('cache-control')).toContain('private'); + const body=await res.json();expect(body).toMatchObject({historyComplete:false,replace:true,hasEarlier:true,gap:false,rows:[{text:'recent'}],carriers:{queue:['next']}}); + expect(JSON.stringify(body)).not.toContain('secret');expect(mocks.read).toHaveBeenCalledOnce(); + }); + it('corrupt/missing head retains cached paint',async()=>{ + for(const raw of [null,'invalid',JSON.stringify({id:'wrong',messages:[]})]){mocks.read.mockResolvedValue(raw);expect(await (await request()).json()).toMatchObject({replace:false,source:'unavailable'});} + }); + it('auth/ownership run mismatch before body access; unknown page params rejected until phase4',async()=>{ + mocks.auth.mockResolvedValue({ok:false,response:Response.json({}, {status:401})});expect((await request()).status).toBe(401); + mocks.auth.mockResolvedValue({ok:true,user:{id:'u'}});expect((await request('?runId=foreign')).status).toBe(404); + expect((await request('?objectId=anything')).status).toBe(400);expect(mocks.read).not.toHaveBeenCalled(); + }); + it('foreign planted pointer is not read and missing store fails closed',async()=>{ + mocks.envelope.mockResolvedValue({meta:{transcriptPointer:newBlobObjectId({...scope,userId:'foreign'})}}); + expect(await (await request()).json()).toMatchObject({replace:false});expect(mocks.read).not.toHaveBeenCalled(); + mocks.store.mockResolvedValue({ok:false});expect((await request()).status).toBe(503); + }); +}); diff --git a/app/api/sessions/[id]/viewport/route.ts b/app/api/sessions/[id]/viewport/route.ts new file mode 100644 index 00000000..72241266 --- /dev/null +++ b/app/api/sessions/[id]/viewport/route.ts @@ -0,0 +1,44 @@ +/** Read-only, head-only display view. Never a canonical transcript or model seed. */ +import { createProdServices } from '../../../../../lib/di'; +import { requireSessionUser } from '../../../../../lib/tenancy/session'; +import { resolveSessionStore, sessionKeyFor } from '../../../../../lib/tenancy/harnessSessionsRedis'; +import { isEnvelopeStore } from '../../../../../lib/sessions/sessionStore'; +import { readViewportHead, emptyViewport } from '../../../../../lib/sessions/viewportRead'; +import { isRedisSafeOpaqueId, sanitizeTurnRunId, VIEWPORT_RECOVERY_MAX_MS } from '../../../../../lib/sessionCloudCaps'; +import { viewportWait } from '../../../../../lib/workflows/viewportRunReader'; + +const services = createProdServices(); +export const runtime = 'nodejs'; +/** JSON 5 s recovery budget — not a long-lived SSE attach. */ +export const maxDuration = 15; +const headers = { 'Cache-Control': 'private, no-store, no-transform' }; + +export async function GET(req: Request, ctx: { params: Promise<{ id: string }> }): Promise { + const auth = await requireSessionUser(); + if (!auth.ok) return auth.response; + if (!auth.user?.id) return Response.json({ error: 'Authentication required.' }, { status: 401 }); + const { id } = await ctx.params; + if (!isRedisSafeOpaqueId(id)) return Response.json({ error: 'Invalid session id.' }, { status: 400 }); + const q = new URL(req.url).searchParams; + const rawRun = q.get('runId'); + if (q.getAll('runId').length > 1 || (rawRun !== null && sanitizeTurnRunId(rawRun) !== rawRun) || + [...q.keys()].some(k => k !== 'runId')) return Response.json({ error: 'Invalid viewport query.' }, { status: 400 }); + try { + const tenant = await services.harnessSessionsRedis.resolveTenantIdForUser(auth.user.id); + const stored = await resolveSessionStore(); + if (!tenant.ok || !stored.ok || !isEnvelopeStore(stored.value)) throw new Error('Store unavailable'); + const scope = { tenantId: tenant.value, userId: auth.user.id, sessionId: id }; + const key = sessionKeyFor(scope.tenantId, scope.userId, id); + const envelope = await stored.value.readEnvelope(key); + if (!envelope || (rawRun !== null && envelope.meta.turnRunId !== rawRun)) + return Response.json({ error: 'Session not found.' }, { status: 404 }); + const head = await viewportWait((async () => { + try { return await readViewportHead({ scope, meta: envelope.meta, + blob: services.createBlobTranscriptStore(), signal: req.signal }); } + catch { return emptyViewport(id, envelope.meta); } + })(), VIEWPORT_RECOVERY_MAX_MS, req.signal).catch(() => emptyViewport(id, envelope.meta)); + return Response.json(head, { headers }); + } catch { + return Response.json({ error: 'Viewport store unavailable.' }, { status: 503, headers }); + } +} diff --git a/app/api/turns/[runId]/stream/route.ts b/app/api/turns/[runId]/stream/route.ts index 3ac6a857..421c0a27 100644 --- a/app/api/turns/[runId]/stream/route.ts +++ b/app/api/turns/[runId]/stream/route.ts @@ -96,6 +96,15 @@ export async function GET( return Response.json({ error: 'Invalid runId' }, { status: 400 }); } + const query = new URL(req.url); + let viewportMode: import('../../../../../lib/viewportStreamProtocol').ViewportMode = { kind: 'legacy' }; + if (query.searchParams.has('viewportVersion') || query.searchParams.has('hydrate')) { + const { parseViewportMode } = await import('../../../../../lib/viewportStreamProtocol'); + const mode = parseViewportMode(query, 'GET'); + if (!mode) return Response.json({ error: 'Invalid viewport negotiation.' }, { status: 400 }); + viewportMode = mode; + } + // Parse and validate startIndex query param. // - absent → default 0 (full replay) // - present → non-negative integer ≤ TURN_STREAM_CURSOR_MAX @@ -179,10 +188,12 @@ export async function GET( ); } + let envelopeMeta: Record = {}; // Read the session envelope — fail-closed (503) on read throw, 404 on // miss or turnRunId mismatch. Previously fail-open on read throw. try { const envelope = await envelopeStore.readEnvelope(sessionKey); + envelopeMeta = envelope?.meta ?? {}; if (!envelope || envelope.meta?.turnRunId !== cleanRunId) { return Response.json( { error: `Run not found: ${cleanRunId}` }, @@ -220,11 +231,43 @@ export async function GET( 'x-workflow-run-id': cleanRunId, }; + if (viewportMode.kind !== 'legacy') { + const { createViewportRunReader } = await import('../../../../../lib/workflows/viewportRunReader'); + const { viewportStream } = await import('../../../../../lib/agent/viewportStream'); + const { readViewportHead, emptyViewport } = await import('../../../../../lib/sessions/viewportRead'); + const runReader = createViewportRunReader(run); + const cold = viewportMode.kind === 'cold'; + const body = viewportStream({ + runId: cleanRunId, sessionId, run: runReader, + startIndex: viewportMode.kind === 'indexed' ? viewportMode.startIndex : 0, + signal: req.signal, + ...(cold ? { cold: { + status: envelopeMeta.turnStatus === 'cancelling' ? 'cancelling' : 'running', + readHead: async (deadline: number) => { + try { + return await readViewportHead({ scope: { tenantId: tenantRes.value, userId, sessionId }, + meta: envelopeMeta, blob: services.createBlobTranscriptStore(), signal: req.signal, deadline }); + } catch { return emptyViewport(sessionId, envelopeMeta); } + }, + } } : {}), + stillOwned: async () => { + const current = await envelopeStore.readEnvelope(sessionKey); + return current?.meta?.turnRunId === cleanRunId; + }, + }); + return new Response(body, { headers: { ...headers, 'Cache-Control': 'private, no-store, no-transform', 'x-viewport-version': '1' } }); + } + return new Response(await bodyForRun(run, { startIndex }), { status: 200, headers, }); } catch (err) { + if (viewportMode.kind !== 'legacy') { + return Response.json({ error: 'Viewport stream unavailable.' }, { + status: 503, headers: { 'Cache-Control': 'private, no-store, no-transform' }, + }); + } const msg = err instanceof Error ? err.message : String(err); return Response.json( { error: `Unable to attach to run stream (fail closed): ${msg}` }, diff --git a/app/api/turns/[runId]/stream/viewport.test.ts b/app/api/turns/[runId]/stream/viewport.test.ts new file mode 100644 index 00000000..5c6f3afd --- /dev/null +++ b/app/api/turns/[runId]/stream/viewport.test.ts @@ -0,0 +1,102 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { newBlobObjectId } from '../../../../../lib/sessions/blobStore'; +import { ViewportStreamDecoder, type ViewportRecord } from '../../../../../lib/viewportStreamProtocol'; +const mocks=vi.hoisted(()=>({auth:vi.fn(),tenant:vi.fn(),store:vi.fn(),getRun:vi.fn(),read:vi.fn(),envelope:vi.fn()})); +vi.mock('workflow/api',()=>({getRun:mocks.getRun})); +vi.mock('../../../../../lib/tenancy/session',()=>({requireSessionUser:mocks.auth})); +vi.mock('../../../../../lib/di',()=>({createProdServices:()=>({harnessSessionsRedis:{resolveTenantIdForUser:mocks.tenant},createBlobTranscriptStore:()=>({read:mocks.read})})})); +vi.mock('../../../../../lib/tenancy/harnessSessionsRedis',()=>({resolveSessionStore:mocks.store,sessionKeyFor:(tenantId:string,userId:string,sessionId:string)=>({tenantId,userId,sessionId})})); +vi.mock('../../../../../lib/sessions/sessionStore',()=>({isEnvelopeStore:()=>true})); +import { GET } from './route'; +const scope={tenantId:'t1',userId:'u1',sessionId:'s1'}; +const line=(event:object)=>`data: ${JSON.stringify(event)}\n\n`; +let open:ReturnType, cancel:ReturnType; +beforeEach(()=>{ + vi.clearAllMocks(); + mocks.auth.mockResolvedValue({ok:true,user:{id:'u1'}});mocks.tenant.mockResolvedValue({ok:true,value:'t1'}); + mocks.store.mockResolvedValue({ok:true,value:{readEnvelope:mocks.envelope}}); + mocks.envelope.mockResolvedValue({meta:{turnRunId:'run',turnStatus:'cancelling',transcriptPointer:newBlobObjectId(scope)}}); + mocks.read.mockResolvedValue(JSON.stringify({id:'s1',messages:[{role:'assistant',text:'head'}]})); + cancel=vi.fn(); + open=vi.fn((opts?:{startIndex?:number})=>Object.assign(new ReadableStream({start(c){ + if(opts?.startIndex===0){c.enqueue(line({type:'reasoning_delta',text:'old-thinking'}));c.enqueue(line({type:'text_delta',text:'recent'}));c.close();} + else if(opts?.startIndex===2){c.enqueue(line({type:'reasoning_delta',text:'new-thinking'}));c.enqueue(line({type:'done',text:'done'}));c.close();} + },cancel}),{getTailIndex:async()=>1})); + mocks.getRun.mockReturnValue({exists:Promise.resolve(true),status:Promise.resolve('running'),getReadable:open}); +}); +afterEach(()=>vi.useRealTimers()); +const request=(query='sessionId=s1&viewportVersion=1&hydrate=tail')=>GET(new Request(`https://example.com/api/turns/run/stream?${query}`),{params:Promise.resolve({runId:'run'})}); +async function decode(res:Response,initial?:number){const parser=new ViewportStreamDecoder('run',initial), out:ViewportRecord[]=[]; const reader=res.body!.getReader();for(;;){const r=await reader.read();if(r.done)break;out.push(...parser.push(r.value));}parser.push(new Uint8Array(),true);return out;} +describe('negotiated route uses real bounded service/codec',()=>{ + it('authorizes, samples, returns cancelling state and snapshot, then new live reasoning',async()=>{ + const res=await request();expect(res.status).toBe(200);expect(res.headers.get('x-viewport-version')).toBe('1'); + expect(res.headers.get('cache-control')).toBe('private, no-store, no-transform'); + const records=await decode(res);expect(records[0]).toMatchObject({type:'viewport_state',status:'cancelling'}); + expect(records[1]).toMatchObject({type:'viewport_snapshot',resumeIndex:2,source:'stream_tail'}); + expect(JSON.stringify(records)).not.toContain('old-thinking');expect(JSON.stringify(records)).toContain('new-thinking'); + expect(mocks.read).toHaveBeenCalledOnce();expect(mocks.envelope).toHaveBeenCalledTimes(2); + expect(open).toHaveBeenCalledWith({startIndex:-1}); + }); + it('negotiated hot attach never probes head or tail',async()=>{ + const records=await decode(await request('sessionId=s1&viewportVersion=1&startIndex=2'),2); + expect(records[0]).toMatchObject({type:'turn_event',nextIndex:3}); + expect(mocks.read).not.toHaveBeenCalled();expect(open).toHaveBeenCalledExactlyOnceWith({startIndex:2}); + }); + it.each(['sessionId=s1&viewportVersion=2','sessionId=s1&viewportVersion=1','sessionId=s1&viewportVersion=1&hydrate=tail&startIndex=0','sessionId=s1&viewportVersion=1&startIndex=1e2','sessionId=s1&sessionId=s2&viewportVersion=1'])('rejects bad negotiation before SDK access: %s',async q=>{ + expect((await request(q)).status).toBe(400);expect(mocks.getRun).not.toHaveBeenCalled(); + }); + it('unauth and foreign run never reach SDK/Blob',async()=>{ + mocks.auth.mockResolvedValue({ok:false,response:Response.json({error:'auth'},{status:401})});expect((await request()).status).toBe(401); + mocks.auth.mockResolvedValue({ok:true,user:{id:'u1'}});mocks.envelope.mockResolvedValue({meta:{turnRunId:'foreign'}}); + expect((await request()).status).toBe(404);expect(mocks.getRun).not.toHaveBeenCalled();expect(mocks.read).not.toHaveBeenCalled(); + }); + it('foreign planted Blob pointer does not read it, live sample still useful',async()=>{ + mocks.envelope.mockResolvedValue({meta:{turnRunId:'run',transcriptPointer:newBlobObjectId({...scope,userId:'foreign'})}}); + const records=await decode(await request());expect(records[1]).toMatchObject({source:'stream_tail'});expect(mocks.read).not.toHaveBeenCalled(); + }); + it('missing head is a partial-history miss; revoked session before handoff is not released',async()=>{ + mocks.read.mockResolvedValue(null);mocks.envelope.mockResolvedValueOnce({meta:{turnRunId:'run'}}).mockResolvedValue(null); + const records=await decode(await request());expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_error']); + }); + it('does not expose SDK failure details in negotiated pre-stream errors', async () => { + mocks.getRun.mockImplementation(() => { throw new Error('provider URL/token-private-detail'); }); + const res = await request(); + expect(res.status).toBe(503); + expect(await res.json()).toEqual({ error: 'Viewport stream unavailable.' }); + expect(res.headers.get('cache-control')).toContain('private'); + }); + it('unavailable initial tail still emits recovering state and head; never origin-replays',async()=>{ + const getReadable=vi.fn(()=>new ReadableStream()); + mocks.getRun.mockReturnValue({exists:Promise.resolve(true),status:'running',getReadable}); + const res=await request();expect(res.status).toBe(200); + const records=await decode(res); + expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_snapshot','viewport_error']); + expect(records[0]).toMatchObject({type:'viewport_state',phase:'recovering'}); + expect(records[1]).toMatchObject({source:'stored_head'}); + expect(records[1]).not.toHaveProperty('resumeIndex'); + expect(JSON.stringify(records)).toContain('head'); + expect(records[2]).toMatchObject({code:'STREAM_UNAVAILABLE'}); + expect(getReadable).toHaveBeenCalledWith({startIndex:-1}); + expect(getReadable).not.toHaveBeenCalledWith({startIndex:0}); + }); + it.each(['cancelled','failed'])('%s hanging getReadable is never opened on hydrate=tail; head snapshot still ships',async status=>{ + const getReadable=vi.fn(()=>new Promise>(()=>{})); + mocks.getRun.mockReturnValue({exists:Promise.resolve(true),status:Promise.resolve(status),getReadable}); + const res=await request();expect(res.status).toBe(200); + const records=await decode(res); + expect(getReadable).not.toHaveBeenCalled();expect(mocks.read).toHaveBeenCalledOnce(); + expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_snapshot','viewport_end']); + expect(records[0]).toMatchObject({status}); + expect(records[1]).toMatchObject({source:'stored_head'}); + expect(records[1]).not.toHaveProperty('resumeIndex'); + expect(JSON.stringify(records)).toContain('head'); + expect(records[2]).toMatchObject({status}); + }); + it.each(['cancelled','failed'])('%s hanging getReadable is never opened on indexed GET',async status=>{ + const getReadable=vi.fn(()=>new Promise>(()=>{})); + mocks.getRun.mockReturnValue({exists:Promise.resolve(true),status:Promise.resolve(status),getReadable}); + const records=await decode(await request('sessionId=s1&viewportVersion=1&startIndex=2'),2); + expect(getReadable).not.toHaveBeenCalled();expect(mocks.read).not.toHaveBeenCalled(); + expect(records).toEqual([{type:'viewport_end',version:1,runId:'run',status}]); + }); +}); diff --git a/app/api/turns/route.test.ts b/app/api/turns/route.test.ts index ad3dac33..d8277fc2 100644 --- a/app/api/turns/route.test.ts +++ b/app/api/turns/route.test.ts @@ -314,6 +314,66 @@ describe('POST /api/turns', () => { expect(patchCall.key).toEqual({ tenantId: 't1', userId: 'u1', sessionId: 's1' }); }); + it('negotiated POST uses indexed events without changing start args or replaying a snapshot', async () => { + standardHarness(); mockAuthedSession(); + const { getReadable } = mockStart(); + getReadable.mockImplementation(() => new ReadableStream({ start(c) { + c.enqueue('data: {"type":"reasoning_delta","text":"new thinking"}\n\n'); + c.enqueue('data: {"type":"done","text":"new answer"}\n\n'); c.close(); + } })); + ({ POST } = await import('./route')); + const res = await POST(new Request('https://x/api/turns?viewportVersion=1', { + method: 'POST', headers: { accept: 'text/event-stream', 'content-type': 'application/json' }, + body: JSON.stringify({ sessionId: 's1', prompt: 'hi' }), + })); + expect(res.status).toBe(200); expect(res.headers.get('x-viewport-version')).toBe('1'); + const { ViewportStreamDecoder } = await import('../../../lib/viewportStreamProtocol'); + const records = new ViewportStreamDecoder('wf_turn_123', 0).push(new Uint8Array(await res.arrayBuffer()), true); + expect(records).toHaveLength(2); expect(records[0]).toMatchObject({ type: 'turn_event', nextIndex: 1, event: { type: 'reasoning_delta' } }); + expect(startMock).toHaveBeenCalledOnce(); expect(startMock.mock.calls[0][1][0].userMessage).toBe('hi'); + expect(getReadable).toHaveBeenCalledWith({ startIndex: 0 }); + }); + + it('invalid viewport negotiation rejects before starting a durable run', async () => { + standardHarness(); mockAuthedSession(); mockStart(); ({ POST } = await import('./route')); + const res = await POST(new Request('https://x/api/turns?viewportVersion=1&hydrate=tail', { + method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ sessionId: 's1', prompt: 'hi' }), + })); + expect(res.status).toBe(400); expect(startMock).not.toHaveBeenCalled(); + }); + + it('negotiated POST start throw is sanitized 503, never SDK details', async () => { + standardHarness(); mockAuthedSession(); mockStart(); + startMock.mockRejectedValueOnce(new Error('provider URL/token-private-detail')); + ({ POST } = await import('./route')); + const res = await POST(new Request('https://x/api/turns?viewportVersion=1', { + method: 'POST', headers: { accept: 'text/event-stream', 'content-type': 'application/json' }, + body: JSON.stringify({ sessionId: 's1', prompt: 'hi' }), + })); + expect(res.status).toBe(503); + expect(res.headers.get('cache-control')).toBe('private, no-store, no-transform'); + expect(await res.json()).toEqual({ error: 'Viewport stream unavailable.' }); + expect(startMock).toHaveBeenCalledOnce(); + }); + + it('negotiated POST live-guard getRun throw is sanitized 503, never SDK details', async () => { + standardHarness(); mockAuthedSession(); mockStart(); + getRunMock.mockImplementation(() => { throw new Error('provider URL/token-private-detail'); }); + readEnvelopeMock.mockResolvedValue({ + updatedAt: FUTURE_UPDATED_AT, + meta: { turnStatus: 'running', turnRunId: 'wf_live_1' }, + }); + ({ POST } = await import('./route')); + const res = await POST(new Request('https://x/api/turns?viewportVersion=1', { + method: 'POST', headers: { accept: 'text/event-stream', 'content-type': 'application/json' }, + body: JSON.stringify({ sessionId: 's1', prompt: 'hi' }), + })); + expect(res.status).toBe(503); + expect(res.headers.get('cache-control')).toBe('private, no-store, no-transform'); + expect(await res.json()).toEqual({ error: 'Viewport stream unavailable.' }); + expect(startMock).not.toHaveBeenCalled(); + }); + it('body reasoning is passed to start() (plan #897)', async () => { standardHarness(); mockAuthedSession(); diff --git a/app/api/turns/route.ts b/app/api/turns/route.ts index aabdd411..dab435c6 100644 --- a/app/api/turns/route.ts +++ b/app/api/turns/route.ts @@ -133,6 +133,14 @@ function failClosed(err: unknown): string { return `Unable to start durable turn (fail closed): ${msg}`; } +/** Negotiated v1 never interpolates SDK/connection details (GET stream parity). */ +function viewportStreamUnavailable(): Response { + return Response.json( + { error: 'Viewport stream unavailable.' }, + { status: 503, headers: { 'Cache-Control': 'private, no-store, no-transform' } }, + ); +} + /** Hard deny when resolve returned no client AND there is no soft-path fallback. */ function isHardSandboxDeny( res: ResolveAgentSandboxResult, @@ -178,6 +186,15 @@ export async function POST(req: Request): Promise { return Response.json({ error: AUTH_REQUIRED_ERROR }, { status: 401 }); } + let indexedViewport = false; + const query = new URL(req.url); + if (query.searchParams.has('viewportVersion') || query.searchParams.has('hydrate')) { + const { parseViewportMode } = await import('../../../lib/viewportStreamProtocol'); + const mode = parseViewportMode(query, 'POST'); + if (!mode) return Response.json({ error: 'Invalid viewport negotiation.' }, { status: 400 }); + indexedViewport = mode.kind === 'indexed'; + } + let body: unknown; try { body = await req.json(); @@ -377,6 +394,7 @@ export async function POST(req: Request): Promise { } // exists === false or terminal status → not live; allow start. } catch (err) { + if (indexedViewport) return viewportStreamUnavailable(); return Response.json({ error: failClosed(err) }, { status: 503 }); } } @@ -783,6 +801,15 @@ export async function POST(req: Request): Promise { runHeaders['x-workflow-run-warning'] = runWarning; } if (wantsAgentStream(req)) { + if (indexedViewport) { + const { createViewportRunReader } = await import('../../../lib/workflows/viewportRunReader'); + const { viewportStream } = await import('../../../lib/agent/viewportStream'); + return new Response(viewportStream({ runId: run.runId, sessionId, startIndex: 0, + run: createViewportRunReader(run), signal: req.signal }), { headers: { + ...runHeaders, 'content-type': AGENT_STREAM_CONTENT_TYPE, 'x-viewport-version': '1', + 'Cache-Control': 'private, no-store, no-transform', 'X-Accel-Buffering': 'no', + } }); + } return new Response(await bodyForRun(run), { status: 200, headers: { @@ -805,6 +832,7 @@ export async function POST(req: Request): Promise { // Ignore close errors. } } + if (indexedViewport) return viewportStreamUnavailable(); return Response.json({ error: failClosed(err) }, { status: 503 }); } finally { // Clear the in-flight flag on EVERY path — success, throw, or any early diff --git a/docs/agent-stream.md b/docs/agent-stream.md index 0f29ea95..f6886d0f 100644 --- a/docs/agent-stream.md +++ b/docs/agent-stream.md @@ -17,6 +17,59 @@ Early failures (auth, grants, bad body, BYOK) always use **JSON** status respons Response hints: `Cache-Control: no-cache, no-transform`, `X-Accel-Buffering: no`. +## Negotiated viewport stream (version 1) + +This is additive backend support; the current production host still uses the +unnegotiated event stream. Clients must opt in and use `ViewportStreamDecoder` +(`lib/viewportStreamProtocol.ts`), not feed these controls to the legacy parser. + +| Request | Behavior | +|---|---| +| GET existing run stream with `sessionId`, `viewportVersion=1`, `hydrate=tail` | Recovering state → one best-effort snapshot → live indexed events; `startIndex` is forbidden with hydrate | +| GET with `sessionId`, `viewportVersion=1`, `startIndex=N` | Indexed hot resume from a same-heap applied raw cursor; no history read. `N=0` is an explicit origin replay (not the default). | +| GET `viewportVersion=1` without `hydrate` and without `startIndex` | **400** — not an implicit `startIndex=0` origin replay | +| POST `/api/turns?viewportVersion=1`, `Accept: text/event-stream` | New run's indexed events from origin; inference/start args unchanged | +| No version | Existing SSE/JSON behavior unchanged | + +Unknown versions, duplicate/conflicting selectors and noncanonical indexed +cursor syntax reject with 400. Negotiated responses include `x-viewport-version: 1` +and `Cache-Control: private, no-store, no-transform`. + +Each block has an SSE `event:` name and JSON `data:` (all data carries `version:1` +and `runId`). Only stored `turn_event` records have `id:`: + +| event | Additional fields / meaning | +|---|---| +| `viewport_state` | `status`, `phase:'recovering'`; restore Busy/Stop without erasing cached paint | +| `viewport_snapshot` | `sessionId`, optional `resumeIndex`, `rows`, `replace`, `source`, `sampledRange:{start,end}`, `historyComplete:false`, `incomplete:true`, `gap`, `hasEarlier`, safe `carriers`. `resumeIndex` is omitted when the live tail is unknown (display-only; the decoder does not jump). A missing tail is never encoded as `resumeIndex:0`. | +| `turn_event` | `nextIndex`, `event` (allowlisted AgentStreamEvent) and matching SSE `id: nextIndex`; a malformed known stored frame instead carries `skipped:true` with no event. A stored `done` or `error` **is producer-terminal**: the iterator closes after that record and **does not** emit `viewport_end` (cancelled inject is an `error` event, not a synthetic `failed`). Consumers must treat stored `done`/`error`, `viewport_end`, `viewport_error`, **and** reader EOF as terminal. | +| `viewport_end` | Synthetic terminal `status` (`completed`, `failed`, `cancelled`) when the wrapper stops from `run.status` or readable EOF **without** a stored `done`/`error` (hang-class attach, completed drain that never wrote `done`). No raw index, no transcript-completeness promise | +| `viewport_error` | Sanitized `code`; preserve last applied cursor and detach, never restart/cancel inference | + +Recovery samples only the recent raw-frame interval and chooses sampled display +**or** the latest stored head, not an exact merge. Historical thinking is omitted. +A single final tail probe establishes the handoff; new data produced while sampling +may be skipped and marked `gap:true`. Post-handoff reasoning is live. No loop chases +the producer. See [harness-limits.md](harness-limits.md) for recovery budgets. + +An explicit snapshot may jump the cursor; subsequent stored ids must be contiguous. +UTF-8/CRLF network fragmentation does not create new stored positions. A transport +or decoder failure that loses position closes the reader, rather than inventing a +cursor. Synthetic terminal status never consumes an index. Completed-run buffered +frames drain before a hung read is resolved from a **1 s** terminal status poll; +the 0-delay first poll unsticks already-`cancelled`/`failed` only (not `completed`). +Already-`cancelled`/`failed` attach never calls `getReadable` (same C16 gate as +`bodyForRun`); cold hydrate still returns a head snapshot then `viewport_end`. +Cold GET emits `viewport_state` **before** capturing H0 (under the same 5 s +recovery clock as head+sample). Unavailable or hanging live-tail metadata is an +in-band `viewport_error` after that head snapshot — never an empty 503, never +`open(0)` origin replay, and never a guessed `resumeIndex: 0` cursor jump. + +All view rows are disposable and incomplete. Missing history may preserve the +cached ring (`replace:false`); oversized rows carry a visible excerpt marker. +Never turn these rows into a canonical full-snapshot upload or inference seed. +No new producer metadata, checkpoint certificates, or storage writes are required. + ## Request body (cwd) | Field | Required | Notes | diff --git a/docs/harness-limits.md b/docs/harness-limits.md index 1a3d1a4a..e0c8937a 100644 --- a/docs/harness-limits.md +++ b/docs/harness-limits.md @@ -41,6 +41,42 @@ click/tap scrolls the transcript so the message is back in view near the top. | Cache | `public, max-age=3600, stale-while-revalidate=86400` | | Build id | Baked short git SHA (`-Dbuild-id`); file `public/harness/build-id.txt`; shown as `h:…` in status bar line 1 — must match after deploy | +## Bounded viewport recovery + +The opt-in `viewportVersion=1` backend read path prioritizes useful recent display +over exact archive reconstruction. The current host's default restore remains +unnegotiated. These **new recovery budgets do not change durable storage, model +context, queue, turn or bridge limits**. + +| Recovery cap (`lib/sessionCloudCaps.ts`) | Value | Behavior | +|---|---:|---| +| `VIEWPORT_TAIL_MAX_FRAMES` | 2048 | Sample only the recent stored-frame interval, never full-origin history on a large run | +| `VIEWPORT_RECOVERY_MAX_BYTES` | 8 MiB | Decoded sampled frame bytes, including reasoning that is discarded | +| `VIEWPORT_RECOVERY_MAX_MS` | 5000 ms | Shared optional H0 + head/sample deadline; not a timeout/cancel of live inference | +| `VIEWPORT_FINAL_PROBE_MAX_MS` | 1000 ms | One final tail probe; on timeout keep initial known tail, no chase loop | +| `VIEWPORT_HEAD_READ_MAX_OBJECTS` | 1 | Current scoped head only; no `prev` reconstruction | +| `VIEWPORT_RESPONSE_MAX_BYTES` | 2 MiB | Entire escaped JSON snapshot/control incl. carriers/rows, below real 4.5 MB Function ceiling; not lifetime SSE bytes | +| Output ring/row | existing 2048 /262144 UTF-8 bytes | Newest rows and labeled UTF-8-safe display excerpts; source objects unchanged | + +Budget exhaustion omits history and resumes at a known tail, with incomplete/gap +flags. Sample rows or stored head are selected, not expensively aligned/merged. +One neutral note explains omissions; historical reasoning is never in the snapshot. +A missing optional head does not turn into an origin replay or run cancellation. +Missing/hanging live-tail metadata (`getTailIndex`) fails the live attach in-band +after `viewport_state` + head snapshot — never an empty 503, never `open(0)`, and +never a guessed `resumeIndex: 0` (the snapshot is display-only; the decoder does +not jump). + +The SDK may decode one oversized frame before the byte check, and a Blob read +returns one complete object (raw + parsed JS overhead is larger than wire bytes). +Cancellation stops readers and app work; a provider call without abort support can +finish late but cannot launch more recovery work or replace the selected view. +Only finite head/probe/sample calls are launched. Time bounds cover asynchronous +waits and checked loop boundaries, not preemption of one synchronous JSON parse. +This is bounded best effort, not an exact heap/latency guarantee or a new physical +history index. The real-Wasm durable int project tests the recovery/parser/ring +path separately from the default unit suite. + ## Keyboard & focus All harness chords are rows of one keymap table diff --git a/docs/session-model.md b/docs/session-model.md index 5ee15a52..48d18f9d 100644 --- a/docs/session-model.md +++ b/docs/session-model.md @@ -10,6 +10,39 @@ How harness continuity works: **local-first** browser restore plus - Client must not use Node `fs` - Persistence I/O stays on the **DOM host** — Wasm never talks to storage or `/api/sessions*` +## Bounded viewport read API + +The backend offers an explicitly **partial, read-only** display path; the current +`HarnessHost` still uses the unnegotiated restore path described below. + +- `GET /api/sessions/:id/viewport` reads only the current session-bound Blob + **head**, keeps recent rows and safe host carriers, and never follows `prev`. + Missing/corrupt history returns `replace:false` so a consumer keeps cached paint. +- `GET /api/turns/:runId/stream?sessionId=:id&viewportVersion=1&hydrate=tail` + sends a recovering-state record, one bounded recent snapshot, then indexed live + events on the **same response**. Recovering state is emitted before H0 / tail + metadata; a missing tail is in-band `viewport_error` after the head snapshot + (not an empty 503) and does **not** publish `resumeIndex: 0`. It samples at most 2048 recent stored frames, + 8 MiB, and 5 seconds of optional recovery work. One final tail probe (1 second) + can skip a growing backlog; omitted gaps are explicit, never replayed at UI speed. +- Sampled reasoning is discarded. Post-handoff reasoning is live, including a + continuation of an older segment. Older prompt/tool context may be absent, + fragmented or stale; this is a useful recent view, not exact reconstruction. +- `historyComplete:false` is unconditional. `resumeIndex` is a transport position + when the tail is known, **not** proof that every preceding message is present, + and is omitted entirely when the tail cannot be captured. Consumers must not upload + these rows as a full transcript or use them as model seed/history. +- No transcript/envelope writes, model/tool execution, run start or run cancel + occur during recovery. SDK failure or disconnect detaches the reader only. + A stored `done`/`error` on the negotiated live tail closes that reader + **without** a following `viewport_end`; EOF, `viewport_end`, and + `viewport_error` are the other terminal signals (see + [agent-stream.md](agent-stream.md)). + +The worker's existing durable transcript/model projections remain unchanged. +See [agent-stream.md](agent-stream.md) for version negotiation and record grammar, +and [harness-limits.md](harness-limits.md) for the bounded-work caveats. + ## Local session (always) | Piece | Location | Notes | diff --git a/int/viewport-read.int.test.ts b/int/viewport-read.int.test.ts new file mode 100644 index 00000000..83a63f10 --- /dev/null +++ b/int/viewport-read.int.test.ts @@ -0,0 +1,95 @@ +import { describe, expect, it, vi } from 'vitest'; +import { loadBridge } from './loadBridge'; +import { ringTexts } from './driver'; +import { MessageKind, Lifecycle } from '../lib/harnessBridge'; +import { MemoryBlobTranscriptStore } from '../lib/sessions/blobStores'; +import { newBlobObjectId } from '../lib/sessions/blobStore'; +import { readViewportHead } from '../lib/sessions/viewportRead'; +import { createViewportRunReader, type ViewportRun } from '../lib/workflows/viewportRunReader'; +import { viewportStream } from '../lib/agent/viewportStream'; +import { ViewportStreamDecoder } from '../lib/viewportStreamProtocol'; +import { VIEWPORT_HISTORY_NOTE } from '../lib/sessionViewport'; + +const scope = { tenantId: 'int_tenant', userId: 'int_user', sessionId: 'int_session' }; +const line = (event: object) => `data: ${JSON.stringify(event)}\n\n`; + +describe('bounded recovery → protocol → real Wasm display (backend foundation)', () => { + it('one-head/late-tail recovery never paints historical thinking; subsequent live text/thinking survives', async () => { + const bridge = await loadBridge(); + bridge.pushMessage(MessageKind.Assistant, 'cached tail must not rewind to prompt'); + const blob = new MemoryBlobTranscriptStore(), pointer = newBlobObjectId(scope), prev = newBlobObjectId(scope); + await blob.writeSegment({ objectId: pointer, maxBytes: 8 * 1024 * 1024, + content: JSON.stringify({ id: scope.sessionId, prev, messages: [{ role: 'assistant', text: 'older head' }] }) }); + const read = vi.spyOn(blob, 'read'), write = vi.spyOn(blob, 'writeSegment'); + const opens: number[] = []; + const H0 = 100000, H = H0 + 7; + let probes = 0; + const run: ViewportRun = { status: 'running', getReadable(opts) { + const start = opts?.startIndex ?? 0; opens.push(start); + let i = start; + return Object.assign(new ReadableStream({ pull(c) { + if (start === -1) return; // public tail metadata handle, cancelled by adapter + if (start < H0) { + if (i >= H0) { c.close(); return; } + c.enqueue(line(i++ === H0 - 1 ? { type: 'text_delta', text: 'most recent sampled assistant' } + : { type: 'reasoning_delta', text: 'HISTORICAL_THINKING_MUST_NOT_PAINT' })); + } else { + const events = [ + { type: 'reasoning_delta', text: 'new live thinking' }, + { type: 'text_delta', text: 'new live assistant' }, + { type: 'done', text: 'aggregate (not appended by int consumer)' }, + ]; + if (i - H >= events.length) { c.close(); return; } + c.enqueue(line(events[i++ - H])); + } + } }), { getTailIndex: async () => (probes++ === 0 ? H0 : H) - 1 }); + } }; + // Same composition as GET hydrate=tail: startIndex 0, H0 captured inside viewportStream. + const stream = viewportStream({ runId: 'run', sessionId: scope.sessionId, run: createViewportRunReader(run), startIndex: 0, + cold: { status: 'cancelling', readHead: deadline => readViewportHead({ scope, meta: { transcriptPointer: pointer }, blob, deadline }) } }); + const parser = new ViewportStreamDecoder('run'), reader = stream.getReader(); + let sawSnapshot = false, nextIndex = 0; + for (;;) { + const { value, done } = await reader.read(); if (done) break; + for (const record of parser.push(value)) { + if (record.type === 'viewport_state') { + bridge.setLifecycle(Lifecycle.Busy); + expect(ringTexts(bridge)).toContain('cached tail must not rewind to prompt'); + } else if (record.type === 'viewport_snapshot') { + sawSnapshot = true; + if (record.resumeIndex !== undefined) nextIndex = record.resumeIndex; + expect(record.gap).toBe(true); expect(record.historyComplete).toBe(false); + if (record.replace) bridge.hydrateMessages(record.rows.map(row => ({ + kind: row.role === 'assistant' ? MessageKind.Assistant : row.role === 'tool_run' ? MessageKind.ToolRun : MessageKind.System, + text: row.text, + }))); + expect(ringTexts(bridge)).toContain('most recent sampled assistant'); + expect(ringTexts(bridge)).toContain(VIEWPORT_HISTORY_NOTE); + } else if (record.type === 'turn_event') { + expect(record.nextIndex).toBe(nextIndex + 1); nextIndex = record.nextIndex; + if (record.event?.type === 'reasoning_delta') bridge.pushMessage(MessageKind.Thinking, record.event.text); + if (record.event?.type === 'text_delta') bridge.pushMessage(MessageKind.Assistant, record.event.text); + } + expect(ringTexts(bridge).join('\n')).not.toContain('HISTORICAL_THINKING'); + expect(bridge.messageCount()).toBeLessThanOrEqual(2048); + } + } + parser.push(new Uint8Array(), true); + expect(sawSnapshot).toBe(true); expect(nextIndex).toBe(H + 3); + expect(ringTexts(bridge)).toContain('new live thinking'); + expect(ringTexts(bridge)).toContain('new live assistant'); + expect(opens).toContain(H0 - 2048); expect(opens).toContain(H); expect(opens).not.toContain(0); + expect(read).toHaveBeenCalledExactlyOnceWith(pointer); expect(write).not.toHaveBeenCalled(); + }); + + it('head-only fallback paints the latest real ring without reading its ancestors', async () => { + const bridge = await loadBridge(), blob = new MemoryBlobTranscriptStore(), pointer = newBlobObjectId(scope); + await blob.writeSegment({ objectId: pointer, maxBytes: 8 * 1024 * 1024, content: JSON.stringify({ id: scope.sessionId, + prev: newBlobObjectId(scope), messages: Array.from({ length: 3000 }, (_, i) => ({ role: 'assistant', text: `row-${i}` })) }) }); + const read = vi.spyOn(blob, 'read'); + const view = await readViewportHead({ scope, meta: { transcriptPointer: pointer }, blob }); + bridge.hydrateMessages(view.rows.map(r => ({ kind: MessageKind.Assistant, text: r.text }))); + expect(bridge.messageCount()).toBe(2048); expect(ringTexts(bridge).at(-1)).toBe('row-2999'); + expect(view.historyComplete).toBe(false); expect(view.hasEarlier).toBe(true); expect(read).toHaveBeenCalledOnce(); + }); +}); diff --git a/lib/agent/viewportStream.test.ts b/lib/agent/viewportStream.test.ts new file mode 100644 index 00000000..24c541e8 --- /dev/null +++ b/lib/agent/viewportStream.test.ts @@ -0,0 +1,206 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { viewportStream } from './viewportStream'; +import { emptyViewport } from '../sessions/viewportRead'; +import { VIEWPORT_HISTORY_NOTE } from '../sessionViewport'; +import { VIEWPORT_RECOVERY_MAX_MS, VIEWPORT_TAIL_MAX_FRAMES } from '../sessionCloudCaps'; +import { ViewportStreamDecoder, type ViewportRecord } from '../viewportStreamProtocol'; +import type { ViewportRunReader } from '../workflows/viewportRunReader'; +const line = (e: object) => `data: ${JSON.stringify(e)}\n\n`; +async function collect(stream: ReadableStream, initial?:number) { + const reader = stream.getReader(), parser = new ViewportStreamDecoder('run',initial); + const records: ViewportRecord[] = []; + for (;;) { const {done,value}=await reader.read(); if(done)break; records.push(...parser.push(value)); } + parser.push(new Uint8Array(),true); return records; +} +afterEach(()=>vi.useRealTimers()); + +describe('snapshot-first and indexed live transport',()=>{ + it('emits state/snapshot, skips historical thinking/backlog and paints only post-H live events',async()=>{ + const run:ViewportRunReader={status:async()=> 'running',nextIndex:async()=>100, + open:vi.fn(start=>new ReadableStream({start(c){ + const events=start===0 ? [{type:'reasoning_delta',text:'old-secret-think'},{type:'text_delta',text:'sampled tail'}] + :[{type:'reasoning_delta',text:'new live thinking'},{type:'text_delta',text:'new text'},{type:'done',text:'all text'}]; + for(const e of events)c.enqueue(line(e)); c.close(); + }}))}; + const records=await collect(viewportStream({runId:'run',sessionId:'session',startIndex:2,run, + cold:{status:'cancelling',readHead:async()=>emptyViewport('session')}})); + expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_snapshot','turn_event','turn_event','turn_event']); + expect(records[0]).toMatchObject({status:'cancelling'}); + expect(records[1]).toMatchObject({resumeIndex:100,gap:true,source:'stream_tail'}); + expect(JSON.stringify(records[1])).toContain(VIEWPORT_HISTORY_NOTE); + expect(JSON.stringify(records)).not.toContain('old-secret-think'); + expect(records[2]).toMatchObject({nextIndex:101,event:{type:'reasoning_delta',text:'new live thinking'}}); + expect(records[4]).toMatchObject({type:'turn_event',event:{type:'done'}}); + expect(records.some(r=>r.type==='viewport_end')).toBe(false); + expect(run.open).toHaveBeenNthCalledWith(1,0); expect(run.open).toHaveBeenNthCalledWith(2,100); + }); + it('stored error is producer-terminal and does not synthesize viewport_end', async () => { + const run: ViewportRunReader = { + status: async () => 'running', nextIndex: async () => 0, + open: () => new ReadableStream({ start(c) { + c.enqueue(line({ type: 'error', error: 'Request cancelled.' })); c.close(); + } }), + }; + const records = await collect(viewportStream({ runId: 'run', sessionId: 'session', run, startIndex: 0 }), 0); + expect(records).toEqual([{ + type: 'turn_event', version: 1, runId: 'run', nextIndex: 1, + event: { type: 'error', error: 'Request cancelled.' }, + }]); + }); + it('known malformed frame is skipped at its raw index; EOF synthesizes no stored index',async()=>{ + const run:ViewportRunReader={status:async()=> 'completed',nextIndex:async()=>0,open:()=>new ReadableStream({start(c){c.enqueue('broken');c.close();}})}; + const records=await collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:9}),9); + expect(records).toEqual([{type:'turn_event',version:1,runId:'run',nextIndex:10,skipped:true},{type:'viewport_end',version:1,runId:'run',status:'completed'}]); + }); + it('invalid UTF-8 means unknown position and closes with error, not skipped index',async()=>{ + const run:ViewportRunReader={status:async()=> 'running',nextIndex:async()=>0,open:()=>new ReadableStream({start(c){c.enqueue(new Uint8Array([255]));}})}; + expect(await collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:9}),9)) + .toEqual([{type:'viewport_error',version:1,runId:'run',code:'STREAM_UNAVAILABLE'}]); + }); + it('session revoked before handoff cannot release snapshot or live data',async()=>{ + const run:ViewportRunReader={status:async()=> 'running',nextIndex:async()=>0,open:vi.fn()}; + const records=await collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:0, + cold:{status:'running',readHead:async()=>emptyViewport('session')},stillOwned:async()=>false})); + expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_error']);expect(run.open).not.toHaveBeenCalled(); + }); + it('cancel/disconnect unblocks pending read and cleans timers without cancelling the run',async()=>{ + vi.useFakeTimers();const cancel=vi.fn();const run:ViewportRunReader={status:vi.fn(async()=> 'running'),nextIndex:async()=>0, + open:()=>new ReadableStream({cancel})}; + const controller=new AbortController(); + const result=collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:0,signal:controller.signal}),0); + await vi.advanceTimersByTimeAsync(0); + controller.abort(); + expect(await result).toEqual([]);expect(cancel).toHaveBeenCalledOnce();expect(vi.getTimerCount()).toBe(0); + }); + it('shares a stalled status probe across frames and EOF instead of piling up requests', async () => { + vi.useFakeTimers(); + let source!: ReadableStreamDefaultController; + const status = vi.fn(() => new Promise(() => {})); + const run: ViewportRunReader = { + status, nextIndex: async () => 0, + open: () => new ReadableStream({ start(controller) { source = controller; } }), + }; + const result = collect(viewportStream({ runId: 'run', sessionId: 'session', run, startIndex: 0 }), 0); + await vi.advanceTimersByTimeAsync(1000); + expect(status).toHaveBeenCalledOnce(); + source.enqueue(line({ type: 'text_delta', text: 'late frame' })); + await vi.advanceTimersByTimeAsync(1000); + expect(status).toHaveBeenCalledOnce(); + source.close(); + await vi.advanceTimersByTimeAsync(1000); + expect((await result).map(record => record.type)).toEqual(['turn_event', 'viewport_error']); + expect(status).toHaveBeenCalledOnce(); + expect(vi.getTimerCount()).toBe(0); + }); + it.each(['completed'])('polls a hung readable and emits synthetic %s with no fake index',async status=>{ + vi.useFakeTimers();const cancel=vi.fn();const run:ViewportRunReader={status:async()=>status,nextIndex:async()=>0, + open:()=>new ReadableStream({cancel})}; + const result=collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:42}),42); + await vi.advanceTimersByTimeAsync(1000); + expect(await result).toEqual([{type:'viewport_end',version:1,runId:'run',status}]); + expect(cancel).toHaveBeenCalledOnce();expect(vi.getTimerCount()).toBe(0); + }); + it('completed + buffered frames drain; 0-delay poll does not inject completed ahead of a later chunk', async () => { + let pulls = 0; + const run: ViewportRunReader = { + status: async () => 'completed', nextIndex: async () => 0, + open: () => new ReadableStream({ + pull(c) { + pulls++; + if (pulls === 1) { c.enqueue(line({ type: 'text_delta', text: 'Hi' })); return; } + if (pulls === 2) { + return new Promise(resolve => { + setTimeout(() => { + c.enqueue(line({ type: 'text_delta', text: ' there' })); + c.close(); + resolve(); + }, 20); + }); + } + }, + }), + }; + const records = await collect(viewportStream({ runId: 'run', sessionId: 'session', run, startIndex: 0 }), 0); + expect(records.map(r => r.type)).toEqual(['turn_event', 'turn_event', 'viewport_end']); + expect(records[0]).toMatchObject({ event: { type: 'text_delta', text: 'Hi' } }); + expect(records[1]).toMatchObject({ event: { type: 'text_delta', text: ' there' } }); + expect(records[2]).toMatchObject({ status: 'completed' }); + }); + it.each(['failed','cancelled'])('already-%s never opens getReadable and emits viewport_end',async status=>{ + const open=vi.fn();const nextIndex=vi.fn(); + const run:ViewportRunReader={status:async()=>status,nextIndex,open}; + expect(await collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:42}),42)) + .toEqual([{type:'viewport_end',version:1,runId:'run',status}]); + expect(open).not.toHaveBeenCalled();expect(nextIndex).not.toHaveBeenCalled(); + }); + it.each(['failed','cancelled'])('cold already-%s is head-only: no SDK sample, snapshot then viewport_end',async status=>{ + const open=vi.fn();const nextIndex=vi.fn(async()=>9); + const run:ViewportRunReader={status:async()=>status,nextIndex,open}; + const head={...emptyViewport('session'),source:'stored_head' as const,replace:true, + rows:[{id:'h',role:'assistant' as const,text:'head-only',at:0}]}; + const records=await collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:4, + cold:{status,readHead:async()=>head}})); + expect(open).not.toHaveBeenCalled();expect(nextIndex).not.toHaveBeenCalled(); + expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_snapshot','viewport_end']); + expect(records[1]).toMatchObject({source:'stored_head',gap:true}); + expect(records[1]).not.toHaveProperty('resumeIndex'); + expect(JSON.stringify(records)).toContain('head-only'); + expect(records[2]).toMatchObject({status}); + }); + it('cold unavailable H0 emits state and head, never opens origin live',async()=>{ + const open=vi.fn(); + const run:ViewportRunReader={status:async()=>'running',nextIndex:async()=>{throw new Error('unavailable');},open}; + const head={...emptyViewport('session'),source:'stored_head' as const,replace:true, + rows:[{id:'h',role:'assistant' as const,text:'head-only',at:0}]}; + const records=await collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:0, + cold:{status:'running',readHead:async()=>head}})); + expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_snapshot','viewport_error']); + expect(records[1]).toMatchObject({source:'stored_head'}); + expect(records[1]).not.toHaveProperty('resumeIndex'); + expect(JSON.stringify(records)).toContain('head-only'); + expect(open).not.toHaveBeenCalled(); + }); + it('cold hanging H0 emits viewport_state before the recovery deadline and never origin-replays',async()=>{ + vi.useFakeTimers(); + const open=vi.fn(); + const run:ViewportRunReader={status:async()=>'running',nextIndex:()=>new Promise(()=>{}),open}; + const result=collect(viewportStream({runId:'run',sessionId:'session',run,startIndex:0, + cold:{status:'running',readHead:async()=>emptyViewport('session')}})); + await vi.advanceTimersByTimeAsync(0); + await vi.advanceTimersByTimeAsync(VIEWPORT_RECOVERY_MAX_MS); + const records=await result; + expect(records[0]).toMatchObject({type:'viewport_state',phase:'recovering'}); + expect(records.map(r=>r.type)).toEqual(['viewport_state','viewport_snapshot','viewport_error']); + expect(records[1]).not.toHaveProperty('resumeIndex'); + expect(open).not.toHaveBeenCalled();expect(vi.getTimerCount()).toBe(0); + }); + it('samples from the first H0, not a later final-probe tail', async () => { + let probes = 0; + const H0 = 5000, H = 5070; + const run: ViewportRunReader = { + status: async () => 'running', + nextIndex: async () => (probes++ === 0 ? H0 : H), + open: vi.fn(start => new ReadableStream({ start(c) { + if (start === H0 - VIEWPORT_TAIL_MAX_FRAMES) { + c.enqueue(line({ type: 'text_delta', text: 'sampled-from-h0' })); + c.close(); + } else if (start === H) { + c.enqueue(line({ type: 'text_delta', text: 'live-after-probe' })); + c.close(); + } else { c.close(); } + } })), + }; + const records = await collect(viewportStream({ + runId: 'run', sessionId: 'session', startIndex: 0, run, + cold: { status: 'running', readHead: async () => emptyViewport('session') }, + })); + expect(run.open).toHaveBeenCalledWith(H0 - VIEWPORT_TAIL_MAX_FRAMES); + expect(run.open).toHaveBeenCalledWith(H); + expect(run.open).not.toHaveBeenCalledWith(0); + expect(run.open).not.toHaveBeenCalledWith(H - VIEWPORT_TAIL_MAX_FRAMES); + expect(records[1]).toMatchObject({ type: 'viewport_snapshot', resumeIndex: H, source: 'stream_tail' }); + expect(JSON.stringify(records)).toContain('sampled-from-h0'); + expect(records.at(-2)).toMatchObject({ type: 'turn_event', nextIndex: H + 1, event: { type: 'text_delta', text: 'live-after-probe' } }); + expect(probes).toBe(2); + }); +}); diff --git a/lib/agent/viewportStream.ts b/lib/agent/viewportStream.ts new file mode 100644 index 00000000..21b8a628 --- /dev/null +++ b/lib/agent/viewportStream.ts @@ -0,0 +1,167 @@ +/** Negotiated snapshot-first/indexed stream. Default legacy transport remains untouched. */ +import { TURN_STREAM_STATUS_POLL_MS, sanitizeTurnStreamCursor, VIEWPORT_RECOVERY_MAX_MS } from '../sessionCloudCaps'; +import { parseViewportEvent, type ViewportView } from '../sessionViewport'; +import { encodeViewportRecord, type ViewportRecord } from '../viewportStreamProtocol'; +import { recoverViewport } from '../sessions/viewportRead'; +import { viewportWait, type ViewportRunReader } from '../workflows/viewportRunReader'; +import { isTerminalRunStatus } from './pipeRunReadable'; + +type FrameResult = { kind: 'frame'; result: ReadableStreamReadResult } | { kind: 'terminal'; status: string }; + +function isHangClassStatus(status: string | undefined): status is 'cancelled' | 'failed' { + return status === 'cancelled' || status === 'failed'; +} + +/** One pending raw read, one pending status probe at most; stop all timers/listeners on settle. */ +function nextFrame(reader: ReadableStreamDefaultReader, statusProbe: () => Promise, signal: AbortSignal): Promise { + return new Promise((resolve, reject) => { + let settled = false; + let timer: ReturnType | undefined; + const finish = (fn: () => void) => { + if (settled) return; + settled = true; clearTimeout(timer); signal.removeEventListener('abort', abort); fn(); + }; + const abort = () => finish(() => reject(new Error('Viewport detached'))); + const consider = (status: string, allowCompleted: boolean): boolean => { + if (isHangClassStatus(status) || (allowCompleted && status === 'completed')) { + finish(() => resolve({ kind: 'terminal', status })); + return true; + } + return false; + }; + const poll = (allowCompleted: boolean) => { + if (settled) return; + // No overlapping status calls, even if the provider never settles. + void statusProbe().then(status => { + if (settled) return; + if (!consider(status, allowCompleted)) timer = setTimeout(() => poll(true), TURN_STREAM_STATUS_POLL_MS); + }, () => { if (!settled) timer = setTimeout(() => poll(true), TURN_STREAM_STATUS_POLL_MS); }); + }; + signal.addEventListener('abort', abort, { once: true }); + void reader.read().then(result => finish(() => resolve({ kind: 'frame', result })), error => finish(() => reject(error))); + // 0-delay: hang-class cancelled/failed only (unstick C16 hung readables). + // completed drains buffered/in-flight frames; synthetic end after a 1s hung + // read or EOF — same discipline as pipeRunReadable. + timer = setTimeout(() => poll(false), 0); + if (signal.aborted) abort(); + }); +} + +export function viewportStream(opts: { + runId: string; sessionId: string; run: ViewportRunReader; startIndex: number; + cold?: { readHead: (deadline: number) => Promise; status: string }; + signal?: AbortSignal; + /** Recheck session ownership before releasing the recovered snapshot/live stream. */ + stillOwned?: () => Promise; +}): ReadableStream { + const aborter = new AbortController(); + const onAbort = () => aborter.abort(); + opts.signal?.addEventListener('abort', onAbort, { once: true }); + if (opts.signal?.aborted) onAbort(); + let reader: ReadableStreamDefaultReader | undefined; + // Share a stalled provider probe across raw reads; a late frame must not + // cause the next pull to launch another metadata request beside it. + let pendingStatus: Promise | undefined; + const statusProbe = (): Promise => { + if (!pendingStatus) { + pendingStatus = Promise.resolve().then(() => opts.run.status()); + void pendingStatus.then(() => { pendingStatus = undefined; }, () => { pendingStatus = undefined; }); + } + return pendingStatus; + }; + const cleanup = () => { + aborter.abort(); opts.signal?.removeEventListener('abort', onAbort); + if (reader) void reader.cancel().catch(() => {}); + }; + async function* records(): AsyncGenerator { + let index = opts.startIndex; + try { + if (aborter.signal.aborted) return; + // Same C16 gate as bodyForRun: already-cancelled/failed never touches getReadable. + const liveStatus = await viewportWait(statusProbe(), TURN_STREAM_STATUS_POLL_MS, aborter.signal).catch(() => undefined); + const skipReadable = isHangClassStatus(liveStatus); + if (opts.cold) { + const stateStatus = isHangClassStatus(liveStatus) ? liveStatus + : liveStatus === 'running' && opts.cold.status === 'cancelling' ? 'cancelling' + : (liveStatus ?? opts.cold.status); + yield { type: 'viewport_state', version: 1, runId: opts.runId, status: stateStatus, phase: 'recovering' }; + const deadline = Date.now() + VIEWPORT_RECOVERY_MAX_MS; + const headInFlight = opts.cold.readHead(deadline); + let initialIndex = opts.startIndex; + let skipStream = skipReadable; + let h0Failed = false; + let resumeKnown = false; + if (!skipStream) { + try { + initialIndex = await viewportWait(opts.run.nextIndex(), Math.max(0, deadline - Date.now()), aborter.signal); + resumeKnown = true; + } catch { + skipStream = true; + h0Failed = true; + } + } + const snapshot = await recoverViewport({ runId: opts.runId, sessionId: opts.sessionId, + initialIndex, run: opts.run, readHead: async () => headInFlight, signal: aborter.signal, + skipStream, deadline }); + if (opts.stillOwned && !(await viewportWait(opts.stillOwned(), TURN_STREAM_STATUS_POLL_MS, aborter.signal))) + throw new Error('Viewport session changed'); + if (!resumeKnown) delete snapshot.resumeIndex; + else index = snapshot.resumeIndex ?? initialIndex; + yield { type: 'viewport_snapshot', ...snapshot }; + if (skipReadable) { + yield { type: 'viewport_end', version: 1, runId: opts.runId, status: liveStatus }; + return; + } + if (h0Failed) { + yield { type: 'viewport_error', version: 1, runId: opts.runId, code: 'STREAM_UNAVAILABLE' }; + return; + } + } else if (skipReadable) { + yield { type: 'viewport_end', version: 1, runId: opts.runId, status: liveStatus }; + return; + } + if (aborter.signal.aborted) return; + reader = opts.run.open(index).getReader(); + for (;;) { + const next = await nextFrame(reader, statusProbe, aborter.signal); + if (next.kind === 'terminal') { + yield { type: 'viewport_end', version: 1, runId: opts.runId, status: next.status }; + return; + } + if (next.result.done) { + const status = await viewportWait(statusProbe(), TURN_STREAM_STATUS_POLL_MS, aborter.signal).catch(() => 'unknown'); + if (isTerminalRunStatus(status)) yield { type: 'viewport_end', version: 1, runId: opts.runId, status }; + else yield { type: 'viewport_error', version: 1, runId: opts.runId, code: 'STREAM_ENDED' }; + return; + } + const value = next.result.value; + const text = typeof value === 'string' ? value : new TextDecoder('utf-8', { fatal: true }).decode(value); + const event = parseViewportEvent(text); + const nextIndex = sanitizeTurnStreamCursor(index + 1); + if (nextIndex === undefined) throw new Error('Viewport cursor exhausted'); + index = nextIndex; + yield event ? { type: 'turn_event', version: 1, runId: opts.runId, nextIndex, event } + : { type: 'turn_event', version: 1, runId: opts.runId, nextIndex, skipped: true }; + // Stored done/error are producer-terminal. Do not synthesize viewport_end + // (cancelled inject is an error event, not status failed). + if (event?.type === 'done' || event?.type === 'error') return; + } + } catch { + if (!aborter.signal.aborted) yield { type: 'viewport_error', version: 1, runId: opts.runId, code: 'STREAM_UNAVAILABLE' }; + } finally { cleanup(); } + } + const iterator = records(); + let closed = false; + return new ReadableStream({ + async pull(controller) { + const next = await iterator.next(); + if (closed) return; + if (next.done) { closed = true; controller.close(); } + else controller.enqueue(encodeViewportRecord(next.value)); + }, + cancel() { + closed = true; cleanup(); + void iterator.return(undefined).catch(() => {}); + }, + }); +} \ No newline at end of file diff --git a/lib/sessionCloudCaps.ts b/lib/sessionCloudCaps.ts index 0bb580ad..9f9e6563 100644 --- a/lib/sessionCloudCaps.ts +++ b/lib/sessionCloudCaps.ts @@ -39,6 +39,15 @@ export const HARNESS_SESSION_MAX_FUNCTION_BODY_BYTES = 2 * 1024 * 1024; */ export const HARNESS_SESSION_MAX_BODY_BYTES = 8 * 1024 * 1024; +/** Best-effort viewport recovery only; none of these change durable retention. */ +export const VIEWPORT_TAIL_MAX_FRAMES = 2048; +export const VIEWPORT_RECOVERY_MAX_BYTES = 8 * 1024 * 1024; +export const VIEWPORT_RECOVERY_MAX_MS = 5000; +export const VIEWPORT_FINAL_PROBE_MAX_MS = 1000; +export const VIEWPORT_HEAD_READ_MAX_OBJECTS = 1; +/** Entire escaped JSON snapshot/control, not the lifetime of a live SSE stream. */ +export const VIEWPORT_RESPONSE_MAX_BYTES = HARNESS_SESSION_MAX_FUNCTION_BODY_BYTES; + /** * Max objects a transcript `prev` walk may visit (plan #886). Loop/DoS bound * on reconstruct **and** worker persist — not a message cap and not a change diff --git a/lib/sessionViewport.test.ts b/lib/sessionViewport.test.ts new file mode 100644 index 00000000..2dff0c57 --- /dev/null +++ b/lib/sessionViewport.test.ts @@ -0,0 +1,101 @@ +import { describe, expect, it } from 'vitest'; +import { byteLength, fitViewportRows, parseViewportCarriers, parseViewportEvent, parseViewportRow, viewportCarriers, viewportExcerpt, ViewportReducer } from './sessionViewport'; +import { HARNESS_SESSION_MAX_MSG_BYTES, VIEWPORT_RESPONSE_MAX_BYTES } from './sessionCloudCaps'; +import { decodeToolRun } from './toolRun'; + +describe('bounded disposable viewport', () => { + it('keeps the newest UTF-8 excerpt within the existing bridge rail', () => { + const text = '漢🙂'.repeat(100000) + ' latest'; + const excerpt = viewportExcerpt(text); + expect(excerpt.startsWith('[Earlier text omitted]')).toBe(true); + expect(excerpt.endsWith(' latest')).toBe(true); + expect(excerpt).not.toContain('\ufffd'); + expect(byteLength(excerpt)).toBeLessThanOrEqual(HARNESS_SESSION_MAX_MSG_BYTES); + expect(viewportExcerpt('small')).toBe('small'); + }); + it('fits complete escaped JSON including scaffold and keeps newest rows without mutation', () => { + const rows = Array.from({ length: 2048 }, (_, i) => ({ id: `row_${i}`, role: 'assistant' as const, + text: '\u0001"漢'.repeat(500), at: i })); + const scaffold = { rows: [], source: 'stored_head', carriers: { queue: ['"'.repeat(5000)] } }; + const kept = fitViewportRows(rows, scaffold); + expect(byteLength(JSON.stringify({ ...scaffold, rows: kept }))).toBeLessThanOrEqual(VIEWPORT_RESPONSE_MAX_BYTES); + expect(kept.at(-1)?.id).toBe('row_2047'); + expect(kept.length).toBeLessThan(rows.length); + expect(rows).toHaveLength(2048); + }); + it('rejects invalid rows and excludes arbitrary meta, keys and model bodies', () => { + expect(parseViewportRow({ role: 'reasoning', text: 'secret reasoning' }, 0)).toBeUndefined(); + expect(parseViewportRow({ role: 'assistant', text: 123 }, 0)).toBeUndefined(); + const carriers = viewportCarriers({ logicalCwd: 'src', selectedModel: 'test/model', workingNotes: 'private notes', + personaSnapshot: 'persona body', compactionPointer: 'private', randomSecret: 'key', + turnRunId: 'run_1', activeSandboxId: '*', usage: 'bad-json' }, ['next']); + expect(carriers.cwd).toBe('src'); + expect(carriers.queue).toEqual(['next']); + expect(JSON.stringify(carriers)).not.toMatch(/private|persona body|randomSecret/); + expect(carriers.activeSandboxId).toBeUndefined(); + const encoded = parseViewportCarriers({ + cwd: carriers.cwd, selectedModel: 'test/model', queue: ['next'], + workingNotes: 'private notes', personaSnapshot: 'persona body', + cwdHost: '/etc/passwd', + }); + expect(encoded).toMatchObject({ cwd: 'src', selectedModel: 'test/model', queue: ['next'] }); + expect(JSON.stringify(encoded)).not.toMatch(/private|persona body|\/etc/); + expect(parseViewportCarriers({ cwd: '/etc/passwd', workingNotes: 'secret' }).cwd).toBeUndefined(); + expect(JSON.stringify(parseViewportCarriers({ cwd: '/etc/passwd', workingNotes: 'secret' }))).not.toContain('secret'); + }); + it('parses one stored frame, ignores unknown fields, rejects multiple/malformed frames', () => { + expect(parseViewportEvent('data: {"type":"text_delta","text":"hi","secret":"hidden"}\r\n\r\n')) + .toEqual({ type: 'text_delta', text: 'hi' }); + expect(parseViewportEvent('data: nope\n\n')).toBeUndefined(); + expect(parseViewportEvent('data: {"type":"new_unrecognized"}\n\n')).toBeUndefined(); + expect(parseViewportEvent('data: {"type":"text_delta","text":"a"}\n\ndata: {"type":"text_delta","text":"b"}\n\n')).toBeUndefined(); + expect(parseViewportEvent('data: {"type":"skill_attached","slug":"test","action":"attach","ok":true}\n\n')?.type).toBe('skill_attached'); + }); + it('never retains historical thinking or appends aggregate done over text', () => { + const reducer = new ViewportReducer(); + reducer.apply({ type: 'reasoning_delta', text: 'historical thinking' }); + reducer.apply({ type: 'text_delta', text: 'a ' }); + reducer.apply({ type: 'text_delta', text: 'b' }); + reducer.apply({ type: 'done', text: 'earlier a b' }); + expect(reducer.snapshot().map(r => r.text)).toEqual(['a b']); + expect(JSON.stringify(reducer)).not.toContain('historical thinking'); + const across = new ViewportReducer(); + across.apply({ type: 'text_delta', text: 'Hello' }); + across.apply({ type: 'reasoning_delta', text: 'hidden thinking' }); + across.apply({ type: 'text_delta', text: ' world' }); + expect(across.snapshot().map(r => r.text)).toEqual(['Hello world']); + expect(JSON.stringify(across)).not.toContain('hidden thinking'); + const skills = new ViewportReducer(); + skills.apply({ type: 'skill_attached', slug: 'create-plan', action: 'attach', ok: true }); + skills.apply({ type: 'text_delta', text: 'after skill' }); + expect(skills.snapshot().map(r => ({ role: r.role, text: r.text }))).toEqual([ + { role: 'skill_attached', text: 'Skill attached: create-plan' }, + { role: 'assistant', text: 'after skill' }, + ]); + skills.apply({ type: 'skill_attached', slug: 'create-plan', action: 'detach', ok: true }); + expect(skills.snapshot().at(-1)).toMatchObject({ role: 'skill_attached', text: 'Skill detached: create-plan' }); + }); + it('pairs tool ids within the visible group but does not name-pair an unrelated result', () => { + const reducer = new ViewportReducer(); + reducer.apply({ type: 'tool_start', name: 'read', id: 'a' }); + reducer.apply({ type: 'tool_start', name: 'read', id: 'b' }); + reducer.apply({ type: 'tool_result', name: 'read', id: 'b', ok: true, summary: 'result b' }); + const group = decodeToolRun(reducer.snapshot()[0].text)!; + expect(group.pending).toBe(1); expect(group.ok).toBe(1); + reducer.apply({ type: 'tool_result', name: 'read', id: 'unknown', ok: true, summary: 'orphan' }); + expect(reducer.snapshot()).toHaveLength(2); + expect(decodeToolRun(reducer.snapshot()[1].text)?.ok).toBe(1); + }); + it('rolls groups and retains bounded state for very long text and many rows', () => { + const reducer = new ViewportReducer(); + for (let i = 0; i < 2500; i++) { + reducer.apply({ type: 'text_delta', text: 'x'.repeat(1024) }); + reducer.apply({ type: 'error', error: `err${i}` }); + } + expect(reducer.snapshot().length).toBeLessThanOrEqual(2048); + expect(byteLength(JSON.stringify(reducer.snapshot()))).toBeLessThanOrEqual(VIEWPORT_RESPONSE_MAX_BYTES + 2); + const last = new ViewportReducer(); + last.apply({ type: 'done', text: 'x'.repeat(400000) + 'newest' }); + expect(last.snapshot()[0].text.endsWith('newest')).toBe(true); + }); +}); diff --git a/lib/sessionViewport.ts b/lib/sessionViewport.ts new file mode 100644 index 00000000..701d3b1c --- /dev/null +++ b/lib/sessionViewport.ts @@ -0,0 +1,206 @@ +/** Disposable display windows. Never use these rows as a transcript replacement or model seed. */ +import type { AgentStreamEvent } from './agent/agentStream'; +import type { SessionMessage, SessionRole } from './sessionStore'; +import { sanitizeUsageSummary, decodeUsageMetaString } from './agent/usageSummary'; +import { + HARNESS_SESSION_MAX_MSG_BYTES, VIEWPORT_RESPONSE_MAX_BYTES, + isRedisSafeOpaqueId, normalizeSessionCwd, parseAttachedSkills, + sanitizeModelId, sanitizeReasoningEffort, sanitizeResolvedProvider, + sanitizeTurnRunId, sanitizeTurnStatus, HARNESS_SESSION_MAX_META_BYTES, + HARNESS_SESSION_MAX_ATTACHED_SKILLS, SKILL_SLUG_RE, +} from './sessionCloudCaps'; +import { HARNESS_RING_MAX } from './sessionWindow'; +import { sanitizeQueue } from './turnQueue'; +import { addToolStart, addToolResult, createToolRunGroup, encodeToolRun, toolRunIsFull } from './toolRun'; + +const encoder = new TextEncoder(); +export const VIEWPORT_HISTORY_NOTE = 'Recent history only; some earlier activity may be omitted.'; +export const byteLength = (s: string): number => encoder.encode(s).byteLength; +export const objectRecord = (v: unknown): Record | undefined => + v !== null && typeof v === 'object' && !Array.isArray(v) ? v as Record : undefined; + +/** Bound before encoding; keep newest text and don't split UTF-8 code points. */ +export function viewportExcerpt(text: string, max = HARNESS_SESSION_MAX_MSG_BYTES): string { + const suffix = text.slice(-max); + const bytes = encoder.encode(suffix); + if (text.length === suffix.length && bytes.length <= max) return text; + const marker = '[Earlier text omitted]\n'; + let start = Math.max(0, bytes.length - Math.max(0, max - byteLength(marker))); + while (start < bytes.length && (bytes[start] & 0xc0) === 0x80) start++; + return marker + new TextDecoder().decode(bytes.subarray(start)); +} + +const roles = new Set(['user', 'assistant', 'system', 'error', 'tool_run', 'skill_attached']); +export function parseViewportRow(value: unknown, index: number): SessionMessage | undefined { + const o = objectRecord(value); + if (!o || !roles.has(o.role as SessionRole) || typeof o.text !== 'string') return; + return { + id: isRedisSafeOpaqueId(o.id) ? o.id : `view_${index}`, + role: o.role as SessionRole, text: viewportExcerpt(o.text), + at: typeof o.at === 'number' && Number.isFinite(o.at) ? o.at : 0, + }; +} + +export type ViewportCarriers = { + cwd?: string; activeSandboxId?: string; selectedModel?: string; reasoningEffort?: string; + resolvedProvider?: string; personaId?: string; attachedSlugs?: string[]; + usage?: ReturnType; turnRunId?: string; + turnStatus?: ReturnType; queue?: string[]; +}; +export function viewportCarriers(meta: unknown, queue?: unknown): ViewportCarriers { + const m = objectRecord(meta) ?? {}; + return { + cwd: normalizeSessionCwd(m.logicalCwd), + activeSandboxId: isRedisSafeOpaqueId(m.activeSandboxId) ? m.activeSandboxId : undefined, + selectedModel: sanitizeModelId(m.selectedModel), reasoningEffort: sanitizeReasoningEffort(m.reasoningEffort), + resolvedProvider: sanitizeResolvedProvider(m.resolvedProvider), + personaId: isRedisSafeOpaqueId(m.personaId) ? m.personaId : undefined, + attachedSlugs: typeof m.attachedSkills === 'string' && m.attachedSkills.length <= HARNESS_SESSION_MAX_META_BYTES + ? parseAttachedSkills(m.attachedSkills).slice(0, HARNESS_SESSION_MAX_ATTACHED_SKILLS) : undefined, + usage: typeof m.usage === 'string' ? decodeUsageMetaString(m.usage) : undefined, + turnRunId: sanitizeTurnRunId(m.turnRunId), turnStatus: sanitizeTurnStatus(m.turnStatus), + queue: sanitizeQueue(queue), + }; +} + +/** Re-sanitize encoded carriers (output shape: `cwd`, not envelope `logicalCwd`). */ +export function parseViewportCarriers(raw: unknown): ViewportCarriers { + const c = objectRecord(raw) ?? {}; + const slugs = Array.isArray(c.attachedSlugs) + ? c.attachedSlugs.slice(0, HARNESS_SESSION_MAX_ATTACHED_SKILLS) + .filter((s): s is string => typeof s === 'string' && SKILL_SLUG_RE.test(s)) + : undefined; + return { + cwd: normalizeSessionCwd(c.cwd), + activeSandboxId: isRedisSafeOpaqueId(c.activeSandboxId) ? c.activeSandboxId : undefined, + selectedModel: sanitizeModelId(c.selectedModel), reasoningEffort: sanitizeReasoningEffort(c.reasoningEffort), + resolvedProvider: sanitizeResolvedProvider(c.resolvedProvider), + personaId: isRedisSafeOpaqueId(c.personaId) ? c.personaId : undefined, + attachedSlugs: slugs?.length ? slugs : undefined, + usage: sanitizeUsageSummary(c.usage), + turnRunId: sanitizeTurnRunId(c.turnRunId), turnStatus: sanitizeTurnStatus(c.turnStatus), + queue: sanitizeQueue(c.queue), + }; +} + +export type ViewportView = { + version: 1; sessionId: string; rows: SessionMessage[]; replace: boolean; + source: 'stream_tail' | 'stored_head' | 'unavailable'; + historyComplete: false; incomplete: true; gap: boolean; hasEarlier: boolean; + carriers: ViewportCarriers; +}; + +/** Linear row accounting; final caller supplies its actual serialized control scaffold. */ +export function fitViewportRows(rows: readonly SessionMessage[], scaffold: object): SessionMessage[] { + let remaining = VIEWPORT_RESPONSE_MAX_BYTES - byteLength(JSON.stringify(scaffold)); + const kept: SessionMessage[] = []; + for (let i = rows.length - 1; i >= 0 && kept.length < HARNESS_RING_MAX; i--) { + const row = rows[i]; + const size = byteLength(JSON.stringify(row)) + (kept.length > 0 ? 1 : 0); + if (size > remaining) break; + kept.push(row); remaining -= size; + } + return kept.reverse(); +} + +/** Parse ONE decoded stored frame, with an allowlist (never forward arbitrary JSON properties). */ +export function parseViewportEvent(text: string): AgentStreamEvent | undefined { + const normalized = text.replace(/\r\n/g, '\n'); + const blocks = normalized.trim().split('\n\n'); + if (blocks.length !== 1) return; + const data = blocks[0].split('\n').filter(l => l.startsWith('data:')).map(l => l.slice(5).trimStart()).join('\n'); + let raw: unknown; + try { raw = JSON.parse(data); } catch { return; } + return validateViewportEvent(raw); +} + +export function validateViewportEvent(raw: unknown): AgentStreamEvent | undefined { + const o = objectRecord(raw); + if (!o) return; + const str = (key: string): string | undefined => typeof o[key] === 'string' ? o[key] as string : undefined; + switch (o.type) { + case 'text_delta': case 'reasoning_delta': + return typeof o.text === 'string' ? { type: o.type, text: o.text } : undefined; + case 'tool_start': + return typeof o.name === 'string' ? { type: o.type, name: o.name, id: str('id') } : undefined; + case 'tool_result': + return typeof o.name === 'string' && typeof o.ok === 'boolean' && typeof o.summary === 'string' + ? { type: o.type, name: o.name, ok: o.ok, summary: o.summary, preview: str('preview'), id: str('id'), + changeDirCwd: normalizeSessionCwd(o.changeDirCwd), + activeSandboxId: isRedisSafeOpaqueId(o.activeSandboxId) ? o.activeSandboxId : undefined } : undefined; + case 'provider': { + const provider = sanitizeResolvedProvider(o.provider); + return provider ? { type: o.type, provider } : undefined; + } + case 'usage': { + const usage = sanitizeUsageSummary(o.usage); + return usage ? { type: o.type, usage } : undefined; + } + case 'done': return typeof o.text === 'string' ? { + type: o.type, text: o.text, finishReason: str('finishReason'), cwd: normalizeSessionCwd(o.cwd), + activeSandboxId: isRedisSafeOpaqueId(o.activeSandboxId) ? o.activeSandboxId : undefined, + sandboxId: isRedisSafeOpaqueId(o.sandboxId) ? o.sandboxId : undefined, + usage: sanitizeUsageSummary(o.usage), resolvedProvider: sanitizeResolvedProvider(o.resolvedProvider), + } : undefined; + case 'error': return typeof o.error === 'string' ? { type: o.type, error: o.error, + ...(typeof o.status === 'number' && Number.isInteger(o.status) ? { status: o.status } : {}) } : undefined; + case 'skill_attached': return typeof o.slug === 'string' && SKILL_SLUG_RE.test(o.slug) && + (o.action === 'attach' || o.action === 'detach') && typeof o.ok === 'boolean' + ? { type: o.type, slug: o.slug, action: o.action, ok: o.ok, reason: str('reason'), + ...(Array.isArray(o.attachedSlugs) ? { attachedSlugs: o.attachedSlugs + .slice(0, HARNESS_SESSION_MAX_ATTACHED_SKILLS).filter((s): s is string => typeof s === 'string' && SKILL_SLUG_RE.test(s)) } : {}) } : undefined; + default: return; + } +} + +/** Bounded recent display accumulator; reasoning never retained, done is not a duplicate append. */ +export class ViewportReducer { + private rows: SessionMessage[] = []; + private sizes: number[] = []; + private bytes = 0; + private serial = 0; + private group = createToolRunGroup(); + private lastKind = ''; + private sawAssistant = false; + private row(role: SessionRole, text: string, replace = false): void { + if (replace && this.rows.length) { + this.rows.pop(); this.bytes -= this.sizes.pop() ?? 0; + } + const row = { id: `view_${this.serial++}`, role, text: viewportExcerpt(text), at: 0 }; + const size = byteLength(JSON.stringify(row)) + 1; + this.rows.push(row); this.sizes.push(size); this.bytes += size; + while (this.bytes > VIEWPORT_RESPONSE_MAX_BYTES || this.rows.length > HARNESS_RING_MAX) { + this.rows.shift(); this.bytes -= this.sizes.shift() ?? 0; + } + } + apply(ev: AgentStreamEvent): void { + if (ev.type === 'reasoning_delta') return; + if (ev.type === 'text_delta') { + const grow = this.lastKind === 'assistant' && this.rows.at(-1)?.role === 'assistant'; + this.row('assistant', (grow ? this.rows.at(-1)!.text : '') + viewportExcerpt(ev.text), grow); + this.lastKind = 'assistant'; this.sawAssistant = true; + } else if (ev.type === 'tool_start' || ev.type === 'tool_result') { + const matches = ev.id ? this.group.items.filter(i => i.status === 'running' && i.callId === ev.id) : []; + // No name-only pairing across a missing/ambiguous sample start. + const reuse = this.lastKind === 'tool' && this.rows.at(-1)?.role === 'tool_run' && + !toolRunIsFull(this.group) && (ev.type === 'tool_start' || matches.length === 1); + if (!reuse) this.group = createToolRunGroup(); + if (ev.type === 'tool_start') addToolStart(this.group, viewportExcerpt(ev.name), ev.id ? viewportExcerpt(ev.id) : undefined); + else addToolResult(this.group, viewportExcerpt(ev.name), ev.ok, viewportExcerpt(ev.summary), + ev.preview ? viewportExcerpt(ev.preview) : undefined, ev.id ? viewportExcerpt(ev.id) : undefined); + this.row('tool_run', encodeToolRun(this.group) ?? 'Tool activity (incomplete)', reuse); + this.lastKind = 'tool'; + } else if (ev.type === 'error') { this.row('error', ev.error); this.lastKind = ''; } + else if (ev.type === 'skill_attached') { + // Same display strings as skillRowText (lib/harnessChat.ts) — do not import the host. + const text = ev.ok + ? (ev.action === 'detach' ? `Skill detached: ${ev.slug}` : `Skill attached: ${ev.slug}`) + : `Skill not attached: ${ev.slug}`; + this.row('skill_attached', text); this.lastKind = ''; + } + else if (ev.type === 'done' && !this.sawAssistant && ev.text) { + this.row('assistant', ev.text); this.sawAssistant = true; this.lastKind = ''; + } + } + snapshot(): SessionMessage[] { return this.rows.slice(); } +} diff --git a/lib/sessions/viewportRead.test.ts b/lib/sessions/viewportRead.test.ts new file mode 100644 index 00000000..8edb8213 --- /dev/null +++ b/lib/sessions/viewportRead.test.ts @@ -0,0 +1,119 @@ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { MemoryBlobTranscriptStore } from './blobStores'; +import { newBlobObjectId } from './blobStore'; +import { emptyViewport, readViewportHead, recoverViewport } from './viewportRead'; +import { VIEWPORT_RECOVERY_MAX_BYTES, VIEWPORT_RECOVERY_MAX_MS, VIEWPORT_FINAL_PROBE_MAX_MS } from '../sessionCloudCaps'; +import type { ViewportRunReader } from '../workflows/viewportRunReader'; +import { VIEWPORT_HISTORY_NOTE } from '../sessionViewport'; + +const scope = { tenantId: 'tenant', userId: 'user', sessionId: 'session' }; +const line = (event: object) => `data: ${JSON.stringify(event)}\n\n`; +const head = () => Promise.resolve(emptyViewport(scope.sessionId)); +function runWith(frames: string[], final = frames.length): ViewportRunReader { + return { nextIndex: vi.fn(async () => final), status: async () => 'running', + open: vi.fn(start => new ReadableStream({ start(c) { for (const f of frames.slice(start)) c.enqueue(f); c.close(); } })) }; +} +afterEach(() => vi.useRealTimers()); + +describe('one-head viewport source', () => { + it('reads only scoped head, never prev, validates body id, does not write', async () => { + const blob = new MemoryBlobTranscriptStore(); + const id = newBlobObjectId(scope); const prev = newBlobObjectId(scope); + await blob.writeSegment({ objectId: id, maxBytes: 8*1024*1024, content: JSON.stringify({ id: 'session', prev, queue: ['next'], + messages: Array.from({length: 2500}, (_, i) => ({role:'assistant', text:`message${i}`})) }) }); + const read = vi.spyOn(blob, 'read'), write = vi.spyOn(blob, 'writeSegment'); + const view = await readViewportHead({scope, meta:{transcriptPointer:id, workingNotes:'hidden'}, blob}); + expect(read).toHaveBeenCalledExactlyOnceWith(id); expect(write).not.toHaveBeenCalled(); + expect(view.rows.at(-1)?.text).toBe('message2499'); expect(view.rows).toHaveLength(2048); + expect(view.hasEarlier).toBe(true); expect(view.gap).toBe(false); expect(view.carriers.queue).toEqual(['next']); + expect(JSON.stringify(view)).not.toContain('hidden'); + }); + it('a complete small head is stored_head with gap false and hasEarlier false', async () => { + const blob = new MemoryBlobTranscriptStore(); + const id = newBlobObjectId(scope); + await blob.writeSegment({ objectId: id, maxBytes: 8*1024*1024, content: JSON.stringify({ + id: 'session', messages: [{ role: 'assistant', text: 'only' }] }) }); + const view = await readViewportHead({ scope, meta: { transcriptPointer: id }, blob }); + expect(view).toMatchObject({ source: 'stored_head', replace: true, gap: false, hasEarlier: false }); + expect(view.rows).toEqual([expect.objectContaining({ text: 'only' })]); + }); + it('foreign pointer is not read; corrupt/wrong-id/missing heads are optional misses', async () => { + const read = vi.fn(async () => null as string|null); + const pointer = newBlobObjectId({...scope, userId:'foreign'}); + expect((await readViewportHead({scope, meta:{transcriptPointer:pointer}, blob:{read}})).replace).toBe(false); + expect(read).not.toHaveBeenCalled(); + for (const raw of ['bad-json', JSON.stringify({id:'wrong',messages:[{role:'assistant',text:'foreign'}]}), null]) { + read.mockResolvedValue(raw); + const view = await readViewportHead({scope, meta:{transcriptPointer:newBlobObjectId(scope)},blob:{read}}); + expect(view.replace).toBe(false); expect(view.gap).toBe(true); expect(view.source).toBe('unavailable'); + } + }); +}); + +describe('bounded recent stream recovery', () => { + it('samples 2048 recent frames of a huge run, discards thinking, final probe once and no catch-up chase', async () => { + const cancel = vi.fn(); let reads = 0; + const run: ViewportRunReader = { status: async () => 'running', nextIndex: vi.fn(async () => 1000050), + open: vi.fn(start => new ReadableStream({ pull(c) { + reads++; + c.enqueue(line(start + reads === 1000000 ? {type:'text_delta',text:'latest'} : {type:'reasoning_delta',text:'historic'})); + }, cancel })) }; + const view = await recoverViewport({runId:'run',sessionId:'session',initialIndex:1000000,run,readHead:head}); + expect(run.open).toHaveBeenCalledExactlyOnceWith(1000000-2048); + // Web Streams may prefetch one raw frame; the application consumes only 2048. + expect(reads).toBeLessThanOrEqual(2049); expect(cancel).toHaveBeenCalledOnce(); + expect(view.sampledRange).toEqual({start:1000000-2048,end:1000000}); + expect(run.nextIndex).toHaveBeenCalledOnce(); expect(view.resumeIndex).toBe(1000050); expect(view.gap).toBe(true); + expect(view.rows.some(r=>r.text==='latest')).toBe(true); expect(view.rows.some(r=>r.text===VIEWPORT_HISTORY_NOTE)).toBe(true); + expect(JSON.stringify(view)).not.toContain('historic"'); + }); + it('chooses sample OR head, never merges ambiguous history; thinking-only sample falls back', async () => { + const h = {...emptyViewport('session'), source:'stored_head' as const, replace:true, + rows:[{id:'h',role:'assistant' as const,text:'head-only',at:0}]}; + for (const event of [{type:'text_delta',text:'sample-only'}, {type:'reasoning_delta',text:'never-paint'}]) { + const run = runWith([line(event)]); + const view = await recoverViewport({runId:'run',sessionId:'session',initialIndex:1,run,readHead:async()=>h}); + expect(view.rows[0].text).toBe(event.type==='text_delta'?'sample-only':'head-only'); + expect(view.source).toBe(event.type==='text_delta'?'stream_tail':'stored_head'); + expect(view.rows.some(r=>r.text===VIEWPORT_HISTORY_NOTE)).toBe(true); + } + }); + it('oversized decoded frame ends optional sample but still attaches at tail', async () => { + const run = runWith(['x'.repeat(VIEWPORT_RECOVERY_MAX_BYTES+1)], 3); + const view = await recoverViewport({runId:'run',sessionId:'session',initialIndex:1,run,readHead:head}); + expect(view.replace).toBe(false); expect(view.resumeIndex).toBe(3); expect(view.gap).toBe(true); + }); + it('shares the 5s optional deadline and uses H0 when final metadata probe hangs', async () => { + vi.useFakeTimers(); const cancel = vi.fn(); + const run: ViewportRunReader = {status:async()=> 'running',nextIndex:vi.fn(()=>new Promise(()=>{})),open:()=>new ReadableStream({cancel})}; + const promise = recoverViewport({runId:'run',sessionId:'session',initialIndex:7,run,readHead:()=>new Promise(()=>{})}); + await vi.advanceTimersByTimeAsync(VIEWPORT_RECOVERY_MAX_MS + VIEWPORT_FINAL_PROBE_MAX_MS); + const view = await promise; + expect(view.resumeIndex).toBe(7); expect(view.replace).toBe(false); expect(cancel).toHaveBeenCalledOnce(); + expect(run.nextIndex).toHaveBeenCalledOnce(); expect(vi.getTimerCount()).toBe(0); + }); + it('abort releases pending sample and never probes live tail afterward', async () => { + const controller = new AbortController(), cancel = vi.fn(); + const run: ViewportRunReader = {status:async()=> 'running',nextIndex:vi.fn(async()=>4),open:()=>new ReadableStream({cancel})}; + const promise = recoverViewport({runId:'run',sessionId:'session',initialIndex:3,run,readHead:head,signal:controller.signal}); + controller.abort(); await expect(promise).rejects.toThrow('aborted'); + expect(cancel).toHaveBeenCalledOnce(); expect(run.nextIndex).not.toHaveBeenCalled(); + }); + it('malformed/empty/early EOF are best effort, not a retry loop', async () => { + const run = runWith(['bad frame'], 5); + const view = await recoverViewport({runId:'run',sessionId:'session',initialIndex:4,run,readHead:head}); + expect(view.resumeIndex).toBe(5); expect(view.replace).toBe(false); expect(run.open).toHaveBeenCalledOnce(); + }); + it('skipStream is head-only and never opens or probes the run', async () => { + const run: ViewportRunReader = { status: async () => 'cancelled', nextIndex: vi.fn(async () => 99), open: vi.fn() }; + const stored = { ...emptyViewport('session'), source: 'stored_head' as const, replace: true, + rows: [{ id: 'h', role: 'assistant' as const, text: 'head-only', at: 0 }] }; + const view = await recoverViewport({ runId: 'run', sessionId: 'session', initialIndex: 50, run, + readHead: async () => stored, skipStream: true }); + expect(run.open).not.toHaveBeenCalled(); expect(run.nextIndex).not.toHaveBeenCalled(); + expect(view.source).toBe('stored_head'); expect(view.gap).toBe(true); + expect(view).not.toHaveProperty('resumeIndex'); + expect(view.rows.some(r => r.text === 'head-only')).toBe(true); + expect(view.rows.some(r => r.text === VIEWPORT_HISTORY_NOTE)).toBe(true); + }); +}); diff --git a/lib/sessions/viewportRead.ts b/lib/sessions/viewportRead.ts new file mode 100644 index 00000000..cee56577 --- /dev/null +++ b/lib/sessions/viewportRead.ts @@ -0,0 +1,113 @@ +/** Read-only head + bounded stream tail. No ancestry reconstruction and no writes. */ +import { + HARNESS_SESSION_MAX_BODY_BYTES, VIEWPORT_RECOVERY_MAX_MS, VIEWPORT_RECOVERY_MAX_BYTES, + VIEWPORT_TAIL_MAX_FRAMES, VIEWPORT_FINAL_PROBE_MAX_MS, VIEWPORT_HEAD_READ_MAX_OBJECTS, +} from '../sessionCloudCaps'; +import { HARNESS_RING_MAX } from '../sessionWindow'; +import { isObjectIdBoundTo, type ObjectScope, type BlobTranscriptStore } from './blobStore'; +import { + byteLength, objectRecord, parseViewportRow, parseViewportEvent, fitViewportRows, + viewportCarriers, ViewportReducer, VIEWPORT_HISTORY_NOTE, type ViewportView, +} from '../sessionViewport'; +import type { SessionMessage } from '../sessionStore'; +import type { ViewportSnapshot } from '../viewportStreamProtocol'; +import { type ViewportRunReader, viewportWait } from '../workflows/viewportRunReader'; + +export function emptyViewport(sessionId: string, meta?: unknown): ViewportView { + return { version: 1, sessionId, replace: false, rows: [], source: 'unavailable', + historyComplete: false, incomplete: true, gap: true, hasEarlier: false, carriers: viewportCarriers(meta) }; +} + +export async function readViewportHead(opts: { + scope: ObjectScope; meta: Record; blob: Pick; + signal?: AbortSignal; deadline?: number; +}): Promise { + const fallback = emptyViewport(opts.scope.sessionId, opts.meta); + const pointer = opts.meta.transcriptPointer; + if (typeof pointer !== 'string' || !isObjectIdBoundTo(pointer, opts.scope) || VIEWPORT_HEAD_READ_MAX_OBJECTS < 1 || opts.signal?.aborted) return fallback; + try { + const deadline = opts.deadline ?? Date.now() + VIEWPORT_RECOVERY_MAX_MS; + if (Date.now() >= deadline) return fallback; + const raw = await viewportWait(opts.blob.read(pointer), deadline - Date.now(), opts.signal); + if (raw === null || raw.length > HARNESS_SESSION_MAX_BODY_BYTES || byteLength(raw) > HARNESS_SESSION_MAX_BODY_BYTES || Date.now() >= deadline) return fallback; + const body = objectRecord(JSON.parse(raw)); + if (!body || body.id !== opts.scope.sessionId || !Array.isArray(body.messages)) return fallback; + const rows: SessionMessage[] = []; + // Only newest candidates; no full-message mapping (including malformed histories). + for (let i = body.messages.length - 1; i >= Math.max(0, body.messages.length - HARNESS_RING_MAX) && Date.now() < deadline; i--) { + const row = parseViewportRow(body.messages[i], i); + if (row) rows.push(row); + } + rows.reverse(); + const view: ViewportView = { ...fallback, source: rows.length ? 'stored_head' : 'unavailable', + carriers: viewportCarriers(opts.meta, body.queue), replace: rows.length > 0, gap: false, + hasEarlier: body.messages.length > rows.length || + (typeof body.prev === 'string' && isObjectIdBoundTo(body.prev, opts.scope)) }; + view.rows = fitViewportRows(rows, view); + view.hasEarlier ||= view.rows.length < rows.length; + view.replace = view.rows.length > 0; + return view; + } catch { return fallback; } +} + +export async function recoverViewport(opts: { + runId: string; sessionId: string; initialIndex: number; run: ViewportRunReader; + readHead: (deadline: number) => Promise; signal?: AbortSignal; + /** Skip SDK sample + tail probe (cancelled/failed C16 hang class). Head-only. */ + skipStream?: boolean; + /** Shared recovery clock (H0 + head + sample). Defaults to now + VIEWPORT_RECOVERY_MAX_MS. */ + deadline?: number; +}): Promise { + const deadline = opts.deadline ?? Date.now() + VIEWPORT_RECOVERY_MAX_MS; + const remaining = () => Math.max(0, deadline - Date.now()); + const headPromise = viewportWait(opts.readHead(deadline), remaining(), opts.signal) + .catch(() => emptyViewport(opts.sessionId)); + const start = Math.max(0, opts.initialIndex - VIEWPORT_TAIL_MAX_FRAMES); + let end = start, bytes = 0, gap = opts.skipStream || start > 0; + const reducer = new ViewportReducer(); + let reader: ReadableStreamDefaultReader | undefined; + try { + if (!opts.skipStream && start < opts.initialIndex && !opts.signal?.aborted) { + reader = opts.run.open(start).getReader(); + while (end < opts.initialIndex && end - start < VIEWPORT_TAIL_MAX_FRAMES && Date.now() < deadline) { + const chunk = await viewportWait(reader.read(), deadline - Date.now(), opts.signal); + if (chunk.done) { gap = true; break; } + const size = typeof chunk.value === 'string' ? byteLength(chunk.value) : chunk.value.byteLength; + if (size > VIEWPORT_RECOVERY_MAX_BYTES - bytes) { gap = true; break; } + bytes += size; + const text = typeof chunk.value === 'string' ? chunk.value : new TextDecoder('utf-8', { fatal: true }).decode(chunk.value); + const event = parseViewportEvent(text); + end++; + if (event) reducer.apply(event); else gap = true; + } + } + } catch { gap = true; } + finally { if (reader) void reader.cancel().catch(() => {}); } + const head = await headPromise; + if (opts.signal?.aborted) throw new Error('Viewport read aborted'); + let resumeIndex = opts.initialIndex; + if (!opts.skipStream) { + try { + resumeIndex = Math.max(resumeIndex, await viewportWait(opts.run.nextIndex(), VIEWPORT_FINAL_PROBE_MAX_MS, opts.signal)); + } catch { gap = true; } + } + if (opts.signal?.aborted) throw new Error('Viewport read aborted'); + gap ||= end < opts.initialIndex || resumeIndex > opts.initialIndex; + const sampled = reducer.snapshot(); + const rows = sampled.length ? sampled : head.rows; + const snapshot: ViewportSnapshot = { + ...head, runId: opts.runId, sampledRange: { start, end }, + source: sampled.length ? 'stream_tail' : head.source, rows: [], + gap, hasEarlier: gap || head.hasEarlier, replace: rows.length > 0, + // skipStream never captured a tail; do not mint a guessed transport cursor. + ...(opts.skipStream ? {} : { resumeIndex }), + }; + if (rows.length) { + const note: SessionMessage = { id: 'viewport_history_note', role: 'system', text: VIEWPORT_HISTORY_NOTE, at: 0 }; + // Include actual record envelope/framing in accounting, not just the rows. + snapshot.rows = fitViewportRows([...rows, note], { ...snapshot, type: 'viewport_snapshot', framing: 'event: viewport_snapshot\ndata: \n\n' }); + snapshot.replace = snapshot.rows.length > 1; + if (!snapshot.replace) snapshot.rows = []; + } + return snapshot; +} diff --git a/lib/viewportStreamProtocol.test.ts b/lib/viewportStreamProtocol.test.ts new file mode 100644 index 00000000..e43f7ec9 --- /dev/null +++ b/lib/viewportStreamProtocol.test.ts @@ -0,0 +1,90 @@ +import { describe, expect, it } from 'vitest'; +import { encodeViewportRecord, parseViewportMode, ViewportStreamDecoder, type ViewportRecord } from './viewportStreamProtocol'; +import { emptyViewport } from './sessions/viewportRead'; + +const snap: ViewportRecord = { type: 'viewport_snapshot', ...emptyViewport('session'), runId: 'run', resumeIndex: 42, sampledRange: { start: 0, end: 40 } }; +describe('negotiated viewport codec', () => { + it('decodes fragmented UTF-8/CRLF and separates explicit snapshot jumps from stored indices', () => { + const decoder = new ViewportStreamDecoder('run'); + const text = new TextDecoder().decode(encodeViewportRecord(snap)) + new TextDecoder().decode(encodeViewportRecord({ type: 'turn_event', version: 1, runId: 'run', nextIndex: 43, event: { type: 'text_delta', text: '漢🙂' } })); + const bytes = new TextEncoder().encode(text.replace(/\n/g, '\r\n')); + const records: ViewportRecord[] = []; + for (const byte of bytes) records.push(...decoder.push(new Uint8Array([byte]))); + decoder.push(new Uint8Array(), true); + expect(records).toHaveLength(2); + expect(records[1]).toMatchObject({ nextIndex: 43, event: { text: '漢🙂' } }); + }); + it('display-only snapshot (unknown tail) does not jump the cursor to origin', () => { + const decoder = new ViewportStreamDecoder('run'); + const { resumeIndex: _ignored, ...display } = snap; + const rec = decoder.push(encodeViewportRecord({ ...display, type: 'viewport_snapshot' }))[0]; + expect(rec).toMatchObject({ type: 'viewport_snapshot', sampledRange: { start: 0, end: 40 } }); + expect(rec).not.toHaveProperty('resumeIndex'); + expect(() => decoder.push(encodeViewportRecord({ + type: 'turn_event', version: 1, runId: 'run', nextIndex: 1, event: { type: 'text_delta', text: 'x' }, + }))).toThrow('Invalid viewport cursor'); + }); + it('drops extra snapshot keys and re-sanitizes encoded carriers', () => { + const decoder = new ViewportStreamDecoder('run'); + const hostile = [ + 'event: viewport_snapshot', + `data: ${JSON.stringify({ + version: 1, runId: 'run', sessionId: 'session', resumeIndex: 42, + sampledRange: { start: 0, end: 40 }, rows: [], replace: false, + source: 'unavailable', historyComplete: false, incomplete: true, + gap: true, hasEarlier: false, + carriers: { cwd: '/etc/passwd', workingNotes: 'secret notes', queue: ['next'] }, + workingNotes: 'secret notes', personaSnapshot: 'persona body', + })}`, + '', + '', + ].join('\n'); + const rec = decoder.push(new TextEncoder().encode(hostile))[0]; + expect(rec).toMatchObject({ type: 'viewport_snapshot', resumeIndex: 42, carriers: { queue: ['next'] } }); + expect(rec).not.toHaveProperty('workingNotes'); + expect(rec).not.toHaveProperty('personaSnapshot'); + expect(JSON.stringify(rec)).not.toMatch(/secret|persona body|\/etc/); + expect((rec as { carriers: { cwd?: string } }).carriers.cwd).toBeUndefined(); + }); + it('rejects a non-opaque snapshot sessionId', () => { + const decoder = new ViewportStreamDecoder('run'); + const hostile = [ + 'event: viewport_snapshot', + `data: ${JSON.stringify({ + version: 1, runId: 'run', sessionId: '../other', resumeIndex: 42, + sampledRange: { start: 0, end: 40 }, rows: [], replace: false, + source: 'unavailable', historyComplete: false, incomplete: true, + gap: true, hasEarlier: false, carriers: {}, + })}`, + '', '', + ].join('\n'); + expect(() => decoder.push(new TextEncoder().encode(hostile))).toThrow('Invalid viewport snapshot'); + }); + it('known skipped frame advances; synthetic terminal does not need an index', () => { + const decoder = new ViewportStreamDecoder('run', 3); + expect(decoder.push(encodeViewportRecord({ type: 'turn_event', version: 1, runId: 'run', nextIndex: 4, skipped: true }))[0]).toMatchObject({ skipped: true }); + expect(decoder.push(encodeViewportRecord({ type: 'viewport_end', version: 1, runId: 'run', status: 'cancelled' }))[0].type).toBe('viewport_end'); + }); + it.each([ + 'event: turn_event\nid: 8\ndata: {"version":1,"runId":"run","nextIndex":8,"skipped":true}\n\n', + 'event: viewport_end\nid: 4\ndata: {"version":1,"runId":"run","status":"completed"}\n\n', + 'event: viewport_end\ndata: {"version":1,"runId":"foreign","status":"completed"}\n\n', + 'event: turn_event\nid: 4\ndata: {"version":1,"runId":"run","nextIndex":4,"event":{"type":"unknown"}}\n\n', + ])('rejects inconsistent or hostile record %s', text => { + expect(() => new ViewportStreamDecoder('run', 3).push(new TextEncoder().encode(text))).toThrow(); + }); + it('fails closed on incomplete/malformed UTF-8 rather than guessing raw position', () => { + expect(() => new ViewportStreamDecoder('run').push(new TextEncoder().encode('event: unfinished'), true)).toThrow(); + expect(() => new ViewportStreamDecoder('run').push(new Uint8Array([0xff]), true)).toThrow(); + }); + it('negotiates explicitly, preserves old mode, rejects duplicate/conflicting selectors', () => { + const mode = (query: string, method: 'GET'|'POST' = 'GET') => parseViewportMode(new URL(`https://example.com/?${query}`), method); + expect(mode('startIndex=2')).toEqual({ kind: 'legacy' }); + expect(mode('viewportVersion=1&hydrate=tail')).toEqual({ kind: 'cold' }); + expect(mode('viewportVersion=1&startIndex=42')).toEqual({ kind: 'indexed', startIndex: 42 }); + expect(mode('viewportVersion=1&startIndex=0')).toEqual({ kind: 'indexed', startIndex: 0 }); + expect(mode('viewportVersion=1', 'POST')).toEqual({ kind: 'indexed', startIndex: 0 }); + for (const query of ['hydrate=tail', 'viewportVersion=2', 'viewportVersion=1', 'viewportVersion=1&sessionId=s1', 'viewportVersion=1&viewportVersion=1', 'viewportVersion=1&hydrate=tail&startIndex=0', 'viewportVersion=1&startIndex=01', 'viewportVersion=1&startIndex=1e3', 'viewportVersion=1&startIndex=1000000001']) expect(mode(query)).toBeNull(); + expect(mode('viewportVersion=1&startIndex=1', 'POST')).toBeNull(); + }); +}); diff --git a/lib/viewportStreamProtocol.ts b/lib/viewportStreamProtocol.ts new file mode 100644 index 00000000..6b3b104c --- /dev/null +++ b/lib/viewportStreamProtocol.ts @@ -0,0 +1,126 @@ +/** Versioned viewport records. Synthetic lifecycle controls never occupy a stored cursor. */ +import type { AgentStreamEvent } from './agent/agentStream'; +import { HARNESS_SESSION_MAX_MSG_BYTES, VIEWPORT_RESPONSE_MAX_BYTES, sanitizeTurnRunId, sanitizeTurnStreamCursor, isRedisSafeOpaqueId } from './sessionCloudCaps'; +import { HARNESS_RING_MAX } from './sessionWindow'; +import { byteLength, objectRecord, parseViewportRow, parseViewportCarriers, validateViewportEvent, type ViewportView } from './sessionViewport'; + +export type ViewportSnapshot = ViewportView & { + runId: string; resumeIndex?: number; sampledRange: { start: number; end: number }; +}; +export type ViewportRecord = + | { type: 'viewport_state'; version: 1; runId: string; status: string; phase: 'recovering' } + | ({ type: 'viewport_snapshot' } & ViewportSnapshot) + | { type: 'turn_event'; version: 1; runId: string; nextIndex: number; event: AgentStreamEvent; skipped?: never } + | { type: 'turn_event'; version: 1; runId: string; nextIndex: number; skipped: true; event?: never } + | { type: 'viewport_end'; version: 1; runId: string; status: string } + | { type: 'viewport_error'; version: 1; runId: string; code: string }; + +export function encodeViewportRecord(record: ViewportRecord): Uint8Array { + const { type, ...data } = record; + return new TextEncoder().encode(`event: ${type}\n${type === 'turn_event' ? `id: ${record.nextIndex}\n` : ''}data: ${JSON.stringify(data)}\n\n`); +} + +export type ViewportMode = { kind: 'legacy' } | { kind: 'cold' } | { kind: 'indexed'; startIndex: number }; +export function parseViewportMode(url: URL, method: 'GET' | 'POST'): ViewportMode | null { + const q = url.searchParams; + if (!q.has('viewportVersion')) return q.has('hydrate') ? null : { kind: 'legacy' }; + if (['viewportVersion', 'hydrate', 'startIndex', 'sessionId'].some(k => q.getAll(k).length > 1) || q.get('viewportVersion') !== '1') return null; + if (method === 'POST') return q.has('hydrate') || q.has('startIndex') ? null : { kind: 'indexed', startIndex: 0 }; + if (q.has('hydrate')) return q.get('hydrate') === 'tail' && !q.has('startIndex') ? { kind: 'cold' } : null; + // GET v1 must pick hydrate=tail (bounded recovery) or an explicit startIndex. + // Omitting both is not origin replay — that was the #924 class this path exists to avoid. + if (!q.has('startIndex')) return null; + const raw = q.get('startIndex'); + const index = raw !== null && /^(0|[1-9]\d*)$/.test(raw) ? sanitizeTurnStreamCursor(Number(raw)) : undefined; + return index === undefined ? null : { kind: 'indexed', startIndex: index }; +} + +/** Incremental UTF-8/SSE decoder; unknown or inconsistent positions are not guessed. */ +export class ViewportStreamDecoder { + private decoder = new TextDecoder('utf-8', { fatal: true }); + private pending = ''; + private nextIndex: number | 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 }); + // Normalize complete CRLF pairs only: a trailing CR may join next network chunk. + this.pending = this.pending.replace(/\r\n/g, '\n'); + const out: ViewportRecord[] = []; + for (;;) { + const end = this.pending.indexOf('\n\n'); + if (end < 0) break; + const block = this.pending.slice(0, end); + this.pending = this.pending.slice(end + 2); + if (!block.trim() || block.split('\n').every(l => l.startsWith(':'))) continue; + out.push(this.parse(block)); + } + if (final && this.pending.trim()) throw new Error('Incomplete viewport record'); + return out; + } + 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'); + if (types.length !== 1 || ids.length > 1) throw new Error('Invalid viewport framing'); + const o = objectRecord(JSON.parse(fields('data').join('\n'))); + if (!o || o.version !== 1 || sanitizeTurnRunId(o.runId) !== this.runId || o.runId !== this.runId) throw new Error('Invalid viewport identity'); + const type = types[0]; + if (type !== 'turn_event' && ids.length) throw new Error('Synthetic viewport index'); + if (type === 'turn_event') { + const index = sanitizeTurnStreamCursor(o.nextIndex); + if (index === undefined || this.nextIndex === undefined || index !== this.nextIndex + 1 || ids[0] !== String(index)) throw new Error('Invalid viewport cursor'); + const event = o.skipped === true ? undefined : validateViewportEvent(o.event); + if (!event && o.skipped !== true) throw new Error('Invalid viewport event'); + this.nextIndex = index; + return event ? { type, version: 1, runId: this.runId, nextIndex: index, event } + : { type, version: 1, runId: this.runId, nextIndex: index, skipped: true }; + } + if (type === 'viewport_snapshot') { + const hasResume = Object.prototype.hasOwnProperty.call(o, 'resumeIndex'); + const index = hasResume ? sanitizeTurnStreamCursor(o.resumeIndex) : undefined; + const range = objectRecord(o.sampledRange); + if ((hasResume && index === undefined) || !range || sanitizeTurnStreamCursor(range.start) === undefined || + sanitizeTurnStreamCursor(range.end) === undefined || (range.start as number) > (range.end as number) || + (index !== undefined && (range.end as number) > index) || + typeof o.sessionId !== 'string' || !isRedisSafeOpaqueId(o.sessionId) || !Array.isArray(o.rows) || o.rows.length > HARNESS_RING_MAX || + o.historyComplete !== false || o.incomplete !== true || typeof o.replace !== 'boolean' || + typeof o.gap !== 'boolean' || typeof o.hasEarlier !== 'boolean' || + !['stream_tail', 'stored_head', 'unavailable'].includes(String(o.source)) || + !objectRecord(o.carriers) || byteLength(block) > VIEWPORT_RESPONSE_MAX_BYTES) throw new Error('Invalid viewport snapshot'); + const rows = o.rows.map((r, i) => { + const raw = objectRecord(r); + const parsed = parseViewportRow(r, i); + if (!parsed || typeof raw?.text !== 'string' || byteLength(raw.text) > HARNESS_SESSION_MAX_MSG_BYTES) throw new Error('Invalid viewport row'); + return parsed; + }); + if (index !== undefined) { + if (this.nextIndex !== undefined && index < this.nextIndex) throw new Error('Rewound viewport'); + this.nextIndex = index; + } + const start = sanitizeTurnStreamCursor(range.start)!; + const end = sanitizeTurnStreamCursor(range.end)!; + return { + type, + version: 1, + runId: this.runId, + sessionId: o.sessionId as string, + rows, + replace: o.replace as boolean, + source: String(o.source) as ViewportView['source'], + historyComplete: false, + incomplete: true, + gap: o.gap as boolean, + hasEarlier: o.hasEarlier as boolean, + carriers: parseViewportCarriers(o.carriers), + sampledRange: { start, end }, + ...(index !== undefined ? { resumeIndex: index } : {}), + }; + } + if (type === 'viewport_state' && o.phase === 'recovering' && typeof o.status === 'string') + return { type, version: 1, runId: this.runId, status: o.status, phase: 'recovering' }; + if (type === 'viewport_end' && ['completed', 'failed', 'cancelled'].includes(String(o.status))) + return { type, version: 1, runId: this.runId, status: o.status as string }; + if (type === 'viewport_error' && typeof o.code === 'string') + return { type, version: 1, runId: this.runId, code: o.code }; + throw new Error('Unknown viewport record'); + } +} diff --git a/lib/workflows/viewportRunReader.test.ts b/lib/workflows/viewportRunReader.test.ts new file mode 100644 index 00000000..0238cdf3 --- /dev/null +++ b/lib/workflows/viewportRunReader.test.ts @@ -0,0 +1,22 @@ +import { describe, expect, it, vi } from 'vitest'; +import { createViewportRunReader, type ViewportRun } from './viewportRunReader'; + +describe('public decoded Workflow adapter',()=>{ + it('probes next raw index without consuming frames and closes the probe stream',async()=>{ + const cancel=vi.fn(), getTailIndex=vi.fn(async()=>41); + const run:ViewportRun={status:Promise.resolve('running'),getReadable:vi.fn(()=>Object.assign(new ReadableStream({cancel}),{getTailIndex}))}; + const adapter=createViewportRunReader(run); + expect(await adapter.nextIndex()).toBe(42);expect(getTailIndex).toHaveBeenCalledOnce();expect(cancel).toHaveBeenCalledOnce(); + expect(await adapter.status()).toBe('running'); + const readable=adapter.open(42);expect(run.getReadable).toHaveBeenLastCalledWith({startIndex:42});await readable.cancel(); + }); + it.each([-2,1e9,NaN])('rejects invalid tail %s',async tail=>{ + const run:ViewportRun={status:'running',getReadable:()=>Object.assign(new ReadableStream(),{getTailIndex:async()=>tail})}; + await expect(createViewportRunReader(run).nextIndex()).rejects.toThrow(); + }); + it('empty tail -1 is next index0; absent helper is unavailable rather than origin fallback',async()=>{ + const run:ViewportRun={status:'running',getReadable:()=>Object.assign(new ReadableStream(),{getTailIndex:async()=>-1})}; + expect(await createViewportRunReader(run).nextIndex()).toBe(0); + run.getReadable=()=>new ReadableStream();await expect(createViewportRunReader(run).nextIndex()).rejects.toThrow('unavailable'); + }); +}); diff --git a/lib/workflows/viewportRunReader.ts b/lib/workflows/viewportRunReader.ts new file mode 100644 index 00000000..51a5994b --- /dev/null +++ b/lib/workflows/viewportRunReader.ts @@ -0,0 +1,50 @@ +/** Structural SDK adapter: routes inject their authorized public getRun handle. No keys or run-input reads. */ +import { sanitizeTurnStreamCursor } from '../sessionCloudCaps'; + +export type ViewportRun = { + readonly status: PromiseLike | string; + getReadable(opts?: { startIndex?: number }): ReadableStream & { + getTailIndex?: () => Promise; + }; +}; +export type ViewportRunReader = { + status(): Promise; + nextIndex(): Promise; + open(startIndex: number): ReadableStream; +}; +export function createViewportRunReader(run: ViewportRun): ViewportRunReader { + return { + status: async () => await run.status, + open: startIndex => run.getReadable({ startIndex }), + async nextIndex() { + // The public helper rides a readable. Start at the tail (never origin), + // cancel its data side immediately, and only await the metadata request. + // A hung metadata probe must not leave an unconsumed SDK reader running. + const stream = run.getReadable({ startIndex: -1 }); + let tail: Promise; + try { + if (!stream.getTailIndex) throw new Error('Viewport tail unavailable'); + tail = stream.getTailIndex(); + } finally { void stream.cancel().catch(() => {}); } + const next = sanitizeTurnStreamCursor((await tail) + 1); + if (next === undefined) throw new Error('Invalid viewport tail'); + return next; + }, + }; +} + +/** Deadline/abort race with listener and timer cleanup; late rejections are always observed. */ +export function viewportWait(promise: PromiseLike, ms: number, signal?: AbortSignal): Promise { + return new Promise((resolve, reject) => { + let finished = false; + const finish = (fn: () => void) => { + if (finished) return; + finished = true; clearTimeout(timer); signal?.removeEventListener('abort', abort); fn(); + }; + const abort = () => finish(() => reject(new Error('Viewport read aborted'))); + const timer = setTimeout(() => finish(() => reject(new Error('Viewport read timed out'))), Math.max(0, ms)); + signal?.addEventListener('abort', abort, { once: true }); + Promise.resolve(promise).then(v => finish(() => resolve(v)), e => finish(() => reject(e))); + if (signal?.aborted) abort(); + }); +}