| 1 | /** Pinned Tier-0 orchestration. Rust holds commands, environment, approval and process trees. */ |
| 2 | import type { OwnerRef, HarnessRunParams, Json } from '../protocol.ts' |
| 3 | import { matchesMatcher } from '../dsh/upstream/hooks/hook-protocol/src/matcher.ts' |
| 4 | import { parseHookOutput } from '../dsh/upstream/hooks/hook-protocol/src/codec.ts' |
| 5 | import { mergeHookOutputs } from '../dsh/upstream/hooks/hook-protocol/src/merge.ts' |
| 6 | import { isJson } from '../json.ts' |
| 7 | import { transformStockSnapshot } from './shared/stock-adapters.ts' |
| 8 | |
| 9 | interface BrokerRpc { request<T = unknown>(method: string, params: unknown, signal?: AbortSignal): Promise<T> } |
| 10 | interface Output { success: boolean; stdout: string; stderr: string } |
| 11 | interface Result { hook_completed?:boolean;hook_skipped?:boolean;proposal?:string;ok: boolean; result?: { content: string; success: boolean; metadata?: Json }; error?: string } |
| 12 | const MAX_PENDING = 32 |
| 13 | const MAX_STDOUT = 1024 * 1024 |
| 14 | const MAX_STDERR = 64 * 1024 |
| 15 | |
| 16 | /** Legacy script JSON/plain-text semantics; no approval or launch facts cross this seam. */ |
| 17 | export function normalizeOutput(value: unknown): Result { |
| 18 | const out = value as Partial<Output> | null |
| 19 | if (!out || typeof out.success !== 'boolean' || typeof out.stdout !== 'string' || typeof out.stderr !== 'string' || Buffer.byteLength(out.stdout) > MAX_STDOUT || Buffer.byteLength(out.stderr) > MAX_STDERR) throw new Error('invalid or oversized execution output') |
| 20 | if (!out.success) return { ok: false, error: out.stderr ? (out.stdout ? `${out.stdout}\n${out.stderr}` : out.stderr) : out.stdout } |
| 21 | try { |
| 22 | const result = JSON.parse(out.stdout) as Record<string, unknown> | null |
| 23 | if (result && typeof result.content === 'string' && typeof result.success === 'boolean' && (result.metadata === undefined || isJson(result.metadata))) { |
| 24 | return { ok: true, result: { content: result.content, success: result.success, ...(result.metadata === undefined ? {} : { metadata: result.metadata as Json }) } } |
| 25 | } |
| 26 | } catch { /* Non-JSON stdout is a successful plain-text result, as in Rust's legacy caller. */ } |
| 27 | return { ok: true, result: { content: out.stdout, success: true } } |
| 28 | } |
| 29 | |
| 30 | export function createHarnessModule(rpc: BrokerRpc, reference: OwnerRef) { |
| 31 | if (reference.plugin_id !== 'host:harness') throw new Error('execution runner requires its builtin owner') |
| 32 | const owner = Object.freeze(structuredClone(reference)) |
| 33 | const stop = new AbortController() |
| 34 | const pending = new Map<string, { abort: AbortController; done: Promise<Result> }>() |
| 35 | let disposed = false |
| 36 | return { |
| 37 | async run(params: HarnessRunParams, signal: AbortSignal): Promise<Result> { |
| 38 | if (disposed || signal.aborted || params.owner.plugin_id !== owner.plugin_id || params.owner.generation !== owner.generation || params.owner.owner_token !== owner.owner_token) throw new Error('execution runner owner is stale or cancelled') |
| 39 | if (!params.execution_id || !params.ticket || !Number.isSafeInteger(params.deadline_ms) || params.deadline_ms < 1 || params.deadline_ms > (params.hook ? 2147483647 : 120_000) || pending.size >= MAX_PENDING || pending.has(params.execution_id)) throw new Error('execution is not admitted') |
| 40 | if(params.hook) { |
| 41 | const {event,dialect,point,matcher,query}=params.hook |
| 42 | if(typeof event!=='string' || typeof point!=='string' || typeof query!=='string' || query.length>1024 || (matcher!==undefined && (typeof matcher!=='string' || matcher.length>1024)) || !['codewhale','claude-code','codex'].includes(dialect))throw new Error('invalid hook projection') |
| 43 | if(dialect!=='codewhale' && !matchesMatcher(matcher,query,dialect as 'claude-code'|'codex'))return {ok:true,hook_skipped:true} |
| 44 | } |
| 45 | const abort = new AbortController() |
| 46 | const cancel = () => abort.abort() |
| 47 | signal.addEventListener('abort', cancel, { once: true }); stop.signal.addEventListener('abort', cancel, { once: true }) |
| 48 | const timer = setTimeout(cancel, params.deadline_ms) |
| 49 | const done = (async () => { |
| 50 | const cancelled = new Promise<never>((_, reject) => { abort.signal.addEventListener('abort', () => reject(new Error('execution cancelled')), { once: true }) }) |
| 51 | let value = await Promise.race([rpc.request('exec/redeem', { owner, execution_id: params.execution_id, ticket: params.ticket }, abort.signal), cancelled]) |
| 52 | if (disposed || abort.signal.aborted) throw new Error('execution cancelled') |
| 53 | if(params.hook) { |
| 54 | const out=value as any |
| 55 | if(!out || out.kind!=='hook' || out.event!==params.hook.event || typeof out.success!=='boolean' || (out.exit_code!==null && !Number.isInteger(out.exit_code)))throw new Error('invalid hook completion projection') |
| 56 | if(params.hook.event==='shell_env') { |
| 57 | if('stdout' in out || 'stderr' in out || Object.keys(out).some(key=>!['kind','event','success','exit_code','keys'].includes(key)))throw new Error('ShellEnv projection contains private process data') |
| 58 | return {ok:true,hook_completed:true} |
| 59 | } |
| 60 | if(params.hook.dialect==='codewhale')return {ok:true,hook_completed:true} |
| 61 | if(typeof out.stdout!=='string' || typeof out.stderr!=='string' || Buffer.byteLength(out.stdout)>65536 || Buffer.byteLength(out.stderr)>65536)throw new Error('invalid dialect hook output') |
| 62 | const decoded=parseHookOutput(out.exit_code===null?undefined:out.exit_code,out.stdout,out.stderr,params.hook.point) |
| 63 | const merged=mergeHookOutputs([decoded]) |
| 64 | if(params.hook.event==='tool_call_before')return {ok:true,hook_completed:true,proposal:JSON.stringify({...(merged.decision==='deny'?{decision:'deny',reason:merged.reason ?? 'Blocked by hook'}:merged.decision==='ask' && params.hook.dialect!=='codex'?{decision:'ask',reason:merged.reason}:{}),...(merged.additionalContext.length?{additional_context:merged.additionalContext.join('\n\n')}:{})})} |
| 65 | if(params.hook.event==='message_submit' && (merged.stop || decoded.updatedInput!==undefined || merged.additionalContext.length))throw new Error('hook requested unsupported prompt steering') |
| 66 | if(params.hook.event==='message_submit')return {ok:true,hook_completed:true,proposal:JSON.stringify(merged.decision==='deny'?{block:true,reason:merged.reason ?? 'Blocked by hook'}:{})} |
| 67 | if(merged.stop || merged.decision==='deny' || decoded.updatedInput!==undefined || merged.additionalContext.length)throw new Error('hook requested steering unavailable at this observer boundary') |
| 68 | return {ok:true,hook_completed:true} |
| 69 | } |
| 70 | if (value && typeof value === 'object' && (value as Record<string, unknown>).kind === 'ocr_process') { |
| 71 | if ((value as Record<string, unknown>).state !== 'native') throw new Error('OCR must begin with its admitted Native step') |
| 72 | const first = transformStockSnapshot(value) |
| 73 | if ((first.result.metadata as Record<string, unknown>).code !== 'fallback') return first |
| 74 | const next = (value as Record<string, unknown>).next_ticket as string |
| 75 | // A separate, single-use grant exists only after Core captured Native |
| 76 | // failure/unavailability. No command/path/stage choice crosses IPC. |
| 77 | value = await Promise.race([rpc.request('exec/redeem', { owner, execution_id: params.execution_id, ticket: next }, abort.signal), cancelled]) |
| 78 | if (disposed || abort.signal.aborted) throw new Error('execution cancelled') |
| 79 | if (!value || typeof value !== 'object' || (value as Record<string, unknown>).kind !== 'ocr_process' || (value as Record<string, unknown>).state !== 'tesseract') throw new Error('OCR continuation returned the wrong stage') |
| 80 | return transformStockSnapshot(value) |
| 81 | } |
| 82 | if (value && typeof value === 'object' && ['stock_adapter', 'pdf_process'].includes(String((value as Record<string, unknown>).kind))) return transformStockSnapshot(value) |
| 83 | return normalizeOutput(value) |
| 84 | })() |
| 85 | pending.set(params.execution_id, { abort, done }) |
| 86 | try { return await done } finally { clearTimeout(timer); signal.removeEventListener('abort', cancel); stop.signal.removeEventListener('abort', cancel); pending.delete(params.execution_id) } |
| 87 | }, |
| 88 | async dispose(): Promise<void> { if (disposed) return; disposed = true; stop.abort(); await Promise.allSettled([...pending.values()].map(row => row.done)) }, |
| 89 | } |
| 90 | } |
| 91 |