diff --git a/packages/core/src/events.ts b/packages/core/src/events.ts index d103137354..4c14c169e7 100644 --- a/packages/core/src/events.ts +++ b/packages/core/src/events.ts @@ -782,6 +782,7 @@ type ShellRunResultMetadata = { kind: 'shell_run'; ref: string; status: ShellRunStatus; + pid?: number; cwd: string; cmd: string; startedAt: number; diff --git a/packages/core/src/shell-run-result.ts b/packages/core/src/shell-run-result.ts index e10aa47194..cf3bb564c6 100644 --- a/packages/core/src/shell-run-result.ts +++ b/packages/core/src/shell-run-result.ts @@ -111,6 +111,7 @@ const CURRENT_TERMINAL_RESULT_SHAPE = defineObjectShape()( const CURRENT_SHELL_RUN_RESULT_SHAPE = defineObjectShape()( ['kind', 'ref', 'mode', 'status', 'cwd', 'cmd', 'startedAt', 'updatedAt', 'revision'], [ + 'pid', 'completedAt', 'exitCode', 'failureMessage', diff --git a/packages/core/src/shell-run.ts b/packages/core/src/shell-run.ts index 7c5ae2da33..e5b937719b 100644 --- a/packages/core/src/shell-run.ts +++ b/packages/core/src/shell-run.ts @@ -137,6 +137,8 @@ export interface ShellRunRecord { cwd: string; command: string; status: ShellRunStatus; + /** Native root process id, when admitted by the process driver. */ + pid?: number; exitCode?: number; failureMessage?: string; startedAt: number; @@ -159,7 +161,14 @@ export interface ShellRunRecord { export type ShellRunPatch = Partial< Pick< ShellRunRecord, - 'status' | 'exitCode' | 'failureMessage' | 'updatedAt' | 'completedAt' | 'observedAt' | 'output' + | 'status' + | 'pid' + | 'exitCode' + | 'failureMessage' + | 'updatedAt' + | 'completedAt' + | 'observedAt' + | 'output' > >; @@ -307,6 +316,7 @@ const SHELL_RUN_SESSION_ID_PATTERN = /^[A-Za-z0-9_-]{1,128}$/; const SHELL_RUN_PATCH_KEYS: ReadonlySet = new Set([ 'status', + 'pid', 'exitCode', 'failureMessage', 'updatedAt', @@ -325,6 +335,7 @@ const SHELL_RUN_RECORD_KEYS: ReadonlySet = new Set([ 'cwd', 'command', 'status', + 'pid', 'startedAt', 'updatedAt', 'completedAt', @@ -380,6 +391,7 @@ export function normalizeShellRunRecord( record.sessionId === sessionId && record.shellRunId === shellRunId && isShellRunStatus(record.status) && + (record.pid === undefined || isPositiveInteger(record.pid)) && isFiniteNumber(record.startedAt) && isFiniteNumber(record.updatedAt) && isPositiveInteger(record.revision) && @@ -515,6 +527,7 @@ function isShellRunSandboxEscalation(value: unknown, execution: unknown): boolea function canonicalShellRunRecord(record: ShellRunRecord): ShellRunRecord { return { + ...(record.pid !== undefined ? { pid: record.pid } : {}), shellRunId: record.shellRunId, sessionId: record.sessionId, ...(record.sourceRunId !== undefined ? { sourceRunId: record.sourceRunId } : {}), diff --git a/packages/runtime-host/src/__tests__/web-fetch-tool.test.ts b/packages/runtime-host/src/__tests__/web-fetch-tool.test.ts index c0fe8521da..34bf6d76c6 100644 --- a/packages/runtime-host/src/__tests__/web-fetch-tool.test.ts +++ b/packages/runtime-host/src/__tests__/web-fetch-tool.test.ts @@ -18,6 +18,7 @@ */ import assert from 'node:assert/strict'; +import { createServer } from 'node:http'; import { test } from 'node:test'; import { createDefaultRuntimePolicy } from '@maka/core/runtime-policy'; import type { MakaToolContext } from '@maka/runtime/tool-runtime'; @@ -26,7 +27,67 @@ import type { ResolveHostOutboundExecutionResult, RuntimePolicyOperationCoordinator, } from '@maka/storage/runtime-policy-stores'; -import { createHostWebFetchTool } from '../server/web-fetch-tool.js'; +import { createHostWebFetchService, createHostWebFetchTool } from '../server/web-fetch-tool.js'; + +test('health probes reject metadata before creating a transport', async () => { + const service = createHostWebFetchService({ + policy: resolver({ + kind: 'ready', + networkProxy: createDefaultRuntimePolicy().networkProxy, + secretMaterial: {}, + }), + createFetchTransport: () => { + throw new Error('must not create transport'); + }, + }); + for (const url of [ + 'http://169.254.169.254/latest/meta-data/', + 'http://metadata.google.internal/', + ]) { + await assert.rejects( + service.probe({ url, sessionId: 'session-1', abortSignal: new AbortController().signal }), + /metadata/, + ); + } +}); + +test('real health probe falls back to GET and bounds stalled responses', async () => { + const methods: string[] = []; + const server = createServer((req, res) => { + methods.push(req.method!); + if (req.url === '/stalled') return; + res.writeHead(req.method === 'HEAD' ? 405 : 200); + res.end('ready'); + }); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + assert.ok(address && typeof address !== 'string'); + const service = createHostWebFetchService({ + policy: resolver({ + kind: 'ready', + networkProxy: createDefaultRuntimePolicy().networkProxy, + secretMaterial: {}, + }), + probeTimeoutMs: 100, + }); + const input = { sessionId: 'session-1', abortSignal: new AbortController().signal }; + try { + assert.equal( + (await service.probe({ ...input, url: `http://127.0.0.1:${address.port}/ready` })).status, + 200, + ); + assert.deepEqual(methods, ['HEAD', 'GET']); + await assert.rejects( + service.probe({ ...input, url: `http://127.0.0.1:${address.port}/stalled` }), + /timed out/, + ); + } finally { + server.closeAllConnections(); + await new Promise((resolve, reject) => + server.close((error) => (error ? reject(error) : resolve())), + ); + } +}); test('Host WebFetch uses the resolved proxy snapshot and closes its transport', async () => { const networkProxy = { diff --git a/packages/runtime-host/src/server/execution-composition.ts b/packages/runtime-host/src/server/execution-composition.ts index b989d93587..bba8259bbf 100644 --- a/packages/runtime-host/src/server/execution-composition.ts +++ b/packages/runtime-host/src/server/execution-composition.ts @@ -251,6 +251,7 @@ import { shouldResolveHostTavilyWebSearchReadiness, } from './web-search-tool.js'; import { createHostWebFetchService, createHostWebFetchToolFromService } from './web-fetch-tool.js'; +import { buildBackgroundTaskHealthTool } from '@maka/runtime/background-task-health-tool'; import { createHostExecutionArtifactServices } from './execution-artifacts.js'; import { openToolResultArchiveEvidenceReader } from '@maka/storage/tool-result-archive-evidence'; import { @@ -653,6 +654,10 @@ export async function createExecutionRuntimeHostComposition( const webFetchService = createHostWebFetchService({ policy: runtimePolicyStores.operations, }); + const backgroundTaskHealthTool = buildBackgroundTaskHealthTool( + runtimeResources!, + webFetchService, + ); pluginWeb.bindRuntime({ search: ({ query, limit, abortSignal }) => webSearchService.search({ query, limit, ...(abortSignal ? { abortSignal } : {}) }), @@ -675,6 +680,7 @@ export async function createExecutionRuntimeHostComposition( const childHostTools = [ createHostWebSearchToolFromService(webSearchService), createHostWebFetchToolFromService(webFetchService), + backgroundTaskHealthTool, ...runtimePolicy.modelTools, ]; const hostTools = [...childHostTools, ...historyTools]; diff --git a/packages/runtime-host/src/server/web-fetch-tool.ts b/packages/runtime-host/src/server/web-fetch-tool.ts index 956af4b517..76b17b7460 100644 --- a/packages/runtime-host/src/server/web-fetch-tool.ts +++ b/packages/runtime-host/src/server/web-fetch-tool.ts @@ -18,7 +18,7 @@ */ import { buildWebFetchTool } from '@maka/runtime/web-fetch-tool'; -import { createLocalWebFetchExecutor } from '@maka/runtime/local-web-fetch'; +import { assertAllowedTarget, createLocalWebFetchExecutor } from '@maka/runtime/local-web-fetch'; import { createProxiedFetchTransport, type ProxiedFetchProxy, @@ -29,6 +29,7 @@ import type { RuntimePolicyOperationCoordinator } from '@maka/storage/runtime-po import { toRuntimePolicyProxy } from './runtime-policy-proxy.js'; interface HostWebFetchServiceInput { + readonly probeTimeoutMs?: number; readonly policy: Pick; readonly createFetchTransport?: (proxy: ProxiedFetchProxy | null) => ProxiedFetchTransport; } @@ -39,6 +40,11 @@ export interface HostWebFetchService { readonly sessionId: string; readonly abortSignal?: AbortSignal; }): Promise; + probe(input: { + url: string; + sessionId: string; + abortSignal: AbortSignal; + }): Promise<{ status: number; statusText?: string; elapsedMs: number }>; } export function createHostWebFetchService(input: HostWebFetchServiceInput): HostWebFetchService { @@ -65,6 +71,46 @@ export function createHostWebFetchService(input: HostWebFetchServiceInput): Host await transport.close(); } }, + probe: async ({ url, abortSignal }) => { + const parsed = new URL(url); + assertAllowedTarget(parsed); + abortSignal.throwIfAborted(); + const resolved = await input.policy.resolveHostOutboundExecution(); + if (resolved.kind === 'privacy_mode') + throw new Error('Endpoint health checks are disabled while privacy mode is active.'); + if (resolved.kind === 'credential_not_configured') + throw new Error('Configure the network proxy credential before checking an endpoint.'); + const transport = createFetchTransport( + toRuntimePolicyProxy(resolved.networkProxy, resolved.secretMaterial.networkProxy?.secret), + ); + const started = Date.now(); + const timeout = new AbortController(); + const timer = setTimeout( + () => timeout.abort(new Error('Endpoint health probe timed out.')), + input.probeTimeoutMs ?? 30_000, + ); + const signal = AbortSignal.any([abortSignal, timeout.signal]); + try { + let response = await transport.fetch(parsed, { + method: 'HEAD', + redirect: 'manual', + signal, + }); + await response.body?.cancel(); + if (response.status === 405 || response.status === 501) { + response = await transport.fetch(parsed, { method: 'GET', redirect: 'manual', signal }); + await response.body?.cancel(); + } + return { + status: response.status, + ...(response.statusText ? { statusText: response.statusText } : {}), + elapsedMs: Date.now() - started, + }; + } finally { + clearTimeout(timer); + await transport.close(); + } + }, }; } diff --git a/packages/runtime/README.md b/packages/runtime/README.md index cd3bc72a32..c6cd22f429 100644 --- a/packages/runtime/README.md +++ b/packages/runtime/README.md @@ -42,6 +42,18 @@ The main integration points are: Shared execution composition — where `BackendRegistry` and `SessionManager` are constructed — lives in the Runtime Host at [`packages/runtime-host/src/server/execution-composition.ts`](../runtime-host/src/server/execution-composition.ts). Clients, including Desktop, execute Maka through Runtime Host rather than composing Runtime directly. +## Background-task readiness + +`Bash` background runs return a durable runtime-task ref and expose the native +process id when available. A process id published after startup is captured on +the next task observation, output flush, or finalization. +`BackgroundTaskHealth` deliberately keeps +the process lifecycle (`starting`, `running`, or terminal, with timestamps and +captured output) separate from endpoint readiness. An endpoint is `healthy` +only after an explicit HTTP(S) probe succeeds; an omitted probe is +`not_checked`, and a failed or policy-blocked probe is `unknown`. Consumers must +not infer HTTP readiness from the process status alone. + ## Extension rules - Add backend behavior behind `AgentBackend` and register it through the existing registry. diff --git a/packages/runtime/package.json b/packages/runtime/package.json index 5c27fdefa1..915c7c8f4d 100644 --- a/packages/runtime/package.json +++ b/packages/runtime/package.json @@ -13,6 +13,7 @@ "./builtin-tools": "./dist/builtin-tools.js", "./shell-tools": "./dist/shell-tools.js", "./shell-run-manager": "./dist/shell-run-manager.js", + "./background-task-health-tool": "./dist/background-task-health-tool.js", "./deep-research-tools": "./dist/deep-research-tools.js", "./durable-tool-result-projection": "./dist/durable-tool-result-projection.js", "./tool-artifacts": "./dist/tool-artifacts.js", diff --git a/packages/runtime/src/__tests__/background-task-health-tool.test.ts b/packages/runtime/src/__tests__/background-task-health-tool.test.ts new file mode 100644 index 0000000000..d13c0815b1 --- /dev/null +++ b/packages/runtime/src/__tests__/background-task-health-tool.test.ts @@ -0,0 +1,119 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import assert from 'node:assert/strict'; +import { test } from 'node:test'; +import { buildBackgroundTaskHealthTool } from '../background-task-health-tool.js'; + +const context = { + sessionId: 'session-1', + turnId: 'turn-1', + toolCallId: 'tool-1', + cwd: '/tmp', + abortSignal: new AbortController().signal, +} as any; +const shell = (status: string, pid?: number) => ({ + kind: 'shell_run', + ref: 'maka://runtime/background-tasks/run-1', + status, + mode: 'pipes', + cwd: '/tmp', + cmd: 'python -m http.server', + startedAt: 1, + updatedAt: 2, + revision: 2, + ...(pid ? { pid } : {}), +}); + +test('reports process tracking separately when endpoint is not checked', async () => { + const tool = buildBackgroundTaskHealthTool( + { readRuntimeResource: async () => shell('running', 1234) } as any, + { + probe: async () => { + throw new Error('must not probe'); + }, + }, + ); + assert.deepEqual( + JSON.parse(String(await tool.impl({ ref: 'maka://runtime/background-tasks/run-1' }, context))), + { + process: { status: 'running', tracked: true, startedAt: 1, updatedAt: 2, pid: 1234 }, + endpoint: { state: 'not_checked' }, + }, + ); +}); + +test('reports endpoint health only from the probe result', async () => { + let called = 0; + const tool = buildBackgroundTaskHealthTool( + { readRuntimeResource: async () => shell('running', 1234) } as any, + { + probe: async () => { + called += 1; + return { status: 204, statusText: 'No Content', elapsedMs: 4 }; + }, + }, + ); + assert.deepEqual( + JSON.parse( + String( + await tool.impl( + { ref: 'maka://runtime/background-tasks/run-1', url: 'http://127.0.0.1:8765/' }, + context, + ), + ), + ), + { + process: { status: 'running', tracked: true, startedAt: 1, updatedAt: 2, pid: 1234 }, + endpoint: { + state: 'checked', + httpStatus: 204, + elapsedMs: 4, + target: 'http://127.0.0.1:8765/', + health: 'healthy', + }, + }, + ); + assert.equal(called, 1); +}); + +test('does not convert a failed probe into a ready claim', async () => { + const tool = buildBackgroundTaskHealthTool( + { readRuntimeResource: async () => shell('running', 1234) } as any, + { + probe: async () => { + throw new Error('connection refused'); + }, + }, + ); + assert.deepEqual( + JSON.parse( + String( + await tool.impl( + { ref: 'maka://runtime/background-tasks/run-1', url: 'http://127.0.0.1:8765/' }, + context, + ), + ), + ), + { + process: { status: 'running', tracked: true, startedAt: 1, updatedAt: 2, pid: 1234 }, + endpoint: { state: 'unknown', target: 'http://127.0.0.1:8765/', error: 'connection refused' }, + }, + ); +}); diff --git a/packages/runtime/src/__tests__/shell-run-manager.test.ts b/packages/runtime/src/__tests__/shell-run-manager.test.ts index 212ea60474..d798e9b54a 100644 --- a/packages/runtime/src/__tests__/shell-run-manager.test.ts +++ b/packages/runtime/src/__tests__/shell-run-manager.test.ts @@ -38,6 +38,7 @@ import { type ShellRunUpdate, type ToolResultContent } from '@maka/core/events'; import { createSqliteShellRunStore } from '@maka/storage/shell-run-store'; import { ShellRunProcessManager } from '../shell-run-manager.js'; +import { buildBackgroundTaskHealthTool } from '../background-task-health-tool.js'; import { ShellRunPtyControlClosedError, type ShellRunPtyDataEvent, @@ -328,6 +329,94 @@ describe('ShellRunProcessManager', () => { ); }); + for (const observation of ['read', 'exit'] as const) { + test(`persists a late ConPTY PID on ${observation} without new output`, async (t) => { + const cwd = await workspace(); + const exitGate = join(cwd, 'exit-gate'); + const store = sqliteShellRunStore(cwd); + const flushes = manualFlushScheduler(); + const manager = createManager(store, undefined, { scheduleFlush: flushes.schedule }); + const nativePid = Object.getOwnPropertyDescriptor(PtyProcessDriver.prototype, 'pid')!.get!; + let publishPid = false; + let driver: PtyProcessDriver | undefined; + t.mock.getter(PtyProcessDriver.prototype, 'pid', function (this: PtyProcessDriver) { + driver = this; + return publishPid ? nativePid.call(this) : 0; + }); + let ref: string | undefined; + try { + const initial = await manager.runBackgroundBash( + shellInput({ + cwd, + command: nodeCommand(` + const { existsSync } = require('node:fs'); + process.stdout.write('READY\\n'); + setInterval(() => { + if (existsSync(${JSON.stringify(exitGate)})) process.exit(0); + }, 10); + `), + pty: true, + timeoutMs: 30_000, + }), + ); + ref = initial.ref; + assert.equal(initial.status, 'running'); + assert.equal(initial.pid, undefined); + await waitForPtyText(manager, ref, /READY/, 15_000); + const before = await store.readShellRun('session-1', 'shell-run-1'); + assert.equal(before.pid, undefined); + assert.ok(driver); + const expectedPid = nativePid.call(driver); + assert.ok(expectedPid > 0, 'the native PTY has published its real PID'); + publishPid = true; + + if (observation === 'exit') { + await writeFile(exitGate, 'exit'); + await waitUntil(() => manager.liveCount() === 0, 15_000); + } + const result = await manager.readRuntimeResource('session-1', ref, NO_ABORT); + assertShellRun(result); + assert.equal(result.pid, expectedPid); + assert.equal(result.status, observation === 'exit' ? 'completed' : 'running'); + const stored = await store.readShellRun('session-1', 'shell-run-1'); + assert.equal(stored.pid, expectedPid); + assert.deepEqual(stored.output, before.output); + const tool = buildBackgroundTaskHealthTool(manager, { + probe: async () => { + throw new Error('must not probe'); + }, + }); + const health = JSON.parse( + String( + await tool.impl( + { ref }, + { + sessionId: 'session-1', + turnId: 'turn-1', + toolCallId: 'health-1', + cwd, + abortSignal: NO_ABORT, + emitOutput: () => {}, + }, + ), + ), + ); + assert.equal(health.process.pid, expectedPid); + assert.deepEqual(health.endpoint, { state: 'not_checked' }); + if (observation === 'read') { + const repeated = await manager.readRuntimeResource('session-1', ref, NO_ABORT); + assertShellRun(repeated); + assert.equal(repeated.revision, stored.revision); + } + } finally { + t.mock.restoreAll(); + if (ref && manager.liveCount() > 0) { + await manager.stopBackgroundTask('session-1', ref, NO_ABORT); + } + } + }); + } + test('hands off a long pipe command without output and publishes monotonic revisions', async () => { const updates: ShellRunUpdate[] = []; const store = sqliteShellRunStore(await workspace()); @@ -342,6 +431,7 @@ describe('ShellRunProcessManager', () => { assert.equal(initial.kind, 'shell_run'); assert.equal(initial.mode, 'pipes'); assert.equal(initial.output, undefined); + assert.ok(initial.pid !== undefined && initial.pid > 0); assert.equal((await store.readShellRun('session-1', 'shell-run-1')).timeoutMs, undefined); await waitForShellRun( manager, @@ -357,6 +447,7 @@ describe('ShellRunProcessManager', () => { assert.ok(runningUpdate); const running = await manager.readRuntimeResource('session-1', initial.ref, NO_ABORT); assertShellRun(running); + assert.equal(running.pid, initial.pid); assert.equal(running.output?.mode, 'pipes'); if (running.output?.mode !== 'pipes') throw new Error('expected pipes output'); assert.equal(running.output.stdout, 'start'); diff --git a/packages/runtime/src/background-task-health-tool.ts b/packages/runtime/src/background-task-health-tool.ts new file mode 100644 index 0000000000..959cdfd8a3 --- /dev/null +++ b/packages/runtime/src/background-task-health-tool.ts @@ -0,0 +1,137 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import { z } from 'zod'; +import type { ToolResultContent } from '@maka/core/events'; +import type { MakaTool } from './tool-runtime.js'; + +export interface BackgroundTaskHealthReader { + readRuntimeResource( + sessionId: string, + ref: string, + abortSignal: AbortSignal, + ): Promise; +} + +export interface BackgroundTaskEndpointProbe { + probe(input: { url: string; sessionId: string; abortSignal: AbortSignal }): Promise<{ + status: number; + statusText?: string; + elapsedMs: number; + }>; +} + +/** + * Produces an explicit two-axis result. A tracked process is never described + * as endpoint-ready unless the caller supplied a URL and the probe succeeded. + */ +export function buildBackgroundTaskHealthTool( + reader: BackgroundTaskHealthReader, + probe: BackgroundTaskEndpointProbe, +): MakaTool { + return { + name: 'BackgroundTaskHealth', + displayName: 'Background task health', + categoryHint: 'web_read', + description: + 'Check a tracked background task and an optional HTTP endpoint. Uses HEAD with one GET fallback for 405/501; discards the body. Reports HTTP status only, not browser loading or ownership of the listener. Redirects are not followed. Logs are omitted by default; use Read(ref) for full logs.', + parameters: z + .object({ + ref: z.string().describe('The maka://runtime/background-tasks/ ref returned by Bash'), + include_logs: z.boolean().optional().describe('Include captured task logs in this report'), + url: z + .string() + .url() + .refine( + (value) => ['http:', 'https:'].includes(new URL(value).protocol), + 'Health endpoint must use HTTP or HTTPS', + ) + .optional() + .describe('The HTTP or HTTPS endpoint to probe'), + }) + .strict(), + impl: async ({ ref, url, include_logs }, context) => { + const resource = await reader.readRuntimeResource( + context.sessionId, + ref, + context.abortSignal, + ); + if ( + !resource || + typeof resource !== 'object' || + Array.isArray(resource) || + resource.kind !== 'shell_run' + ) { + throw new Error('BackgroundTaskHealth requires a shell_run runtime resource'); + } + const shell = resource as Extract; + const process = { + status: shell.status, + tracked: true, + startedAt: shell.startedAt, + updatedAt: shell.updatedAt, + ...(shell.pid !== undefined ? { pid: shell.pid } : {}), + ...(shell.completedAt !== undefined ? { completedAt: shell.completedAt } : {}), + ...(shell.failureMessage !== undefined ? { failureMessage: shell.failureMessage } : {}), + ...(include_logs && shell.output + ? { + logs: + shell.output.mode === 'pipes' + ? { stdout: shell.output.stdout, stderr: shell.output.stderr } + : { screen: shell.output.screen, scrollback: shell.output.scrollback }, + } + : {}), + }; + if (!url) return JSON.stringify({ process, endpoint: { state: 'not_checked' } }); + let endpoint; + try { + endpoint = await probe.probe({ + url, + sessionId: context.sessionId, + abortSignal: context.abortSignal, + }); + } catch (error) { + context.abortSignal.throwIfAborted(); + return JSON.stringify({ + process, + endpoint: { + state: 'unknown', + target: new URL(url).href, + error: error instanceof Error ? error.message : String(error), + }, + }); + } + return JSON.stringify({ + process, + endpoint: { + state: 'checked', + httpStatus: endpoint.status, + elapsedMs: endpoint.elapsedMs, + target: new URL(url).href, + health: + endpoint.status >= 200 && endpoint.status < 300 + ? 'healthy' + : endpoint.status < 400 + ? 'unknown' + : 'unhealthy', + }, + }); + }, + }; +} diff --git a/packages/runtime/src/local-web-fetch.ts b/packages/runtime/src/local-web-fetch.ts index 069644d1c0..a465049e5b 100644 --- a/packages/runtime/src/local-web-fetch.ts +++ b/packages/runtime/src/local-web-fetch.ts @@ -132,7 +132,7 @@ function responseLimitError(): Error { return new Error('WebFetch response exceeds the 5 MB response limit.'); } -function assertAllowedTarget(url: URL): void { +export function assertAllowedTarget(url: URL): void { if (url.protocol !== 'http:' && url.protocol !== 'https:') { throw new Error('WebFetch URL must use HTTP or HTTPS.'); } diff --git a/packages/runtime/src/shell-run-manager.ts b/packages/runtime/src/shell-run-manager.ts index f9ad47bd31..5fc60c9839 100644 --- a/packages/runtime/src/shell-run-manager.ts +++ b/packages/runtime/src/shell-run-manager.ts @@ -972,6 +972,7 @@ export class ShellRunProcessManager private async markRunning(live: LiveShellRun): Promise { live.record = await this.input.store.updateShellRun(live.sessionId, live.shellRunId, { status: 'running', + ...this.processPidPatch(live), output: (await this.snapshotAtCut(live, false)).output, updatedAt: this.input.now(), }); @@ -982,6 +983,11 @@ export class ShellRunProcessManager } } + private processPidPatch(live: LiveShellRun): Pick { + const pid = live.driver.pid; + return pid !== undefined && Number.isSafeInteger(pid) && pid > 0 ? { pid } : {}; + } + private onPipeData(live: LivePipeShellRun, stream: 'stdout' | 'stderr', data: string): void { if (live.driverExit || live.finalizeOnce) return; live.collector.accept(stream, data); @@ -1141,12 +1147,13 @@ export class ShellRunProcessManager failureStage = 'persist'; if (live.persistFailure && !options.bestEffort) throw live.persistFailure; const current = live.record; - const candidate: ShellRunRecord = { ...current, ...patch, output: snapshot.output }; + // ConPTY can publish its PID after admission, even without new output. + const update = { ...patch, ...this.processPidPatch(live), output: snapshot.output }; + const candidate: ShellRunRecord = { ...current, ...update }; let updated = current; if (!isDeepStrictEqual(candidate, current)) { updated = await this.input.store.updateShellRun(live.sessionId, live.shellRunId, { - ...patch, - output: snapshot.output, + ...update, updatedAt: this.input.now(), }); live.record = updated; diff --git a/packages/runtime/src/shell-run-tool-result.ts b/packages/runtime/src/shell-run-tool-result.ts index c7def411c2..884c6ae9ec 100644 --- a/packages/runtime/src/shell-run-tool-result.ts +++ b/packages/runtime/src/shell-run-tool-result.ts @@ -146,6 +146,7 @@ function shellRunStateContent(record: ShellRunRecord): ShellRunCompactResult { cmd: record.command, startedAt: record.startedAt, updatedAt: record.updatedAt, + ...(record.pid !== undefined ? { pid: record.pid } : {}), ...(record.completedAt !== undefined ? { completedAt: record.completedAt } : {}), ...(record.timeoutMs !== undefined ? { timeoutMs: record.timeoutMs } : {}), ...(record.exitCode !== undefined ? { exitCode: record.exitCode } : {}),