Browse Source

Merge pull request #652 from mhenke: decompose/task-session-manager

refactor: decompose task-session-manager into focused modules
Alvin 1 month ago
parent
commit
e24d9bac93

+ 67 - 248
src/hooks/task-session-manager/index.ts

@@ -1,9 +1,7 @@
-import path from 'node:path';
 import type { PluginInput } from '@opencode-ai/plugin';
 import {
   BackgroundJobBoard,
   type BackgroundJobRecord,
-  type ContextFile,
   deriveTaskSessionLabel,
   parseTaskIdFromTaskOutput,
   parseTaskLaunchOutput,
@@ -18,6 +16,12 @@ import {
   type MessagePart,
   type MessageWithParts,
 } from '../types';
+import type { PendingTaskCall } from './pending-call-tracker';
+import { createPendingCallTracker } from './pending-call-tracker';
+import {
+  createTaskContextTracker,
+  extractReadFiles,
+} from './task-context-tracker';
 
 interface TaskArgs {
   description?: unknown;
@@ -26,124 +30,52 @@ interface TaskArgs {
   task_id?: unknown;
 }
 
-interface PendingTaskCall {
-  callId: string;
-  parentSessionId: string;
-  agentType: string;
-  label: string;
-  resumedTaskId?: string;
-}
-
-const MAX_PENDING_TASK_CALLS = 100;
-
-interface PendingContextFile {
-  path: string;
-  lines: Set<number>;
-  lastReadAt: number;
-}
-
 const BACKGROUND_JOB_BOARD_SENTINEL = 'SENTINEL: background-job-board-v2';
 const BACKGROUND_COMPLETION_COMPLETED = /^Background task completed: /;
 const BACKGROUND_COMPLETION_FAILED = /^Background task failed: /;
 const MAX_PROCESSED_INJECTED_COMPLETIONS = 500;
 const RAW_SESSION_ID_PATTERN = /^ses_[A-Za-z0-9_-]+$/;
 
-/**
- * Simple deterministic string hash for stable occurrence IDs.
- * Uses DJB2 algorithm - fast and good distribution for short strings.
- */
 function djb2Hash(str: string): string {
   let hash = 5381;
   for (let i = 0; i < str.length; i++) {
-    hash = (hash << 5) + hash + str.charCodeAt(i); // hash * 33 + char
+    hash = (hash << 5) + hash + str.charCodeAt(i);
   }
-  // Convert to unsigned 32-bit and then to hex
   return (hash >>> 0).toString(16).padStart(8, '0');
 }
 
-/**
- * Create a stable occurrence ID for synthetic completion deduplication.
- * Prefers part.id, then message.info.id + partIndex, then content-derived hash.
- */
 function createOccurrenceId(
   part: MessagePart,
   message: MessageWithParts,
   partIndex: number,
 ): string {
-  // Prefer explicit part.id if available
   if (typeof part.id === 'string') {
     return part.id;
   }
 
-  // Fall back to message.info.id + partIndex
   if (typeof message.info.id === 'string') {
     return `${message.info.id}:${partIndex}`;
   }
 
-  // Final fallback: content-derived hash from sessionID + parsed taskID/state/result
-  // This ensures the same anonymous synthetic completion is deduped
-  // even when its message index changes between transform calls
   const sessionID = message.info.sessionID ?? 'unknown';
   const content = typeof part.text === 'string' ? part.text : '';
 
-  // Parse task status to get stable identifiers
   const status = parseTaskStatusOutput(content);
   if (status) {
-    // Use taskID + state + result for a stable hash
     const stableKey = `${sessionID}:${status.taskID}:${status.state}:${status.result ?? ''}`;
     const hash = djb2Hash(stableKey);
     return `anon:${hash}`;
   }
 
-  // Fallback to hashing the full content if parsing fails
   const hash = djb2Hash(`${sessionID}:${content}`);
   return `anon:${hash}`;
 }
 
-function extractPath(output: string): string | undefined {
-  return /<path>([^<]+)<\/path>/.exec(output)?.[1];
-}
-
 function extractTaskSummary(output: string): string | undefined {
   const summary = /<summary>\s*([\s\S]*?)\s*<\/summary>/i.exec(output)?.[1];
   return summary?.trim() || undefined;
 }
 
-function normalizePath(root: string, file: string): string {
-  const relative = path.relative(root, file);
-  if (!relative || relative.startsWith('..') || path.isAbsolute(relative)) {
-    return file;
-  }
-  return relative;
-}
-
-function extractReadFiles(
-  root: string,
-  output: { output: unknown; metadata?: unknown },
-): ContextFile[] {
-  if (typeof output.output !== 'string') return [];
-
-  const file = extractPath(output.output);
-  if (!file) return [];
-
-  return [
-    {
-      path: normalizePath(root, file),
-      lineCount: countReadLines(output.output).length,
-      lineNumbers: countReadLines(output.output),
-      lastReadAt: Date.now(),
-    },
-  ];
-}
-
-function countReadLines(output: string): number[] {
-  const lines = new Set<number>();
-  for (const match of output.matchAll(/^([0-9]+):/gm)) {
-    lines.add(Number(match[1]));
-  }
-  return [...lines];
-}
-
 export function createTaskSessionManagerHook(
   _ctx: PluginInput,
   options: {
@@ -161,65 +93,13 @@ export function createTaskSessionManagerHook(
       readContextMinLines: options.readContextMinLines,
       readContextMaxFiles: options.readContextMaxFiles,
     });
-  const pendingCalls = new Map<string, PendingTaskCall>();
-  const pendingCallOrder: string[] = [];
-  const contextByTask = new Map<string, Map<string, PendingContextFile>>();
-  const pendingManagedTaskIds = new Set<string>();
-  const terminalJobsInjectedByParent = new Map<string, Set<string>>();
-  const processedInjectedCompletions = new Set<string>();
-  const processedInjectedCompletionOrder: string[] = [];
-  let anonymousPendingCallId = 0;
 
-  function addTaskContext(taskId: string, files: ContextFile[]): void {
-    if (files.length === 0) return;
-
-    let context = contextByTask.get(taskId);
-    if (!context) {
-      context = new Map();
-      contextByTask.set(taskId, context);
-    }
-    for (const file of files) {
-      const pending = context.get(file.path) ?? {
-        path: file.path,
-        lines: new Set<number>(),
-        lastReadAt: file.lastReadAt,
-      };
-      for (const line of file.lineNumbers ?? []) {
-        pending.lines.add(line);
-      }
-      pending.lastReadAt = Math.max(pending.lastReadAt, file.lastReadAt);
-      context.set(file.path, pending);
-    }
-
-    backgroundJobBoard.addContext(taskId, contextFilesForPrompt(context));
-  }
-
-  function contextFilesForPrompt(
-    context: Map<string, PendingContextFile> | undefined,
-  ): ContextFile[] {
-    if (!context) return [];
-    return [...context.values()].map((file) => ({
-      path: file.path,
-      lineCount: file.lines.size,
-      lastReadAt: file.lastReadAt,
-    }));
-  }
-
-  function canTrackTaskContext(taskId: string): boolean {
-    return (
-      pendingManagedTaskIds.has(taskId) ||
-      backgroundJobBoard.taskIDs().has(taskId)
-    );
-  }
+  const pendingCallTracker = createPendingCallTracker();
+  const taskContextTracker = createTaskContextTracker();
 
-  function pruneContext(): void {
-    const remembered = backgroundJobBoard.taskIDs();
-    for (const taskId of contextByTask.keys()) {
-      if (!pendingManagedTaskIds.has(taskId) && !remembered.has(taskId)) {
-        contextByTask.delete(taskId);
-      }
-    }
-  }
+  const processedInjectedCompletions = new Set<string>();
+  const processedInjectedCompletionOrder: string[] = [];
+  const terminalJobsInjectedByParent = new Map<string, Set<string>>();
 
   function updateBackgroundJobFromOutput(
     output: unknown,
@@ -241,7 +121,8 @@ export function createTaskSessionManagerHook(
       log('[task-session-manager] suppressed late cancelled task error', {
         taskID: status.taskID,
         alias: existing?.alias,
-        state: existing?.state,
+        parsedState: status.state,
+        boardState: existing?.state,
         terminalState: existing?.terminalState,
         result: status.result,
       });
@@ -271,13 +152,13 @@ export function createTaskSessionManagerHook(
       timedOut: updated.timedOut,
     });
 
-    if (updated.terminalUnreconciled) {
-      pendingManagedTaskIds.delete(updated.taskID);
+    if (backgroundJobBoard.isTerminalUnreconciled(updated.taskID)) {
+      taskContextTracker.pendingManagedTaskIds.delete(updated.taskID);
       backgroundJobBoard.addContext(
         updated.taskID,
-        contextFilesForPrompt(contextByTask.get(updated.taskID)),
+        taskContextTracker.contextFilesForPrompt(updated.taskID),
       );
-      pruneContext();
+      taskContextTracker.prune(backgroundJobBoard);
     }
 
     return updated;
@@ -316,12 +197,13 @@ export function createTaskSessionManagerHook(
     if (isFailed && isLateCancelledTaskError(existing, status.state)) {
       part.text = formatCancelledTaskStatusOutput(
         status.taskID,
-        existing?.resultSummary,
+        backgroundJobBoard.getResultSummary(status.taskID),
       );
       log('[task-session-manager] normalized late cancelled injected failure', {
         taskID: status.taskID,
         alias: existing?.alias,
-        state: existing?.state,
+        parsedState: status.state,
+        boardState: existing?.state,
         terminalState: existing?.terminalState,
         result: status.result,
       });
@@ -329,14 +211,9 @@ export function createTaskSessionManagerHook(
       return existing;
     }
 
-    // Enforce summary/state consistency when upstream includes a completion
-    // summary. Current upstream renders synthetic completions as task XML with
-    // the completion/failure label inside <summary> rather than as the first
-    // line of text.
     if (isCompleted && status.state !== 'completed') return undefined;
     if (isFailed && status.state !== 'error') return undefined;
 
-    // Dedupe by synthetic message occurrence using stable occurrence ID
     if (processedInjectedCompletions.has(occurrenceId)) return undefined;
 
     const updated = updateBackgroundJobFromOutput(part.text);
@@ -377,61 +254,6 @@ export function createTaskSessionManagerHook(
     );
   }
 
-  function pendingCallId(input: {
-    callID?: string;
-    sessionID?: string;
-  }): string {
-    return (
-      input.callID ??
-      `${input.sessionID ?? 'unknown'}:anonymous-${++anonymousPendingCallId}`
-    );
-  }
-
-  function rememberPendingCall(call: PendingTaskCall): void {
-    const existingIndex = pendingCallOrder.indexOf(call.callId);
-    if (existingIndex >= 0) {
-      pendingCallOrder.splice(existingIndex, 1);
-    }
-
-    pendingCalls.set(call.callId, call);
-    pendingCallOrder.push(call.callId);
-
-    while (pendingCallOrder.length > MAX_PENDING_TASK_CALLS) {
-      const evictedCallId = pendingCallOrder.shift();
-      if (!evictedCallId) {
-        break;
-      }
-      pendingCalls.delete(evictedCallId);
-    }
-  }
-
-  function takePendingCall(
-    callId?: string,
-    parentSessionId?: string,
-  ): PendingTaskCall | undefined {
-    const resolvedCallId = callId ?? firstPendingCallForParent(parentSessionId);
-    if (!resolvedCallId) return undefined;
-
-    const pending = pendingCalls.get(resolvedCallId);
-    pendingCalls.delete(resolvedCallId);
-
-    const orderIndex = pendingCallOrder.indexOf(resolvedCallId);
-    if (orderIndex >= 0) {
-      pendingCallOrder.splice(orderIndex, 1);
-    }
-
-    return pending;
-  }
-
-  function firstPendingCallForParent(
-    parentSessionId?: string,
-  ): string | undefined {
-    if (!parentSessionId) return undefined;
-    return pendingCallOrder.find(
-      (callId) => pendingCalls.get(callId)?.parentSessionId === parentSessionId,
-    );
-  }
-
   function rememberInjectedTerminalJobs(parentSessionID: string): void {
     const taskIDs = backgroundJobBoard
       .list(parentSessionID)
@@ -500,15 +322,12 @@ export function createTaskSessionManagerHook(
       });
 
       const pendingCall: PendingTaskCall = {
-        callId: pendingCallId({
-          callID: input.callID,
-          sessionID: input.sessionID,
-        }),
+        callId: pendingCallTracker.pendingCallId(input.sessionID, input.callID),
         parentSessionId: input.sessionID,
         agentType,
         label,
       };
-      rememberPendingCall(pendingCall);
+      pendingCallTracker.add(pendingCall);
 
       if (typeof args.task_id !== 'string' || args.task_id.trim() === '') {
         return;
@@ -524,7 +343,7 @@ export function createTaskSessionManagerHook(
       if (!remembered) {
         if (RAW_SESSION_ID_PATTERN.test(requested)) {
           pendingCall.resumedTaskId = requested;
-          rememberPendingCall(pendingCall);
+          pendingCallTracker.add(pendingCall);
           return;
         }
         delete args.task_id;
@@ -532,10 +351,10 @@ export function createTaskSessionManagerHook(
       }
 
       args.task_id = remembered.taskID;
-      pendingManagedTaskIds.add(remembered.taskID);
+      taskContextTracker.pendingManagedTaskIds.add(remembered.taskID);
       backgroundJobBoard.markUsed(input.sessionID, remembered.taskID);
       pendingCall.resumedTaskId = remembered.taskID;
-      rememberPendingCall(pendingCall);
+      pendingCallTracker.add(pendingCall);
     },
 
     'tool.execute.after': async (
@@ -543,18 +362,23 @@ export function createTaskSessionManagerHook(
       output: { output: unknown; metadata?: unknown },
     ): Promise<void> => {
       if (input.tool.toLowerCase() === 'read') {
-        if (input.sessionID && canTrackTaskContext(input.sessionID)) {
-          addTaskContext(
-            input.sessionID,
-            extractReadFiles(_ctx.directory, output),
-          );
+        if (input.sessionID) {
+          const canTrack =
+            taskContextTracker.pendingManagedTaskIds.has(input.sessionID) ||
+            backgroundJobBoard.taskIDs().has(input.sessionID);
+          if (canTrack) {
+            taskContextTracker.addContext(
+              input.sessionID,
+              extractReadFiles(_ctx.directory, output),
+            );
+          }
         }
         return;
       }
 
       if (input.tool.toLowerCase() !== 'task') return;
 
-      const pending = takePendingCall(input.callID, input.sessionID);
+      const pending = pendingCallTracker.take(input.callID, input.sessionID);
 
       if (!pending || typeof output.output !== 'string') return;
       const launch = parseTaskLaunchOutput(output.output);
@@ -574,11 +398,11 @@ export function createTaskSessionManagerHook(
           description: record.description,
           state: record.state,
         });
+        taskContextTracker.pendingManagedTaskIds.add(launch.taskID);
         backgroundJobBoard.addContext(
           launch.taskID,
-          contextFilesForPrompt(contextByTask.get(launch.taskID)),
+          taskContextTracker.contextFilesForPrompt(launch.taskID),
         );
-        pendingManagedTaskIds.add(launch.taskID);
         return;
       }
 
@@ -611,12 +435,12 @@ export function createTaskSessionManagerHook(
         if (pending.resumedTaskId && pending.resumedTaskId !== status.taskID) {
           backgroundJobBoard.drop(pending.resumedTaskId);
         }
-        pendingManagedTaskIds.delete(status.taskID);
-        const contextFiles = contextFilesForPrompt(
-          contextByTask.get(status.taskID),
+        taskContextTracker.pendingManagedTaskIds.delete(status.taskID);
+        backgroundJobBoard.addContext(
+          status.taskID,
+          taskContextTracker.contextFilesForPrompt(status.taskID),
         );
-        backgroundJobBoard.addContext(status.taskID, contextFiles);
-        pruneContext();
+        taskContextTracker.prune(backgroundJobBoard);
         return;
       }
 
@@ -635,10 +459,12 @@ export function createTaskSessionManagerHook(
         backgroundJobBoard.drop(pending.resumedTaskId);
       }
 
-      pendingManagedTaskIds.delete(taskId);
-      const contextFiles = contextFilesForPrompt(contextByTask.get(taskId));
-      backgroundJobBoard.addContext(taskId, contextFiles);
-      pruneContext();
+      taskContextTracker.pendingManagedTaskIds.delete(taskId);
+      backgroundJobBoard.addContext(
+        taskId,
+        taskContextTracker.contextFilesForPrompt(taskId),
+      );
+      taskContextTracker.prune(backgroundJobBoard);
     },
 
     'experimental.chat.messages.transform': async (
@@ -720,7 +546,7 @@ export function createTaskSessionManagerHook(
           info.parentID &&
           options.shouldManageSession(info.parentID)
         ) {
-          pendingManagedTaskIds.add(info.id);
+          taskContextTracker.pendingManagedTaskIds.add(info.id);
         }
         return;
       }
@@ -788,8 +614,6 @@ export function createTaskSessionManagerHook(
             previousTerminalState: before.terminalState,
             terminalUnreconciled: before.terminalUnreconciled,
             resultSummary: before.resultSummary,
-            updatedState: updated?.state,
-            updatedCancellationRequested: updated?.cancellationRequested,
           });
         }
         log('[task-session-manager] busy/status busy observed', {
@@ -799,10 +623,10 @@ export function createTaskSessionManagerHook(
             : false,
           previousState: before?.state,
           previousTerminalState: before?.terminalState,
-          previousCancellationRequested: before?.cancellationRequested,
+          previousCancellationRequested: before?.cancellationRequested ?? false,
           previousLastLiveBusyAt: before?.lastLiveBusyAt,
           updatedState: updated?.state,
-          updatedCancellationRequested: updated?.cancellationRequested,
+          updatedCancellationRequested: updated?.cancellationRequested ?? false,
           updatedLastLiveBusyAt: updated?.lastLiveBusyAt,
         });
         return;
@@ -817,14 +641,16 @@ export function createTaskSessionManagerHook(
         '[task-session-manager] session.deleted observed; clearing job state',
         {
           sessionID: sessionId,
-          deletedJob: backgroundJobBoard.get(sessionId)
-            ? {
-                state: backgroundJobBoard.get(sessionId)?.state,
-                parentSessionID:
-                  backgroundJobBoard.get(sessionId)?.parentSessionID,
-                alias: backgroundJobBoard.get(sessionId)?.alias,
-              }
-            : undefined,
+          deletedJob: (() => {
+            const record = backgroundJobBoard.get(sessionId);
+            return record
+              ? {
+                  state: record.state,
+                  parentSessionID: record.parentSessionID,
+                  alias: record.alias,
+                }
+              : undefined;
+          })(),
           childJobCount: backgroundJobBoard.list(sessionId).length,
           managesSession: options.shouldManageSession(sessionId),
         },
@@ -833,16 +659,9 @@ export function createTaskSessionManagerHook(
       backgroundJobBoard.drop(sessionId);
       backgroundJobBoard.clearParent(sessionId);
       terminalJobsInjectedByParent.delete(sessionId);
-      contextByTask.delete(sessionId);
-      pendingManagedTaskIds.delete(sessionId);
-      pruneContext();
-
-      for (const [callId, pending] of pendingCalls.entries()) {
-        if (pending.parentSessionId !== sessionId) {
-          continue;
-        }
-        takePendingCall(callId);
-      }
+      taskContextTracker.clearSession(sessionId);
+      taskContextTracker.prune(backgroundJobBoard);
+      pendingCallTracker.clearSession(sessionId);
     },
   };
 
@@ -864,7 +683,7 @@ export function createTaskSessionManagerHook(
     });
     output.output = formatCancelledTaskStatusOutput(
       status.taskID,
-      existing?.resultSummary,
+      backgroundJobBoard.getResultSummary(status.taskID),
     );
     if (isObjectRecord(output) && isObjectRecord(output.metadata)) {
       output.metadata.state = 'cancelled';

+ 57 - 0
src/hooks/task-session-manager/pending-call-tracker.ts

@@ -0,0 +1,57 @@
+export interface PendingTaskCall {
+  callId: string;
+  parentSessionId: string;
+  agentType: string;
+  label: string;
+  resumedTaskId?: string;
+}
+
+const MAX_PENDING_TASK_CALLS = 100;
+
+export function createPendingCallTracker() {
+  const pendingCalls = new Map<string, PendingTaskCall>();
+  let anonymousPendingCallId = 0;
+
+  return {
+    add(call: PendingTaskCall) {
+      pendingCalls.delete(call.callId);
+      pendingCalls.set(call.callId, call);
+      while (pendingCalls.size > MAX_PENDING_TASK_CALLS) {
+        const firstKey = pendingCalls.keys().next().value;
+        if (firstKey === undefined) break;
+        pendingCalls.delete(firstKey);
+      }
+    },
+
+    take(callId?: string, parentSessionId?: string) {
+      if (!callId && parentSessionId) {
+        for (const id of pendingCalls.keys()) {
+          const call = pendingCalls.get(id);
+          if (call && call.parentSessionId === parentSessionId) {
+            callId = id;
+            break;
+          }
+        }
+      }
+      if (!callId) return undefined;
+      const pending = pendingCalls.get(callId);
+      pendingCalls.delete(callId);
+      return pending;
+    },
+
+    clearSession(sessionId: string) {
+      for (const [callId, pending] of pendingCalls.entries()) {
+        if (pending.parentSessionId === sessionId) {
+          pendingCalls.delete(callId);
+        }
+      }
+    },
+
+    pendingCallId(sessionID?: string, callID?: string) {
+      return (
+        callID ??
+        `${sessionID ?? 'unknown'}:anonymous-${++anonymousPendingCallId}`
+      );
+    },
+  };
+}

+ 100 - 0
src/hooks/task-session-manager/task-context-tracker.ts

@@ -0,0 +1,100 @@
+import path from 'node:path';
+import type { ContextFile } from '../../utils';
+
+interface PendingContextFile {
+  path: string;
+  lines: Set<number>;
+  lastReadAt: number;
+}
+
+export function createTaskContextTracker() {
+  const contextByTask = new Map<string, Map<string, PendingContextFile>>();
+  const pendingManagedTaskIds = new Set<string>();
+
+  return {
+    pendingManagedTaskIds,
+
+    addContext(taskId: string, files: ContextFile[]) {
+      if (files.length === 0) return;
+      let context = contextByTask.get(taskId);
+      if (!context) {
+        context = new Map();
+        contextByTask.set(taskId, context);
+      }
+      for (const file of files) {
+        const pending = context.get(file.path) ?? {
+          path: file.path,
+          lines: new Set<number>(),
+          lastReadAt: file.lastReadAt,
+        };
+        for (const line of file.lineNumbers ?? []) {
+          pending.lines.add(line);
+        }
+        pending.lastReadAt = Math.max(pending.lastReadAt, file.lastReadAt);
+        context.set(file.path, pending);
+      }
+    },
+
+    canTrack(taskId: string, backgroundJobBoard: { taskIDs(): Set<string> }) {
+      return (
+        pendingManagedTaskIds.has(taskId) ||
+        backgroundJobBoard.taskIDs().has(taskId)
+      );
+    },
+
+    prune(backgroundJobBoard: { taskIDs(): Set<string> }) {
+      const remembered = backgroundJobBoard.taskIDs();
+      for (const taskId of contextByTask.keys()) {
+        if (!pendingManagedTaskIds.has(taskId) && !remembered.has(taskId)) {
+          contextByTask.delete(taskId);
+        }
+      }
+    },
+
+    clearSession(sessionId: string) {
+      contextByTask.delete(sessionId);
+      pendingManagedTaskIds.delete(sessionId);
+    },
+
+    contextFilesForPrompt(taskId: string): ContextFile[] {
+      const context = contextByTask.get(taskId);
+      if (!context) return [];
+      return [...context.values()].map((file) => ({
+        path: file.path,
+        lineCount: file.lines.size,
+        lastReadAt: file.lastReadAt,
+      }));
+    },
+  };
+}
+
+export function extractReadFiles(
+  root: string,
+  output: { output: unknown; metadata?: unknown },
+): ContextFile[] {
+  if (typeof output.output !== 'string') return [];
+
+  const extractPath = /<path>([^<]+)<\/path>/.exec(output.output)?.[1];
+  if (!extractPath) return [];
+
+  const relative = path.relative(root, extractPath);
+  const normalized =
+    !relative || relative.startsWith('..') || path.isAbsolute(relative)
+      ? extractPath
+      : relative;
+
+  const matchedLines = new Set<number>();
+  for (const match of output.output.matchAll(/^([0-9]+):/gm)) {
+    matchedLines.add(Number(match[1]));
+  }
+  const lineNumbers = [...matchedLines];
+
+  return [
+    {
+      path: normalized,
+      lineCount: lineNumbers.length,
+      lineNumbers,
+      lastReadAt: Date.now(),
+    },
+  ];
+}

+ 17 - 8
src/multiplexer/session-manager.ts

@@ -6,7 +6,10 @@ import {
   isServerRunning,
   type Multiplexer,
 } from '../multiplexer';
-import type { BackgroundJobBoard } from '../utils/background-job-board';
+import type {
+  BackgroundJobBoard,
+  BackgroundJobState,
+} from '../utils/background-job-board';
 import { log } from '../utils/logger';
 
 interface TrackedSession {
@@ -252,7 +255,7 @@ export class MultiplexerSessionManager {
         tracked: this.sessions.has(sessionId),
         known: this.knownSessions.has(sessionId),
         ownerInstanceId: this.sessions.get(sessionId)?.ownerInstanceId,
-        backgroundJobState: this.backgroundJobBoard?.get(sessionId)?.state,
+        backgroundJobState: this.backgroundJobState(sessionId),
       });
 
       await this.closeSession(sessionId, 'idle');
@@ -273,7 +276,7 @@ export class MultiplexerSessionManager {
         tracked: this.sessions.has(sessionId),
         known: this.knownSessions.has(sessionId),
         ownerInstanceId: this.sessions.get(sessionId)?.ownerInstanceId,
-        backgroundJobState: this.backgroundJobBoard?.get(sessionId)?.state,
+        backgroundJobState: this.backgroundJobState(sessionId),
       });
       await this.closeSession(sessionId, 'idle');
       return;
@@ -290,7 +293,7 @@ export class MultiplexerSessionManager {
         tracked: this.sessions.has(sessionId),
         known: this.knownSessions.has(sessionId),
         ownerInstanceId: this.sessions.get(sessionId)?.ownerInstanceId,
-        backgroundJobState: this.backgroundJobBoard?.get(sessionId)?.state,
+        backgroundJobState: this.backgroundJobState(sessionId),
       });
       await this.respawnIfKnown(sessionId);
     }
@@ -309,7 +312,7 @@ export class MultiplexerSessionManager {
       tracked: this.sessions.has(sessionId),
       known: this.knownSessions.has(sessionId),
       ownerInstanceId: this.sessions.get(sessionId)?.ownerInstanceId,
-      backgroundJobState: this.backgroundJobBoard?.get(sessionId)?.state,
+      backgroundJobState: this.backgroundJobState(sessionId),
     });
 
     this.deferredIdleCloses.delete(sessionId);
@@ -456,7 +459,7 @@ export class MultiplexerSessionManager {
           sessionId,
           paneId: tracked.paneId,
           reason,
-          backgroundJobState: this.backgroundJobBoard?.get(sessionId)?.state,
+          backgroundJobState: this.backgroundJobState(sessionId),
         },
       );
       return;
@@ -470,7 +473,7 @@ export class MultiplexerSessionManager {
       sessionId,
       paneId: tracked.paneId,
       reason,
-      backgroundJobState: this.backgroundJobBoard?.get(sessionId)?.state,
+      backgroundJobState: this.backgroundJobState(sessionId),
       parentId: tracked.parentId,
       title: tracked.title,
     });
@@ -606,8 +609,14 @@ export class MultiplexerSessionManager {
     return event.properties?.info?.id ?? event.properties?.sessionID;
   }
 
+  private backgroundJobState(
+    sessionId: string,
+  ): BackgroundJobState | undefined {
+    return this.backgroundJobBoard?.getState(sessionId);
+  }
+
   private isRunningBackgroundJob(sessionId: string): boolean {
-    return this.backgroundJobBoard?.get(sessionId)?.state === 'running';
+    return this.backgroundJobBoard?.isRunning(sessionId) ?? false; // ponytail: intent-revealing query
   }
 
   async retryDeferredIdleClose(sessionId: string): Promise<void> {

+ 30 - 0
src/tools/cancel-task.test.ts

@@ -426,6 +426,36 @@ describe('cancel_task tool', () => {
     });
   });
 
+  test('cancelSessionByID returns state: error when abort throws non-SessionStillRunningError, even if board shows running', async () => {
+    const { board, abort, cancelTask } = createTool({
+      includeDelete: false,
+      abort: async () => {
+        throw new Error('network timeout');
+      },
+    });
+    // Register a running job so that isRunning(taskID) would be true
+    // if the function incorrectly checks it.
+    board.registerLaunch({
+      taskID: 'ses_running',
+      parentSessionID: 'parent-1',
+      agent: 'fixer',
+    });
+    // Override resolve to return undefined, forcing the cancelSessionByID
+    // raw session path instead of the tracked task path.
+    board.resolve = mock(() => undefined);
+
+    const output = await cancelTask.execute(
+      { task_id: 'ses_running', reason: 'regression guard' },
+      context,
+    );
+
+    expect(abort).toHaveBeenCalledWith({ path: { id: 'ses_running' } });
+    // cancelSessionByID must return state: error for non-SessionStillRunningError,
+    // NOT state: running (which would happen if || isRunning() were present).
+    expect(String(output)).toContain('state: error');
+    expect(String(output)).not.toContain('state: running');
+  });
+
   test('denies non-orchestrator agents', async () => {
     const { cancelTask } = createTool();
 

+ 28 - 16
src/tools/cancel-task.ts

@@ -61,9 +61,15 @@ Use only for obsolete, wrong, conflicting, or user-requested cancellation. Accep
         parentSessionID,
         requested,
         resolvedTaskID: job?.taskID,
-        alias: job?.alias,
-        state: job?.state,
-        terminalState: job?.terminalState,
+        alias: job
+          ? options.backgroundJobBoard.field(job.taskID, 'alias')
+          : undefined,
+        state: job
+          ? options.backgroundJobBoard.field(job.taskID, 'state')
+          : undefined,
+        terminalState: job
+          ? options.backgroundJobBoard.field(job.taskID, 'terminalState')
+          : undefined,
         cancellationRequested: job?.cancellationRequested,
       });
       if (!job) {
@@ -77,11 +83,13 @@ Use only for obsolete, wrong, conflicting, or user-requested cancellation. Accep
           }
 
           const knownJob = options.backgroundJobBoard.get(requested);
-          if (knownJob && knownJob.parentSessionID !== parentSessionID) {
+          const ownerParentSessionID =
+            options.backgroundJobBoard.getParentSessionID(requested);
+          if (knownJob && ownerParentSessionID !== parentSessionID) {
             log('[cancel-task] rejected unowned tracked raw session', {
               parentSessionID,
               taskID: requested,
-              ownerParentSessionID: knownJob.parentSessionID,
+              ownerParentSessionID,
             });
             return unknownTaskOutput(
               requested,
@@ -119,9 +127,11 @@ Use only for obsolete, wrong, conflicting, or user-requested cancellation. Accep
         await abortAndVerifySession(options, job.taskID);
       } catch (error) {
         const stillRunning = error instanceof SessionStillRunningError;
+        const boardRunning = options.backgroundJobBoard.isRunning(job.taskID);
         log('[cancel-task] abort failed', {
           taskID: job.taskID,
           stillRunning,
+          boardRunning,
           error: error instanceof Error ? error.message : String(error),
         });
         options.backgroundJobBoard.updateStatus({
@@ -141,26 +151,29 @@ Use only for obsolete, wrong, conflicting, or user-requested cancellation. Accep
         ].join('\n');
       }
 
-      const cancelled = options.backgroundJobBoard.markCancelled(
+      options.backgroundJobBoard.markCancelled(
         job.taskID,
         args.reason,
         Date.now(),
         { force: true },
       );
+      const state = options.backgroundJobBoard.getState(job.taskID);
       log('[cancel-task] marked job cancelled after verified abort', {
         taskID: job.taskID,
-        alias: job.alias,
-        previousState: job.state,
-        state: cancelled?.state,
-        cancellationRequested: cancelled?.cancellationRequested,
+        alias: options.backgroundJobBoard.field(job.taskID, 'alias'),
+        state,
+        cancellationRequested: options.backgroundJobBoard.field(
+          job.taskID,
+          'cancellationRequested',
+        ),
       });
 
       return [
         `task_id: ${job.taskID}`,
-        `state: ${cancelled?.state ?? 'cancelled'}`,
+        `state: ${state ?? 'cancelled'}`,
         '',
         '<task_error>',
-        cancelled?.resultSummary ?? 'cancelled',
+        options.backgroundJobBoard.getResultSummary(job.taskID) ?? 'cancelled',
         '</task_error>',
       ].join('\n');
     },
@@ -257,12 +270,11 @@ async function abortAndVerifySession(
       stableStoppedForMs: stableStoppedSince
         ? Date.now() - stableStoppedSince
         : 0,
-      boardState: options.backgroundJobBoard.get(taskID)?.state,
-      boardLastLiveBusyAt:
-        options.backgroundJobBoard.get(taskID)?.lastLiveBusyAt,
+      boardState: options.backgroundJobBoard.getState(taskID),
+      boardLastLiveBusyAt: options.backgroundJobBoard.getLastLiveBusyAt(taskID),
     });
     const boardLastLiveBusyAt =
-      options.backgroundJobBoard.get(taskID)?.lastLiveBusyAt;
+      options.backgroundJobBoard.getLastLiveBusyAt(taskID);
     if (boardLastLiveBusyAt && boardLastLiveBusyAt >= abortStartedAt) {
       log('[cancel-task] abort verification saw board busy after abort', {
         taskID,

+ 92 - 0
src/utils/background-job-board.test.ts

@@ -732,4 +732,96 @@ describe('BackgroundJobBoard', () => {
     const prompt = board.formatForPrompt('parent-1', 9_000);
     expect(prompt).toContain('running [resumed, 4s ago]');
   });
+
+  describe('intent-revealing query methods', () => {
+    test('isRunning: true for running jobs, false for terminal/reconciled/unknown', () => {
+      const board = new BackgroundJobBoard();
+      board.registerLaunch({
+        taskID: 'running-1',
+        parentSessionID: 'parent-1',
+        agent: 'fixer',
+        now: 100,
+      });
+      board.registerLaunch({
+        taskID: 'terminal-1',
+        parentSessionID: 'parent-1',
+        agent: 'fixer',
+        now: 100,
+      });
+      board.updateStatus({
+        taskID: 'terminal-1',
+        state: 'completed',
+        now: 200,
+      });
+      board.markReconciled('terminal-1', 300);
+
+      expect(board.isRunning('running-1')).toBe(true);
+      expect(board.isRunning('terminal-1')).toBe(false);
+      expect(board.isRunning('unknown-1')).toBe(false);
+    });
+
+    test('isTerminalUnreconciled: true after updateStatus to terminal, false after markReconciled', () => {
+      const board = new BackgroundJobBoard();
+      board.registerLaunch({
+        taskID: 'job-1',
+        parentSessionID: 'parent-1',
+        agent: 'fixer',
+        now: 100,
+      });
+
+      expect(board.isTerminalUnreconciled('job-1')).toBe(false);
+      board.updateStatus({ taskID: 'job-1', state: 'completed', now: 200 });
+      expect(board.isTerminalUnreconciled('job-1')).toBe(true);
+      board.markReconciled('job-1', 300);
+      expect(board.isTerminalUnreconciled('job-1')).toBe(false);
+      expect(board.isTerminalUnreconciled('unknown-1')).toBe(false);
+    });
+
+    test('getResultSummary: returns summary after updateStatus with result', () => {
+      const board = new BackgroundJobBoard();
+      board.registerLaunch({
+        taskID: 'job-1',
+        parentSessionID: 'parent-1',
+        agent: 'fixer',
+        now: 100,
+      });
+      board.updateStatus({
+        taskID: 'job-1',
+        state: 'completed',
+        resultSummary: 'all good',
+        now: 200,
+      });
+
+      expect(board.getResultSummary('job-1')).toBe('all good');
+      expect(board.getResultSummary('unknown-1')).toBeUndefined();
+    });
+
+    test('getLastLiveBusyAt: returns timestamp after markRunningFromLiveSession', () => {
+      const board = new BackgroundJobBoard();
+      board.registerLaunch({
+        taskID: 'job-1',
+        parentSessionID: 'parent-1',
+        agent: 'fixer',
+        now: 100,
+      });
+
+      expect(board.getLastLiveBusyAt('job-1')).toBe(100);
+      board.markRunningFromLiveSession('job-1', 200);
+      expect(board.getLastLiveBusyAt('job-1')).toBe(200);
+      expect(board.getLastLiveBusyAt('unknown-1')).toBeUndefined();
+    });
+
+    test('getParentSessionID: returns parentSessionID after registerLaunch', () => {
+      const board = new BackgroundJobBoard();
+      board.registerLaunch({
+        taskID: 'job-1',
+        parentSessionID: 'parent-1',
+        agent: 'fixer',
+        now: 100,
+      });
+
+      expect(board.getParentSessionID('job-1')).toBe('parent-1');
+      expect(board.getParentSessionID('unknown-1')).toBeUndefined();
+    });
+  });
 });

+ 38 - 2
src/utils/background-job-board.ts

@@ -327,6 +327,39 @@ export class BackgroundJobBoard {
     return this.jobs.get(taskID);
   }
 
+  field<K extends keyof BackgroundJobRecord>(
+    taskID: string,
+    key: K,
+  ): BackgroundJobRecord[K] | undefined {
+    return this.get(taskID)?.[key];
+  }
+
+  isRunning(taskID: string): boolean {
+    const job = this.get(taskID);
+    return job?.state === 'running';
+  }
+
+  isTerminalUnreconciled(taskID: string): boolean {
+    const job = this.get(taskID);
+    return !!job?.terminalUnreconciled;
+  }
+
+  getResultSummary(taskID: string): string | undefined {
+    return this.field(taskID, 'resultSummary');
+  }
+
+  getLastLiveBusyAt(taskID: string): number | undefined {
+    return this.field(taskID, 'lastLiveBusyAt');
+  }
+
+  getParentSessionID(taskID: string): string | undefined {
+    return this.field(taskID, 'parentSessionID');
+  }
+
+  getState(taskID: string): BackgroundJobState | undefined {
+    return this.field(taskID, 'state');
+  }
+
   resolve(
     parentSessionID: string,
     taskIDOrAlias: string,
@@ -366,8 +399,11 @@ export class BackgroundJobBoard {
     for (const file of files) {
       const previous = existing.get(file.path);
       if (previous) {
-        previous.lineCount = Math.max(previous.lineCount, file.lineCount);
-        previous.lastReadAt = Math.max(previous.lastReadAt, file.lastReadAt);
+        existing.set(file.path, {
+          ...previous,
+          lineCount: Math.max(previous.lineCount, file.lineCount),
+          lastReadAt: Math.max(previous.lastReadAt, file.lastReadAt),
+        });
       } else {
         existing.set(file.path, { ...file });
       }