Jelajahi Sumber

Merge pull request #1038 from deepseek-harness/fix/command-execute-transport-timeout

fix(apiproxy): let commands outlive unary deadline
imccyu 1 bulan lalu
induk
melakukan
c48a50d553

+ 2 - 2
.agents/notes/implemented/architecture/2026-07-19-gui-layering-and-rpc-protocol.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write .agents/notes/implemented/architecture/2026-07-19-gui-layering-and-rpc-protocol.md
-2026-07-19-gui-layering-and-rpc-protocol.md: b9718da4725316c64686adef24827e2984d8723d
-2026-07-19-gui-layering-and-rpc-protocol.zh.md: 2add148054e8f97c65600cd719fb4f8e0283f52d
+2026-07-19-gui-layering-and-rpc-protocol.md: b7081591cf7e5e3c586c74c5a71b4317376135cb
+2026-07-19-gui-layering-and-rpc-protocol.zh.md: 89557182ca7781f4fb59b8daf866aaca96cf20ee

+ 3 - 2
.agents/notes/implemented/architecture/2026-07-19-gui-layering-and-rpc-protocol.md

@@ -204,7 +204,7 @@ The same domain tree as `ApiProxy`, but unary methods **take the business payloa
 | `callUnary` | mint → tap → POST full form → `serverResponseSchema` parse → **rpcId echo check** (mismatch throws) → tap → emit narrow form |
 | `readSse` | streaming fetch (not EventSource), `\n\n` framing, `data:` concatenation, ServerRequest full-form parse, tap, emit narrow `RpcRequest<frame>` |
 | `respond` | client-response passthrough (rpcId is an echo — never minted here); response body parsed by `rpcReceiptSchema` |
-| unary timeout | `AbortSignal.timeout` (default 30s, constructor-tunable); streams have no timeout (long-lived by nature) |
+| unary deadline | Ordinary unary calls use `AbortSignal.timeout` (default 30s, constructor-tunable); user-paced `host.pickDirectory` and `command.execute` omit that deadline but keep caller/connection cancellation; streams have no deadline |
 | `resolveBase` | browser = same-origin origin; no-location environment (Node) = the `http://dsh.internal` fake authority |
 
 ### The instance-level envelope observation aspect
@@ -234,7 +234,7 @@ All four quadrant full forms pass through `onEnvelope`; the base implementation
 
 ## Consequences
 
-Every client shape consumes one contract: adding a unary method is a five-step mechanical change radiating from a single signature, swapping a carrier touches only a `doFetch` subclass, and every wire message is zod-validated, observable through the envelope tap, and reconcilable by rpcId. The accepted costs: two groups of packages need explicit tsconfig paths entries, and the reserved seams (fork/inject/task.list/listModels/hostInstanceId) stay dormant until a real consumer arrives.
+Every client shape consumes one contract: adding a unary method is a five-step mechanical change radiating from a single signature, swapping a carrier touches only a `doFetch` subclass, and every wire message is zod-validated, observable through the envelope tap, and reconcilable by rpcId. Ordinary unary calls remain bounded, while `host.pickDirectory` and `command.execute` may stay pending until the operation finishes or caller/connection cancellation arrives; this accepts that a non-cooperative user-paced operation can hang its request rather than treating valid operation duration as transport failure. The other accepted costs: two groups of packages need explicit tsconfig paths entries, and the reserved seams (fork/inject/task.list/listModels/hostInstanceId) stay dormant until a real consumer arrives.
 
 ## Alternatives considered
 
@@ -252,3 +252,4 @@ Every client shape consumes one contract: adding a unary method is a five-step m
 | A DTO layer (a second wire-only structure set) | Core types reach the browser type-only at zero cost; a DTO is a permanent two-way synchronization tax |
 | Cursor resumption (implementing mux since) | Reconnect = rebuild (opencode-style) covers all v1 needs; the signature keeps the seat, implementation waits for a real consumer |
 | A createApiClient factory function (the original implementation) | Platform differences (transport/observation) are inheritance aspects, not parameters; the class family lets the fixture substitute at the protocol layer instead of wrapping a fake envelope |
+| Applying the 30-second transport deadline to `command.execute` | Command duration is operation work, not a transport-health budget; the deadline kills valid long-running handlers, while caller/connection cancellation already supplies the required stop path |

+ 3 - 2
.agents/notes/implemented/architecture/2026-07-19-gui-layering-and-rpc-protocol.zh.md

@@ -202,7 +202,7 @@ export type ResponseValue<K> =
 | `callUnary` | mint → tap → POST 全形 → `serverResponseSchema` parse → **rpcId 回显校验**(不符即 throw)→ tap → 吐窄形 |
 | `readSse` | streaming fetch(非 EventSource)、`\n\n` 分帧、`data:` 拼接、ServerRequest 全形 parse、tap、吐窄形 `RpcRequest<帧>` |
 | `respond` | client-response 透传(rpcId 是回填,此处不 mint);应答体 `rpcReceiptSchema` parse |
-| unary 超时 | `AbortSignal.timeout`(默认 30s,构造参数可调);流不设超时(长连接本性) |
+| unary 时限 | 普通 unary 调用使用 `AbortSignal.timeout`(默认 30s,构造参数可调);由用户掌控节奏的 `host.pickDirectory` 和 `command.execute` 不设该时限,但保留调用方/连接取消;流不设时限 |
 | `resolveBase` | 浏览器=同源 origin;无 location 环境(Node)=`http://dsh.internal` 假 authority |
 
 ### 实例级 envelope 观测切面
@@ -232,7 +232,7 @@ export type ResponseValue<K> =
 
 ## Consequences
 
-所有 client 形态消费同一契约:加一个 unary 方法是从单一签名辐射的五步机械改动,换载体只动一个 `doFetch` 子类,wire 上每条消息可 zod 校验、可经 envelope tap 观测、可按 rpcId 对账。接受的代价:两组包需要显式 tsconfig paths 条目;预留接缝(fork/inject/task.list/listModels/hostInstanceId)在真实消费者出现前保持休眠。
+所有 client 形态消费同一契约:加一个 unary 方法是从单一签名辐射的五步机械改动,换载体只动一个 `doFetch` 子类,wire 上每条消息可 zod 校验、可经 envelope tap 观测、可按 rpcId 对账。普通 unary 调用仍受时限约束,而 `host.pickDirectory` 与 `command.execute` 可保持挂起,直到操作完成或调用方/连接取消到来;若由用户掌控节奏的操作不自行结束,请求可能一直挂起,这是为避免把合理的操作时长视为传输失败而接受的代价。其余接受的代价:两组包需要显式 tsconfig paths 条目;预留接缝(fork/inject/task.list/listModels/hostInstanceId)在真实消费者出现前保持休眠。
 
 ## Alternatives considered
 
@@ -250,3 +250,4 @@ export type ResponseValue<K> =
 | DTO 层(wire 专用第二套结构) | core 类型 type-only 直达浏览器零成本;DTO 是永久的双向同步税 |
 | cursor 续传(mux since 实装) | 重连=重建(opencode 同款)覆盖 v1 全部需求;签名留座,实装等真实消费者 |
 | createApiClient 工厂函数(原实现) | 平台差异(传输/观测)是继承切面不是参数;类体系让 fixture 在协议层替换而不是包一层假信封 |
+| 对 `command.execute` 应用 30 秒传输时限 | 命令耗时属于操作本身,而非传输健康预算;该时限会终止本应继续运行的长时处理器,调用方/连接取消已提供所需的停止路径 |

+ 2 - 2
packages/host/apiproxy/README.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write packages/host/apiproxy/README.md
-README.md: 73d8afb32f868ca82dfa2d350df089a5d0b9b358
-README.zh.md: 47af18f76302e261e18f682e0d3cf0ee903933db
+README.md: 926c631e4784d5ca9b52b1af93487ce4887511b1
+README.zh.md: 030e6f9dab7e659e8ea6bc2543c676f35ade4b39

File diff ditekan karena terlalu besar
+ 2 - 2
packages/host/apiproxy/README.md


File diff ditekan karena terlalu besar
+ 2 - 2
packages/host/apiproxy/README.zh.md


+ 22 - 12
packages/host/apiproxy/src/fetch/client.ts

@@ -59,9 +59,10 @@ import { llmModelsValueSchema, llmProvidersValueSchema } from '../api/llm.schema
  * Client consumption face of the contract (shape a): same domain tree as ApiProxy, but unary
  * methods take the business payload directly — the carrier mints the rpcId and wraps the
  * envelope. Business code needing the call's rpcId reads it from the RpcResponse echo.
- * Unary methods and respond accept an optional external AbortSignal as the last parameter
- * (merged with the instance timeout via AbortSignal.any; same "signal rides beside the
- * request, never on the wire" discipline as the stream signatures).
+ * Unary methods and respond accept an optional external AbortSignal as the last parameter.
+ * Bounded calls merge it with the instance timeout via AbortSignal.any; user-paced calls
+ * carry only that external signal. In both cases the signal rides beside the request, never
+ * on the wire, like the stream signatures.
  * Stream methods accept an optional onOpen callback: it fires once the SSE transport is
  * readable (response headers received, before any frame) — the "stream established" signal
  * connection controllers need for the readiness handshake. Generators are lazy, so the
@@ -182,9 +183,12 @@ const UNARY_VALUE_SCHEMAS: { [K in keyof RpcMethodMap]: z.ZodType<Wire<ResponseV
   'llm.models': llmModelsValueSchema,
 }
 
-/** Default unary timeout (rpc-compare 2026-07-19: a hung host must not leave callers pending forever). */
+/** Default timeout for bounded unary calls (rpc-compare 2026-07-19: a hung host must not leave callers pending forever). */
 const DEFAULT_TIMEOUT_MS = 30_000
 
+/** Whether a unary call uses the transport health deadline or only caller/connection cancellation. */
+type UnaryTimeoutPolicy = 'default' | 'caller-signal-only'
+
 /** URL base for in-process handler injection (fake authority, opencode precedent). */
 const INTERNAL_BASE = 'http://dsh.internal'
 
@@ -202,7 +206,7 @@ export abstract class AbstractApiClient implements IApiClient {
   private flushScheduled = false
   private readonly envelopeListeners = new Set<(batch: readonly RpcMessage[]) => void>()
 
-  /** @param timeoutMs - unary timeout; streams never time out (long-lived by nature). */
+  /** @param timeoutMs - timeout for bounded unary calls; user-paced calls and streams do not use it. */
   constructor(protected readonly timeoutMs: number = DEFAULT_TIMEOUT_MS) {}
 
   /** Transport aspect: browser fetch, injected handler.fetch, IPC bridge, ... */
@@ -257,15 +261,15 @@ export abstract class AbstractApiClient implements IApiClient {
 
   /**
    * Shared POST leg of both C→S carriers (callUnary/respond): JSON body,
-   * timeout merged with the caller's optional external signal, non-2xx → transport throw.
+   * optional default timeout merged with the caller's external signal, non-2xx → transport throw.
    */
   private async postJson(
     path: string,
     body: ClientRequest | ClientResponse,
     signal: AbortSignal | undefined,
-    useDefaultTimeout = true,
+    timeoutPolicy: UnaryTimeoutPolicy = 'default',
   ): Promise<Response> {
-    const requestSignal = useDefaultTimeout
+    const requestSignal = timeoutPolicy === 'default'
       ? signal === undefined
         ? AbortSignal.timeout(this.timeoutMs)
         : AbortSignal.any([AbortSignal.timeout(this.timeoutMs), signal])
@@ -289,11 +293,11 @@ export abstract class AbstractApiClient implements IApiClient {
     method: K,
     payload: RequestPayload<K>,
     signal?: AbortSignal,
-    useDefaultTimeout = true,
+    timeoutPolicy: UnaryTimeoutPolicy = 'default',
   ): Promise<RpcResponse<ResponseValue<K>>> {
     const message: ClientRequest = { type: 'client-request', rpcId: this.mintRpcId(), method, payload }
     this.onEnvelope(message)
-    const response = await this.postJson(`/api/${method}`, message, signal, useDefaultTimeout)
+    const response = await this.postJson(`/api/${method}`, message, signal, timeoutPolicy)
     const full = serverResponseSchema.parse(await response.json())
     this.onEnvelope(full)
     if (full.rpcId !== message.rpcId) throw new Error(`rpcId mismatch for ${method}: sent ${message.rpcId}, got ${full.rpcId}`)
@@ -382,7 +386,9 @@ export abstract class AbstractApiClient implements IApiClient {
     describe: (payload, signal) => this.callUnary('host.describe', payload, signal),
     // A native system dialog is user-paced and may legitimately stay open
     // longer than the normal unary deadline. Caller/connection aborts remain.
-    pickDirectory: (payload, signal) => this.callUnary('host.pickDirectory', payload, signal, false),
+    pickDirectory: (payload, signal) => this.callUnary(
+      'host.pickDirectory', payload, signal, 'caller-signal-only',
+    ),
     listDirectory: (payload, signal) => this.callUnary('host.listDirectory', payload, signal),
     createDirectory: (payload, signal) => this.callUnary('host.createDirectory', payload, signal),
     openPath: (payload, signal) => this.callUnary('host.openPath', payload, signal),
@@ -398,7 +404,11 @@ export abstract class AbstractApiClient implements IApiClient {
 
   readonly commands: IApiClient['commands'] = {
     list: (payload, signal) => this.callUnary('command.list', payload, signal),
-    execute: (payload, signal) => this.callUnary('command.execute', payload, signal),
+    // Command handlers are user-driven operations and may legitimately exceed
+    // the transport health deadline. Caller/connection aborts remain.
+    execute: (payload, signal) => this.callUnary(
+      'command.execute', payload, signal, 'caller-signal-only',
+    ),
   }
 
   readonly skills: IApiClient['skills'] = {

+ 62 - 0
packages/host/apiproxy/tests/fetch-carrier.spec.ts

@@ -349,6 +349,68 @@ describe('unary round trip (handler ⇄ client, no network)', () => {
     expect(skills.result).toEqual({ ok: true, value: { skills: [{ name: 'commit-helper', description: 'Git commits' }] } })
   })
 
+  it('lets command.execute finish after the 30-second default unary deadline', async () => {
+    vi.useFakeTimers()
+    const timeoutSpy = vi.spyOn(AbortSignal, 'timeout').mockImplementation((milliseconds) => {
+      const controller = new AbortController()
+      setTimeout(() => {
+        controller.abort(new DOMException('The operation was aborted due to timeout', 'TimeoutError'))
+      }, milliseconds)
+      return controller.signal
+    })
+    try {
+      const api = fakeApi()
+      api.commands.execute = async (request) => {
+        await new Promise(resolve => setTimeout(resolve, 30_001))
+        return {
+          rpcId: request.rpcId,
+          result: { ok: true, value: { matched: true, commandId: CommandId('cmd-slow') } },
+        }
+      }
+      const execution = client(api).commands.execute({ sessionId: 's' as never, line: '/slow' })
+      const assertion = expect(execution).resolves.toMatchObject({
+        result: { ok: true, value: { matched: true, commandId: 'cmd-slow' } },
+      })
+
+      await Promise.all([
+        vi.advanceTimersByTimeAsync(30_001),
+        assertion,
+      ])
+      expect(timeoutSpy).not.toHaveBeenCalled()
+    } finally {
+      timeoutSpy.mockRestore()
+      vi.useRealTimers()
+    }
+  })
+
+  it('keeps caller and connection aborts on command.execute', async () => {
+    const api = fakeApi()
+    const started = Promise.withResolvers<AbortSignal>()
+    api.commands.execute = async (request, signal) => {
+      started.resolve(signal)
+      if (!signal.aborted) {
+        await new Promise<void>((resolve) => {
+          signal.addEventListener('abort', () => { resolve() }, { once: true })
+        })
+      }
+      return {
+        rpcId: request.rpcId,
+        result: { ok: false, error: { code: 'cancelled', message: 'aborted', details: {} } },
+      }
+    }
+    const controller = new AbortController()
+    const execution = client(api).commands.execute(
+      { sessionId: 's' as never, line: '/hang' },
+      controller.signal,
+    )
+    const handlerSignal = await started.promise
+
+    controller.abort(new Error('connection closed'))
+
+    await expect(execution).rejects.toThrow('connection closed')
+    expect(handlerSignal.aborted).toBe(true)
+  })
+
   it('propagates the carrier Request signal into command.execute', async () => {
     const handler = toFetchHandler(fakeApi())
     const controller = new AbortController()

Beberapa file tidak ditampilkan karena terlalu banyak file yang berubah dalam diff ini