Skip to content
This repository was archived by the owner on Sep 8, 2026. It is now read-only.
Draft
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
5 changes: 5 additions & 0 deletions lib/sessionRepository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -827,10 +827,13 @@ function cloudBodyBytes(body: CloudPutBody): number {
* #515 envelope/Blob path, where the transcript object is ferried **client→Blob**
* and never through a Function body.
*/
import { assertCompleteForCloudPut } from './viewportFirewall';

export function trimForCloudPut(
snapshot: SessionSnapshot,
maxBytes: number = HARNESS_SESSION_MAX_FUNCTION_BODY_BYTES,
): CloudPutBody {
assertCompleteForCloudPut(snapshot);
const messages = snapshot.messages.map((m) => ({
id: m.id,
role: m.role,
Expand Down Expand Up @@ -1306,6 +1309,8 @@ export function createHttpSessionRepository(
// Canonical identity (parent #415): a snapshot never stores under a different
// resource id than its own persisted id.
if (snapshot.id !== id) return;
// Plan #960 firewall — never queue a partial view snapshot for cloud persist.
if (snapshot.historyComplete === false) return;
const c = channel(id);
c.pending = snapshot;
void drain(id, c);
Expand Down
8 changes: 8 additions & 0 deletions lib/sessionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,14 @@ export type SessionSnapshot = {
* host-observable without a protocol bump (documented residual on #815).
*/
queue?: string[];
/**
* Plan #960 — explicit view marker: this snapshot is a partial (view-only) window
* rebuilt from a viewport snapshot (`historyComplete: false`). Cloud transcript
* persist (put/mint/upload/flatten) MUST skip it; metadata-only preferences
* (model/reasoning/queue mirror) ride the narrow PATCH channel. `undefined`/`true`
* = a fully restored session. Never inferred from length.
*/
historyComplete?: boolean;
};

import {
Expand Down
17 changes: 11 additions & 6 deletions lib/sessions/sessionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,8 @@ export const RESERVED_META_KEYS = [
'workingNotes',
'freshnessReminderPointer',
'compactionPointer',
/** Plan #960 reserved carrier — JSON-array string (each entry a host-known prompt). `[]` means an explicit empty tombstone, not "unset". */
'queueMirror',
] as const;
export type HarnessSessionMetaKey = (typeof RESERVED_META_KEYS)[number];

Expand Down Expand Up @@ -636,12 +638,15 @@ export function validateMeta(value: unknown): SessionStoreResult<HarnessSessionM
if (isRedisSafeOpaqueId(v)) meta.compactionPointer = v;
continue;
}
if (typeof v !== 'string' && typeof v !== 'number' && typeof v !== 'boolean') {
return {
ok: false,
code: 'invalid_meta',
error: `meta.${key} must be a string, number, or boolean.`,
};
if (key === 'queueMirror') {
// Plan #960 — JSON-array string; invalid JSON drops, never a transfy path.
if (typeof v !== 'string') {
return {
ok: false,
code: 'invalid_meta',
error: 'meta.queueMirror must be a JSON-encoded string.',
};
}
}
meta[key as HarnessSessionMetaKey] = v as string | number | boolean;
}
Expand Down
73 changes: 73 additions & 0 deletions lib/viewportAttach.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
import { afterEach, describe, expect, it, vi } from 'vitest';
import { viewportAttachTurnStream } from './viewportAttach';

function sseResponse(records: string[], viewHeader = true): Response {
const body = new ReadableStream<Uint8Array>({
start(controller) {
const enc = new TextEncoder();
for (const e of records) controller.enqueue(enc.encode(e));
controller.close();
},
});
return new Response(body, {
status: 200,
headers: { 'Content-Type': 'text/event-stream; charset=utf-8', ...(viewHeader ? { 'x-viewport-version': '1' } : {}) },
});
}

describe('viewportAttachTurnStream', () => {
afterEach(() => vi.unstubAllGlobals());

it('GETs cold hydrate=tail (no startIndex) and folds the decoded records', async () => {
const fetchMock = vi.fn(() =>
Promise.resolve(
sseResponse([
'event: viewport_state\ndata: {"type":"viewport_state","version":1,"runId":"wr_1","phase":"recovering","status":"running"}\n\n',
]),
),
);
vi.stubGlobal('fetch', fetchMock);
const events: unknown[] = [];
const result = await viewportAttachTurnStream('wr_1', {
sessionId: 's_1',
onEvent: (rec) => { events.push(rec); },
});
expect(result.ok).toBe(true);
if ('status' in result) expect(result.status).toBe(200);
const url = (fetchMock.mock.calls as unknown as Array<[RequestInfo | URL, RequestInit]>)[0]?.[0] as RequestInfo | URL;
expect(String(url as string).includes('viewportVersion=1&hydrate=tail')).toBe(true);
});

it('GETs indexed startIndex=N when explicitly supplied', async () => {
const fetchMock = vi.fn(() =>
Promise.resolve(
sseResponse([
'event: viewport_state\ndata: {"type":"viewport_state","version":1,"runId":"wr_1","phase":"recovering","status":"running"}\n\n',
]),
),
);
vi.stubGlobal('fetch', fetchMock);
const result = await viewportAttachTurnStream('wr_1', { sessionId: 's_1', startIndex: 7 });
expect(result.ok).toBe(true);
const url = (fetchMock.mock.calls as unknown as Array<[RequestInfo | URL, RequestInit]>)[0]?.[0] as RequestInfo | URL;
expect(String(url as string).includes('startIndex=7')).toBe(true);
});

it('rejects an unknown/wrong viewportVersion (never legacy consume)', async () => {
const fetchMock = vi.fn(() => Promise.resolve(new Response('bad', { status: 400 })));
vi.stubGlobal('fetch', fetchMock);
const result = await viewportAttachTurnStream('wr_1', { sessionId: 's_1' });
expect(result.ok).toBe(false);
if ('status' in result) expect(result.status).toBe(400);
});

it('a legacy 200 SSE body without x-viewport-version is a client error', async () => {
const fetchMock = vi.fn(() =>
Promise.resolve(sseResponse(['data: {"type":"done","text":"x"}'], false)),
);
vi.stubGlobal('fetch', fetchMock);
const result = await viewportAttachTurnStream('wr_1', { sessionId: 's_1' });
expect(result.ok).toBe(false);
if ('error' in result) expect(result.error).toMatch(/Viewport negotiation not accepted/);
});
});
167 changes: 167 additions & 0 deletions lib/viewportAttach.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
import { normalizePrompt } from './chatApi';
import {
AGENT_STREAM_ACCEPT,
type AgentStreamEvent,
} from './agent/agentStream';
import type { AgentFailure } from './agentApi';
import { sanitizeUsageSummary, type UsageSummary } from './agent/usageSummary';
import {
isRedisSafeOpaqueId,
sanitizeResolvedProvider,
sanitizeTurnRunId,
sanitizeTurnStreamCursor,
} from './sessionCloudCaps';
import {
ViewportStreamDecoder,
type ViewportRecord,
type ViewportSnapshot,
} from './viewportStreamProtocol';

function parseTurnRunId(res: Response): string | undefined {
const raw = res.headers.get('x-workflow-run-id');
if (!raw) return undefined;
const trimmed = raw.trim();
return trimmed || undefined;
}

function isAbortError(err: unknown): boolean {
return err instanceof Error && err.name === 'AbortError';
}

export type ViewportNegotiation = 'cold' | 'indexed';

export type ViewportAttachInit = {
sessionId: string;
startIndex?: number;
negotiation?: ViewportNegotiation;
onTurnStarted?: (info: { turnRunId: string }) => Promise<void> | void;
onEvent?: (event: ViewportRecord) => Promise<void> | void;
signal?: AbortSignal;
};

export type ViewportAttachResult =
| AgentFailure
| { ok: true; text: string; turnRunId?: string; cursor?: number; snapshot?: ViewportSnapshot; usage?: UsageSummary; resolvedProvider?: string; }
| { ok: true; text: string; turnRunId?: string; cursor?: number; snapshot?: ViewportSnapshot; usage?: UsageSummary; resolvedProvider?: string; completedAt: 'done' };

/** Read a viewport body and dispatch each decoded record to opts.onEvent. */
export async function readViewportBody(
body: ReadableStream<Uint8Array>,
runId: string,
opts: ViewportAttachInit,
): Promise<{ ok: boolean; error?: string; status?: number; turnRunId?: string; cursor?: number; snapshot?: ViewportSnapshot; text?: string; usage?: UsageSummary; resolvedProvider?: string }> {
const decoder = new ViewportStreamDecoder(runId);
let lastError: AgentFailure | undefined;
let lastTurnEventNextIndex: number | undefined;
let sawViewportEnd = false;
let doneEvent: Extract<AgentStreamEvent, { type: 'done' }> | undefined;
let sawDoneText: string | undefined;
let streamUsage: UsageSummary | undefined;
let streamProvider: string | undefined;
const reader = body.getReader();
for (;;) {
const { value, done } = await reader.read();
if (done) break;
for (const rec of decoder.push(value)) {
if (opts.onEvent) {
await Promise.resolve(opts.onEvent(rec));
}
if (rec.type === 'turn_event' && !rec.skipped && rec.event) {
const ev = rec.event;
if (ev.type === 'done') {
if (doneEvent === undefined) doneEvent = ev;
if (typeof ev.text === 'string') sawDoneText = ev.text;
} else if (ev.type === 'usage') {
streamUsage = sanitizeUsageSummary(ev.usage) ?? streamUsage;
} else if (ev.type === 'provider') {
streamProvider = sanitizeResolvedProvider(ev.provider) ?? streamProvider;
}
lastTurnEventNextIndex = rec.nextIndex;
} else if (rec.type === 'viewport_end') {
sawViewportEnd = true;
} else if (rec.type === 'viewport_error') {
lastError = { ok: false, error: `[viewport_error ${rec.code}] Viewport attach failed.` };
}
if (sawViewportEnd || lastError) break;
}
}
if (!sawViewportEnd && !lastError) {
for (const rec of decoder.push(new Uint8Array(), true)) {
if (opts.onEvent) {
await Promise.resolve(opts.onEvent(rec));
}
}
}
if (lastError) return lastError;
if (!sawViewportEnd) {
return {
ok: true,
text: sawDoneText,
cursor: lastTurnEventNextIndex,
...(doneEvent !== undefined ? { doneEvent } : {}),
...(streamUsage ? { usage: streamUsage } : {}),
...(streamProvider ? { resolvedProvider: streamProvider } : {}),
...(decoder.lastSnapshot !== undefined ? { snapshot: decoder.lastSnapshot } : {}),
};
}
if (doneEvent !== undefined && sawDoneText !== undefined) {
return {
ok: true,
text: sawDoneText,
...(streamUsage ? { usage: streamUsage } : {}),
...(streamProvider ? { resolvedProvider: streamProvider } : {}),
};
}
return { ok: true };
}

/** GET `/api/turns/:runId/stream` with negotiated viewport transport. */
export async function viewportAttachTurnStream(
runId: string,
opts: ViewportAttachInit,
): Promise<AgentFailure | ViewportAttachResult> {
const cleanRunId = sanitizeTurnRunId(runId);
if (!cleanRunId) return { ok: false, status: 400, error: 'Invalid run id' };
if (!normalizePrompt(opts.sessionId ?? '') || !isRedisSafeOpaqueId(opts.sessionId)) {
return { ok: false, status: 400, error: 'Invalid session id' };
}
const params = new URLSearchParams();
params.set('sessionId', opts.sessionId);
params.set('viewportVersion', '1');
if (opts.negotiation === 'indexed' || opts.startIndex !== undefined) {
if (opts.startIndex === undefined) {
return { ok: false, status: 400, error: 'Indexed attach needs startIndex.' };
}
const startIndex = sanitizeTurnStreamCursor(opts.startIndex);
if (startIndex === undefined) {
return { ok: false, status: 400, error: 'Invalid startIndex.' };
}
params.set('startIndex', String(startIndex));
} else {
params.set('hydrate', 'tail');
}
const path = `/api/turns/${encodeURIComponent(cleanRunId)}/stream?${params.toString()}`;
let res: Response;
try {
res = await fetch(path, {
method: 'GET',
headers: { Accept: AGENT_STREAM_ACCEPT },
signal: opts.signal,
});
} catch (err) {
if (isAbortError(err)) return { ok: false, error: 'Aborted.' };
return { ok: false, error: err instanceof Error ? err.message : 'Network request failed.' };
}
const headerRunId = parseTurnRunId(res) ?? cleanRunId;
const contentType = res.headers.get('content-type') ?? '';
if (!res.body || !contentType.includes('text/event-stream') || res.headers.get('x-viewport-version') !== '1') {
const status = res.status;
const error = !contentType.includes('text/event-stream')
? 'View attach failed.'
: 'Viewport negotiation not accepted.';
return { ok: false, status, error: !res.body ? 'Empty viewport stream body.' : error, turnRunId: headerRunId };
}
await opts.onTurnStarted?.({ turnRunId: headerRunId });
const bodyResult = await readViewportBody(res.body!, headerRunId, opts);
return bodyResult as AgentFailure | ViewportAttachResult;
}
13 changes: 13 additions & 0 deletions lib/viewportFirewall.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
/**
* Plan #960 firewall — explicit view/partial guard on CLOUD transcript paths.
* Meta-only preferences (model / reasoning / queue mirror) ride the narrow
* PATCH channel instead of a transcript write. The marker is the snapshot's
* explicit `historyComplete` field, never length inference.
*/
export function assertCompleteForCloudPut(snapshot: { historyComplete?: boolean }): void {
if (snapshot.historyComplete === false) {
throw new Error(
'snapshot has historyComplete:false — view only (blocked from transcript put by the partial-view firewall)',
);
}
}
10 changes: 8 additions & 2 deletions lib/viewportStreamProtocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ export class ViewportStreamDecoder {
private decoder = new TextDecoder('utf-8', { fatal: true });
private pending = '';
private nextIndex: number | undefined;
private lastViewportSnapshot: ViewportSnapshot | undefined;
constructor(private readonly runId: string, startIndex?: number) { this.nextIndex = startIndex; }
push(bytes: Uint8Array, final = false): ViewportRecord[] {
this.pending += this.decoder.decode(bytes, { stream: !final });
Expand All @@ -57,6 +58,9 @@ export class ViewportStreamDecoder {
if (final && this.pending.trim()) throw new Error('Incomplete viewport record');
return out;
}
public get lastSnapshot(): ViewportSnapshot | undefined {
return this.lastViewportSnapshot;
}
private parse(block: string): ViewportRecord {
const fields = (name: string) => block.split('\n').filter(l => l.startsWith(`${name}:`)).map(l => l.slice(name.length + 1).trimStart());
const types = fields('event'), ids = fields('id');
Expand Down Expand Up @@ -98,8 +102,7 @@ export class ViewportStreamDecoder {
}
const start = sanitizeTurnStreamCursor(range.start)!;
const end = sanitizeTurnStreamCursor(range.end)!;
return {
type,
const recSnapshotData: ViewportSnapshot = {
version: 1,
runId: this.runId,
sessionId: o.sessionId as string,
Expand All @@ -114,6 +117,9 @@ export class ViewportStreamDecoder {
sampledRange: { start, end },
...(index !== undefined ? { resumeIndex: index } : {}),
};
const recSnapshot = { type, ...recSnapshotData };
if (this.lastViewportSnapshot === undefined) this.lastViewportSnapshot = recSnapshot;
return recSnapshot as ({ type: 'viewport_snapshot' } & ViewportSnapshot);
}
if (type === 'viewport_state' && o.phase === 'recovering' && typeof o.status === 'string')
return { type, version: 1, runId: this.runId, status: o.status, phase: 'recovering' };
Expand Down
Loading