Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 74 additions & 0 deletions apps/electron/src/main/lib/adapters/pi-agent-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,40 @@ const PI_NATIVE_MAX_TOTAL_DELAY_MS = 5 * 60_000
const PI_NATIVE_RETRY_JITTER_RATIO = 0.2
const MAX_AUTOMATIC_COMPACTION_CONTINUATIONS = 20

/**
* Bash 流式输出缓冲:toolCallId → 已透传的累计文本。
* 用于计算增量(只推送新增部分),并在 tool_execution_end 时清理。
*
* 生命周期:Pi SDK 契约保证 tool_execution_end 在正常/错误/abort/超时/截断
* 所有路径都发射(bash 工具 throw → error result → emitToolExecutionEnd),
* 因此清理可靠。仅进程被杀等极端场景会残留少量条目(key 全局唯一、不重用,
* 单条几 KB,可接受)。
*/
const bashOutputBuffer = new Map<string, string>()

/**
* 从 Pi 的 partialResult(AgentToolUpdateCallback 参数)中提取文本内容。
*
* partialResult 形如 { content: [{ type: 'text', text: '...' }], details },
* 与工具最终 result 的 content 结构一致。只提取 text 块并拼接。
* 导出供单元测试使用。
*/
export function extractPartialResultText(
partialResult: unknown,
): string | undefined {
if (!partialResult || typeof partialResult !== 'object') return undefined
const content = (partialResult as { content?: unknown }).content
if (!Array.isArray(content)) return undefined
const parts: string[] = []
for (const block of content) {
if (!block || typeof block !== 'object') continue
const b = block as { type?: unknown; text?: unknown }
if (b.type === 'text' && typeof b.text === 'string') parts.push(b.text)
}
const text = parts.join('')
return text.length > 0 ? text : undefined
}

/** Pi SDK 查询选项(扩展通用 AgentQueryInput) */
export interface PiAgentQueryOptions extends AgentQueryInput {
apiKey: string
Expand Down Expand Up @@ -1728,6 +1762,46 @@ export class PiAgentAdapter implements AgentProviderAdapter {
tool_name: displayToolName(event.toolName, event.args as Record<string, unknown> | undefined),
parent_tool_use_id: null,
} as unknown as SDKMessage)
// Bash 工具的流式输出:Pi 的 onUpdate 推送累计输出快照,
// 提取文本内容并计算增量后透传给渲染进程,驱动实时终端化渲染。
// 仅处理 Bash:当前只有 Bash 工具调用 onUpdate,白名单避免未来
// 其他工具意外接入时产生无谓的 buffer + IPC 开销。
if (event.toolName === 'Bash' || event.toolName === 'bash') {
const partialOutput = extractPartialResultText(event.partialResult)
if (partialOutput) {
const prev = bashOutputBuffer.get(event.toolCallId)
if (prev === undefined || !partialOutput.startsWith(prev)) {
// 首帧或快照被截断重置(SDK 只保留 tail):整帧替换,避免渲染层拼出脏内容
queue.push({
type: 'tool_output',
session_id: session.sessionId,
tool_use_id: event.toolCallId,
tool_name: displayToolName(event.toolName, event.args as Record<string, unknown> | undefined),
parent_tool_use_id: null,
output: partialOutput,
replace: true,
} as unknown as SDKMessage)
bashOutputBuffer.set(event.toolCallId, partialOutput)
} else if (partialOutput.length > prev.length) {
// 常规增量:仅推送新增文本,减少 IPC 与渲染压力
queue.push({
type: 'tool_output',
session_id: session.sessionId,
tool_use_id: event.toolCallId,
tool_name: displayToolName(event.toolName, event.args as Record<string, unknown> | undefined),
parent_tool_use_id: null,
output: partialOutput.slice(prev.length),
} as unknown as SDKMessage)
bashOutputBuffer.set(event.toolCallId, partialOutput)
}
}
}
break
case 'tool_execution_end':
// 工具执行结束,清理流式输出缓冲,避免长会话内存泄漏
if (bashOutputBuffer.has(event.toolCallId)) {
bashOutputBuffer.delete(event.toolCallId)
}
break
case 'compaction_start':
// 压缩开始(手动 /compact 或自动阈值/溢出触发):发前端已识别的 compacting system 消息,
Expand Down
36 changes: 35 additions & 1 deletion apps/electron/src/main/lib/adapters/pi-agent-bash.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { describe, expect, test } from 'bun:test'
import { buildWslBashArgs, windowsPathToWslPath } from './pi-agent-adapter'
import { buildWslBashArgs, windowsPathToWslPath, extractPartialResultText } from './pi-agent-adapter'

describe('Pi WSL Bash', () => {
test('Given a Windows workspace path When building WSL Bash arguments Then uses its mounted Linux path', () => {
Expand All @@ -24,3 +24,37 @@ describe('Pi WSL Bash', () => {
expect(windowsPathToWslPath('/home/alice/project')).toBe('/home/alice/project')
})
})

describe('extractPartialResultText', () => {
test('Given Pi SDK partialResult with text content When extracting Then returns the text', () => {
expect(extractPartialResultText({
content: [{ type: 'text', text: 'Compiling...\n' }],
})).toBe('Compiling...\n')
})

test('Given multiple text blocks When extracting Then joins them', () => {
expect(extractPartialResultText({
content: [{ type: 'text', text: 'a' }, { type: 'text', text: 'b\n' }],
})).toBe('ab\n')
})

test('Given empty content When extracting Then returns undefined', () => {
expect(extractPartialResultText({ content: [] })).toBeUndefined()
})

test('Given no content field When extracting Then returns undefined', () => {
expect(extractPartialResultText({ details: {} })).toBeUndefined()
})

test('Given non-text blocks only When extracting Then returns undefined', () => {
expect(extractPartialResultText({ content: [{ type: 'image', source: 'x' }] })).toBeUndefined()
})

test('Given null When extracting Then returns undefined', () => {
expect(extractPartialResultText(null)).toBeUndefined()
})

test('Given empty string text When extracting Then returns undefined', () => {
expect(extractPartialResultText({ content: [{ type: 'text', text: '' }] })).toBeUndefined()
})
})
68 changes: 68 additions & 0 deletions apps/electron/src/renderer/atoms/agent-atoms.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -287,3 +287,71 @@ describe('Agent 流式错误状态', () => {
expect(clearAgentStreamError(errors, 'retried-session')).toBe(errors)
})
})

describe('Agent Bash 流式输出状态', () => {
function stateWithBashActivity(streamingOutput?: string): AgentStreamState {
return createStreamState({
toolActivities: [{
toolUseId: 'tool-bash-1',
toolName: 'Bash',
input: { command: 'npm run build' },
done: false,
streamingOutput,
}],
})
}

test('given Bash 增量 chunk when 收到 tool_output then 追加到 streamingOutput', () => {
const result = applyAgentEvent(stateWithBashActivity('line1\n'), {
type: 'tool_output',
toolUseId: 'tool-bash-1',
output: 'line2\n',
})

expect(result.toolActivities[0]?.streamingOutput).toBe('line1\nline2\n')
expect(result.toolActivities[0]?.done).toBe(false)
})

test('given 无已有输出 when 收到首个 tool_output then 初始化为 chunk 内容', () => {
const result = applyAgentEvent(stateWithBashActivity(), {
type: 'tool_output',
toolUseId: 'tool-bash-1',
output: 'build start\n',
})

expect(result.toolActivities[0]?.streamingOutput).toBe('build start\n')
})

test('given 快照被截断重置 when 收到 replace=true then 整体替换而非追加', () => {
const result = applyAgentEvent(stateWithBashActivity('old-tail-content'), {
type: 'tool_output',
toolUseId: 'tool-bash-1',
output: 'new-tail-content',
replace: true,
})

expect(result.toolActivities[0]?.streamingOutput).toBe('new-tail-content')
})

test('given 空增量 when 收到 tool_output then 保持原引用避免重渲染', () => {
const state = stateWithBashActivity('same')
const result = applyAgentEvent(state, {
type: 'tool_output',
toolUseId: 'tool-bash-1',
output: '',
})

expect(result).toBe(state)
})

test('given 不匹配的 toolUseId when 收到 tool_output then 不修改任何 activity', () => {
const state = stateWithBashActivity('keep')
const result = applyAgentEvent(state, {
type: 'tool_output',
toolUseId: 'tool-other',
output: 'ignored',
})

expect(result).toBe(state)
})
})
47 changes: 47 additions & 0 deletions apps/electron/src/renderer/atoms/agent-atoms.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ export interface ToolActivity {
isBackground?: boolean
/** MCP 工具返回的图片附件 */
imageAttachments?: Array<{ localPath: string; filename: string; mediaType: string }>
/** Bash 等工具执行期间的流式输出(实时终端化渲染用) */
streamingOutput?: string
}

/** 活动分组(Task 子代理) */
Expand Down Expand Up @@ -305,6 +307,31 @@ export const agentSessionStreamingStateAtomFamily = atomFamily((sessionId: strin
atom((get) => get(agentStreamingStatesAtom).get(sessionId)),
)

/**
* 单个 toolUseId 的流式输出派生 atom — 按 toolUseId 切片订阅。
*
* ContentBlock 遍历所有 session 的 toolActivities 查找 streamingOutput,
* 若直接订阅 agentStreamingStatesAtom,任意 session 的 Bash 流式 tick(10Hz)
* 都会让消息树所有 ContentBlock 重渲染。本 family 让订阅者只在目标 toolUseId
* 的 streamingOutput 引用变化时重渲染——其他工具/会话的更新虽然让 base atom
* 变化,但派生 atom 输出引用未变,jotai 自动跳过通知。
*
* toolUseId 由 SDK 生成全局唯一,跨 session 遍历定位是安全的。
*/
export const agentToolStreamingOutputAtomFamily = atomFamily((toolUseId: string) =>
atom<string | undefined>((get) => {
const states = get(agentStreamingStatesAtom)
for (const state of states.values()) {
for (const activity of state.toolActivities) {
if (activity.toolUseId === toolUseId && activity.streamingOutput) {
return activity.streamingOutput
}
}
}
return undefined
}),
)

/**
* 实时 SDKMessage 累积 Map — Phase 2 新增
*
Expand Down Expand Up @@ -760,6 +787,26 @@ export function applyAgentEvent(
}
}

case 'tool_output': {
// Bash 等工具的实时输出 chunk:
// - replace=true:快照被截断重置,直接整体替换缓冲
// - 否则:增量追加到 streamingOutput
// 无匹配 activity 或输出无变化时返回原引用,避免高频事件导致整个状态树重渲染。
const resumed = clearFinishedCompactionForResumedWork(prev)
let changed = false
const toolActivities = resumed.toolActivities.map((t) => {
if (t.toolUseId !== event.toolUseId) return t
const nextOutput = event.replace
? event.output
: `${t.streamingOutput ?? ''}${event.output}`
if (nextOutput === t.streamingOutput) return t
changed = true
return { ...t, streamingOutput: nextOutput }
})
if (!changed) return prev
return { ...resumed, toolActivities }
}

case 'task_backgrounded': {
const resumed = clearFinishedCompactionForResumedWork(prev)
return {
Expand Down
51 changes: 48 additions & 3 deletions apps/electron/src/renderer/components/agent/ContentBlock.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
*/

import * as React from 'react'
import { useAtomValue } from 'jotai'
import {
ChevronRight,
ChevronDown,
Expand All @@ -26,6 +27,7 @@ import { PreviewOpenButton } from './tool-result-renderers/preview-open-button'
import { getTaskGetStatusLabel, parseTaskGetResult, type ParsedTaskGetResult } from './tool-result-renderers/task-get-result'
import { parseTaskListResult, type ParsedTaskListItem } from './tool-result-renderers/task-list-result'
import { formatDuration } from './AgentMessages'
import { agentToolStreamingOutputAtomFamily } from '@/atoms/agent-atoms'
import type {
SDKContentBlock,
SDKMessage,
Expand Down Expand Up @@ -75,6 +77,18 @@ function useToolResult(toolUseId: string, allMessages: SDKMessage[]): ToolResult
}, [toolUseId, allMessages])
}

// ===== useToolStreamingOutput Hook =====

/**
* 按全局唯一 toolUseId 查找工具执行期间的流式输出。
*
* 通过 atomFamily 按 toolUseId 切片订阅:只有目标工具的输出变化时
* 本组件才重渲染,避免任意 session 的 Bash tick(10Hz)污染消息树。
*/
function useToolStreamingOutput(toolUseId: string): string | undefined {
return useAtomValue(agentToolStreamingOutputAtomFamily(toolUseId))
}

// ===== useSubAgentMeta Hook =====

interface SubAgentMeta {
Expand Down Expand Up @@ -338,6 +352,11 @@ function ToolUseBlock({ block, allMessages, animate = false, index = 0, dimmed =
const resultText = toolResult?.result
const isError = toolResult?.isError === true
const shouldShowResult = !!resultText
// Bash 等工具的实时流式输出:从流式状态中按全局唯一 toolUseId 查找 activity
const streamingOutput = useToolStreamingOutput(block.id)
// 渲染条件派生布尔值(避免 JSX 中复杂组合表达式)
const hasResult = shouldShowResult && !!resultText
const hasStreamFallback = !!streamingOutput
const taskGetSummary = React.useMemo(() => {
if (block.name !== 'TaskGet' || !resultText || isError) return null
return parseTaskGetResult(resultText)
Expand All @@ -358,6 +377,25 @@ function ToolUseBlock({ block, allMessages, animate = false, index = 0, dimmed =

const isCompleted = toolResult !== null

// Bash 工具在流式期间命令启动即显示终端(即使当前尚无输出);
// tool_result 一到(isCompleted)立即切结束模式,避免显示假的“执行中”脉冲。
const isStreamingOutput = isStreaming && block.name === 'Bash' && !isCompleted
const showStreaming = isStreamingOutput

// 记录是否经历过流式输出:结束后结果保持展开,避免展开/收起跳动。
// 用户手动收起时清除标记,恢复常规折叠行为。
const [keepResultExpanded, setKeepResultExpanded] = React.useState(false)
React.useEffect(() => {
if (isStreamingOutput) setKeepResultExpanded(true)
}, [isStreamingOutput])
const toggleResult = React.useCallback(() => {
setExpanded((prev) => {
const next = !prev
if (!next) setKeepResultExpanded(false)
return next
})
}, [])

// 运行中显示进行时短语,完成或非流式(已终止)显示完成态短语
const displayLabel = (isCompleted || !isStreaming) ? phrase.label : phrase.loadingLabel
const filePath = extractFilePath(block.input)
Expand Down Expand Up @@ -481,7 +519,7 @@ function ToolUseBlock({ block, allMessages, animate = false, index = 0, dimmed =
'inline-flex max-w-full items-center gap-2 py-0.5 text-left transition-opacity group',
'hover:opacity-70',
)}
onClick={() => setExpanded(!expanded)}
onClick={toggleResult}
>
{!isCompleted && isStreaming ? (
<Loader2 className="size-3.5 animate-spin text-primary/50 shrink-0" />
Expand Down Expand Up @@ -533,17 +571,24 @@ function ToolUseBlock({ block, allMessages, animate = false, index = 0, dimmed =
)}
</button>

{shouldShowResult && resultText && expanded && (
{/* 渲染条件:
* - hasResult:正常结束态有完整结果
* - showStreaming:Bash 流式执行中(命令启动即显示)
* - hasStreamFallback:流式输出缓冲仍存在(如结束但 result 为空兜底)
* 展开条件:用户展开 / 流式中强制显示 / 曾流式且仍有内容(避免结束后内容消失) */}
{(hasResult || showStreaming || hasStreamFallback) && (expanded || showStreaming || (keepResultExpanded && (hasResult || hasStreamFallback))) && (
<div className={cn(
'ml-5.5 mt-1 mb-2 pl-3 border-l-2 border-border/30',
animate && 'animate-in fade-in slide-in-from-top-1 duration-150',
)}>
<ToolResultRenderer
toolName={block.name}
input={block.input}
result={resultText}
result={resultText ?? ''}
isError={isError}
basePath={basePath}
streamingOutput={streamingOutput}
isStreamingOutput={isStreamingOutput}
/>
</div>
)}
Expand Down
Loading