From db9cc72a3851f40c40aecd65f0d29c9e094bff8d Mon Sep 17 00:00:00 2001 From: DJJ Date: Tue, 15 Sep 2026 19:37:48 +0800 Subject: [PATCH 1/5] =?UTF-8?q?feat(device):=20=E8=AE=A9=E9=95=BF=E5=91=BD?= =?UTF-8?q?=E4=BB=A4=E5=9C=A8=E8=AF=B7=E6=B1=82=E7=BB=93=E6=9D=9F=E5=90=8E?= =?UTF-8?q?=E4=BB=8D=E5=8F=AF=E8=A7=82=E5=AF=9F=E5=B9=B6=E6=8C=89=E8=B0=83?= =?UTF-8?q?=E7=94=A8=E6=96=B9=E9=9A=94=E7=A6=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 会话显式绑定现有命令配置,复用可信 caller 和 realtime 通道;独立日志游标、短期启动去重与 runtime 关闭清理避免重新执行和跨 owner 读取。设备摘要沿现有 help 投影,HTTP/MCP 用真实输入输出检查契约。 升级行为变化:受限 shell allowlist 现在拒绝 LF/CR,堵住换行执行第二条命令的旁路;单值 * 仍保留显式任意命令语义。新增能力按 0.x minor 提升 app/cli/sdk/server,server 内联 app 所以产物也变化。 --- packages/app/package.json | 2 +- packages/app/src/deviceHello.ts | 14 +- packages/app/src/federation.ts | 1 + packages/app/src/helpModel.ts | 8 +- .../contracts/processSessionCoverage.test.ts | 210 ++++++++ .../test/deviceContext.integration.test.ts | 108 ++++ .../cli/examples/process-session-profile.json | 21 + packages/cli/package.json | 2 +- packages/cli/src/commands/connect.ts | 8 + packages/cli/src/commands/daemon.ts | 1 + packages/cli/src/commands/device.ts | 2 + packages/cli/src/commands/deviceSession.ts | 179 +++++++ packages/cli/src/daemon.ts | 8 + packages/cli/src/deviceRuntime.ts | 48 +- packages/cli/src/deviceSessions.ts | 50 ++ packages/cli/test/deviceRuntime.test.ts | 3 + packages/cli/test/deviceSession.test.ts | 132 +++++ packages/cli/test/deviceSessions.test.ts | 127 +++++ packages/cli/test/strictParsing.test.ts | 4 + packages/core/src/device/environment.ts | 11 + packages/core/src/device/frames.ts | 2 + packages/core/src/device/helpModel.ts | 10 +- .../core/src/device/processSessionContract.ts | 49 ++ packages/core/src/device/public.ts | 11 + packages/core/src/device/shellAllow.ts | 4 +- packages/core/src/htbp/helpDsl.ts | 2 + packages/core/src/htbp/helpMarkdown.ts | 5 + packages/core/src/htbp/model.ts | 12 +- packages/core/src/index.ts | 1 + packages/core/src/node/index.ts | 1 + packages/core/src/node/processSessions.ts | 464 ++++++++++++++++++ packages/core/src/node/shellExecutor.ts | 30 ++ packages/core/src/node/structuredCommand.ts | 39 ++ packages/core/src/protocol/wire.ts | 13 + packages/core/src/tree/registry.ts | 5 +- packages/core/src/types.ts | 16 +- packages/core/test/device/shellAllow.test.ts | 5 + .../core/test/node/processSessions.test.ts | 240 +++++++++ packages/core/test/node/shellExecutor.test.ts | 10 + packages/sdk/package.json | 2 +- packages/sdk/src/client/index.ts | 10 + packages/sdk/src/device/connection.ts | 4 + packages/sdk/src/device/index.ts | 1 + packages/sdk/test/client/client.test.ts | 13 + packages/sdk/test/device.test.ts | 16 + packages/server/package.json | 2 +- 46 files changed, 1890 insertions(+), 16 deletions(-) create mode 100644 packages/app/test/contracts/processSessionCoverage.test.ts create mode 100644 packages/app/test/deviceContext.integration.test.ts create mode 100644 packages/cli/examples/process-session-profile.json create mode 100644 packages/cli/src/commands/deviceSession.ts create mode 100644 packages/cli/src/deviceSessions.ts create mode 100644 packages/cli/test/deviceSession.test.ts create mode 100644 packages/cli/test/deviceSessions.test.ts create mode 100644 packages/core/src/device/environment.ts create mode 100644 packages/core/src/device/processSessionContract.ts create mode 100644 packages/core/src/node/processSessions.ts create mode 100644 packages/core/test/node/processSessions.test.ts diff --git a/packages/app/package.json b/packages/app/package.json index 05985b9f..a7b9061d 100644 --- a/packages/app/package.json +++ b/packages/app/package.json @@ -1,6 +1,6 @@ { "name": "@tool-bridge/app", - "version": "0.22.0", + "version": "0.23.0", "description": "Host-neutral HTBP application: assemble the tool-bridge tree on any runtime by injecting state, object, secret, device-channel and search adapters", "type": "module", "license": "MIT", diff --git a/packages/app/src/deviceHello.ts b/packages/app/src/deviceHello.ts index 86978890..ddf5400f 100644 --- a/packages/app/src/deviceHello.ts +++ b/packages/app/src/deviceHello.ts @@ -4,6 +4,7 @@ import { type CallContext, check, checkRegisterPath, + deviceEnvironmentSchema, type DeviceExpose, type DeviceNodeInput, identify, @@ -199,6 +200,12 @@ export async function processDeviceHello(opts: { } if (hello.expose.fs !== undefined) assertFsRoots(hello.expose.fs.roots) + const environment = hello.expose.environment === undefined + ? undefined + : deviceEnvironmentSchema.safeParse(hello.expose.environment) + if (environment !== undefined && !environment.success) { + throw new TBError('invalid_argument', 'device environment contains invalid or unsupported fields') + } const inputs = nodesForHello(mountPath, hello.deviceId, hello.expose) const registry = new NodeRegistryStore(store) for (const input of inputs) { @@ -226,7 +233,12 @@ export async function processDeviceHello(opts: { authCtx.keyId, now, input.path === mountPath - ? { deviceId: hello.deviceId, online: true, lastSeenAt: now } + ? { + deviceId: hello.deviceId, + online: true, + lastSeenAt: now, + ...(environment === undefined ? {} : { deviceEnvironment: environment.data, deviceReportedAt: now }), + } : {}, ) } catch (error) { diff --git a/packages/app/src/federation.ts b/packages/app/src/federation.ts index 0f84a470..686841ec 100644 --- a/packages/app/src/federation.ts +++ b/packages/app/src/federation.ts @@ -150,6 +150,7 @@ export class RemotePathProjector { node: { ...help.node, path: nodePath }, cmds, ...(children === undefined ? {} : { children }), + ...(help.deviceContext === undefined ? {} : { deviceContext: help.deviceContext }), ...(help.feedback === undefined ? {} : { feedback: help.feedback }), ...(help.hint === undefined ? {} : { hint: help.hint }), ...(help.note === undefined ? {} : { note: help.note }), diff --git a/packages/app/src/helpModel.ts b/packages/app/src/helpModel.ts index 599c4f31..a599e9d2 100644 --- a/packages/app/src/helpModel.ts +++ b/packages/app/src/helpModel.ts @@ -79,7 +79,13 @@ export async function helpModelFor( now: opts.now, }) return deviceDirectoryHelpModel( - { path: node.path, description: node.description, presence }, + { + path: node.path, + description: node.description, + presence, + ...(node.deviceEnvironment === undefined ? {} : { environment: node.deviceEnvironment }), + ...(node.deviceReportedAt === undefined ? {} : { reportedAt: node.deviceReportedAt }), + }, children.map(n => ({ path: n.path, kind: n.kind, description: n.description })), ) } diff --git a/packages/app/test/contracts/processSessionCoverage.test.ts b/packages/app/test/contracts/processSessionCoverage.test.ts new file mode 100644 index 00000000..69e4647e --- /dev/null +++ b/packages/app/test/contracts/processSessionCoverage.test.ts @@ -0,0 +1,210 @@ +import { createProcessSessionManager, type ProcessSessionBinding, processSessionCommands, type ProcessSessionManager } from '@tool-bridge/core/node' +import { type HelpJson, isTBError, type Scope, TBError } from '@tool-bridge/core' +import { afterEach, describe, expect, it, vi } from 'vitest' +import { helpJsonSchema } from '@tool-bridge/core/protocol' +import { execPath } from 'node:process' +import { z } from 'zod/v4' +import type { DeviceInvokeRequest } from '../../src/deps' +import { ToolJsonSchemaValidator } from '../../src/jsonSchemaValidator' +import { bearer, createTestApp, type TestApp } from '../harness' +import { processDeviceHello } from '../../src/deviceHello' +import { connectTestMcpClient } from '../mcpClient' +import { TEST_ADMIN_SK } from '../fixtures' + +const NODE_PATH = 'device/contract-host/sessions/command' +const BINDING_PATH = 'sessions/command' +const binding: ProcessSessionBinding = { + path: BINDING_PATH, + description: 'Controlled contract-test process', + effect: 'read', + inputSchema: z.strictObject({}), + maxRuntimeMs: 10_000, + prepare: () => ({ + executable: execPath, + argv: ['-e', 'process.stdout.write("contract-ready\\n"); setInterval(() => {}, 1000)'], + }), +} +const managers = new Set() +afterEach(async () => { + await Promise.all([...managers].map(manager => manager.close())) + managers.clear() +}) + +async function post(tb: TestApp, path: string, body: unknown, secret = TEST_ADMIN_SK): Promise { + return await tb.request(`https://tb.test/${path}`, bearer(secret, { + method: 'POST', headers: { 'content-type': 'application/json', 'accept': 'application/json' }, + body: JSON.stringify(body), + })) +} +async function issue(tb: TestApp, owner: string, scopes: Scope[]): Promise { + const response = await post(tb, 'system/sk/write', { owner, scopes }) + expect(response.status).toBe(200) + return (await response.json() as { secret: string }).secret +} +async function setup() { + const manager = createProcessSessionManager({ platform: 'linux', terminationGraceMs: 25 }) + managers.add(manager) + manager.register(binding) + const invoke = vi.fn(async (_deviceId: string, request: DeviceInvokeRequest) => { + const parts = request.path.split('/') + const action = parts.pop()! + try { + const value = await manager.invoke(parts.join('/'), action, request.arguments, + request.context === undefined ? {} : { caller: request.context.caller }) + return { disposition: 'completed' as const, result: { ok: true as const, value } } + } catch (error) { + const failure = isTBError(error) ? error : new TBError('internal', 'test device failed') + return { disposition: 'completed' as const, result: { ok: false as const, error: failure.toJSON() } } + } + }) + const tb = await createTestApp({ objects: null, device: { ws: async () => new Response(null, { status: 501 }), invoke } }) + await processDeviceHello({ + store: tb.state, authorization: `Bearer ${TEST_ADMIN_SK}`, deviceIdHint: 'contract-host', + hello: { deviceId: 'contract-host', expose: { nodes: [{ + path: BINDING_PATH, kind: 'tool', description: binding.description!, cmds: processSessionCommands(binding), + }] } }, + }) + const scopes: Scope[] = [{ pattern: 'device/contract-host/**', actions: ['read', 'call'] }] + const alice = await issue(tb, 'agent:alice', scopes) + const bob = await issue(tb, 'agent:bob', scopes) + const reader = await issue(tb, 'agent:reader', [{ pattern: 'device/contract-host/**', actions: ['read'] }]) + const response = await tb.request(`https://tb.test/${NODE_PATH}/~help?schemas=1`, bearer(alice, { headers: { accept: 'application/json' } })) + expect(response.status).toBe(200) + const help = helpJsonSchema.parse(await response.json()) + return { tb, manager, invoke, help, alice, bob, reader } +} +type Harness = Awaited> +function startArgs(manager: ProcessSessionManager, requestId: string) { + return { input: {}, expectedRuntimeId: manager.runtimeId, requestId, startBefore: new Date(Date.now() + 60_000).toISOString() } +} +async function call(h: Harness, name: string, args: Record, secret = h.alice): Promise> { + const response = await post(h.tb, `${NODE_PATH}/${name}`, args, secret) + expect(response.status).toBe(200) + return await response.json() as Record +} +function validateOutput(help: HelpJson, name: string, output: unknown): void { + const command = help.cmds.find(candidate => candidate.name === name) + expect(command, `missing advertised command ${name}`).toBeDefined() + expect(command!.outputSchema, `${name} needs a public output schema`).toBeDefined() + const validate = new ToolJsonSchemaValidator().getValidator(command!.outputSchema as Record) + expect(validate(output)).toMatchObject({ valid: true }) +} + +interface ScenarioContext { harness: Harness, sessionId?: string } +type Scenario = (context: ScenarioContext) => Promise +// The same executable map drives coverage comparison and actual HTTP execution. +// Adding a command requires a scenario; merely registering a name cannot satisfy the gate. +const scenarios: Record = { + async start(context) { + const result = await call(context.harness, 'start', startArgs(context.harness.manager, 'coverage-start')) + expect(result.state).toBe('running') + context.sessionId = String(result.sessionId) + return result + }, + async list(context) { + const result = await call(context.harness, 'list', {}) + expect(result.items).toEqual([expect.objectContaining({ sessionId: context.sessionId })]) + return result + }, + async observe(context) { + const result = await call(context.harness, 'observe', { sessionId: context.sessionId, waitMs: 1000 }) + expect(result.chunks).toEqual(expect.arrayContaining([expect.objectContaining({ stream: 'stdout', text: 'contract-ready\n' })])) + return result + }, + async stop(context) { + const result = await call(context.harness, 'stop', { sessionId: context.sessionId }) + expect(result.state).toBe('stopped') + return result + }, +} +function assertCoverage(help: HelpJson, cases: Record): void { + expect(Object.keys(cases).sort(), 'every advertised process command needs an executed output-contract scenario') + .toEqual(help.cmds.map(command => command.name).sort()) +} + +describe('process session public command contract coverage', () => { + it('executes every advertised command through HTTP and validates actual output against actual help', async () => { + const harness = await setup() + expect(harness.help.cmds.map(command => command.name)).toEqual(harness.manager.cmds(BINDING_PATH).map(command => command.name)) + assertCoverage(harness.help, scenarios) + const context: ScenarioContext = { harness } + const executed: string[] = [] + for (const [name, run] of Object.entries(scenarios)) { + validateOutput(harness.help, name, await run(context)) + executed.push(name) + } + expect(executed.sort()).toEqual(harness.help.cmds.map(command => command.name).sort()) + expect(harness.invoke).toHaveBeenCalledTimes(executed.length) + expect(() => validateOutput(harness.help, 'list', { items: 'invalid output' })).toThrow() + }) + + it('fails the gate when any public command scenario is deleted', async () => { + const { help } = await setup() + for (const command of help.cmds) { + const incomplete = { ...scenarios } + delete incomplete[command.name] + expect(() => assertCoverage(help, incomplete), command.name).toThrow() + } + }) + + it('enforces current gateway scopes, owner isolation, strict arguments and realtime-only delivery', async () => { + const h = await setup() + const start = await call(h, 'start', startArgs(h.manager, 'owner-check')) + const sessionId = String(start.sessionId) + expect((await call(h, 'list', {}, h.bob)).items).toEqual([]) + for (const command of ['observe', 'stop']) { + const response = await post(h.tb, `${NODE_PATH}/${command}`, { sessionId }, h.bob) + expect(response.status).toBe(404) + expect(await response.json()).toMatchObject({ code: 'not_found' }) + } + const beforeUnauthorized = h.invoke.mock.calls.length + const denied = await post(h.tb, `${NODE_PATH}/list`, {}, h.reader) + expect(denied.status).toBe(403) + expect(h.invoke.mock.calls).toHaveLength(beforeUnauthorized) + for (const [command, args] of [ + ['start', { ...startArgs(h.manager, 'bad-input'), input: { executable: 'arbitrary' } }], + ['list', { owner: 'agent:alice' }], + ['observe', { sessionId, cursor: -1 }], + ['stop', { sessionId, owner: 'agent:alice' }], + ] as const) { + const response = await post(h.tb, `${NODE_PATH}/${command}`, args, h.alice) + expect(response.status).toBe(400) + expect(await response.json()).toMatchObject({ code: 'invalid_argument' }) + } + const beforeMailbox = h.invoke.mock.calls.length + const mailboxArguments: Record> = { + start: startArgs(h.manager, 'mailbox-must-not-start'), list: {}, observe: { sessionId }, stop: { sessionId }, + } + for (const command of h.help.cmds) { + expect(command.delivery).toBe('realtime') + const response = await post(h.tb, `${NODE_PATH}/${command.name}`, { ...mailboxArguments[command.name], '~delivery': 'mailbox' }, h.alice) + expect(response.status).toBe(400) + expect(await response.json()).toMatchObject({ code: 'invalid_argument', message: 'device command does not support mailbox delivery' }) + } + expect(h.invoke.mock.calls).toHaveLength(beforeMailbox) + expect((await call(h, 'observe', { sessionId })).state).toBe('running') + await expect(h.manager.invoke(BINDING_PATH, 'list', {})).rejects.toMatchObject({ code: 'unavailable' }) + }) + + it('MCP discovers the same command contract and observes the same owner session', async () => { + const h = await setup() + const start = await call(h, 'start', startArgs(h.manager, 'mcp-owner')) + const client = await connectTestMcpClient('https://tb.test/~mcp', h.alice, (input, init) => h.tb.request(input, init)) + try { + const result = await client.callTool({ name: 'tb_help', arguments: { path: NODE_PATH, format: 'json', schemas: true } }) + expect(result.isError).not.toBe(true) + const help = helpJsonSchema.parse(result.structuredContent) + assertCoverage(help, scenarios) + const listed = await client.callTool({ name: 'tb_call', arguments: { path: `${NODE_PATH}/list`, args: {} } }) + expect(listed.isError).not.toBe(true) + validateOutput(help, 'list', listed.structuredContent) + expect(listed.structuredContent).toMatchObject({ items: [expect.objectContaining({ sessionId: start.sessionId })] }) + const observed = await client.callTool({ name: 'tb_call', arguments: { path: `${NODE_PATH}/observe`, args: { sessionId: start.sessionId, waitMs: 1000 } } }) + expect(observed.isError).not.toBe(true) + validateOutput(help, 'observe', observed.structuredContent) + expect(observed.structuredContent).toMatchObject({ sessionId: start.sessionId, state: 'running' }) + } finally { + await client.close() + } + }) +}) diff --git a/packages/app/test/deviceContext.integration.test.ts b/packages/app/test/deviceContext.integration.test.ts new file mode 100644 index 00000000..86476adb --- /dev/null +++ b/packages/app/test/deviceContext.integration.test.ts @@ -0,0 +1,108 @@ +import { decodeDeviceFrame, type DeviceEnvironment, encodeDeviceFrame, type HelpJson, NodeRegistryStore } from '@tool-bridge/core' +import { helpJsonSchema, registryNodeSchema } from '@tool-bridge/core/protocol' +import { describe, expect, it } from 'vitest' +import { bearer, createTestApp, type TestApp } from './harness' +import { processDeviceHello } from '../src/deviceHello' +import { RemotePathProjector } from '../src/federation' +import { TEST_ADMIN_SK } from './fixtures' + +const environment: DeviceEnvironment = { + platform: 'linux', arch: 'arm64', runtime: 'node', runtimeVersion: '22.12.0', runtimeId: 'runtime-1', +} +const now = '2026-09-15T00:00:00.000Z' + +async function hello(app: TestApp, report: DeviceEnvironment | undefined = environment): Promise { + const decoded = decodeDeviceFrame(encodeDeviceFrame({ + type: 'hello', deviceId: 'context-test', expose: { + ...(report === undefined ? {} : { environment: report }), + nodes: ['public', 'private'].map(path => ({ path, kind: 'tool', description: path, cmds: [{ name: 'read', effect: 'read' }] })), + }, + })) + if (decoded.type !== 'hello') throw new Error('expected hello') + await processDeviceHello({ store: app.state, hello: decoded, deviceIdHint: decoded.deviceId, authorization: `Bearer ${TEST_ADMIN_SK}`, now }) +} + +async function post(app: TestApp, path: string, body: unknown): Promise { + return await app.request(`https://tb.test/${path}`, bearer(TEST_ADMIN_SK, { + method: 'POST', headers: { 'content-type': 'application/json', 'accept': 'application/json' }, body: JSON.stringify(body), + })) +} + +async function help(app: TestApp, sk = TEST_ADMIN_SK): Promise { + const response = await app.request('https://tb.test/device/context-test/~help', bearer(sk, { headers: { accept: 'application/json' } })) + expect(response.status).toBe(200) + return helpJsonSchema.parse(await response.json()) +} + +describe('device context hello → storage → authorized help representations', () => { + it('keeps the hello snapshot across heartbeat/offline and filters children with current scopes', async () => { + const app = await createTestApp() + await hello(app) + const registry = new NodeRegistryStore(app.state) + const node = registryNodeSchema.parse(await registry.get('device/context-test')) + expect(node.deviceEnvironment).toEqual(environment) + expect(node.deviceReportedAt).toBe(now) + expect((await help(app)).deviceContext?.presence?.state).toBe('stale') + const heartbeatAt = '2026-09-15T00:01:00.000Z' + await registry.touchSeen('device/context-test', heartbeatAt) + expect((await registry.get('device/context-test')).deviceReportedAt).toBe(now) + await registry.setOnline('device/context-test', false, heartbeatAt) + const issued = await post(app, 'system/sk/write', { + owner: 'agent:context-reader', scopes: [ + { pattern: 'device/context-test/**', actions: ['read', 'call'] }, + { pattern: 'device/context-test/private/**', actions: ['read', 'call'], effect: 'deny' }, + ], + }) + expect(issued.status).toBe(200) + const { secret } = await issued.json() as { secret: string } + const result = await help(app, secret) + expect(result.deviceContext).toEqual({ environment, reportedAt: now, presence: { state: 'offline', lastSeenAt: heartbeatAt } }) + expect(result.children?.map(child => child.path)).toEqual(['device/context-test/public']) + expect(JSON.stringify(result)).not.toContain('private') + expect(result.hint).toContain('last hello report (offline)') + for (const accept of ['text/plain', 'text/markdown']) { + const response = await app.request('https://tb.test/device/context-test/~help', bearer(secret, { headers: { accept } })) + const text = await response.text() + expect(text).toContain('runtime-1') + expect(text).toContain(now) + expect(text).not.toContain('private') + } + const projected = new RemotePathProjector('remote').projectHelp(result, 'remote/device/context-test') + expect(projected.deviceContext).toEqual(result.deviceContext) + expect(projected.children?.[0]?.path).toBe('remote/device/context-test/public') + }) + + it('legacy hello removes an old snapshot instead of inventing current metadata', async () => { + const app = await createTestApp() + await hello(app) + await processDeviceHello({ store: app.state, authorization: `Bearer ${TEST_ADMIN_SK}`, deviceIdHint: 'context-test', hello: { deviceId: 'context-test', expose: {} }, now }) + expect((await help(app)).deviceContext).toBeUndefined() + }) + + it('rejects forged gateway fields from both registry authoring surfaces and custom hello nodes', async () => { + const app = await createTestApp() + for (const path of ['system/registry/write', 'forged/~register']) { + const response = await post(app, path, { + path: 'forged', kind: 'directory', description: 'forged', deviceEnvironment: environment, deviceReportedAt: now, + }) + expect(response.status).toBe(400) + } + await expect(processDeviceHello({ + store: app.state, authorization: `Bearer ${TEST_ADMIN_SK}`, deviceIdHint: 'forged', + hello: { deviceId: 'forged', expose: { nodes: [{ path: 'node', kind: 'tool', description: 'forged', ...{ deviceEnvironment: environment } }] } }, + })).rejects.toMatchObject({ code: 'invalid_argument' }) + expect(await app.state.get('node:device/forged')).toBeNull() + }) + + it.each([ + { platform: 'linux', cwd: '/home/private' }, + { platform: 'linux', runtimeVersion: '22\nsecret' }, + { platform: 'linux', arch: 'a'.repeat(33) }, + { platform: 'unknown' }, + ])('rejects unsupported environment fields before writing nodes: %j', async (invalid) => { + const app = await createTestApp() + expect(() => decodeDeviceFrame(JSON.stringify({ type: 'hello', deviceId: 'invalid', expose: { environment: invalid } }))).toThrow() + await expect(processDeviceHello({ store: app.state, authorization: `Bearer ${TEST_ADMIN_SK}`, deviceIdHint: 'invalid', hello: { deviceId: 'invalid', expose: { environment: invalid as DeviceEnvironment } } })).rejects.toMatchObject({ code: 'invalid_argument' }) + expect(await app.state.get('node:device/invalid')).toBeNull() + }) +}) diff --git a/packages/cli/examples/process-session-profile.json b/packages/cli/examples/process-session-profile.json new file mode 100644 index 00000000..2fe9e046 --- /dev/null +++ b/packages/cli/examples/process-session-profile.json @@ -0,0 +1,21 @@ +{ + "version": 1, + "path": "ops/demo", + "description": "A one-minute process session example (requires node on the device PATH)", + "commands": [ + { + "name": "wait", + "description": "Print a start message and complete after one minute", + "executable": "node", + "argv": [ + "-e", + "console.log('started'); setTimeout(() => console.log('completed'), 60000)" + ], + "effect": "read", + "session": { + "path": "sessions/demo", + "maxRuntimeMs": 120000 + } + } + ] +} diff --git a/packages/cli/package.json b/packages/cli/package.json index 12f03799..467179f1 100644 --- a/packages/cli/package.json +++ b/packages/cli/package.json @@ -1,6 +1,6 @@ { "name": "@tool-bridge/cli", - "version": "0.32.0", + "version": "0.33.0", "description": "tb \u2014 Tool Bridge CLI: one gateway for HTBP/MCP/HTTP tools, contexts, and device shells", "type": "module", "license": "MIT", diff --git a/packages/cli/src/commands/connect.ts b/packages/cli/src/commands/connect.ts index 9adbaeec..4215d70d 100644 --- a/packages/cli/src/commands/connect.ts +++ b/packages/cli/src/commands/connect.ts @@ -7,6 +7,7 @@ import { import { Command, type OptionValues } from 'commander' import { readFileSync } from 'node:fs' import { resolve } from 'node:path' +import { assertDeviceSessionPaths, deviceSessionBindings, sessionExposeNodes } from '../deviceSessions' import { collect, resolveTarget, withGlobalOpts } from '../args' import { asArray, printJson, printLine } from '../output' import { runDeviceConnection } from '../deviceRuntime' @@ -24,6 +25,7 @@ export interface ConnectArgs { path?: string /** `--no-shell` → false;缺省(undefined)= 暴露 shell。 */ shell?: boolean + shellSessionPath?: string sk?: string timeout?: string url?: string @@ -35,6 +37,7 @@ export interface PreparedConnect { deviceId: string expose: DeviceExpose mountPath?: string + shellSessionPath?: string sk: string } @@ -112,6 +115,9 @@ export function buildExpose( } }) } + const bindings = deviceSessionBindings(commandProfiles, args.shellSessionPath, expose.shell) + assertDeviceSessionPaths(commandProfiles, bindings) + if (bindings.length > 0) expose.nodes = [...(expose.nodes ?? []), ...sessionExposeNodes(bindings)] if (expose.shell === undefined && expose.fs === undefined && expose.nodes === undefined) { throw new CliError( 'nothing to expose: omit --no-shell, pass --fs, or pass --command-profile', @@ -148,6 +154,7 @@ export function prepareConnect(args: ConnectArgs): PreparedConnect { return { baseUrl: target.baseUrl, sk: target.sk, + ...(args.shellSessionPath === undefined ? {} : { shellSessionPath: args.shellSessionPath }), deviceId, expose, ...(commandProfiles.length > 0 ? { commandProfiles } : {}), @@ -198,6 +205,7 @@ export function connectCommand() { collect, [], ) + .option('--shell-session-path ', 'Opt in to Linux shell process sessions at this relative device path') .option('--no-shell', 'Do not expose shell; mutually exclusive with --allow') .addHelpText( 'after', diff --git a/packages/cli/src/commands/daemon.ts b/packages/cli/src/commands/daemon.ts index 03ffc195..48fa7696 100644 --- a/packages/cli/src/commands/daemon.ts +++ b/packages/cli/src/commands/daemon.ts @@ -85,6 +85,7 @@ function daemonInstallCommand() { collect, [], ) + .option('--shell-session-path ', 'Opt in to Linux shell process sessions at this relative device path') .option('--no-shell', 'Do not expose shell; mutually exclusive with --allow') .option('--yes', 'Confirm persistent arbitrary command execution when --allow "*" is used') .addHelpText( diff --git a/packages/cli/src/commands/device.ts b/packages/cli/src/commands/device.ts index 4cd261ad..f601c08d 100644 --- a/packages/cli/src/commands/device.ts +++ b/packages/cli/src/commands/device.ts @@ -8,6 +8,7 @@ import { collect, parsePageOpts, resolveTarget, withGlobalOpts, withPageOpts } f import { operationMeaning, printDeviceOperation } from '../deviceOutput' import { callDirect, CliError, withClient } from '../http' import { printJson, printLine, table } from '../output' +import { deviceSessionCommand } from './deviceSession' const OPERATION_STATES = new Set([ 'queued', @@ -155,4 +156,5 @@ export function deviceCommand() { .description('Manage reverse-connected devices') .addCommand(deviceLsCommand()) .addCommand(deviceOperationCommand()) + .addCommand(deviceSessionCommand()) } diff --git a/packages/cli/src/commands/deviceSession.ts b/packages/cli/src/commands/deviceSession.ts new file mode 100644 index 00000000..abf96b89 --- /dev/null +++ b/packages/cli/src/commands/deviceSession.ts @@ -0,0 +1,179 @@ +import { processSessionListSchema, processSessionObservationSchema, type ProcessSessionSummary, processSessionSummarySchema } from '@tool-bridge/sdk/client' +import { stripVTControlCharacters } from 'node:util' +import { randomUUID } from 'node:crypto' +import { Command } from 'commander' +import { collect, parsePageOpts, parsePositiveInt, resolveTarget, withGlobalOpts, withPageOpts } from '../args' +import { CliError, type Target, withClient } from '../http' +import { printJson, printLine, table } from '../output' +import { confirmDestructive } from '../confirm' +import { parseCallArgs } from './call' + +const STATES = ['running', 'stopping', 'stopped', 'exited', 'timed_out'] + +function nodePath(value: string): string { + const path = value.trim().replace(/^\/+|\/+$/g, '') + if (!path || path.split('/').some(part => !part || part === '.' || part === '..')) { + throw new CliError('a complete session node path is required') + } + return path +} + +function integer(value: string | undefined, flag: string, max: number): number | undefined { + if (value === undefined) return undefined + const n = Number(value) + if (!Number.isSafeInteger(n) || n < 0 || n > max) throw new CliError(`${flag} must be an integer between 0 and ${max}`) + return n +} + +function printSummary(session: ProcessSessionSummary): void { + printLine(`session: ${session.sessionId}`) + printLine(`runtime: ${session.runtimeId}`) + printLine(`request: ${session.requestId}`) + printLine(`state: ${session.state}${session.state === 'stopping' ? ' (termination requested; exit not yet confirmed)' : ''}`) + if (session.exitCode !== undefined) printLine(`exit code: ${session.exitCode}`) + if (session.signal !== undefined) printLine(`signal: ${session.signal}`) +} + +async function sessionCall( + target: Target, + path: string, + input: Record, + schema: { safeParse: (value: unknown) => { data: T, success: true } | { success: false } }, + signal?: AbortSignal, +): Promise { + const result = schema.safeParse(await withClient(target, client => client.invokeJson(path, input, { signal }))) + if (result.success) return result.data + const error = new CliError('invalid process session response; request outcome is unknown', 'internal') + error.kind = 'protocol' + error.outcome = 'unknown' + throw error +} + +async function confirmCommand(target: Target, path: string, name: string, yes?: boolean): Promise { + const help = await withClient(target, client => client.getHelp(path)) + const command = help.cmds.find(cmd => cmd.name === name) + if (!command) throw new CliError(`session command ${name} is not available`) + if (command.confirm) await confirmDestructive({ yes }, `Execute ${path}/${name}${command.effect ? ` (${command.effect})` : ''}?`) +} + +export function deviceSessionCommand() { + return new Command('session') + .description('Start, observe and stop realtime device process sessions (separate from daemon service logs)') + .addCommand(withGlobalOpts(new Command('start')) + .argument('', 'Full session node path') + .argument('[args]', 'Bound command input as a JSON object') + .option('--args ', 'Bound command input as a JSON object') + .option('--args-file ', 'Input from a JSON file, or - for stdin') + .option('--arg ', 'One flat input argument (repeatable)', collect, []) + .option('--run-timeout ', 'Shorten the configured process runtime limit') + .option('--yes', 'Skip interactive confirmation') + .description('Start once; on uncertain results, find the original request with session list --request-id') + .action(async (pathArg, positional, opts) => { + const path = nodePath(pathArg) + const input = await parseCallArgs(opts.args, opts.argsFile, positional, opts.arg) + const timeoutMs = parsePositiveInt(opts.runTimeout, '--run-timeout') + const target = resolveTarget(opts) + await confirmCommand(target, path, 'start', opts.yes) + const page = await sessionCall(target, `${path}/list`, { limit: 1 }, processSessionListSchema) + const time = Date.parse(page.now) + if (!page.runtimeId || !Number.isFinite(time)) throw new CliError('invalid session runtime response', 'internal') + const requestId = randomUUID() + try { + const session = await sessionCall(target, `${path}/start`, { + input, requestId, expectedRuntimeId: page.runtimeId, + startBefore: new Date(time + 60_000).toISOString(), + ...(timeoutMs === undefined ? {} : { timeoutMs }), + }, processSessionSummarySchema) + if (opts.json) printJson(session) + else printSummary(session) + } catch (error) { + if (error instanceof CliError) error.hint = `requestId: ${requestId}; find it with tb device session list ${path} --request-id ${requestId}. Do not blindly start again.` + throw error + } + })) + .addCommand(withPageOpts(withGlobalOpts(new Command('list'))) + .alias('ls') + .argument('', 'Full session node path') + .option('--request-id ', 'Find an original start request after a lost response') + .option('--state ', `Filter: ${STATES.join(', ')}`) + .description('List this caller’s sessions; an empty page does not prove a command never ran') + .action(async (pathArg, opts) => { + const path = nodePath(pathArg) + const pageOpts = parsePageOpts(opts) + if (opts.state !== undefined && !STATES.includes(opts.state)) throw new CliError('invalid --state') + if (opts.requestId !== undefined && !opts.requestId.trim()) throw new CliError('--request-id must not be empty') + const page = await sessionCall(resolveTarget(opts), `${path}/list`, { + ...pageOpts, + ...(opts.state === undefined ? {} : { state: opts.state }), + ...(opts.requestId === undefined ? {} : { requestId: opts.requestId }), + }, processSessionListSchema) + if (opts.json) printJson(page) + else { + printLine(`runtime: ${page.runtimeId} (reported at ${page.now})`) + printLine(page.items.length ? table(['SESSION', 'REQUEST', 'STATE', 'STARTED'], page.items.map(item => [item.sessionId, item.requestId, item.state, item.startedAt])) : '(no visible sessions on this page; this does not prove a command never ran)') + if (page.cursor) printLine(`next cursor: ${page.cursor}`) + } + })) + .addCommand(withGlobalOpts(new Command('observe')) + .argument('', 'Full session node path') + .argument('', 'Session identifier') + .option('--cursor ', 'Resume after this log sequence') + .option('--limit-bytes ', 'Maximum log bytes per response (4-262144)') + .option('--wait ', 'Wait for logs up to 20000 ms per request') + .option('--follow', 'Follow logs; Ctrl-C only stops observation. JSON output is one object per line') + .description('Read incremental process logs; daemon restart invalidates old sessions') + .action(async (pathArg, sessionId, opts) => { + const path = nodePath(pathArg) + let cursor = integer(opts.cursor, '--cursor', Number.MAX_SAFE_INTEGER) ?? 0 + const limitBytes = parsePositiveInt(opts.limitBytes, '--limit-bytes') + if (limitBytes !== undefined && (limitBytes < 4 || limitBytes > 262_144)) throw new CliError('--limit-bytes must be between 4 and 262144') + const waitMs = integer(opts.wait, '--wait', 20_000) ?? (opts.follow ? 20_000 : 0) + if (opts.follow && waitMs === 0) throw new CliError('--follow requires a positive --wait') + const target = resolveTarget(opts) + const controller = new AbortController() + const interrupt = () => controller.abort() + if (opts.follow) process.on('SIGINT', interrupt) + try { + while (!controller.signal.aborted) { + const page = await sessionCall(target, `${path}/observe`, { + sessionId, cursor, waitMs, ...(limitBytes === undefined ? {} : { limitBytes }), + }, processSessionObservationSchema, controller.signal) + if (opts.json) { + if (opts.follow) process.stdout.write(`${JSON.stringify(page)}\n`) + else printJson(page) + } else { + if (page.gap) process.stderr.write(`log gap: ${page.droppedBytes} bytes dropped; remaining logs follow\n`) + for (const chunk of page.chunks) { + const stream = chunk.stream === 'stderr' ? process.stderr : process.stdout + stream.write(stripVTControlCharacters(chunk.text)) + } + if (!opts.follow || (page.state !== 'running' && page.state !== 'stopping' && cursor === page.nextCursor)) printSummary(page) + } + const previousCursor = cursor + cursor = page.nextCursor + // Terminal responses can still contain a partial page of buffered logs. + if (!opts.follow || (page.state !== 'running' && page.state !== 'stopping' && cursor === previousCursor)) break + } + } catch (error) { + if (!controller.signal.aborted) { + if (error instanceof CliError) error.hint = `Resume observation with tb device session observe ${path} ${sessionId} --cursor ${cursor} --follow. A daemon restart invalidates the old session.` + throw error + } + } finally { + if (opts.follow) process.off('SIGINT', interrupt) + } + })) + .addCommand(withGlobalOpts(new Command('stop')) + .argument('', 'Full session node path') + .argument('', 'Session identifier') + .option('--yes', 'Skip interactive confirmation when declared by command help') + .description('Explicitly terminate a process session; reports the actual state') + .action(async (pathArg, sessionId, opts) => { + const path = nodePath(pathArg) + const target = resolveTarget(opts) + await confirmCommand(target, path, 'stop', opts.yes) + const session = await sessionCall(target, `${path}/stop`, { sessionId }, processSessionSummarySchema) + if (opts.json) printJson(session) + else printSummary(session) + })) +} diff --git a/packages/cli/src/daemon.ts b/packages/cli/src/daemon.ts index c51beef5..cbf09f14 100644 --- a/packages/cli/src/daemon.ts +++ b/packages/cli/src/daemon.ts @@ -25,6 +25,7 @@ export interface DaemonConfig { expose: DeviceExpose mountPath?: string revision: string + shellSessionPath?: string sk: string version: 1 } @@ -64,6 +65,7 @@ export interface DaemonInstallInput { deviceId: string expose: DeviceExpose mountPath?: string + shellSessionPath?: string sk: string } @@ -116,6 +118,7 @@ function assertDaemonConfig(value: unknown): asserts value is DaemonConfig { || typeof value.sk !== 'string' || typeof value.deviceId !== 'string' || typeof value.revision !== 'string' + || (value.shellSessionPath !== undefined && typeof value.shellSessionPath !== 'string') || !isRecord(value.expose) || (value.commandProfiles !== undefined && !Array.isArray(value.commandProfiles)) || (value.mountPath !== undefined && typeof value.mountPath !== 'string')) { @@ -183,6 +186,9 @@ ExecStart=${command} Restart=always RestartSec=5s KillSignal=SIGTERM +KillMode=control-group +TimeoutStopSec=15s +SendSIGKILL=yes [Install] WantedBy=default.target @@ -379,6 +385,7 @@ export async function installDaemon( sk: input.sk, deviceId: input.deviceId, expose: input.expose, + ...(input.shellSessionPath === undefined ? {} : { shellSessionPath: input.shellSessionPath }), ...(input.commandProfiles !== undefined ? { commandProfiles: input.commandProfiles } : {}), ...(input.mountPath !== undefined ? { mountPath: input.mountPath } : {}), } @@ -503,6 +510,7 @@ export async function runDaemon(configPath: string): Promise { sk: config.sk, deviceId: config.deviceId, expose: config.expose, + ...(config.shellSessionPath === undefined ? {} : { shellSessionPath: config.shellSessionPath }), ...(config.commandProfiles !== undefined ? { commandProfiles: config.commandProfiles } : {}), ...(config.mountPath !== undefined ? { mountPath: config.mountPath } : {}), onMailboxError: error => process.stderr.write(`device mailbox: ${error.message}\n`), diff --git a/packages/cli/src/deviceRuntime.ts b/packages/cli/src/deviceRuntime.ts index b01d5da0..bcb62151 100644 --- a/packages/cli/src/deviceRuntime.ts +++ b/packages/cli/src/deviceRuntime.ts @@ -1,5 +1,6 @@ import { createDeviceMailboxProcessor, + type DeviceCallContext, type DeviceConnection, type DeviceConnectionState, type DeviceCredentialProvider, @@ -20,12 +21,14 @@ import { type SearchOptions, } from '@tool-bridge/core' import { + createProcessSessionManager, createShellExecutor, createStructuredCommandRuntime, FsObjectStore, type StructuredCommandProfile, } from '@tool-bridge/core/node' import WS, { type ClientOptions } from 'ws' +import { assertDeviceSessionPaths, deviceSessionBindings, sessionExposeNodes } from './deviceSessions' import { createFileDeviceOperationJournal } from './deviceMailboxJournal' import { CliError } from './http' @@ -40,6 +43,7 @@ export interface DeviceConnectionOptions { onMailboxError?: (error: Error) => void onReady?: (mountPath: string) => void onStateChange?: (state: DeviceConnectionState) => void + shellSessionPath?: string sk: string } @@ -143,11 +147,39 @@ export function startDeviceConnection(opts: DeviceConnectionOptions): DeviceConn } structured.set(runtime.path, runtime) } + const bindings = deviceSessionBindings(opts.commandProfiles ?? [], opts.shellSessionPath, opts.expose.shell) + const sessionPaths = new Set(bindings.map(binding => binding.path)) + const customPaths = (opts.expose.nodes ?? []).map(node => node.path) + .filter(path => !sessionPaths.has(path) && !structured.has(path)) + assertDeviceSessionPaths(opts.commandProfiles ?? [], bindings, customPaths) + const sessions = createProcessSessionManager() + for (const binding of bindings) sessions.register(binding) + const expose: DeviceExpose = { + ...opts.expose, + environment: { + platform: process.platform === 'linux' || process.platform === 'darwin' || process.platform === 'win32' + ? process.platform + : 'other', + arch: process.arch, + runtime: typeof process.versions.bun === 'string' ? 'bun' : 'node', + runtimeVersion: process.versions.bun ?? process.versions.node, + ...(bindings.length === 0 ? {} : { runtimeId: sessions.runtimeId }), + }, + ...(bindings.length === 0 + ? {} + : { nodes: [ + ...(opts.expose.nodes ?? []).filter(node => !sessionPaths.has(node.path)), + ...sessionExposeNodes(bindings), + ] }), + } + let sessionCleanup: Promise | undefined + const cleanupSessions = (): Promise => sessionCleanup ??= sessions.close() let activeMountPath = opts.mountPath ?? `device/${opts.deviceId}` let files = store === undefined ? undefined : fsProvider(store, activeMountPath, readOnly) const handler = async (call: { arguments: Record + context?: DeviceCallContext path: string signal: DeviceAbortSignal }): Promise => { @@ -155,6 +187,12 @@ export function startDeviceConnection(opts: DeviceConnectionOptions): DeviceConn const mount = slash < 0 ? call.path : call.path.slice(0, slash) const cmd = slash < 0 ? '' : call.path.slice(slash + 1) try { + if (sessionPaths.has(mount)) { + return await sessions.invoke(mount, cmd, call.arguments, { + ...(call.context === undefined ? {} : { caller: call.context.caller }), + signal: call.signal, + }) + } if (mount === 'shell') { if (cmd !== 'exec') throw new TBError('invalid_argument', `unknown shell cmd '${cmd}'`) if (shell === undefined) throw TBError.notFound('shell not exposed') @@ -196,7 +234,10 @@ export function startDeviceConnection(opts: DeviceConnectionOptions): DeviceConn const fail = (error: unknown): void => { if (settled || userClosed) return settled = true - rejectClosed(cliError(error)) + void cleanupSessions().then( + () => rejectClosed(cliError(error)), + cleanupError => rejectClosed(cliError(cleanupError)), + ) } const connectionControl: { close?: () => void } = {} const credentialProvider: DeviceCredentialProvider = { @@ -255,7 +296,7 @@ export function startDeviceConnection(opts: DeviceConnectionOptions): DeviceConn const connection = openPortableDeviceConnection({ baseUrl: opts.baseUrl, deviceId: opts.deviceId, - expose: async () => opts.expose, + expose: async () => expose, mountPath: opts.mountPath, webSocketFactory: nodeWebSocketFactory, credentialProvider, @@ -282,7 +323,7 @@ export function startDeviceConnection(opts: DeviceConnectionOptions): DeviceConn connection.closed.then(() => { if (settled) return settled = true - resolveClosed() + void cleanupSessions().then(resolveClosed, error => rejectClosed(cliError(error))) }, fail) return { @@ -294,6 +335,7 @@ export function startDeviceConnection(opts: DeviceConnectionOptions): DeviceConn close() { userClosed = true mailboxDrainController?.abort() + void cleanupSessions().catch(fail) connection.close() }, restart() { diff --git a/packages/cli/src/deviceSessions.ts b/packages/cli/src/deviceSessions.ts new file mode 100644 index 00000000..0c7fb9dd --- /dev/null +++ b/packages/cli/src/deviceSessions.ts @@ -0,0 +1,50 @@ +import { + createShellSessionBinding, + createStructuredCommandRuntime, + type ProcessSessionBinding, + processSessionCommands, + type StructuredCommandProfile, +} from '@tool-bridge/core/node' +import { type DeviceExpose, normalizePath, validatePath } from '@tool-bridge/core' +import { CliError } from './http' + +export function deviceSessionBindings( + profiles: readonly StructuredCommandProfile[], + shellSessionPath?: string, + shell?: DeviceExpose['shell'], +): ProcessSessionBinding[] { + const bindings = profiles.flatMap(profile => createStructuredCommandRuntime(profile).sessionBindings) + if (shellSessionPath !== undefined) { + if (shell === undefined) throw new CliError('--shell-session-path requires shell exposure') + const path = normalizePath(shellSessionPath) + if (path === '') throw new CliError('shell session path must not be root') + bindings.push(createShellSessionBinding(path, shell.allow)) + } + return bindings +} + +/** Reused by installation and live runtime so frozen profiles cannot bypass collisions. */ +export function assertDeviceSessionPaths( + profiles: readonly StructuredCommandProfile[], + bindings: readonly ProcessSessionBinding[], + extraPaths: readonly string[] = [], +): void { + const paths = ['shell', 'fs', ...profiles.map(profile => profile.path), ...bindings.map(binding => binding.path), ...extraPaths] + for (const [index, path] of paths.entries()) { + if (validatePath(path) !== null || normalizePath(path) !== path || path === '') throw new CliError(`invalid device path '${path}'`) + for (const other of paths.slice(0, index)) { + if (path === other || path.startsWith(`${other}/`) || other.startsWith(`${path}/`)) { + throw new CliError(`device path '${path}' conflicts with '${other}'`) + } + } + } +} + +export function sessionExposeNodes(bindings: readonly ProcessSessionBinding[]): NonNullable { + return bindings.map(binding => ({ + path: binding.path, + kind: 'tool', + description: binding.description ?? 'Device process sessions', + cmds: processSessionCommands(binding), + })) +} diff --git a/packages/cli/test/deviceRuntime.test.ts b/packages/cli/test/deviceRuntime.test.ts index a7ebf11d..06c67a68 100644 --- a/packages/cli/test/deviceRuntime.test.ts +++ b/packages/cli/test/deviceRuntime.test.ts @@ -222,6 +222,9 @@ describe('SDK device supervisor 的 Node adapter', () => { expose: { shell: { allow: ['echo'], description: 'CI shell' }, fs: { roots: [root], readOnly: true }, + environment: { platform: process.platform, arch: process.arch, + runtime: typeof process.versions.bun === 'string' ? 'bun' : 'node', + runtimeVersion: process.versions.bun ?? process.versions.node }, }, }]) socket?.dispatch('message', { diff --git a/packages/cli/test/deviceSession.test.ts b/packages/cli/test/deviceSession.test.ts new file mode 100644 index 00000000..3f55fd8d --- /dev/null +++ b/packages/cli/test/deviceSession.test.ts @@ -0,0 +1,132 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { resetFetch, setFetch } from '../src/http' +import { runCli } from './cliHarness' + +const path = 'device/build/sessions/build' +const gateway = ['--base-url', 'https://gw', '--sk', 'tbk_admin'] +const summary = { + sessionId: 'runtime:session', runtimeId: 'runtime', requestId: 'original-request', + state: 'running', startedAt: '2026-09-15T00:00:00.000Z', +} +const list = { runtimeId: 'runtime', now: summary.startedAt, items: [summary] } +const help = { htbp: '0.1', node: { path, kind: 'tool', description: '' }, cmds: ['start', 'list', 'observe', 'stop'].map(name => ({ name, path: `${path}/${name}`, scope: 'call', effect: 'write', method: 'POST' })) } +const output = () => vi.mocked(process.stdout.write).mock.calls.map(call => String(call[0])).join('') + +function transport(handler: (url: string, body: Record, signal?: AbortSignal | null) => unknown | Promise) { + const fetcher = vi.fn(async (url: string | URL | Request, init?: RequestInit) => new Response(JSON.stringify(await handler(String(url), JSON.parse(String(init?.body ?? '{}')), init?.signal)), { headers: { 'content-type': 'application/json' } })) + setFetch(fetcher as typeof fetch) + return fetcher +} + +beforeEach(() => { + process.exitCode = 0 + vi.spyOn(process.stdout, 'write').mockReturnValue(true) + vi.spyOn(process.stderr, 'write').mockReturnValue(true) +}) +afterEach(() => { + process.exitCode = 0 + resetFetch() + vi.restoreAllMocks() +}) + +describe('device process session CLI', () => { + it('start uses device time and generation and sends the bound input once', async () => { + const fetcher = transport((url, body) => { + if (url.endsWith('/~help')) return help + if (url.endsWith('/list')) return list + return { ...summary, requestId: body.requestId } + }) + await runCli(['device', 'session', 'start', path, '--arg', 'message=hello', '--run-timeout', '90000', '--json', ...gateway]) + expect(process.exitCode).toBe(0) + const [, init] = fetcher.mock.calls[2]! + expect(JSON.parse(String(init?.body))).toEqual({ input: { message: 'hello' }, requestId: expect.any(String), expectedRuntimeId: 'runtime', startBefore: '2026-09-15T00:01:00.000Z', timeoutMs: 90_000 }) + expect(JSON.parse(output()).sessionId).toBe(summary.sessionId) + expect(fetcher).toHaveBeenCalledTimes(3) + }) + + it('lost start response retains its original request id and never retries', async () => { + let requestId = '' + const fetcher = transport((url, body) => { + if (url.endsWith('/~help')) return help + if (url.endsWith('/list')) return list + requestId = String(body.requestId) + throw new TypeError('network disconnected') + }) + await runCli(['device', 'session', 'start', path, '--json', ...gateway]) + expect(fetcher).toHaveBeenCalledTimes(3) + expect(JSON.parse(output())).toMatchObject({ ok: false, outcome: 'unknown', hint: expect.stringContaining(`--request-id ${requestId}`) }) + expect(process.exitCode).toBe(1) + }) + + it('follow advances independent cursors and drains terminal buffered output as JSON lines', async () => { + const cursors: unknown[] = [] + transport((_url, body) => { + cursors.push(body.cursor) + return { ...summary, state: 'exited', exitCode: 3, chunks: body.cursor === 0 ? [{ seq: 1, stream: 'stderr', text: 'failure' }] : [], nextCursor: 1, gap: body.cursor === 0, droppedBytes: 12 } + }) + await runCli(['device', 'session', 'observe', path, summary.sessionId, '--follow', '--json', ...gateway]) + expect(cursors).toEqual([0, 1]) + expect(output().trim().split('\n').map(line => JSON.parse(line))).toEqual([ + expect.objectContaining({ state: 'exited', exitCode: 3, gap: true }), + expect.objectContaining({ nextCursor: 1, chunks: [] }), + ]) + }) + + it('Ctrl-C aborts only the current observer and removes its signal listener', async () => { + const before = process.listenerCount('SIGINT') + const fetcher = transport(async (_url, _body, signal) => await new Promise((_resolve, reject) => { + signal?.addEventListener('abort', () => reject(new DOMException('aborted', 'AbortError')), { once: true }) + queueMicrotask(() => process.emit('SIGINT')) + })) + await runCli(['device', 'session', 'observe', path, summary.sessionId, '--follow', ...gateway]) + expect(fetcher).toHaveBeenCalledOnce() + expect(String(fetcher.mock.calls[0]?.[0])).toMatch(/\/observe$/) + expect(process.listenerCount('SIGINT')).toBe(before) + expect(process.exitCode).toBe(0) + }) + + it('renders VT-free logs and accurately reports a pending stop', async () => { + transport(url => url.endsWith('/~help') ? help : { ...summary, state: 'stopping' }) + await runCli(['device', 'session', 'stop', path, summary.sessionId, ...gateway]) + expect(output()).toContain('exit not yet confirmed') + vi.mocked(process.stdout.write).mockClear() + transport(() => ({ ...summary, chunks: [{ seq: 1, stream: 'stdout', text: '\u001b[31mhello\u001b[0m' }], nextCursor: 1, gap: false, droppedBytes: 0 })) + await runCli(['device', 'session', 'observe', path, summary.sessionId, ...gateway]) + expect(output()).toContain('hello') + expect(output()).not.toContain('\u001b') + }) + + it.each([ + ['list', path, '--state', 'unknown'], + ['observe', path, summary.sessionId, '--cursor', '-1'], + ['observe', path, summary.sessionId, '--wait', '20001'], + ['observe', path, summary.sessionId, '--follow', '--wait', '0'], + ['observe', path, summary.sessionId, '--limit-bytes', '262145'], + ['observe', path, summary.sessionId, '--limit-bytes', '1'], + ['observe', path, summary.sessionId, '--limit-bytes', '2'], + ['observe', path, summary.sessionId, '--limit-bytes', '3'], + ['start', path, '--args', '[]'], + ['start', path, '--run-timeout', '0'], + ['stop', path, summary.sessionId, 'extra'], + ['list', path, '--unknown'], + ])('rejects invalid arguments before requests: %j', async (...args) => { + const fetcher = transport(() => list) + await runCli(['device', 'session', ...args, ...gateway]) + expect(process.exitCode).toBe(1) + expect(fetcher).not.toHaveBeenCalled() + }) + + it('forwards filters and does not describe an empty cursor page as no sessions globally', async () => { + const fetcher = transport(() => ({ ...list, items: [], cursor: 'next-page' })) + await runCli(['device', 'session', 'list', path, '--request-id', 'recover', '--state', 'running', '--limit', '2', ...gateway]) + expect(JSON.parse(String(fetcher.mock.calls[0]?.[1]?.body))).toEqual({ requestId: 'recover', state: 'running', limit: 2 }) + expect(output()).toContain('next cursor: next-page') + expect(output()).toContain('does not prove a command never ran') + }) + + it('rejects malformed successful responses instead of inventing session state', async () => { + transport(() => ({ sessionId: 'incomplete' })) + await runCli(['device', 'session', 'observe', path, summary.sessionId, '--json', ...gateway]) + expect(JSON.parse(output())).toMatchObject({ ok: false, kind: 'protocol', outcome: 'unknown' }) + }) +}) diff --git a/packages/cli/test/deviceSessions.test.ts b/packages/cli/test/deviceSessions.test.ts new file mode 100644 index 00000000..9c38a97c --- /dev/null +++ b/packages/cli/test/deviceSessions.test.ts @@ -0,0 +1,127 @@ +import type { OpenPortableDeviceConnectionOptions } from '@tool-bridge/sdk/device' +import { processSessionListSchema, processSessionObservationSchema, processSessionSummarySchema } from '@tool-bridge/sdk/client' +import { parseStructuredCommandProfile } from '@tool-bridge/core/node' +import { afterEach, describe, expect, it, vi } from 'vitest' +import { startDeviceConnection } from '../src/deviceRuntime' +import { buildExpose } from '../src/commands/connect' +import { renderSystemdUnit } from '../src/daemon' + +const transport = vi.hoisted(() => ({ + options: undefined as OpenPortableDeviceConnectionOptions | undefined, + finish: () => {}, + close: vi.fn(), +})) +vi.mock('@tool-bridge/sdk/device', async original => ({ + ...await original(), + openPortableDeviceConnection: (options: OpenPortableDeviceConnectionOptions) => { + transport.options = options + return { + ready: Promise.resolve('device/test'), + closed: new Promise((resolve) => { transport.finish = resolve }), + close: () => { + transport.close() + transport.finish() + }, + restart: vi.fn(), resume: vi.fn(), suspend: vi.fn(), state: 'ready', + } + }, +})) +// Linux is the enabled production platform; these tests also exercise POSIX groups on macOS. +vi.mock('@tool-bridge/core/node', async (original) => { + const actual = await original() + return { ...actual, createProcessSessionManager: () => actual.createProcessSessionManager({ platform: 'linux', terminationGraceMs: 30 }) } +}) + +const handles: ReturnType[] = [] +afterEach(async () => { + for (const handle of handles.splice(0)) { + handle.close() + await handle.closed.catch(() => {}) + } + transport.options = undefined + vi.clearAllMocks() +}) +const profile = () => parseStructuredCommandProfile({ + version: 1, path: 'ops', description: 'test command', + commands: [{ name: 'build', description: 'test build', executable: process.execPath, + argv: ['-e', 'console.log(\'ready\');setInterval(()=>console.log(\'tick\'),20)'], effect: 'read', + session: { path: 'sessions/build', maxRuntimeMs: 60_000 } }], +}) +function connect() { + const commandProfiles = [profile()] + const handle = startDeviceConnection({ baseUrl: 'https://gateway.example', sk: 'test-key', deviceId: 'test', + commandProfiles, expose: buildExpose({ shell: false }, commandProfiles) }) + handles.push(handle) + return handle +} +async function call(action: string, args: Record, owner = 'agent:alice') { + if (!transport.options) throw new Error('transport not initialized') + return transport.options.handler({ id: crypto.randomUUID(), path: `sessions/build/${action}`, arguments: args, + context: { caller: { owner, keyId: 'key-1' }, createdAt: new Date().toISOString(), expiresAt: new Date(Date.now() + 60_000).toISOString(), traceId: 'test' }, + signal: new AbortController().signal, + uploadObject: async () => { throw new Error('not used') } }) +} + +describe('device process session host integration', () => { + it('opts in through profile, preserves sync command and blocks path collisions', () => { + const expose = buildExpose({ shell: false }, [profile()]) + expect(expose.nodes?.map(node => node.path)).toEqual(['ops', 'sessions/build']) + expect(expose.nodes?.[0]?.cmds?.map(cmd => cmd.name)).toEqual(['build']) + expect(expose.nodes?.[1]?.cmds?.map(cmd => cmd.name)).toEqual(['start', 'list', 'observe', 'stop']) + expect(() => buildExpose({ shell: false, shellSessionPath: 'sessions/shell' })).toThrow('requires shell') + expect(() => buildExpose({ shellSessionPath: 'ops' }, [profile()])).toThrow('conflicts') + expect(() => buildExpose({ shellSessionPath: 'shell/sessions' })).toThrow('conflicts') + expect(() => parseStructuredCommandProfile({ ...profile(), commands: [{ ...profile().commands[0], session: { path: 'ops/nested' } }] })).toThrow('conflicts') + }) + + it('retains one runtime across reconnect and waits for child cleanup on close', async () => { + const handle = connect() + await handle.ready + const listed = processSessionListSchema.parse(await call('list', {})) + const started = processSessionSummarySchema.parse(await call('start', { + input: {}, requestId: 'start-once', expectedRuntimeId: listed.runtimeId, + startBefore: new Date(Date.parse(listed.now) + 60_000).toISOString(), + })) + expect(started.state).toBe('running') + transport.options?.onStateChange?.('reconnecting') + transport.options?.onStateChange?.('ready') + const afterReconnect = processSessionListSchema.parse(await call('list', {})) + expect(afterReconnect.runtimeId).toBe(listed.runtimeId) + expect(afterReconnect.items[0]?.sessionId).toBe(started.sessionId) + await vi.waitFor(async () => { + const output = processSessionObservationSchema.parse(await call('observe', { sessionId: started.sessionId })) + expect(output.chunks.map(chunk => chunk.text).join('')).toContain('ready') + }) + const first = processSessionObservationSchema.parse(await call('observe', { sessionId: started.sessionId })) + const replay = processSessionObservationSchema.parse(await call('observe', { sessionId: started.sessionId })) + expect(replay.chunks[0]).toEqual(first.chunks[0]) + await expect(call('observe', { sessionId: started.sessionId }, 'agent:bob')).rejects.toMatchObject({ code: 'not_found' }) + const wireExpose = await transport.options!.expose() + expect(wireExpose.environment).toMatchObject({ runtimeId: listed.runtimeId, arch: process.arch }) + expect(wireExpose.environment).not.toHaveProperty('cwd') + handle.close() + await handle.closed + await expect(call('list', {})).rejects.toMatchObject({ code: 'unavailable' }) + }) + + it('invalidates old IDs after a new runtime and fails closed without caller identity', async () => { + const first = connect() + const old = processSessionListSchema.parse(await call('list', {})) + first.close() + await first.closed + connect() + const next = processSessionListSchema.parse(await call('list', {})) + expect(next.runtimeId).not.toBe(old.runtimeId) + await expect(call('start', { input: {}, requestId: 'old', expectedRuntimeId: old.runtimeId, + startBefore: new Date(Date.now() + 60_000).toISOString() })).rejects.toMatchObject({ code: 'conflict' }) + await expect(transport.options!.handler({ id: 'no-owner', path: 'sessions/build/list', arguments: {}, signal: new AbortController().signal, + uploadObject: async () => { throw new Error('unused') } })).rejects.toMatchObject({ code: 'unavailable' }) + }) + + it('pins systemd group cleanup and a bounded stop deadline', () => { + const unit = renderSystemdUnit(['/usr/bin/node', '/opt/tb/index.js'], '/home/test/device.json') + expect(unit).toContain('KillMode=control-group') + expect(unit).toContain('TimeoutStopSec=15s') + expect(unit).toContain('SendSIGKILL=yes') + }) +}) diff --git a/packages/cli/test/strictParsing.test.ts b/packages/cli/test/strictParsing.test.ts index 34c13456..d7cbbedd 100644 --- a/packages/cli/test/strictParsing.test.ts +++ b/packages/cli/test/strictParsing.test.ts @@ -143,6 +143,10 @@ describe('未知 flag 必须报错(事故回归)', () => { ['daemon', 'uninstall', '--bogus'], ['daemon', '_run', '--config', '/tmp/device.json', '--bogus'], ['device', 'ls', '--bogus'], + ['device', 'session', 'start', 'p', '--bogus'], + ['device', 'session', 'list', 'p', '--bogus'], + ['device', 'session', 'observe', 'p', 's', '--bogus'], + ['device', 'session', 'stop', 'p', 's', '--bogus'], ['device', 'op', 'ls', 'd', '--bogus'], ['device', 'op', 'get', 'd', 'dop_x', '--bogus'], ['device', 'op', 'cancel', 'd', 'dop_x', '--bogus'], diff --git a/packages/core/src/device/environment.ts b/packages/core/src/device/environment.ts new file mode 100644 index 00000000..f1527038 --- /dev/null +++ b/packages/core/src/device/environment.ts @@ -0,0 +1,11 @@ +import { z } from 'zod' +import type { DeviceEnvironment } from '../types' + +/** A small explicit allowlist: never accept arbitrary process/env metadata. */ +export const deviceEnvironmentSchema: z.ZodType = z.strictObject({ + platform: z.enum(['darwin', 'linux', 'win32', 'android', 'ios', 'other']), + arch: z.string().min(1).max(32).regex(/^[a-zA-Z0-9._-]+$/).optional(), + runtime: z.enum(['node', 'bun', 'react-native', 'other']).optional(), + runtimeVersion: z.string().min(1).max(64).regex(/^[a-zA-Z0-9.+_-]+$/).optional(), + runtimeId: z.string().min(1).max(128).regex(/^[a-zA-Z0-9._-]+$/).optional(), +}) diff --git a/packages/core/src/device/frames.ts b/packages/core/src/device/frames.ts index 6e77c062..c73dbf2c 100644 --- a/packages/core/src/device/frames.ts +++ b/packages/core/src/device/frames.ts @@ -9,6 +9,7 @@ import { z } from 'zod' import { type DeviceExpose, NODE_KINDS, type OwnerRef, type Timestamp, type TreePath } from '../types' import { tbErrorBodySchema } from '../protocol/errorWire' +import { deviceEnvironmentSchema } from './environment' import { TBError, type TBErrorBody } from '../errors' // ---------- 帧类型(TS 定义为真源) ---------- @@ -135,6 +136,7 @@ const nodeInputSchema = z .passthrough() const deviceExposeSchema = z.object({ + environment: deviceEnvironmentSchema.optional(), shell: z .object({ description: z.string().optional(), diff --git a/packages/core/src/device/helpModel.ts b/packages/core/src/device/helpModel.ts index 639e7273..f3830efb 100644 --- a/packages/core/src/device/helpModel.ts +++ b/packages/core/src/device/helpModel.ts @@ -6,9 +6,9 @@ * directory(mountPath)节点:description 呈现三态 presence(online/stale/offline)。 */ +import type { DeviceEnvironment, Timestamp, TreePath } from '../types' import type { ChildRef, CmdSpec, HelpModel } from '../htbp/model' import type { Presence } from './presence' -import type { TreePath } from '../types' import { contextHelpModel, type ContextHelpOptions } from '../context/help' import { cmdPath, withCommandPaths } from '../builtin/util' import { describeAllow } from './shellAllow' @@ -57,7 +57,7 @@ export function deviceFsHelpModel( /** `` directory 节点的 ~help;description 附三态 presence(online/stale/offline)。 */ export function deviceDirectoryHelpModel( - node: { description: string, path: TreePath, presence: Presence }, + node: { description: string, environment?: DeviceEnvironment, path: TreePath, presence: Presence, reportedAt?: Timestamp }, children: ChildRef[] = [], ): HelpModel { return { @@ -68,5 +68,11 @@ export function deviceDirectoryHelpModel( }, cmds: [], children, + ...(node.environment === undefined || node.reportedAt === undefined + ? {} + : { + deviceContext: { environment: node.environment, reportedAt: node.reportedAt, presence: node.presence }, + hint: `Environment is the last hello report (${node.presence.state}); inspect authorized child ~help for commands, execution limits and session recovery rules.`, + }), } } diff --git a/packages/core/src/device/processSessionContract.ts b/packages/core/src/device/processSessionContract.ts new file mode 100644 index 00000000..073f99ce --- /dev/null +++ b/packages/core/src/device/processSessionContract.ts @@ -0,0 +1,49 @@ +import { z } from 'zod/v4' +/** Portable result contract for device-owned process sessions. */ +export type ProcessSessionState = 'running' | 'stopping' | 'stopped' | 'exited' | 'timed_out' + +export interface ProcessSessionSummary { + completedAt?: string + exitCode?: number + requestId: string + runtimeId: string + sessionId: string + signal?: string + startedAt: string + state: ProcessSessionState +} + +export interface ProcessSessionList { + cursor?: string + items: ProcessSessionSummary[] + now: string + runtimeId: string +} + +export interface ProcessSessionChunk { + seq: number + stream: 'stdout' | 'stderr' + text: string +} + +export interface ProcessSessionObservation extends ProcessSessionSummary { + chunks: ProcessSessionChunk[] + droppedBytes: number + gap: boolean + nextCursor: number +} + +/** Public schemas shared by device metadata, consumer parsers, and contract tests. */ +export const processSessionSummarySchema = z.strictObject({ + sessionId: z.string(), runtimeId: z.string(), requestId: z.string(), + state: z.enum(['running', 'stopping', 'stopped', 'exited', 'timed_out']), + startedAt: z.iso.datetime(), completedAt: z.iso.datetime().optional(), + exitCode: z.number().int().optional(), signal: z.string().optional(), +}) +export const processSessionListSchema = z.strictObject({ + runtimeId: z.string(), now: z.iso.datetime(), items: z.array(processSessionSummarySchema), cursor: z.string().optional(), +}) +export const processSessionObservationSchema = processSessionSummarySchema.extend({ + chunks: z.array(z.strictObject({ seq: z.number().int().nonnegative(), stream: z.enum(['stdout', 'stderr']), text: z.string() })), + nextCursor: z.number().int().nonnegative(), gap: z.boolean(), droppedBytes: z.number().int().nonnegative(), +}) diff --git a/packages/core/src/device/public.ts b/packages/core/src/device/public.ts index 2bc72d22..5686c754 100644 --- a/packages/core/src/device/public.ts +++ b/packages/core/src/device/public.ts @@ -11,12 +11,14 @@ export { } from '../errors' export { normalizePath, validatePath } from '../tree/path' export { + type DeviceEnvironment, type DeviceExpose, type DeviceNodeCmd, type DeviceNodeInput, type TreePath, } from '../types' export * from './client' +export * from './environment' export * from './frames' // mailbox 是 server 侧权威(DeviceMailboxService 等),不属于设备客户端窄入口; // 只显式导出设备执行侧真实消费的完成载荷类型。 @@ -29,3 +31,12 @@ export { PRESENCE_STALE_AFTER_MS, type PresenceState, } from './presence' +export type { + ProcessSessionChunk, + ProcessSessionList, + ProcessSessionObservation, + ProcessSessionState, + ProcessSessionSummary, +} from './processSessionContract' + +export { processSessionListSchema, processSessionObservationSchema, processSessionSummarySchema } from './processSessionContract' diff --git a/packages/core/src/device/shellAllow.ts b/packages/core/src/device/shellAllow.ts index 633447a7..f1c296a1 100644 --- a/packages/core/src/device/shellAllow.ts +++ b/packages/core/src/device/shellAllow.ts @@ -3,12 +3,12 @@ * * 规则:allow 缺省/空数组 = 拒绝一切(默认拒);单值 ['*'] = 放行全部;其余对 command * 做 shell-word 切分取 argv[0] 的 basename,与条目精确匹配;白名单非 ['*'] 时 command - * 含 shell 元字符(; | & $( 反引号 > <)→ 直接拒——不封元字符则 `echo hi; rm -rf` + * 含 shell 元字符(; | & $( 反引号 > < 换行)→ 直接拒——不封元字符则 `echo hi; rm -rf` * 可绕过任何 argv[0] 判定。判定在设备侧执行前完成(shellExecutor 调用本函数)。 */ /** 单字符元字符;`$(` 是双字符序列,单独判。 */ -const SHELL_METACHARS = [';', '|', '&', '`', '>', '<'] as const +const SHELL_METACHARS = [';', '|', '&', '`', '>', '<', '\n', '\r'] as const function hasShellMetachar(command: string): boolean { if (command.includes('$(')) return true diff --git a/packages/core/src/htbp/helpDsl.ts b/packages/core/src/htbp/helpDsl.ts index a69cd038..3d5f4371 100644 --- a/packages/core/src/htbp/helpDsl.ts +++ b/packages/core/src/htbp/helpDsl.ts @@ -56,6 +56,7 @@ function attrLines(key: string, value: string): string[] { export function renderHelpDsl(model: HelpModel): string { const lines: string[] = [HTBP_HELP_HEADER] lines.push(nodeLine(model.node.path, model.node.kind, model.node.description)) + if (model.deviceContext !== undefined) lines.push(`deviceContext ${JSON.stringify(model.deviceContext)}`) if (model.hint !== undefined) lines.push(`hint ${collapseToOneLine(model.hint)}`) if (model.note !== undefined) lines.push(`note "${collapseToOneLine(model.note)}"`) for (const cmd of model.cmds) { @@ -119,6 +120,7 @@ export function renderHelpJson(model: HelpModel): HelpJson { node: { path: model.node.path, kind: model.node.kind, description: model.node.description }, cmds, } + if (model.deviceContext !== undefined) json.deviceContext = model.deviceContext if (model.hint !== undefined) json.hint = model.hint if (model.note !== undefined) json.note = model.note if (model.feedback !== undefined && model.feedback.length > 0) { diff --git a/packages/core/src/htbp/helpMarkdown.ts b/packages/core/src/htbp/helpMarkdown.ts index 88a60d04..8dfc1528 100644 --- a/packages/core/src/htbp/helpMarkdown.ts +++ b/packages/core/src/htbp/helpMarkdown.ts @@ -75,6 +75,11 @@ export function renderHelpMarkdown(model: HelpModel): string { if (model.node.description.trim() !== '') { out.push(model.node.description.trim(), '') } + if (model.deviceContext !== undefined) { + out.push('## Device context', '') + out.push('Last reported environment (self-reported; presence does not refresh this snapshot):', '') + out.push('```json', JSON.stringify(model.deviceContext, null, 2), '```', '') + } if (model.hint !== undefined) { out.push(`> **Next step**: ${collapseToOneLine(model.hint)}`, '') } diff --git a/packages/core/src/htbp/model.ts b/packages/core/src/htbp/model.ts index 01d4b50a..fe58f00b 100644 --- a/packages/core/src/htbp/model.ts +++ b/packages/core/src/htbp/model.ts @@ -6,7 +6,15 @@ * 不各自持有数据,故两种表现不可能字段漂移。 */ -import type { Action, NodeKind, TreePath } from '../types' +import type { Action, DeviceEnvironment, NodeKind, Timestamp, TreePath } from '../types' +import type { Presence } from '../device/presence' + +/** Last hello snapshot; presence is derived separately from the current connection. */ +export interface DeviceContext { + environment: DeviceEnvironment + presence?: Presence + reportedAt: Timestamp +} /** ~help 默认 feedback 区块的单条形态(只露 id+title+score,详情经 system/feedback get 下钻)。 */ export interface HelpFeedbackItem { @@ -61,6 +69,7 @@ export interface HelpModel { /** directory 节点携带:上级/自身 `~help` 列出的子节点。 */ children?: ChildRef[] cmds: CmdSpec[] + deviceContext?: DeviceContext /** * Agent feedback 默认区块(该 path 头部可见条目,网关 ~help 注入;空数组不注入)。 * DSL 渲染为 `feedback` 头行 + 缩进条目行(未知行忽略通道);JSON 同名字段;Markdown Feedback 节。 @@ -94,6 +103,7 @@ export interface HelpJson { /** directory 节点携带。 */ children?: ChildRef[] cmds: CmdSpec[] + deviceContext?: DeviceContext /** Agent feedback 默认区块,对应 DSL 的 `feedback` 块(有条目才出现)。 */ feedback?: HelpFeedbackItem[] /** 下一步指引,对应 DSL 的 `hint` 行(有值才出现)。 */ diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 49cb0ccd..f6114960 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -22,6 +22,7 @@ export * from './context/ttl' export * from './context/types' export * from './deployment' export * from './device/client' +export * from './device/environment' export * from './device/frames' export * from './device/helpModel' export * from './device/mailbox' diff --git a/packages/core/src/node/index.ts b/packages/core/src/node/index.ts index 35d6f148..8ecc8a84 100644 --- a/packages/core/src/node/index.ts +++ b/packages/core/src/node/index.ts @@ -7,5 +7,6 @@ export * from './fsObjectStore' export * from './processExecution' +export * from './processSessions' export * from './shellExecutor' export * from './structuredCommand' diff --git a/packages/core/src/node/processSessions.ts b/packages/core/src/node/processSessions.ts new file mode 100644 index 00000000..efe3bda8 --- /dev/null +++ b/packages/core/src/node/processSessions.ts @@ -0,0 +1,464 @@ +/** Runtime-owned POSIX process sessions. Transport cancellation never owns the child. */ +import { type ChildProcess, spawn } from 'node:child_process' +import { createHash, randomUUID } from 'node:crypto' +import { StringDecoder } from 'node:string_decoder' +import { z } from 'zod/v4' +import type { DeviceAbortSignal } from '../device/client' +import type { ToolSpec } from '../tool/types' +import { type ProcessSessionChunk, type ProcessSessionList, processSessionListSchema, type ProcessSessionObservation, processSessionObservationSchema, type ProcessSessionSummary, processSessionSummarySchema } from '../device/processSessionContract' +import { OperationRegistry } from '../operation/registry' +import { canonicalizePath } from '../tree/path' +import { TBError } from '../errors' + +export const PROCESS_SESSION_DEFAULT_RUNTIME_MS = 30 * 60_000 +export const PROCESS_SESSION_MAX_RUNTIME_MS = 24 * 60 * 60_000 +export const PROCESS_SESSION_STREAM_LIMIT_BYTES = 1024 * 1024 +export const PROCESS_SESSION_MAX_START_WINDOW_MS = 5 * 60_000 +export interface ProcessSessionBinding { + confirm?: boolean + description?: string + effect: 'read' | 'write' | 'destructive' + inputSchema: unknown + maxRuntimeMs?: number + path: string + prepare(input: Record): { + argv: string[] + cwd?: string + env?: NodeJS.ProcessEnv + executable: string + shell?: boolean + } +} +export interface ProcessSessionContext { + caller?: { keyId: string, owner: string } + signal?: DeviceAbortSignal +} +export interface ProcessSessionManagerOptions { + maxActive?: number + maxRetained?: number + now?: () => number + /** Only Linux is enabled in production; tests may exercise POSIX on a different host. */ + platform?: NodeJS.Platform + retentionMs?: number + runtimeId?: string + streamLimitBytes?: number + terminationGraceMs?: number +} +export interface ProcessSessionManager { + close(): Promise + cmds(path: string): ToolSpec[] + invoke(path: string, action: string, args: Record, context?: ProcessSessionContext): Promise + register(binding: ProcessSessionBinding): void + runtimeId: string +} + +const outputSchemas = { + start: z.toJSONSchema(processSessionSummarySchema), + list: z.toJSONSchema(processSessionListSchema), + observe: z.toJSONSchema(processSessionObservationSchema), + stop: z.toJSONSchema(processSessionSummarySchema), +} +const requestIdSchema = z.string().min(1).max(128) +const listSchema = z.strictObject({ + requestId: requestIdSchema.optional(), state: processSessionSummarySchema.shape.state.optional(), + cursor: z.string().max(256).optional(), limit: z.number().int().min(1).max(32).optional(), +}) +const observeSchema = z.strictObject({ + sessionId: z.string().min(1).max(256), cursor: z.number().int().nonnegative().optional(), + limitBytes: z.number().int().min(4).max(256 * 1024).optional(), + waitMs: z.number().int().min(0).max(20_000).optional(), +}) +const stopSchema = z.strictObject({ sessionId: z.string().min(1).max(256) }) + +function startSchema(binding: ProcessSessionBinding) { + if (!(binding.inputSchema instanceof z.ZodType)) { + throw new TBError('invalid_argument', 'process session binding requires a Zod input schema') + } + return z.strictObject({ + input: binding.inputSchema, + requestId: requestIdSchema, + expectedRuntimeId: z.string().min(1).max(128), + startBefore: z.iso.datetime(), + timeoutMs: z.number().int().positive().max(binding.maxRuntimeMs ?? PROCESS_SESSION_DEFAULT_RUNTIME_MS).optional(), + }) +} +function makeRegistry(binding: ProcessSessionBinding, handler: (action: string, args: unknown, context: ProcessSessionContext) => Promise) { + const registry = new OperationRegistry() + registry.register('start', { + description: `${binding.description ?? 'Start the bound command'}. No stdin/PTY; reconnect within this runtime; daemon restart invalidates sessions.`, + effect: binding.effect, confirm: binding.effect === 'destructive' ? true : binding.confirm, + delivery: 'realtime', inputSchema: startSchema(binding), outputSchema: outputSchemas.start, + }, (args, context) => handler('start', args, context)) + registry.register('list', { + description: 'List this owner’s sessions for this command; includes runtimeId and device time.', + effect: 'read', delivery: 'realtime', inputSchema: listSchema, outputSchema: outputSchemas.list, + }, (args, context) => handler('list', args, context)) + registry.register('observe', { + description: 'Read incremental plain-text logs with an independent UTF-8 byte cursor. Cancelling observation leaves the process running.', + effect: 'read', delivery: 'realtime', inputSchema: observeSchema, outputSchema: outputSchemas.observe, + }, (args, context) => handler('observe', args, context)) + registry.register('stop', { + description: 'Stop the process group; running/stopping states do not claim termination has completed.', + effect: 'write', delivery: 'realtime', inputSchema: stopSchema, outputSchema: outputSchemas.stop, + }, (args, context) => handler('stop', args, context)) + return registry +} +/** Pure metadata assembly: safe without Linux, timers, or child processes. */ +export function processSessionCommands(binding: ProcessSessionBinding): ToolSpec[] { + return makeRegistry(binding, async () => undefined).list() +} + +interface Session { + bytes: { stderr: number, stdout: number } + child: ChildProcess + chunks: ProcessSessionChunk[] + cleanupError?: TBError + dedupeKey: string + done: Promise + drainTimer?: ReturnType + droppedBytes: number + fingerprint: string + keyId: string + killTimer?: ReturnType + nextCursor: number + notify: Set<() => void> + owner: string + path: string + ready: Promise + resolveDone(): void + startBefore: number + summary: ProcessSessionSummary + termination?: 'stopped' | 'timed_out' + timeout?: ReturnType +} +function canonical(value: unknown): string { + if (Array.isArray(value)) return `[${value.map(canonical).join(',')}]` + if (typeof value === 'object' && value !== null) { + return `{${Object.keys(value).sort().map(key => `${JSON.stringify(key)}:${canonical((value as Record)[key])}`).join(',')}}` + } + return JSON.stringify(value) ?? 'null' +} +function active(session: Session): boolean { + return session.summary.state === 'running' || session.summary.state === 'stopping' +} +function copySummary(session: Session): ProcessSessionSummary { + return { ...session.summary } +} +/** Cut only between decoded UTF-8 code points. */ +function prefixBytes(buffer: Buffer, maxBytes: number): Buffer { + let end = Math.min(buffer.length, maxBytes) + while (end > 0 && end < buffer.length && (buffer[end]! & 0xc0) === 0x80) end-- + return buffer.subarray(0, end) +} +function wake(session: Session): void { + for (const notify of [...session.notify]) notify() +} +function append(session: Session, stream: 'stdout' | 'stderr', text: string, limit: number): void { + if (text === '') return + const bytes = Buffer.byteLength(text) + session.nextCursor += bytes + session.chunks.push({ seq: session.nextCursor, stream, text }) + session.bytes[stream] += bytes + while (session.bytes[stream] > limit) { + const index = session.chunks.findIndex(chunk => chunk.stream === stream) + const chunk = session.chunks[index]! + const buffer = Buffer.from(chunk.text) + const excess = session.bytes[stream] - limit + let drop = Math.min(excess, buffer.length) + while (drop < buffer.length && (buffer[drop]! & 0xc0) === 0x80) drop++ + session.bytes[stream] -= drop + session.droppedBytes += drop + if (drop === buffer.length) session.chunks.splice(index, 1) + else chunk.text = buffer.subarray(drop).toString('utf8') + } + // Byte budgets also need a fragment-count bound: alternating one-byte writes + // otherwise retain millions of JS objects despite a small text budget. + while (session.chunks.length > 2048) { + const removed = session.chunks.shift()! + const removedBytes = Buffer.byteLength(removed.text) + session.bytes[removed.stream] -= removedBytes + session.droppedBytes += removedBytes + } + wake(session) +} +function groupSignal(session: Session, signal: NodeJS.Signals): boolean { + const pid = session.child.pid + if (pid === undefined) return true + try { + process.kill(-pid, signal) + if (signal === 'SIGKILL') session.cleanupError = undefined + return true + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ESRCH') { + session.cleanupError = undefined + return true + } + // Timer and child-event callbacks must not throw and crash the daemon. + session.cleanupError = new TBError('unavailable', 'could not signal the process group; termination is unconfirmed') + session.summary.state = 'stopping' + wake(session) + return false + } +} +function waitForChange(session: Session, waitMs: number, signal?: DeviceAbortSignal): Promise { + return new Promise((resolve, reject) => { + const timers: { wait?: ReturnType } = {} + const handlers = { + finish: () => { + clearTimeout(timers.wait) + session.notify.delete(handlers.finish) + signal?.removeEventListener('abort', handlers.abort) + resolve() + }, + abort: () => { + clearTimeout(timers.wait) + session.notify.delete(handlers.finish) + signal?.removeEventListener('abort', handlers.abort) + reject(new TBError('unavailable', 'session observation cancelled')) + }, + } + timers.wait = setTimeout(handlers.finish, waitMs) + const { finish, abort } = handlers + session.notify.add(finish) + signal?.addEventListener('abort', abort, { once: true }) + if (signal?.aborted === true) abort() + }) +} +function observation(session: Session, cursor: number, limit: number): ProcessSessionObservation { + if (cursor > session.nextCursor) throw new TBError('invalid_argument', 'cursor is ahead of this session') + let nextCursor = cursor + let remaining = limit + let gap = false + const chunks: ProcessSessionChunk[] = [] + for (const chunk of session.chunks) { + if (chunk.seq <= nextCursor) continue + const buffer = Buffer.from(chunk.text) + const start = chunk.seq - buffer.length + if (start > nextCursor) { + gap = true + nextCursor = start + } + let offset = Math.max(0, nextCursor - start) + if (offset < buffer.length && (buffer[offset]! & 0xc0) === 0x80) { + throw new TBError('invalid_argument', 'cursor must be a UTF-8 boundary returned by observe') + } + const selected = prefixBytes(buffer.subarray(offset), remaining) + if (selected.length === 0) break + offset += selected.length + nextCursor = start + offset + remaining -= selected.length + chunks.push({ seq: nextCursor, stream: chunk.stream, text: selected.toString('utf8') }) + if (remaining === 0 || offset < buffer.length) break + } + if (chunks.length === 0 && nextCursor < session.nextCursor) { + gap = true + nextCursor = session.nextCursor + } + return { ...copySummary(session), chunks, nextCursor, gap, droppedBytes: session.droppedBytes } +} + +export function createProcessSessionManager(opts: ProcessSessionManagerOptions = {}): ProcessSessionManager { + const runtimeId = opts.runtimeId ?? randomUUID() + const now = opts.now ?? Date.now + const grace = opts.terminationGraceMs ?? 1000 + const sessions = new Map() + const dedupe = new Map() + const registries = new Map>() + let closed = false + let closing: Promise | undefined + const cleanup = () => { + for (const session of sessions.values()) { + if (!active(session) && now() > session.startBefore + && now() - Date.parse(session.summary.completedAt!) >= (opts.retentionMs ?? 30 * 60_000)) { + sessions.delete(session.summary.sessionId) + dedupe.delete(session.dedupeKey) + } + } + } + const terminate = (session: Session, reason: 'stopped' | 'timed_out') => { + if (!active(session) || (session.termination !== undefined && session.cleanupError === undefined)) return + clearTimeout(session.killTimer) + session.termination ??= reason + session.summary.state = 'stopping' + groupSignal(session, 'SIGTERM') + session.killTimer = setTimeout(() => { + groupSignal(session, 'SIGKILL') + }, grace) + wake(session) + } + const owned = (path: string, owner: string, id: string): Session => { + const session = sessions.get(id) + if (session === undefined || session.path !== path || session.owner !== owner) { + throw TBError.notFound('session expired or not found in this runtime') + } + return session + } + const start = async (binding: ProcessSessionBinding, raw: unknown, context: ProcessSessionContext) => { + const args = raw as z.infer> + if (context.signal?.aborted === true) throw new TBError('unavailable', 'session start cancelled before acceptance') + if ((opts.platform ?? process.platform) !== 'linux') throw TBError.unimplemented('process sessions currently require Linux') + if (args.expectedRuntimeId !== runtimeId) throw new TBError('conflict', 'daemon runtime changed; previous sessions are invalid') + const deadline = Date.parse(args.startBefore) + if (deadline <= now() || deadline > now() + PROCESS_SESSION_MAX_START_WINDOW_MS) { + throw new TBError('invalid_argument', 'startBefore must be in the next five minutes') + } + const key = JSON.stringify([binding.path, context.caller!.owner, args.requestId]) + const timeoutMs = args.timeoutMs ?? binding.maxRuntimeMs ?? PROCESS_SESSION_DEFAULT_RUNTIME_MS + const fingerprint = createHash('sha256').update(canonical({ input: args.input, timeoutMs })).digest('hex') + const existing = dedupe.get(key) + if (existing !== undefined) { + if (existing.fingerprint !== fingerprint) throw new TBError('conflict', 'requestId already belongs to different input') + // A retry may renew its short receive deadline; never discard its dedupe + // record while any accepted replay window can still arrive. + existing.startBefore = Math.max(existing.startBefore, deadline) + await existing.ready + return copySummary(existing) + } + if ([...sessions.values()].filter(active).length >= (opts.maxActive ?? 8) + || sessions.size >= (opts.maxRetained ?? 32)) { + throw new TBError('rate_limited', 'process session capacity reached; wait for retained sessions to expire') + } + const prepared = binding.prepare(args.input as Record) + let child: ChildProcess + try { + child = spawn(prepared.executable, prepared.argv, { cwd: prepared.cwd, env: prepared.env ?? {}, shell: prepared.shell ?? false, detached: true, stdio: ['ignore', 'pipe', 'pipe'] }) + } catch { + throw new TBError('internal', 'could not start bound command') + } + let resolveDone!: () => void + let resolveReady!: () => void + let rejectReady!: (reason: Error) => void + const session: Session = { + summary: { sessionId: `${runtimeId}:${randomUUID()}`, runtimeId, requestId: args.requestId, state: 'running', startedAt: new Date(now()).toISOString() }, + owner: context.caller!.owner, keyId: context.caller!.keyId, path: binding.path, + fingerprint, dedupeKey: key, startBefore: deadline, child, + ready: new Promise((resolve, reject) => { + resolveReady = resolve + rejectReady = reject + }), + done: new Promise((resolve) => { + resolveDone = resolve + }), resolveDone: () => resolveDone(), + chunks: [], bytes: { stdout: 0, stderr: 0 }, droppedBytes: 0, nextCursor: 0, notify: new Set(), + } + sessions.set(session.summary.sessionId, session) + dedupe.set(key, session) + const decoders = { stdout: new StringDecoder('utf8'), stderr: new StringDecoder('utf8') } + const limit = opts.streamLimitBytes ?? PROCESS_SESSION_STREAM_LIMIT_BYTES + for (const stream of ['stdout', 'stderr'] as const) { + child[stream]?.on('data', (chunk: Buffer) => append(session, stream, decoders[stream].write(chunk), limit)) + } + const finish = (code: number | null, signal: NodeJS.Signals | null) => { + if (!active(session)) return + clearTimeout(session.timeout) + clearTimeout(session.killTimer) + clearTimeout(session.drainTimer) + // Descendants may close stdio before exiting; never abandon the process group. + if (!groupSignal(session, 'SIGKILL')) { + session.resolveDone() + return + } + for (const stream of ['stdout', 'stderr'] as const) append(session, stream, decoders[stream].end(), limit) + session.summary = { + ...session.summary, state: session.termination ?? 'exited', completedAt: new Date(now()).toISOString(), + ...(code !== null ? { exitCode: code } : {}), ...(signal !== null ? { signal } : {}), + } + wake(session) + session.resolveDone() + } + child.once('spawn', resolveReady) + child.once('error', () => { + clearTimeout(session.timeout) + sessions.delete(session.summary.sessionId) + dedupe.delete(key) + rejectReady(new TBError('internal', 'could not start bound command')) + session.resolveDone() + }) + child.once('exit', (code, signal) => { + // A shell descendant may keep the pipes open after its leader exits. + groupSignal(session, 'SIGTERM') + session.drainTimer = setTimeout(() => { + groupSignal(session, 'SIGKILL') + child.stdout?.destroy() + child.stderr?.destroy() + finish(code, signal) + }, grace) + }) + child.once('close', finish) + session.timeout = setTimeout(() => terminate(session, 'timed_out'), timeoutMs) + await session.ready + return copySummary(session) + } + return { + runtimeId, + register(rawBinding) { + const path = canonicalizePath(rawBinding.path) + if (path === '' || registries.has(path)) throw new TBError('invalid_argument', 'duplicate or empty process session path') + if (rawBinding.maxRuntimeMs !== undefined + && (!Number.isInteger(rawBinding.maxRuntimeMs) || rawBinding.maxRuntimeMs < 1 || rawBinding.maxRuntimeMs > PROCESS_SESSION_MAX_RUNTIME_MS)) { + throw new TBError('invalid_argument', 'session maxRuntimeMs must be between 1 ms and 24 hours') + } + const binding = { ...rawBinding, path } + registries.set(path, makeRegistry(binding, async (action, raw, context) => { + if (action === 'start') return start(binding, raw, context) + if (action === 'list') { + const args = { limit: 32, ...raw as z.infer } + const candidates = [...sessions.values()].filter(session => session.owner === context.caller!.owner && session.path === path + && (args.requestId === undefined || session.summary.requestId === args.requestId) + && (args.state === undefined || session.summary.state === args.state)) + const offset = args.cursor === undefined ? 0 : candidates.findIndex(session => session.summary.sessionId === args.cursor) + 1 + if (args.cursor !== undefined && offset === 0) throw new TBError('invalid_argument', 'list cursor expired or not found') + const items = candidates.slice(offset, offset + args.limit).map(copySummary) + return { + runtimeId, now: new Date(now()).toISOString(), items, + ...(offset + items.length < candidates.length ? { cursor: items.at(-1)!.sessionId } : {}), + } satisfies ProcessSessionList + } + const args = { cursor: 0, limitBytes: 64 * 1024, waitMs: 0, ...raw as z.infer } + const session = owned(path, context.caller!.owner, args.sessionId) + if (action === 'stop') { + terminate(session, 'stopped') + if (active(session)) await waitForChange(session, grace + 500) + if (active(session)) await waitForChange(session, grace + 500) + if (session.cleanupError !== undefined) throw session.cleanupError + return copySummary(session) + } + if (args.cursor === session.nextCursor && active(session) && args.waitMs > 0) { + await waitForChange(session, args.waitMs, context.signal) + } + return observation(session, args.cursor, args.limitBytes) + })) + }, + cmds(path) { + const registry = registries.get(canonicalizePath(path)) + if (registry === undefined) throw TBError.notFound('unknown process session binding') + return registry.list() + }, + async invoke(path, action, args, context = {}) { + if (closed) throw new TBError('unavailable', 'daemon runtime is closed; sessions are invalid') + if (!context.caller?.owner || !context.caller.keyId) throw new TBError('unavailable', 'authenticated caller context required for process sessions') + cleanup() + const registry = registries.get(canonicalizePath(path)) + if (registry === undefined) throw TBError.notFound('unknown process session binding') + return registry.invoke(action, args, context) + }, + close() { + if (closing !== undefined) return closing + closed = true + closing = (async () => { + const pending = [...sessions.values()].filter(active) + for (const session of pending) terminate(session, 'stopped') + let timer: ReturnType | undefined + try { + await Promise.race([ + Promise.all(pending.map(session => session.done)), + new Promise((_, reject) => { timer = setTimeout(() => reject(new TBError('unavailable', 'process shutdown did not complete')), grace + 3000) }), + ]) + const failure = pending.find(session => session.cleanupError !== undefined) + if (failure?.cleanupError !== undefined) throw failure.cleanupError + } finally { clearTimeout(timer) } + })() + return closing + }, + } +} diff --git a/packages/core/src/node/shellExecutor.ts b/packages/core/src/node/shellExecutor.ts index bada8502..6f1ef935 100644 --- a/packages/core/src/node/shellExecutor.ts +++ b/packages/core/src/node/shellExecutor.ts @@ -10,6 +10,8 @@ */ import { spawn as nodeSpawn } from 'node:child_process' +import { z } from 'zod/v4' +import type { ProcessSessionBinding } from './processSessions' import { executeProcess, type ProcessExecutionResult, @@ -77,3 +79,31 @@ export function createShellExecutor(opts: ShellExecutorOptions = {}): ShellExecu ) } } + +/** Explicit session sibling of shell/exec; never widens the configured allowlist. */ +export function createShellSessionBinding(path: string, allow?: string[]): ProcessSessionBinding { + return { + path, + description: 'Shell process sessions; uses the shell allowlist and daemon working directory unless cwd is supplied', + effect: 'destructive', + confirm: true, + inputSchema: z.strictObject({ + command: z.string().min(1).refine(value => !value.includes('\0'), 'must not contain NUL'), + cwd: z.string().min(1).refine(value => !value.includes('\0'), 'must not contain NUL').optional(), + }), + prepare: (input) => { + const command = input.command as string + if (!isCommandAllowed(command, allow)) { + throw new TBError('permission_denied', 'command not in shell allowlist') + } + return { + executable: command, + argv: [], + shell: true, + env: Object.fromEntries(['HOME', 'PATH', 'LANG', 'LC_ALL', 'LC_CTYPE', 'TMPDIR', 'USER', 'LOGNAME', 'TZ'] + .flatMap(name => process.env[name] === undefined ? [] : [[name, process.env[name]!]])), + ...(input.cwd === undefined ? {} : { cwd: input.cwd as string }), + } + }, + } +} diff --git a/packages/core/src/node/structuredCommand.ts b/packages/core/src/node/structuredCommand.ts index ca076453..89b8b5b7 100644 --- a/packages/core/src/node/structuredCommand.ts +++ b/packages/core/src/node/structuredCommand.ts @@ -5,6 +5,7 @@ import { spawn as nodeSpawn } from 'node:child_process' import { z } from 'zod/v4' +import type { ProcessSessionBinding } from './processSessions' import type { DeviceAbortSignal } from '../device/client' import type { ToolSpec } from '../tool/types' import { @@ -56,6 +57,8 @@ export interface StructuredCommandDefinition { inheritEnv?: string[] maxOutputBytes?: number name: string + /** Opt-in session tool; path is relative to the device mount. */ + session?: { maxRuntimeMs?: number, path: string } timeoutMs?: number } @@ -82,6 +85,7 @@ export interface StructuredCommandRuntime { ): Promise path: string profile: StructuredCommandProfile + sessionBindings: ProcessSessionBinding[] } export interface StructuredCommandRuntimeOptions { @@ -137,6 +141,10 @@ const commandSchema = z.strictObject({ timeoutMs: z.number().int().positive().max(SHELL_EXEC_DEFAULT_TIMEOUT_MS).optional(), maxOutputBytes: z.number().int().positive().max(STRUCTURED_COMMAND_MAX_OUTPUT_BYTES).optional(), inheritEnv: z.array(envNameSchema).optional(), + session: z.strictObject({ + path: z.string().min(1), + maxRuntimeMs: z.number().int().positive().max(86_400_000).optional(), + }).optional(), }) const profileSchema = z.strictObject({ version: z.literal(STRUCTURED_COMMAND_PROFILE_VERSION), @@ -216,9 +224,23 @@ export function parseStructuredCommandProfile(value: unknown): StructuredCommand ...raw, name, argv, + ...(raw.session === undefined + ? {} + : { session: { + ...raw.session, path: canonicalizePath(raw.session.path), + } }), ...(raw.effect === 'destructive' ? { confirm: true } : {}), } }) + const paths = [path, ...commands.flatMap(command => command.session === undefined ? [] : [command.session.path])] + for (const [index, candidate] of paths.entries()) { + if (candidate === '') throw invalidProfile('session path must not be root') + for (const previous of paths.slice(0, index)) { + if (candidate === previous || candidate.startsWith(`${previous}/`) || previous.startsWith(`${candidate}/`)) { + throw invalidProfile(`session path '${candidate}' conflicts with '${previous}'`) + } + } + } return { version: STRUCTURED_COMMAND_PROFILE_VERSION, path, @@ -350,6 +372,23 @@ export function createStructuredCommandRuntime( path: profile.path, description: profile.description, cmds: registry.list(), + sessionBindings: profile.commands.flatMap(definition => definition.session === undefined + ? [] + : [{ + path: definition.session.path, + description: `${definition.description} (process sessions; ${definition.cwd === undefined ? 'daemon working directory' : 'fixed profile working directory'})`, + effect: definition.effect, + ...(definition.confirm === undefined ? {} : { confirm: definition.confirm }), + ...(definition.session.maxRuntimeMs === undefined ? {} : { maxRuntimeMs: definition.session.maxRuntimeMs }), + inputSchema: inputSchema(definition), + prepare: (args: Record) => ({ + executable: definition.executable, + argv: argvFor(definition, args), + ...(definition.cwd === undefined ? {} : { cwd: definition.cwd }), + env: environmentFor(definition, sourceEnv), + shell: false, + }), + }]), async invoke(command, args, invokeOpts = {}) { const name = canonicalizeSegment(command) return await registry.invoke(name, args, invokeOpts) as ProcessExecutionResult diff --git a/packages/core/src/protocol/wire.ts b/packages/core/src/protocol/wire.ts index da8b5269..95db4fbd 100644 --- a/packages/core/src/protocol/wire.ts +++ b/packages/core/src/protocol/wire.ts @@ -24,12 +24,14 @@ import { import { type Action, ACTIONS, + type DeviceEnvironment, NODE_KINDS, type NodeInput, type NodeKind, type Page, type TreeNode, } from '../types' +import { deviceEnvironmentSchema } from '../device/environment' import { tbErrorBodySchema } from './errorWire' export { tbErrorCodeSchema, @@ -101,9 +103,16 @@ export const helpCommandSchema = z.object({ scope: actionSchema, }) +export const deviceContextSchema = z.strictObject({ + environment: deviceEnvironmentSchema, + reportedAt: z.iso.datetime({ offset: true }), + presence: presenceSchema.optional(), +}) + export const helpJsonSchema: z.ZodType = z.object({ children: z.array(helpChildSchema).optional(), cmds: z.array(helpCommandSchema), + deviceContext: deviceContextSchema.optional(), feedback: z.array(helpFeedbackItemSchema).optional(), hint: z.string().optional(), htbp: z.string(), @@ -405,6 +414,8 @@ export const registryNodeSchema = z.object({ createdAt: z.string().optional(), description: z.string(), deviceId: z.string().optional(), + deviceEnvironment: deviceEnvironmentSchema.optional(), + deviceReportedAt: z.iso.datetime({ offset: true }).optional(), kind: nodeKindSchema, lastSeenAt: z.string().optional(), online: z.boolean().optional(), @@ -434,7 +445,9 @@ export interface WireRegistryNode { config?: Record createdAt?: string description: string + deviceEnvironment?: DeviceEnvironment deviceId?: string + deviceReportedAt?: string kind: WireNodeKind lastSeenAt?: string online?: boolean diff --git a/packages/core/src/tree/registry.ts b/packages/core/src/tree/registry.ts index 3f193a8b..5b60c28d 100644 --- a/packages/core/src/tree/registry.ts +++ b/packages/core/src/tree/registry.ts @@ -10,6 +10,7 @@ */ import { + type DeviceEnvironment, LIST_LIMIT_DEFAULT, LIST_LIMIT_MAX, type ListOptions, @@ -221,7 +222,7 @@ export class NodeRegistryStore { node: NodeInput, registeredBy: string, now: Timestamp, - opts: { deviceId?: string, lastSeenAt?: Timestamp, online?: boolean } = {}, + opts: { deviceEnvironment?: DeviceEnvironment, deviceId?: string, deviceReportedAt?: Timestamp, lastSeenAt?: Timestamp, online?: boolean } = {}, ): Promise { const invalid = validatePath(node.path) if (invalid) throw invalid @@ -250,6 +251,8 @@ export class NodeRegistryStore { ...(node.config !== undefined ? { config: node.config } : {}), ...(node.virtualize !== undefined ? { virtualize: node.virtualize } : {}), ...(opts.deviceId !== undefined ? { deviceId: opts.deviceId } : {}), + ...(opts.deviceEnvironment !== undefined ? { deviceEnvironment: opts.deviceEnvironment } : {}), + ...(opts.deviceReportedAt !== undefined ? { deviceReportedAt: opts.deviceReportedAt } : {}), ...(opts.online !== undefined ? { online: opts.online } : {}), ...(opts.lastSeenAt !== undefined ? { lastSeenAt: opts.lastSeenAt } : {}), registeredBy, diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index 55c82cd6..653dc090 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -120,8 +120,11 @@ export interface TreeNode { createdAt: Timestamp /** 一句话;上级 ~help 列子节点时展示。 */ description: string + /** 设备最后一次 hello 的环境快照;仅网关写入,不作为权限依据。 */ + deviceEnvironment?: DeviceEnvironment /** 仅设备挂载根:hello 中声明的稳定设备身份。由设备注册流程写入,不接受普通 NodeInput。 */ deviceId?: string + deviceReportedAt?: Timestamp kind: NodeKind /** 仅 device:最近一次观察到设备存活(hello / 心跳 / 成功调用)的时刻。缺省表示从未观察或旧数据; * freshness 判定见 device/presence.ts。写路径专用,不经普通注册面。 */ @@ -170,7 +173,18 @@ export interface McpOAuthClientConfig { clientSecretRef?: string } +/** 有界的设备自报环境;禁止自动携带路径、环境变量或凭据。 */ +export interface DeviceEnvironment { + arch?: string + platform: 'darwin' | 'linux' | 'win32' | 'android' | 'ios' | 'other' + runtime?: 'node' | 'bun' | 'react-native' | 'other' + runtimeId?: string + runtimeVersion?: string +} + export interface DeviceExpose { + /** 宿主显式提供,neutral SDK 不采集本机信息。 */ + environment?: DeviceEnvironment /** 挂 `/fs` context 节点(file provider);支持多根。 */ fs?: { readOnly?: boolean, roots: string[] } /** SDK 自定义节点(路径相对 mountPath)。 */ @@ -262,7 +276,7 @@ export type NodeConfig export type NodeInput = Omit< TreeNode, - 'registeredBy' | 'online' | 'lastSeenAt' | 'deviceId' | 'createdAt' | 'updatedAt' + 'registeredBy' | 'online' | 'lastSeenAt' | 'deviceId' | 'deviceEnvironment' | 'deviceReportedAt' | 'createdAt' | 'updatedAt' > /** 自动物化中间 directory 的 registeredBy 标记。 */ diff --git a/packages/core/test/device/shellAllow.test.ts b/packages/core/test/device/shellAllow.test.ts index 4f0b4c5f..0136bb48 100644 --- a/packages/core/test/device/shellAllow.test.ts +++ b/packages/core/test/device/shellAllow.test.ts @@ -17,6 +17,7 @@ describe('isCommandAllowed:[\'*\'] 全放行', () => { it('含元字符也放行(复合命令由用户显式授权)', () => { expect(isCommandAllowed('echo hi; rm -rf /', ['*'])).toBe(true) expect(isCommandAllowed('cat a | grep b > c', ['*'])).toBe(true) + expect(isCommandAllowed('echo first\nprintf second', ['*'])).toBe(true) }) it('\'*\' 混在列表里不算全放行(仅单值)', () => { expect(isCommandAllowed('anything', ['*', 'echo'])).toBe(false) @@ -54,6 +55,10 @@ describe('isCommandAllowed:argv[0] basename 精确匹配', () => { describe('isCommandAllowed:非 [\'*\'] 时元字符直接拒', () => { const injections = [ 'echo hi; rm -rf /', + 'echo safe\nprintf second-command', + 'echo safe\r\nprintf second-command', + 'echo safe\rprintf second-command', + 'echo "quoted\nnewline"', 'echo a | cat', 'echo a & whoami', 'echo $(whoami)', diff --git a/packages/core/test/node/processSessions.test.ts b/packages/core/test/node/processSessions.test.ts new file mode 100644 index 00000000..2d589f1d --- /dev/null +++ b/packages/core/test/node/processSessions.test.ts @@ -0,0 +1,240 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { setTimeout as delay } from 'node:timers/promises' +import { z } from 'zod/v4' +import { + processSessionListSchema, type ProcessSessionObservation, processSessionObservationSchema, + type ProcessSessionSummary, processSessionSummarySchema, +} from '../../src/device/processSessionContract' +import { + createProcessSessionManager, processSessionCommands, + type ProcessSessionManager, type ProcessSessionManagerOptions, +} from '../../src/node/processSessions' + +const context = { caller: { owner: 'alice', keyId: 'key-1' } } +const managers: ProcessSessionManager[] = [] +const codeSchema = z.strictObject({ code: z.string(), note: z.string().optional() }) +function manager(options: ProcessSessionManagerOptions = {}) { + const value = createProcessSessionManager({ platform: 'linux', terminationGraceMs: 50, ...options }) + value.register({ + path: 'sessions/test', effect: 'destructive', inputSchema: codeSchema, + prepare: input => ({ executable: process.execPath, argv: ['-e', input.code as string], env: {} }), + }) + managers.push(value) + return value +} +function args(value: ProcessSessionManager, code: string, requestId = 'request-1') { + return { input: { code }, expectedRuntimeId: value.runtimeId, requestId, startBefore: new Date(Date.now() + 60_000).toISOString() } +} +async function start(value: ProcessSessionManager, code: string, requestId?: string) { + return await value.invoke('sessions/test', 'start', args(value, code, requestId), context) as ProcessSessionSummary +} +async function observe(value: ProcessSessionManager, sessionId: string, options: Record = {}) { + return await value.invoke('sessions/test', 'observe', { sessionId, ...options }, context) as ProcessSessionObservation +} +async function terminal(value: ProcessSessionManager, sessionId: string) { + let result = await observe(value, sessionId) + for (let i = 0; i < 50 && ['running', 'stopping'].includes(result.state); i++) { + result = await observe(value, sessionId, { cursor: result.nextCursor, waitMs: 100 }) + } + expect(['running', 'stopping']).not.toContain(result.state) + return result +} +afterEach(async () => { + await Promise.all(managers.splice(0).map(value => value.close())) +}) + +describe('runtime-owned process sessions', () => { + it('starts once under concurrent replay, fingerprints canonical input and keeps owner/path isolation', async () => { + const value = manager() + const request = args(value, 'setInterval(() => {}, 1000)') + const [first, second] = await Promise.all([ + value.invoke('sessions/test', 'start', request, context), + value.invoke('sessions/test', 'start', { ...request, input: { code: request.input.code } }, context), + ]) as ProcessSessionSummary[] + expect(first!.sessionId).toBe(second!.sessionId) + await expect(value.invoke('sessions/test', 'start', { ...request, input: { code: 'process.exit()' } }, context)).rejects.toMatchObject({ code: 'conflict' }) + await expect(value.invoke('sessions/test', 'observe', { sessionId: first!.sessionId }, { caller: { owner: 'bob', keyId: 'key-2' } })).rejects.toMatchObject({ code: 'not_found' }) + value.register({ path: 'sessions/other', effect: 'read', inputSchema: z.strictObject({}), prepare: () => ({ executable: process.execPath, argv: [] }) }) + await expect(value.invoke('sessions/other', 'stop', { sessionId: first!.sessionId }, context)).rejects.toMatchObject({ code: 'not_found' }) + const sameOwner = await value.invoke('sessions/test', 'list', {}, { caller: { owner: 'alice', keyId: 'rotated' } }) + expect(processSessionListSchema.parse(sameOwner).items).toHaveLength(1) + }) + + it('fails closed for missing identity, invalid inputs, expired starts and old runtimes', async () => { + const value = manager() + const request = args(value, 'process.exit()') + await expect(value.invoke('sessions/test', 'start', request)).rejects.toMatchObject({ code: 'unavailable' }) + await expect(value.invoke('sessions/test', 'start', { ...request, owner: 'alice' }, context)).rejects.toMatchObject({ code: 'invalid_argument' }) + await expect(value.invoke('sessions/test', 'start', { ...request, input: { code: 'ok', executable: 'sh' } }, context)).rejects.toMatchObject({ code: 'invalid_argument' }) + await expect(value.invoke('sessions/test', 'start', { ...request, expectedRuntimeId: 'old' }, context)).rejects.toMatchObject({ code: 'conflict' }) + await expect(value.invoke('sessions/test', 'start', { ...request, startBefore: new Date(0).toISOString() }, context)).rejects.toMatchObject({ code: 'invalid_argument' }) + await expect(value.invoke('sessions/test', 'start', { ...request, startBefore: new Date(Date.now() + 600_000).toISOString() }, context)).rejects.toMatchObject({ code: 'invalid_argument' }) + await expect(value.invoke('sessions/test', 'start', { ...request, timeoutMs: 86_400_001 }, context)).rejects.toMatchObject({ code: 'invalid_argument' }) + }) + + it('preserves nonzero exit, incremental independent readers and UTF-8 fragments', async () => { + const value = manager() + const session = await start(value, 'const b=Buffer.from(\'你😀好\');process.stdout.write(b.subarray(0,2));setTimeout(()=>{process.stdout.write(b.subarray(2));process.stderr.write(\'err\');process.exitCode=7},20)') + await terminal(value, session.sessionId) + const first = await observe(value, session.sessionId, { limitBytes: 4 }) + const second = await observe(value, session.sessionId, { limitBytes: 4 }) + expect(first).toEqual(second) + let cursor = 0 + const text = { stdout: '', stderr: '' } + while (true) { + const fragment = await observe(value, session.sessionId, { cursor, limitBytes: 4 }) + for (const chunk of fragment.chunks) text[chunk.stream] += chunk.text + if (fragment.nextCursor === cursor) break + cursor = fragment.nextCursor + } + expect(text).toEqual({ stdout: '你😀好', stderr: 'err' }) + expect(first.exitCode).toBe(7) + expect(first.state).toBe('exited') + await expect(observe(value, session.sessionId, { cursor: 1 })).rejects.toMatchObject({ code: 'invalid_argument' }) + await expect(observe(value, session.sessionId, { cursor: 1000 })).rejects.toMatchObject({ code: 'invalid_argument' }) + }) + + it('reports dropped output with bounded per-stream storage, including a huge UTF-8 write', async () => { + const value = manager({ streamLimitBytes: 32 }) + const session = await start(value, 'process.stdout.write(\'你\'.repeat(100000));process.stderr.write(\'😀\'.repeat(100000))') + await terminal(value, session.sessionId) + const output = await observe(value, session.sessionId) + expect(output.gap).toBe(true) + expect(output.droppedBytes).toBeGreaterThan(600_000) + expect(Buffer.byteLength(output.chunks.filter(chunk => chunk.stream === 'stdout').map(chunk => chunk.text).join(''))).toBeLessThanOrEqual(32) + expect(Buffer.byteLength(output.chunks.filter(chunk => chunk.stream === 'stderr').map(chunk => chunk.text).join(''))).toBeLessThanOrEqual(32) + expect(output.chunks.map(chunk => chunk.text).join('')).not.toContain('\ufffd') + expect((await observe(value, session.sessionId, { cursor: output.nextCursor })).chunks).toEqual([]) + }) + + it('observation cancellation leaves the child running, and explicit stop waits for a real exit', async () => { + const value = manager() + const controller = new AbortController() + const session = await value.invoke('sessions/test', 'start', args(value, 'setInterval(() => {}, 1000)'), { ...context, signal: controller.signal }) as ProcessSessionSummary + const waiting = value.invoke('sessions/test', 'observe', { sessionId: session.sessionId, waitMs: 20_000 }, { ...context, signal: controller.signal }) + controller.abort() + await expect(waiting).rejects.toMatchObject({ code: 'unavailable' }) + expect((await observe(value, session.sessionId)).state).toBe('running') + const stopped = await value.invoke('sessions/test', 'stop', { sessionId: session.sessionId }, context) + expect(processSessionSummarySchema.parse(stopped).state).toBe('stopped') + expect(await value.invoke('sessions/test', 'stop', { sessionId: session.sessionId }, context)).toEqual(stopped) + }) + + it('kills a SIGTERM-ignoring process on timeout and distinguishes the terminal state', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] }) + const value = manager() + let session: ProcessSessionSummary + try { + session = await value.invoke('sessions/test', 'start', { + ...args(value, 'process.on(\'SIGTERM\',()=>{});process.stdout.write(\'ready\');setInterval(()=>{},1000)'), timeoutMs: 200, + }, context) as ProcessSessionSummary + // Advance the runtime deadline only after the real child installed its handler. + // CPU contention must not turn this escalation test into a startup-speed test. + let ready = false + for (let attempt = 0; attempt < 200 && !ready; attempt++) { + ready = (await observe(value, session.sessionId)).chunks.some(chunk => chunk.text.includes('ready')) + if (!ready) await delay(10) + } + expect(ready).toBe(true) + await vi.advanceTimersByTimeAsync(251) + } finally { + vi.useRealTimers() + } + const result = await terminal(value, session!.sessionId) + expect(result.state).toBe('timed_out') + expect(result.signal).toBe('SIGKILL') + }) + + it('close cleans the process group including a descendant and invalidates the runtime', async () => { + const value = manager() + const session = await start(value, 'const {spawn}=require(\'node:child_process\');const child=spawn(process.execPath,[\'-e\',\'setInterval(()=>{},1000)\'],{stdio:\'ignore\'});process.stdout.write(String(child.pid));setInterval(()=>{},1000)') + let output = await observe(value, session.sessionId, { waitMs: 1000 }) + if (output.chunks.length === 0) output = await observe(value, session.sessionId, { waitMs: 1000 }) + const descendantPid = Number(output.chunks.map(chunk => chunk.text).join('')) + expect(descendantPid).toBeGreaterThan(0) + await value.close() + await expect(value.invoke('sessions/test', 'list', {}, context)).rejects.toMatchObject({ code: 'unavailable' }) + // POSIX delivery can precede OS reaping by a scheduling tick. + await vi.waitFor(() => expect(() => process.kill(descendantPid, 0)).toThrow(), { timeout: 2000, interval: 10 }) + }) + + it('holds dedupe records until the acceptance deadline and rejects over-capacity starts', async () => { + let clock = Date.now() + const value = manager({ maxActive: 1, maxRetained: 1, retentionMs: 0, now: () => clock }) + const request = args(value, 'process.exit()') + const session = await value.invoke('sessions/test', 'start', request, context) as ProcessSessionSummary + await terminal(value, session.sessionId) + expect(await value.invoke('sessions/test', 'start', request, context)).toMatchObject({ sessionId: session.sessionId }) + await expect(start(value, 'process.exit()', 'other')).rejects.toMatchObject({ code: 'rate_limited' }) + clock += 61_000 + await expect(value.invoke('sessions/test', 'start', request, context)).rejects.toMatchObject({ code: 'invalid_argument' }) + expect(processSessionListSchema.parse(await value.invoke('sessions/test', 'list', {}, context)).items).toEqual([]) + }) + + it('spawn failures release capacity; metadata is pure and platform gating is start-only', async () => { + const value = manager({ maxActive: 1, maxRetained: 1 }) + value.register({ path: 'bad', effect: 'read', inputSchema: z.strictObject({}), prepare: () => ({ executable: '/nonexistent/tool-bridge-test', argv: [] }) }) + await expect(value.invoke('bad', 'start', { ...args(value, ''), input: {} }, context)).rejects.toMatchObject({ code: 'internal' }) + const valid = await start(value, 'process.exit()') + await terminal(value, valid.sessionId) + const disabled = manager({ platform: 'darwin' }) + expect(disabled.cmds('sessions/test')).toHaveLength(4) + await expect(start(disabled, 'process.exit()')).rejects.toMatchObject({ code: 'unavailable' }) + }) + + it('contains signal errors from child/timer callbacks and still escalates to SIGKILL', async () => { + const value = manager() + const session = await start(value, 'setInterval(()=>{},1000)') + const kill = process.kill.bind(process) + const spy = vi.spyOn(process, 'kill').mockImplementation((pid, signal) => { + if (pid < 0 && signal === 'SIGTERM') { + throw Object.assign(new Error('synthetic signal permission failure'), { code: 'EPERM' }) + } + return kill(pid, signal) + }) + try { + const stopped = await value.invoke('sessions/test', 'stop', { sessionId: session.sessionId }, context) + expect(stopped).toMatchObject({ state: 'stopped', signal: 'SIGKILL' }) + } finally { + spy.mockRestore() + } + }) + + it('does not apply the synchronous 55-second deadline, and restart invalidates session IDs', async () => { + vi.useFakeTimers() + const value = manager() + let session: ProcessSessionSummary + try { + session = await start(value, 'setInterval(()=>{},1000)') + await vi.advanceTimersByTimeAsync(55_001) + expect((await observe(value, session.sessionId)).state).toBe('running') + } finally { + vi.useRealTimers() + } + const restarted = manager() + await expect(observe(restarted, session!.sessionId)).rejects.toMatchObject({ code: 'not_found' }) + await expect(restarted.invoke('sessions/test', 'start', args(value, 'process.exit()'), context)).rejects.toMatchObject({ code: 'conflict' }) + const cancelled = new AbortController() + cancelled.abort() + await expect(restarted.invoke('sessions/test', 'start', args(restarted, 'process.exit()'), { ...context, signal: cancelled.signal })).rejects.toMatchObject({ code: 'unavailable' }) + expect(processSessionListSchema.parse(await restarted.invoke('sessions/test', 'list', {}, context)).items).toEqual([]) + }) + + it('contract coverage is derived from public commands and every registered scenario executes', async () => { + const value = manager() + const session = await start(value, 'setInterval(()=>{},1000)') + const scenarios = { + start: async () => processSessionSummarySchema.parse(await value.invoke('sessions/test', 'start', args(value, 'setInterval(()=>{},1000)'), context)), + list: async () => processSessionListSchema.parse(await value.invoke('sessions/test', 'list', {}, context)), + observe: async () => processSessionObservationSchema.parse(await observe(value, session.sessionId)), + stop: async () => processSessionSummarySchema.parse(await value.invoke('sessions/test', 'stop', { sessionId: session.sessionId }, context)), + } + const commands = processSessionCommands({ path: 'test', effect: 'destructive', inputSchema: codeSchema, prepare: () => ({ executable: '', argv: [] }) }) + expect(commands.map(command => command.name).sort()).toEqual(Object.keys(scenarios).sort()) + expect(commands.find(command => command.name === 'start')).toMatchObject({ confirm: true, effect: 'destructive', delivery: 'realtime' }) + for (const command of commands) { + expect(command.outputSchema).toBeDefined() + await scenarios[command.name as keyof typeof scenarios]() + } + }) +}) diff --git a/packages/core/test/node/shellExecutor.test.ts b/packages/core/test/node/shellExecutor.test.ts index 691f70e1..54c4770d 100644 --- a/packages/core/test/node/shellExecutor.test.ts +++ b/packages/core/test/node/shellExecutor.test.ts @@ -4,6 +4,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import { createShellExecutor, + createShellSessionBinding, SHELL_EXEC_DEFAULT_TIMEOUT_MS, SHELL_OUTPUT_LIMIT_BYTES, SHELL_TIMEOUT_EXIT_CODE, @@ -77,6 +78,15 @@ describe('白名单前置判定(执行前完成)', () => { expect(spawn).not.toHaveBeenCalled() }) + it.each(['echo safe\nprintf second-command', 'echo safe\r\nprintf second-command'])('rejects multiple shell lines before either execution path: %s', async (command) => { + const spawn = vi.fn() + const exec = createShellExecutor({ allow: ['echo'], spawn: spawn as unknown as SpawnFn }) + await expect(exec(command)).rejects.toMatchObject({ code: 'permission_denied' }) + expect(spawn).not.toHaveBeenCalled() + expect(() => createShellSessionBinding('sessions/shell', ['echo']).prepare({ command })) + .toThrow('command not in shell allowlist') + }) + it('缺省 allow = [] → 一切拒(与默认拒对齐)', async () => { const exec = createShellExecutor({}) await expect(exec('echo hi')).rejects.toMatchObject({ code: 'permission_denied' }) diff --git a/packages/sdk/package.json b/packages/sdk/package.json index bfd514f4..47ae622e 100644 --- a/packages/sdk/package.json +++ b/packages/sdk/package.json @@ -1,6 +1,6 @@ { "name": "@tool-bridge/sdk", - "version": "0.24.0", + "version": "0.25.0", "description": "tool-bridge SDK: embed a TB instance (createToolBridge), register local providers, and connect to a remote gateway", "type": "module", "license": "MIT", diff --git a/packages/sdk/src/client/index.ts b/packages/sdk/src/client/index.ts index 4ffafd13..27e058b2 100644 --- a/packages/sdk/src/client/index.ts +++ b/packages/sdk/src/client/index.ts @@ -102,9 +102,19 @@ export { PRESENCE_STALE_AFTER_MS, } from '@tool-bridge/core/device' +export type { + ProcessSessionChunk, + ProcessSessionList, + ProcessSessionObservation, + ProcessSessionState, + ProcessSessionSummary, +} from '@tool-bridge/core/device' + +export { processSessionListSchema, processSessionObservationSchema, processSessionSummarySchema } from '@tool-bridge/core/device' export { fixedControlPlaneOpenApi } from '@tool-bridge/core/protocol' export type { FixedControlPlaneOpenApi } from '@tool-bridge/core/protocol' + export type { WireAction as Action, WireDeviceOperationDetail as DeviceOperationDetail, diff --git a/packages/sdk/src/device/connection.ts b/packages/sdk/src/device/connection.ts index 25703a5d..e27e196d 100644 --- a/packages/sdk/src/device/connection.ts +++ b/packages/sdk/src/device/connection.ts @@ -9,6 +9,8 @@ import { type DeviceCallContext as CoreDeviceCallContext, type DeviceCallHandler as CoreDeviceCallHandler, DeviceClient, + type DeviceEnvironment, + deviceEnvironmentSchema, type DeviceNodeCmd, normalizePath, PING_FRAME_JSON, @@ -61,6 +63,7 @@ export interface DeviceNodeDefinition { } export interface DeviceClientExpose { + environment?: DeviceEnvironment nodes: readonly DeviceNodeDefinition[] } @@ -273,6 +276,7 @@ function portableExpose(expose: DeviceClientExpose): WireDeviceExpose { throw new TBError('invalid_argument', 'device expose.nodes 至少需要一个节点') } return { + ...(expose.environment === undefined ? {} : { environment: deviceEnvironmentSchema.parse(expose.environment) }), nodes: expose.nodes.map((node) => { if (node.kind !== 'tool' && node.kind !== 'context') { throw new TBError('invalid_argument', 'device 节点只支持 tool/context') diff --git a/packages/sdk/src/device/index.ts b/packages/sdk/src/device/index.ts index 1dd5c4b1..ae41d67a 100644 --- a/packages/sdk/src/device/index.ts +++ b/packages/sdk/src/device/index.ts @@ -55,6 +55,7 @@ export { export type { CallFrame, CancelFrame, + DeviceEnvironment, DeviceFrame, DeviceNodeCmd, ErrorFrame, diff --git a/packages/sdk/test/client/client.test.ts b/packages/sdk/test/client/client.test.ts index 74f123e2..9ec2dde0 100644 --- a/packages/sdk/test/client/client.test.ts +++ b/packages/sdk/test/client/client.test.ts @@ -94,6 +94,19 @@ describe('@tool-bridge/sdk/client', () => { }) }) + it.each(['online', 'stale', 'offline'] as const)('preserves %s device context through getHelp parsing', async (state) => { + const deviceContext = { + environment: { platform: 'linux', runtime: 'node', runtimeVersion: '22.12.0', runtimeId: 'runtime-1' }, + reportedAt: '2026-09-15T00:00:00.000Z', + presence: { state, lastSeenAt: '2026-09-15T00:01:00.000Z' }, + } + const client = createToolBridgeClient({ + baseUrl: 'https://gw.example', sk: 'tbk_fixture', + fetcher: (async () => json({ ...fixture.help, deviceContext })) as typeof fetch, + }) + expect((await client.getHelp('device/build')).deviceContext).toEqual(deviceContext) + }) + it('uses invoke delivery and fixed device operation management routes', async () => { const operation = { attempt: 0, diff --git a/packages/sdk/test/device.test.ts b/packages/sdk/test/device.test.ts index 7405ae8d..4d5cfa64 100644 --- a/packages/sdk/test/device.test.ts +++ b/packages/sdk/test/device.test.ts @@ -121,6 +121,22 @@ describe('@tool-bridge/sdk/device neutral connection', () => { await connection.closed }) + it('neutral device API preserves explicitly supplied environment in hello', async () => { + const harness = factoryHarness() + const environment = { platform: 'ios' as const, arch: 'arm64', runtime: 'react-native' as const, runtimeVersion: '0.82.0' } + const connection = connectDevice({ + baseUrl: 'https://tb.example', deviceId: 'phone-context', + expose: { environment, nodes: [{ path: 'camera', kind: 'tool', description: 'camera' }] }, + credentialProvider: { prepare: () => ({ headers: {} }) }, + webSocketFactory: harness.factory, handler: async () => null, + }) + const socket = await connectAttempt(harness, 1) + socket.open() + expect(helloFrames(socket)[0]).toMatchObject({ expose: { environment } }) + connection.close() + await connection.closed + }) + it('注入 RN transport 与 Authorization,完成 hello/ready/call/result', async () => { const harness = factoryHarness() const calls: unknown[] = [] diff --git a/packages/server/package.json b/packages/server/package.json index 66e1a325..af89f644 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -1,6 +1,6 @@ { "name": "@tool-bridge/server", - "version": "0.23.0", + "version": "0.24.0", "description": "Self-hosted Tool Bridge Node server with PostgreSQL, S3 object storage and WebSocket device channels", "type": "module", "license": "MIT", From ea7132f4efcc0181d6ca3b4924b439b621cb6499 Mon Sep 17 00:00:00 2001 From: DJJ Date: Tue, 15 Sep 2026 19:53:22 +0800 Subject: [PATCH 2/5] =?UTF-8?q?fix(sdk):=20=E5=9C=A8=E5=85=AC=E5=BC=80?= =?UTF-8?q?=E4=BC=9A=E8=AF=9D=E6=A0=A1=E9=AA=8C=E6=8E=A5=E5=8F=A3=E9=9A=94?= =?UTF-8?q?=E7=A6=BB=E6=89=93=E5=8C=85=E5=99=A8=E5=86=85=E8=81=94=E7=9A=84?= =?UTF-8?q?=20Zod=20=E7=B1=BB=E5=9E=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 最终 tarball 检查发现直接导出 Zod schema 会使 neutral client 声明携带实现类型。改为结构化 parse/safeParse 接口,继续复用 core 真源,让 CLI/Dashboard 保持相同校验行为并满足干净消费者约束。 --- packages/sdk/src/client/index.ts | 4 +-- packages/sdk/src/client/processSessions.ts | 32 ++++++++++++++++++++++ 2 files changed, 34 insertions(+), 2 deletions(-) create mode 100644 packages/sdk/src/client/processSessions.ts diff --git a/packages/sdk/src/client/index.ts b/packages/sdk/src/client/index.ts index 27e058b2..0a633d80 100644 --- a/packages/sdk/src/client/index.ts +++ b/packages/sdk/src/client/index.ts @@ -49,6 +49,8 @@ export type { PresignedPutGrant, PutPresignedOptions, } from './presignedPut' +export { processSessionListSchema, processSessionObservationSchema, type ProcessSessionParser, processSessionSummarySchema } from './processSessions' + /** * builtin/system 管理面视图类型的唯一对外出口(真源在 core;core 是 private 包, * CLI/Dashboard 经此消费,不再各自手抄 PluginManifest/CatalogListItem/SecretKeyView 等)。 @@ -109,8 +111,6 @@ export type { ProcessSessionState, ProcessSessionSummary, } from '@tool-bridge/core/device' - -export { processSessionListSchema, processSessionObservationSchema, processSessionSummarySchema } from '@tool-bridge/core/device' export { fixedControlPlaneOpenApi } from '@tool-bridge/core/protocol' export type { FixedControlPlaneOpenApi } from '@tool-bridge/core/protocol' diff --git a/packages/sdk/src/client/processSessions.ts b/packages/sdk/src/client/processSessions.ts new file mode 100644 index 00000000..a3a7472e --- /dev/null +++ b/packages/sdk/src/client/processSessions.ts @@ -0,0 +1,32 @@ +import { + processSessionListSchema as listSchema, + processSessionObservationSchema as observationSchema, + type ProcessSessionList, + type ProcessSessionObservation, + type ProcessSessionSummary, + processSessionSummarySchema as summarySchema, +} from '@tool-bridge/core/device' + +/** Consumer validation stays portable without exposing the bundled validator's types. */ +export interface ProcessSessionParser { + parse(value: unknown): T + safeParse(value: unknown): { data: T, success: true } | { success: false } +} + +function parser(schema: { safeParse(value: unknown): { data: T, success: true } | { success: false } }): ProcessSessionParser { + return { + parse(value) { + const result = schema.safeParse(value) + if (!result.success) throw new Error('invalid device process session response') + return result.data + }, + safeParse(value) { + const result = schema.safeParse(value) + return result.success ? { success: true, data: result.data } : { success: false } + }, + } +} + +export const processSessionListSchema: ProcessSessionParser = parser(listSchema) +export const processSessionObservationSchema: ProcessSessionParser = parser(observationSchema) +export const processSessionSummarySchema: ProcessSessionParser = parser(summarySchema) From 714a068550404e0b2528b1738d086996a45cbd32 Mon Sep 17 00:00:00 2001 From: DJJ Date: Tue, 15 Sep 2026 19:54:46 +0800 Subject: [PATCH 3/5] =?UTF-8?q?feat(dashboard):=20=E7=94=A8=E4=BC=9A?= =?UTF-8?q?=E8=AF=9D=E7=8A=B6=E6=80=81=E4=B8=8E=E7=8B=AC=E7=AB=8B=E6=97=A5?= =?UTF-8?q?=E5=BF=97=E6=B8=B8=E6=A0=87=E6=89=BF=E6=8E=A5=E8=BF=9C=E7=A8=8B?= =?UTF-8?q?=E9=95=BF=E5=91=BD=E4=BB=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 从 compact 帮助按需读取 start 契约后展示会话面板,让关闭日志与停止进程成为独立动作;保留未知启动的请求标识,避免连接中断后重复执行。同步提高 dashboard minor 并覆盖渐进发现和交互回归。 --- packages/dashboard/package.json | 2 +- .../src/components/node/CommandWorkspace.tsx | 6 + .../components/node/DeviceSessionPanel.tsx | 264 ++++++++++++++++++ packages/dashboard/src/lib/api.ts | 1 + packages/dashboard/src/lib/deviceSession.ts | 16 ++ .../src/lib/useProcessSessionCommands.ts | 17 ++ packages/dashboard/src/pages/ToolPage.tsx | 28 +- .../dashboard/test/deviceSession.dom.test.tsx | 114 ++++++++ .../test/deviceSessionDiscovery.dom.test.tsx | 87 ++++++ 9 files changed, 523 insertions(+), 12 deletions(-) create mode 100644 packages/dashboard/src/components/node/DeviceSessionPanel.tsx create mode 100644 packages/dashboard/src/lib/deviceSession.ts create mode 100644 packages/dashboard/src/lib/useProcessSessionCommands.ts create mode 100644 packages/dashboard/test/deviceSession.dom.test.tsx create mode 100644 packages/dashboard/test/deviceSessionDiscovery.dom.test.tsx diff --git a/packages/dashboard/package.json b/packages/dashboard/package.json index 6a644097..fdc3d34f 100644 --- a/packages/dashboard/package.json +++ b/packages/dashboard/package.json @@ -1,6 +1,6 @@ { "name": "@tool-bridge/dashboard", - "version": "0.29.0", + "version": "0.30.0", "description": "Tool Bridge self-hosted Dashboard as prebuilt static assets", "type": "module", "license": "MIT", diff --git a/packages/dashboard/src/components/node/CommandWorkspace.tsx b/packages/dashboard/src/components/node/CommandWorkspace.tsx index c19c3aeb..a3323b4a 100644 --- a/packages/dashboard/src/components/node/CommandWorkspace.tsx +++ b/packages/dashboard/src/components/node/CommandWorkspace.tsx @@ -3,9 +3,12 @@ import { useLocation, useNavigate } from 'react-router' import { useMemo, useState } from 'react' import type { HelpCmd } from '@/lib/types' import { Dialog, DialogContent, DialogTitle } from '@/components/ui/dialog' +import { useProcessSessionCommands } from '@/lib/useProcessSessionCommands' import { safeToolReturnPath, toolHref } from '@/lib/toolNavigation' import { CmdPanel } from '@/components/node/CmdPanel' +import { sessionCommands } from '@/lib/deviceSession' import { cn } from '@/lib/utils' +import { DeviceSessionPanel } from './DeviceSessionPanel' /** * 命令目录打开独立调用页;旧 ?tool 深链接继续在 Inspector 内自动打开弹窗。 @@ -21,6 +24,7 @@ export function CommandWorkspace({ lazySchema: boolean path: string }) { + const sessionCmds = useProcessSessionCommands(path, cmds) const navigate = useNavigate() const location = useLocation() const [query, setQuery] = useState('') @@ -39,6 +43,8 @@ export function CommandWorkspace({ const active = openTool ? cmds.find(cmd => cmd.name === openTool) : undefined + if (sessionCommands(sessionCmds)) return + return (
state === 'running' || state === 'stopping' +const stateLabel: Record = { + running: '运行中', stopping: '正在停止,尚未确认退出', stopped: '已停止', exited: '已退出', timed_out: '运行超时', +} + +function SessionPanel({ commands }: { commands: NonNullable> }) { + const conn = useConn() + const inputId = useId() + const readLifetime = useRef(new AbortController()) + const [page, setPage] = useState(null) + const [input, setInput] = useState('{}') + const [requestFilter, setRequestFilter] = useState('') + const [startRequest, setStartRequest] = useState(null) + const [selected, setSelected] = useState(null) + const [logs, setLogs] = useState('') + const [gap, setGap] = useState(false) + const [following, setFollowing] = useState(true) + const [readError, setReadError] = useState(null) + const [error, setError] = useState(null) + const [busy, setBusy] = useState(false) + const [pending, setPending] = useState<'start' | 'stop' | null>(null) + const cursor = useRef(0) + const alive = useRef(true) + + useEffect(() => { + alive.current = true + readLifetime.current = new AbortController() + return () => { + alive.current = false + readLifetime.current.abort() + } + }, []) + + const request = useCallback(async (command: HelpCmd, args: unknown, signal?: AbortSignal): Promise => { + const response = await invoke(conn, command.path, args, 'json', { signal: signal ?? (command.name === 'list' || command.name === 'observe' ? readLifetime.current.signal : undefined) }) + const schema = command.name === 'list' ? processSessionListSchema : command.name === 'observe' ? processSessionObservationSchema : processSessionSummarySchema + const result = schema.safeParse(response.json) + if (!result.success) throw new Error('invalid process session response') + return result.data as T + }, [conn]) + + const load = useCallback(async (next?: string, signal?: AbortSignal) => { + const result = await request(commands.list, { + limit: 20, + ...(next ? { cursor: next } : {}), + ...(requestFilter.trim() ? { requestId: requestFilter.trim() } : {}), + }, signal) + if (!signal?.aborted && alive.current) setPage(result) + return result + }, [commands.list, request, requestFilter]) + + useEffect(() => { + const controller = new AbortController() + void load(undefined, controller.signal).catch(() => { + if (!controller.signal.aborted) setError('无法读取会话。设备可能离线,或当前连接没有访问权限。') + }) + return () => controller.abort() + }, [load]) + + const sessionId = selected?.sessionId + const runtimeId = selected?.runtimeId + useEffect(() => { + if (!sessionId || !following) return + const controller = new AbortController() + const observe = async () => { + while (!controller.signal.aborted) { + const result = await request(commands.observe, { sessionId, cursor: cursor.current, limitBytes: 65_536, waitMs: 20_000 }, controller.signal) + if (controller.signal.aborted) return + if (result.runtimeId !== runtimeId) throw new Error('runtime changed') + const previous = cursor.current + cursor.current = result.nextCursor + setSelected(result) + setLogs((current) => { + const next = current + result.chunks.map(chunk => chunk.text).join('') + return next.slice(-1_048_576) + }) + if (result.gap) setGap(true) + setReadError(null) + if (!running(result.state) && previous === result.nextCursor) { + setFollowing(false) + return + } + } + } + void observe().catch(() => { + if (!controller.signal.aborted) { + setReadError('读取中断。连接恢复后可从当前日志位置继续;daemon 重启后旧会话失效。') + setFollowing(false) + } + }) + return () => controller.abort() + }, [commands.observe, following, request, runtimeId, sessionId]) + + const choose = (session: ProcessSessionSummary) => { + if (selected?.sessionId === session.sessionId) { + setFollowing(true) + return + } + cursor.current = 0 + setLogs('') + setGap(false) + setReadError(null) + setSelected(session) + setFollowing(true) + } + + const execute = async (action: 'start' | 'stop') => { + setPending(null) + setBusy(true) + setError(null) + try { + if (action === 'start') { + const parsed: unknown = JSON.parse(input) + if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) throw new Error('input') + const context = await request(commands.list, { limit: 1 }) + if (!alive.current) return + const now = Date.parse(context.now) + if (!Number.isFinite(now) || !context.runtimeId) throw new Error('runtime') + const requestId = crypto.randomUUID() + setStartRequest(requestId) + const result = await request(commands.start, { + input: parsed, requestId, expectedRuntimeId: context.runtimeId, + startBefore: new Date(now + 60_000).toISOString(), + }) + if (!alive.current) return + setStartRequest(null) + choose(result) + } else if (selected) { + const result = await request(commands.stop, { sessionId: selected.sessionId }) + if (!alive.current) return + setSelected(result) + } + await load().catch(() => { + if (alive.current) setError('操作结果已显示,但会话列表刷新失败。') + }) + } catch { + if (alive.current) setError(action === 'start' + ? '未取得启动结果。请检查输入;若已有请求标识,请先用它查找原会话,避免重复启动。' + : '未取得停止结果。请继续观察会话的实际状态。') + } finally { + if (alive.current) setBusy(false) + } + } + + const submit = (action: 'start' | 'stop') => { + if (commands[action].confirm) setPending(action) + else void execute(action) + } + + return ( +
+
+

设备进程会话

+

长命令可持续运行。关闭页面或暂停日志只停止查看;停止进程需显式操作。断线可恢复,daemon 重启后会话失效。

+
+
+ +