diff --git a/.env.example b/.env.example index 48a6fd1..b133765 100644 --- a/.env.example +++ b/.env.example @@ -5,6 +5,9 @@ SEARXNG_URL= SEARXNG_TOKEN= +# Optional public URL reported by pi-agent-public.ps1 (e.g. your own tunnel domain). +PI_AGENT_PUBLIC_URL= + # Optional higher Context7 quota. Context7 also works at its public unauthenticated limit. CONTEXT7_API_KEY= diff --git a/.gitignore b/.gitignore index 2cf1766..44ed145 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,7 @@ -/data/ +/data/* +# Local tools are code (shareable) — supervisor / public-access scripts for +# the managed Web service. Everything else under data/ stays untracked. +!/data/local-tools/ # Local configuration and credentials. Keep the shareable template only. .env diff --git a/data/local-tools/pi-agent-public.ps1 b/data/local-tools/pi-agent-public.ps1 new file mode 100644 index 0000000..1083194 --- /dev/null +++ b/data/local-tools/pi-agent-public.ps1 @@ -0,0 +1,66 @@ +[CmdletBinding()] +param( + [ValidateSet('status', 'start', 'restart', 'stop')] + [string]$Action = 'status' +) + +$ErrorActionPreference = 'Stop' +$webTaskName = 'PiAgentIntegratedWeb' +$tunnelTaskName = 'PiAgentIntegratedTunnel' +$localUrl = 'http://127.0.0.1:30141/' +# 公网地址从环境变量读取(.env 中设置 PI_AGENT_PUBLIC_URL),不在代码里写死个人域名 +$publicUrl = $env:PI_AGENT_PUBLIC_URL +if ([string]::IsNullOrWhiteSpace($publicUrl)) { + $publicUrl = '(未配置 PI_AGENT_PUBLIC_URL)' +} + +function Get-PublicAccessStatus { + $webTask = Get-ScheduledTask -TaskName $webTaskName -ErrorAction SilentlyContinue + $tunnelTask = Get-ScheduledTask -TaskName $tunnelTaskName -ErrorAction SilentlyContinue + $listener = Get-NetTCPConnection -State Listen -LocalPort 30141 -ErrorAction SilentlyContinue + + $httpStatus = 'unavailable' + try { + $response = Invoke-WebRequest -Uri $localUrl -Method Head -TimeoutSec 10 + $httpStatus = [string][int]$response.StatusCode + } + catch { + if ($_.Exception.Response) { + $httpStatus = [string][int]$_.Exception.Response.StatusCode + } + } + + [pscustomobject]@{ + WebTask = if ($webTask) { [string]$webTask.State } else { 'missing' } + TunnelTask = if ($tunnelTask) { [string]$tunnelTask.State } else { 'missing' } + Port30141 = if ($listener) { 'listening' } else { 'closed' } + LocalHttp = $httpStatus + PublicUrl = $publicUrl + } | Format-List +} + +switch ($Action) { + 'status' { + Get-PublicAccessStatus + } + 'start' { + Start-ScheduledTask -TaskName $tunnelTaskName + Start-ScheduledTask -TaskName $webTaskName + Start-Sleep -Seconds 3 + Get-PublicAccessStatus + } + 'restart' { + Stop-ScheduledTask -TaskName $webTaskName -ErrorAction SilentlyContinue + Stop-ScheduledTask -TaskName $tunnelTaskName -ErrorAction SilentlyContinue + Start-Sleep -Seconds 2 + Start-ScheduledTask -TaskName $tunnelTaskName + Start-ScheduledTask -TaskName $webTaskName + Start-Sleep -Seconds 3 + Get-PublicAccessStatus + } + 'stop' { + Stop-ScheduledTask -TaskName $webTaskName -ErrorAction SilentlyContinue + Stop-ScheduledTask -TaskName $tunnelTaskName -ErrorAction SilentlyContinue + Get-PublicAccessStatus + } +} diff --git a/data/local-tools/pi-agent-web-supervisor.ps1 b/data/local-tools/pi-agent-web-supervisor.ps1 new file mode 100644 index 0000000..c7db677 --- /dev/null +++ b/data/local-tools/pi-agent-web-supervisor.ps1 @@ -0,0 +1,213 @@ +[CmdletBinding()] +param() + +$ErrorActionPreference = 'Stop' +$projectRoot = [IO.Path]::GetFullPath((Join-Path $PSScriptRoot '..\..')) +$logsDirectory = Join-Path $projectRoot 'data\logs' +$standardOutputLog = Join-Path $logsDirectory 'pi-web-service.out.log' +$standardErrorLog = Join-Path $logsDirectory 'pi-web-service.err.log' +$supervisorLog = Join-Path $logsDirectory 'pi-web-supervisor.log' +$restartRequestPath = Join-Path $projectRoot 'data\agent\restart-request.json' +$webUrl = 'http://127.0.0.1:30141/' +$restartRequestVersion = 1 + +New-Item -ItemType Directory -Path $logsDirectory -Force | Out-Null +New-Item -ItemType Directory -Path (Split-Path -Parent $restartRequestPath) -Force | Out-Null +Set-Content -LiteralPath $supervisorLog -Value "$(Get-Date -Format o) supervisor started" + +function Write-SupervisorLog { + param([string]$Message) + + Add-Content -LiteralPath $supervisorLog -Value "$(Get-Date -Format o) $Message" +} + +function Get-RestartRequest { + if (-not (Test-Path -LiteralPath $restartRequestPath -PathType Leaf)) { + return $null + } + + try { + # The extension writes the request with Node (UTF-8 without BOM); + # Windows PowerShell 5.1 Get-Content defaults to ANSI and would corrupt + # non-ASCII testInstructions, breaking ConvertFrom-Json silently. + $request = Get-Content -LiteralPath $restartRequestPath -Raw -Encoding UTF8 | ConvertFrom-Json + if ([int]$request.version -ne $restartRequestVersion) { + return $null + } + if ([string]::IsNullOrWhiteSpace([string]$request.requestId)) { + return $null + } + if ([string]::IsNullOrWhiteSpace([string]$request.sessionId)) { + return $null + } + if ([string]::IsNullOrWhiteSpace([string]$request.testInstructions)) { + return $null + } + + if ([string]$request.state -eq 'ready') { + return $request + } + + # session_shutdown normally changes requested -> ready. If that event + # was interrupted, allow a stale request to proceed after a grace period. + if ([string]$request.state -eq 'requested') { + $createdAt = [DateTimeOffset]::MinValue + if ([DateTimeOffset]::TryParse([string]$request.createdAt, [ref]$createdAt)) { + $ageSeconds = ([DateTimeOffset]::UtcNow - $createdAt.ToUniversalTime()).TotalSeconds + if ($ageSeconds -ge 15) { + return $request + } + } + } + } + catch { + # The extension writes the request atomically. A transient parse failure + # means the file is still being replaced; retry on the next poll. + } + + return $null +} + +function Remove-RestartRequest { + Remove-Item -LiteralPath $restartRequestPath -Force -ErrorAction SilentlyContinue +} + +function Stop-ProcessTree { + param([System.Diagnostics.Process]$Process) + + if (-not $Process -or $Process.HasExited) { + return + } + + Write-SupervisorLog "stopping Web process tree for restart request (PID $($Process.Id))" + & taskkill.exe /PID $Process.Id /T /F 2>$null | Out-Null +} + +function Test-WebReady { + try { + $response = Invoke-WebRequest -Uri $webUrl -Method Get -TimeoutSec 3 -UseBasicParsing + return $response.StatusCode -ge 200 -and $response.StatusCode -lt 400 + } + catch { + return $false + } +} + +function Get-RequestField { + param( + [object]$Request, + [string]$Name + ) + + $property = $Request.PSObject.Properties[$Name] + if ($property) { + return [string]$property.Value + } + + if ($Request -is [System.Collections.IDictionary] -and $Request.Contains($Name)) { + return [string]$Request[$Name] + } + + return '' +} + +function Resume-AgentSession { + param([object]$Request) + + $requestId = Get-RequestField $Request 'requestId' + $sessionId = Get-RequestField $Request 'sessionId' + $testInstructions = Get-RequestField $Request 'testInstructions' + $encodedSessionId = [Uri]::EscapeDataString($sessionId) + $resumeMessage = @" +[Pi Web restart complete] + +Pi Web was restarted by the external supervisor and the previous Agent session was restored. Request ID: $requestId +Continue the previous task without repeating completed edits. + +Post-restart verification: +$testInstructions +"@ + $body = @{ + type = 'prompt' + message = $resumeMessage.Trim() + } | ConvertTo-Json -Compress + $bodyBytes = [Text.Encoding]::UTF8.GetBytes($body) + + $response = Invoke-RestMethod ` + -Uri "http://127.0.0.1:30141/api/agent/$encodedSessionId" ` + -Method Post ` + -ContentType 'application/json; charset=utf-8' ` + -Body $bodyBytes ` + -TimeoutSec 30 + + if (-not $response.success) { + throw "Agent resume API returned an unsuccessful response" + } + + Write-SupervisorLog "resumed Agent session $sessionId for request $requestId" + Remove-RestartRequest +} + +$pendingRequest = Get-RestartRequest + +while ($true) { + $process = $null + try { + $npmCommand = (Get-Command npm.cmd -ErrorAction Stop).Source + Write-SupervisorLog 'starting npm run restart' + + $process = Start-Process ` + -FilePath $npmCommand ` + -ArgumentList @('run', 'restart') ` + -WorkingDirectory $projectRoot ` + -WindowStyle Hidden ` + -RedirectStandardOutput $standardOutputLog ` + -RedirectStandardError $standardErrorLog ` + -PassThru + + $startedAt = Get-Date + $nextResumeAttempt = Get-Date + + while (-not $process.HasExited) { + if ($pendingRequest -and (Get-Date) -ge $nextResumeAttempt) { + # Give npm/Next a moment to finish binding the port before the + # first resume attempt. Failed attempts remain pending and retry. + if (((Get-Date) - $startedAt).TotalSeconds -ge 2 -and (Test-WebReady)) { + try { + Resume-AgentSession $pendingRequest + $pendingRequest = $null + } + catch { + Write-SupervisorLog "Agent resume failed: $($_.Exception.Message)" + $nextResumeAttempt = (Get-Date).AddSeconds(3) + } + } + } + + if (-not $pendingRequest) { + $candidate = Get-RestartRequest + if ($candidate) { + $pendingRequest = $candidate + Write-SupervisorLog "restart request $($candidate.requestId) detected" + Stop-ProcessTree $process + break + } + } + + $null = $process.WaitForExit(500) + } + + if (-not $process.HasExited) { + $null = $process.WaitForExit(10000) + } + + if ($process.HasExited) { + Write-SupervisorLog "web process exited with code $($process.ExitCode)" + } + } + catch { + Write-SupervisorLog "supervisor error: $($_.Exception.Message)" + } + + Start-Sleep -Seconds 5 +} diff --git a/resources/extensions/package.json b/resources/extensions/package.json index dbfccfb..7c1663f 100644 --- a/resources/extensions/package.json +++ b/resources/extensions/package.json @@ -4,7 +4,8 @@ "pi": { "extensions": [ "./searxng-search.ts", - "./vision.ts" + "./vision.ts", + "./runtime-control.ts" ] } } diff --git a/resources/extensions/runtime-control.ts b/resources/extensions/runtime-control.ts new file mode 100644 index 0000000..964674c --- /dev/null +++ b/resources/extensions/runtime-control.ts @@ -0,0 +1,303 @@ +import type { ExtensionAPI } from "@earendil-works/pi-coding-agent"; +import { randomUUID } from "node:crypto"; +import { mkdir, readFile, rename, unlink, writeFile } from "node:fs/promises"; +import { join } from "node:path"; +import { Type } from "typebox"; + +const WEB_BASE_URL = + process.env.PI_AGENT_WEB_URL?.trim() || "http://127.0.0.1:30141"; +const REQUEST_FILE_NAME = "restart-request.json"; +const REQUEST_VERSION = 1; +const IDLE_TIMEOUT_MS = 15_000; +const POLL_INTERVAL_MS = 150; + +type RestartRequest = { + version: number; + requestId: string; + sessionId: string; + cwd: string; + testInstructions: string; + state: "requested" | "ready"; + createdAt: string; + readyAt?: string; +}; + +type AgentStateResponse = { + running?: boolean; + state?: { + isStreaming?: boolean; + isPromptRunning?: boolean; + isCompacting?: boolean; + isBashRunning?: boolean; + }; +}; + +function projectRoot(): string { + return process.env.PI_AGENT_APP_ROOT?.trim() || process.cwd(); +} + +function restartRequestPath(): string { + return join(projectRoot(), "data", "agent", REQUEST_FILE_NAME); +} + +function sleep(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); +} + +async function writeJsonAtomically(filePath: string, value: unknown): Promise { + await mkdir(join(projectRoot(), "data", "agent"), { recursive: true }); + const temporaryPath = `${filePath}.${randomUUID()}.tmp`; + await writeFile(temporaryPath, JSON.stringify(value, null, 2), "utf8"); + try { + await rename(temporaryPath, filePath); + } catch (error) { + await unlink(filePath).catch(() => undefined); + await rename(temporaryPath, filePath).catch(() => { + throw error; + }); + } +} + +async function readRestartRequest(): Promise { + try { + const parsed = JSON.parse( + await readFile(restartRequestPath(), "utf8"), + ) as Partial; + if ( + parsed.version !== REQUEST_VERSION || + typeof parsed.requestId !== "string" || + typeof parsed.sessionId !== "string" || + typeof parsed.cwd !== "string" || + typeof parsed.testInstructions !== "string" || + (parsed.state !== "requested" && parsed.state !== "ready") + ) { + return null; + } + return parsed as RestartRequest; + } catch { + return null; + } +} + +async function markRestartReady(sessionId: string): Promise { + const request = await readRestartRequest(); + if (!request || request.sessionId !== sessionId || request.state !== "requested") { + return; + } + await writeJsonAtomically(restartRequestPath(), { + ...request, + state: "ready", + readyAt: new Date().toISOString(), + }); +} + +function sessionIdFromContext(ctx: { + sessionManager: { getSessionId(): string }; +}): string { + const sessionId = ctx.sessionManager.getSessionId(); + if (!sessionId) throw new Error("The current Pi session has no session ID."); + return sessionId; +} + +function requireRpcMode(ctx: { mode: string }): void { + if (ctx.mode !== "rpc") { + throw new Error( + "Runtime control is available in Pi Web only; use the local /reload command in TUI mode.", + ); + } +} + +async function postAgentCommand( + sessionId: string, + command: Record, +): Promise { + const response = await fetch( + `${WEB_BASE_URL}/api/agent/${encodeURIComponent(sessionId)}`, + { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(command), + }, + ); + if (!response.ok) { + const detail = (await response.text().catch(() => "")).slice(0, 300); + throw new Error( + `Pi Web command failed with HTTP ${response.status}${detail ? `: ${detail}` : ""}`, + ); + } + return response.json().catch(() => undefined); +} + +async function waitForAgentIdle(sessionId: string): Promise { + const deadline = Date.now() + IDLE_TIMEOUT_MS; + while (Date.now() < deadline) { + const response = await fetch( + `${WEB_BASE_URL}/api/agent/${encodeURIComponent(sessionId)}`, + ); + if (!response.ok) throw new Error(`Pi Web state failed with HTTP ${response.status}`); + const data = (await response.json()) as AgentStateResponse; + const state = data.state; + if ( + !data.running || + !state || + (!state.isStreaming && + !state.isPromptRunning && + !state.isCompacting && + !state.isBashRunning) + ) { + return; + } + await sleep(POLL_INTERVAL_MS); + } + throw new Error("Timed out waiting for the current Agent turn to finish."); +} + +function scheduleRuntimeReload( + sessionId: string, + testInstructions: string, +): void { + void (async () => { + try { + await waitForAgentIdle(sessionId); + await postAgentCommand(sessionId, { type: "reload" }); + if (testInstructions) { + await postAgentCommand(sessionId, { + type: "prompt", + message: + "[Pi 运行时已重新加载]\n\n" + + "请继续验证刚才的修改,不要重复编辑。测试要求:\n" + + testInstructions, + }); + } + } catch (error) { + const cause = error instanceof Error ? error.message : String(error); + try { + await postAgentCommand(sessionId, { + type: "prompt", + message: + "[Pi 运行时重载失败]\n\n" + + `请先检查运行时状态。失败原因:${cause}`, + }); + } catch { + // The Web process may be unavailable; the failure is still visible in + // the extension/runtime logs when the session can be recovered. + } + } + })(); +} + +export default function runtimeControlExtension(pi: ExtensionAPI): void { + pi.registerCommand("reload-runtime", { + description: "重新加载 Pi 扩展、skills、prompts 和配置", + handler: async (_args, ctx) => { + await ctx.reload(); + }, + }); + + pi.on("session_shutdown", async (_event, ctx) => { + await markRestartReady(ctx.sessionManager.getSessionId()); + }); + + pi.registerTool({ + name: "reload_runtime", + label: "Reload Pi runtime", + description: + "重新加载 Pi 的扩展、skills、prompts 和运行时配置,不重启 Web 服务。" + + "修改 resources/extensions、skills、prompts 或视觉配置后使用。" + + "可选地在重载完成后自动继续执行测试。", + promptSnippet: "Reload Pi runtime after extension or resource changes", + promptGuidelines: [ + "修改 Pi 扩展、skills、prompts 或视觉配置后,优先使用 reload_runtime。", + "reload_runtime 不会重启 Web 服务;如果需要测试,可填写 testInstructions。", + "不要为了扩展修改直接执行 npm run restart 或 pi-agent-public.ps1。", + ], + parameters: Type.Object({ + testInstructions: Type.Optional( + Type.String({ + description: "重载完成后需要继续执行的测试说明", + maxLength: 4000, + }), + ), + }), + async execute(_toolCallId, params, _signal, _onUpdate, ctx) { + requireRpcMode(ctx); + const sessionId = sessionIdFromContext(ctx); + const testInstructions = + typeof params.testInstructions === "string" + ? params.testInstructions.trim() + : ""; + scheduleRuntimeReload(sessionId, testInstructions); + return { + content: [ + { + type: "text", + text: testInstructions + ? "已安排运行时重载。重载完成后会自动继续执行测试。" + : "已安排运行时重载。", + }, + ], + details: { mode: "reload", sessionId }, + terminate: true, + }; + }, + }); + + pi.registerTool({ + name: "restart_web_and_test", + label: "Restart Pi Web and test", + description: + "通过外部 supervisor 重启 Pi Web,并在服务恢复后自动恢复当前会话、继续执行测试。" + + "这是终止当前 Agent 回合的操作;只用于 Web 服务、依赖、环境变量或 Pi 核心代码必须重启的情况。", + promptSnippet: "Schedule a supervised Pi Web restart and resume testing", + promptGuidelines: [ + "只有修改 Web 服务、Node 依赖、环境变量或 Pi 核心代码时才使用 restart_web_and_test。", + "restart_web_and_test 会短暂断开 Web,但会保存并恢复当前会话。", + "不要直接通过 bash 执行 npm run restart 或 pi-agent-public.ps1。", + "必须在 testInstructions 中说明 Web 恢复后需要执行的验证步骤。", + ], + parameters: Type.Object({ + testInstructions: Type.String({ + description: "Pi Web 重启完成后,Agent 应继续执行的测试说明", + minLength: 1, + maxLength: 4000, + }), + }), + async execute(_toolCallId, params, _signal, _onUpdate, ctx) { + requireRpcMode(ctx); + const sessionId = sessionIdFromContext(ctx); + const testInstructions = params.testInstructions.trim(); + const existing = await readRestartRequest(); + if (existing && (existing.state === "requested" || existing.state === "ready")) { + throw new Error( + `A Web restart is already pending (${existing.requestId}). Wait for it to finish.`, + ); + } + + const request: RestartRequest = { + version: REQUEST_VERSION, + requestId: randomUUID(), + sessionId, + cwd: ctx.cwd, + testInstructions, + state: "requested", + createdAt: new Date().toISOString(), + }; + await writeJsonAtomically(restartRequestPath(), request); + ctx.shutdown(); + return { + content: [ + { + type: "text", + text: + "已请求外部 supervisor 重启 Pi Web。当前回合将结束,服务恢复后会自动继续测试。", + }, + ], + details: { + mode: "restart", + requestId: request.requestId, + }, + terminate: true, + }; + }, + }); +} diff --git a/resources/extensions/vision.ts b/resources/extensions/vision.ts index ec7fea9..f85e0e8 100644 --- a/resources/extensions/vision.ts +++ b/resources/extensions/vision.ts @@ -517,76 +517,52 @@ export default function visionExtension(pi: ExtensionAPI) { if (!Array.isArray(messages) || messages.length === 0) return; if (!isTextOnlyModel(ctx.model)) return; // vision-capable models pass through untouched - // Collect images and replace each with a numbered text placeholder. - const images: LoadedImage[] = []; + // Replace each image in place with its stable transcription. Do not move + // historical descriptions to a later message: DeepSeek caches exact prompt + // prefixes, so rewriting an old image message invalidates everything after it. let dirty = false; + let imageNumber = 0; for (const msg of messages) { if (msg?.role !== "user" || !Array.isArray(msg.content)) continue; const newContent: unknown[] = []; + let messageDirty = false; for (const part of msg.content) { const p = part as { type?: string; image_url?: { url?: string } }; - if (p?.type === "image_url" && typeof p.image_url?.url === "string") { - const img = dataUrlToLoadedImage(p.image_url.url); - if (img) { - images.push(img); - newContent.push({ type: "text", text: `[图片 ${images.length}]` }); - dirty = true; - } else { - newContent.push(part); - } - } else { + if (p?.type !== "image_url" || typeof p.image_url?.url !== "string") { newContent.push(part); + continue; } - } - msg.content = newContent; - } - if (!dirty || images.length === 0) return; - // Transcribe per image: only genuinely new images hit the vision model - // (history resolves from cache), so one request never re-batches the - // whole conversation's images and exceeds Ollama's context window. - const transcribed: string[] = []; - for (let i = 0; i < images.length; i++) { - try { - const desc = await getImageDescription( - images[i], - HOOK_PROMPT, - ctx.signal, - ); - transcribed.push(`【图片${i + 1}】\n${desc}`); - } catch (error) { - const cause = error instanceof Error ? error.message : String(error); - transcribed.push(`【图片${i + 1}】\n[图片处理失败:${cause}]`); - } - } - const description = transcribed.join("\n\n"); - - // Place the full transcription on the LAST image-carrying message (the one - // the model is actively processing); earlier image messages reference it. - // Putting it on the first message instead hid it in history. - const targets: Array<{ part: { type?: string; text?: string } }> = []; - for (const msg of messages) { - if (msg?.role !== "user" || !Array.isArray(msg.content)) continue; - for (const part of msg.content) { - const p = part as { type?: string; text?: string }; - if ( - p?.type === "text" && - typeof p.text === "string" && - p.text.startsWith("[图片 ") - ) { - targets.push({ part: p }); + const img = dataUrlToLoadedImage(p.image_url.url); + if (!img) { + newContent.push(part); + continue; } + + imageNumber += 1; + let description: string; + try { + description = await getImageDescription( + img, + HOOK_PROMPT, + ctx.signal, + ); + } catch { + // Keep failures deterministic. Variable error details in the prompt + // would themselves invalidate the cache on every retry. + description = "[图片转录失败]"; + } + + newContent.push({ + type: "text", + text: `\n\n【图片${imageNumber}】\n${description}\n\n`, + }); + messageDirty = true; + dirty = true; } + if (messageDirty) msg.content = newContent; } - const lastTarget = targets[targets.length - 1]; - if (lastTarget) { - lastTarget.part.text = `[用户上传了 ${images.length} 张图片,以下为视觉模型转录的文本描述]\n\n${description}`; - } - for (const { part } of targets) { - if (part !== lastTarget?.part) { - part.text = "(图片描述见最新消息中的综合转录)"; - } - } + if (!dirty) return; return payload; }); }