- {/* Step progress header */}
+ {/* Step progress header — sequential */}
{pipeline.currentStep && (
@@ -577,8 +593,65 @@ const PipelineRunner = ({ taskDirName, context, cwd }: PipelineRunnerProps) => {
)}
- {/* Messages */}
+ {/* Step progress header — parallel */}
+ {ps && (
+
+
+ {ps.taskName}
+
+ {pDone}/{pTotal} done
+ {pRunning > 0 && · {pRunning} running}
+ {pError > 0 && · {pError} failed}
+
+ {pDone + pError < pTotal && (
+
+ )}
+ {pDone + pError === pTotal && pTotal > 0 && (
+ done
+ )}
+
+
+ )}
+
+ {/* Running stats */}
+ {pipeline.phase === 'running' && (
+
+ {formatElapsed(pipeline.elapsed)}
+ {totalTokens > 0 && {totalTokens.toLocaleString()} tok}
+ {rc.totalUSD > 0 && ${rc.totalUSD.toFixed(3)}}
+
+ )}
+
+ {/* Content area */}
+ {/* Parallel iteration grid */}
+ {ps && (
+
+ {ps.iterations.map((it) => (
+
+ {it.status === 'pending' &&
}
+ {it.status === 'running' &&
}
+ {it.status === 'complete' &&
}
+ {it.status === 'error' &&
}
+
+ {it.label}
+
+ {it.cost && (
+
+ ${it.cost.totalUSD.toFixed(3)}
+
+ )}
+ {it.error && (
+
+ {it.error}
+
+ )}
+
+ ))}
+
+ )}
+
+ {/* Sequential messages */}
{pipeline.messages.map((msg, i) => (
{}} />
@@ -614,8 +687,8 @@ const PipelineRunner = ({ taskDirName, context, cwd }: PipelineRunnerProps) => {
)}
{pipeline.totalCost && (
-
- ${pipeline.totalCost.totalUSD.toFixed(3)} · {pipeline.totalCost.inputTokens + pipeline.totalCost.outputTokens} tokens
+
+ {formatElapsed(pipeline.elapsed)} · {(pipeline.totalCost.inputTokens + pipeline.totalCost.outputTokens).toLocaleString()} tokens · ${pipeline.totalCost.totalUSD.toFixed(3)}
)}
{pipeline.skippedItems.length > 0 && (
@@ -623,6 +696,15 @@ const PipelineRunner = ({ taskDirName, context, cwd }: PipelineRunnerProps) => {
{pipeline.skippedItems.length} skipped ({pipeline.skippedItems.map((s) => s.label).join(', ')})
)}
+ {pipeline.jobId && (
+
+
+ View in Jobs
+
+ )}
);
@@ -656,8 +738,9 @@ export const TaskRunnerModal = ({ open, onOpenChange, task, entryName, entryFull
const taskInfo: TaskInfo = { taskName: task.name, taskDirName: task.dirName, entryName: entryName ?? '', entryType: entryType ?? 'file' };
// Context values for autofill
- // Build a ~/relative path for the agent (works inside bwrap sandbox)
- const entryRelPath = entryName && cwd.path ? `~/${cwd.path}/${entryName}` : entryName ? `~/${entryName}` : undefined;
+ // Build absolute path the agent sees (sandboxed: /data/home/..., non-sandboxed: ~/...)
+ const homePrefix = sandboxed ? '/data/home' : '~';
+ const entryRelPath = entryName && cwd.path ? `${homePrefix}/${cwd.path}/${entryName}` : entryName ? `${homePrefix}/${entryName}` : undefined;
const autofillContext: Record
= {};
if (entryName) autofillContext.entry_name = entryName;
if (entryRelPath) autofillContext.entry_path = entryRelPath;
diff --git a/src/workspaces/officerdev/src/apps/FileBrowser/FileBrowserApp/components/usePipelineRunner.ts b/src/workspaces/officerdev/src/apps/FileBrowser/FileBrowserApp/components/usePipelineRunner.ts
index 1aa955e7..b228de35 100644
--- a/src/workspaces/officerdev/src/apps/FileBrowser/FileBrowserApp/components/usePipelineRunner.ts
+++ b/src/workspaces/officerdev/src/apps/FileBrowser/FileBrowserApp/components/usePipelineRunner.ts
@@ -6,6 +6,7 @@ type Phase = 'ready' | 'running' | 'done';
type StepDef = {
task: string;
foreach?: string;
+ concurrency?: number;
};
type StepStatus = {
@@ -15,31 +16,60 @@ type StepStatus = {
cost?: { inputTokens: number; outputTokens: number; totalUSD: number };
};
+export type IterationStatus = {
+ label: string;
+ status: 'pending' | 'running' | 'complete' | 'error';
+ error?: string;
+ cost?: { inputTokens: number; outputTokens: number; totalUSD: number };
+};
+
+type ParallelStep = {
+ stepIndex: number;
+ taskName: string;
+ concurrency: number;
+ iterations: IterationStatus[];
+};
+
type ServerMessage =
- | { type: 'pipeline:init'; steps: StepDef[] }
- | { type: 'step:start'; stepIndex: number; taskName: string; iteration?: { current: number; total: number; label: string } }
- | { type: 'step:complete'; stepIndex: number; cost?: { inputTokens: number; outputTokens: number; totalUSD: number } }
- | { type: 'step:skip'; stepIndex: number; label: string; reason: string }
- | { type: 'assistant:delta'; text: string }
- | { type: 'assistant:text'; text: string }
- | { type: 'tool:start'; toolCallId: string; toolName: string; toolInput: Record }
- | { type: 'tool:result'; toolCallId: string; output: string; isError: boolean }
- | { type: 'pipeline:complete'; totalCost: { inputTokens: number; outputTokens: number; totalUSD: number } }
- | { type: 'error'; message: string }
- | { type: 'stopped' };
+ | { jobId: string; type: 'pipeline:init'; steps: StepDef[] }
+ | { jobId: string; type: 'step:start'; stepIndex: number; taskName: string; iteration?: { current: number; total: number; label: string } }
+ | { jobId: string; type: 'step:complete'; stepIndex: number; cost?: { inputTokens: number; outputTokens: number; totalUSD: number } }
+ | { jobId: string; type: 'step:skip'; stepIndex: number; label: string; reason: string }
+ | { jobId: string; type: 'step:parallel'; stepIndex: number; taskName: string; iterations: string[]; concurrency: number }
+ | { jobId: string; type: 'iteration:start'; stepIndex: number; label: string }
+ | { jobId: string; type: 'iteration:complete'; stepIndex: number; label: string; cost?: { inputTokens: number; outputTokens: number; totalUSD: number } }
+ | { jobId: string; type: 'iteration:error'; stepIndex: number; label: string; error: string }
+ | { jobId: string; type: 'assistant:delta'; text: string; iterationLabel?: string }
+ | { jobId: string; type: 'assistant:text'; text: string; iterationLabel?: string }
+ | { jobId: string; type: 'tool:start'; toolCallId: string; toolName: string; toolInput: Record; iterationLabel?: string }
+ | { jobId: string; type: 'tool:result'; toolCallId: string; output: string; isError: boolean; iterationLabel?: string }
+ | { jobId: string; type: 'pipeline:complete'; totalCost: { inputTokens: number; outputTokens: number; totalUSD: number } }
+ | { jobId: string; type: 'error'; message: string }
+ | { jobId: string; type: 'stopped' }
+ | { type: 'job:created'; jobId: string }
+ | { type: 'job:state'; jobId: string; status: string; progress: unknown; cost: unknown }
+ | { type: 'error'; message: string };
export function usePipelineRunner() {
const [phase, setPhase] = useState('ready');
const [isConnected, setIsConnected] = useState(false);
+ const [jobId, setJobId] = useState(null);
const [steps, setSteps] = useState([]);
const [currentStep, setCurrentStep] = useState(null);
+ const [parallelStep, setParallelStep] = useState(null);
const [messages, setMessages] = useState([]);
const [streamingText, setStreamingText] = useState('');
const [totalCost, setTotalCost] = useState<{ inputTokens: number; outputTokens: number; totalUSD: number } | null>(null);
+ const [runningCost, setRunningCost] = useState({ inputTokens: 0, outputTokens: 0, totalUSD: 0 });
const [hasError, setHasError] = useState(false);
const [skippedItems, setSkippedItems] = useState>([]);
+ const [elapsed, setElapsed] = useState(0);
const wsRef = useRef(null);
const streamBufferRef = useRef('');
+ const startTimeRef = useRef(0);
+ const timerRef = useRef | null>(null);
+ const jobIdRef = useRef(null);
+ const inParallelRef = useRef(false);
const flushStream = useCallback(() => {
const text = streamBufferRef.current;
@@ -50,6 +80,183 @@ export function usePipelineRunner() {
}
}, []);
+ const stopTimer = useCallback(() => {
+ if (timerRef.current) { clearInterval(timerRef.current); timerRef.current = null; }
+ }, []);
+
+ const addCost = useCallback((cost: { inputTokens: number; outputTokens: number; totalUSD: number }) => {
+ setRunningCost((prev) => ({
+ inputTokens: prev.inputTokens + cost.inputTokens,
+ outputTokens: prev.outputTokens + cost.outputTokens,
+ totalUSD: prev.totalUSD + cost.totalUSD,
+ }));
+ }, []);
+
+ const handleEvent = useCallback((msg: ServerMessage) => {
+ // Filter events by jobId (ignore events from other jobs)
+ if ('jobId' in msg && msg.jobId && jobIdRef.current && msg.jobId !== jobIdRef.current) return;
+
+ switch (msg.type) {
+ case 'job:created':
+ jobIdRef.current = msg.jobId;
+ setJobId(msg.jobId);
+ break;
+
+ case 'job:state':
+ // Reconnection to a completed/failed job
+ if (msg.status === 'completed' || msg.status === 'failed' || msg.status === 'stopped' || msg.status === 'interrupted') {
+ setPhase('done');
+ if (msg.cost) setTotalCost(msg.cost as { inputTokens: number; outputTokens: number; totalUSD: number });
+ if (msg.status === 'failed' || msg.status === 'interrupted') setHasError(true);
+ stopTimer();
+ }
+ break;
+
+ case 'pipeline:init':
+ setSteps(msg.steps);
+ break;
+
+ case 'step:start':
+ flushStream();
+ setMessages([]);
+ setParallelStep(null);
+ inParallelRef.current = false;
+ setCurrentStep({
+ taskName: msg.taskName,
+ iteration: msg.iteration,
+ status: 'running',
+ });
+ break;
+
+ case 'step:complete':
+ flushStream();
+ setCurrentStep((prev) => prev ? { ...prev, status: 'complete', cost: msg.cost } : null);
+ if (msg.cost) addCost(msg.cost);
+ break;
+
+ case 'step:skip':
+ setSkippedItems((prev) => [...prev, { label: msg.label, reason: msg.reason }]);
+ break;
+
+ case 'step:parallel':
+ flushStream();
+ setMessages([]);
+ setCurrentStep(null);
+ inParallelRef.current = true;
+ setParallelStep({
+ stepIndex: msg.stepIndex,
+ taskName: msg.taskName,
+ concurrency: msg.concurrency,
+ iterations: msg.iterations.map((label) => ({ label, status: 'pending' })),
+ });
+ break;
+
+ case 'iteration:start':
+ setParallelStep((prev) => {
+ if (!prev) return prev;
+ return {
+ ...prev,
+ iterations: prev.iterations.map((it) =>
+ it.label === msg.label ? { ...it, status: 'running' } : it,
+ ),
+ };
+ });
+ break;
+
+ case 'iteration:complete':
+ setParallelStep((prev) => {
+ if (!prev) return prev;
+ return {
+ ...prev,
+ iterations: prev.iterations.map((it) =>
+ it.label === msg.label ? { ...it, status: 'complete', cost: msg.cost } : it,
+ ),
+ };
+ });
+ if (msg.cost) addCost(msg.cost);
+ break;
+
+ case 'iteration:error':
+ setParallelStep((prev) => {
+ if (!prev) return prev;
+ return {
+ ...prev,
+ iterations: prev.iterations.map((it) =>
+ it.label === msg.label ? { ...it, status: 'error', error: msg.error } : it,
+ ),
+ };
+ });
+ break;
+
+ case 'assistant:delta':
+ // Skip messages from parallel sub-agents (shown in iteration grid instead)
+ if (inParallelRef.current && msg.iterationLabel) break;
+ streamBufferRef.current += msg.text;
+ setStreamingText(streamBufferRef.current);
+ break;
+
+ case 'assistant:text': {
+ if (inParallelRef.current && msg.iterationLabel) break;
+ const text = msg.text || streamBufferRef.current;
+ if (text) {
+ setMessages((prev) => [...prev, { role: 'assistant', id: crypto.randomUUID(), text }]);
+ }
+ streamBufferRef.current = '';
+ setStreamingText('');
+ break;
+ }
+
+ case 'tool:start':
+ if (inParallelRef.current && msg.iterationLabel) break;
+ flushStream();
+ setMessages((prev) => [
+ ...prev,
+ {
+ role: 'tool' as const,
+ id: crypto.randomUUID(),
+ toolCallId: msg.toolCallId,
+ toolName: msg.toolName,
+ toolInput: msg.toolInput,
+ output: undefined,
+ isError: false,
+ },
+ ]);
+ break;
+
+ case 'tool:result':
+ if (inParallelRef.current && msg.iterationLabel) break;
+ setMessages((prev) =>
+ prev.map((m) =>
+ m.role === 'tool' && 'toolCallId' in m && m.toolCallId === msg.toolCallId
+ ? { ...m, output: msg.output, isError: msg.isError }
+ : m,
+ ),
+ );
+ break;
+
+ case 'pipeline:complete':
+ flushStream();
+ setTotalCost(msg.totalCost);
+ setPhase('done');
+ stopTimer();
+ break;
+
+ case 'error':
+ flushStream();
+ setMessages((prev) => [...prev, { role: 'error' as const, id: crypto.randomUUID(), text: msg.message }]);
+ setHasError(true);
+ setPhase('done');
+ stopTimer();
+ break;
+
+ case 'stopped':
+ flushStream();
+ setPhase('done');
+ stopTimer();
+ break;
+ }
+ }, [flushStream, stopTimer, addCost]);
+
useEffect(() => {
const token = localStorage.getItem('BEARER_TOKEN');
if (!token) return;
@@ -59,98 +266,18 @@ export function usePipelineRunner() {
const ws = new WebSocket(url);
wsRef.current = ws;
- ws.addEventListener('open', () => setIsConnected(true));
+ ws.addEventListener('open', () => {
+ setIsConnected(true);
+ if (jobIdRef.current) {
+ ws.send(JSON.stringify({ type: 'attach', jobId: jobIdRef.current }));
+ }
+ });
ws.addEventListener('close', () => setIsConnected(false));
ws.addEventListener('message', (ev) => {
try {
const msg = JSON.parse(ev.data) as ServerMessage;
-
- switch (msg.type) {
- case 'pipeline:init':
- setSteps(msg.steps);
- break;
-
- case 'step:start':
- // Flush any streaming text from the previous step
- flushStream();
- // Clear messages for the new step iteration
- setMessages([]);
- setCurrentStep({
- taskName: msg.taskName,
- iteration: msg.iteration,
- status: 'running',
- });
- break;
-
- case 'step:complete':
- flushStream();
- setCurrentStep((prev) => prev ? { ...prev, status: 'complete', cost: msg.cost } : null);
- break;
-
- case 'step:skip':
- setSkippedItems((prev) => [...prev, { label: msg.label, reason: msg.reason }]);
- break;
-
- case 'assistant:delta':
- streamBufferRef.current += msg.text;
- setStreamingText(streamBufferRef.current);
- break;
-
- case 'assistant:text': {
- const text = msg.text || streamBufferRef.current;
- if (text) {
- setMessages((prev) => [...prev, { role: 'assistant', id: crypto.randomUUID(), text }]);
- }
- streamBufferRef.current = '';
- setStreamingText('');
- break;
- }
-
- case 'tool:start':
- flushStream();
- setMessages((prev) => [
- ...prev,
- {
- role: 'tool' as const,
- id: crypto.randomUUID(),
- toolCallId: msg.toolCallId,
- toolName: msg.toolName,
- toolInput: msg.toolInput,
- output: undefined,
- isError: false,
- },
- ]);
- break;
-
- case 'tool:result':
- setMessages((prev) =>
- prev.map((m) =>
- m.role === 'tool' && 'toolCallId' in m && m.toolCallId === msg.toolCallId
- ? { ...m, output: msg.output, isError: msg.isError }
- : m,
- ),
- );
- break;
-
- case 'pipeline:complete':
- flushStream();
- setTotalCost(msg.totalCost);
- setPhase('done');
- break;
-
- case 'error':
- flushStream();
- setMessages((prev) => [...prev, { role: 'error' as const, id: crypto.randomUUID(), text: msg.message }]);
- setHasError(true);
- setPhase('done');
- break;
-
- case 'stopped':
- flushStream();
- setPhase('done');
- break;
- }
+ handleEvent(msg);
} catch {
// ignore
}
@@ -159,6 +286,7 @@ export function usePipelineRunner() {
return () => {
ws.close();
wsRef.current = null;
+ stopTimer();
};
}, []);
@@ -169,17 +297,32 @@ export function usePipelineRunner() {
setMessages([]);
setStreamingText('');
setTotalCost(null);
+ setRunningCost({ inputTokens: 0, outputTokens: 0, totalUSD: 0 });
setHasError(false);
setSkippedItems([]);
+ setCurrentStep(null);
+ setParallelStep(null);
+ setElapsed(0);
+ setJobId(null);
+ jobIdRef.current = null;
streamBufferRef.current = '';
+ startTimeRef.current = Date.now();
+ if (timerRef.current) clearInterval(timerRef.current);
+ timerRef.current = setInterval(() => {
+ setElapsed(Math.floor((Date.now() - startTimeRef.current) / 1000));
+ }, 1000);
+
wsRef.current.send(JSON.stringify({ type: 'run', taskDirName, inputs, cwd }));
}, []);
const stop = useCallback(() => {
- if (!wsRef.current || wsRef.current.readyState !== WebSocket.OPEN) return;
- wsRef.current.send(JSON.stringify({ type: 'stop' }));
+ if (!wsRef.current || wsRef.current.readyState !== WebSocket.OPEN || !jobIdRef.current) return;
+ wsRef.current.send(JSON.stringify({ type: 'stop', jobId: jobIdRef.current }));
}, []);
- return { phase, isConnected, steps, currentStep, messages, streamingText, totalCost, hasError, skippedItems, run, stop };
+ return {
+ phase, isConnected, jobId, steps, currentStep, parallelStep, messages, streamingText,
+ totalCost, runningCost, hasError, skippedItems, elapsed, run, stop,
+ };
}