From 859db0510837d368ebc99e756a799577761be0fc Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sat, 1 Aug 2026 20:23:22 +0100 Subject: [PATCH 1/4] feat(overseer): hub-owned converse context from convo_turns Assemble budgeted prior convo_turn history on every converse call so text/voice transports stop restarting cold. Add GET /converse/recent for talk-to hydrate; UI sends only the latest operator line. Fixes #105. Co-authored-by: Cursor --- hub/src/overseer/converseContext.test.ts | 104 ++++++++++ hub/src/overseer/converseContext.ts | 195 ++++++++++++++++++ hub/src/web/routes/overseer.test.ts | 20 ++ hub/src/web/routes/overseer.ts | 47 ++++- shared/src/overseerConverse.ts | 14 ++ web/src/api/client.ts | 8 + .../settings/OverseerChatDebugControls.tsx | 64 +++++- 7 files changed, 439 insertions(+), 13 deletions(-) create mode 100644 hub/src/overseer/converseContext.test.ts create mode 100644 hub/src/overseer/converseContext.ts diff --git a/hub/src/overseer/converseContext.test.ts b/hub/src/overseer/converseContext.test.ts new file mode 100644 index 0000000000..061bf6d1c7 --- /dev/null +++ b/hub/src/overseer/converseContext.test.ts @@ -0,0 +1,104 @@ +import { describe, expect, it } from 'bun:test' +import { Store } from '../store' +import { SyncEngine } from '../sync/syncEngine' +import { RpcRegistry } from '../socket/rpcRegistry' +import { + assembleOverseerConverseMessages, + budgetConvoTurns, + DEFAULT_CONVERSE_HISTORY_MAX_CHARS, + listRecentConvoTurns, + parseConvoTurnPayload +} from './converseContext' + +function buildEngine(store: Store): SyncEngine { + const io = { of: () => ({ to: () => ({ emit: () => {}, timeout: () => ({ emit: () => {} }) }) }) } as never + return new SyncEngine(store, io, new RpcRegistry(), { broadcast: () => {} } as never) +} + +describe('converseContext assembler', () => { + it('parseConvoTurnPayload reads operator/overseer/toolCalls', () => { + const parsed = parseConvoTurnPayload(JSON.stringify({ + operatorText: 'What needs me?', + overseerText: 'Three PRs.', + toolCalls: [{ tool: 'query_inbox', argsSummary: '{"limit":10}' }] + })) + expect(parsed.operatorText).toBe('What needs me?') + expect(parsed.overseerText).toBe('Three PRs.') + expect(parsed.toolCalls[0]?.tool).toBe('query_inbox') + }) + + it('hydrates prior convo_turns and appends only the latest client operator line', () => { + const store = new Store(':memory:') + const engine = buildEngine(store) + const overseer = engine.getOverseer() + + overseer.recordConvoTurn({ + operatorText: 'Which agents are blocked?', + overseerText: 'None are blocked.', + ts: 1000 + }) + overseer.recordConvoTurn({ + operatorText: 'What about expenses?', + overseerText: 'Item #56 is waiting.', + ts: 2000 + }) + + const assembled = assembleOverseerConverseMessages({ + overseer, + // Client still has stale local history + a new ask — hub ignores prior client turns. + clientMessages: [ + { role: 'operator', content: 'stale local only' }, + { role: 'overseer', content: 'stale reply' }, + { role: 'operator', content: 'Ok next?' } + ] + }) + + expect(assembled.hydratedTurns).toBe(2) + expect(assembled.truncated).toBe(false) + expect(assembled.messages.map((m) => m.role)).toEqual([ + 'operator', 'overseer', 'operator', 'overseer', 'operator' + ]) + expect(assembled.messages.map((m) => m.content)).toEqual([ + 'Which agents are blocked?', + 'None are blocked.', + 'What about expenses?', + 'Item #56 is waiting.', + 'Ok next?' + ]) + // Client stale lines never enter the brain context. + expect(assembled.messages.some((m) => m.content.includes('stale'))).toBe(false) + }) + + it('budget drops oldest turns when over char budget (kill criterion)', () => { + const turns = Array.from({ length: 8 }, (_, i) => ({ + id: i + 1, + ts: 1000 + i, + operatorText: `q${i} ${'x'.repeat(800)}`, + overseerText: `a${i} ${'y'.repeat(800)}`, + relatedSessionId: null, + toolCalls: [] as Array<{ tool: import('@hapi/protocol').OverseerToolName; argsSummary?: string }> + })) + const { turns: kept, truncated } = budgetConvoTurns(turns, { + maxTurns: 16, + maxChars: 4_000 + }) + expect(truncated).toBe(true) + expect(kept.length).toBeGreaterThan(0) + expect(kept.length).toBeLessThan(8) + expect(kept[0]!.id).toBeGreaterThan(1) + + // Unbounded dump must not survive default budget either when forced small. + const tiny = budgetConvoTurns(turns, { maxTurns: 2, maxChars: DEFAULT_CONVERSE_HISTORY_MAX_CHARS }) + expect(tiny.turns.length).toBe(2) + expect(tiny.truncated).toBe(true) + }) + + it('listRecentConvoTurns returns chronological views for UI hydrate', () => { + const store = new Store(':memory:') + const overseer = buildEngine(store).getOverseer() + overseer.recordConvoTurn({ operatorText: 'first', overseerText: 'one', ts: 1 }) + overseer.recordConvoTurn({ operatorText: 'second', overseerText: 'two', ts: 2 }) + const list = listRecentConvoTurns(overseer, { limit: 10 }) + expect(list.map((t) => t.operatorText)).toEqual(['first', 'second']) + }) +}) diff --git a/hub/src/overseer/converseContext.ts b/hub/src/overseer/converseContext.ts new file mode 100644 index 0000000000..ec2d22ac06 --- /dev/null +++ b/hub/src/overseer/converseContext.ts @@ -0,0 +1,195 @@ +/** + * Hub-owned Overseer converse context. + * + * Transports (text / voice / XR) send the latest operator utterance; the hub + * hydrates prior `convo_turn` events into `messages` before `runOverseerConverse`. + * No transport-local chat DB. Disposition tombstones stay available via the + * existing `query_dispositions` tool (tombstone-as-summarizer) — this module + * does not invent a second summarize engine. + */ + +import { + OVERSEER_CONVO_TURN_EVENT_TYPE, + type OverseerConverseMessage, + type OverseerToolName +} from '@hapi/protocol' +import type { OverseerEntity } from '../sync/overseerEntity' +import type { StoredSystemEvent } from '../store' + +/** Max prior operator↔overseer pairs to load from the events table. */ +export const DEFAULT_CONVERSE_HISTORY_MAX_TURNS = 16 + +/** + * Soft char budget for hydrated history (excludes the latest operator line). + * Keeps the brain window from filling with old tool-trace prose; drop oldest first. + */ +export const DEFAULT_CONVERSE_HISTORY_MAX_CHARS = 24_000 + +export type StoredConvoTurnView = { + id: number + ts: number + operatorText: string + overseerText: string + relatedSessionId: string | null + toolCalls: Array<{ tool: OverseerToolName; argsSummary?: string }> +} + +export type AssembleConverseContextResult = { + /** Oldest-first messages for the brain, ending with the latest operator line. */ + messages: OverseerConverseMessage[] + /** How many prior convo_turn rows contributed (after budget trim). */ + hydratedTurns: number + /** True when older turns were dropped to stay under maxChars / maxTurns. */ + truncated: boolean +} + +function isObj(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +export function parseConvoTurnPayload(payloadJson: string | null): { + operatorText: string + overseerText: string + toolCalls: Array<{ tool: OverseerToolName; argsSummary?: string }> +} { + if (!payloadJson) { + return { operatorText: '', overseerText: '', toolCalls: [] } + } + try { + const parsed: unknown = JSON.parse(payloadJson) + if (!isObj(parsed)) { + return { operatorText: '', overseerText: '', toolCalls: [] } + } + const operatorText = typeof parsed.operatorText === 'string' ? parsed.operatorText : '' + const overseerText = typeof parsed.overseerText === 'string' ? parsed.overseerText : '' + const rawCalls = Array.isArray(parsed.toolCalls) ? parsed.toolCalls : [] + const toolCalls: Array<{ tool: OverseerToolName; argsSummary?: string }> = [] + for (const call of rawCalls) { + if (!isObj(call) || typeof call.tool !== 'string') continue + toolCalls.push({ + tool: call.tool as OverseerToolName, + argsSummary: typeof call.argsSummary === 'string' ? call.argsSummary : undefined + }) + } + return { operatorText, overseerText, toolCalls } + } catch { + return { operatorText: '', overseerText: '', toolCalls: [] } + } +} + +export function eventToConvoTurnView(event: StoredSystemEvent): StoredConvoTurnView | null { + if (event.eventType !== OVERSEER_CONVO_TURN_EVENT_TYPE) return null + const payload = parseConvoTurnPayload(event.payloadJson) + if (!payload.operatorText.trim() && !payload.overseerText.trim()) return null + return { + id: event.id, + ts: event.ts, + operatorText: payload.operatorText, + overseerText: payload.overseerText, + relatedSessionId: event.relatedSessionId, + toolCalls: payload.toolCalls + } +} + +/** Newest-first from the store; returned oldest-first for display / assemble. */ +export function listRecentConvoTurns( + overseer: OverseerEntity, + opts: { limit?: number } = {} +): StoredConvoTurnView[] { + const limit = Math.min(Math.max(opts.limit ?? DEFAULT_CONVERSE_HISTORY_MAX_TURNS, 1), 50) + const events = overseer.queryEvents({ + eventType: OVERSEER_CONVO_TURN_EVENT_TYPE, + limit + }) + const views: StoredConvoTurnView[] = [] + for (const event of events) { + const view = eventToConvoTurnView(event) + if (view) views.push(view) + } + // events.query is id DESC — reverse to chronological. + return views.reverse() +} + +function turnsToMessages(turns: StoredConvoTurnView[]): OverseerConverseMessage[] { + const messages: OverseerConverseMessage[] = [] + for (const turn of turns) { + const op = turn.operatorText.trim() + const ov = turn.overseerText.trim() + if (op) messages.push({ role: 'operator', content: op }) + if (ov) messages.push({ role: 'overseer', content: ov }) + } + return messages +} + +function messagesCharCount(messages: OverseerConverseMessage[]): number { + return messages.reduce((sum, m) => sum + m.content.length, 0) +} + +/** + * Apply turn + char budgets by dropping the oldest turns first. + * Returns the kept turns (oldest-first) and whether anything was dropped. + */ +export function budgetConvoTurns( + turnsOldestFirst: StoredConvoTurnView[], + opts: { maxTurns?: number; maxChars?: number } = {} +): { turns: StoredConvoTurnView[]; truncated: boolean } { + const maxTurns = opts.maxTurns ?? DEFAULT_CONVERSE_HISTORY_MAX_TURNS + const maxChars = opts.maxChars ?? DEFAULT_CONVERSE_HISTORY_MAX_CHARS + + let kept = turnsOldestFirst + let truncated = false + if (kept.length > maxTurns) { + kept = kept.slice(kept.length - maxTurns) + truncated = true + } + + while (kept.length > 0 && messagesCharCount(turnsToMessages(kept)) > maxChars) { + kept = kept.slice(1) + truncated = true + } + return { turns: kept, truncated } +} + +/** + * Hub-owned assemble: prior `convo_turn`s (budgeted) + the client's latest operator line. + * Ignores any prior client history so transports cannot fork memory. + */ +export function assembleOverseerConverseMessages(params: { + overseer: OverseerEntity + clientMessages: OverseerConverseMessage[] + maxTurns?: number + maxChars?: number +}): AssembleConverseContextResult { + const { overseer, clientMessages, maxTurns, maxChars } = params + if (clientMessages.length === 0) { + throw new Error('clientMessages must include the latest operator utterance') + } + const latest = clientMessages[clientMessages.length - 1]! + if (latest.role !== 'operator') { + throw new Error('Last client message must be from the operator') + } + + const fetched = listRecentConvoTurns(overseer, { + limit: Math.max(maxTurns ?? DEFAULT_CONVERSE_HISTORY_MAX_TURNS, 1) + }) + const { turns, truncated } = budgetConvoTurns(fetched, { maxTurns, maxChars }) + const history = turnsToMessages(turns) + + // If the operator re-sent the exact last logged question without a reply yet + // (unlikely — we only record after reply), avoid duplicating. Normal path: + // current utterance is not in the store until after this turn completes. + const lastHistory = history[history.length - 1] + if ( + lastHistory?.role === 'operator' + && lastHistory.content === latest.content + && (history.length < 2 || history[history.length - 2]?.role !== 'overseer') + ) { + return { messages: history, hydratedTurns: turns.length, truncated } + } + + return { + messages: [...history, latest], + hydratedTurns: turns.length, + truncated + } +} diff --git a/hub/src/web/routes/overseer.test.ts b/hub/src/web/routes/overseer.test.ts index da82940727..349bfded84 100644 --- a/hub/src/web/routes/overseer.test.ts +++ b/hub/src/web/routes/overseer.test.ts @@ -122,6 +122,26 @@ describe('overseer routes', () => { expect(res.status).toBe(400) }) + it('GET /overseer/converse/recent returns chronological hub-owned turns', async () => { + const store = new Store(':memory:') + const app = buildApp(store) + await app.request('/api/overseer/convo-turns', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ operatorText: 'first', overseerText: 'one' }) + }) + await app.request('/api/overseer/convo-turns', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ operatorText: 'second', overseerText: 'two' }) + }) + const res = await app.request('/api/overseer/converse/recent?limit=10') + expect(res.status).toBe(200) + const body = await res.json() as { turns: Array<{ operatorText: string; overseerText: string }> } + expect(body.turns.map((t) => t.operatorText)).toEqual(['first', 'second']) + expect(body.turns.map((t) => t.overseerText)).toEqual(['one', 'two']) + }) + it('GET /overseer/brains reports profiles + a null active until one is set', async () => { const prev = process.env.OVERSEER_BRAIN_URL process.env.OVERSEER_BRAIN_URL = 'http://brain.test/v1' diff --git a/hub/src/web/routes/overseer.ts b/hub/src/web/routes/overseer.ts index b4ef79a8a8..2ceff92dc0 100644 --- a/hub/src/web/routes/overseer.ts +++ b/hub/src/web/routes/overseer.ts @@ -12,6 +12,7 @@ import type { WebAppEnv } from '../middleware/auth' import { requireSyncEngine } from './guards' import { isOverseerToolName, OverseerWriteNotAllowedError, runOverseerTool } from '../../overseer/runOverseerTool' import { runOverseerConverse } from '../../overseer/converse' +import { assembleOverseerConverseMessages, listRecentConvoTurns } from '../../overseer/converseContext' import { BrainUnavailableError, filterChatModels, isKnownBrainProfile, listBrainModels, listBrainProfiles, resolveBrainConfig, resolveBrainSelection } from '../../overseer/brainClient' import type { ActiveBrainSetting } from '../../store/settingsStore' @@ -28,6 +29,11 @@ const convoTurnBodySchema = z.object({ }) const converseBodySchema = z.object({ + /** + * Transport may send full local history or just the latest operator line. + * Hub hydrates prior `convo_turn`s and keeps only the last operator utterance + * from this array (hub-owned memory — transports do not fork the thread). + */ messages: z.array(z.object({ role: z.enum(['operator', 'overseer']), content: z.string().max(8000) @@ -187,11 +193,25 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho }) + // Recent convo_turns for transport hydrate (talk-to reload, voice attach). + // Durable memory lives in events — this is a thin read, not a chat DB. + app.get('/overseer/converse/recent', (c) => { + const engine = requireSyncEngine(c, getSyncEngine) + if (engine instanceof Response) return engine + const rawLimit = Number(c.req.query('limit') ?? '20') + const limit = Number.isFinite(rawLimit) ? Math.min(Math.max(Math.trunc(rawLimit), 1), 50) : 20 + const { turns } = listRecentConvoTurns(engine.getOverseer(c.get('namespace')), { limit }) + return c.json({ turns }) + }) + // 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 // the brain is offline (GPU pulled for VR), returns brainOnline:false with a // friendly message rather than an error. + // + // Continuity: hub assembles prior `convo_turn`s (budgeted) + latest operator + // line. Transports do not own the thread. app.post('/overseer/converse', async (c) => { const engine = requireSyncEngine(c, getSyncEngine) if (engine instanceof Response) return engine @@ -207,12 +227,18 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho if (!parsed.success) { return c.json({ error: 'Invalid body', issues: parsed.error.flatten() }, 400) } - const messages = parsed.data.messages as OverseerConverseMessage[] - if (messages[messages.length - 1]?.role !== 'operator') { + const clientMessages = parsed.data.messages as OverseerConverseMessage[] + if (clientMessages[clientMessages.length - 1]?.role !== 'operator') { return c.json({ error: 'Last message must be from the operator' }, 400) } const overseer = engine.getOverseer(c.get('namespace')) + const assembled = assembleOverseerConverseMessages({ + overseer, + clientMessages + }) + const messages = assembled.messages + const active = getSanitizedActiveBrain(engine, c.get('namespace')) const config = resolveBrainConfig(process.env, resolveBrainSelection(active, { profile: parsed.data.profile, @@ -223,7 +249,9 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho reply: 'The Overseer brain is not configured on this hub (set OVERSEER_BRAIN_URL). I can still show raw events and inbox items, but I cannot answer in conversation yet.', toolTrace: [], model: null, - brainOnline: false + brainOnline: false, + hydratedTurns: assembled.hydratedTurns, + truncated: assembled.truncated }) } @@ -245,7 +273,14 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho .map((t) => ({ tool: t.tool, argsSummary: JSON.stringify(t.args).slice(0, 500) })) }) - return c.json({ reply, toolTrace, model: config.model, brainOnline: true }) + return c.json({ + reply, + toolTrace, + model: config.model, + brainOnline: true, + hydratedTurns: assembled.hydratedTurns, + truncated: assembled.truncated + }) } catch (error) { if (error instanceof BrainUnavailableError) { // Reachable-but-failed (http 4xx/5xx, malformed body) is a converse @@ -257,7 +292,9 @@ export function createOverseerRoutes(getSyncEngine: () => SyncEngine | null): Ho reply, toolTrace: [], model: config.model, - brainOnline: error.reachable + brainOnline: error.reachable, + hydratedTurns: assembled.hydratedTurns, + truncated: assembled.truncated }) } throw error diff --git a/shared/src/overseerConverse.ts b/shared/src/overseerConverse.ts index 6d6a6321ea..e98702f522 100644 --- a/shared/src/overseerConverse.ts +++ b/shared/src/overseerConverse.ts @@ -75,6 +75,20 @@ export type OverseerConverseResponse = { * treat this as an error. */ brainOnline: boolean + /** How many prior `convo_turn` rows the hub hydrated into this request. */ + hydratedTurns?: number + /** True when older turns were dropped to stay under the history budget. */ + truncated?: boolean +} + +/** One durable operator↔Overseer exchange for transport hydrate (UI / voice attach). */ +export type OverseerRecentConvoTurn = { + id: number + ts: number + operatorText: string + overseerText: string + relatedSessionId: string | null + toolCalls: Array<{ tool: OverseerToolName; argsSummary?: string }> } // --------------------------------------------------------------------------- diff --git a/web/src/api/client.ts b/web/src/api/client.ts index f594ad1218..de2df5dd50 100644 --- a/web/src/api/client.ts +++ b/web/src/api/client.ts @@ -834,6 +834,14 @@ export class ApiClient { }) } + /** Hub-owned recent Overseer turns for talk-to / voice hydrate after reload. */ + async fetchOverseerConverseRecent( + limit = 20 + ): Promise<{ turns: import('@hapi/protocol').OverseerRecentConvoTurn[] }> { + const q = new URLSearchParams({ limit: String(limit) }) + return await this.request(`/api/overseer/converse/recent?${q}`) + } + async fetchOverseerBrains(): Promise<{ profiles: import('@hapi/protocol').OverseerBrainProfileInfo[] active: { profile: string; model: string | null } | null diff --git a/web/src/components/settings/OverseerChatDebugControls.tsx b/web/src/components/settings/OverseerChatDebugControls.tsx index 815d7c5e81..863b7d3882 100644 --- a/web/src/components/settings/OverseerChatDebugControls.tsx +++ b/web/src/components/settings/OverseerChatDebugControls.tsx @@ -9,10 +9,9 @@ type ChatTurn = { brainOnline?: boolean } -// Debug-only text transport for the modality-agnostic Overseer converse core. -// This is deliberately a Settings/debug affordance, not a top-level surface: -// voice/XR are the intended first-class modalities and reuse the same -// /api/overseer/converse endpoint. Text is here only to exercise the loop. +// Text transport for the modality-agnostic Overseer converse core. +// Durable memory is hub-owned (`convo_turn` events); this panel hydrates on open +// and sends only the latest operator line — voice/XR reuse the same core. const STARTER_QUESTIONS = [ 'What needs my attention?', 'Which agents are blocked?', @@ -24,6 +23,7 @@ export function OverseerChatDebugControls() { const { api } = useAppContext() const [open, setOpen] = useState(false) const [turns, setTurns] = useState([]) + const [hydrated, setHydrated] = useState(false) const [input, setInput] = useState('') const [loading, setLoading] = useState(false) const [error, setError] = useState(null) @@ -43,6 +43,42 @@ export function OverseerChatDebugControls() { .catch(() => { /* brains list is optional chrome */ }) }, [open, api, profiles.length]) + // Hydrate from hub-owned convo_turns once per open (reload-safe continuity). + useEffect(() => { + if (!open || !api || hydrated) return + let cancelled = false + void api.fetchOverseerConverseRecent(24) + .then((res) => { + if (cancelled) return + const next: ChatTurn[] = [] + for (const turn of res.turns) { + if (turn.operatorText.trim()) { + next.push({ role: 'operator', content: turn.operatorText }) + } + if (turn.overseerText.trim()) { + next.push({ + role: 'overseer', + content: turn.overseerText, + toolTrace: turn.toolCalls.map((t) => ({ + tool: t.tool, + args: t.argsSummary ? safeJsonArgs(t.argsSummary) : {}, + ok: true + })) + }) + } + } + setTurns(next) + setHydrated(true) + requestAnimationFrame(() => { + scrollRef.current?.scrollTo({ top: scrollRef.current.scrollHeight }) + }) + }) + .catch(() => { + if (!cancelled) setHydrated(true) + }) + return () => { cancelled = true } + }, [open, api, hydrated]) + // Populate the model dropdown live from the selected profile's endpoint // (server proxies GET /models so the api key never reaches the browser). useEffect(() => { @@ -70,8 +106,8 @@ export function OverseerChatDebugControls() { setError(null) setInput('') - const history = turns.map((turn): OverseerConverseMessage => ({ role: turn.role, content: turn.content })) - const nextHistory: OverseerConverseMessage[] = [...history, { role: 'operator', content: trimmed }] + // Hub owns prior context — send only the new operator line. + const nextHistory: OverseerConverseMessage[] = [{ role: 'operator', content: trimmed }] setTurns((prev) => [...prev, { role: 'operator', content: trimmed }]) setLoading(true) try { @@ -94,7 +130,7 @@ export function OverseerChatDebugControls() { } finally { setLoading(false) } - }, [api, loading, turns, selectedProfile, selectedModel]) + }, [api, loading, selectedProfile, selectedModel]) return (
@@ -110,7 +146,7 @@ export function OverseerChatDebugControls() { {open && (

- 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…). + Fleet chief-of-staff (Stage 1.5 — read + dispositions + relay). Hub-owned memory via convo_turns — reload keeps the thread. 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…).

@@ -237,6 +273,7 @@ export function OverseerChatDebugControls() {
) } + +function safeJsonArgs(raw: string): Record { + try { + const parsed: unknown = JSON.parse(raw) + return typeof parsed === 'object' && parsed !== null && !Array.isArray(parsed) + ? parsed as Record + : { raw } + } catch { + return { raw } + } +} From 7795e64b6365d8fe87510f11a9fe8d6b73f66185 Mon Sep 17 00:00:00 2001 From: HeavyGee <133152184+heavygee@users.noreply.github.com> Date: Sat, 1 Aug 2026 20:46:32 +0100 Subject: [PATCH 2/4] fix(overseer): address Codex P2s on converse context hydrate - Rehydrate talk-to on every panel open; block send until hydrate settles - Report truncated when store query clips older turns (limit+1 probe) - Dedupe dangling operator retry after a completed pair without the broken preceding-role check --- hub/src/overseer/converseContext.test.ts | 37 ++++++++++++++++++- hub/src/overseer/converseContext.ts | 37 ++++++++++--------- .../settings/OverseerChatDebugControls.tsx | 36 ++++++++++-------- 3 files changed, 77 insertions(+), 33 deletions(-) diff --git a/hub/src/overseer/converseContext.test.ts b/hub/src/overseer/converseContext.test.ts index 061bf6d1c7..bee748ba23 100644 --- a/hub/src/overseer/converseContext.test.ts +++ b/hub/src/overseer/converseContext.test.ts @@ -98,7 +98,42 @@ describe('converseContext assembler', () => { const overseer = buildEngine(store).getOverseer() overseer.recordConvoTurn({ operatorText: 'first', overseerText: 'one', ts: 1 }) overseer.recordConvoTurn({ operatorText: 'second', overseerText: 'two', ts: 2 }) - const list = listRecentConvoTurns(overseer, { limit: 10 }) + const { turns: list, clippedByLimit } = listRecentConvoTurns(overseer, { limit: 10 }) expect(list.map((t) => t.operatorText)).toEqual(['first', 'second']) + expect(clippedByLimit).toBe(false) + }) + + it('reports truncated when store query clipped older turns', () => { + const store = new Store(':memory:') + const overseer = buildEngine(store).getOverseer() + for (let i = 0; i < 5; i++) { + overseer.recordConvoTurn({ + operatorText: `q${i}`, + overseerText: `a${i}`, + ts: 1000 + i + }) + } + const assembled = assembleOverseerConverseMessages({ + overseer, + clientMessages: [{ role: 'operator', content: 'now' }], + maxTurns: 2 + }) + expect(assembled.hydratedTurns).toBe(2) + expect(assembled.truncated).toBe(true) + expect(assembled.messages.map((m) => m.content)).toEqual(['q3', 'a3', 'q4', 'a4', 'now']) + }) + + it('dedupes dangling operator retry after a completed pair', () => { + const store = new Store(':memory:') + const overseer = buildEngine(store).getOverseer() + overseer.recordConvoTurn({ operatorText: 'done pair', overseerText: 'answered', ts: 1 }) + // Dangling operator-only row (empty overseer reply permitted by write path). + overseer.recordConvoTurn({ operatorText: 'retry me', overseerText: '', ts: 2 }) + const assembled = assembleOverseerConverseMessages({ + overseer, + clientMessages: [{ role: 'operator', content: 'retry me' }] + }) + expect(assembled.messages.filter((m) => m.content === 'retry me')).toHaveLength(1) + expect(assembled.messages.at(-1)).toEqual({ role: 'operator', content: 'retry me' }) }) }) diff --git a/hub/src/overseer/converseContext.ts b/hub/src/overseer/converseContext.ts index ec2d22ac06..f80d9576fc 100644 --- a/hub/src/overseer/converseContext.ts +++ b/hub/src/overseer/converseContext.ts @@ -91,23 +91,26 @@ export function eventToConvoTurnView(event: StoredSystemEvent): StoredConvoTurnV } } -/** Newest-first from the store; returned oldest-first for display / assemble. */ +/** Newest-first from the store; returned oldest-first for display / assemble. + * Fetches `limit + 1` so callers can detect that older history was clipped. */ export function listRecentConvoTurns( overseer: OverseerEntity, opts: { limit?: number } = {} -): StoredConvoTurnView[] { +): { turns: StoredConvoTurnView[]; clippedByLimit: boolean } { const limit = Math.min(Math.max(opts.limit ?? DEFAULT_CONVERSE_HISTORY_MAX_TURNS, 1), 50) const events = overseer.queryEvents({ eventType: OVERSEER_CONVO_TURN_EVENT_TYPE, - limit + limit: limit + 1 }) const views: StoredConvoTurnView[] = [] for (const event of events) { const view = eventToConvoTurnView(event) if (view) views.push(view) } + const clippedByLimit = views.length > limit + const kept = clippedByLimit ? views.slice(0, limit) : views // events.query is id DESC — reverse to chronological. - return views.reverse() + return { turns: kept.reverse(), clippedByLimit } } function turnsToMessages(turns: StoredConvoTurnView[]): OverseerConverseMessage[] { @@ -131,13 +134,13 @@ function messagesCharCount(messages: OverseerConverseMessage[]): number { */ export function budgetConvoTurns( turnsOldestFirst: StoredConvoTurnView[], - opts: { maxTurns?: number; maxChars?: number } = {} + opts: { maxTurns?: number; maxChars?: number; alreadyClipped?: boolean } = {} ): { turns: StoredConvoTurnView[]; truncated: boolean } { const maxTurns = opts.maxTurns ?? DEFAULT_CONVERSE_HISTORY_MAX_TURNS const maxChars = opts.maxChars ?? DEFAULT_CONVERSE_HISTORY_MAX_CHARS let kept = turnsOldestFirst - let truncated = false + let truncated = opts.alreadyClipped === true if (kept.length > maxTurns) { kept = kept.slice(kept.length - maxTurns) truncated = true @@ -169,21 +172,21 @@ export function assembleOverseerConverseMessages(params: { throw new Error('Last client message must be from the operator') } - const fetched = listRecentConvoTurns(overseer, { - limit: Math.max(maxTurns ?? DEFAULT_CONVERSE_HISTORY_MAX_TURNS, 1) + const max = Math.max(maxTurns ?? DEFAULT_CONVERSE_HISTORY_MAX_TURNS, 1) + const { turns: fetched, clippedByLimit } = listRecentConvoTurns(overseer, { limit: max }) + const { turns, truncated } = budgetConvoTurns(fetched, { + maxTurns: max, + maxChars, + alreadyClipped: clippedByLimit }) - const { turns, truncated } = budgetConvoTurns(fetched, { maxTurns, maxChars }) const history = turnsToMessages(turns) - // If the operator re-sent the exact last logged question without a reply yet - // (unlikely — we only record after reply), avoid duplicating. Normal path: - // current utterance is not in the store until after this turn completes. + // Dedup a dangling operator line (same text already last in history with no + // overseer reply after it). Do not require the previous message to be non-overseer — + // after a completed pair the prior message IS overseer, and a dangling operator + // turn is still the common retry shape. const lastHistory = history[history.length - 1] - if ( - lastHistory?.role === 'operator' - && lastHistory.content === latest.content - && (history.length < 2 || history[history.length - 2]?.role !== 'overseer') - ) { + if (lastHistory?.role === 'operator' && lastHistory.content === latest.content) { return { messages: history, hydratedTurns: turns.length, truncated } } diff --git a/web/src/components/settings/OverseerChatDebugControls.tsx b/web/src/components/settings/OverseerChatDebugControls.tsx index 863b7d3882..f02e610509 100644 --- a/web/src/components/settings/OverseerChatDebugControls.tsx +++ b/web/src/components/settings/OverseerChatDebugControls.tsx @@ -23,7 +23,7 @@ export function OverseerChatDebugControls() { const { api } = useAppContext() const [open, setOpen] = useState(false) const [turns, setTurns] = useState([]) - const [hydrated, setHydrated] = useState(false) + const [hydrating, setHydrating] = useState(false) const [input, setInput] = useState('') const [loading, setLoading] = useState(false) const [error, setError] = useState(null) @@ -35,6 +35,7 @@ export function OverseerChatDebugControls() { const [modelsError, setModelsError] = useState(null) const [selectedModel, setSelectedModel] = useState('') const scrollRef = useRef(null) + const hydrateGenRef = useRef(0) useEffect(() => { if (!open || !api || profiles.length > 0) return @@ -43,13 +44,16 @@ export function OverseerChatDebugControls() { .catch(() => { /* brains list is optional chrome */ }) }, [open, api, profiles.length]) - // Hydrate from hub-owned convo_turns once per open (reload-safe continuity). + // Rehydrate from hub every time the panel opens (other transports may have + // written turns while closed). Block send until this settles. useEffect(() => { - if (!open || !api || hydrated) return + if (!open || !api) return + const gen = ++hydrateGenRef.current let cancelled = false + setHydrating(true) void api.fetchOverseerConverseRecent(24) .then((res) => { - if (cancelled) return + if (cancelled || gen !== hydrateGenRef.current) return const next: ChatTurn[] = [] for (const turn of res.turns) { if (turn.operatorText.trim()) { @@ -68,16 +72,16 @@ export function OverseerChatDebugControls() { } } setTurns(next) - setHydrated(true) requestAnimationFrame(() => { scrollRef.current?.scrollTo({ top: scrollRef.current.scrollHeight }) }) }) - .catch(() => { - if (!cancelled) setHydrated(true) + .catch(() => { /* empty local view; hub still owns memory on send */ }) + .finally(() => { + if (!cancelled && gen === hydrateGenRef.current) setHydrating(false) }) return () => { cancelled = true } - }, [open, api, hydrated]) + }, [open, api]) // Populate the model dropdown live from the selected profile's endpoint // (server proxies GET /models so the api key never reaches the browser). @@ -99,10 +103,11 @@ export function OverseerChatDebugControls() { }, [open, api, selectedProfile]) const profileDefaultModel = profiles.find((p) => p.id === selectedProfile)?.model ?? null + const sendBlocked = loading || hydrating const send = useCallback(async (text: string) => { const trimmed = text.trim() - if (!trimmed || !api || loading) return + if (!trimmed || !api || sendBlocked) return setError(null) setInput('') @@ -130,7 +135,7 @@ export function OverseerChatDebugControls() { } finally { setLoading(false) } - }, [api, loading, selectedProfile, selectedModel]) + }, [api, sendBlocked, selectedProfile, selectedModel]) return (
@@ -201,7 +206,7 @@ export function OverseerChatDebugControls() {
)) )} + {hydrating ?

Loading hub thread…

: null} {loading ?

Overseer is thinking…

: null}
@@ -258,13 +264,13 @@ export function OverseerChatDebugControls() { type="text" value={input} onChange={(e) => setInput(e.target.value)} - placeholder="Ask the Overseer…" - disabled={loading} + placeholder={hydrating ? 'Loading thread…' : 'Ask the Overseer…'} + disabled={sendBlocked} className="flex-1 rounded-md border border-[var(--app-border)] bg-[var(--app-bg)] px-2 py-1.5 text-[13px] text-[var(--app-fg)] disabled:opacity-50" />