Skip to content
This repository was archived by the owner on Sep 8, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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` |
Expand Down
23 changes: 23 additions & 0 deletions SECURITY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
37 changes: 37 additions & 0 deletions app/api/sessions/[id]/viewport/route.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
44 changes: 44 additions & 0 deletions app/api/sessions/[id]/viewport/route.ts
Original file line number Diff line number Diff line change
@@ -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<Response> {
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 });
}
}
43 changes: 43 additions & 0 deletions app/api/turns/[runId]/stream/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -179,10 +188,12 @@ export async function GET(
);
}

let envelopeMeta: Record<string, unknown> = {};
// 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}` },
Expand Down Expand Up @@ -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}` },
Expand Down
Loading
Loading