diff --git a/AGENTS.md b/AGENTS.md index 96e0a4d..8dab051 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -37,7 +37,8 @@ duck. Check direct inspector controls, the connect box, and stale-feed recovery. | `src/interaction.mjs` / `src/conversation.mjs` | Door crossings and browser message capability checks | | `src/feed.mjs` | HTTP server, snapshot cache, adapter integration, activity routes | | `src/transcripts.mjs` / `src/activity.mjs` | Incremental transcript parsing and public activity; optional Hub timeline | -| `src/hub.mjs` | Optional readonly sqlite `agent_runs` adapter | +| `src/hub.mjs` | Optional standalone readonly sqlite `agent_runs` adapter | +| `src/hub-inventory.mjs` | Shared `/v1/tasks` inventory when `AUTOHUB_HUB_BASE` is configured; replaces local discovery | | `src/todos.mjs` / `src/messages.mjs` | Explicit checklist snapshots and user-submitted AutoHub messages | | `src/pr.mjs` / `src/github.mjs` | Shared PR classifier/counts and cached read-only GitHub observations | | `src/occupancy.mjs` | live / recent / settled, letters, conservative PR identity inference | diff --git a/README.md b/README.md index b655bdf..163b6de 100644 --- a/README.md +++ b/README.md @@ -130,6 +130,35 @@ Optional fields add task identity, the original request, timestamps, activity, w The local process can also populate `/agents` without you writing a server: +For a shared AutoHub inventory, start the feed with: + +```sh +AUTOHUB_HUB_BASE=http://127.0.0.1:8767 npm start +``` + +For a persistent setting in this checkout, add this line to its local `.env`: + +```dotenv +AUTOHUB_HUB_BASE=http://127.0.0.1:8767 +``` + +Then use plain `npm start`. Both `npm start` and `npm run once` load `.env` +when present; existing shell variables take precedence. The file is ignored by +Git, and no global shell export is needed. With no Hub setting, standalone +behavior remains the default. Direct `node src/feed.mjs` launches use the +process environment; add `--env-file-if-exists=.env` before the script path +when launching that way. + +Set `AUTOHUB_HUB_TOKEN` in `.env` or the server environment if that Hub requires a bearer +token. In this mode CottageCode follows every `/v1/tasks` cursor page and reads +canonical task detail, timelines, checklist data, and control capabilities. +It does not open SQLite or scan Claude transcripts. The Hub owns discovery; +CottageCode keeps its existing towns, cottages, inspectors, and message receipts. +If the Hub is unavailable, the last complete feed stays visible as stale. +Credentials and raw upstream errors are never returned to the browser. + +Without `AUTOHUB_HUB_BASE`, the standalone adapters remain: + 1. Readonly SQLite at `AGENT_DB_PATH` (an `agent_runs` table), if set 2. Claude Code session JSONL under `~/.claude/projects` (or `CLAUDE_PROJECTS_DIR`) 3. Cached, read-only GitHub metadata through an authenticated `gh` CLI diff --git a/docs/COTTAGE_FEED.md b/docs/COTTAGE_FEED.md index 051a490..2ab0569 100644 --- a/docs/COTTAGE_FEED.md +++ b/docs/COTTAGE_FEED.md @@ -354,6 +354,8 @@ Replay steps through recorded milestones and highlights their cottages. It does | Variable | Default | Behavior | |---|---|---| | `CLAUDE_PROJECTS_DIR` | `~/.claude/projects` | Directory scanned for local session JSONL. | +| `AUTOHUB_HUB_BASE` | Unset | Select shared Hub inventory mode, e.g. `http://127.0.0.1:8767` (optional `/v1`). Replaces local SQLite inventory and Claude scanning; follows all `/v1/tasks` pages and uses canonical ids for detail/activity/control. | +| `AUTOHUB_HUB_TOKEN` | Unset | Optional server-side bearer token for shared inventory, timelines, and explicit user messages. Never published. Falls back to `COTTAGE_HUB_TOKEN` if absent. | | `AGENT_DB_PATH` | Unset | Optional readonly SQLite database containing `agent_runs`. | | `AGENT_STALE_THRESHOLD_MS` | `900000` | Hub running-task inactivity threshold, in milliseconds. | | `COTTAGE_GITHUB` | Enabled | Set to `0` to disable GitHub enrichment. | @@ -361,6 +363,11 @@ Replay steps through recorded milestones and highlights their cottages. It does | `COTTAGE_HUB_TOKEN` | Unset | Optional server-side bearer token for Hub requests. Never returned to the browser. | | `COTTAGE_MESSAGE_LEDGER` | `~/.cottagecode/message-receipts.jsonl` | Private local request hashes and delivery outcomes, used to prevent duplicate writes. | +`npm start` and `npm run once` optionally load the checkout’s Git-ignored `.env`. +A local `AUTOHUB_HUB_BASE=http://127.0.0.1:8767` setting persists Hub mode across +launches without a global shell export. Existing process variables take precedence, +and an absent `.env` preserves standalone defaults. + GitHub enrichment uses the authenticated local `gh` CLI with read-only `pr view`, `pr list`, and `repo view` calls. It accepts an explicit GitHub.com PR link or exact repository/number. Branch discovery requires an exact repository and a non-default branch with one unambiguous same-repository PR match. The local adapter can resolve `owner/repo` from a worktree's GitHub `origin`; town names never supply repository identity. GitHub requests run in the background, at most 3 at a time, with a 30-second refresh interval and failure backoff from 30 seconds to 5 minutes. Failed refreshes retain old observations marked stale. Cold snapshots can contain unknown PRs while those queries finish. Existing custom JSON feeds provide their own PR metadata; the browser does not run `gh` against them. @@ -437,3 +444,20 @@ npx --yes http-server . -p 9999 --cors ``` **Pause feed** stops browser refreshes and the demo simulator. It never writes back to the feed or pauses real agents. Walking and direct inspector controls remain available. + +### Shared Hub inventory mode + +The Node process caches successful inventory reads for five seconds; it follows +all cursor pages before replacing a snapshot. Failed pages keep the previous +complete snapshot and mark the feed stale. `id`, `parent`, and `taskId` retain +Hub canonical task identities, including observed external sessions. Their +`sessionId` stays separate. Town labels still derive from recorded project +paths; labels and worktree slugs never determine task identity or control. + +The existing local activity route refreshes `/v1/tasks/:id` and reads the Hub +`/timeline` and `/todo` routes. Source `thinking` events remain excluded from +public journals. Explicit Send rechecks detail and the live `controlTargetId` +and capabilities; a null target is read-only. Replies also use Hub's optional +`expected_attention_version` guard when supplied by detail. The existing +receipt ledger and uncertain-delivery behavior remain in effect. No new cancel +or autonomous messaging behavior is added. diff --git a/package.json b/package.json index e2d2af9..b8a490b 100644 --- a/package.json +++ b/package.json @@ -27,8 +27,8 @@ "cottagecode": "./src/feed.mjs" }, "scripts": { - "start": "node src/feed.mjs", - "once": "node src/feed.mjs --once", + "start": "node --env-file-if-exists=.env src/feed.mjs", + "once": "node --env-file-if-exists=.env src/feed.mjs --once", "build:demo": "node scripts/build-demo.mjs", "deploy:demo": "node scripts/deploy-demo.mjs", "test": "node --test", diff --git a/src/activity.mjs b/src/activity.mjs index 24a3a5c..7131a9b 100644 --- a/src/activity.mjs +++ b/src/activity.mjs @@ -107,8 +107,8 @@ export function pageActivity(events, { after, before, limit = 100, source = "non /** AutoHub only supports backwards cursors, so retain safe events for forward polls. */ export function createHubTimelineReader({ - baseUrl = process.env.COTTAGE_HUB_URL || "", - token = process.env.COTTAGE_HUB_TOKEN || "", + baseUrl = process.env.AUTOHUB_HUB_BASE || process.env.COTTAGE_HUB_URL || "", + token = process.env.AUTOHUB_HUB_TOKEN || process.env.COTTAGE_HUB_TOKEN || "", fetchFn = globalThis.fetch, now = Date.now, ttl = 1500, @@ -162,6 +162,9 @@ export function createHubTimelineReader({ const old = cache.get(taskId); if (old && !options.before && now() - old.checkedAt < ttl) return old; const entry = old ? { ...old } : { events: [], hasOlder: false, checkedAt: 0, source: "none" }; + // Historical pages are request-local windows; they must neither compete + // with the recent-event cap nor evict the cursor used by forward polling. + if (options.before) entry.events = []; let before = options.before || ""; try { for (let page = 0; page < 4; page++) { @@ -191,8 +194,10 @@ export function createHubTimelineReader({ // Never copy fetch errors: their messages may include deployment credentials. entry.checkedAt = now(); } - cache.set(taskId, entry); - if (cache.size > 200) cache.delete(cache.keys().next().value); + if (!options.before) { + cache.set(taskId, entry); + if (cache.size > 200) cache.delete(cache.keys().next().value); + } return entry; } @@ -204,9 +209,14 @@ export function createHubTimelineReader({ const key = `${taskId}\n${options.before || ""}`; if (!inFlight.has(key)) inFlight.set(key, refresh(taskId, options).finally(() => inFlight.delete(key))); const [entry, checkpointTodos] = await Promise.all([inFlight.get(key), readTodos(taskId)]); - const result = pageActivity(entry.events, { ...options, source: entry.source }); + // The upstream query already applied `before`, so its anchor is outside + // this window. Apply only the local display limit to that result. + const pageOptions = options.before ? { ...options, before: undefined, after: undefined } : options; + const result = pageActivity(entry.events, { ...pageOptions, source: entry.source }); + if (options.before) result.cursor = result.events[0]?.id || options.before; if (!options.after && entry.hasOlder) result.hasMore = true; - let todos = latestTodos(entry.events, { source: entry.source }); + const todoEntry = options.before ? cache.get(taskId) || entry : entry; + let todos = latestTodos(todoEntry.events, { source: todoEntry.source }); if (todos && entry.stale) todos = { ...todos, stale: true }; const checkpointWins = checkpointTodos && (!todos || (checkpointTodos.updatedAt !== null diff --git a/src/feed.mjs b/src/feed.mjs index ed90262..b00d631 100644 --- a/src/feed.mjs +++ b/src/feed.mjs @@ -10,6 +10,7 @@ import { join, dirname, resolve, sep } from "node:path"; import { fileURLToPath } from "node:url"; import { repoOf, worktreeOf, townName } from "./towns.mjs"; import { readHubAgents, shortModel } from "./hub.mjs"; +import { createHubInventoryReader } from "./hub-inventory.mjs"; import { classifyOccupancy, inferPr, lettersOf, liveCostOf, stampOccupancy } from "./occupancy.mjs"; import { createTranscriptReader } from "./transcripts.mjs"; import { createHubTimelineReader, pageActivity, mergeActivityEvents } from "./activity.mjs"; @@ -219,8 +220,13 @@ function newestTodos(current, incoming) { export function createFeed({ scanClaude = createClaudeScanner(), readHub = readHubAgents, enrich = enrichAgents, + hubInventory = createHubInventoryReader(), resolveRepos = createRepoResolver(), timeline = createHubTimelineReader(), now = Date.now, } = {}) { + if (hubInventory.configured) { + scanClaude = async () => ({ok:true,agents:[],sessions:[]}); + readHub = () => hubInventory.read(); + } let cache = []; let claude = { agents: [], sessions: [] }; let hub = { agents: [], keys: new Set(), links: new Map() }; @@ -240,7 +246,7 @@ export function createFeed({ if (claude.stale) nextErrors.push(claude.error || "Some transcripts are temporarily unavailable"); } else nextErrors.push("Local transcripts are temporarily unavailable"); if (results[1].status === "fulfilled" && results[1].value.ok !== false) hub = results[1].value; - else nextErrors.push("Hub database is temporarily unavailable"); + else nextErrors.push(hubInventory.configured ? "Hub inventory is temporarily unavailable" : "Hub database is temporarily unavailable"); const localById = new Map(claude.agents.map(agent => [agent.id, agent])); const previousById = new Map(cache.map(agent => [agent.id, agent])); @@ -418,21 +424,33 @@ export function createFeed({ return doScan(); }, async getActivity(id, options = {}) { - const agent = cache.find(cottage => cottage.id === id); + let agent = cache.find(cottage => cottage.id === id); if (!agent) return null; + let refreshedDetail = false, detailFailed = false; + const activityResponse = page => detailFailed ? { + ...page, stale: true, error: "Hub task detail temporarily unavailable", + inputRequest: page.inputRequest ? {...page.inputRequest,stale:true} : null, + ...(page.cursorReset && !page.events?.length ? {cursorReset:false,cursor:options.after || options.before || null} : {}), + } : page; + const inputForActivity = () => refreshedDetail ? normalizeInputRequest(agent.inputRequest) : currentInputRequest(agent); + if (hubInventory.configured && typeof hubInventory.detail === "function") { + // Detail serves this activity request; only scans replace enriched inventory. + try { agent = { ...agent, ...await hubInventory.detail(id) }; refreshedDetail = true; } + catch { detailFailed = true; } + } const local = activity.get(id) || []; if (local.length) { const todos = agent.source === "hub" && timeline.configured && typeof timeline.readTodos === "function" ? rememberActivityTodos(agent, await timeline.readTodos(agent.id)) : normalizeTodos(agent.todos); - return { ...pageActivity(mergeActivityEvents([], local), { ...options, source: "claude-transcript" }), todos, inputRequest: currentInputRequest(agent), ...(stale ? { stale: true } : {}) }; + return activityResponse({ ...pageActivity(mergeActivityEvents([], local), { ...options, source: "claude-transcript" }), todos, inputRequest: inputForActivity(), ...(stale ? { stale: true } : {}) }); } if (agent.source === "hub" && timeline.configured) { const result = await timeline.read(agent.id, options); const todos = rememberActivityTodos(agent, result.todos); - return { ...result, todos, inputRequest: currentInputRequest(agent) }; + return activityResponse({ ...result, todos, inputRequest: inputForActivity() }); } - return { ...pageActivity([], options), todos: normalizeTodos(agent.todos), inputRequest: currentInputRequest(agent), unavailable: true }; + return activityResponse({ ...pageActivity([], options), todos: normalizeTodos(agent.todos), inputRequest: inputForActivity(), unavailable: true }); }, }; } diff --git a/src/hub-inventory.mjs b/src/hub-inventory.mjs new file mode 100644 index 0000000..281e5c0 --- /dev/null +++ b/src/hub-inventory.mjs @@ -0,0 +1,112 @@ +/** Canonical AutoHub task inventory, adapted to the existing cottage contract. + * Only this server holds the optional bearer token; it is never in feed JSON. */ +import {toCottage} from './hub.mjs'; +import {timestampMs,safeActivityUrl} from './activity.mjs'; +import {taskText} from './task-text.mjs'; +import {classifyOccupancy} from './occupancy.mjs'; + +export function toInventoryCottage(task, now=Date.now()) { + const context=typeof task.context==='object' && task.context ? task.context : {}; + const resultLinks=Array.isArray(task.resultLinks) ? task.resultLinks.flatMap(link=>{ + const url=safeActivityUrl(link?.url); + return url?[{...link,url}]:[]; + }) : []; + const prLink=resultLinks.find(link=>link.kind==='pull_request'); + const cottage=toCottage({...task, + model:task.model || context.agentKernel?.route?.model || + (task.backend==='codex'?'openai':task.provider || task.backend || 'unknown'), + record_kind:task.recordKind,parent_id:task.parentId,session_id:task.sessionId, + status:task.normalizedStatus==='completed_without_report'?'completed':task.normalizedStatus || task.status, + queued_at:task.queuedAt,started_at:task.startedAt,completed_at:task.completedAt,updated_at:task.updatedAt, + input_tokens:task.inputTokens,output_tokens:task.outputTokens,cache_write_tokens:task.cacheWriteTokens,cache_read_tokens:task.cacheReadTokens,total_cost:task.totalCost, + result_summary:task.resultSummary || task.displayResult, + context:{...context,projectPath:task.projectPath || context.projectPath,workFolder:task.worktreePath || task.projectPath, + branch:task.branch,originalAsk:task.originalRequest, + ...(prLink ? {pr:{url:prLink.url,state:'unknown',source:'context'}} : {})}, + },now); + const request=taskText(task.originalRequest); + cottage.taskId=task.id; + cottage.originalAsk=request.slice(0,32768); + cottage.originalAskTruncated=request.length>32768; + cottage.originalAskSource=task.recordKind==='external_session'?'session':'task'; + cottage.task=(request || taskText(task.title || task.task) || 'Task').slice(0,150); + cottage.activity=task.currentActivity || task.currentStep || cottage.activity; + cottage.lastLine=cottage.activity || cottage.result || cottage.task; + cottage.worktreePath=task.worktreePath || task.projectPath || ''; + cottage.branch=task.branch || ''; + cottage.updatedAt=timestampMs(task.updatedAt) || cottage.updatedAt; + cottage.inventorySource=task.source || null; + cottage.dispatchedBy=task.provenance?.caller || task.provenance?.platform || task.platform || task.provider || 'hub'; + cottage.freshness=task.freshness || {lastSeenAt:task.lastSeenAt || null,isStale:task.isStale===true}; + cottage.resultLinks=resultLinks.map(({kind,url})=>({kind,url})); + cottage.artifacts=resultLinks.map(link=>({url:link.url,title:taskText(link.label || link.title) || (link.kind==='pull_request'?'Pull request':'Result')})); + cottage.conversationTarget={...cottage.conversationTarget, + taskId:task.id,taskStatus:task.normalizedStatus || task.status,recordKind:task.recordKind, + controlTargetId:typeof task.controlTargetId==='string'?task.controlTargetId:null, + capabilities:{canRespond:task.capabilities?.canRespond===true,canSteer:task.capabilities?.canSteer===true}, + }; + const canonicalFresh = task.freshness?.isStale===false || task.isStale===false; + if((task.normalizedStatus || task.status)==='running' && canonicalFresh && !cottage.inputRequest) cottage.status='working'; + if(task.isStale || cottage.freshness.isStale) { + if(cottage.status==='working')cottage.status='offline'; + if(cottage.inputRequest)cottage.inputRequest={...cottage.inputRequest,stale:true}; + } + cottage.occupancy=classifyOccupancy(cottage,now); + return cottage; +} + +export function createHubInventoryReader({ + baseUrl=process.env.AUTOHUB_HUB_BASE || '',token=process.env.AUTOHUB_HUB_TOKEN || process.env.COTTAGE_HUB_TOKEN || '', + fetchFn=globalThis.fetch,now=Date.now,ttl=5000,pageSize=250, +}={}) { + let base=null; + if(baseUrl) { + try { + base=new URL(baseUrl); + if(!['http:','https:'].includes(base.protocol)||base.username||base.password)throw new Error(); + base.pathname=base.pathname.replace(/\/$/,'').replace(/\/v1$/,'')+'/v1/';base.search='';base.hash=''; + } catch {throw new Error('AUTOHUB_HUB_BASE must be an HTTP(S) URL without credentials.');} + } + let previous=null,checkedAt=0,inFlight=null; + async function get(path) { + const response=await fetchFn(new URL(path,base),{headers:{accept:'application/json',...(token?{authorization:`Bearer ${token}`}:{})},redirect:'error',signal:AbortSignal.timeout(8000)}); + if(!response.ok)throw new Error('Hub inventory temporarily unavailable'); + return response.json(); + } + async function refresh() { + try { + const tasks=new Map(),seen=new Set();let cursor=''; + do { + const query=new URLSearchParams({scope:'history',limit:String(pageSize)}); + if(cursor)query.set('cursor',cursor); + const page=await get(`tasks?${query}`); + if(!Array.isArray(page.tasks))throw new Error(); + for(const task of page.tasks) { + if(typeof task?.id!=='string'||!task.id)throw new Error(); + tasks.set(task.id,toInventoryCottage(task,now())); + } + cursor=page.has_more ? page.next_cursor : ''; + if(page.has_more && (typeof cursor!=='string'||!cursor||seen.has(cursor)))throw new Error(); + if(cursor)seen.add(cursor); + }while(cursor); + previous={ok:true,agents:[...tasks.values()],keys:new Set(tasks.keys()),links:new Map([...tasks.keys()].map(id=>[id,new Set([id])]))}; + checkedAt=now();return previous; + }catch {throw new Error('Hub inventory temporarily unavailable');} + } + return {configured:Boolean(base), + async read(){ + if(!base)return {ok:true,missing:true,agents:[],keys:new Set(),links:new Map()}; + if(previous && now()-checkedAt{inFlight=null;}); + return inFlight; + }, + async detail(id){ + if(!base||!id)return null; + try { + const task=await get(`tasks/${encodeURIComponent(id)}`); + if(task.id!==id)throw new Error(); + return toInventoryCottage(task,now()); + }catch {throw new Error('Hub task detail temporarily unavailable');} + }, + }; +} diff --git a/src/messages.mjs b/src/messages.mjs index 0de5be9..0774ccb 100644 --- a/src/messages.mjs +++ b/src/messages.mjs @@ -27,7 +27,7 @@ export function acceptsMessageOrigin(req){ } export function createHubMessenger({ - baseUrl=process.env.COTTAGE_HUB_URL||'',token=process.env.COTTAGE_HUB_TOKEN||'', + baseUrl=process.env.AUTOHUB_HUB_BASE||process.env.COTTAGE_HUB_URL||'',token=process.env.AUTOHUB_HUB_TOKEN||process.env.COTTAGE_HUB_TOKEN||'', fetchImpl=globalThis.fetch,now=Date.now,timeoutMs=8000, ledgerPath=process.env.COTTAGE_MESSAGE_LEDGER||join(homedir(),'.cottagecode','message-receipts.jsonl'), }={}){ @@ -47,18 +47,20 @@ export function createHubMessenger({ function capability(agent,{stale=false,checkedAt=now()}={}){ const unavailable=reason=>({available:false,reason,source:'autohub',checkedAt}); if(!base)return unavailable('Messaging needs a configured AutoHub connection. This source currently provides activity only.'); - if(stale||!checkedAt||now()-checkedAt>120000)return unavailable('The task feed is stale. Reconnect before sending.'); + if(stale||agent.isStale||agent.freshness?.isStale||!checkedAt||now()-checkedAt>120000)return unavailable('The task feed is stale. Reconnect before sending.'); if(agent.inputRequest?.stale)return unavailable('The input request is stale. Refresh before replying.'); const target=agent.conversationTarget; - if(agent.source!=='hub'||target?.recordKind!=='logical_task'||!agent.taskId||target.taskId!==agent.taskId)return unavailable('This observed session has no supported task messaging route.'); + const canonical=target && Object.hasOwn(target,'controlTargetId'); + if(agent.source!=='hub'||!agent.taskId||target?.taskId!==agent.taskId || + (canonical ? target.controlTargetId!==agent.taskId : target.recordKind!=='logical_task'))return unavailable('This observed session has no supported task messaging route.'); const shownInput=normalizeInputRequest(agent.inputRequest); let mode=''; - if(['awaiting_input','needs_input'].includes(target.taskStatus))mode='respond'; + if(['awaiting_input','needs_input'].includes(target.taskStatus) && (!canonical || target.capabilities?.canRespond===true))mode='respond'; // A transitional Hub row can expose an explicit question before its raw // status flips from running. Do not turn that answer into a free-form // redirect: wait until Hub offers the matching response route. else if(shownInput)return unavailable('This task has a pending input request. Wait for AutoHub to expose its response route.'); - else if(target.taskStatus==='running'&&target.transport==='tmux'&&target.supportsRedirection===true)mode='redirect'; + else if(target.taskStatus==='running'&&(canonical ? target.capabilities?.canSteer===true : target.transport==='tmux'&&target.supportsRedirection===true))mode='redirect'; else return unavailable(['completed','failed','cancelled','interrupted'].includes(target.taskStatus)?'This task has finished. Start any follow-up in its original workflow.':target.transport==='direct'?'This direct session does not support mid-task messages.':'This task has no supported live message route.'); const resolution=agent.inputRequestResolution; const newerThanUnidentifiedResolution=!resolution?.id&&shownInput?.updatedAt&&resolution?.resolvedAt&&shownInput.updatedAt>resolution.resolvedAt; @@ -117,30 +119,33 @@ export function createHubMessenger({ if(supported.mode==='respond'&&(shownInput?.id||null)!==(payload.inputRequestId||null))return rejected(409,'The input request changed. Reopen the current question before replying.'); if(supported.mode==='respond'&&inputRequestVersion(shownInput)!==(payload.inputRequestVersion||null))return rejected(409,'The question or its choices changed. Refresh before replying.'); const path='tasks/'+encodeURIComponent(agent.taskId); + let expectedAttentionVersion; try{ // Native Hub question rounds live in orchestrator context. Request it // only server-side so identical wording in a new round stays distinct. const response=await request(path+'?context=raw'); if(!response.ok)return rejected(502,'AutoHub could not verify the current task. Nothing was sent.'); const task=await response.json(),context=parse(task.context),execution=context.lifecycle?.execution||{}; - const current={...agent,conversationTarget:{taskId:task.id,taskStatus:task.status,recordKind:task.recordKind,transport:execution.sessionMode,supportsRedirection:execution.supportsRedirection}}; - if(task.id!==agent.taskId||task.isStale||task.archived)return rejected(409,'The task is no longer available for messages.'); + const canonical=Object.hasOwn(agent.conversationTarget,'controlTargetId'); + const current={...agent,inputRequest:inputRequestFromHub(task),freshness:task.freshness,isStale:task.isStale,conversationTarget:{taskId:task.id,taskStatus:task.normalizedStatus || task.status,recordKind:task.recordKind,transport:execution.sessionMode,supportsRedirection:execution.supportsRedirection, ...(canonical ? {controlTargetId:task.controlTargetId ?? null,capabilities:task.capabilities || {}} : {})}}; + if(task.id!==agent.taskId||task.isStale||task.freshness?.isStale||task.archived)return rejected(409,'The task is no longer available for messages.'); const verified=capability(current); if(!verified.available||verified.mode!==supported.mode|| (supported.mode==='respond'&&task.canRespond===false))return rejected(409,'The task’s input state changed. Refresh before sending.'); if(supported.mode==='respond'&&inputIdentity(inputRequestFromHub(task))!==inputIdentity(shownInput))return rejected(409,'The agent is now asking a different question. Refresh before replying.'); + if(typeof task.attentionVersion==='string')expectedAttentionVersion=task.attentionVersion; }catch{return rejected(502,'AutoHub could not verify the current task. Nothing was sent.');} const entry={id:payload.requestId,hash,at:now(),response:null}; try{await remember(entry);}catch{return rejected(503,'The message receipt could not be saved. Nothing was sent.');} let receipt; try{ - const response=await request(path+'/'+supported.mode,{method:'POST',headers:{'content-type':'application/json'},body:JSON.stringify(supported.mode==='redirect'?{instruction:message}:{response:message})}); + const response=await request(path+'/'+supported.mode,{method:'POST',headers:{'content-type':'application/json'},body:JSON.stringify(supported.mode==='redirect'?{instruction:message}:{response:message,...(expectedAttentionVersion?{expected_attention_version:expectedAttentionVersion}:{})})}); const body=await response.json(); if(response.ok&&((supported.mode==='redirect'&&body.ok===true&&body.taskId===agent.taskId)||(supported.mode==='respond'&&body.id===agent.taskId))){ receipt=result(200,{ok:true,delivery:supported.mode==='redirect'?'submitted':'accepted',requestId:payload.requestId,source:'autohub',timestamp:now()}); }else if(response.status===401 || (response.status===403&&body.error==='owner_approval_required') || (response.status===404&&['task_not_found','Task not found'].includes(body.error)) || - (response.status===409&&['no_running_session','not_awaiting_input'].includes(body.error)) || + (response.status===409&&['no_running_session','not_awaiting_input','attention_changed'].includes(body.error)) || (response.status===400&&body.error==='instruction_required')){ // Only known pre-delivery errors are safe to retry. AutoHub can return // already_answered/cancellation_in_progress after affecting the task. diff --git a/src/observatory.mjs b/src/observatory.mjs index febbb6e..b774a97 100644 --- a/src/observatory.mjs +++ b/src/observatory.mjs @@ -39,17 +39,39 @@ export function appendInlineHandoffs(events,append){ } export function applyActivityPage(cache,data,{older=false}={}){ if(!Array.isArray(data?.events))throw new Error('Activity response has no events'); + if(older&&!cache.browsingHistory){ + cache.recentEvents=cache.events.slice(-1500);cache.recentHasMore=cache.hasMore;cache.browsingHistory=true; + } + let current=!older&&cache.browsingHistory?cache.recentEvents:cache.events; if(data.cursorReset){ - if(!older)cache.events=[]; - cache.hasMore=!!data.hasMore;cache.warning='An activity cursor expired. Only retained source history is available.'; + if(!older)current=[]; + cache.warning='An activity cursor expired. Only retained source history is available.'; } - cache.events=(older?mergeActivity(data.events,cache.events):mergeActivity(cache.events,data.events)).slice(-1500); + const merged=older?mergeActivity(data.events,current):mergeActivity(current,data.events); + const events=older?merged.slice(0,1500):merged.slice(-1500); + if(!older&&cache.browsingHistory)cache.recentEvents=events; + else cache.events=events; cache.source=String(data.source||'feed'); - if(older||!cache.cursor)cache.hasMore=!!data.hasMore; - if(!older)cache.cursor=data.cursor||cache.events.at(-1)?.id||cache.cursor; + if(older)cache.hasMore=!!data.hasMore; + else if(!cache.cursor||data.cursorReset){ + if(cache.browsingHistory)cache.recentHasMore=!!data.hasMore; + else cache.hasMore=!!data.hasMore; + } + if(!older)cache.cursor=data.cursor||events.at(-1)?.id||cache.cursor; cache.stale=!!data.stale;cache.error=data.error||'';cache.unavailable=!!data.unavailable; return cache; } +export function returnToLatestActivity(cache){ + if(!cache.browsingHistory)return cache; + cache.events=cache.recentEvents;cache.hasMore=cache.recentHasMore; + cache.browsingHistory=false;cache.activityViewVersion=(cache.activityViewVersion||0)+1; + delete cache.recentEvents;delete cache.recentHasMore; + return cache; +} +export function activityPagingButtons(cache){ + return button('older','Earlier entries',cache.hasMore?'':'disabled')+ + (cache.browsingHistory?button('latest','Latest entries'):''); +} export function activityJournalPresentation(cache={}){ const source=String(cache.source||'feed'); if(cache.unavailable)return {state:'unavailable',text:source+' · activity unavailable from this source'}; @@ -339,18 +361,18 @@ export function createObservatory(api){ const url=activityAddress(a,api.getEndpoint());if(!url||inflight.has(a.id))return; if(older&&cache.events.length)url.searchParams.set('before',cache.events[0].id); else if(cache.events.length)url.searchParams.set('after',cache.cursor||cache.events.at(-1).id); - const controller=new AbortController();inflight.set(a.id,controller);const epoch=sourceKey; + const controller=new AbortController();inflight.set(a.id,controller);const epoch=sourceKey,viewVersion=cache.activityViewVersion||0; const timeout=setTimeout(()=>controller.abort(),7000); try{ const res=await fetch(url,{signal:controller.signal,headers:{accept:'application/json'}}); if(!res.ok)throw new Error('Activity HTTP '+res.status); const data=await res.json(); - if(epoch!==sourceKey||activity.get(a.id)!==cache)return; + if(epoch!==sourceKey||activity.get(a.id)!==cache||(older&&viewVersion!==(cache.activityViewVersion||0)))return; applyActivityPage(cache,data,{older}); if(Object.hasOwn(data,'todos'))cache.todos=normalizeTodos(data.todos); recordActivity(a,data.events); for(const e of data.events)if(e.kind==='handoff'&&e.from&&e.to)appendHandoff(e); - }catch(err){if(epoch===sourceKey&&activity.get(a.id)===cache){cache.stale=true;cache.error=err.name==='AbortError'?'Activity request timed out':err.message;}} + }catch(err){if(epoch===sourceKey&&activity.get(a.id)===cache&&(!older||viewVersion===(cache.activityViewVersion||0))){cache.stale=true;cache.error=err.name==='AbortError'?'Activity request timed out':err.message;}} finally{clearTimeout(timeout);if(inflight.get(a.id)===controller)inflight.delete(a.id);if((selected===a.id||interiorId===a.id)&&epoch===sourceKey)renderPanel(selected);} } function eventHtml(events){ @@ -358,7 +380,7 @@ export function createObservatory(api){ return events.map(e=>'
  • '+esc(e.kind||'progress')+'

    '+esc(e.text)+'

    '+link(e.url,'Open artifact')+'
  • ').join(''); } function latestLine(a){ - const events=activity.get(a.id)?.events||a.events||[]; + const cache=activity.get(a.id),events=cache?.recentEvents||cache?.events||a.events||[]; return [...events].reverse().find(e=>['progress','summary','tool','result'].includes(e.kind))?.text||a.activity||a.lastLine||''; } function todosHtml(a,cache){ @@ -384,7 +406,7 @@ export function createObservatory(api){ return '

    A word with '+esc(a.name)+'

    '+(sound.enabled?'Your host answers with a little murmur.':'Enable sound to hear your host’s little murmur.')+' Activity below comes from the task’s recorded updates.

    '+input+composer+ (receipts?'
      '+receipts+'
    ':'')+ '
    Agent’s to-do list'+todosHtml(a,cache)+'
    '+ - '

    From the workbench

    '+button('older','Earlier entries',cache.hasMore?'':'disabled')+'

    '+esc(cache.source)+(cache.stale?' · stale — '+esc(cache.error||'connection interrupted'):' · recorded activity')+'

      '+eventHtml(cache.events)+'
    '; + '

    From the workbench

    '+activityPagingButtons(cache)+'

    '+esc(cache.source)+(cache.stale?' · stale — '+esc(cache.error||'connection interrupted'):' · recorded activity')+'

      '+eventHtml(cache.events)+'
    '; } async function sendMessage(){ const a=byId(interiorId||selected);if(!a)return; @@ -430,7 +452,7 @@ export function createObservatory(api){ content='

    Pinned request

    '+(a.originalAskSource==='session'?'First request recorded in this session. Current task boundaries are unavailable.':'Original request supplied by the task source.')+'

    '+esc(a.originalAsk||'The feed has not supplied the original request.')+'
    '; }else if(tab==='journal'){ const presentation=activityJournalPresentation(cache); - content='

    At the workbench

    '+button('older','Earlier entries',cache.hasMore?'':'disabled')+'

    '+esc(presentation.text)+'

    '+(cache.warning?'

    '+esc(cache.warning)+'

    ':'')+(cache.unavailable?'

    This source has not made a task journal available.

    ':'
      '+eventHtml(cache.events)+'
    '); + content='

    At the workbench

    '+activityPagingButtons(cache)+'

    '+esc(presentation.text)+'

    '+(cache.warning?'

    '+esc(cache.warning)+'

    ':'')+(cache.unavailable?'

    This source has not made a task journal available.

    ':'
      '+eventHtml(cache.events)+'
    '); }else if(tab==='review'){ content='

    Review desk

    '+prSnapshotHtml(pr)+'

    The parcel opens the PR. Review and merge stay in your existing workflow.

    '; }else if(tab==='artifacts'){ @@ -491,6 +513,7 @@ export function createObservatory(api){ } else if(action.startsWith('tab:')){tab=action.slice(4);selectedObject=({request:'request',journal:'workbench',review:'review',artifacts:'shelf',overview:'clock'})[tab];renderPanel(selected);if(['journal','todos'].includes(tab))ensureActivity(byId(interiorId||selected));} else if(action==='older')await ensureActivity(byId(interiorId||selected),{older:true}); + else if(action==='latest'){const a=byId(interiorId||selected);if(a){returnToLatestActivity(cacheFor(a));renderPanel(selected);const log=panel.querySelector('.journal');if(log)log.scrollTop=log.scrollHeight;}} else if(action==='copy'){try{await navigator.clipboard.writeText(byId(interiorId||selected).worktreePath);e.target.textContent='Copied';}catch{e.target.textContent='Copy unavailable';}} else if(action==='handoff'){const url=safeHttpsUrl(byId(interiorId||selected)?.handoffUrl);if(url)window.open(url,'_blank','noopener,noreferrer');} else if(action.startsWith('jump:'))focusCottage(action.slice(5)); diff --git a/test/activity.test.mjs b/test/activity.test.mjs index 08fc9af..1fe0d36 100644 --- a/test/activity.test.mjs +++ b/test/activity.test.mjs @@ -488,3 +488,24 @@ test("unconfigured and unsupported Hub timelines have explicit unavailable resul assert.equal(unsupported.stale, true); assert.deepEqual(unsupported.events, []); }); + +test("Hub history traverses past the live cache limit without evicting forward polling", async () => { + const records=Array.from({length:2600},(_,index)=>event(index+1)); + const timeline=createHubTimelineReader({baseUrl:"http://hub.local",ttl:0,now:()=>start,fetchFn:async url=>{ + if(url.pathname.endsWith("/todo"))return {ok:true,json:async()=>({todos:[]})}; + const before=url.searchParams.get("before"),end=before?records.findIndex(item=>item.id===before):records.length; + return {ok:true,json:async()=>({source:"claude-transcript",events:records.slice(Math.max(0,end-500),end),has_more:end>500})}; + }}); + const head=await timeline.read("long-task",{limit:500}); + let page=head;const all=[...head.events]; + for(let count=0;page.hasMore&&count<10;count++){ + page=await timeline.read("long-task",{before:page.events[0]?.id,limit:500}); + assert.ok(page.events.length,"an available older page must survive the retention cap"); + all.unshift(...page.events); + } + assert.equal(page.hasMore,false);assert.deepEqual(all.map(item=>item.id),records.map(item=>item.id)); + records.push(event(2601)); + const increment=await timeline.read("long-task",{after:head.cursor}); + assert.notEqual(increment.cursorReset,true); + assert.deepEqual(increment.events.map(item=>item.id),["event-2601"]); +}); diff --git a/test/hub-inventory.test.mjs b/test/hub-inventory.test.mjs new file mode 100644 index 0000000..201019e --- /dev/null +++ b/test/hub-inventory.test.mjs @@ -0,0 +1,143 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; +import {createHubInventoryReader,toInventoryCottage} from '../src/hub-inventory.mjs'; +import {createFeed} from '../src/feed.mjs'; +const now=Date.parse('2026-09-22T06:00:00Z'); +const task=(id,extra={})=>({id,recordKind:'logical_task',task:'Short task',originalRequest:'Full request',status:'running',updatedAt:new Date(now).toISOString(),projectPath:'/projects/autohub',worktreePath:'/projects/autohub/.worktrees/fix',branch:'fix/cursor',source:'codex_app_server',provider:'openai',backend:'codex',sessionId:'session',currentActivity:'Testing pagination',controlTargetId:null,capabilities:{canInspect:true,canRespond:false,canSteer:false},...extra}); +test('canonical model hints preserve known providers without inventing a model family',()=>{ + assert.equal(toInventoryCottage(task('codex'),now).model,'gpt'); + assert.equal(toInventoryCottage(task('backend',{provider:null}),now).model,'gpt'); + assert.equal(toInventoryCottage(task('explicit',{model:'claude-opus-4'}),now).model,'opus'); + assert.equal(toInventoryCottage(task('route',{context:{agentKernel:{route:{model:'claude-haiku-4'}}}}),now).model,'haiku'); + assert.equal(toInventoryCottage(task('anthropic',{provider:'anthropic',backend:'claude-code'}),now).model,'anthropic'); + assert.equal(toInventoryCottage(task('unknown',{provider:null,backend:null}),now).model,'unknown'); +}); +test('inventory adapter follows all cursor pages and preserves canonical parent/control identities',async()=>{ + const calls=[]; + const reader=createHubInventoryReader({baseUrl:'http://hub.test',token:'secret',now:()=>now,ttl:0,fetchFn:async(url,options)=>{ + calls.push({url:String(url),options}); + return {ok:true,json:async()=>new URL(url).searchParams.get('cursor') ? {tasks:[task('child',{recordKind:'external_session',parentId:'canonical-parent'})],has_more:false,next_cursor:null} : {tasks:[task('canonical-parent')],has_more:true,next_cursor:'page-2'}}; + }}); + const value=await reader.read(); + assert.equal(calls.length,2); + assert.equal(calls[0].options.headers.authorization,'Bearer secret'); + assert.equal(new URL(calls[0].url).pathname,'/v1/tasks'); + assert.deepEqual(value.agents.map(a=>a.id),['canonical-parent','child']); + assert.equal(value.agents[1].parent,'canonical-parent'); + assert.equal(value.agents[1].conversationTarget.controlTargetId,null); + assert.equal(value.agents[0].originalAsk,'Full request'); + assert.equal(value.agents[0].activity,'Testing pagination'); + assert.equal(value.agents[0].branch,'fix/cursor'); + assert.equal(JSON.stringify(value.agents).includes('secret'),false); +}); +test('configured hub mode never calls standalone scanner or SQLite reader',async()=>{ + let scans=0,reads=0; + const inventory={configured:true,read:async()=>({ok:true,agents:[toInventoryCottage(task('canonical'),now)],keys:new Set(),links:new Map()})}; + const feed=createFeed({hubInventory:inventory,scanClaude:async()=>{scans++;throw new Error('scanner forbidden');},readHub:()=>{reads++;throw new Error('DB forbidden');},resolveRepos:async a=>a,enrich:async a=>a,now:()=>now}); + await feed.scan();assert.equal(scans,0);assert.equal(reads,0);assert.equal(feed.snapshot().agents[0].id,'canonical');assert.equal(feed.snapshot().source,'hub'); +}); +test('failed continuation retains the previous complete feed without exposing token/errors',async()=>{ + let fail=false; + const reader=createHubInventoryReader({baseUrl:'http://hub.test',token:'private-token',ttl:0,fetchFn:async()=>{ + if(fail)throw new Error('private-token connection failed'); + return {ok:true,json:async()=>({tasks:[task('stable')],has_more:false})}; + }}); + await reader.read();fail=true; + await assert.rejects(reader.read(),/Hub inventory temporarily unavailable/); +}); + +test('Hub mode retains last complete snapshot during failure and reads canonical detail/timeline',async()=>{ + let failed=false;const details=[],timelines=[]; + const inventory={configured:true,read:async()=>{if(failed)throw new Error('private failure');return {ok:true,agents:[toInventoryCottage(task('canonical'),now)],keys:new Set(),links:new Map()};},detail:async id=>{details.push(id);return toInventoryCottage(task(id,{currentActivity:'Updated activity'}),now);}}; + const feed=createFeed({hubInventory:inventory,resolveRepos:async a=>a,enrich:async a=>a,now:()=>now,timeline:{configured:true,read:async id=>{timelines.push(id);return {events:[],source:'hub:claude-transcript'};}}}); + await feed.scan();await feed.getActivity('canonical'); + assert.deepEqual(details,['canonical']);assert.deepEqual(timelines,['canonical']); + failed=true;await feed.scan(); + assert.equal(feed.snapshot().stale,true); + assert.equal(feed.snapshot().agents[0].id,'canonical'); + assert.equal(JSON.stringify(feed.snapshot()).includes('private failure'),false); +}); + +test('canonical result links populate the existing artifact shelves contract',()=>{ + const cottage=toInventoryCottage(task('outputs',{resultLinks:[ + {kind:'result',url:'https://example.com/report',label:'Verification report'}, + {kind:'result',url:'https://example.com/build'}, + {kind:'pull_request',url:'https://github.com/acme/repo/pull/1'}, + {kind:'result',url:'javascript:alert(1)'}, + ]}),now); + assert.deepEqual(cottage.artifacts,[ + {url:'https://example.com/report',title:'Verification report'}, + {url:'https://example.com/build',title:'Result'}, + {url:'https://github.com/acme/repo/pull/1',title:'Pull request'}, + ]); +}); + +test('activity detail refresh preserves the enriched feed and returns the current question',async()=>{ + const raw=toInventoryCottage(task('canonical'),now); + const question={id:'new-question',kind:'question',prompt:'Choose scope',detail:'',questions:[]}; + const inventory={configured:true,read:async()=>({ok:true,agents:[raw],keys:new Set(),links:new Map()}),detail:async()=>({...raw,activity:'Waiting for input',inputRequest:question,repo:'',pr:{state:'unknown'}})}; + const feed=createFeed({hubInventory:inventory,resolveRepos:async agents=>agents.map(agent=>({...agent,repo:'acme/repo'})),enrich:async agents=>agents.map(agent=>({...agent,pr:{state:'ready',url:'https://github.com/acme/repo/pull/1',checkedAt:now}})),timeline:{configured:true,read:async()=>({events:[],source:'hub:claude-transcript'})},now:()=>now}); + await feed.scan(); + const before=structuredClone(feed.snapshot()); + const activity=await feed.getActivity('canonical'); + assert.deepEqual(feed.snapshot(),before,'activity polling does not replace enriched inventory state'); + assert.equal(activity.inputRequest.id,'new-question'); +}); + +test('detail failure keeps cached journal pages and marks the input evidence stale',async()=>{ + const {pageActivity}=await import('../src/activity.mjs'); + const {applyActivityPage}=await import('../src/observatory.mjs'); + let failed=false,reads=0; + const raw={...toInventoryCottage(task('canonical'),now),inputRequest:{id:'q',kind:'question',prompt:'Proceed?',questions:[]}}; + const events=[{id:'event-1',kind:'progress',text:'Retained progress',timestamp:now}]; + const inventory={configured:true,read:async()=>({ok:true,agents:[raw],keys:new Set(),links:new Map()}),detail:async()=>{if(failed)throw new Error('detail unavailable');return raw;}}; + const feed=createFeed({hubInventory:inventory,resolveRepos:async a=>a,enrich:async a=>a,timeline:{configured:true,read:async(id,options)=>{reads++;return pageActivity(events,options);}},now:()=>now}); + await feed.scan(); + const cache={events:[],cursor:null};applyActivityPage(cache,await feed.getActivity('canonical')); + failed=true; + const next=await feed.getActivity('canonical',{after:'event-1'}); + assert.equal(next.stale,true);assert.equal(next.inputRequest.stale,true); + assert.equal(next.cursor,'event-1');assert.notEqual(next.cursorReset,true); + applyActivityPage(cache,next); + assert.deepEqual(cache.events,events);assert.equal(reads,2); +}); + +test('canonical freshness keeps healthy running tasks live despite an older task update',()=>{ + const old=new Date(now-3600e3).toISOString(); + for(const fields of [{freshness:{lastSeenAt:new Date(now).toISOString(),isStale:false}},{isStale:false}]){ + const fresh=toInventoryCottage(task('healthy',{updatedAt:old,...fields}),now); + assert.equal(fresh.status,'working');assert.equal(fresh.occupancy,'live'); + assert.equal(fresh.updatedAt,Date.parse(old),'retain the actual task update timestamp'); + } + const stale=toInventoryCottage(task('stale',{freshness:{isStale:true}}),now); + assert.equal(stale.status,'offline'); + const unknown=toInventoryCottage(task('unknown',{updatedAt:old}),now); + assert.equal(unknown.status,'offline','legacy heuristic remains when canonical freshness is absent'); + const question=toInventoryCottage(task('question',{updatedAt:old,freshness:{isStale:false},attentionType:'question',attentionMessage:'Proceed?'}),now); + assert.equal(question.status,'blocked','freshness never hides a pending question'); +}); + +test('canonical requests honor the public original request limit and truncation flag',()=>{ + const text='x'.repeat(32769); + const truncated=toInventoryCottage(task('large',{originalRequest:text}),now); + assert.equal(truncated.originalAsk.length,32768); + assert.equal(truncated.originalAsk,text.slice(0,32768)); + assert.equal(truncated.originalAskTruncated,true); + const exact=toInventoryCottage(task('exact',{originalRequest:text.slice(1)}),now); + assert.equal(exact.originalAskTruncated,false); + const absent=toInventoryCottage(task('missing',{originalRequest:null,task:text}),now); + assert.equal(absent.originalAsk,'');assert.equal(absent.originalAskTruncated,false); +}); + +test('canonical result links reject credentials and malformed URLs before public projection',()=>{ + const cottage=toInventoryCottage(task('links',{resultLinks:[ + {kind:'result',url:'https://user:secret@example.com/report'}, + {kind:'pull_request',url:'https://secret@github.com/acme/repo/pull/1'}, + {kind:'result',url:'https://'}, + {kind:'result',url:'javascript:alert(1)'}, + {kind:'result',url:'https://example.com/report',label:'Report'}, + ]}),now); + assert.deepEqual(cottage.resultLinks,[{kind:'result',url:'https://example.com/report'}]); + assert.deepEqual(cottage.artifacts,[{url:'https://example.com/report',title:'Report'}]); + assert.doesNotMatch(JSON.stringify(cottage),/secret/); +}); diff --git a/test/messages.test.mjs b/test/messages.test.mjs index 612cab8..45ffa87 100644 --- a/test/messages.test.mjs +++ b/test/messages.test.mjs @@ -238,3 +238,32 @@ test('HTTP server exposes capability and accepts only explicit same-origin messa assert.equal((await fetch(base+'/agents',{method:'POST'})).status,405); }finally{server.closeAllConnections();await new Promise(resolve=>server.close(resolve));} }); + +test('canonical inventory controls use refreshed capabilities and reject revoked targets',async()=>{ + const inventory={...agent,conversationTarget:{taskId:'task-1',taskStatus:'running',recordKind:'external_session',controlTargetId:'task-1',capabilities:{canSteer:true,canRespond:false}}}; + const supported=fixture({current:{...task,recordKind:'external_session',controlTargetId:'task-1',capabilities:{canSteer:true,canRespond:false}}}); + assert.equal(supported.messenger.capability(inventory).available,true); + assert.equal((await supported.messenger.send(inventory,message)).body.delivery,'submitted'); + const revoked=fixture({current:{...task,recordKind:'external_session',controlTargetId:null,capabilities:{canSteer:false,canRespond:false}}}); + assert.equal((await revoked.messenger.send(inventory,message)).status,409); + assert.equal(revoked.calls.length,1); + assert.equal(supported.messenger.capability({...inventory,conversationTarget:{...inventory.conversationTarget,controlTargetId:null}}).available,false); +}); + +test('canonical inventory rejects nested stale controls before a write',async()=>{ + const inventory={...agent,conversationTarget:{taskId:'task-1',taskStatus:'running',recordKind:'external_session',controlTargetId:'task-1',capabilities:{canSteer:true}}}; + const stale=fixture({current:{...task,recordKind:'external_session',controlTargetId:'task-1',capabilities:{canSteer:true},freshness:{isStale:true}}}); + const sent=await stale.messenger.send(inventory,message); + assert.equal(sent.status,409); + assert.equal(sent.body.delivery,'not_sent'); + assert.equal(stale.calls.length,1,'fresh detail is inspected but no POST occurs'); + assert.equal(stale.messenger.capability({...inventory,freshness:{isStale:true}}).available,false); +}); + +test('freshly pending canonical input cannot receive a redirect from an older running snapshot',async()=>{ + const inventory={...agent,conversationTarget:{taskId:'task-1',taskStatus:'running',recordKind:'external_session',controlTargetId:'task-1',capabilities:{canSteer:true}}}; + const pending=fixture({current:{...task,recordKind:'external_session',controlTargetId:'task-1',capabilities:{canSteer:true},attentionType:'question',attentionMessage:'Which scope should I use?'}}); + const sent=await pending.messenger.send(inventory,message); + assert.equal(sent.status,409);assert.equal(sent.body.delivery,'not_sent'); + assert.equal(pending.calls.length,1,'only verifies detail; never redirects through a new question'); +}); diff --git a/test/observatory.test.mjs b/test/observatory.test.mjs index f191705..761803a 100644 --- a/test/observatory.test.mjs +++ b/test/observatory.test.mjs @@ -1,6 +1,6 @@ import test from 'node:test'; import assert from 'node:assert/strict'; -import {activityJournalPresentation,activityCacheFor,applyActivityPage,isPracticeDemo,button,handoffAction,appendInlineHandoffs,isApprenticeArrivalActive,apprenticeArrivalPosition,apprenticeResidentPosition,apprenticeResidentTarget,prSnapshotHtml} from '../src/observatory.mjs'; +import {activityJournalPresentation,activityCacheFor,applyActivityPage,returnToLatestActivity,activityPagingButtons,isPracticeDemo,button,handoffAction,appendInlineHandoffs,isApprenticeArrivalActive,apprenticeArrivalPosition,apprenticeResidentPosition,apprenticeResidentTarget,prSnapshotHtml} from '../src/observatory.mjs'; const pending = { id: 'practice-question', prompt: 'Which scope should I use?' }; @@ -100,3 +100,43 @@ test('activity cache clears task artifacts when a session-only cottage advances assert.equal(next.todos,undefined,'an omitted new-session todo snapshot cannot reuse the prior task checklist'); assert.equal(next.cursor,null); }); + +test('older journal pages remain visible within the bounded browser window',()=>{ + const records=Array.from({length:2000},(_,index)=>({id:'history-'+(index+1),kind:'progress',text:'Event '+(index+1),timestamp:Date.parse('2026-09-22T00:00:00Z')+index})); + const cache={events:records.slice(-1500),cursor:'history-2000',hasMore:true}; + applyActivityPage(cache,{events:records.slice(0,500),hasMore:false,source:'hub'}, {older:true}); + assert.equal(cache.events.length,1500); + assert.equal(cache.events[0].id,'history-1');assert.equal(cache.events.at(-1).id,'history-1500'); + assert.equal(cache.cursor,'history-2000','live polling keeps its own cursor'); +}); + +test('historical browsing keeps a recent window and returns to live entries without gaps',()=>{ + const records=Array.from({length:2500},(_,index)=>({id:'history-'+(index+1),kind:'progress',text:'Event '+(index+1),timestamp:Date.parse('2026-09-22T00:00:00Z')+index})); + const cache={events:records.slice(-1500),cursor:'history-2500',hasMore:true}; + applyActivityPage(cache,{events:records.slice(500,1000),hasMore:true,source:'hub'},{older:true}); + applyActivityPage(cache,{events:records.slice(0,500),hasMore:false,source:'hub'},{older:true}); + const historical=cache.events.map(event=>event.id); + const live={id:'history-2501',kind:'progress',text:'New live event',timestamp:records.at(-1).timestamp+1}; + applyActivityPage(cache,{events:[live],cursor:live.id,hasMore:false,source:'hub'}); + assert.deepEqual(cache.events.map(event=>event.id),historical,'forward polling does not move the historical window'); + returnToLatestActivity(cache); + assert.equal(cache.events.length,1500); + assert.deepEqual(cache.events.map(event=>event.id),[...records.slice(-1499),live].map(event=>event.id)); + assert.equal(cache.cursor,live.id); + assert.equal(cache.hasMore,true,'the recent window can page into history again'); + assert.equal(cache.browsingHistory,false); +}); + +test('latest journal action survives a live source cursor reset while browsing history',()=>{ + const row=id=>({id,kind:'progress',text:id,timestamp:Date.parse('2026-09-22T00:00:00Z')}); + const cache={events:[row('recent')],cursor:'recent',hasMore:true}; + assert.doesNotMatch(activityPagingButtons(cache),/data-action="latest"/); + applyActivityPage(cache,{events:[row('old')],hasMore:false},{older:true}); + assert.match(activityPagingButtons(cache),/data-action="latest"\s*>Latest entries/); + applyActivityPage(cache,{events:[row('new-source')],cursor:'new-source',hasMore:true,cursorReset:true}); + assert.equal(cache.events[0].id,'old'); + returnToLatestActivity(cache); + assert.deepEqual(cache.events.map(event=>event.id),['new-source']); + assert.equal(cache.cursor,'new-source');assert.equal(cache.hasMore,true); + assert.doesNotMatch(activityPagingButtons(cache),/data-action="latest"/); +});