From 3ad9bb74139ee4e89b729da4ef7d8aca9c5d98f6 Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sat, 1 Aug 2026 11:47:45 +0100 Subject: [PATCH 1/6] feat(overseer): ping_session relay write-tool (Stage 1.5) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Gives the Overseer an operator-directed way to ping an individual project session — the delegation brick from the action architecture (R5). - New write tool `ping_session` (sessionId and/or itemId → related session, then resume-if-inactive + enqueue). Wired through SyncEngine's existing resumeSession/sendMessage primitives — never shells out to hapi-ping-peer. - Identity gains `canRelay`; catalog now has 11 tools (2 writes). System prompt updated for Stage 1.5 relay discipline (imperative-only, tombstone). - runOverseerTool is async; HTTP tool surface still 403s writes (converse allowWrites remains the only write path). Also: operator hub URL is now https://hapi-gc-oos.forest-adder.ts.net (svc:hapi-gc-oos; old *.tail9944ee.ts.net MagicDNS is dead). Co-authored-by: Cursor --- hub/src/overseer/converse.ts | 2 +- hub/src/overseer/runOverseerTool.ts | 13 +- hub/src/sync/overseerEntity.test.ts | 62 +++++++- hub/src/sync/overseerEntity.ts | 133 ++++++++++++++++++ hub/src/sync/syncEngine.ts | 30 +++- hub/src/web/routes/overseer.test.ts | 12 +- hub/src/web/routes/overseer.ts | 2 +- shared/src/overseerConverse.ts | 9 +- shared/src/overseerEntity.test.ts | 16 ++- shared/src/overseerEntity.ts | 79 +++++++++-- .../settings/OverseerChatDebugControls.tsx | 2 +- web/src/routes/overseer/index.tsx | 1 + 12 files changed, 326 insertions(+), 35 deletions(-) diff --git a/hub/src/overseer/converse.ts b/hub/src/overseer/converse.ts index e9108a7fcd..053e6453c4 100644 --- a/hub/src/overseer/converse.ts +++ b/hub/src/overseer/converse.ts @@ -124,7 +124,7 @@ export async function runOverseerConverse(params: { try { // The conversational surface is the operator-directed write-path, so dispositions // are allowed here (gated off on the raw HTTP tool-dispatch endpoint). - const result = runOverseerTool(overseer, name, args, true) + const result = await runOverseerTool(overseer, name, args, true) toolTrace.push({ tool: name, args, ok: true }) // The brain opts into 'full' per call when it needs depth; default lean. const detail = args.detail === 'full' ? 'full' : 'lean' diff --git a/hub/src/overseer/runOverseerTool.ts b/hub/src/overseer/runOverseerTool.ts index 03d2b69320..09330c7715 100644 --- a/hub/src/overseer/runOverseerTool.ts +++ b/hub/src/overseer/runOverseerTool.ts @@ -6,7 +6,7 @@ import { } from '@hapi/protocol' import type { OverseerEntity } from '../sync/overseerEntity' -/** Thrown when a write tool (`record_disposition`) is dispatched on a read-only surface (R2 gate). */ +/** Thrown when a write tool is dispatched on a read-only surface (R2 gate). */ export class OverseerWriteNotAllowedError extends Error { constructor(tool: string) { super(`Tool "${tool}" writes and is not allowed on this surface`) @@ -17,15 +17,16 @@ export class OverseerWriteNotAllowedError extends Error { /** * Execute one Overseer tool by name against the entity. Shared by the HTTP tool-dispatch route and * the converse tool-calling loop so both go through exactly one place. Throws `ZodError` on invalid - * args. Every tool is read-only EXCEPT `record_disposition`; writes are gated behind `allowWrites` - * (the conversational path sets it; the raw HTTP dispatch does not). + * args. Write tools (`record_disposition`, `ping_session`) are gated behind `allowWrites` + * (the conversational path sets it; the raw HTTP dispatch does not). Async because `ping_session` + * may resume a worker before enqueueing. */ -export function runOverseerTool( +export async function runOverseerTool( overseer: OverseerEntity, tool: OverseerToolName, args: unknown, allowWrites = false -): unknown { +): Promise { if (isOverseerWriteTool(tool) && !allowWrites) { throw new OverseerWriteNotAllowedError(tool) } @@ -58,6 +59,8 @@ export function runOverseerTool( return overseer.queryDispositions(overseerToolArgsSchemas.query_dispositions.parse(args)) case 'record_disposition': return overseer.recordDisposition(overseerToolArgsSchemas.record_disposition.parse(args)) + case 'ping_session': + return overseer.pingSession(overseerToolArgsSchemas.ping_session.parse(args)) default: { const exhaustive: never = tool throw new Error(`Unknown overseer tool: ${String(exhaustive)}`) diff --git a/hub/src/sync/overseerEntity.test.ts b/hub/src/sync/overseerEntity.test.ts index 487a2dbc29..ddb8b53fb1 100644 --- a/hub/src/sync/overseerEntity.test.ts +++ b/hub/src/sync/overseerEntity.test.ts @@ -4,7 +4,7 @@ import { Store } from '../store' import { SyncEngine } from './syncEngine' import { RpcRegistry } from '../socket/rpcRegistry' import { OverseerWriteNotAllowedError, runOverseerTool } from '../overseer/runOverseerTool' -import type { OverseerEntity } from './overseerEntity' +import { OverseerEntity } from './overseerEntity' function makeEngine(): { store: Store; engine: SyncEngine } { const store = new Store(':memory:') @@ -400,17 +400,71 @@ describe('OverseerEntity dispositions (Stage 1 keystone)', () => { expect(o.queryDispositions({}).total).toBe(0) }) - it('record_disposition is gated: runOverseerTool refuses the write unless allowWrites', () => { + it('record_disposition is gated: runOverseerTool refuses the write unless allowWrites', async () => { const store = new Store(':memory:') const { itemId } = promoteItem(store, { key: 'gate', eventType: 'blocked', project: 'hapi' }) const o = overseer(buildEngine(store)) - expect(() => runOverseerTool(o, 'record_disposition', { itemId, action: 'done' })).toThrow( + await expect(runOverseerTool(o, 'record_disposition', { itemId, action: 'done' })).rejects.toBeInstanceOf( OverseerWriteNotAllowedError ) // With writes allowed (the conversational path) it lands. - const res = runOverseerTool(o, 'record_disposition', { itemId, action: 'done' }, true) as { + const res = await runOverseerTool(o, 'record_disposition', { itemId, action: 'done' }, true) as { ok: boolean } expect(res.ok).toBe(true) }) + + it('ping_session resolves by sessionId / itemId and returns a tombstone via injected relay', async () => { + const store = new Store(':memory:') + const { itemId } = promoteItem(store, { + key: 'expenses', + eventType: 'blocked', + project: 'expenses' + }) + const item = store.inbox.getById(itemId)! + const sessionId = item.relatedSessionId! + expect(sessionId).toBeTruthy() + + let lastRelay: { sessionId: string; message: string } | undefined + const o = new OverseerEntity({ + events: store.events, + inbox: store.inbox, + messages: store.messages, + getSession: (id) => { + const s = store.sessions.getSession(id) + if (!s) return undefined + return { ...s, active: true, namespace: s.namespace || 'default' } as never + }, + getSessions: () => { + const s = store.sessions.getSession(sessionId) + return s ? [{ ...s, active: true, namespace: s.namespace || 'default' } as never] : [] + }, + relayToSession: async ({ sessionId: sid, message }) => { + lastRelay = { sessionId: sid, message } + return { ok: true, resumed: false } + } + }) + + const byItem = await o.pingSession({ + itemId, + message: 'Please draft the Cursor Pro June note.' + }) + expect(byItem.ok).toBe(true) + expect(byItem.sessionId).toBe(sessionId) + expect(byItem.tombstone).toContain('Relayed') + expect(lastRelay).toEqual({ + sessionId, + message: 'Please draft the Cursor Pro June note.' + }) + const byPrefix = await o.pingSession({ + sessionId: sessionId.slice(0, 8), + message: 'Second ping' + }) + expect(byPrefix.ok).toBe(true) + expect(byPrefix.sessionId).toBe(sessionId) + + await expect(runOverseerTool(o, 'ping_session', { sessionId, message: 'x' })).rejects.toBeInstanceOf( + OverseerWriteNotAllowedError + ) + }) }) diff --git a/hub/src/sync/overseerEntity.ts b/hub/src/sync/overseerEntity.ts index a638176a32..9184409aad 100644 --- a/hub/src/sync/overseerEntity.ts +++ b/hub/src/sync/overseerEntity.ts @@ -39,6 +39,8 @@ import { type QueryInboxArgs, type QueryOpenLoopsArgs, type RecordDispositionArgs, + type PingSessionArgs, + type OverseerPingResult, type ListActiveWorkersArgs } from '@hapi/protocol' import { buildOverseerSessionIdentity } from '@hapi/protocol' @@ -54,6 +56,15 @@ export type OverseerEntityDeps = { messages: MessageStore getSession: (sessionId: string) => Session | undefined getSessions: () => Session[] + /** + * Stage 1.5 relay (R5): resume if inactive, then enqueue a user message. + * Injected from SyncEngine — never shell out to `hapi-ping-peer`. + */ + relayToSession?: (args: { + sessionId: string + message: string + namespace?: string + }) => Promise<{ ok: boolean; resumed: boolean; error?: string }> now?: () => number staleSilenceMs?: number } @@ -87,6 +98,7 @@ export class OverseerEntity { private readonly messages: MessageStore private readonly getSession: (sessionId: string) => Session | undefined private readonly getSessions: () => Session[] + private readonly relayToSession: OverseerEntityDeps['relayToSession'] private readonly now: () => number private readonly staleSilenceMs: number @@ -96,6 +108,7 @@ export class OverseerEntity { this.messages = deps.messages this.getSession = deps.getSession this.getSessions = deps.getSessions + this.relayToSession = deps.relayToSession this.now = deps.now ?? (() => Date.now()) this.staleSilenceMs = deps.staleSilenceMs ?? OVERSEER_STALE_SILENCE_MS } @@ -532,6 +545,126 @@ export class OverseerEntity { } } + /** + * Stage 1.5 write — relay an operator-directed message to one worker session (R5). + * Resolves sessionId (full or unique prefix) and/or itemId → relatedSessionId, then + * calls the injected SyncEngine resume+send primitive. + */ + async pingSession(args: PingSessionArgs): Promise { + const resolved = this.resolvePingTarget(args) + if (!resolved.ok) { + return { + ok: false, + sessionId: resolved.sessionId ?? '', + sessionName: null, + project: null, + resumed: false, + tombstone: resolved.error, + error: resolved.error + } + } + + if (!this.relayToSession) { + return { + ok: false, + sessionId: resolved.sessionId, + sessionName: resolved.sessionName, + project: resolved.project, + resumed: false, + tombstone: 'Relay not wired on this hub — nothing sent.', + error: 'relay_not_configured' + } + } + + const result = await this.relayToSession({ + sessionId: resolved.sessionId, + message: args.message, + namespace: resolved.namespace + }) + + const label = resolved.sessionName ?? resolved.project ?? resolved.sessionId.slice(0, 8) + const snippet = args.message.length > 80 ? `${args.message.slice(0, 77)}…` : args.message + if (!result.ok) { + return { + ok: false, + sessionId: resolved.sessionId, + sessionName: resolved.sessionName, + project: resolved.project, + resumed: result.resumed, + tombstone: `Failed to relay to ${label}: ${result.error ?? 'unknown error'}`, + error: result.error + } + } + + return { + ok: true, + sessionId: resolved.sessionId, + sessionName: resolved.sessionName, + project: resolved.project, + resumed: result.resumed, + tombstone: `Relayed to ${label} (${resolved.sessionId.slice(0, 8)})${result.resumed ? ' [resumed]' : ''}: "${snippet}"` + } + } + + private resolvePingTarget(args: PingSessionArgs): { + ok: true + sessionId: string + sessionName: string | null + project: string | null + namespace: string + } | { ok: false; sessionId?: string; error: string } { + let sessionId = args.sessionId?.trim() || null + + if (!sessionId && args.itemId != null) { + const item = this.inbox.getById(args.itemId) + if (!item) { + return { ok: false, error: `No inbox item #${args.itemId} — nothing relayed.` } + } + sessionId = item.relatedSessionId?.trim() || null + if (!sessionId) { + return { + ok: false, + error: `Inbox item #${args.itemId} has no related session — nothing relayed.` + } + } + } + + if (!sessionId) { + return { ok: false, error: 'sessionId or itemId is required — nothing relayed.' } + } + + // Exact match first, then unique prefix (same as hapi-ping-peer). + let session = this.getSession(sessionId) + if (!session) { + const matches = this.getSessions().filter((s) => s.id.startsWith(sessionId!)) + if (matches.length === 1) { + session = matches[0] + sessionId = session.id + } else if (matches.length > 1) { + return { + ok: false, + sessionId, + error: `Ambiguous session prefix "${sessionId}" (${matches.length} matches) — nothing relayed.` + } + } else { + return { + ok: false, + sessionId, + error: `No session matching "${sessionId}" — nothing relayed.` + } + } + } + + const identity = deriveIdentity(session) + return { + ok: true, + sessionId: session.id, + sessionName: identity.name, + project: identity.project, + namespace: session.namespace || 'default' + } + } + private buildTombstone(action: string, statusAfter: string, item: StoredInboxItem): string { const verb = action === 'done' diff --git a/hub/src/sync/syncEngine.ts b/hub/src/sync/syncEngine.ts index 8a4d75e920..a4989a9f40 100644 --- a/hub/src/sync/syncEngine.ts +++ b/hub/src/sync/syncEngine.ts @@ -172,7 +172,35 @@ export class SyncEngine { inbox: store.inbox, messages: store.messages, getSession: (sessionId) => this.getSession(sessionId), - getSessions: () => this.getSessions() + getSessions: () => this.getSessions(), + relayToSession: async ({ sessionId, message, namespace = 'default' }) => { + const existing = this.getSession(sessionId) + if (!existing) { + return { ok: false, resumed: false, error: 'session_not_found' } + } + let resumed = false + if (!existing.active) { + const resume = await this.resumeSession(sessionId, namespace) + if (resume.type !== 'success') { + return { + ok: false, + resumed: false, + error: resume.code ?? resume.message ?? 'resume_failed' + } + } + resumed = true + } + try { + await this.sendMessage(sessionId, { text: message, sentFrom: 'webapp' }) + return { ok: true, resumed } + } catch (error) { + return { + ok: false, + resumed, + error: error instanceof Error ? error.message : 'send_failed' + } + } + } }) this.reloadAll() this.inactivityTimer = setInterval(() => this.expireInactive(), 5_000) diff --git a/hub/src/web/routes/overseer.test.ts b/hub/src/web/routes/overseer.test.ts index 6b953b5228..da82940727 100644 --- a/hub/src/web/routes/overseer.test.ts +++ b/hub/src/web/routes/overseer.test.ts @@ -29,7 +29,7 @@ describe('overseer routes', () => { systemPrompt: string } expect(body.identity.canDispatch).toBe(false) - expect(body.identity.tools.length).toBe(10) + expect(body.identity.tools.length).toBe(11) expect(body.systemPrompt).toContain('Overseer') }) @@ -44,6 +44,16 @@ describe('overseer routes', () => { expect(res.status).toBe(403) }) + it('POST /overseer/tools/ping_session is gated off on the read-only HTTP surface (403)', async () => { + const app = buildApp(new Store(':memory:')) + const res = await app.request('/api/overseer/tools/ping_session', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ sessionId: 'abc', message: 'hi' }) + }) + expect(res.status).toBe(403) + }) + it('GET /overseer/voice returns prompt + backend descriptor', async () => { const app = buildApp(new Store(':memory:')) const res = await app.request('/api/overseer/voice') diff --git a/hub/src/web/routes/overseer.ts b/hub/src/web/routes/overseer.ts index 7fc498320e..0a4c23f143 100644 --- a/hub/src/web/routes/overseer.ts +++ b/hub/src/web/routes/overseer.ts @@ -170,7 +170,7 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho } try { - const result = runOverseerTool(engine.getOverseer(), tool, body ?? {}) + const result = await runOverseerTool(engine.getOverseer(), tool, body ?? {}) return c.json({ tool, result }) } catch (error) { if (error instanceof z.ZodError) { diff --git a/shared/src/overseerConverse.ts b/shared/src/overseerConverse.ts index 6283c306f4..34e9fa5018 100644 --- a/shared/src/overseerConverse.ts +++ b/shared/src/overseerConverse.ts @@ -175,10 +175,15 @@ const OVERSEER_TOOL_PARAMS: Record = { action: { type: 'string', enum: [...OVERSEER_DISPOSITION_ACTIONS], description: 'done=resolve, dismiss=tombstone, snooze (needs snoozedUntil), open=reopen.' }, feedback: { type: 'string', description: 'Optional operator note / learning label to freeze with the disposition.' }, snoozedUntil: { type: 'integer', minimum: 1, description: 'Required for snooze: epoch ms to sleep the item until.' } - }, ['itemId', 'action']) + }, ['itemId', 'action']), + ping_session: obj({ + sessionId: { type: 'string', description: 'Worker session id (full UUID or unique prefix).' }, + itemId: { type: 'integer', minimum: 1, description: 'Inbox item id — resolves its relatedSessionId when sessionId omitted.' }, + message: { type: 'string', description: 'Operator-directed message to relay to that session.' } + }, ['message']) } -/** The Overseer tool catalog (read-only + the single disposition write) as an OpenAI `tools` array. */ +/** The Overseer tool catalog (reads + disposition + relay writes) as an OpenAI `tools` array. */ export function buildOverseerOpenAiTools(): OverseerOpenAiTool[] { return OVERSEER_TOOL_CATALOG.map((entry) => ({ type: 'function', diff --git a/shared/src/overseerEntity.test.ts b/shared/src/overseerEntity.test.ts index 3d97c895b6..ed1da61848 100644 --- a/shared/src/overseerEntity.test.ts +++ b/shared/src/overseerEntity.test.ts @@ -19,33 +19,35 @@ import { const STALE = 30 * 60 * 1000 describe('overseer entity protocol', () => { - it('catalog covers every tool name; only record_disposition writes (R2)', () => { + it('catalog covers every tool name; only write tools are non-readonly (R2)', () => { const catalogNames = OVERSEER_TOOL_CATALOG.map((t) => t.name).sort() expect(catalogNames).toEqual([...OVERSEER_TOOL_NAMES].sort()) - const writeTools = OVERSEER_TOOL_CATALOG.filter((t) => !t.readonly).map((t) => t.name) - expect(writeTools).toEqual(['record_disposition']) + const writeTools = OVERSEER_TOOL_CATALOG.filter((t) => !t.readonly).map((t) => t.name).sort() + expect(writeTools).toEqual(['ping_session', 'record_disposition']) expect(isOverseerWriteTool('record_disposition')).toBe(true) + expect(isOverseerWriteTool('ping_session')).toBe(true) expect(isOverseerWriteTool('query_events')).toBe(false) }) - it('exposes a cannot-dispatch identity that CAN record dispositions (Stage 1)', () => { + it('exposes a cannot-dispatch identity that CAN disposition + relay (Stage 1.5)', () => { const identity = buildOverseerIdentity() expect(identity.id).toBe(OVERSEER_ENTITY_ID) expect(identity.kind).toBe(OVERSEER_SOURCE_KIND) expect(identity.canDispatch).toBe(false) expect(identity.canDisposition).toBe(true) + expect(identity.canRelay).toBe(true) expect(identity.tools).toHaveLength(OVERSEER_TOOL_NAMES.length) }) - it('system prompt frames chief-of-staff + read-only tools + disposition write discipline', () => { + it('system prompt frames chief-of-staff + read tools + disposition/relay write discipline', () => { const prompt = buildOverseerSystemPrompt() expect(prompt).toContain('chief-of-staff') expect(prompt).toContain('Read-only tools') expect(prompt).toContain('record_disposition') - expect(prompt).toContain('CANNOT dispatch') + expect(prompt).toContain('ping_session') + expect(prompt).toContain('CANNOT spawn') expect(prompt).toContain('Show receipts') }) - describe('worker state derivation', () => { it('maps notify status and event type to worker state', () => { expect(mapNotifyStatusToWorkerState('blocked')).toBe('blocked') diff --git a/shared/src/overseerEntity.ts b/shared/src/overseerEntity.ts index 34ca135b29..ca514b66f9 100644 --- a/shared/src/overseerEntity.ts +++ b/shared/src/overseerEntity.ts @@ -266,17 +266,18 @@ export const OVERSEER_TOOL_NAMES = [ 'list_active_workers', 'query_open_loops', 'query_dispositions', - 'record_disposition' + 'record_disposition', + 'ping_session' ] as const export type OverseerToolName = typeof OVERSEER_TOOL_NAMES[number] /** - * The one tool that writes (Stage 0→1 keystone). Every other tool is read-only against the - * substrate; `record_disposition` is the single explicit, operator-directed mutation. It is - * clearly marked non-`readonly` in the catalog so the write surface stays enumerable (R2). + * Write tools (Stage 1+). Every other tool is read-only against the substrate. + * `record_disposition` = operator decision on an inbox item; `ping_session` = relay + * a message to a worker session (resume + enqueue). Both are operator-directed only (R2/R5). */ -export const OVERSEER_WRITE_TOOL_NAMES = ['record_disposition'] as const +export const OVERSEER_WRITE_TOOL_NAMES = ['record_disposition', 'ping_session'] as const export type OverseerWriteToolName = typeof OVERSEER_WRITE_TOOL_NAMES[number] export function isOverseerWriteTool(name: string): name is OverseerWriteToolName { @@ -405,6 +406,33 @@ export const recordDispositionArgsSchema = z.object({ }) export type RecordDispositionArgs = z.infer +/** + * `ping_session` — relay an operator-directed message to one worker session (R5). + * Resolves by `sessionId` (full id or unique prefix) and/or `itemId` (inbox item → + * relatedSessionId). Hub resumes if inactive, then enqueues the message — same + * primitives as `hapi-ping-peer`, called in-process (never shell out). + */ +export const pingSessionArgsSchema = z.object({ + sessionId: z.string().min(1).max(128).optional(), + itemId: z.number().int().positive().optional(), + message: z.string().min(1).max(8000), +}).refine((v) => Boolean(v.sessionId?.trim()) || typeof v.itemId === 'number', { + message: 'sessionId or itemId is required' +}) +export type PingSessionArgs = z.infer + +export type OverseerPingResult = { + ok: boolean + sessionId: string + sessionName: string | null + project: string | null + /** True when the hub had to resume before sending. */ + resumed: boolean + /** One-line human confirmation to read back ("Relayed to Expenses (a492…): …"). */ + tombstone: string + error?: string +} + /** * One cold open loop: a session whose latest worker status is not `done`, * carrying how long it has sat and which lens bucket it belongs to. @@ -486,7 +514,8 @@ export const overseerToolArgsSchemas = { list_active_workers: listActiveWorkersArgsSchema, query_open_loops: queryOpenLoopsArgsSchema, query_dispositions: queryDispositionsArgsSchema, - record_disposition: recordDispositionArgsSchema + record_disposition: recordDispositionArgsSchema, + ping_session: pingSessionArgsSchema } as const satisfies Record export type OverseerToolCatalogEntry = { @@ -496,7 +525,7 @@ export type OverseerToolCatalogEntry = { readonly: boolean } -/** Catalog surfaced to the voice/system layer; exactly one entry (`record_disposition`) writes. */ +/** Catalog surfaced to the voice/system layer; write tools are marked `readonly: false` (R2). */ export const OVERSEER_TOOL_CATALOG: OverseerToolCatalogEntry[] = [ { name: 'query_events', @@ -547,6 +576,11 @@ export const OVERSEER_TOOL_CATALOG: OverseerToolCatalogEntry[] = [ name: 'record_disposition', description: 'WRITE: record the operator\'s explicit decision on one inbox item — done (resolve), dismiss (tombstone), snooze (needs snoozedUntil), or open (reopen). Only call this when the operator has clearly directed it in the conversation; never on your own judgement. Returns a tombstone to read back.', readonly: false + }, + { + name: 'ping_session', + description: 'WRITE: relay an operator-directed message to one worker session (resume if inactive, then enqueue). Pass sessionId (full id or unique prefix) and/or itemId (inbox item → its related session). Only call when the operator has clearly said to tell / ping / ask that project or agent something. Irreversible once delivered — never invent a ping.', + readonly: false } ] @@ -555,6 +589,11 @@ export function overseerCanDisposition(): boolean { return OVERSEER_TOOL_CATALOG.some((t) => t.name === 'record_disposition' && !t.readonly) } +/** Whether the Overseer may relay to a worker session — Stage 1.5 delegation gate (R5). */ +export function overseerCanRelay(): boolean { + return OVERSEER_TOOL_CATALOG.some((t) => t.name === 'ping_session' && !t.readonly) +} + export type OverseerIdentity = { id: string kind: typeof OVERSEER_SOURCE_KIND @@ -562,6 +601,8 @@ export type OverseerIdentity = { canDispatch: false /** Stage 1 keystone: the Overseer may record operator-directed dispositions on inbox items. */ canDisposition: boolean + /** Stage 1.5: the Overseer may relay operator-directed messages to a worker session. */ + canRelay: boolean tools: OverseerToolCatalogEntry[] } @@ -571,6 +612,7 @@ export function buildOverseerIdentity(): OverseerIdentity { kind: OVERSEER_SOURCE_KIND, canDispatch: false, canDisposition: overseerCanDisposition(), + canRelay: overseerCanRelay(), tools: OVERSEER_TOOL_CATALOG } } @@ -594,9 +636,10 @@ export function buildOverseerSystemPrompt(): string { 'You hold a continuous view of the whole fleet and speak to the operator about it. You are not', 'any single worker, and you never speak as one.', '', - '# What you can do (Stage 1 — read + record dispositions)', + '# What you can do (Stage 1.5 — read + dispositions + relay)', '', - 'You can READ the fleet, ANSWER questions, and RECORD the operator\'s decisions on inbox items.', + 'You can READ the fleet, ANSWER questions, RECORD the operator\'s decisions on inbox items,', + 'and RELAY an operator-directed message to one worker session.', '', 'Read-only tools:', '- query_events — the events stream (blockers, completions, decisions, progress, errors).', @@ -609,12 +652,14 @@ export function buildOverseerSystemPrompt(): string { '- query_open_loops — the "what am I forgetting?" lens: cold threads whose latest status is not done.', '- query_dispositions — past operator decisions (list, or groupBy+minCount to cluster them).', '', - 'Write tool (the ONLY thing you can change):', + 'Write tools (the ONLY things you can change):', '- record_disposition — record the operator\'s decision on ONE inbox item: done (resolve),', ' dismiss (tombstone), snooze (needs snoozedUntil), or open (reopen).', + '- ping_session — relay a message to ONE worker session (resume if needed, then enqueue).', + ' Pass sessionId and/or itemId. Irreversible once delivered.', '', - 'You still CANNOT dispatch, message workers, spawn, or confirm anything on a worker. If the', - 'operator asks you to act on a worker, say you can advise but cannot dispatch yet.', + 'You still CANNOT spawn new workers, invent work, or confirm anything on a worker\'s behalf.', + 'If the operator asks you to "just handle it" without naming a target and an intent, ask.', '', '# Recording a disposition (be careful — this writes)', '', @@ -625,6 +670,16 @@ export function buildOverseerSystemPrompt(): string { ' / explain_priority) so you pass the right itemId.', '- After it lands, read the returned tombstone back in one line so the operator knows it stuck.', '', + '# Relaying to a project / session (be careful — this writes and is not undoable)', + '', + '- Call ping_session ONLY when the operator has clearly directed a relay: "tell that expenses', + ' session…", "ping the hapi peer about…", "ask that project to…". Never invent a ping.', + '- Prefer itemId when the conversation is about a specific inbox item (it resolves the related', + ' session). Otherwise use sessionId (full id or unique short prefix).', + '- Keep the message short and complete — the worker rehydrates its own context; you are a', + ' secretary passing intent, not transferring your whole briefing.', + '- After it lands, read the returned tombstone back in one line.', + '', '# How to answer', '', '- Lead with the answer, not the method. "Peer 15 is blocked on CI auth" — not "let me check".', diff --git a/web/src/components/settings/OverseerChatDebugControls.tsx b/web/src/components/settings/OverseerChatDebugControls.tsx index be19d13a78..815d7c5e81 100644 --- a/web/src/components/settings/OverseerChatDebugControls.tsx +++ b/web/src/components/settings/OverseerChatDebugControls.tsx @@ -110,7 +110,7 @@ export function OverseerChatDebugControls() { {open && (

- Fleet chief-of-staff (Stage 1). Text transport over the same converse core voice will use. The brain can read fleet state and record operator-directed dispositions (done / dismiss / snooze / open) on inbox items — it still cannot dispatch or message workers. + Fleet chief-of-staff (Stage 1.5 — read + dispositions + relay). Text transport over the same converse core voice will use. The brain can read fleet state, record operator-directed dispositions (done / dismiss / snooze / open), and relay a message to a worker when you explicitly ask (ping / tell / snooze…).

diff --git a/web/src/routes/overseer/index.tsx b/web/src/routes/overseer/index.tsx index 3326fd2b9a..301988ecfa 100644 --- a/web/src/routes/overseer/index.tsx +++ b/web/src/routes/overseer/index.tsx @@ -42,6 +42,7 @@ function OverseerIdentityPanel() { {identity.id} · {identity.tools.length} tool{identity.tools.length === 1 ? '' : 's'} {identity.canDisposition ? ' · can disposition' : ''} + {identity.canRelay ? ' · can relay' : ''} {open ? ( From 6a6391cdfb4b6f1d418090527270976c2eb872ba Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sat, 1 Aug 2026 12:55:50 +0100 Subject: [PATCH 2/6] fix(overseer): resolve unique session-id prefixes on read tools MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit get_session_state/get_worker_health/query_events only did exact UUID match, so the brain's truncated 8-char ids returned null and Overseer falsely called live inbox rows "ghost sessions". Shared unique-prefix resolve (same as ping_session); lean inbox keeps full relatedSessionId + summary; prompt clarifies null ≠ deleted. Co-authored-by: Cursor --- hub/src/overseer/toolProjection.test.ts | 15 +++- hub/src/overseer/toolProjection.ts | 12 ++- hub/src/sync/overseerEntity.test.ts | 52 +++++++++++ hub/src/sync/overseerEntity.ts | 111 +++++++++++++++--------- shared/src/overseerEntity.ts | 31 +++++-- 5 files changed, 170 insertions(+), 51 deletions(-) diff --git a/hub/src/overseer/toolProjection.test.ts b/hub/src/overseer/toolProjection.test.ts index a744575700..74e9709121 100644 --- a/hub/src/overseer/toolProjection.test.ts +++ b/hub/src/overseer/toolProjection.test.ts @@ -14,12 +14,13 @@ describe('projectToolResultForBrain', () => { expect(lean.counts).toEqual({ candidates: 1, surfaced: 1, held: 0 }) }) - it('thins inbox items to id/what/status/priority and drops the fat', () => { + it('thins inbox items to id/what/summary/category/session/status/priority and drops the fat', () => { const full = { items: [ { id: 7, title: 'CI auth blocking 3 workers', status: 'surfaced', priority: 90, category: 'blocker', summary: 'long summary…', reasonForPriority: 'shared root cause', + relatedSessionId: '96f67085-5dd3-4a10-aa7c-785f72a227c2', sourceEventIds: [1, 2, 3], artifactRefs: ['a'.repeat(400)], createdAt: 1, updatedAt: 2 }, { id: 8, title: 'needs a decision', status: 'new', priority: 40, artifactRefs: ['x'.repeat(400)] } @@ -28,8 +29,16 @@ describe('projectToolResultForBrain', () => { const lean = projectToolResultForBrain('query_inbox', full) as { total: number; items: unknown[] } expect(lean.total).toBe(2) expect(lean.items).toEqual([ - { id: 7, what: 'CI auth blocking 3 workers', status: 'surfaced', priority: 90 }, - { id: 8, what: 'needs a decision', status: 'new', priority: 40 } + { + id: 7, + what: 'CI auth blocking 3 workers', + summary: 'long summary…', + category: 'blocker', + session: '96f67085-5dd3-4a10-aa7c-785f72a227c2', + status: 'surfaced', + priority: 90 + }, + { id: 8, what: 'needs a decision', summary: undefined, category: undefined, session: undefined, status: 'new', priority: 40 } ]) // the fat is gone expect(JSON.stringify(lean)).not.toContain('artifactRefs') diff --git a/hub/src/overseer/toolProjection.ts b/hub/src/overseer/toolProjection.ts index 0e491201aa..60cc9298de 100644 --- a/hub/src/overseer/toolProjection.ts +++ b/hub/src/overseer/toolProjection.ts @@ -37,6 +37,8 @@ function truncate(text: unknown, max: number): unknown { * Inbox item → the minimum for triage: * - `id` — to reference it (explain_priority, follow-ups) * - `what` — the title (the "what") + * - `summary` — action line (often better than a bare URL title) + * - `session` — full relatedSessionId (so get_session_state can be called without truncating) * - `status` — new / surfaced / held (has the operator seen it?) * - `priority` — explicit rank (also implied by order, but explicit lets the * brain speak with confidence) @@ -44,7 +46,15 @@ function truncate(text: unknown, max: number): unknown { */ function projectInboxItem(item: unknown): Record { const o = isObj(item) ? item : {} - return { id: o.id, what: o.title, status: o.status, priority: o.priority } + return { + id: o.id, + what: o.title, + summary: truncate(o.summary, 160), + category: o.category, + session: o.relatedSessionId, + status: o.status, + priority: o.priority + } } /** Disposition list row → the predicate keys + as-seen title (the rest is one query away). */ diff --git a/hub/src/sync/overseerEntity.test.ts b/hub/src/sync/overseerEntity.test.ts index ddb8b53fb1..401a6545aa 100644 --- a/hub/src/sync/overseerEntity.test.ts +++ b/hub/src/sync/overseerEntity.test.ts @@ -126,6 +126,58 @@ describe('OverseerEntity read-only tools', () => { expect(overseer(engine).getWorkerHealth('does-not-exist')).toBeNull() }) + it('session-scoped tools resolve unique id prefixes (inactive still returns state)', () => { + const store = new Store(':memory:') + const fullId = '96f67085-5dd3-4a10-aa7c-785f72a227c2' + const collisionId = '96f67085-aaaa-bbbb-cccc-ddddeeeeffff' + store.sessions.getOrCreateSession( + 'prefix-a', + { flavor: 'cursor', host: 'local', path: '/tmp/inline-model-error-detect', name: 'cursor inline model-error detect' }, + null, + 'default', + undefined, + undefined, + undefined, + fullId + ) + store.events.insert({ + ts: Date.now(), sourceKind: 'worker', eventType: 'completed', attentionCandidate: 0, severity: 2, + summary: 'done', relatedSessionId: fullId, + payloadJson: JSON.stringify({ session: { project: 'inline-model-error-detect', name: 'cursor inline model-error detect' } }) + }) + store.messages.addMessage(fullId, agentMessage('shipped the bridge')) + + // Unique 8-char prefix resolves (the live bug: brain truncates, tool used to return null). + const oUnique = overseer(buildEngine(store)) + const byPrefix = oUnique.getSessionState('96f67085') + expect(byPrefix).not.toBeNull() + expect(byPrefix!.sessionId).toBe(fullId) + expect(byPrefix!.name).toBe('cursor inline model-error detect') + expect(byPrefix!.workerReportedState).toBe('complete') + expect(oUnique.getWorkerHealth('96f67085')!.sessionId).toBe(fullId) + expect(oUnique.getSessionRecentOutput('96f67085').some((c) => c.text.includes('shipped'))).toBe(true) + expect(oUnique.queryEvents({ sessionId: '96f67085' }).length).toBe(1) + + // Ambiguous prefix must NOT silently pick a winner (rebuild engine so cache sees both). + store.sessions.getOrCreateSession( + 'prefix-b', + { flavor: 'codex', path: '/tmp/other', name: 'collision' }, + null, + 'default', + undefined, + undefined, + undefined, + collisionId + ) + const oAmbiguous = overseer(buildEngine(store)) + expect(oAmbiguous.getSessionState('96f67085')).toBeNull() + expect(oAmbiguous.getWorkerHealth('96f67085')).toBeNull() + expect(oAmbiguous.getSessionRecentOutput('96f67085')).toEqual([]) + expect(oAmbiguous.queryEvents({ sessionId: '96f67085' })).toEqual([]) + // Longer unique prefix still works after the collision appears. + expect(oAmbiguous.getSessionState('96f67085-5dd3')!.sessionId).toBe(fullId) + }) + it('get_session_recent_output returns last transcript chunks with roles', () => { const store = new Store(':memory:') const session = store.sessions.getOrCreateSession('out', { flavor: 'codex', path: '/tmp/web' }, null, 'default') diff --git a/hub/src/sync/overseerEntity.ts b/hub/src/sync/overseerEntity.ts index 9184409aad..4809da58d2 100644 --- a/hub/src/sync/overseerEntity.ts +++ b/hub/src/sync/overseerEntity.ts @@ -124,8 +124,17 @@ export class OverseerEntity { // --- Tool 1: query_events ------------------------------------------------ queryEvents(args: QueryEventsArgs): StoredSystemEvent[] { + // Resolve unique prefixes before querying — brain (and operators) often + // pass the short form. Ambiguous / unknown prefix → empty result (same + // as "no events for that session"), never a partial match. + let sessionId = args.sessionId ?? null + if (sessionId) { + const resolved = this.resolveSession(sessionId) + if (!resolved) return [] + sessionId = resolved.id + } return this.events.query({ - sessionId: args.sessionId ?? null, + sessionId, project: args.project ?? null, eventType: args.eventType ?? null, sourceKind: args.sourceKind ?? null, @@ -173,12 +182,13 @@ export class OverseerEntity { // --- Tool 3: get_session_state ------------------------------------------ getSessionState(sessionId: string): OverseerSessionStateView | null { - const session = this.getSession(sessionId) + const session = this.resolveSession(sessionId) if (!session) return null + const resolvedId = session.id const now = this.now() const { name, project, flavor } = deriveIdentity(session) - const latestEvent = this.events.query({ sessionId, limit: 1 })[0] ?? null + const latestEvent = this.events.query({ sessionId: resolvedId, limit: 1 })[0] ?? null const lastActivityAt = this.computeLastActivityAt(session, latestEvent) const silenceMs = lastActivityAt !== null ? Math.max(0, now - lastActivityAt) : null const pending = pendingRequestCount(session) @@ -191,19 +201,19 @@ export class OverseerEntity { staleSilenceMs: this.staleSilenceMs }) - const lastToolCall = this.events.query({ sessionId, eventType: 'tool_call', limit: 1 })[0] - ?? this.events.query({ sessionId, eventType: 'tool_result', limit: 1 })[0] + const lastToolCall = this.events.query({ sessionId: resolvedId, eventType: 'tool_call', limit: 1 })[0] + ?? this.events.query({ sessionId: resolvedId, eventType: 'tool_result', limit: 1 })[0] ?? null return { - sessionId, + sessionId: resolvedId, name, project, flavor, active: session.active, thinking: session.thinking, observedState, - workerReportedState: this.deriveReportedState(sessionId), + workerReportedState: this.deriveReportedState(resolvedId), lastActivityAt, silenceMs, lastToolCallAgeMs: lastToolCall ? Math.max(0, now - lastToolCall.ts) : null, @@ -214,8 +224,10 @@ export class OverseerEntity { // --- Tool 4: get_session_recent_output ---------------------------------- getSessionRecentOutput(sessionId: string, n = 10): OverseerRecentOutputChunk[] { + const session = this.resolveSession(sessionId) + if (!session) return [] const limit = Math.min(Math.max(n, 1), 50) - const messages = this.messages.getMessages(sessionId, limit) + const messages = this.messages.getMessages(session.id, limit) const chunks: OverseerRecentOutputChunk[] = [] for (const message of messages) { const text = this.extractMessageText(message.content) @@ -233,12 +245,13 @@ export class OverseerEntity { // --- Tool 5: get_worker_health ------------------------------------------ getWorkerHealth(sessionId: string): OverseerWorkerHealth | null { - const session = this.getSession(sessionId) + const session = this.resolveSession(sessionId) if (!session) return null + const resolvedId = session.id const now = this.now() const { name, project, flavor } = deriveIdentity(session) - const latestEvent = this.events.query({ sessionId, limit: 1 })[0] ?? null + const latestEvent = this.events.query({ sessionId: resolvedId, limit: 1 })[0] ?? null const lastActivityAt = this.computeLastActivityAt(session, latestEvent) const silenceMs = lastActivityAt !== null ? Math.max(0, now - lastActivityAt) : null const pending = pendingRequestCount(session) @@ -250,7 +263,7 @@ export class OverseerEntity { pendingRequestCount: pending, staleSilenceMs: this.staleSilenceMs }) - const reportedState = this.deriveReportedState(sessionId) + const reportedState = this.deriveReportedState(resolvedId) const inferred = inferWorkerState({ reported: reportedState, observed: observedState, @@ -273,7 +286,7 @@ export class OverseerEntity { signals.push(inferred.note) return { - sessionId, + sessionId: resolvedId, name, project, flavor, @@ -306,6 +319,10 @@ export class OverseerEntity { sourceKind: event.sourceKind })) + const relatedSessionId = item.relatedSessionId + const related = relatedSessionId ? this.resolveSession(relatedSessionId) : null + const identity = related ? deriveIdentity(related) : null + return { inboxItemId: item.id, title: item.title, @@ -318,7 +335,10 @@ export class OverseerEntity { // Recite the stored provenance — do NOT recompute (substrate authored it). reasonForPriority: item.reasonForPriority, sourceEventIds: item.sourceEventIds, - relatedSessionId: item.relatedSessionId, + relatedSessionId, + relatedSessionName: identity?.name ?? null, + relatedSessionActive: related ? related.active : null, + summary: item.summary, sourceEvents } } @@ -633,35 +653,30 @@ export class OverseerEntity { return { ok: false, error: 'sessionId or itemId is required — nothing relayed.' } } - // Exact match first, then unique prefix (same as hapi-ping-peer). - let session = this.getSession(sessionId) - if (!session) { - const matches = this.getSessions().filter((s) => s.id.startsWith(sessionId!)) - if (matches.length === 1) { - session = matches[0] - sessionId = session.id - } else if (matches.length > 1) { - return { - ok: false, - sessionId, - error: `Ambiguous session prefix "${sessionId}" (${matches.length} matches) — nothing relayed.` - } - } else { - return { - ok: false, - sessionId, - error: `No session matching "${sessionId}" — nothing relayed.` - } + // Exact match first, then unique prefix (same as hapi-ping-peer / resolveSession). + const matches = this.matchSessions(sessionId) + if (matches.length === 1) { + const session = matches[0]! + const identity = deriveIdentity(session) + return { + ok: true, + sessionId: session.id, + sessionName: identity.name, + project: identity.project, + namespace: session.namespace || 'default' + } + } + if (matches.length > 1) { + return { + ok: false, + sessionId, + error: `Ambiguous session prefix "${sessionId}" (${matches.length} matches) — nothing relayed.` } } - - const identity = deriveIdentity(session) return { - ok: true, - sessionId: session.id, - sessionName: identity.name, - project: identity.project, - namespace: session.namespace || 'default' + ok: false, + sessionId, + error: `No session matching "${sessionId}" — nothing relayed.` } } @@ -695,6 +710,24 @@ export class OverseerEntity { // --- internals ----------------------------------------------------------- + /** + * Exact session id, else unique prefix (hapi-ping-peer / loomux / pi pattern). + * Ambiguous or unknown → undefined. Never silently picks among collisions. + */ + private resolveSession(sessionId: string): Session | undefined { + const matches = this.matchSessions(sessionId) + return matches.length === 1 ? matches[0] : undefined + } + + /** Exact hit as a singleton, else all prefix matches (may be 0/1/many). */ + private matchSessions(sessionId: string): Session[] { + const trimmed = sessionId.trim() + if (!trimmed) return [] + const exact = this.getSession(trimmed) + if (exact) return [exact] + return this.getSessions().filter((s) => s.id.startsWith(trimmed)) + } + private parseEventPayload(payloadJson: string | null): { notify_summary?: unknown suggested_action?: unknown diff --git a/shared/src/overseerEntity.ts b/shared/src/overseerEntity.ts index ca514b66f9..997ef48a89 100644 --- a/shared/src/overseerEntity.ts +++ b/shared/src/overseerEntity.ts @@ -118,6 +118,12 @@ export type OverseerExplainPriority = { reasonForPriority: string | null sourceEventIds: number[] relatedSessionId: string | null + /** Session display name when the related session still exists in the hub. */ + relatedSessionName: string | null + /** Hub `active` flag for the related session; null when the session row is gone. */ + relatedSessionActive: boolean | null + /** Inbox item summary (action line) — often more useful than a bare URL title. */ + summary: string /** Lightweight detail for each contributing event (provenance trail). */ sourceEvents: Array<{ id: number @@ -539,17 +545,17 @@ export const OVERSEER_TOOL_CATALOG: OverseerToolCatalogEntry[] = [ }, { name: 'get_session_state', - description: 'Hub-observed state for one session: activity, tool-call recency, pending requests, and the worker-reported state when available.', + description: 'Hub-observed state for one session (full id or unique short prefix): activity, tool-call recency, pending requests, and the worker-reported state when available. Inactive sessions still return a state object (active:false) — null means no matching session, not "deleted".', readonly: true }, { name: 'get_session_recent_output', - description: 'Last N transcript chunks for one session, for context.', + description: 'Last N transcript chunks for one session (full id or unique short prefix).', readonly: true }, { name: 'get_worker_health', - description: 'Combined worker health: reported + observed + inferred state with a signal trail (never collapses a reported/observed conflict).', + description: 'Combined worker health for one session (full id or unique short prefix): reported + observed + inferred state with a signal trail (never collapses a reported/observed conflict).', readonly: true }, { @@ -644,9 +650,11 @@ export function buildOverseerSystemPrompt(): string { 'Read-only tools:', '- query_events — the events stream (blockers, completions, decisions, progress, errors).', '- query_inbox — what currently needs the operator: candidates, surfaced items, held items.', - '- get_session_state — one session\'s observed state, activity, and reported state.', - '- get_session_recent_output — the last few transcript chunks of a session.', - '- get_worker_health — reported + observed + inferred state for one worker.', + '- get_session_state — one session\'s observed state, activity, and reported state', + ' (sessionId: full UUID or unique short prefix; inactive sessions still return a state', + ' object with active:false — null means no match, NOT that the session was deleted).', + '- get_session_recent_output — the last few transcript chunks of a session (same id rules).', + '- get_worker_health — reported + observed + inferred state for one worker (same id rules).', '- explain_priority — why an inbox item sits where it does, with its provenance.', '- list_active_workers — the current roster, filterable by project / state / age.', '- query_open_loops — the "what am I forgetting?" lens: cold threads whose latest status is not done.', @@ -692,8 +700,15 @@ export function buildOverseerSystemPrompt(): string { '- Prioritize. Surface the root cause, not five symptoms ("GitHub auth is blocking 5 workers",', ' not a roll-call of each blocked worker).', '- When the operator asks about a SPECIFIC inbox item, first call explain_priority for its', - ' provenance, then query_events with that item\'s sessionId to pull the rest of that session\'s', - ' recorded activity as context/salience before answering — do not answer from the item alone.', + ' provenance, then query_events with that item\'s relatedSessionId (full UUID from the tool', + ' result — do not truncate it yourself) to pull the rest of that session\'s recorded activity', + ' as context/salience before answering — do not answer from the item alone.', + '- Prefer relatedSessionId / session fields from tool results over inventing short prefixes.', + ' Unique short prefixes are accepted by session tools, but truncating a UUID is how you get', + ' false "session missing" answers.', + '- An inbox item can still need a decision after its worker goes idle/complete — that is not', + ' a ghost. Check get_session_state (expect active:false + reported/observed) before claiming', + ' the session is gone.', '', '# Two questions, two axes', '', From 05cb8f6d376e11accc1383e0553df0aca182e9f7 Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sat, 1 Aug 2026 20:52:37 +0100 Subject: [PATCH 3/6] fix(overseer): deliver ping to resumed session id + honest toolTrace Resume can remap the hub session id after spawn+merge; relay now sends to that id. Failed relays mark toolTrace.ok false, and a brain failure after a successful write returns the write audit trail instead of an empty failure. Co-authored-by: Cursor --- hub/src/overseer/converse.test.ts | 77 ++++++++++++++++++++++++++++ hub/src/overseer/converse.ts | 84 ++++++++++++++++++++++++++++--- hub/src/sync/overseerEntity.ts | 12 +++-- hub/src/sync/sessionRelay.test.ts | 66 ++++++++++++++++++++++++ hub/src/sync/sessionRelay.ts | 63 +++++++++++++++++++++++ hub/src/sync/syncEngine.ts | 37 ++++---------- shared/src/overseerConverse.ts | 20 ++++++-- 7 files changed, 315 insertions(+), 44 deletions(-) create mode 100644 hub/src/sync/sessionRelay.test.ts create mode 100644 hub/src/sync/sessionRelay.ts diff --git a/hub/src/overseer/converse.test.ts b/hub/src/overseer/converse.test.ts index e476bd98cf..6a606bd7bc 100644 --- a/hub/src/overseer/converse.test.ts +++ b/hub/src/overseer/converse.test.ts @@ -175,4 +175,81 @@ describe('runOverseerConverse', () => { expect((e as BrainUnavailableError).reachable).toBe(true) } }) + + it('marks ping_session ok:false in the tool trace when relay fails', async () => { + const overseer = { + ...fakeOverseer, + pingSession: async () => ({ + ok: false, + sessionId: 'sess-1', + sessionName: null, + project: null, + resumed: false, + tombstone: 'Failed to relay: session_not_found', + error: 'session_not_found' + }) + } as unknown as OverseerEntity + const fetchMock = vi.fn() + .mockResolvedValueOnce(chatResponse({ + role: 'assistant', + content: '', + tool_calls: [{ + id: 'c1', + type: 'function', + function: { name: 'ping_session', arguments: '{"sessionId":"sess-1","message":"hi"}' } + }] + })) + .mockResolvedValueOnce(chatResponse({ role: 'assistant', content: 'Could not reach that session.' })) + setFetch(fetchMock) + + const { toolTrace } = await runOverseerConverse({ + overseer, + config, + messages: [{ role: 'operator', content: 'ping sess-1: hi' }] + }) + + expect(toolTrace[0]).toMatchObject({ + tool: 'ping_session', + ok: false, + error: 'session_not_found' + }) + }) + + it('keeps successful relay in the tool trace when the follow-up brain call fails', async () => { + const overseer = { + ...fakeOverseer, + pingSession: async () => ({ + ok: true, + sessionId: 'new-id', + sessionName: 'Worker', + project: 'hapi', + resumed: true, + tombstone: 'Relayed to Worker (new-id00) [resumed]: "please continue"' + }) + } as unknown as OverseerEntity + const fetchMock = vi.fn() + .mockResolvedValueOnce(chatResponse({ + role: 'assistant', + content: '', + tool_calls: [{ + id: 'c1', + type: 'function', + function: { name: 'ping_session', arguments: '{"sessionId":"old-id","message":"please continue"}' } + }] + })) + .mockRejectedValueOnce(new Error('ECONNREFUSED')) + setFetch(fetchMock) + + const { reply, toolTrace } = await runOverseerConverse({ + overseer, + config, + messages: [{ role: 'operator', content: 'tell that session to continue' }] + }) + + expect(toolTrace).toHaveLength(1) + expect(toolTrace[0]).toMatchObject({ tool: 'ping_session', ok: true }) + expect(reply).toContain('already succeeded') + expect(reply).toContain('Relayed to Worker') + expect(reply).toContain('Do not retry') + }) }) diff --git a/hub/src/overseer/converse.ts b/hub/src/overseer/converse.ts index 053e6453c4..6c37bfc98e 100644 --- a/hub/src/overseer/converse.ts +++ b/hub/src/overseer/converse.ts @@ -10,7 +10,9 @@ import { buildOverseerOpenAiTools, buildOverseerSystemPrompt, + isOverseerWriteTool, type OverseerConverseMessage, + type OverseerToolName, type OverseerToolTraceEntry } from '@hapi/protocol' import type { OverseerEntity } from '../sync/overseerEntity' @@ -62,6 +64,43 @@ function parseToolArgs(raw: string): Record { } } +/** Write tools return `{ ok: boolean, … }`; propagate that into the converse audit trail. */ +export function toolResultOk(result: unknown): boolean { + if (result && typeof result === 'object' && 'ok' in result) { + return (result as { ok: unknown }).ok !== false + } + return true +} + +function toolResultError(result: unknown): string | undefined { + if (!result || typeof result !== 'object') return undefined + const record = result as { error?: unknown; tombstone?: unknown } + if (typeof record.error === 'string' && record.error.trim()) return record.error + if (typeof record.tombstone === 'string' && record.tombstone.trim()) return record.tombstone + return 'tool returned ok:false' +} + +function writeResultTombstone(result: unknown): string | null { + if (!result || typeof result !== 'object') return null + const tombstone = (result as { tombstone?: unknown }).tombstone + return typeof tombstone === 'string' && tombstone.trim() ? tombstone.trim() : null +} + +function fallbackReplyAfterWriteSuccess(confirmations: string[]): string { + if (confirmations.length === 0) { + return 'A follow-up brain call failed after tools ran. Check the tool trace before retrying.' + } + return [ + 'Write tool(s) already succeeded; the brain failed while composing the confirmation.', + 'Do not retry the same write unless you intend a duplicate.', + ...confirmations.map((line) => `- ${line}`) + ].join('\n') +} + +function hasSuccessfulWrite(toolTrace: OverseerToolTraceEntry[]): boolean { + return toolTrace.some((entry) => entry.ok && isOverseerWriteTool(entry.tool as OverseerToolName)) +} + export async function runOverseerConverse(params: { overseer: OverseerEntity config: BrainConfig @@ -81,6 +120,8 @@ export async function runOverseerConverse(params: { ] const toolTrace: OverseerToolTraceEntry[] = [] + /** Tombstones from successful write tools — used if a later brain call fails. */ + const writeConfirmations: string[] = [] // The brain (llama-server) does not honor tool_choice:'required', so it will // sometimes answer a fleet question from nothing (e.g. "the inbox is empty" // when it never called query_inbox). Guardrail: if the very first answer @@ -89,7 +130,17 @@ export async function runOverseerConverse(params: { let nudged = false for (let iter = 0; iter < maxIterations; iter++) { - const message = await callBrain({ config, messages: convo, tools, signal }) + let message: OpenAiChatMessage + try { + message = await callBrain({ config, messages: convo, tools, signal }) + } catch (error) { + // Irreversible writes already landed — return their audit trail so the + // route can record the turn and the operator does not duplicate-retry. + if (hasSuccessfulWrite(toolTrace)) { + return { reply: fallbackReplyAfterWriteSuccess(writeConfirmations), toolTrace } + } + throw error + } const calls = message.tool_calls ?? [] if (calls.length === 0) { @@ -125,7 +176,17 @@ export async function runOverseerConverse(params: { // The conversational surface is the operator-directed write-path, so dispositions // are allowed here (gated off on the raw HTTP tool-dispatch endpoint). const result = await runOverseerTool(overseer, name, args, true) - toolTrace.push({ tool: name, args, ok: true }) + const ok = toolResultOk(result) + toolTrace.push({ + tool: name, + args, + ok, + ...(ok ? {} : { error: toolResultError(result) }) + }) + if (ok && isOverseerWriteTool(name)) { + const tombstone = writeResultTombstone(result) + writeConfirmations.push(tombstone ?? `${name} succeeded`) + } // The brain opts into 'full' per call when it needs depth; default lean. const detail = args.detail === 'full' ? 'full' : 'lean' const projected = projectToolResultForBrain(name, result, detail) @@ -143,10 +204,17 @@ export async function runOverseerConverse(params: { } // Iteration cap hit while still calling tools — ask once more for a plain answer. - const finalMsg = await callBrain({ - config, - messages: [...convo, { role: 'user', content: 'Answer now in plain text, no more tools.' }], - signal - }) - return { reply: (finalMsg.content ?? '').trim() || 'I gathered the data but could not compose an answer.', toolTrace } + try { + const finalMsg = await callBrain({ + config, + messages: [...convo, { role: 'user', content: 'Answer now in plain text, no more tools.' }], + signal + }) + return { reply: (finalMsg.content ?? '').trim() || 'I gathered the data but could not compose an answer.', toolTrace } + } catch (error) { + if (hasSuccessfulWrite(toolTrace)) { + return { reply: fallbackReplyAfterWriteSuccess(writeConfirmations), toolTrace } + } + throw error + } } diff --git a/hub/src/sync/overseerEntity.ts b/hub/src/sync/overseerEntity.ts index 4809da58d2..09570dbf20 100644 --- a/hub/src/sync/overseerEntity.ts +++ b/hub/src/sync/overseerEntity.ts @@ -64,7 +64,7 @@ export type OverseerEntityDeps = { sessionId: string message: string namespace?: string - }) => Promise<{ ok: boolean; resumed: boolean; error?: string }> + }) => Promise<{ ok: boolean; resumed: boolean; sessionId?: string; error?: string }> now?: () => number staleSilenceMs?: number } @@ -602,12 +602,14 @@ export class OverseerEntity { namespace: resolved.namespace }) - const label = resolved.sessionName ?? resolved.project ?? resolved.sessionId.slice(0, 8) + // After resume+merge the hub may have remapped to a replacement session id. + const deliveredSessionId = result.sessionId ?? resolved.sessionId + const label = resolved.sessionName ?? resolved.project ?? deliveredSessionId.slice(0, 8) const snippet = args.message.length > 80 ? `${args.message.slice(0, 77)}…` : args.message if (!result.ok) { return { ok: false, - sessionId: resolved.sessionId, + sessionId: deliveredSessionId, sessionName: resolved.sessionName, project: resolved.project, resumed: result.resumed, @@ -618,11 +620,11 @@ export class OverseerEntity { return { ok: true, - sessionId: resolved.sessionId, + sessionId: deliveredSessionId, sessionName: resolved.sessionName, project: resolved.project, resumed: result.resumed, - tombstone: `Relayed to ${label} (${resolved.sessionId.slice(0, 8)})${result.resumed ? ' [resumed]' : ''}: "${snippet}"` + tombstone: `Relayed to ${label} (${deliveredSessionId.slice(0, 8)})${result.resumed ? ' [resumed]' : ''}: "${snippet}"` } } diff --git a/hub/src/sync/sessionRelay.test.ts b/hub/src/sync/sessionRelay.test.ts new file mode 100644 index 0000000000..cd5523d489 --- /dev/null +++ b/hub/src/sync/sessionRelay.test.ts @@ -0,0 +1,66 @@ +import { describe, expect, it, vi } from 'vitest' +import { executeSessionRelay } from './sessionRelay' + +describe('executeSessionRelay', () => { + it('sends to the requested id when the session is already active', async () => { + const sendMessage = vi.fn().mockResolvedValue(undefined) + const result = await executeSessionRelay( + { + getSession: () => ({ active: true }), + resumeSession: vi.fn(), + sendMessage + }, + { sessionId: 'old-id', message: 'hello', namespace: 'default' } + ) + expect(result).toEqual({ ok: true, resumed: false, sessionId: 'old-id' }) + expect(sendMessage).toHaveBeenCalledWith('old-id', { text: 'hello', sentFrom: 'webapp' }) + }) + + it('sends to the resumed session id when resume remaps', async () => { + const sendMessage = vi.fn().mockResolvedValue(undefined) + const resumeSession = vi.fn().mockResolvedValue({ type: 'success', sessionId: 'new-id' }) + const result = await executeSessionRelay( + { + getSession: () => ({ active: false }), + resumeSession, + sendMessage + }, + { sessionId: 'old-id', message: 'wake up', namespace: 'ns-a' } + ) + expect(resumeSession).toHaveBeenCalledWith('old-id', 'ns-a') + expect(result).toEqual({ ok: true, resumed: true, sessionId: 'new-id' }) + expect(sendMessage).toHaveBeenCalledWith('new-id', { text: 'wake up', sentFrom: 'webapp' }) + expect(sendMessage).not.toHaveBeenCalledWith('old-id', expect.anything()) + }) + + it('returns resume failure without sending', async () => { + const sendMessage = vi.fn() + const result = await executeSessionRelay( + { + getSession: () => ({ active: false }), + resumeSession: async () => ({ type: 'error', code: 'no_machine_online', message: 'No machine online' }), + sendMessage + }, + { sessionId: 'old-id', message: 'x' } + ) + expect(result).toEqual({ + ok: false, + resumed: false, + sessionId: 'old-id', + error: 'no_machine_online' + }) + expect(sendMessage).not.toHaveBeenCalled() + }) + + it('returns session_not_found when missing', async () => { + const result = await executeSessionRelay( + { + getSession: () => undefined, + resumeSession: vi.fn(), + sendMessage: vi.fn() + }, + { sessionId: 'ghost', message: 'x' } + ) + expect(result).toMatchObject({ ok: false, error: 'session_not_found', sessionId: 'ghost' }) + }) +}) diff --git a/hub/src/sync/sessionRelay.ts b/hub/src/sync/sessionRelay.ts new file mode 100644 index 0000000000..9f6e439b51 --- /dev/null +++ b/hub/src/sync/sessionRelay.ts @@ -0,0 +1,63 @@ +/** + * In-process session relay used by Overseer `ping_session`. + * + * Resume may mint a replacement hub session ID (spawn + merge). Delivery must + * target that returned ID — sending to the pre-resume id hits a deleted row. + */ + +export type SessionRelayResumeResult = + | { type: 'success'; sessionId: string } + | { type: 'error'; message?: string; code?: string } + +export type SessionRelayResult = { + ok: boolean + resumed: boolean + /** Session id that received (or would receive) the message after remaps. */ + sessionId: string + error?: string +} + +export type SessionRelayDeps = { + getSession: (sessionId: string) => { active: boolean } | undefined + resumeSession: (sessionId: string, namespace: string) => Promise + sendMessage: (sessionId: string, payload: { text: string; sentFrom: 'webapp' }) => Promise +} + +export async function executeSessionRelay( + deps: SessionRelayDeps, + args: { sessionId: string; message: string; namespace?: string } +): Promise { + const namespace = args.namespace ?? 'default' + const existing = deps.getSession(args.sessionId) + if (!existing) { + return { ok: false, resumed: false, sessionId: args.sessionId, error: 'session_not_found' } + } + + let deliverySessionId = args.sessionId + let resumed = false + if (!existing.active) { + const resume = await deps.resumeSession(args.sessionId, namespace) + if (resume.type !== 'success') { + return { + ok: false, + resumed: false, + sessionId: args.sessionId, + error: resume.code ?? resume.message ?? 'resume_failed' + } + } + deliverySessionId = resume.sessionId + resumed = true + } + + try { + await deps.sendMessage(deliverySessionId, { text: args.message, sentFrom: 'webapp' }) + return { ok: true, resumed, sessionId: deliverySessionId } + } catch (error) { + return { + ok: false, + resumed, + sessionId: deliverySessionId, + error: error instanceof Error ? error.message : 'send_failed' + } + } +} diff --git a/hub/src/sync/syncEngine.ts b/hub/src/sync/syncEngine.ts index a4989a9f40..69aea6a1f9 100644 --- a/hub/src/sync/syncEngine.ts +++ b/hub/src/sync/syncEngine.ts @@ -44,6 +44,7 @@ import { import { SessionCache } from './sessionCache' import { OverseerEventRecorder, toSessionSnapshot } from './overseerEventRecorder' import { OverseerEntity } from './overseerEntity' +import { executeSessionRelay } from './sessionRelay' import { extractAssistantPlainText } from '@hapi/protocol/messages' import type { InboxOperatorAction } from '@hapi/protocol' import type { ListSystemEventsOptions, StoredSystemEvent } from '../store' @@ -174,32 +175,16 @@ export class SyncEngine { getSession: (sessionId) => this.getSession(sessionId), getSessions: () => this.getSessions(), relayToSession: async ({ sessionId, message, namespace = 'default' }) => { - const existing = this.getSession(sessionId) - if (!existing) { - return { ok: false, resumed: false, error: 'session_not_found' } - } - let resumed = false - if (!existing.active) { - const resume = await this.resumeSession(sessionId, namespace) - if (resume.type !== 'success') { - return { - ok: false, - resumed: false, - error: resume.code ?? resume.message ?? 'resume_failed' - } - } - resumed = true - } - try { - await this.sendMessage(sessionId, { text: message, sentFrom: 'webapp' }) - return { ok: true, resumed } - } catch (error) { - return { - ok: false, - resumed, - error: error instanceof Error ? error.message : 'send_failed' - } - } + // Resume may return a replacement hub session id after spawn+merge; + // always deliver to that id (see executeSessionRelay). + return executeSessionRelay( + { + getSession: (id) => this.getSession(id), + resumeSession: (id, ns) => this.resumeSession(id, ns), + sendMessage: (id, payload) => this.sendMessage(id, payload) + }, + { sessionId, message, namespace } + ) } }) this.reloadAll() diff --git a/shared/src/overseerConverse.ts b/shared/src/overseerConverse.ts index 34e9fa5018..6d6a6321ea 100644 --- a/shared/src/overseerConverse.ts +++ b/shared/src/overseerConverse.ts @@ -176,11 +176,21 @@ const OVERSEER_TOOL_PARAMS: Record = { feedback: { type: 'string', description: 'Optional operator note / learning label to freeze with the disposition.' }, snoozedUntil: { type: 'integer', minimum: 1, description: 'Required for snooze: epoch ms to sleep the item until.' } }, ['itemId', 'action']), - ping_session: obj({ - sessionId: { type: 'string', description: 'Worker session id (full UUID or unique prefix).' }, - itemId: { type: 'integer', minimum: 1, description: 'Inbox item id — resolves its relatedSessionId when sessionId omitted.' }, - message: { type: 'string', description: 'Operator-directed message to relay to that session.' } - }, ['message']) + // anyOf mirrors runtime Zod: message plus at least one of sessionId|itemId. + ping_session: { + type: 'object', + properties: { + sessionId: { type: 'string', description: 'Worker session id (full UUID or unique prefix).' }, + itemId: { type: 'integer', minimum: 1, description: 'Inbox item id — resolves its relatedSessionId when sessionId omitted.' }, + message: { type: 'string', description: 'Operator-directed message to relay to that session.' } + }, + required: ['message'], + anyOf: [ + { required: ['sessionId'] }, + { required: ['itemId'] } + ], + additionalProperties: false + } } /** The Overseer tool catalog (reads + disposition + relay writes) as an OpenAI `tools` array. */ From c6da6c1d543acc7bd7788646bc1cfb1ab393fe5d Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sat, 1 Aug 2026 21:05:47 +0100 Subject: [PATCH 4/6] fix(overseer): write-intent gate, conflicting ping targets, dispatched audit Server-side authorization for write tools from the latest operator utterance (or explicit allowWrites). Reject sessionId/itemId mismatches. Record idempotent dispatched events after successful relays. Co-authored-by: Cursor --- hub/src/overseer/converse.test.ts | 35 ++++++++++++- hub/src/overseer/converse.ts | 28 ++++++++-- hub/src/sync/overseerEntity.test.ts | 28 ++++++++-- hub/src/sync/overseerEntity.ts | 72 ++++++++++++++++++++++++-- hub/src/web/routes/overseer.ts | 7 ++- shared/src/index.ts | 1 + shared/src/overseerEntity.ts | 64 +++++++++++++++++++++++ shared/src/overseerWriteIntent.test.ts | 49 ++++++++++++++++++ shared/src/overseerWriteIntent.ts | 58 +++++++++++++++++++++ 9 files changed, 326 insertions(+), 16 deletions(-) create mode 100644 shared/src/overseerWriteIntent.test.ts create mode 100644 shared/src/overseerWriteIntent.ts diff --git a/hub/src/overseer/converse.test.ts b/hub/src/overseer/converse.test.ts index 6a606bd7bc..1c12206533 100644 --- a/hub/src/overseer/converse.test.ts +++ b/hub/src/overseer/converse.test.ts @@ -243,7 +243,7 @@ describe('runOverseerConverse', () => { const { reply, toolTrace } = await runOverseerConverse({ overseer, config, - messages: [{ role: 'operator', content: 'tell that session to continue' }] + messages: [{ role: 'operator', content: 'ping session old-id: please continue' }] }) expect(toolTrace).toHaveLength(1) @@ -252,4 +252,37 @@ describe('runOverseerConverse', () => { expect(reply).toContain('Relayed to Worker') expect(reply).toContain('Do not retry') }) + + it('refuses ping_session when the operator message has no write intent', async () => { + const pingSession = vi.fn() + const overseer = { + ...fakeOverseer, + pingSession + } as unknown as OverseerEntity + const fetchMock = vi.fn() + .mockResolvedValueOnce(chatResponse({ + role: 'assistant', + content: '', + tool_calls: [{ + id: 'c1', + type: 'function', + function: { name: 'ping_session', arguments: '{"sessionId":"sess-1","message":"hi"}' } + }] + })) + .mockResolvedValueOnce(chatResponse({ + role: 'assistant', + content: 'I cannot relay without an explicit operator request.' + })) + setFetch(fetchMock) + + const { toolTrace } = await runOverseerConverse({ + overseer, + config, + messages: [{ role: 'operator', content: 'what needs my attention?' }] + }) + + expect(pingSession).not.toHaveBeenCalled() + expect(toolTrace[0]).toMatchObject({ tool: 'ping_session', ok: false }) + expect(toolTrace[0]?.error).toMatch(/not authorized/i) + }) }) diff --git a/hub/src/overseer/converse.ts b/hub/src/overseer/converse.ts index 6c37bfc98e..e44e01e6d1 100644 --- a/hub/src/overseer/converse.ts +++ b/hub/src/overseer/converse.ts @@ -11,9 +11,12 @@ import { buildOverseerOpenAiTools, buildOverseerSystemPrompt, isOverseerWriteTool, + isWriteToolAuthorized, + resolveOverseerWriteAuthorization, type OverseerConverseMessage, type OverseerToolName, - type OverseerToolTraceEntry + type OverseerToolTraceEntry, + type OverseerWriteAuthorization } from '@hapi/protocol' import type { OverseerEntity } from '../sync/overseerEntity' import { isOverseerToolName, runOverseerTool } from './runOverseerTool' @@ -107,10 +110,21 @@ export async function runOverseerConverse(params: { messages: OverseerConverseMessage[] maxIterations?: number signal?: AbortSignal + /** Explicit client opt-in for write tools (admin/voice confirm). */ + allowWrites?: boolean }): Promise<{ reply: string; toolTrace: OverseerToolTraceEntry[] }> { - const { overseer, config, messages, maxIterations = 6, signal } = params - - const tools = buildOverseerOpenAiTools() as OverseerOpenAiToolLike[] + const { overseer, config, messages, maxIterations = 6, signal, allowWrites } = params + + const latestOperatorText = [...messages].reverse().find((m) => m.role === 'operator')?.content ?? '' + const writeAuth: OverseerWriteAuthorization = resolveOverseerWriteAuthorization({ + latestOperatorText, + allowWrites + }) + + const tools = (buildOverseerOpenAiTools() as OverseerOpenAiToolLike[]).filter((tool) => { + const name = tool.function?.name ?? '' + return isWriteToolAuthorized(name, writeAuth) + }) const convo: OpenAiChatMessage[] = [ { role: 'system', content: `${buildOverseerSystemPrompt()}\n\n${GROUNDING_DIRECTIVE}` }, ...messages.map((m): OpenAiChatMessage => ({ @@ -172,6 +186,12 @@ export async function runOverseerConverse(params: { resultLines.push(`${name || 'unknown'}(${argsRaw}) => ${JSON.stringify({ error: `unknown tool: ${name}` })}`) continue } + if (!isWriteToolAuthorized(name, writeAuth)) { + const denied = 'write not authorized by operator message (no explicit write intent)' + toolTrace.push({ tool: name, args, ok: false, error: denied }) + resultLines.push(`${name}(${argsRaw}) => ${JSON.stringify({ error: denied })}`) + continue + } try { // The conversational surface is the operator-directed write-path, so dispositions // are allowed here (gated off on the raw HTTP tool-dispatch endpoint). diff --git a/hub/src/sync/overseerEntity.test.ts b/hub/src/sync/overseerEntity.test.ts index 401a6545aa..4a9b13e612 100644 --- a/hub/src/sync/overseerEntity.test.ts +++ b/hub/src/sync/overseerEntity.test.ts @@ -487,10 +487,10 @@ describe('OverseerEntity dispositions (Stage 1 keystone)', () => { if (!s) return undefined return { ...s, active: true, namespace: s.namespace || 'default' } as never }, - getSessions: () => { - const s = store.sessions.getSession(sessionId) - return s ? [{ ...s, active: true, namespace: s.namespace || 'default' } as never] : [] - }, + getSessions: () => + store.sessions + .getSessions() + .map((s) => ({ ...s, active: true, namespace: s.namespace || 'default' }) as never), relayToSession: async ({ sessionId: sid, message }) => { lastRelay = { sessionId: sid, message } return { ok: true, resumed: false } @@ -518,5 +518,25 @@ describe('OverseerEntity dispositions (Stage 1 keystone)', () => { await expect(runOverseerTool(o, 'ping_session', { sessionId, message: 'x' })).rejects.toBeInstanceOf( OverseerWriteNotAllowedError ) + + const dispatched = store.events.query({ sessionId, eventType: 'dispatched', limit: 10 }) + expect(dispatched.length).toBeGreaterThanOrEqual(1) + expect(dispatched[0]?.summary).toMatch(/Relayed/) + + // Conflicting sessionId + itemId must refuse (not silently prefer sessionId). + const other = store.sessions.getOrCreateSession( + 'other-conflict', + { flavor: 'claude', path: '/tmp/other', name: 'other' }, + null, + 'default' + ) + const conflict = await o.pingSession({ + sessionId: other.id, + itemId, + message: 'wrong target' + }) + expect(conflict.ok).toBe(false) + expect(conflict.error).toMatch(/Conflicting relay targets/) + expect(lastRelay?.message).not.toBe('wrong target') }) }) diff --git a/hub/src/sync/overseerEntity.ts b/hub/src/sync/overseerEntity.ts index 09570dbf20..3679eec268 100644 --- a/hub/src/sync/overseerEntity.ts +++ b/hub/src/sync/overseerEntity.ts @@ -13,6 +13,7 @@ import { OVERSEER_LOOP_CLOSED_EVENT_TYPE, OVERSEER_STALE_SILENCE_MS, buildOverseerConvoTurnEventInput, + buildOverseerDispatchedEventInput, buildOverseerIdentity, buildOverseerSystemPrompt, deriveObservedWorkerState, @@ -618,16 +619,38 @@ export class OverseerEntity { } } + const tombstone = `Relayed to ${label} (${deliveredSessionId.slice(0, 8)})${result.resumed ? ' [resumed]' : ''}: "${snippet}"` + this.recordDispatchedRelay({ + sessionId: deliveredSessionId, + message: args.message, + resumed: result.resumed, + tombstone + }) return { ok: true, sessionId: deliveredSessionId, sessionName: resolved.sessionName, project: resolved.project, resumed: result.resumed, - tombstone: `Relayed to ${label} (${deliveredSessionId.slice(0, 8)})${result.resumed ? ' [resumed]' : ''}: "${snippet}"` + tombstone } } + private recordDispatchedRelay(args: { + sessionId: string + message: string + resumed: boolean + tombstone: string + }): void { + this.events.insert(buildOverseerDispatchedEventInput({ + sessionId: args.sessionId, + message: args.message, + resumed: args.resumed, + tombstone: args.tombstone, + ts: this.now() + })) + } + private resolvePingTarget(args: PingSessionArgs): { ok: true sessionId: string @@ -635,15 +658,16 @@ export class OverseerEntity { project: string | null namespace: string } | { ok: false; sessionId?: string; error: string } { - let sessionId = args.sessionId?.trim() || null + const rawSessionId = args.sessionId?.trim() || null + let sessionIdFromItem: string | null = null - if (!sessionId && args.itemId != null) { + if (args.itemId != null) { const item = this.inbox.getById(args.itemId) if (!item) { return { ok: false, error: `No inbox item #${args.itemId} — nothing relayed.` } } - sessionId = item.relatedSessionId?.trim() || null - if (!sessionId) { + sessionIdFromItem = item.relatedSessionId?.trim() || null + if (!sessionIdFromItem) { return { ok: false, error: `Inbox item #${args.itemId} has no related session — nothing relayed.` @@ -651,6 +675,44 @@ export class OverseerEntity { } } + if (rawSessionId && sessionIdFromItem) { + const fromSession = this.matchSessions(rawSessionId) + const fromItem = this.matchSessions(sessionIdFromItem) + if (fromSession.length !== 1) { + return { + ok: false, + sessionId: rawSessionId, + error: fromSession.length > 1 + ? `Ambiguous session prefix "${rawSessionId}" (${fromSession.length} matches) — nothing relayed.` + : `No session matching "${rawSessionId}" — nothing relayed.` + } + } + if (fromItem.length !== 1) { + return { + ok: false, + sessionId: sessionIdFromItem, + error: `Inbox item #${args.itemId} related session unresolved — nothing relayed.` + } + } + if (fromSession[0]!.id !== fromItem[0]!.id) { + return { + ok: false, + sessionId: fromSession[0]!.id, + error: `Conflicting relay targets: sessionId resolves to ${fromSession[0]!.id.slice(0, 8)}… but item #${args.itemId} points at ${fromItem[0]!.id.slice(0, 8)}… — nothing relayed.` + } + } + const session = fromSession[0]! + const identity = deriveIdentity(session) + return { + ok: true, + sessionId: session.id, + sessionName: identity.name, + project: identity.project, + namespace: session.namespace || 'default' + } + } + + const sessionId = rawSessionId ?? sessionIdFromItem if (!sessionId) { return { ok: false, error: 'sessionId or itemId is required — nothing relayed.' } } diff --git a/hub/src/web/routes/overseer.ts b/hub/src/web/routes/overseer.ts index 0a4c23f143..57c7d1e76d 100644 --- a/hub/src/web/routes/overseer.ts +++ b/hub/src/web/routes/overseer.ts @@ -34,7 +34,9 @@ const converseBodySchema = z.object({ })).min(1).max(40), relatedSessionId: z.string().min(1).optional(), model: z.string().max(100).optional(), - profile: z.string().max(64).optional() + profile: z.string().max(64).optional(), + /** Explicit opt-in for write tools; otherwise server detects intent from the latest operator line. */ + allowWrites: z.boolean().optional() }) const activeBrainBodySchema = z.object({ @@ -226,7 +228,8 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho const { reply, toolTrace } = await runOverseerConverse({ overseer: engine.getOverseer(), config, - messages + messages, + allowWrites: parsed.data.allowWrites }) const lastOperator = [...messages].reverse().find((m) => m.role === 'operator')?.content ?? '' diff --git a/shared/src/index.ts b/shared/src/index.ts index 822e3a60c1..c5f60ab413 100644 --- a/shared/src/index.ts +++ b/shared/src/index.ts @@ -4,6 +4,7 @@ export * from './messages' export * from './overseerEvents' export * from './overseerInbox' export * from './overseerEntity' +export * from './overseerWriteIntent' export * from './overseerConverse' export * from './buildInfo' export * from './effort' diff --git a/shared/src/overseerEntity.ts b/shared/src/overseerEntity.ts index 997ef48a89..2ea45f0798 100644 --- a/shared/src/overseerEntity.ts +++ b/shared/src/overseerEntity.ts @@ -29,6 +29,9 @@ export const OVERSEER_ENTITY_ID = 'overseer' /** `source_kind` / `sink_kind` used for Overseer-authored events. */ export const OVERSEER_SOURCE_KIND = 'overseer' as const +/** Event type for an operator-directed outbound relay (audit / session timeline). */ +export const OVERSEER_DISPATCHED_EVENT_TYPE = 'dispatched' as const + /** Event type for an operator<->Overseer conversation segment (memory-bearing). */ export const OVERSEER_CONVO_TURN_EVENT_TYPE = 'convo_turn' as const @@ -796,3 +799,64 @@ export function buildOverseerConvoTurnEventInput(input: OverseerConvoTurnInput): provenance: 'overseer-convo' } } + +export type OverseerDispatchedEventInput = { + ts: number + sourceKind: typeof OVERSEER_SOURCE_KIND + sourceRef: string + sinkKind: 'worker' + eventType: typeof OVERSEER_DISPATCHED_EVENT_TYPE + attentionCandidate: 0 + operatorActionRequired: 0 + summary: string + payloadJson: string + relatedSessionId: string + relatedEventId: number | null + provenance: string + idempotencyKey: string +} + +export type OverseerDispatchedInput = { + sessionId: string + message: string + resumed: boolean + tombstone: string + ts?: number +} + +/** Build an idempotent `dispatched` event after a successful ping_session relay. */ +export function buildOverseerDispatchedEventInput(input: OverseerDispatchedInput): OverseerDispatchedEventInput { + const ts = input.ts ?? Date.now() + const snippet = input.message.length > 120 ? `${input.message.slice(0, 117)}…` : input.message + // Bucket to the second so identical retries within the same second collapse. + const idempotencyKey = `overseer-dispatched:${input.sessionId}:${ts - (ts % 1000)}:${hashRelaySnippet(input.message)}` + return { + ts, + sourceKind: OVERSEER_SOURCE_KIND, + sourceRef: OVERSEER_ENTITY_ID, + sinkKind: 'worker', + eventType: OVERSEER_DISPATCHED_EVENT_TYPE, + attentionCandidate: 0, + operatorActionRequired: 0, + summary: input.tombstone.slice(0, 240) || `Relayed to ${input.sessionId.slice(0, 8)}: ${snippet}`, + payloadJson: JSON.stringify({ + message: input.message, + resumed: input.resumed, + tombstone: input.tombstone + }), + relatedSessionId: input.sessionId, + relatedEventId: null, + provenance: 'overseer-relay', + idempotencyKey + } +} + +function hashRelaySnippet(message: string): string { + // Short stable fingerprint — not cryptographic; enough for idempotency bucketing. + let h = 2166136261 + for (let i = 0; i < message.length; i++) { + h ^= message.charCodeAt(i) + h = Math.imul(h, 16777619) + } + return (h >>> 0).toString(16) +} diff --git a/shared/src/overseerWriteIntent.test.ts b/shared/src/overseerWriteIntent.test.ts new file mode 100644 index 0000000000..788cfc748d --- /dev/null +++ b/shared/src/overseerWriteIntent.test.ts @@ -0,0 +1,49 @@ +import { describe, expect, it } from 'vitest' +import { + detectOperatorWriteTools, + isWriteToolAuthorized, + resolveOverseerWriteAuthorization +} from './overseerWriteIntent' + +describe('detectOperatorWriteTools', () => { + it('authorizes relay for ping/tell session phrasing', () => { + expect([...detectOperatorWriteTools('ping the expenses session: please continue')]).toEqual([ + 'ping_session' + ]) + expect([...detectOperatorWriteTools('tell that worker to retry the flaky test')]).toContain( + 'ping_session' + ) + }) + + it('authorizes disposition for snooze/done phrasing', () => { + expect([...detectOperatorWriteTools('snooze item 12 until tomorrow')]).toEqual([ + 'record_disposition' + ]) + expect([...detectOperatorWriteTools('mark #7 done')]).toContain('record_disposition') + }) + + it('does not authorize writes for read-only questions', () => { + expect([...detectOperatorWriteTools('what needs my attention?')]).toEqual([]) + expect([...detectOperatorWriteTools('show recent output for session abc')]).toEqual([]) + }) +}) + +describe('resolveOverseerWriteAuthorization', () => { + it('explicit allowWrites unlocks both write tools', () => { + const auth = resolveOverseerWriteAuthorization({ + latestOperatorText: 'what is in the inbox?', + allowWrites: true + }) + expect(auth.explicitClientFlag).toBe(true) + expect(isWriteToolAuthorized('ping_session', auth)).toBe(true) + expect(isWriteToolAuthorized('record_disposition', auth)).toBe(true) + }) + + it('denies write tools when neither flag nor intent matches', () => { + const auth = resolveOverseerWriteAuthorization({ + latestOperatorText: 'summarize the inbox' + }) + expect(isWriteToolAuthorized('ping_session', auth)).toBe(false) + expect(isWriteToolAuthorized('query_inbox', auth)).toBe(true) + }) +}) diff --git a/shared/src/overseerWriteIntent.ts b/shared/src/overseerWriteIntent.ts new file mode 100644 index 0000000000..c84bbdb8bf --- /dev/null +++ b/shared/src/overseerWriteIntent.ts @@ -0,0 +1,58 @@ +/** + * Server-side write authorization for Overseer converse. + * + * Write tools must not run merely because the model asked — untrusted tool + * results (inbox titles, worker output) are fed back as `user` messages and can + * prompt-inject a relay/disposition. Authorization comes from the operator's + * latest utterance and/or an explicit client `allowWrites` flag — never from + * model-selected tools alone. + */ + +import { isOverseerWriteTool, type OverseerWriteToolName } from './overseerEntity' + +export type OverseerWriteAuthorization = { + /** Tools the operator's message (or explicit flag) authorized for this turn. */ + allowed: ReadonlySet + /** True when the client sent allowWrites: true (admin console / voice confirm). */ + explicitClientFlag: boolean +} + +const RELAY_INTENT = + /\b(ping|relay|nudge|wake)\b|\btell\b[\s\S]{0,80}\b(session|worker|peer|agent|him|her|them|it)\b|\b(message|ask|send)\b[\s\S]{0,80}\b(session|worker|peer|agent)\b/i + +const DISPOSITION_INTENT = + /\b(snooze|dismiss|reopen|dispose)\b|\bmark\b[\s\S]{0,40}\bdone\b|\b(resolve|done with)\b/i + +/** Detect which write classes the latest operator message authorizes. */ +export function detectOperatorWriteTools(operatorText: string): Set { + const allowed = new Set() + const text = operatorText.trim() + if (!text) return allowed + if (RELAY_INTENT.test(text)) allowed.add('ping_session') + if (DISPOSITION_INTENT.test(text)) allowed.add('record_disposition') + return allowed +} + +export function resolveOverseerWriteAuthorization(opts: { + latestOperatorText: string + allowWrites?: boolean +}): OverseerWriteAuthorization { + if (opts.allowWrites === true) { + return { + allowed: new Set(['ping_session', 'record_disposition']), + explicitClientFlag: true + } + } + return { + allowed: detectOperatorWriteTools(opts.latestOperatorText), + explicitClientFlag: false + } +} + +export function isWriteToolAuthorized( + tool: string, + auth: OverseerWriteAuthorization +): boolean { + if (!isOverseerWriteTool(tool)) return true + return auth.allowed.has(tool) +} From b3465a9bd4b4f2e0c7eb8e018743ed46eef0018b Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sun, 2 Aug 2026 00:36:02 +0100 Subject: [PATCH 5/6] fix(overseer): snooze default hide, bind write grants, namespace scope Default query_inbox no longer treats sleeping snoozes as visible; write tools bind to operator-named session/item/payload with dedupe; relay refuses fresh-spawn and keeps ok when audit insert fails; OverseerEntity is per-namespace (#107 kill criterion) via SyncEngine.getOverseer(ns). Co-authored-by: Cursor --- hub/src/overseer/converse.test.ts | 4 +- hub/src/overseer/converse.ts | 35 ++++++- hub/src/store/settingsStore.ts | 19 ++-- hub/src/sync/overseerEntity.test.ts | 44 ++++++++ hub/src/sync/overseerEntity.ts | 84 ++++++++++----- hub/src/sync/syncEngine.ts | 67 +++++++----- hub/src/web/routes/overseer.ts | 25 +++-- shared/src/overseerWriteIntent.test.ts | 48 +++++++-- shared/src/overseerWriteIntent.ts | 137 ++++++++++++++++++++++++- 9 files changed, 378 insertions(+), 85 deletions(-) diff --git a/hub/src/overseer/converse.test.ts b/hub/src/overseer/converse.test.ts index 1c12206533..e00b6d5a3d 100644 --- a/hub/src/overseer/converse.test.ts +++ b/hub/src/overseer/converse.test.ts @@ -205,7 +205,7 @@ describe('runOverseerConverse', () => { const { toolTrace } = await runOverseerConverse({ overseer, config, - messages: [{ role: 'operator', content: 'ping sess-1: hi' }] + messages: [{ role: 'operator', content: 'ping session sess-1: "hi"' }] }) expect(toolTrace[0]).toMatchObject({ @@ -243,7 +243,7 @@ describe('runOverseerConverse', () => { const { reply, toolTrace } = await runOverseerConverse({ overseer, config, - messages: [{ role: 'operator', content: 'ping session old-id: please continue' }] + messages: [{ role: 'operator', content: 'ping session old-id: "please continue"' }] }) expect(toolTrace).toHaveLength(1) diff --git a/hub/src/overseer/converse.ts b/hub/src/overseer/converse.ts index e44e01e6d1..668044860c 100644 --- a/hub/src/overseer/converse.ts +++ b/hub/src/overseer/converse.ts @@ -10,8 +10,10 @@ import { buildOverseerOpenAiTools, buildOverseerSystemPrompt, + fingerprintWriteToolCall, isOverseerWriteTool, isWriteToolAuthorized, + isWriteToolCallAuthorized, resolveOverseerWriteAuthorization, type OverseerConverseMessage, type OverseerToolName, @@ -125,8 +127,9 @@ export async function runOverseerConverse(params: { const name = tool.function?.name ?? '' return isWriteToolAuthorized(name, writeAuth) }) + const clockLine = `Server time now: ${new Date().toISOString()} (epoch ms ${Date.now()}, timezone ${Intl.DateTimeFormat().resolvedOptions().timeZone}). Relative snoozes must use absolute snoozedUntil epoch ms from this clock.` const convo: OpenAiChatMessage[] = [ - { role: 'system', content: `${buildOverseerSystemPrompt()}\n\n${GROUNDING_DIRECTIVE}` }, + { role: 'system', content: `${buildOverseerSystemPrompt()}\n\n${GROUNDING_DIRECTIVE}\n\n# Clock\n\n${clockLine}` }, ...messages.map((m): OpenAiChatMessage => ({ role: m.role === 'operator' ? 'user' : 'assistant', content: m.content @@ -136,6 +139,8 @@ export async function runOverseerConverse(params: { const toolTrace: OverseerToolTraceEntry[] = [] /** Tombstones from successful write tools — used if a later brain call fails. */ const writeConfirmations: string[] = [] + /** Successful irreversible call fingerprints — reject duplicates in this turn. */ + const consumedWriteFingerprints = new Set() // The brain (llama-server) does not honor tool_choice:'required', so it will // sometimes answer a fleet question from nothing (e.g. "the inbox is empty" // when it never called query_inbox). Guardrail: if the very first answer @@ -177,6 +182,10 @@ export async function runOverseerConverse(params: { // user/assistant path that all templates render. We also drop the raw // assistant tool-call message from history for the same reason. const resultLines: string[] = [] + const batchHasRead = calls.some((call) => { + const name = call.function?.name ?? '' + return isOverseerToolName(name) && !isOverseerWriteTool(name) + }) for (const call of calls) { const name = call.function?.name ?? '' const argsRaw = call.function?.arguments ?? '' @@ -186,12 +195,27 @@ export async function runOverseerConverse(params: { resultLines.push(`${name || 'unknown'}(${argsRaw}) => ${JSON.stringify({ error: `unknown tool: ${name}` })}`) continue } - if (!isWriteToolAuthorized(name, writeAuth)) { - const denied = 'write not authorized by operator message (no explicit write intent)' - toolTrace.push({ tool: name, args, ok: false, error: denied }) - resultLines.push(`${name}(${argsRaw}) => ${JSON.stringify({ error: denied })}`) + if (batchHasRead && isOverseerWriteTool(name)) { + const deferred = 'write deferred: resolve identifying read tools first, then call the write in a later turn' + toolTrace.push({ tool: name, args, ok: false, error: deferred }) + resultLines.push(`${name}(${argsRaw}) => ${JSON.stringify({ error: deferred })}`) + continue + } + const authz = isWriteToolCallAuthorized(name, args, writeAuth) + if (!authz.ok) { + toolTrace.push({ tool: name, args, ok: false, error: authz.error }) + resultLines.push(`${name}(${argsRaw}) => ${JSON.stringify({ error: authz.error })}`) continue } + if (isOverseerWriteTool(name)) { + const fp = fingerprintWriteToolCall(name, args) + if (consumedWriteFingerprints.has(fp)) { + const dup = 'duplicate irreversible tool call rejected (already executed this turn)' + toolTrace.push({ tool: name, args, ok: false, error: dup }) + resultLines.push(`${name}(${argsRaw}) => ${JSON.stringify({ error: dup })}`) + continue + } + } try { // The conversational surface is the operator-directed write-path, so dispositions // are allowed here (gated off on the raw HTTP tool-dispatch endpoint). @@ -204,6 +228,7 @@ export async function runOverseerConverse(params: { ...(ok ? {} : { error: toolResultError(result) }) }) if (ok && isOverseerWriteTool(name)) { + consumedWriteFingerprints.add(fingerprintWriteToolCall(name, args)) const tombstone = writeResultTombstone(result) writeConfirmations.push(tombstone ?? `${name} succeeded`) } diff --git a/hub/src/store/settingsStore.ts b/hub/src/store/settingsStore.ts index 0da3f17cd1..00cba84914 100644 --- a/hub/src/store/settingsStore.ts +++ b/hub/src/store/settingsStore.ts @@ -24,6 +24,11 @@ export type ActiveBrainSetting = { const ACTIVE_BRAIN_KEY = 'active_brain' +function activeBrainKey(namespace: string): string { + const ns = namespace.trim() || 'default' + return ns === 'default' ? ACTIVE_BRAIN_KEY : `${ACTIVE_BRAIN_KEY}:${ns}` +} + export class SettingsStore { constructor(private readonly db: Database) {} @@ -47,9 +52,9 @@ export class SettingsStore { this.db.prepare('DELETE FROM overseer_settings WHERE key = ?').run(key) } - /** Read the persisted active brain, or null when the operator has never chosen one (use env default). */ - getActiveBrain(): ActiveBrainSetting | null { - const raw = this.get(ACTIVE_BRAIN_KEY) + /** Read the persisted active brain for a namespace, or null when unset. */ + getActiveBrain(namespace = 'default'): ActiveBrainSetting | null { + const raw = this.get(activeBrainKey(namespace)) if (!raw) return null try { const parsed = JSON.parse(raw) as Partial @@ -60,11 +65,11 @@ export class SettingsStore { } } - setActiveBrain(value: ActiveBrainSetting): void { - this.set(ACTIVE_BRAIN_KEY, JSON.stringify({ profile: value.profile, model: value.model ?? null })) + setActiveBrain(value: ActiveBrainSetting, namespace = 'default'): void { + this.set(activeBrainKey(namespace), JSON.stringify({ profile: value.profile, model: value.model ?? null })) } - clearActiveBrain(): void { - this.delete(ACTIVE_BRAIN_KEY) + clearActiveBrain(namespace = 'default'): void { + this.delete(activeBrainKey(namespace)) } } diff --git a/hub/src/sync/overseerEntity.test.ts b/hub/src/sync/overseerEntity.test.ts index 4a9b13e612..3571eb13a3 100644 --- a/hub/src/sync/overseerEntity.test.ts +++ b/hub/src/sync/overseerEntity.test.ts @@ -540,3 +540,47 @@ describe('OverseerEntity dispositions (Stage 1 keystone)', () => { expect(lastRelay?.message).not.toBe('wrong target') }) }) + +describe('OverseerEntity namespace isolation (#107 kill criterion)', () => { + it('refuses to resolve or relay a session that only exists in another namespace', async () => { + const store = new Store(':memory:') + const inA = store.sessions.getOrCreateSession( + 'ns-a-session-aaaaaaaa', + { flavor: 'claude', path: '/tmp/a', name: 'worker-a' }, + null, + 'ns-a' + ) + const inB = store.sessions.getOrCreateSession( + 'ns-b-session-bbbbbbbb', + { flavor: 'claude', path: '/tmp/b', name: 'worker-b' }, + null, + 'ns-b' + ) + store.events.insert({ + ts: Date.now(), + sourceKind: 'worker', + eventType: 'blocked', + attentionCandidate: 1, + severity: 4, + summary: 'secret from B', + relatedSessionId: inB.id, + payloadJson: JSON.stringify({ session: { project: 'secret', name: 'worker-b' } }) + }) + + const engine = buildEngine(store) + const overseerA = engine.getOverseer('ns-a') + const overseerB = engine.getOverseer('ns-b') + + expect(overseerA.getSessionState(inA.id)?.sessionId).toBe(inA.id) + expect(overseerA.getSessionState(inB.id)).toBeNull() + expect(overseerA.queryEvents({}).map((e) => e.summary)).not.toContain('secret from B') + expect(overseerB.queryEvents({}).map((e) => e.summary)).toContain('secret from B') + + const crossRelay = await overseerA.pingSession({ + sessionId: inB.id, + message: 'cross-namespace injection' + }) + expect(crossRelay.ok).toBe(false) + expect(crossRelay.error).toMatch(/No session matching/) + }) +}) diff --git a/hub/src/sync/overseerEntity.ts b/hub/src/sync/overseerEntity.ts index 3679eec268..ddadd39a7a 100644 --- a/hub/src/sync/overseerEntity.ts +++ b/hub/src/sync/overseerEntity.ts @@ -21,6 +21,7 @@ import { isNoOpAction, mapEventTypeToWorkerState, openLoopBucket, + OVERSEER_DISPOSITION_ACTIONS, type OverseerActiveWorker, type OverseerConvoTurnInput, type OverseerExplainPriority, @@ -145,7 +146,7 @@ export class OverseerEntity { untilTs: args.untilTs ?? null, beforeId: args.beforeId ?? null, limit: args.limit ?? 50 - }) + }).filter((event) => this.eventInCallerScope(event)) } // --- Tool 2: query_inbox ------------------------------------------------- @@ -171,7 +172,7 @@ export class OverseerEntity { category: args.category ?? null, limit: args.limit ?? 50, includeSleepingSnoozed: statusesExplicit && statuses.includes('snoozed') - }) + }).filter((item) => this.itemInCallerScope(item)) return { items, candidates: items.filter((item) => item.status === 'new'), @@ -474,6 +475,7 @@ export class OverseerEntity { queryDispositions(args: QueryDispositionsArgs = {}): OverseerDispositionsResult { const filter = { action: args.action ?? null, + actionsAllowlist: args.action ? null : [...OVERSEER_DISPOSITION_ACTIONS], sourceKind: args.sourceKind ?? null, sourceRef: args.sourceRef ?? null, eventType: args.eventType ?? null, @@ -499,24 +501,30 @@ export class OverseerEntity { return { mode: 'cluster', clusters, total: clusters.length } } - const rows = this.inbox.listDispositions(filter).map( - (r): OverseerDispositionRow => ({ - id: r.id, - itemId: r.inboxItemId, - action: r.action, - statusAfter: r.statusAfter, - feedback: r.feedback, - createdAt: r.createdAt, - sourceKind: r.sourceKind, - sourceRef: r.sourceRef, - eventType: r.eventType, - category: r.category, - project: r.project, - artifactKind: r.artifactKind, - repo: r.repo, - title: r.contextSnapshot?.title ?? null + const rows = this.inbox + .listDispositions(filter) + .filter((r) => { + const item = this.inbox.getById(r.inboxItemId) + return item != null && this.itemInCallerScope(item) }) - ) + .map( + (r): OverseerDispositionRow => ({ + id: r.id, + itemId: r.inboxItemId, + action: r.action, + statusAfter: r.statusAfter, + feedback: r.feedback, + createdAt: r.createdAt, + sourceKind: r.sourceKind, + sourceRef: r.sourceRef, + eventType: r.eventType, + category: r.category, + project: r.project, + artifactKind: r.artifactKind, + repo: r.repo, + title: r.contextSnapshot?.title ?? null + }) + ) return { mode: 'list', rows, total: rows.length } } @@ -531,7 +539,7 @@ export class OverseerEntity { */ recordDisposition(args: RecordDispositionArgs): OverseerDispositionResult { const item = this.inbox.getById(args.itemId) - if (!item) { + if (!item || !this.itemInCallerScope(item)) { return { ok: false, itemId: args.itemId, @@ -620,12 +628,16 @@ export class OverseerEntity { } const tombstone = `Relayed to ${label} (${deliveredSessionId.slice(0, 8)})${result.resumed ? ' [resumed]' : ''}: "${snippet}"` - this.recordDispatchedRelay({ - sessionId: deliveredSessionId, - message: args.message, - resumed: result.resumed, - tombstone - }) + try { + this.recordDispatchedRelay({ + sessionId: deliveredSessionId, + message: args.message, + resumed: result.resumed, + tombstone + }) + } catch { + // Delivery already succeeded — audit failure must not flip ok or invite retry. + } return { ok: true, sessionId: deliveredSessionId, @@ -663,7 +675,7 @@ export class OverseerEntity { if (args.itemId != null) { const item = this.inbox.getById(args.itemId) - if (!item) { + if (!item || !this.itemInCallerScope(item)) { return { ok: false, error: `No inbox item #${args.itemId} — nothing relayed.` } } sessionIdFromItem = item.relatedSessionId?.trim() || null @@ -774,6 +786,24 @@ export class OverseerEntity { // --- internals ----------------------------------------------------------- + /** + * Fail-closed for cross-namespace: a related session must resolve uniquely + * inside this entity's injected session list (already namespace-scoped). + * Rows with no relatedSessionId stay visible until #107 adds a namespace + * column (convo_turns / orphan inbox rows are hub-global today). + */ + private itemInCallerScope(item: { relatedSessionId: string | null }): boolean { + const related = item.relatedSessionId?.trim() + if (!related) return true + return this.matchSessions(related).length === 1 + } + + private eventInCallerScope(event: { relatedSessionId: string | null }): boolean { + const related = event.relatedSessionId?.trim() + if (!related) return true + return this.matchSessions(related).length === 1 + } + /** * Exact session id, else unique prefix (hapi-ping-peer / loomux / pi pattern). * Ambiguous or unknown → undefined. Never silently picks among collisions. diff --git a/hub/src/sync/syncEngine.ts b/hub/src/sync/syncEngine.ts index 69aea6a1f9..ca38b82f1e 100644 --- a/hub/src/sync/syncEngine.ts +++ b/hub/src/sync/syncEngine.ts @@ -146,7 +146,7 @@ export class SyncEngine { private readonly messageService: MessageService private readonly rpcGateway: RpcGateway private readonly overseerEvents: OverseerEventRecorder - private readonly overseer: OverseerEntity + private readonly overseerByNamespace: Map private inactivityTimer: NodeJS.Timeout | null = null /** Sessions that emitted `session-ready` (Cursor ACP load/newSession complete). */ private readonly sessionReadyIds = new Set() @@ -168,25 +168,9 @@ export class SyncEngine { ) this.rpcGateway = new RpcGateway(io, rpcRegistry) this.overseerEvents = new OverseerEventRecorder(store.events, store.inbox) - this.overseer = new OverseerEntity({ - events: store.events, - inbox: store.inbox, - messages: store.messages, - getSession: (sessionId) => this.getSession(sessionId), - getSessions: () => this.getSessions(), - relayToSession: async ({ sessionId, message, namespace = 'default' }) => { - // Resume may return a replacement hub session id after spawn+merge; - // always deliver to that id (see executeSessionRelay). - return executeSessionRelay( - { - getSession: (id) => this.getSession(id), - resumeSession: (id, ns) => this.resumeSession(id, ns), - sendMessage: (id, payload) => this.sendMessage(id, payload) - }, - { sessionId, message, namespace } - ) - } - }) + this.overseerByNamespace = new Map() + // Default namespace entity — most tests call getOverseer() with no arg. + this.overseerByNamespace.set('default', this.createOverseerForNamespace('default')) this.reloadAll() this.inactivityTimer = setInterval(() => this.expireInactive(), 5_000) } @@ -350,8 +334,40 @@ export class SyncEngine { this.eventPublisher.emit(event) } - getOverseer(): OverseerEntity { - return this.overseer + getOverseer(namespace = 'default'): OverseerEntity { + const key = namespace.trim() || 'default' + let entity = this.overseerByNamespace.get(key) + if (!entity) { + entity = this.createOverseerForNamespace(key) + this.overseerByNamespace.set(key, entity) + } + return entity + } + + private createOverseerForNamespace(namespace: string): OverseerEntity { + return new OverseerEntity({ + events: this.store.events, + inbox: this.store.inbox, + messages: this.store.messages, + getSession: (sessionId) => { + const access = this.resolveSessionAccess(sessionId, namespace) + return access.ok ? access.session : undefined + }, + getSessions: () => this.getSessionsByNamespace(namespace), + relayToSession: async ({ sessionId, message, namespace: relayNs = namespace }) => { + return executeSessionRelay( + { + getSession: (id) => { + const access = this.resolveSessionAccess(id, relayNs) + return access.ok ? { active: access.session.active } : undefined + }, + resumeSession: (id, ns) => this.resumeSession(id, ns, { allowFreshSpawn: false }), + sendMessage: (id, payload) => this.sendMessage(id, payload) + }, + { sessionId, message, namespace: relayNs } + ) + } + }) } /** Hub settings KV (persisted active brain, etc.) — see SettingsStore. */ @@ -1222,7 +1238,11 @@ export class SyncEngine { return this.store.messages.getFirstMessages(sessionId, 1).length === 0 } - async resumeSession(sessionId: string, namespace: string, opts?: { permissionMode?: PermissionMode }): Promise { + async resumeSession(sessionId: string, namespace: string, opts?: { + permissionMode?: PermissionMode + /** When false, refuse never-started stubs that would fresh-spawn (Overseer relay). Default true. */ + allowFreshSpawn?: boolean + }): Promise { const access = this.sessionCache.resolveSessionAccess(sessionId, namespace) if (!access.ok) { return { @@ -1256,6 +1276,7 @@ export class SyncEngine { directory = targetResult.target.directory } else if ( targetResult.code === 'resume_unavailable' + && opts?.allowFreshSpawn !== false && this.canFreshSpawnNeverStartedSession(session, access.sessionId, namespace) ) { const metadata = session.metadata! diff --git a/hub/src/web/routes/overseer.ts b/hub/src/web/routes/overseer.ts index 57c7d1e76d..b4ef79a8a8 100644 --- a/hub/src/web/routes/overseer.ts +++ b/hub/src/web/routes/overseer.ts @@ -45,12 +45,12 @@ const activeBrainBodySchema = z.object({ }) /** Drop a persisted active brain when its profile was removed from env after restart. */ -function getSanitizedActiveBrain(engine: SyncEngine): ActiveBrainSetting | null { +function getSanitizedActiveBrain(engine: SyncEngine, namespace = 'default'): ActiveBrainSetting | null { const settings = engine.getSettings() - const active = settings.getActiveBrain() + const active = settings.getActiveBrain(namespace) if (!active) return null if (isKnownBrainProfile(active.profile)) return active - settings.clearActiveBrain() + settings.clearActiveBrain(namespace) return null } @@ -87,7 +87,7 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho if (engine instanceof Response) return engine return c.json({ profiles: listBrainProfiles(process.env), - active: getSanitizedActiveBrain(engine) + active: getSanitizedActiveBrain(engine, c.get('namespace')) }) }) @@ -97,7 +97,7 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho app.get('/overseer/brain/active', (c) => { const engine = requireSyncEngine(c, getSyncEngine) if (engine instanceof Response) return engine - const active = getSanitizedActiveBrain(engine) + const active = getSanitizedActiveBrain(engine, c.get('namespace')) const selection = resolveBrainSelection(active) const config = resolveBrainConfig(process.env, selection) return c.json({ @@ -129,7 +129,7 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho } const active = { profile: parsed.data.profile, model: parsed.data.model ?? null } - engine.getSettings().setActiveBrain(active) + engine.getSettings().setActiveBrain(active, c.get('namespace')) return c.json({ active }) }) @@ -172,7 +172,7 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho } try { - const result = await runOverseerTool(engine.getOverseer(), tool, body ?? {}) + const result = await runOverseerTool(engine.getOverseer(c.get('namespace')), tool, body ?? {}) return c.json({ tool, result }) } catch (error) { if (error instanceof z.ZodError) { @@ -185,6 +185,8 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho } }) + }) + // Converse — the modality-agnostic conversation core. Runs the brain LLM // with the read-only tools and returns a human-facing reply + tool trace. // Text is the first transport (debug settings); voice/XR reuse this. When @@ -210,7 +212,8 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho return c.json({ error: 'Last message must be from the operator' }, 400) } - const active = getSanitizedActiveBrain(engine) + const overseer = engine.getOverseer(c.get('namespace')) + const active = getSanitizedActiveBrain(engine, c.get('namespace')) const config = resolveBrainConfig(process.env, resolveBrainSelection(active, { profile: parsed.data.profile, model: parsed.data.model @@ -226,14 +229,14 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho try { const { reply, toolTrace } = await runOverseerConverse({ - overseer: engine.getOverseer(), + overseer, config, messages, allowWrites: parsed.data.allowWrites }) const lastOperator = [...messages].reverse().find((m) => m.role === 'operator')?.content ?? '' - engine.getOverseer().recordConvoTurn({ + overseer.recordConvoTurn({ operatorText: lastOperator, overseerText: reply, relatedSessionId: parsed.data.relatedSessionId ?? null, @@ -282,7 +285,7 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho return c.json({ error: 'operatorText or overseerText is required' }, 400) } - const event = engine.getOverseer().recordConvoTurn({ + const event = engine.getOverseer(c.get('namespace')).recordConvoTurn({ operatorText: parsed.data.operatorText, overseerText: parsed.data.overseerText, relatedSessionId: parsed.data.relatedSessionId ?? null, diff --git a/shared/src/overseerWriteIntent.test.ts b/shared/src/overseerWriteIntent.test.ts index 788cfc748d..240e849d1a 100644 --- a/shared/src/overseerWriteIntent.test.ts +++ b/shared/src/overseerWriteIntent.test.ts @@ -1,7 +1,7 @@ import { describe, expect, it } from 'vitest' import { detectOperatorWriteTools, - isWriteToolAuthorized, + isWriteToolCallAuthorized, resolveOverseerWriteAuthorization } from './overseerWriteIntent' @@ -24,7 +24,6 @@ describe('detectOperatorWriteTools', () => { it('does not authorize writes for read-only questions', () => { expect([...detectOperatorWriteTools('what needs my attention?')]).toEqual([]) - expect([...detectOperatorWriteTools('show recent output for session abc')]).toEqual([]) }) }) @@ -35,15 +34,52 @@ describe('resolveOverseerWriteAuthorization', () => { allowWrites: true }) expect(auth.explicitClientFlag).toBe(true) - expect(isWriteToolAuthorized('ping_session', auth)).toBe(true) - expect(isWriteToolAuthorized('record_disposition', auth)).toBe(true) + expect(isWriteToolCallAuthorized('ping_session', { + sessionId: 'abcdef12-0000-0000-0000-000000000001', + message: 'hi' + }, auth).ok).toBe(true) + }) + + it('binds ping_session to the session id named by the operator', () => { + const auth = resolveOverseerWriteAuthorization({ + latestOperatorText: 'ping session abcdef12: "please continue"' + }) + expect(isWriteToolCallAuthorized('ping_session', { + sessionId: 'abcdef12-ffff-ffff-ffff-ffffffffffff', + message: 'please continue' + }, auth).ok).toBe(true) + expect(isWriteToolCallAuthorized('ping_session', { + sessionId: 'deadbeef-ffff-ffff-ffff-ffffffffffff', + message: 'please continue' + }, auth).ok).toBe(false) + }) + + it('binds short named session tokens after the word session', () => { + const auth = resolveOverseerWriteAuthorization({ + latestOperatorText: 'ping session sess-1: "hi"' + }) + expect(isWriteToolCallAuthorized('ping_session', { + sessionId: 'sess-1', + message: 'hi' + }, auth).ok).toBe(true) + }) + + it('denies ping without a concrete target in the operator message', () => { + const auth = resolveOverseerWriteAuthorization({ + latestOperatorText: 'ping that worker to continue' + }) + const result = isWriteToolCallAuthorized('ping_session', { + sessionId: 'abcdef12', + message: 'continue' + }, auth) + expect(result.ok).toBe(false) }) it('denies write tools when neither flag nor intent matches', () => { const auth = resolveOverseerWriteAuthorization({ latestOperatorText: 'summarize the inbox' }) - expect(isWriteToolAuthorized('ping_session', auth)).toBe(false) - expect(isWriteToolAuthorized('query_inbox', auth)).toBe(true) + expect(isWriteToolCallAuthorized('ping_session', { sessionId: 'x', message: 'y' }, auth).ok).toBe(false) + expect(isWriteToolCallAuthorized('query_inbox', {}, auth).ok).toBe(true) }) }) diff --git a/shared/src/overseerWriteIntent.ts b/shared/src/overseerWriteIntent.ts index c84bbdb8bf..5bba4ba99e 100644 --- a/shared/src/overseerWriteIntent.ts +++ b/shared/src/overseerWriteIntent.ts @@ -6,15 +6,21 @@ * prompt-inject a relay/disposition. Authorization comes from the operator's * latest utterance and/or an explicit client `allowWrites` flag — never from * model-selected tools alone. + * + * Grants are bound to extracted targets/payloads when present so a later + * injected tool call cannot retarget a legitimate "ping session X" grant. */ import { isOverseerWriteTool, type OverseerWriteToolName } from './overseerEntity' export type OverseerWriteAuthorization = { - /** Tools the operator's message (or explicit flag) authorized for this turn. */ allowed: ReadonlySet /** True when the client sent allowWrites: true (admin console / voice confirm). */ explicitClientFlag: boolean + sessionIdPrefixes: readonly string[] + itemIds: readonly number[] + /** Quoted snippets from the operator line that a relay message should match. */ + messageSnippets: readonly string[] } const RELAY_INTENT = @@ -23,6 +29,44 @@ const RELAY_INTENT = const DISPOSITION_INTENT = /\b(snooze|dismiss|reopen|dispose)\b|\bmark\b[\s\S]{0,40}\bdone\b|\b(resolve|done with)\b/i +/** UUID or hex-prefix session ids (production hub shape). */ +const UUID_OR_HEX_SESSION_RE = + /\b([0-9a-f]{8}(?:-[0-9a-f]{4}){3}-[0-9a-f]{12}|[0-9a-f]{8,})\b/gi +/** Explicit `session ` form — covers short test ids like `sess-1` / `old-id`. */ +const NAMED_SESSION_RE = /\bsession\s+([a-z0-9][a-z0-9_-]{1,63})\b/gi +const ITEM_ID_RE = /\b(?:item\s*#?|#)(\d+)\b/gi + +function extractSessionIdPrefixes(text: string): string[] { + const out: string[] = [] + for (const match of text.matchAll(UUID_OR_HEX_SESSION_RE)) { + const value = match[1]?.toLowerCase() + if (value && !out.includes(value)) out.push(value) + } + for (const match of text.matchAll(NAMED_SESSION_RE)) { + const value = match[1]?.toLowerCase() + if (value && !out.includes(value)) out.push(value) + } + return out +} + +function extractItemIds(text: string): number[] { + const out: number[] = [] + for (const match of text.matchAll(ITEM_ID_RE)) { + const id = Number(match[1]) + if (Number.isFinite(id) && id > 0 && !out.includes(id)) out.push(id) + } + return out +} + +function extractQuotedSnippets(text: string): string[] { + const out: string[] = [] + for (const match of text.matchAll(/"([^"]{1,500})"|'([^']{1,500})'/g)) { + const value = (match[1] ?? match[2] ?? '').trim() + if (value && !out.includes(value)) out.push(value) + } + return out +} + /** Detect which write classes the latest operator message authorizes. */ export function detectOperatorWriteTools(operatorText: string): Set { const allowed = new Set() @@ -37,18 +81,98 @@ export function resolveOverseerWriteAuthorization(opts: { latestOperatorText: string allowWrites?: boolean }): OverseerWriteAuthorization { + const text = opts.latestOperatorText if (opts.allowWrites === true) { return { allowed: new Set(['ping_session', 'record_disposition']), - explicitClientFlag: true + explicitClientFlag: true, + sessionIdPrefixes: extractSessionIdPrefixes(text), + itemIds: extractItemIds(text), + messageSnippets: extractQuotedSnippets(text) } } return { - allowed: detectOperatorWriteTools(opts.latestOperatorText), - explicitClientFlag: false + allowed: detectOperatorWriteTools(text), + explicitClientFlag: false, + sessionIdPrefixes: extractSessionIdPrefixes(text), + itemIds: extractItemIds(text), + messageSnippets: extractQuotedSnippets(text) + } +} + +function sessionIdMatchesGrant(sessionId: string, prefixes: readonly string[]): boolean { + const lower = sessionId.trim().toLowerCase() + return prefixes.some((prefix) => lower === prefix || lower.startsWith(prefix)) +} + +function messageMatchesGrant(message: string, snippets: readonly string[]): boolean { + if (snippets.length === 0) return true + return snippets.some((snippet) => message.includes(snippet)) +} + +/** + * Per-call authorization: tool class must be allowed, and when the operator + * named a target, the call args must bind to it (unless explicitClientFlag). + */ +export function isWriteToolCallAuthorized( + tool: string, + args: Record, + auth: OverseerWriteAuthorization +): { ok: true } | { ok: false; error: string } { + if (!isOverseerWriteTool(tool)) return { ok: true } + if (!auth.allowed.has(tool)) { + return { ok: false, error: 'write not authorized by operator message (no explicit write intent)' } } + + if (tool === 'ping_session') { + const sessionId = typeof args.sessionId === 'string' ? args.sessionId.trim() : '' + const itemId = typeof args.itemId === 'number' ? args.itemId : null + const message = typeof args.message === 'string' ? args.message : '' + + if (auth.explicitClientFlag) { + if (!messageMatchesGrant(message, auth.messageSnippets)) { + return { ok: false, error: 'relay message does not match operator-quoted payload' } + } + return { ok: true } + } + + const hasTargetGrant = auth.sessionIdPrefixes.length > 0 || auth.itemIds.length > 0 + if (!hasTargetGrant) { + return { + ok: false, + error: 'relay requires an explicit session id / item id in the operator message (or allowWrites)' + } + } + const sessionOk = sessionId.length > 0 && sessionIdMatchesGrant(sessionId, auth.sessionIdPrefixes) + const itemOk = itemId != null && auth.itemIds.includes(itemId) + if (!sessionOk && !itemOk) { + return { ok: false, error: 'relay target does not match operator-authorized session/item' } + } + if (!messageMatchesGrant(message, auth.messageSnippets)) { + return { ok: false, error: 'relay message does not match operator-quoted payload' } + } + return { ok: true } + } + + if (tool === 'record_disposition') { + const itemId = typeof args.itemId === 'number' ? args.itemId : null + if (auth.explicitClientFlag) return { ok: true } + if (auth.itemIds.length === 0) { + return { + ok: false, + error: 'disposition requires an explicit item id in the operator message (or allowWrites)' + } + } + if (itemId == null || !auth.itemIds.includes(itemId)) { + return { ok: false, error: 'disposition itemId does not match operator-authorized item' } + } + return { ok: true } + } + + return { ok: true } } +/** @deprecated Prefer isWriteToolCallAuthorized — class-only check is insufficient. */ export function isWriteToolAuthorized( tool: string, auth: OverseerWriteAuthorization @@ -56,3 +180,8 @@ export function isWriteToolAuthorized( if (!isOverseerWriteTool(tool)) return true return auth.allowed.has(tool) } + +export function fingerprintWriteToolCall(tool: string, args: Record): string { + const normalized = JSON.stringify(args, Object.keys(args).sort()) + return `${tool}:${normalized}` +} From 483c08815f13756a36e97698a33706f6f7654da6 Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sun, 2 Aug 2026 00:39:25 +0100 Subject: [PATCH 6/6] fix(overseer): remove stray brace from routes conflict resolve Co-authored-by: Cursor --- hub/src/web/routes/overseer.ts | 2 -- 1 file changed, 2 deletions(-) diff --git a/hub/src/web/routes/overseer.ts b/hub/src/web/routes/overseer.ts index b4ef79a8a8..74fa0aa8db 100644 --- a/hub/src/web/routes/overseer.ts +++ b/hub/src/web/routes/overseer.ts @@ -185,8 +185,6 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho } }) - }) - // Converse — the modality-agnostic conversation core. Runs the brain LLM // with the read-only tools and returns a human-facing reply + tool trace. // Text is the first transport (debug settings); voice/XR reuse this. When