Expose audited hardcoded tunables as plugin config

The audit swept every packages/*/* plugin for the new AGENTS.md
convention (no hardcoded tunables in plugins) and exposes each finding
as a defaulted, validated Config field. Defaults are the previously
hardcoded values throughout, so no deployment or golden changes.

- tool-fs (had NO Config): readLimit, readMaxLineLength, readMaxBytes,
  readStreamMinSize. The caps thread through ReadToolCaps/ReadWindow —
  read-render already documented that the consumer applies the caps, so
  they become explicit per-request fields.
- tool-web: searchMaxResults (WEB_SEARCH_MAX_RESULTS stays as the
  schemastery default). Also fixes the stale GREP_LIMIT references in
  search.ts and the web-capability-seam RFC (no such constant exists).
- bash-local: graceMs (SIGTERM->SIGKILL escalation grace). The
  RunInternals.graceMs test seam is gone: graceMs is now a required
  SpawnSpec field filled from config, so tests exercise the real
  config path and the defaults live in exactly one place.
- subagent-acp: disposeEofGraceMs / disposeGraceMs. The AcpRunSpec
  fields become required for the same one-defaulting-layer reason.
- session-persistence-sqlite: journalMode ('wal' default; the
  rollback-journal modes serve filesystems where WAL's shared-memory
  files do not work, e.g. network mounts).
- hooks-claude + hooks-codex: stderrSummaryMaxChars for the persisted
  hook/result stderr summary. The duplicated summarize() helpers merge
  into hook-protocol's summarizeStderr(stderr, maxChars), beside the
  HookResultRecord field it feeds, with the bound parameterized the
  same way runHook's defaultTimeoutMs already is.
- compact-basic: charsPerToken for the token estimator (default 4, the
  English-text heuristic; CJK-heavy deployments need ~1-2 or compaction
  fires far too late). Also corrects the BasicCompactService class doc,
  which claimed defaults the required-field config never had.
- fs-local: deletes the dead STREAM_MIN_SIZE constant and the dead
  FsIoInternals.streamMinSize seam — the read-routing bound lives in
  the consumer (tool-fs), where it is now config. This is item 1 of
  the proposed prune-write-only-fs-surface RFC, annotated accordingly.

Every new field gets range validation (following the existing
assertPositiveFinite pattern), a README row, and tests covering the
configured behavior, the schema default, and load-time rejection.
This commit is contained in:
Tianyi Cui
2026-07-04 17:37:23 +08:00
parent bee8132a7f
commit 774d460889
45 changed files with 592 additions and 171 deletions

View File

@@ -10,7 +10,7 @@ Search and fetch share no request schema and no business logic, but they are del
## Search request and result ## Search request and result
The model-facing tool argument is just a `query`; `maxResults` is a consumer-owned bound (`dsh-tool-web`'s `WEB_SEARCH_MAX_RESULTS`, default `8`) passed through the seam and enforced on the way back — if a provider over-returns, the seam truncates `sources[]` and sets `truncated`. The model-facing tool argument is just a `query`; `maxResults` is a consumer-owned bound (`dsh-tool-web`'s `searchMaxResults` config, default `8`) passed through the seam and enforced on the way back — if a provider over-returns, the seam truncates `sources[]` and sets `truncated`.
```ts type-equiv ```ts type-equiv
interface WebSearchRequest { interface WebSearchRequest {

View File

@@ -175,7 +175,7 @@ The first `web_search` model-facing tool should be small. The only model-facing
- `query`: required string. - `query`: required string.
`max_results` is NOT exposed to the model in the first version. It is a `dsh-tool-web`-layer decision: the tool sets the result bound — a default of `8` (aligning with OpenCode's Exa default), as an exported constant mirroring `dsh-tool-fs`'s `READ_LIMIT` / `GREP_LIMIT` — and passes it to the seam as `maxResults` on the `WebSearchRequest`. Keeping it off the model schema means the model just asks a question and the product controls how much context comes back; the field can be promoted to a model-facing argument later without breaking the seam. `max_results` is NOT exposed to the model in the first version. It is a `dsh-tool-web`-layer decision: the tool sets the result bound — the `searchMaxResults` plugin config, default `8` (aligning with OpenCode's Exa default), mirroring `dsh-tool-fs`'s `readLimit` — and passes it to the seam as `maxResults` on the `WebSearchRequest`. Keeping it off the model schema means the model just asks a question and the product controls how much context comes back; the field can be promoted to a model-facing argument later without breaking the seam.
`maxResults` flows tool → seam → provider, and the bound is enforced on the way back: `maxResults` flows tool → seam → provider, and the bound is enforced on the way back:

View File

@@ -6,7 +6,7 @@ Status: proposed
The [fs seam split](../../implemented/simplification/2026-06-26-fsspec-style-fs-seam.md) moved read routing and policy out of the backend into `dsh-tool-fs` and `dsh-fs-policy`. Four pieces of surface kept the pre-split shape — populated on every call, read by nobody: The [fs seam split](../../implemented/simplification/2026-06-26-fsspec-style-fs-seam.md) moved read routing and policy out of the backend into `dsh-tool-fs` and `dsh-fs-policy`. Four pieces of surface kept the pre-split shape — populated on every call, read by nobody:
1. **`STREAM_MIN_SIZE` + `FsIoInternals.streamMinSize` in `dsh-fs-local`** (`packages/fs/fs-local/src/fsio.ts`, re-exported from `packages/fs/fs-local/src/index.ts`): zero readers anywhere, including fs-local's own source and tests. The backend has no read routing — `readWholeText`/`streamWholeText` are separate primitives the caller chooses between — and the real routing constant lives in the consumer (`packages/fs/tool-fs/src/read.ts`, compared against `info.size`). Two mirrors of the 10 MiB fact; the backend's is dead, and the knob's JSDoc claims a "read routing" override that does not exist. 1. **`STREAM_MIN_SIZE` + `FsIoInternals.streamMinSize` in `dsh-fs-local`** — *already removed by the no-hardcoded-tunables audit (the routing bound became `dsh-tool-fs`'s `readStreamMinSize` config); listed here for the record of the full prune, no work remains.* Originally (`packages/fs/fs-local/src/fsio.ts`, re-exported from `packages/fs/fs-local/src/index.ts`): zero readers anywhere, including fs-local's own source and tests. The backend has no read routing — `readWholeText`/`streamWholeText` are separate primitives the caller chooses between — and the real routing constant lives in the consumer (`packages/fs/tool-fs/src/read.ts`, compared against `info.size`). Two mirrors of the 10 MiB fact; the backend's is dead, and the knob's JSDoc claims a "read routing" override that does not exist.
2. **`FsTarget.inputPath`** (`packages/fs/fs/src/types.ts`): every backend and every test fake must fabricate a "diagnostics only" value with zero production readers — the policy plugin and every error message use `targetKey`/`displayPath`. The `listDir` producer exposes the semantic wobble: directory children get the bare entry name, which was nobody's "input". 2. **`FsTarget.inputPath`** (`packages/fs/fs/src/types.ts`): every backend and every test fake must fabricate a "diagnostics only" value with zero production readers — the policy plugin and every error message use `targetKey`/`displayPath`. The `listDir` producer exposes the semantic wobble: directory children get the bare entry name, which was nobody's "input".
3. **`FsEditOutcome.replacements` + `.replaceAll`** (`packages/fs/fs/src/types.ts`): `replacements` has zero production readers (the single-match policy itself stays — it is enforced by the `FS_AMBIGUOUS_EDIT`/`FS_EDIT_NOT_FOUND` throws inside the backend, whose error message keeps the internal count); `replaceAll` is read only by `formatEditOutput` in `packages/fs/tool-fs/src/edit.ts` — as an echo of the `replace_all` argument the tool already holds. Shrunk, `FsEditOutcome` becomes `{ version, before, after }`, parallel to `FsWriteOutcome`'s genuinely backend-discovered fields. 3. **`FsEditOutcome.replacements` + `.replaceAll`** (`packages/fs/fs/src/types.ts`): `replacements` has zero production readers (the single-match policy itself stays — it is enforced by the `FS_AMBIGUOUS_EDIT`/`FS_EDIT_NOT_FOUND` throws inside the backend, whose error message keeps the internal count); `replaceAll` is read only by `formatEditOutput` in `packages/fs/tool-fs/src/edit.ts` — as an echo of the `replace_all` argument the tool already holds. Shrunk, `FsEditOutcome` becomes `{ version, before, after }`, parallel to `FsWriteOutcome`'s genuinely backend-discovered fields.
4. **`FileReadOutcome.limit` + `.version`** (`packages/fs/tool-fs/src/read-render.ts`): populated by the read tool, but `formatReadOutput` renders `offset`/`lines`/`totalLines`/`truncatedByBytes` only, and the `fs/observed` emit uses `info.version` directly rather than the outcome copy. 4. **`FileReadOutcome.limit` + `.version`** (`packages/fs/tool-fs/src/read-render.ts`): populated by the read tool, but `formatReadOutput` renders `offset`/`lines`/`totalLines`/`truncatedByBytes` only, and the `fs/observed` emit uses `info.version` directly rather than the outcome copy.

View File

@@ -12,6 +12,7 @@ Local-subprocess implementation of the `@deepseek-ai/dsh-bash` executor seam: `L
timeoutMs: 120000 # default foreground timeout timeoutMs: 120000 # default foreground timeout
maxTimeoutMs: 600000 # cap for per-call overrides maxTimeoutMs: 600000 # cap for per-call overrides
maxOutputBytes: 64000 # per-stream in-memory cap; overflow spills to disk maxOutputBytes: 64000 # per-stream in-memory cap; overflow spills to disk
graceMs: 3000 # SIGTERM→SIGKILL escalation grace on kills
``` ```
## Behavior (and where it came from) ## Behavior (and where it came from)
@@ -19,7 +20,7 @@ Local-subprocess implementation of the `@deepseek-ai/dsh-bash` executor seam: `L
Design surveyed against the bash tools of Claude Code, OpenCode, Codex, and pi; the notable choices: Design surveyed against the bash tools of Claude Code, OpenCode, Codex, and pi; the notable choices:
- **Spawn per call, no shell state** — every call is a fresh non-login `bash -c` (deterministic; no rc files). All four surveyed tools spawn per call. `XXX(stateful-shell)` in `src/run.ts` records the two proven stateful designs (Claude Code's cwd-only persistence; Codex's PTY exec sessions) for when real workflows demand them. - **Spawn per call, no shell state** — every call is a fresh non-login `bash -c` (deterministic; no rc files). All four surveyed tools spawn per call. `XXX(stateful-shell)` in `src/run.ts` records the two proven stateful designs (Claude Code's cwd-only persistence; Codex's PTY exec sessions) for when real workflows demand them.
- **Process-group kills with escalation** — children are spawned `detached` (own process group); kills send SIGTERM to the group, then SIGKILL after a 3s grace (OpenCode's escalation; pipelines and subshells die with the parent). ESRCH is tolerated; daemons that re-parent away from the group can still survive — same caveat as the surveyed tools. - **Process-group kills with escalation** — children are spawned `detached` (own process group); kills send SIGTERM to the group, then SIGKILL after the `graceMs` grace (default 3s — OpenCode's escalation; pipelines and subshells die with the parent). ESRCH is tolerated; daemons that re-parent away from the group can still survive — same caveat as the surveyed tools.
- **Tail-keep truncation + spill files** — output beyond `maxOutputBytes` keeps the in-memory TAIL (errors/results cluster at the end — pi/OpenCode rationale) while the FULL stream is appended to a temp file whose path is reported when available. If the final spill close reports a delayed writeback failure, the executor still returns the tail but withholds the path rather than advertising a possibly incomplete file. - **Tail-keep truncation + spill files** — output beyond `maxOutputBytes` keeps the in-memory TAIL (errors/results cluster at the end — pi/OpenCode rationale) while the FULL stream is appended to a temp file whose path is reported when available. If the final spill close reports a delayed writeback failure, the executor still returns the tail but withholds the path rather than advertising a possibly incomplete file.
- **Model-friendly env + credential scrub** — `process.env` minus credential-shaped vars (`*KEY*`/`*SECRET*`/`*TOKEN*`), then `NO_COLOR=1 TERM=dumb PAGER=cat GIT_PAGER=cat` (Codex's hardcoded set) so pagers and ANSI color don't garble results. This scrub is the security control that keeps the harness's *ambient* credentials out of a spawned command. A spec's `env` is merged LAST (after the scrub), so a caller's explicit entry — a value it already holds — wins even on a credential-shaped name. The spec's `stdin`, when supplied, is written to the child and closed; with none supplied, fd 0 is `/dev/null` — the exact pre-seam default, so a command that probes stdin's file type is unaffected. Both `env`/`stdin` are set by in-process plugins (the hooks bridges); the model-facing tool doesn't expose them. See [the bash-stdin-env RFC](../../../docs/rfc/implemented/architecture/2026-06-30-bash-stdin-env-trusted-plugin-surface.md). - **Model-friendly env + credential scrub** — `process.env` minus credential-shaped vars (`*KEY*`/`*SECRET*`/`*TOKEN*`), then `NO_COLOR=1 TERM=dumb PAGER=cat GIT_PAGER=cat` (Codex's hardcoded set) so pagers and ANSI color don't garble results. This scrub is the security control that keeps the harness's *ambient* credentials out of a spawned command. A spec's `env` is merged LAST (after the scrub), so a caller's explicit entry — a value it already holds — wins even on a credential-shaped name. The spec's `stdin`, when supplied, is written to the child and closed; with none supplied, fd 0 is `/dev/null` — the exact pre-seam default, so a command that probes stdin's file type is unaffected. Both `env`/`stdin` are set by in-process plugins (the hooks bridges); the model-facing tool doesn't expose them. See [the bash-stdin-env RFC](../../../docs/rfc/implemented/architecture/2026-06-30-bash-stdin-env-trusted-plugin-surface.md).
- **Background tasks** — `start()` returns immediately, no timeout applies (Claude Code detaches timeouts when backgrounding), `readOutput()` is incremental with whole-stream byte offsets, and disposal kills everything. The spec's opaque `owner` token is stored on the tracked task and returned by `ownerOf(id)` — the executor never interprets it (the consumer's access policy does), and because it lives with the task here it survives a `tool-bash` HMR reload. - **Background tasks** — `start()` returns immediately, no timeout applies (Claude Code detaches timeouts when backgrounding), `readOutput()` is incremental with whole-stream byte offsets, and disposal kills everything. The spec's opaque `owner` token is stored on the tracked task and returned by `ownerOf(id)` — the executor never interprets it (the consumer's access policy does), and because it lives with the task here it survives a `tool-bash` HMR reload.

View File

@@ -17,7 +17,7 @@ import { Context } from 'cordis'
import z from 'schemastery' import z from 'schemastery'
import { BashExecutor, BashTaskId } from '@deepseek-ai/dsh-bash' import { BashExecutor, BashTaskId } from '@deepseek-ai/dsh-bash'
import type { BashExecRequest, BashExecSpec, BashRunResult, BashTask, BashTaskRead, OwnerToken } from '@deepseek-ai/dsh-bash' import type { BashExecRequest, BashExecSpec, BashRunResult, BashTask, BashTaskRead, OwnerToken } from '@deepseek-ai/dsh-bash'
import { runBash } from './run.ts' import { DEFAULT_GRACE_MS, runBash } from './run.ts'
import type { RunInternals, RunningBash } from './run.ts' import type { RunInternals, RunningBash } from './run.ts'
export { DEFAULT_GRACE_MS, ENV_OVERRIDES, killGroup, OutputCollector, runBash } from './run.ts' export { DEFAULT_GRACE_MS, ENV_OVERRIDES, killGroup, OutputCollector, runBash } from './run.ts'
@@ -33,6 +33,8 @@ export interface Config {
maxTimeoutMs?: number maxTimeoutMs?: number
/** Per-stream in-memory output cap; overflow spills to a temp file. */ /** Per-stream in-memory output cap; overflow spills to a temp file. */
maxOutputBytes?: number maxOutputBytes?: number
/** Grace period between the SIGTERM and the SIGKILL escalation on a kill. */
graceMs?: number
} }
/** The shape after schemastery applied the defaults (cwd has none). */ /** The shape after schemastery applied the defaults (cwd has none). */
@@ -57,7 +59,7 @@ interface TrackedTask extends BashTask {
* Local-subprocess bash executor. Defaults follow the agent-tool survey * Local-subprocess bash executor. Defaults follow the agent-tool survey
* consensus: 120s default / 600s max timeout (Claude Code, OpenCode), 64KB * consensus: 120s default / 600s max timeout (Claude Code, OpenCode), 64KB
* in-memory output with full-stream spill files (pi, OpenCode), * in-memory output with full-stream spill files (pi, OpenCode),
* process-group SIGTERM→SIGKILL kills (OpenCode). * process-group SIGTERM→SIGKILL kills with a 3s grace (OpenCode).
*/ */
export class LocalBashExecutor extends BashExecutor { export class LocalBashExecutor extends BashExecutor {
static Config: z<Config> = z.object({ static Config: z<Config> = z.object({
@@ -65,11 +67,12 @@ export class LocalBashExecutor extends BashExecutor {
timeoutMs: z.number().default(120_000), timeoutMs: z.number().default(120_000),
maxTimeoutMs: z.number().default(600_000), maxTimeoutMs: z.number().default(600_000),
maxOutputBytes: z.number().default(64_000), maxOutputBytes: z.number().default(64_000),
graceMs: z.number().default(DEFAULT_GRACE_MS),
}) })
private tasks = new Map<BashTaskId, TrackedTask>() private tasks = new Map<BashTaskId, TrackedTask>()
private nextTaskId = 1 private nextTaskId = 1
/** Test seam: timer/spill knobs forwarded to runBash. */ /** Test seam: spill knobs forwarded to runBash. */
internals: RunInternals = {} internals: RunInternals = {}
/** Validated config (schemastery applied the defaults before construction). */ /** Validated config (schemastery applied the defaults before construction). */
@@ -83,6 +86,7 @@ export class LocalBashExecutor extends BashExecutor {
assertPositiveFinite('timeoutMs', this.config.timeoutMs) assertPositiveFinite('timeoutMs', this.config.timeoutMs)
assertPositiveFinite('maxTimeoutMs', this.config.maxTimeoutMs) assertPositiveFinite('maxTimeoutMs', this.config.maxTimeoutMs)
assertPositiveFinite('maxOutputBytes', this.config.maxOutputBytes) assertPositiveFinite('maxOutputBytes', this.config.maxOutputBytes)
assertPositiveFinite('graceMs', this.config.graceMs)
ctx.effect(() => async () => { ctx.effect(() => async () => {
// Kill every live process group and WAIT for the processes to close so // Kill every live process group and WAIT for the processes to close so
// nothing outlives the fiber (HMR safety) — a TERM-trapping child is // nothing outlives the fiber (HMR safety) — a TERM-trapping child is
@@ -132,6 +136,7 @@ export class LocalBashExecutor extends BashExecutor {
cwd: spec.workdir, cwd: spec.workdir,
timeoutMs: spec.timeoutMs, timeoutMs: spec.timeoutMs,
maxOutputBytes: this.config.maxOutputBytes, maxOutputBytes: this.config.maxOutputBytes,
graceMs: this.config.graceMs,
signal: spec.signal, signal: spec.signal,
stdin: spec.stdin, stdin: spec.stdin,
env: spec.env, env: spec.env,
@@ -150,6 +155,7 @@ export class LocalBashExecutor extends BashExecutor {
cwd: spec.workdir, cwd: spec.workdir,
timeoutMs: 0, timeoutMs: 0,
maxOutputBytes: this.config.maxOutputBytes, maxOutputBytes: this.config.maxOutputBytes,
graceMs: this.config.graceMs,
signal: spec.signal, signal: spec.signal,
stdin: spec.stdin, stdin: spec.stdin,
env: spec.env, env: spec.env,

View File

@@ -73,6 +73,8 @@ export interface SpawnSpec {
timeoutMs: number timeoutMs: number
/** Per-stream in-memory cap; overflow spills to disk (tail kept in memory). */ /** Per-stream in-memory cap; overflow spills to disk (tail kept in memory). */
maxOutputBytes: number maxOutputBytes: number
/** Grace period between the SIGTERM and the SIGKILL escalation on a kill. */
graceMs: number
/** Abort signal — kills the process group when fired. */ /** Abort signal — kills the process group when fired. */
signal?: AbortSignal | undefined signal?: AbortSignal | undefined
/** /**
@@ -100,15 +102,13 @@ export interface SpawnOutcome {
stderr: CollectedOutput stderr: CollectedOutput
} }
/** Injectable knobs so tests can exercise escalation/spill without long waits. */ /** Injectable knobs so tests can exercise spill behavior without the OS tmpdir. */
export interface RunInternals { export interface RunInternals {
/** Grace period between SIGTERM and SIGKILL on the process group. */
graceMs?: number
/** Directory for spill files (defaults to the OS temp dir). */ /** Directory for spill files (defaults to the OS temp dir). */
spillDir?: string spillDir?: string
} }
/** Default SIGTERM→SIGKILL grace period (matches OpenCode's 3s). */ /** Default SIGTERM→SIGKILL grace period (the `graceMs` config; matches OpenCode's 3s). */
export const DEFAULT_GRACE_MS = 3_000 export const DEFAULT_GRACE_MS = 3_000
let spillCounter = 0 let spillCounter = 0
@@ -292,7 +292,6 @@ export interface RunningBash {
* no inherited shell state); revisit when real workflows demand it. * no inherited shell state); revisit when real workflows demand it.
*/ */
export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningBash { export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningBash {
const graceMs = internals.graceMs ?? DEFAULT_GRACE_MS
const spillDir = internals.spillDir ?? privateSpillDir() const spillDir = internals.spillDir ?? privateSpillDir()
if (spec.signal?.aborted) { if (spec.signal?.aborted) {
@@ -331,7 +330,7 @@ export function runBash(spec: SpawnSpec, internals: RunInternals = {}): RunningB
const kill = (): void => { const kill = (): void => {
if (graceTimer !== undefined) return // escalation already in flight if (graceTimer !== undefined) return // escalation already in flight
killGroup(pid, 'SIGTERM') killGroup(pid, 'SIGTERM')
graceTimer = setTimeout(() => { killGroup(pid, 'SIGKILL') }, graceMs) graceTimer = setTimeout(() => { killGroup(pid, 'SIGKILL') }, spec.graceMs)
} }
if (spec.timeoutMs > 0) { if (spec.timeoutMs > 0) {

View File

@@ -11,9 +11,10 @@ const spillDir = mkdtempSync(join(tmpdir(), 'dsh-bash-exec-spec-'))
async function setup(config: ConstructorParameters<typeof LocalBashExecutor>[1] = {}) { async function setup(config: ConstructorParameters<typeof LocalBashExecutor>[1] = {}) {
const ctx = new Context() const ctx = new Context()
await ctx.plugin(LocalBashExecutor, config) // A short kill grace via the REAL config path, so escalation tests stay fast.
await ctx.plugin(LocalBashExecutor, { graceMs: 200, ...config })
const bash = ctx.bash as LocalBashExecutor const bash = ctx.bash as LocalBashExecutor
bash.internals = { spillDir, graceMs: 200 } bash.internals = { spillDir }
return { ctx, bash } return { ctx, bash }
} }
@@ -80,12 +81,22 @@ describe('LocalBashExecutor.run', () => {
await expect(setup({ timeoutMs: Number.NaN })).rejects.toThrow(/timeoutMs/) await expect(setup({ timeoutMs: Number.NaN })).rejects.toThrow(/timeoutMs/)
await expect(setup({ maxTimeoutMs: 0 })).rejects.toThrow(/maxTimeoutMs/) await expect(setup({ maxTimeoutMs: 0 })).rejects.toThrow(/maxTimeoutMs/)
await expect(setup({ maxOutputBytes: -1 })).rejects.toThrow(/maxOutputBytes/) await expect(setup({ maxOutputBytes: -1 })).rejects.toThrow(/maxOutputBytes/)
await expect(setup({ graceMs: 0 })).rejects.toThrow(/graceMs/)
const { bash } = await setup() const { bash } = await setup()
expect(() => bash.resolve({ command: 'true', timeoutMs: Number.NaN })).toThrow(/request\.timeoutMs/) expect(() => bash.resolve({ command: 'true', timeoutMs: Number.NaN })).toThrow(/request\.timeoutMs/)
expect(() => bash.resolve({ command: 'true', timeoutMs: -1 })).toThrow(/request\.timeoutMs/) expect(() => bash.resolve({ command: 'true', timeoutMs: -1 })).toThrow(/request\.timeoutMs/)
}) })
it('kill escalation uses the configured graceMs (a TERM-trapping task dies by SIGKILL)', async () => {
const { bash } = await setup() // setup pins graceMs: 200 via config
const task = bash.start(bash.resolve({ command: 'trap \'\' TERM; sleep 60' }))
await new Promise(resolve => setTimeout(resolve, 100))
bash.kill(task.id)
await task.done
expect(task.signal).toBe('SIGKILL')
})
it('per-call timeout takes precedence under the cap and kills on expiry', async () => { it('per-call timeout takes precedence under the cap and kills on expiry', async () => {
const { bash } = await setup({ timeoutMs: 60_000 }) const { bash } = await setup({ timeoutMs: 60_000 })
const result = await bash.run(bash.resolve({ command: 'sleep 60', timeoutMs: 100 })) const result = await bash.run(bash.resolve({ command: 'sleep 60', timeoutMs: 100 }))
@@ -271,9 +282,9 @@ describe('LocalBashExecutor background tasks', () => {
it('disposing with already-finished tasks only kills the running ones', async () => { it('disposing with already-finished tasks only kills the running ones', async () => {
const ctx = new Context() const ctx = new Context()
const fiber = await ctx.plugin(LocalBashExecutor, {}) const fiber = await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
const bash = ctx.bash as LocalBashExecutor const bash = ctx.bash as LocalBashExecutor
bash.internals = { spillDir, graceMs: 200 } bash.internals = { spillDir }
const finished = bash.start(bash.resolve({ command: 'true' })) const finished = bash.start(bash.resolve({ command: 'true' }))
await finished.done await finished.done
@@ -288,9 +299,9 @@ describe('LocalBashExecutor background tasks', () => {
it('disposing the executor fiber kills running tasks (no orphans)', async () => { it('disposing the executor fiber kills running tasks (no orphans)', async () => {
const ctx = new Context() const ctx = new Context()
const fiber = await ctx.plugin(LocalBashExecutor, {}) const fiber = await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
const bash = ctx.bash as LocalBashExecutor const bash = ctx.bash as LocalBashExecutor
bash.internals = { spillDir, graceMs: 200 } bash.internals = { spillDir }
const listener = vi.fn() const listener = vi.fn()
bash.onTaskDone(listener) bash.onTaskDone(listener)
@@ -337,9 +348,9 @@ describe('review fixes: lifecycle hardening', () => {
it('dispose AWAITS a TERM-trapping process (SIGKILL escalation included)', async () => { it('dispose AWAITS a TERM-trapping process (SIGKILL escalation included)', async () => {
const ctx = new Context() const ctx = new Context()
const fiber = await ctx.plugin(LocalBashExecutor, {}) const fiber = await ctx.plugin(LocalBashExecutor, { graceMs: 200 })
const bash = ctx.bash as LocalBashExecutor const bash = ctx.bash as LocalBashExecutor
bash.internals = { spillDir, graceMs: 200 } bash.internals = { spillDir }
const task = bash.start(bash.resolve({ command: 'trap \'\' TERM; sleep 60' })) const task = bash.start(bash.resolve({ command: 'trap \'\' TERM; sleep 60' }))
await new Promise(resolve => setTimeout(resolve, 100)) await new Promise(resolve => setTimeout(resolve, 100))

View File

@@ -28,6 +28,7 @@ function spec(command: string, overrides: Partial<Parameters<typeof runBash>[0]>
cwd: process.cwd(), cwd: process.cwd(),
timeoutMs: 0, timeoutMs: 0,
maxOutputBytes: 64_000, maxOutputBytes: 64_000,
graceMs: 3_000,
...overrides, ...overrides,
} }
} }
@@ -106,7 +107,7 @@ describe('runBash', () => {
}) })
it('escalates to SIGKILL when SIGTERM is trapped', async () => { it('escalates to SIGKILL when SIGTERM is trapped', async () => {
const running = runBash(spec('trap \'\' TERM; echo ready; sleep 60'), { graceMs: 200 }) const running = runBash(spec('trap \'\' TERM; echo ready; sleep 60', { graceMs: 200 }))
await waitForStdout(running, 'ready\n') await waitForStdout(running, 'ready\n')
running.kill() running.kill()
const result = await running.done const result = await running.done

View File

@@ -21,8 +21,8 @@ async function setup() {
await ctx.plugin(SystemPrompt) await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry) await ctx.plugin(ToolRegistry)
await ctx.plugin(AgentRegistry) await ctx.plugin(AgentRegistry)
await ctx.plugin(LocalBashExecutor, { timeoutMs: 10_000 }) await ctx.plugin(LocalBashExecutor, { timeoutMs: 10_000, graceMs: 200 })
;(ctx.bash as LocalBashExecutor).internals = { spillDir, graceMs: 200 } ;(ctx.bash as LocalBashExecutor).internals = { spillDir }
await ctx.plugin(ToolBash) await ctx.plugin(ToolBash)
return ctx return ctx
} }
@@ -184,8 +184,8 @@ describe('bash tool', () => {
const ctx = new Context() const ctx = new Context()
await ctx.plugin(SystemPrompt) await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry) await ctx.plugin(ToolRegistry)
await ctx.plugin(LocalBashExecutor, { maxOutputBytes: 100 }) await ctx.plugin(LocalBashExecutor, { maxOutputBytes: 100, graceMs: 200 })
;(ctx.bash as LocalBashExecutor).internals = { spillDir, graceMs: 200 } ;(ctx.bash as LocalBashExecutor).internals = { spillDir }
await ctx.plugin(ToolBash) await ctx.plugin(ToolBash)
const result = await call(ctx, 'bash', { command: 'for i in $(seq 1 100); do printf "line-%04d\\n" $i; done', description: 'test command' }) const result = await call(ctx, 'bash', { command: 'for i in $(seq 1 100); do printf "line-%04d\\n" $i; done', description: 'test command' })
expect(text(result)).toContain('[output truncated; full output: ') expect(text(result)).toContain('[output truncated; full output: ')
@@ -316,8 +316,8 @@ describe('background tools', () => {
const ctx = new Context() const ctx = new Context()
await ctx.plugin(SystemPrompt) await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry) await ctx.plugin(ToolRegistry)
await ctx.plugin(LocalBashExecutor, { maxOutputBytes: 100 }) await ctx.plugin(LocalBashExecutor, { maxOutputBytes: 100, graceMs: 200 })
;(ctx.bash as LocalBashExecutor).internals = { spillDir, graceMs: 200 } ;(ctx.bash as LocalBashExecutor).internals = { spillDir }
await ctx.plugin(ToolBash) await ctx.plugin(ToolBash)
const started = await call(ctx, 'bash', { command: 'for i in $(seq 1 200); do printf "line-%04d\\n" $i; done', description: 'test command', run_in_background: true }) const started = await call(ctx, 'bash', { command: 'for i in $(seq 1 200); do printf "line-%04d\\n" $i; done', description: 'test command', run_in_background: true })
@@ -568,8 +568,8 @@ describe('background task ownership (cross-session isolation)', () => {
const ctx = new Context() const ctx = new Context()
await ctx.plugin(SystemPrompt) await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry) await ctx.plugin(ToolRegistry)
await ctx.plugin(LocalBashExecutor, { timeoutMs: 10_000 }) await ctx.plugin(LocalBashExecutor, { timeoutMs: 10_000, graceMs: 200 })
;(ctx.bash as LocalBashExecutor).internals = { spillDir, graceMs: 200 } ;(ctx.bash as LocalBashExecutor).internals = { spillDir }
const fiber = await ctx.plugin(ToolBash) const fiber = await ctx.plugin(ToolBash)
const a = fakeAgent('sess-a') const a = fakeAgent('sess-a')

View File

@@ -32,6 +32,7 @@ Every knob is **required** except `auto` — there is no concrete data yet to ju
| `maxTokens` | yes | Provider generation cap for the summarization call; may include reasoning tokens. | | `maxTokens` | yes | Provider generation cap for the summarization call; may include reasoning tokens. |
| `compactionRetries` | yes | Extra compaction attempts after the first if the compacted surface remains over threshold. | | `compactionRetries` | yes | Extra compaction attempts after the first if the compacted surface remains over threshold. |
| `auto` | no (default `true`) | Register the `agent/pre-step` auto-compaction listener. Set `false` for manual-only. | | `auto` | no (default `true`) | Register the `agent/pre-step` auto-compaction listener. Set `false` for manual-only. |
| `charsPerToken` | no (default `4`) | Token-estimator text density (estimated tokens = chars / `charsPerToken`; may be fractional). The default suits English text; CJK-heavy deployments should set ~1-2 or the estimate undershoots several-fold and compaction fires too late. |
## Usage ## Usage

View File

@@ -148,9 +148,11 @@ function finishError(finish: FinishReason): Error | undefined {
} }
/** /**
* Basic, dependency-light compaction backend. Defaults target a 128K context * Basic, dependency-light compaction backend: estimates the surface's token
* window, compacting at 80% utilization and retaining ~20K tokens of recent * footprint, summarizes the stale prefix through the model, and shadows it
* context. * behind a durable checkpoint. Every threshold/budget knob is required config
* ({@link BasicCompactConfig}); the estimator's text density is the
* `charsPerToken` knob.
*/ */
export class BasicCompactService extends CompactService { export class BasicCompactService extends CompactService {
static inject = ['llm'] static inject = ['llm']
@@ -207,24 +209,27 @@ export class BasicCompactService extends CompactService {
// ---- Token estimation (overridable hooks) ---- // ---- Token estimation (overridable hooks) ----
// TODO: char/4 is a coarse heuristic. Replace with an exact count — a real // TODO: chars/charsPerToken is a coarse heuristic. Replace with an exact
// tokenizer, or the provider's post-response `usage` (input tokens) fed back // count — a real tokenizer, or the provider's post-response `usage` (input
// as a correction — so threshold decisions match the model's actual budget. // tokens) fed back as a correction — so threshold decisions match the
// model's actual budget.
/** /**
* Estimate the token count of content blocks — char/4 with per-block * Estimate the token count of content blocks — chars divided by the
* overhead. Override in a subclass to plug in a real tokenizer. * `charsPerToken` config, with per-block overhead. Override in a subclass to
* plug in a real tokenizer.
*/ */
estimateContentTokens(blocks: readonly ContentBlock[]): number { estimateContentTokens(blocks: readonly ContentBlock[]): number {
const { charsPerToken } = this.config
let tokens = 0 let tokens = 0
for (const block of blocks) { for (const block of blocks) {
switch (block.type) { switch (block.type) {
case 'text': case 'text':
case 'reasoning': case 'reasoning':
tokens += Math.ceil(block.text.length / 4) + BLOCK_OVERHEAD tokens += Math.ceil(block.text.length / charsPerToken) + BLOCK_OVERHEAD
break break
case 'tool-call': case 'tool-call':
tokens += Math.ceil(block.name.length / 4) tokens += Math.ceil(block.name.length / charsPerToken)
+ Math.ceil(block.arguments.length / 4) + Math.ceil(block.arguments.length / charsPerToken)
+ BLOCK_OVERHEAD + BLOCK_OVERHEAD
break break
case 'tool-result': case 'tool-result':
@@ -236,7 +241,7 @@ export class BasicCompactService extends CompactService {
default: default:
// Unknown block types (merge-extensible ContentBlockMap): // Unknown block types (merge-extensible ContentBlockMap):
// estimate conservatively via JSON stringify. // estimate conservatively via JSON stringify.
tokens += BLOCK_OVERHEAD + Math.ceil(JSON.stringify(block).length / 4) tokens += BLOCK_OVERHEAD + Math.ceil(JSON.stringify(block).length / charsPerToken)
} }
} }
return tokens return tokens
@@ -266,7 +271,7 @@ export class BasicCompactService extends CompactService {
total += this.estimateContentTokens(msg.content) total += this.estimateContentTokens(msg.content)
total += ROLE_OVERHEAD total += ROLE_OVERHEAD
} }
if (systemPrompt) total += Math.ceil(systemPrompt.length / 4) if (systemPrompt) total += Math.ceil(systemPrompt.length / this.config.charsPerToken)
return total return total
} }

View File

@@ -10,10 +10,12 @@
*/ */
/** /**
* Backend configuration. Every knob is REQUIRED except `auto`: there is no * Backend configuration. Every knob is REQUIRED except `auto` and
* concrete data yet to justify default thresholds/budgets, so a consumer must * `charsPerToken`: there is no concrete data yet to justify default
* state each value explicitly rather than inherit a guessed default. `auto` * thresholds/budgets, so a consumer must state each value explicitly rather
* alone defaults to `true` (auto-compaction is the intended posture). * than inherit a guessed default. `auto` alone defaults to `true`
* (auto-compaction is the intended posture), and `charsPerToken` defaults to
* the English-text heuristic its estimator was calibrated on.
*/ */
export interface BasicCompactConfig { export interface BasicCompactConfig {
/** Context window size in tokens. */ /** Context window size in tokens. */
@@ -30,13 +32,21 @@ export interface BasicCompactConfig {
compactionRetries: number compactionRetries: number
/** Enable automatic compaction on the `agent/pre-step` seam (default true). */ /** Enable automatic compaction on the `agent/pre-step` seam (default true). */
auto?: boolean auto?: boolean
/**
* Text density for the token estimator: estimated tokens = chars /
* `charsPerToken`. Defaults to 4 (typical English text). A CJK-heavy
* deployment should set ~1-2 — CJK runs at roughly 1-2 chars per token, so
* the default UNDERestimates several-fold and compaction fires far too late.
* May be fractional.
*/
charsPerToken?: number
} }
/** Resolved config with `auto` defaulted. */ /** Resolved config with `auto` and `charsPerToken` defaulted. */
export type ResolvedConfig = Required<BasicCompactConfig> export type ResolvedConfig = Required<BasicCompactConfig>
/** /**
* Default `auto` when unset and reject nonsensical numeric knobs. * Default `auto`/`charsPerToken` when unset and reject nonsensical numeric knobs.
* *
* Convergence is not a static config invariant: provider generation caps can be * Convergence is not a static config invariant: provider generation caps can be
* spent on hidden or surfaced reasoning tokens, and the model may emit a summary * spent on hidden or surfaced reasoning tokens, and the model may emit a summary
@@ -46,13 +56,14 @@ export type ResolvedConfig = Required<BasicCompactConfig>
* throwing if the surface still exceeds the threshold. * throwing if the surface still exceeds the threshold.
*/ */
export function resolveConfig(config: BasicCompactConfig): ResolvedConfig { export function resolveConfig(config: BasicCompactConfig): ResolvedConfig {
const resolved: ResolvedConfig = { auto: true, ...config } const resolved: ResolvedConfig = { auto: true, charsPerToken: 4, ...config }
assertPositiveInteger('contextWindow', resolved.contextWindow) assertPositiveInteger('contextWindow', resolved.contextWindow)
assertRatio('thresholdRatio', resolved.thresholdRatio) assertRatio('thresholdRatio', resolved.thresholdRatio)
assertNonNegativeInteger('retainTokens', resolved.retainTokens) assertNonNegativeInteger('retainTokens', resolved.retainTokens)
assertPositiveInteger('maxTokens', resolved.maxTokens) assertPositiveInteger('maxTokens', resolved.maxTokens)
assertNonNegativeInteger('compactionRetries', resolved.compactionRetries) assertNonNegativeInteger('compactionRetries', resolved.compactionRetries)
assertPositiveFinite('charsPerToken', resolved.charsPerToken)
if (typeof resolved.summarizationModel !== 'string') { if (typeof resolved.summarizationModel !== 'string') {
throw new Error('BasicCompactConfig: summarizationModel must be a string.') throw new Error('BasicCompactConfig: summarizationModel must be a string.')
} }
@@ -74,6 +85,12 @@ function assertNonNegativeInteger(name: string, value: number): void {
} }
} }
function assertPositiveFinite(name: string, value: number): void {
if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0) {
throw new Error(`BasicCompactConfig: ${name} (${value}) must be a positive finite number.`)
}
}
function assertRatio(name: string, value: number): void { function assertRatio(name: string, value: number): void {
if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0 || value > 1) { if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0 || value > 1) {
throw new Error(`BasicCompactConfig: ${name} (${value}) must be a number in (0, 1].`) throw new Error(`BasicCompactConfig: ${name} (${value}) must be a number in (0, 1].`)

View File

@@ -820,6 +820,19 @@ describe('BasicCompactService token estimation (char/4 heuristic)', () => {
const svc = new BasicCompactService(new Context(), cfg({ auto: false })) const svc = new BasicCompactService(new Context(), cfg({ auto: false }))
expect(svc.estimateContentTokens([])).toBe(0) expect(svc.estimateContentTokens([])).toBe(0)
}) })
it('honors a configured charsPerToken (fractional densities included)', () => {
// 'this is a somewhat longer text block' = 36 chars.
const blocks: ContentBlock[] = [{ type: 'text', text: 'this is a somewhat longer text block' }]
// charsPerToken 2: ceil(36/2)+4 = 22 — a CJK-density config doubles the estimate.
const dense = new BasicCompactService(new Context(), cfg({ auto: false, charsPerToken: 2 }))
expect(dense.estimateContentTokens(blocks)).toBe(22)
// Fractional density is legal: ceil(36/1.5)+4 = 28.
const fractional = new BasicCompactService(new Context(), cfg({ auto: false, charsPerToken: 1.5 }))
expect(fractional.estimateContentTokens(blocks)).toBe(28)
// The system-prompt term scales with the same knob: 36-char prompt at density 2 → ceil(36/2) = 18.
expect(dense.estimateTokens([], 'this is a somewhat longer text block')).toBe(18)
})
}) })
describe('BasicCompactService HMR safety', () => { describe('BasicCompactService HMR safety', () => {
@@ -862,6 +875,10 @@ describe('BasicCompactService config validation', () => {
)).toThrow(/summarizationModel must be a string/) )).toThrow(/summarizationModel must be a string/)
expect(() => new BasicCompactService(new Context(), cfg({ auto: 'no' } as unknown as Partial<BasicCompactConfig>))) expect(() => new BasicCompactService(new Context(), cfg({ auto: 'no' } as unknown as Partial<BasicCompactConfig>)))
.toThrow(/auto must be a boolean/) .toThrow(/auto must be a boolean/)
expect(() => new BasicCompactService(new Context(), cfg({ auto: false, charsPerToken: 0 })))
.toThrow(/charsPerToken .* positive finite number/)
expect(() => new BasicCompactService(new Context(), cfg({ auto: false, charsPerToken: Number.NaN })))
.toThrow(/charsPerToken .* positive finite number/)
}) })
it('accepts a large retain budget because convergence is enforced dynamically', () => { it('accepts a large retain budget because convergence is enforced dynamically', () => {

View File

@@ -26,9 +26,6 @@ import { basename, dirname, join, resolve } from 'node:path'
import { TextDecoder } from 'node:util' import { TextDecoder } from 'node:util'
import { FsError, FsTargetKey, FsVersion } from '@deepseek-ai/dsh-fs' import { FsError, FsTargetKey, FsVersion } from '@deepseek-ai/dsh-fs'
/** Files at or above this size stream their text; smaller files read whole. */
export const STREAM_MIN_SIZE = 10 * 1024 * 1024
const BINARY_SAMPLE_BYTES = 8192 const BINARY_SAMPLE_BYTES = 8192
function isENOENT(error: unknown): boolean { function isENOENT(error: unknown): boolean {
@@ -85,13 +82,10 @@ function versionOf(info: Stats): FsVersion {
} }
/** /**
* Test seam: lets specs force the streaming read path (via a small * Test seam: lets specs pin the temp-file name (to prove exclusive-open
* `streamMinSize`) and pin the temp-file name (to prove exclusive-open * behavior) without a name race.
* behavior) without a 10 MB fixture or a name race.
*/ */
export interface FsIoInternals { export interface FsIoInternals {
/** Override {@link STREAM_MIN_SIZE} for read routing. */
streamMinSize?: number
/** Override the generated private staging-dir name (relative to the target dir). */ /** Override the generated private staging-dir name (relative to the target dir). */
tempDirName?: (writePath: string) => string tempDirName?: (writePath: string) => string
/** Override the generated temp-file name (relative to the private staging dir). */ /** Override the generated temp-file name (relative to the private staging dir). */

View File

@@ -41,7 +41,6 @@ import {
import type { FsIoInternals } from './fsio.ts' import type { FsIoInternals } from './fsio.ts'
export { export {
STREAM_MIN_SIZE,
applyLiteralEdit, applyLiteralEdit,
listDirectory, listDirectory,
probe, probe,

View File

@@ -11,11 +11,22 @@ await ctx.plugin(ToolFs) // this package — re
`@deepseek-ai/dsh-fs-policy` is **optional**: omit it and the tools run against the bare provider (unconditional write/overwrite/edit, no observed-state). A deployment that loads these tools is expected to also load it, so the behavior is read-before-write/edit. `@deepseek-ai/dsh-fs-policy` is **optional**: omit it and the tools run against the bare provider (unconditional write/overwrite/edit, no observed-state). A deployment that loads these tools is expected to also load it, so the behavior is read-before-write/edit.
## Config
All keys are optional; the defaults are the shipped read caps.
| Key | Default | Meaning |
|---|---|---|
| `readLimit` | `2000` | Default and maximum lines returned by one `read` call (the tool schema advertises it as the `limit` default). |
| `readMaxLineLength` | `2000` | Characters kept per line before truncation (the suffix names the cap). |
| `readMaxBytes` | `51200` | Byte cap on one `read` call's selected lines; overflow ends the window with a "capped" footer. |
| `readStreamMinSize` | `10485760` | Files at or above this size (or with unknown size) stream instead of loading whole into memory. |
## Tools (schemas per [the filesystem tool schemas RFC](../../../docs/rfc/implemented/feature/2026-06-17-filesystem-tool-schemas.md)) ## Tools (schemas per [the filesystem tool schemas RFC](../../../docs/rfc/implemented/feature/2026-06-17-filesystem-tool-schemas.md))
| Tool | Arguments | Behavior | | Tool | Arguments | Behavior |
|---|---|---| |---|---|---|
| `read` | `file_path`, `offset?`, `limit?` | Line-numbered UTF-8 content with a pagination footer. `offset` is 1-based; `limit` defaults to and caps at 2000 lines. | | `read` | `file_path`, `offset?`, `limit?` | Line-numbered UTF-8 content with a pagination footer. `offset` is 1-based; `limit` defaults to and caps at the configured `readLimit` (2000). |
| `write` | `file_path`, `content` | Create or fully replace a file. With the policy plugin: overwriting an existing file requires a prior `read` at the unchanged version; creating a new file does not. Without it: unconditional. | | `write` | `file_path`, `content` | Create or fully replace a file. With the policy plugin: overwriting an existing file requires a prior `read` at the unchanged version; creating a new file does not. Without it: unconditional. |
| `edit` | `file_path`, non-empty `old_string`, `new_string`, `replace_all?` | Literal replacement; unique match required unless `replace_all` is true. With the policy plugin: requires a prior `read` (any window) and the file unchanged since. Without it: unconditional. | | `edit` | `file_path`, non-empty `old_string`, `new_string`, `replace_all?` | Literal replacement; unique match required unless `replace_all` is true. With the policy plugin: requires a prior `read` (any window) and the file unchanged since. Without it: unconditional. |

View File

@@ -22,7 +22,8 @@
], ],
"license": "BSD-3-Clause", "license": "BSD-3-Clause",
"dependencies": { "dependencies": {
"diff": "^9.0.0" "diff": "^9.0.0",
"schemastery": "^3.18.0"
}, },
"peerDependencies": { "peerDependencies": {
"@deepseek-ai/dsh-fs": "^0.0.1", "@deepseek-ai/dsh-fs": "^0.0.1",

View File

@@ -23,11 +23,14 @@
*/ */
import type { Context } from 'cordis' import type { Context } from 'cordis'
import { applyReadTool } from './read.ts' import z from 'schemastery'
import { applyReadTool, READ_LIMIT, STREAM_MIN_SIZE } from './read.ts'
import { applyWriteTool } from './write.ts' import { applyWriteTool } from './write.ts'
import { applyEditTool } from './edit.ts' import { applyEditTool } from './edit.ts'
import { READ_MAX_BYTES, READ_MAX_LINE_LENGTH } from './read-render.ts'
export { READ_LIMIT, STREAM_MIN_SIZE, applyReadTool, parseReadArgs } from './read.ts' export { READ_LIMIT, STREAM_MIN_SIZE, applyReadTool, parseReadArgs } from './read.ts'
export type { ReadToolCaps } from './read.ts'
export { applyWriteTool, formatWriteOutput, parseWriteArgs } from './write.ts' export { applyWriteTool, formatWriteOutput, parseWriteArgs } from './write.ts'
export { applyEditTool, formatEditOutput, parseEditArgs } from './edit.ts' export { applyEditTool, formatEditOutput, parseEditArgs } from './edit.ts'
export { READ_MAX_BYTES, READ_MAX_LINE_LENGTH, buildWindow, formatReadOutput } from './read-render.ts' export { READ_MAX_BYTES, READ_MAX_LINE_LENGTH, buildWindow, formatReadOutput } from './read-render.ts'
@@ -41,9 +44,49 @@ export const name = 'tool-fs'
/** Services required by the filesystem tool suite. */ /** Services required by the filesystem tool suite. */
export const inject = ['tools', 'fs', 'systemPrompt'] export const inject = ['tools', 'fs', 'systemPrompt']
/** Plugin config (all optional — `Config` supplies the defaults). */
export interface Config {
/** Default and maximum number of lines returned by one `read` call. */
readLimit?: number
/** Maximum characters returned for a single line before truncation. */
readMaxLineLength?: number
/** Maximum bytes returned for the selected lines of one `read` call. */
readMaxBytes?: number
/** Files at or above this size stream instead of loading whole into memory. */
readStreamMinSize?: number
}
export const Config: z<Config> = z.object({
readLimit: z.number().default(READ_LIMIT),
readMaxLineLength: z.number().default(READ_MAX_LINE_LENGTH),
readMaxBytes: z.number().default(READ_MAX_BYTES),
readStreamMinSize: z.number().default(STREAM_MIN_SIZE),
})
/** The shape after schemastery applied the defaults. */
type ResolvedConfig = Required<Config>
/** A read cap must be a positive finite number to bound output and memory. */
function assertPositiveFinite(name: string, value: number): void {
if (!Number.isFinite(value) || value <= 0) {
throw new Error(`tool-fs: ${name} must be a positive finite number`)
}
}
/** Register the full `read`/`write`/`edit` filesystem tool suite. */ /** Register the full `read`/`write`/`edit` filesystem tool suite. */
export function apply(ctx: Context): void { export function apply(ctx: Context, config: Config): void {
applyReadTool(ctx) // schemastery (Config) has already filled every defaulted field.
const resolved = config as ResolvedConfig
assertPositiveFinite('readLimit', resolved.readLimit)
assertPositiveFinite('readMaxLineLength', resolved.readMaxLineLength)
assertPositiveFinite('readMaxBytes', resolved.readMaxBytes)
assertPositiveFinite('readStreamMinSize', resolved.readStreamMinSize)
applyReadTool(ctx, {
limit: resolved.readLimit,
maxLineLength: resolved.readMaxLineLength,
maxBytes: resolved.readMaxBytes,
streamMinSize: resolved.readStreamMinSize,
})
applyWriteTool(ctx) applyWriteTool(ctx)
applyEditTool(ctx) applyEditTool(ctx)
} }

View File

@@ -19,21 +19,22 @@
import { FsError } from '@deepseek-ai/dsh-fs' import { FsError } from '@deepseek-ai/dsh-fs'
import type { FsVersion } from '@deepseek-ai/dsh-fs' import type { FsVersion } from '@deepseek-ai/dsh-fs'
/** Maximum characters returned for a single line. */ /** Default maximum characters returned for a single line (the `readMaxLineLength` config). */
export const READ_MAX_LINE_LENGTH = 2000 export const READ_MAX_LINE_LENGTH = 2000
/** Maximum bytes returned for selected file lines. */ /** Default maximum bytes returned for selected file lines (the `readMaxBytes` config). */
export const READ_MAX_BYTES = 50 * 1024 export const READ_MAX_BYTES = 50 * 1024
const READ_MAX_LINE_SUFFIX = `... (line truncated to ${READ_MAX_LINE_LENGTH} chars)`
const LINE_BUFFER_CAP = READ_MAX_LINE_LENGTH + 1
/** Resolved read window. The consumer applies its defaults/caps before calling. */ /** Resolved read window. The consumer applies its defaults/caps before calling. */
export interface ReadWindow { export interface ReadWindow {
/** 1-based first line to return. */ /** 1-based first line to return. */
offset: number offset: number
/** Maximum number of lines to return. */ /** Maximum number of lines to return. */
limit: number limit: number
/** Maximum characters returned for a single line; overflow is truncated with a suffix. */
maxLineLength: number
/** Maximum bytes of selected output; overflow stops the scan and marks `truncatedByBytes`. */
maxBytes: number
} }
/** One line returned from a text file. */ /** One line returned from a text file. */
@@ -82,8 +83,8 @@ function newAccumulator(): WindowAccumulator {
return { lines: [], totalLines: 0, outputBytes: 0, truncatedByBytes: false, done: false } return { lines: [], totalLines: 0, outputBytes: 0, truncatedByBytes: false, done: false }
} }
function truncateLine(line: string): string { function truncateLine(line: string, maxLineLength: number): string {
return line.length > READ_MAX_LINE_LENGTH ? `${line.substring(0, READ_MAX_LINE_LENGTH)}${READ_MAX_LINE_SUFFIX}` : line return line.length > maxLineLength ? `${line.substring(0, maxLineLength)}... (line truncated to ${maxLineLength} chars)` : line
} }
function lineByteSize(line: string, currentLineCount: number): number { function lineByteSize(line: string, currentLineCount: number): number {
@@ -94,9 +95,9 @@ function consumeLine(acc: WindowAccumulator, rawLine: string, request: ReadWindo
acc.totalLines += 1 acc.totalLines += 1
if (acc.totalLines < request.offset || acc.lines.length >= request.limit) return if (acc.totalLines < request.offset || acc.lines.length >= request.limit) return
const text = truncateLine(rawLine) const text = truncateLine(rawLine, request.maxLineLength)
const bytes = lineByteSize(text, acc.lines.length) const bytes = lineByteSize(text, acc.lines.length)
if (acc.outputBytes + bytes > READ_MAX_BYTES) { if (acc.outputBytes + bytes > request.maxBytes) {
acc.truncatedByBytes = true acc.truncatedByBytes = true
acc.done = true acc.done = true
return return
@@ -121,7 +122,7 @@ function finish(acc: WindowAccumulator, request: ReadWindow, displayPath: string
* Accepts an `AsyncIterable<string>` (a chunked `streamText`) or an * Accepts an `AsyncIterable<string>` (a chunked `streamText`) or an
* `Iterable<string>` (a whole-file `readText` wrapped as `[text]`), so one code * `Iterable<string>` (a whole-file `readText` wrapped as `[text]`), so one code
* path serves both. Scans for newlines with a capped line buffer (a newline-free * path serves both. Scans for newlines with a capped line buffer (a newline-free
* giant line is truncated, never buffered past {@link READ_MAX_LINE_LENGTH}), * giant line is truncated, never buffered past `request.maxLineLength`),
* enforces the byte cap, and throws `FS_NOT_FOUND` for an offset past EOF. * enforces the byte cap, and throws `FS_NOT_FOUND` for an offset past EOF.
*/ */
export async function buildWindow( export async function buildWindow(
@@ -130,12 +131,14 @@ export async function buildWindow(
displayPath: string, displayPath: string,
): Promise<WindowResult> { ): Promise<WindowResult> {
const acc = newAccumulator() const acc = newAccumulator()
// One char past the truncation point is enough to prove a line overflows.
const lineBufferCap = request.maxLineLength + 1
let lineBuffer = '' let lineBuffer = ''
function appendToLineBuffer(segment: string): void { function appendToLineBuffer(segment: string): void {
if (lineBuffer.length >= LINE_BUFFER_CAP) return if (lineBuffer.length >= lineBufferCap) return
lineBuffer += segment lineBuffer += segment
if (lineBuffer.length > LINE_BUFFER_CAP) lineBuffer = lineBuffer.slice(0, LINE_BUFFER_CAP) if (lineBuffer.length > lineBufferCap) lineBuffer = lineBuffer.slice(0, lineBufferCap)
} }
function flushLine(): void { function flushLine(): void {

View File

@@ -23,12 +23,27 @@ import { buildWindow, formatReadOutput } from './read-render.ts'
import type { FileReadOutcome } from './read-render.ts' import type { FileReadOutcome } from './read-render.ts'
import { sessionCwd } from './session-cwd.ts' import { sessionCwd } from './session-cwd.ts'
/** Default and maximum number of lines returned by one `read` call. */ /** Default and maximum number of lines returned by one `read` call (the `readLimit` config). */
export const READ_LIMIT = 2000 export const READ_LIMIT = 2000
/** Files at or above this size stream; smaller files read whole into memory. */ /**
* Default streaming threshold (the `readStreamMinSize` config): files at or
* above this size stream; smaller files read whole into memory.
*/
export const STREAM_MIN_SIZE = 10 * 1024 * 1024 export const STREAM_MIN_SIZE = 10 * 1024 * 1024
/** Resolved read-tool caps — plugin config after defaulting (see `Config` in index.ts). */
export interface ReadToolCaps {
/** Default and maximum number of lines returned by one call. */
limit: number
/** Maximum characters returned for a single line. */
maxLineLength: number
/** Maximum bytes returned for selected file lines. */
maxBytes: number
/** Files at or above this size stream; smaller files read whole into memory. */
streamMinSize: number
}
/** Validated `read` arguments after defaulting. */ /** Validated `read` arguments after defaulting. */
interface ReadInput { interface ReadInput {
filePath: string filePath: string
@@ -43,17 +58,17 @@ function parsePositiveInteger(value: number, name: string): number {
return value return value
} }
/** Validate value constraints the schema DSL can't express. */ /** Validate value constraints the schema DSL can't express. `maxLimit` is the deployment's line cap. */
export function parseReadArgs(args: { file_path: string; offset?: number; limit?: number }): ReadInput { export function parseReadArgs(args: { file_path: string; offset?: number; limit?: number }, maxLimit: number): ReadInput {
if (args.file_path.trim().length === 0) throw new Error('file_path must be a non-empty string') if (args.file_path.trim().length === 0) throw new Error('file_path must be a non-empty string')
const offset = args.offset === undefined ? 1 : parsePositiveInteger(args.offset, 'offset') const offset = args.offset === undefined ? 1 : parsePositiveInteger(args.offset, 'offset')
const limit = args.limit === undefined ? READ_LIMIT : parsePositiveInteger(args.limit, 'limit') const limit = args.limit === undefined ? maxLimit : parsePositiveInteger(args.limit, 'limit')
if (limit > READ_LIMIT) throw new Error(`limit must be less than or equal to ${READ_LIMIT}`) if (limit > maxLimit) throw new Error(`limit must be less than or equal to ${maxLimit}`)
return { filePath: args.file_path, offset, limit } return { filePath: args.file_path, offset, limit }
} }
/** Register the `read` tool and its system-prompt guidance. */ /** Register the `read` tool and its system-prompt guidance. */
export function applyReadTool(ctx: Context): void { export function applyReadTool(ctx: Context, caps: ReadToolCaps): void {
ctx.systemPrompt.section({ ctx.systemPrompt.section({
name: 'tool:read', name: 'tool:read',
order: 100, order: 100,
@@ -66,10 +81,10 @@ export function applyReadTool(ctx: Context): void {
parameters: { parameters: {
file_path: { type: 'string', required: true, description: 'Path to read, resolved by the filesystem backend.' }, file_path: { type: 'string', required: true, description: 'Path to read, resolved by the filesystem backend.' },
offset: { type: 'number', description: '1-based first line to return. Defaults to 1.' }, offset: { type: 'number', description: '1-based first line to return. Defaults to 1.' },
limit: { type: 'number', description: `Maximum number of lines to return. Defaults to ${READ_LIMIT}.` }, limit: { type: 'number', description: `Maximum number of lines to return. Defaults to ${caps.limit}.` },
}, },
async execute(args, exec): Promise<ContentBlock[]> { async execute(args, exec): Promise<ContentBlock[]> {
const input = parseReadArgs(args) const input = parseReadArgs(args, caps.limit)
const cwd = sessionCwd(exec) const cwd = sessionCwd(exec)
const target = await ctx.fs.resolve(input.filePath, cwd !== undefined ? { cwd } : undefined) const target = await ctx.fs.resolve(input.filePath, cwd !== undefined ? { cwd } : undefined)
@@ -83,10 +98,14 @@ export function applyReadTool(ctx: Context): void {
// Stream when the file is large OR size is unknown, so a size-less backend // Stream when the file is large OR size is unknown, so a size-less backend
// never buffers an arbitrarily large file. // never buffers an arbitrarily large file.
const chunks = info.size === undefined || info.size >= STREAM_MIN_SIZE const chunks = info.size === undefined || info.size >= caps.streamMinSize
? await ctx.fs.streamText(target, exec.signal) ? await ctx.fs.streamText(target, exec.signal)
: [await ctx.fs.readText(target, exec.signal)] : [await ctx.fs.readText(target, exec.signal)]
const window = await buildWindow(chunks, { offset: input.offset, limit: input.limit }, target.displayPath) const window = await buildWindow(
chunks,
{ offset: input.offset, limit: input.limit, maxLineLength: caps.maxLineLength, maxBytes: caps.maxBytes },
target.displayPath,
)
const outcome: FileReadOutcome = { const outcome: FileReadOutcome = {
offset: input.offset, offset: input.offset,
@@ -106,7 +125,8 @@ export function applyReadTool(ctx: Context): void {
// appended (`Read foo.txt (5 - 8)`), `read` kind (icon), and a follow-along // appended (`Read foo.txt (5 - 8)`), `read` kind (icon), and a follow-along
// location whose line is the read's offset (defaulting to 1). The window is // location whose line is the read's offset (defaulting to 1). The window is
// derived from the RAW args (offset/limit as the model passed them), NOT the // derived from the RAW args (offset/limit as the model passed them), NOT the
// tool's defaulted 1/READ_LIMIT, so an unbounded read shows a bare title. // tool's defaulted 1/configured limit, so an unbounded read shows a bare
// title (and the presenter stays a pure function of args, config-free).
presentCall(args): GenericCallView { presentCall(args): GenericCallView {
const { offset, limit } = args const { offset, limit } = args
const window = limit !== undefined && limit > 0 const window = limit !== undefined && limit > 0

View File

@@ -6,10 +6,11 @@
*/ */
import { describe, expect, it } from 'vitest' import { describe, expect, it } from 'vitest'
import { buildWindow, READ_MAX_LINE_LENGTH } from '@deepseek-ai/dsh-tool-fs' import { buildWindow, READ_MAX_BYTES, READ_MAX_LINE_LENGTH } from '@deepseek-ai/dsh-tool-fs'
import type { ReadWindow } from '@deepseek-ai/dsh-tool-fs' import type { ReadWindow } from '@deepseek-ai/dsh-tool-fs'
const READ_ALL: ReadWindow = { offset: 1, limit: 2000 } const DEFAULT_CAPS = { maxLineLength: READ_MAX_LINE_LENGTH, maxBytes: READ_MAX_BYTES }
const READ_ALL: ReadWindow = { offset: 1, limit: 2000, ...DEFAULT_CAPS }
/** Yield `text` as one chunk (whole-file read shape). */ /** Yield `text` as one chunk (whole-file read shape). */
async function* whole(text: string): AsyncIterable<string> { async function* whole(text: string): AsyncIterable<string> {
@@ -34,7 +35,7 @@ describe('buildWindow', () => {
}) })
it('applies offset/limit', async () => { it('applies offset/limit', async () => {
const result = await buildWindow(whole('one\ntwo\nthree\nfour'), { offset: 2, limit: 2 }, 'f') const result = await buildWindow(whole('one\ntwo\nthree\nfour'), { offset: 2, limit: 2, ...DEFAULT_CAPS }, 'f')
expect(result.lines.map(l => l.number)).toEqual([2, 3]) expect(result.lines.map(l => l.number)).toEqual([2, 3])
expect(result.totalLines).toBe(4) expect(result.totalLines).toBe(4)
}) })
@@ -62,7 +63,7 @@ describe('buildWindow', () => {
}) })
it('rejects an offset past EOF', async () => { it('rejects an offset past EOF', async () => {
await expect(buildWindow(whole('one\ntwo'), { offset: 9, limit: 1 }, 'f')).rejects.toMatchObject({ code: 'FS_NOT_FOUND' }) await expect(buildWindow(whole('one\ntwo'), { offset: 9, limit: 1, ...DEFAULT_CAPS }, 'f')).rejects.toMatchObject({ code: 'FS_NOT_FOUND' })
}) })
it('flushes a final line with no trailing newline', async () => { it('flushes a final line with no trailing newline', async () => {
@@ -76,9 +77,22 @@ describe('buildWindow', () => {
expect(result.totalLines).toBe(2) expect(result.totalLines).toBe(2)
}) })
describe('caps are per-request (the plugin config reaches the window)', () => {
it('truncates lines at a custom maxLineLength and names it in the suffix', async () => {
const result = await buildWindow(whole('abcdefghij'), { offset: 1, limit: 10, maxLineLength: 5, maxBytes: READ_MAX_BYTES }, 'f')
expect(result.lines[0]?.text).toBe('abcde... (line truncated to 5 chars)')
})
it('caps output at a custom maxBytes', async () => {
const result = await buildWindow(whole('aaaa\nbbbb\ncccc'), { offset: 1, limit: 10, maxLineLength: 2000, maxBytes: 9 }, 'f')
expect(result.lines.map(l => l.text)).toEqual(['aaaa', 'bbbb'])
expect(result.truncatedByBytes).toBe(true)
})
})
describe('chunked input (streamed read shape)', () => { describe('chunked input (streamed read shape)', () => {
it('windows identically when text arrives in small chunks', async () => { it('windows identically when text arrives in small chunks', async () => {
const result = await buildWindow(chunked('one\ntwo\nthree', 2), { offset: 2, limit: 1 }, 'f') const result = await buildWindow(chunked('one\ntwo\nthree', 2), { offset: 2, limit: 1, ...DEFAULT_CAPS }, 'f')
expect(result.lines).toEqual([{ number: 2, text: 'two' }]) expect(result.lines).toEqual([{ number: 2, text: 'two' }])
expect(result.totalLines).toBe(3) expect(result.totalLines).toBe(3)
}) })

View File

@@ -495,3 +495,70 @@ describe('result-time contextual diff (meta + presentResult)', () => {
expect(view).toEqual({ card: 'diff', title: 'Write a.txt', diffs: [{ path: 'a.txt', oldText: null, newText: 'y' }] }) expect(view).toEqual({ card: 'diff', title: 'Write a.txt', diffs: [{ path: 'a.txt', oldText: null, newText: 'y' }] })
}) })
}) })
describe('read caps are plugin config', () => {
async function setupWith(config: ToolFs.Config) {
const ctx = new Context()
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
await ctx.plugin(FakeFs)
await ctx.plugin(FsPolicy)
await ctx.plugin(ToolFs, config)
return { ctx, fs: ctx.fs as FakeFs }
}
it('a configured readLimit is both the default and the cap, and the schema names it', async () => {
const { ctx, fs } = await setupWith({ readLimit: 2 })
fs.files.set('key:a.txt', 'one\ntwo\nthree\nfour')
const result = await call(ctx, 'read', { file_path: 'a.txt' })
expect(text(result)).toContain('(Showing lines 1-2 of 4. Use offset=3 to continue.)')
const overCap = await call(ctx, 'read', { file_path: 'a.txt', limit: 3 })
expect(overCap.isError).toBe(true)
expect(text(overCap)).toContain('less than or equal to 2')
const readSchema = ctx.tools.schemas().find(s => s.name === 'read')
expect(JSON.stringify(readSchema)).toContain('Defaults to 2.')
})
it('a configured readMaxLineLength truncates lines at the configured length', async () => {
const { ctx, fs } = await setupWith({ readMaxLineLength: 4 })
fs.files.set('key:a.txt', 'abcdefgh')
const result = await call(ctx, 'read', { file_path: 'a.txt' })
expect(text(result)).toContain('1: abcd... (line truncated to 4 chars)')
})
it('a configured readMaxBytes caps the window at the configured bytes', async () => {
const { ctx, fs } = await setupWith({ readMaxBytes: 9 })
fs.files.set('key:a.txt', 'aaaa\nbbbb\ncccc')
const result = await call(ctx, 'read', { file_path: 'a.txt' })
expect(text(result)).toContain('Output capped.')
expect(text(result)).not.toContain('cccc')
})
it('a configured readStreamMinSize routes smaller files to the streaming path', async () => {
const { ctx, fs } = await setupWith({ readStreamMinSize: 5 })
fs.files.set('key:a.txt', 'alpha\nbeta')
const readSpy = vi.spyOn(fs, 'readText')
const streamSpy = vi.spyOn(fs, 'streamText')
const result = await call(ctx, 'read', { file_path: 'a.txt' })
expect(result.isError).toBe(false)
expect(streamSpy).toHaveBeenCalled()
expect(readSpy).not.toHaveBeenCalled()
})
it.each([
['readLimit', { readLimit: 0 }],
['readMaxLineLength', { readMaxLineLength: -1 }],
['readMaxBytes', { readMaxBytes: Number.NaN }],
['readStreamMinSize', { readStreamMinSize: 0 }],
] as const)('rejects a non-positive %s at load', async (name, config) => {
const ctx = new Context()
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
await ctx.plugin(FakeFs)
await expect(ctx.plugin(ToolFs, config)).rejects.toThrow(new RegExp(`tool-fs: ${name} must be a positive finite number`))
})
it('has no default export (namespace plugin export shape)', () => {
expect('default' in ToolFs).toBe(false)
})
})

View File

@@ -8,6 +8,7 @@
"references": [ "references": [
{ "path": "../../../vendor/cosmokit" }, { "path": "../../../vendor/cosmokit" },
{ "path": "../../../vendor/cordis" }, { "path": "../../../vendor/cordis" },
{ "path": "../../../vendor/schemastery" },
{ "path": "../../llm/llm" }, { "path": "../../llm/llm" },
{ "path": "../../core/tools" }, { "path": "../../core/tools" },
{ "path": "../../core/system-prompt" }, { "path": "../../core/system-prompt" },

View File

@@ -58,6 +58,18 @@ export function appendHookInvoked(session: Session, invocation: HookInvocation):
}) })
} }
/**
* Truncate a hook's stderr for {@link HookResultRecord.stderrSummary}: trimmed,
* `undefined` when empty, cut at `maxChars` with an ellipsis when over. The
* bound is a parameter — like `runHook`'s `defaultTimeoutMs`, each bridge owns
* the config default and passes it in.
*/
export function summarizeStderr(stderr: string, maxChars: number): string | undefined {
const t = stderr.trim()
if (t.length === 0) return undefined
return t.length > maxChars ? t.slice(0, maxChars) + '…' : t
}
/** Append a `hook/result` outcome event to `session` (pairs with a prior `hook/invoked`). */ /** Append a `hook/result` outcome event to `session` (pairs with a prior `hook/invoked`). */
export function appendHookResult(session: Session, record: HookResultRecord): void { export function appendHookResult(session: Session, record: HookResultRecord): void {
session.append('hook/result', { session.append('hook/result', {

View File

@@ -34,5 +34,5 @@ export { runHook } from './runner.ts'
export type { RunHookOptions, RunHookResult } from './runner.ts' export type { RunHookOptions, RunHookResult } from './runner.ts'
export { mergeHookOutputs } from './merge.ts' export { mergeHookOutputs } from './merge.ts'
export type { MergedDecision, MergedHookOutcome } from './merge.ts' export type { MergedDecision, MergedHookOutcome } from './merge.ts'
export { appendHookInvoked, appendHookResult } from './events.ts' export { appendHookInvoked, appendHookResult, summarizeStderr } from './events.ts'
export type { HookInvocation, HookResultRecord } from './events.ts' export type { HookInvocation, HookResultRecord } from './events.ts'

View File

@@ -1,6 +1,6 @@
import { describe, expect, it } from 'vitest' import { describe, expect, it } from 'vitest'
import { Session, SessionId } from '@deepseek-ai/dsh-session' import { Session, SessionId } from '@deepseek-ai/dsh-session'
import { appendHookInvoked, appendHookResult } from '@deepseek-ai/dsh-hook-protocol' import { appendHookInvoked, appendHookResult, summarizeStderr } from '@deepseek-ai/dsh-hook-protocol'
describe('hook/* session events', () => { describe('hook/* session events', () => {
it('appendHookInvoked records a log-only hook/invoked (with matcher when present)', () => { it('appendHookInvoked records a log-only hook/invoked (with matcher when present)', () => {
@@ -59,3 +59,20 @@ describe('hook/* session events', () => {
expect(result?.type === 'hook/result' && result.data.handlerId).toBe('pair-1') expect(result?.type === 'hook/result' && result.data.handlerId).toBe('pair-1')
}) })
}) })
describe('summarizeStderr', () => {
it('returns undefined for empty/whitespace stderr', () => {
expect(summarizeStderr('', 500)).toBeUndefined()
expect(summarizeStderr(' \n\t ', 500)).toBeUndefined()
})
it('passes through a summary at or under the cap, trimmed', () => {
expect(summarizeStderr(' blocked: bad tool ', 500)).toBe('blocked: bad tool')
expect(summarizeStderr('abc', 3)).toBe('abc')
})
it('truncates past the cap with an ellipsis', () => {
expect(summarizeStderr('abcdef', 4)).toBe('abcd…')
expect(summarizeStderr('x'.repeat(600), 500)).toBe('x'.repeat(500) + '…')
})
})

View File

@@ -13,6 +13,7 @@ const config: Config = {
pluginRoot: '/path/to/plugin', // optional: replaces ${CLAUDE_PLUGIN_ROOT} in command strings pluginRoot: '/path/to/plugin', // optional: replaces ${CLAUDE_PLUGIN_ROOT} in command strings
projectDir: '/path/to/project', // optional: replaces ${CLAUDE_PROJECT_DIR} AND sets the hook env var; defaults to the session cwd when omitted projectDir: '/path/to/project', // optional: replaces ${CLAUDE_PROJECT_DIR} AND sets the hook env var; defaults to the session cwd when omitted
defaultTimeoutMs: 600_000, // optional: per-hook timeout when a hook sets none (CC default) defaultTimeoutMs: 600_000, // optional: per-hook timeout when a hook sets none (CC default)
stderrSummaryMaxChars: 500, // optional: char cap on the hook/result event's persisted stderr summary
} }
``` ```

View File

@@ -34,6 +34,7 @@ import {
matchesMatcher, matchesMatcher,
mergeHookOutputs, mergeHookOutputs,
runHook, runHook,
summarizeStderr,
type HookOutput, type HookOutput,
type MatcherGroup, type MatcherGroup,
type MergedHookOutcome, type MergedHookOutcome,
@@ -73,6 +74,8 @@ export interface Config {
projectDir?: string projectDir?: string
/** Default per-hook timeout in ms when a hook sets none (CC default: 600000). */ /** Default per-hook timeout in ms when a hook sets none (CC default: 600000). */
defaultTimeoutMs?: number defaultTimeoutMs?: number
/** Character cap for the `hook/result` event's persisted stderr summary. */
stderrSummaryMaxChars?: number
} }
export const Config: z<Config> = z.object({ export const Config: z<Config> = z.object({
@@ -80,6 +83,7 @@ export const Config: z<Config> = z.object({
pluginRoot: z.string(), pluginRoot: z.string(),
projectDir: z.string(), projectDir: z.string(),
defaultTimeoutMs: z.number().default(600_000), defaultTimeoutMs: z.number().default(600_000),
stderrSummaryMaxChars: z.number().default(500),
}) })
/** A stable per-handler id so an invoked/result pair correlates in the log. */ /** A stable per-handler id so an invoked/result pair correlates in the log. */
@@ -91,13 +95,6 @@ function nextHandlerId(point: string): string {
/** The `{kind:'plugin'}` source stamped on every context this bridge injects. */ /** The `{kind:'plugin'}` source stamped on every context this bridge injects. */
const PLUGIN_SOURCE: MessageSource = { kind: 'plugin', plugin: 'hooks-claude' } const PLUGIN_SOURCE: MessageSource = { kind: 'plugin', plugin: 'hooks-claude' }
/** Truncate a stderr blob for the `hook/result` summary field. */
function summarize(stderr: string): string | undefined {
const t = stderr.trim()
if (t.length === 0) return undefined
return t.length > 500 ? t.slice(0, 500) + '…' : t
}
export function apply(ctx: Context, config: Config): void { export function apply(ctx: Context, config: Config): void {
// --- Parse the config ONCE at load. A read/parse failure is contained: the // --- Parse the config ONCE at load. A read/parse failure is contained: the
// bridge logs and registers nothing rather than crashing boot (a typo'd path // bridge logs and registers nothing rather than crashing boot (a typo'd path
@@ -119,6 +116,7 @@ export function apply(ctx: Context, config: Config): void {
} }
const defaultTimeoutMs = config.defaultTimeoutMs ?? 600_000 const defaultTimeoutMs = config.defaultTimeoutMs ?? 600_000
const stderrSummaryMaxChars = config.stderrSummaryMaxChars ?? 500
/** /**
* Run every command hook configured for `point` whose matcher selects * Run every command hook configured for `point` whose matcher selects
@@ -182,7 +180,7 @@ export function apply(ctx: Context, config: Config): void {
ctx.logger.warn(`hooks-claude: ${point} hook emitted a systemMessage, which is not yet surfaced (ignored)`) ctx.logger.warn(`hooks-claude: ${point} hook emitted a systemMessage, which is not yet surfaced (ignored)`)
} }
if (session && opts.turn !== undefined) { if (session && opts.turn !== undefined) {
const stderrSummary = summarize(output.stderr) const stderrSummary = summarizeStderr(output.stderr, stderrSummaryMaxChars)
appendHookResult(session, { appendHookResult(session, {
turn: opts.turn, point, handlerId, turn: opts.turn, point, handlerId,
decision: output.decision ?? (output.continue === false ? 'stop' : 'pass'), decision: output.decision ?? (output.continue === false ? 'stop' : 'pass'),

View File

@@ -27,7 +27,7 @@ function hooks(d: string, h: unknown): string {
writeFileSync(join(d, 'hooks.json'), JSON.stringify({ hooks: h })); return join(d, 'hooks.json') writeFileSync(join(d, 'hooks.json'), JSON.stringify({ hooks: h })); return join(d, 'hooks.json')
} }
type HarnessOpts = { pluginRoot?: string; projectDir?: string } type HarnessOpts = { pluginRoot?: string; projectDir?: string; stderrSummaryMaxChars?: number }
async function harness(configPath: string, adapter: MockAdapter, opts: HarnessOpts = {}): Promise<Context> { async function harness(configPath: string, adapter: MockAdapter, opts: HarnessOpts = {}): Promise<Context> {
const ctx = new Context() const ctx = new Context()
await ctx.plugin(LlmService) await ctx.plugin(LlmService)
@@ -139,6 +139,21 @@ describe('hooks-claude coverage — empty/no-op outcomes and no-agent paths', ()
await waitForIdle(ctx, agent) await waitForIdle(ctx, agent)
const res = events(agent).find(e => e.type === 'hook/result') const res = events(agent).find(e => e.type === 'hook/result')
expect(res?.type === 'hook/result' && res.data.stderrSummary?.endsWith('…')).toBe(true) expect(res?.type === 'hook/result' && res.data.stderrSummary?.endsWith('…')).toBe(true)
expect(res?.type === 'hook/result' && res.data.stderrSummary?.length).toBe(501) // default 500-char cap + ellipsis
})
it('the stderr summary cap is plugin config (stderrSummaryMaxChars)', async () => {
const d = dir()
const s = sh(d, 'long.sh', '#!/usr/bin/env bash\nprintf "x%.0s" {1..600} >&2\nexit 2\n')
const path = hooks(d, { PreToolUse: [{ hooks: [{ type: 'command', command: s }] }] })
const adapter = new MockAdapter([toolCallResponse('c1', 'echo', {}), textResponse('done')])
const ctx = await harness(path, adapter, { stderrSummaryMaxChars: 40 })
ctx.tools.register(defineTool({ name: 'echo', description: 'e', parameters: {}, async execute() { return [{ type: 'text', text: 'ok' }] } }))
const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
agent.send([{ type: 'text', text: 'go' }])
await waitForIdle(ctx, agent)
const res = events(agent).find(e => e.type === 'hook/result')
expect(res?.type === 'hook/result' && res.data.stderrSummary).toBe('x'.repeat(40) + '…')
}) })
}) })

View File

@@ -20,6 +20,7 @@ const config: Config = {
configPath: '/path/to/.codex/hooks.json', // required configPath: '/path/to/.codex/hooks.json', // required
model: 'deepseek-v4', // optional: stamped on every payload (Codex includes `model`) model: 'deepseek-v4', // optional: stamped on every payload (Codex includes `model`)
defaultTimeoutMs: 600_000, // optional: per-hook timeout when a hook sets none defaultTimeoutMs: 600_000, // optional: per-hook timeout when a hook sets none
stderrSummaryMaxChars: 500, // optional: char cap on the hook/result event's persisted stderr summary
} }
``` ```

View File

@@ -27,6 +27,7 @@ import {
matchesMatcher, matchesMatcher,
mergeHookOutputs, mergeHookOutputs,
runHook, runHook,
summarizeStderr,
type HookOutput, type HookOutput,
type MatcherGroup, type MatcherGroup,
type MergedHookOutcome, type MergedHookOutcome,
@@ -49,12 +50,15 @@ export interface Config {
model?: string model?: string
/** Default per-hook timeout in ms when a hook sets none (Codex default: 600000). */ /** Default per-hook timeout in ms when a hook sets none (Codex default: 600000). */
defaultTimeoutMs?: number defaultTimeoutMs?: number
/** Character cap for the `hook/result` event's persisted stderr summary. */
stderrSummaryMaxChars?: number
} }
export const Config: z<Config> = z.object({ export const Config: z<Config> = z.object({
configPath: z.string().required(), configPath: z.string().required(),
model: z.string().default(''), model: z.string().default(''),
defaultTimeoutMs: z.number().default(600_000), defaultTimeoutMs: z.number().default(600_000),
stderrSummaryMaxChars: z.number().default(500),
}) })
let handlerCounter = 0 let handlerCounter = 0
@@ -64,12 +68,6 @@ function nextHandlerId(point: string): string {
const PLUGIN_SOURCE: MessageSource = { kind: 'plugin', plugin: 'hooks-codex' } const PLUGIN_SOURCE: MessageSource = { kind: 'plugin', plugin: 'hooks-codex' }
function summarize(stderr: string): string | undefined {
const t = stderr.trim()
if (t.length === 0) return undefined
return t.length > 500 ? t.slice(0, 500) + '…' : t
}
export function apply(ctx: Context, config: Config): void { export function apply(ctx: Context, config: Config): void {
let parsed: CodexHookConfig = {} let parsed: CodexHookConfig = {}
try { try {
@@ -85,6 +83,7 @@ export function apply(ctx: Context, config: Config): void {
} }
const defaultTimeoutMs = config.defaultTimeoutMs ?? 600_000 const defaultTimeoutMs = config.defaultTimeoutMs ?? 600_000
const stderrSummaryMaxChars = config.stderrSummaryMaxChars ?? 500
const model = config.model ?? '' const model = config.model ?? ''
async function runPoint( async function runPoint(
@@ -140,7 +139,7 @@ export function apply(ctx: Context, config: Config): void {
ctx.logger.warn(`hooks-codex: ${point} hook emitted a systemMessage, which is not yet surfaced (ignored)`) ctx.logger.warn(`hooks-codex: ${point} hook emitted a systemMessage, which is not yet surfaced (ignored)`)
} }
if (session && opts.turn !== undefined) { if (session && opts.turn !== undefined) {
const stderrSummary = summarize(output.stderr) const stderrSummary = summarizeStderr(output.stderr, stderrSummaryMaxChars)
appendHookResult(session, { appendHookResult(session, {
turn: opts.turn, point, handlerId, turn: opts.turn, point, handlerId,
decision: output.decision ?? (output.continue === false ? 'stop' : 'pass'), decision: output.decision ?? (output.continue === false ? 'stop' : 'pass'),

View File

@@ -23,12 +23,12 @@ function hooks(d: string, h: unknown): string {
writeFileSync(join(d, 'hooks.json'), JSON.stringify({ hooks: h })); return join(d, 'hooks.json') writeFileSync(join(d, 'hooks.json'), JSON.stringify({ hooks: h })); return join(d, 'hooks.json')
} }
async function harness(configPath: string, adapter: MockAdapter): Promise<Context> { async function harness(configPath: string, adapter: MockAdapter, opts: { stderrSummaryMaxChars?: number } = {}): Promise<Context> {
const ctx = new Context() const ctx = new Context()
await ctx.plugin(LlmService); await ctx.plugin(SessionStore); await ctx.plugin(SystemPrompt) await ctx.plugin(LlmService); await ctx.plugin(SessionStore); await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry); await ctx.plugin(AgentRegistry); await ctx.plugin(AgentLoop, { agents: [] }) await ctx.plugin(ToolRegistry); await ctx.plugin(AgentRegistry); await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(LocalBashExecutor, { timeoutMs: 10_000 }) await ctx.plugin(LocalBashExecutor, { timeoutMs: 10_000 })
await ctx.plugin(HooksCodex, { configPath, model: 'm' }) await ctx.plugin(HooksCodex, { configPath, model: 'm', ...opts })
ctx.llm.registerAdapter(['mock'], adapter) ctx.llm.registerAdapter(['mock'], adapter)
return ctx return ctx
} }
@@ -204,6 +204,19 @@ describe('hooks-codex coverage — decision mapping paths', () => {
agent.send([{ type: 'text', text: 'go' }]); await waitForIdle(ctx, agent) agent.send([{ type: 'text', text: 'go' }]); await waitForIdle(ctx, agent)
const res = events(agent).find(e => e.type === 'hook/result') const res = events(agent).find(e => e.type === 'hook/result')
expect(res?.type === 'hook/result' && res.data.stderrSummary?.endsWith('…')).toBe(true) expect(res?.type === 'hook/result' && res.data.stderrSummary?.endsWith('…')).toBe(true)
expect(res?.type === 'hook/result' && res.data.stderrSummary?.length).toBe(501) // default 500-char cap + ellipsis
})
it('the stderr summary cap is plugin config (stderrSummaryMaxChars)', async () => {
const d = dir()
hooks(d, { PreToolUse: [{ hooks: [{ type: 'command', command: sh(d, 'l.sh', '#!/usr/bin/env bash\nprintf "x%.0s" {1..600} >&2\nexit 2\n') }] }] })
const adapter = new MockAdapter([toolCallResponse('c1', 'Bash', { command: 'x' }), textResponse('done')])
const ctx = await harness(join(d, 'hooks.json'), adapter, { stderrSummaryMaxChars: 40 })
ctx.tools.register(defineTool({ name: 'Bash', description: 'b', parameters: { command: { type: 'string' } }, async execute() { return [{ type: 'text', text: 'ok' }] } }))
const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
agent.send([{ type: 'text', text: 'go' }]); await waitForIdle(ctx, agent)
const res = events(agent).find(e => e.type === 'hook/result')
expect(res?.type === 'hook/result' && res.data.stderrSummary).toBe('x'.repeat(40) + '…')
}) })
it('warns on a skipped async hook and a direct apply() defaults the timeout', async () => { it('warns on a skipped async hook and a direct apply() defaults the timeout', async () => {

View File

@@ -8,7 +8,7 @@ A SQLite durable session-persistence backend — a second `SessionPersistence` i
Each `SessionEvent` maps 1:1 onto a row in an `events` table `(session_id, seq, type, time, data, source_event_seqs, surface_op)` — `data` is the event payload as JSON text, so the row shape is the event verbatim (including `assistant/chunk`, keeping `seq` contiguous). The two `TEXT` columns `source_event_seqs` and `surface_op` are nullable; they store the event's optional surface-metadata fields (see [session surface](../../../docs/rfc/implemented/architecture/2026-06-18-session-surface.md)). Out-of-log metadata (`SessionHeader`) lives in a `sessions` row. A `sessions` row is written only by the first `append` — its existence is the lazy-materialization signal (`list` reports exactly the sessions that have a row), so no separate column is needed. Each `SessionEvent` maps 1:1 onto a row in an `events` table `(session_id, seq, type, time, data, source_event_seqs, surface_op)` — `data` is the event payload as JSON text, so the row shape is the event verbatim (including `assistant/chunk`, keeping `seq` contiguous). The two `TEXT` columns `source_event_seqs` and `surface_op` are nullable; they store the event's optional surface-metadata fields (see [session surface](../../../docs/rfc/implemented/architecture/2026-06-18-session-surface.md)). Out-of-log metadata (`SessionHeader`) lives in a `sessions` row. A `sessions` row is written only by the first `append` — its existence is the lazy-materialization signal (`list` reports exactly the sessions that have a row), so no separate column is needed.
The repo targets Node ≥ 24 (the root `engines` field), which includes the stable `node:sqlite` module. The database opens with `foreign_keys = ON` (so `ON DELETE CASCADE` drops a session's events with its row) and `journal_mode = WAL`. The table-layout version is stored in `PRAGMA user_version` and checked on open: a fresh database is stamped with the current `SCHEMA_VERSION`; a database written by any other, incompatible build (a non-current `user_version`, older or newer) is rejected rather than opened against an unknown layout — there is no migration (unreleased software). The repo targets Node ≥ 24 (the root `engines` field), which includes the stable `node:sqlite` module. The database opens with `foreign_keys = ON` (so `ON DELETE CASCADE` drops a session's events with its row) and the configured `journal_mode` (default `wal`; pick a rollback-journal mode like `delete` on filesystems where WAL's shared-memory files do not work, e.g. network mounts). The table-layout version is stored in `PRAGMA user_version` and checked on open: a fresh database is stamped with the current `SCHEMA_VERSION`; a database written by any other, incompatible build (a non-current `user_version`, older or newer) is rejected rather than opened against an unknown layout — there is no migration (unreleased software).
## Contract semantics over rows ## Contract semantics over rows
@@ -21,6 +21,7 @@ The repo targets Node ≥ 24 (the root `engines` field), which includes the stab
```ts ```ts
interface Config { interface Config {
path: string // SQLite database file path, or ':memory:' for an in-process DB path: string // SQLite database file path, or ':memory:' for an in-process DB
journalMode?: 'wal' | 'delete' | 'truncate' | 'persist' // journal_mode pragma; default 'wal'
} }
``` ```

View File

@@ -28,7 +28,7 @@ import {
} from '@deepseek-ai/dsh-session-persistence' } from '@deepseek-ai/dsh-session-persistence'
import type { Session, SessionEvent, SurfaceEventType, SessionId, SessionHeader } from '@deepseek-ai/dsh-session' import type { Session, SessionEvent, SurfaceEventType, SessionId, SessionHeader } from '@deepseek-ai/dsh-session'
import { import {
openDatabase, rowToMeta, scanRows, type EventRow, type SessionRow, type JournalMode, openDatabase, rowToMeta, scanRows, type EventRow, type SessionRow,
} from './schema.ts' } from './schema.ts'
export { SCHEMA_VERSION } from './schema.ts' export { SCHEMA_VERSION } from './schema.ts'
@@ -54,6 +54,13 @@ export interface Config {
* dirs) on construction. * dirs) on construction.
*/ */
path: string path: string
/**
* SQLite `journal_mode` pragma. `wal` (the default) is the recorded
* durability model; pick a rollback-journal mode (`delete`/`truncate`/
* `persist`) on filesystems where WAL's shared-memory files do not work
* (network mounts). See {@link JournalMode}.
*/
journalMode?: JournalMode
} }
/** /**
@@ -66,6 +73,7 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
static Config: z<Config> = z.object({ static Config: z<Config> = z.object({
path: z.string().required(), path: z.string().required(),
journalMode: z.union(['wal', 'delete', 'truncate', 'persist'] as const).default('wal'),
}) })
/** /**
@@ -83,18 +91,19 @@ export class SessionPersistenceSqlite extends SessionPersistence implements Pers
super(ctx) super(ctx)
// Open the database asynchronously (the parent directory may need creating); // Open the database asynchronously (the parent directory may need creating);
// every hook awaits `ready` first. Opening synchronously would force a sync // every hook awaits `ready` first. Opening synchronously would force a sync
// mkdir and block plugin apply. // mkdir and block plugin apply. schemastery (static Config) has already
this.ready = this.openDb(config.path) // filled `journalMode`; the cast records that runtime fact.
this.ready = this.openDb(config.path, (config as Required<Config>).journalMode)
this.coordinator = new PersistenceCoordinator<number>(this.ctx, this) this.coordinator = new PersistenceCoordinator<number>(this.ctx, this)
} }
private async openDb(path: string): Promise<void> { private async openDb(path: string, journalMode: JournalMode): Promise<void> {
if (path !== ':memory:') { if (path !== ':memory:') {
const abs = resolve(path) const abs = resolve(path)
await mkdir(dirname(abs), { recursive: true, mode: 0o700 }) await mkdir(dirname(abs), { recursive: true, mode: 0o700 })
this.db = openDatabase(abs) this.db = openDatabase(abs, journalMode)
} else { } else {
this.db = openDatabase(path) this.db = openDatabase(path, journalMode)
} }
} }

View File

@@ -45,10 +45,21 @@ export interface EventRow {
surface_op: string | null surface_op: string | null
} }
/**
* Journal modes the backend will run under. `wal` is the default and the
* durability model the persistence ADR records; the rollback-journal modes
* (`delete`/`truncate`/`persist`) exist for filesystems where WAL's
* shared-memory files do not work (network mounts). `memory`/`off` are
* excluded: dropping journal durability silently contradicts what this
* backend promises.
*/
export type JournalMode = 'wal' | 'delete' | 'truncate' | 'persist'
/** /**
* Open the database at `path` and apply the schema + pragmas. `foreign_keys` * Open the database at `path` and apply the schema + pragmas. `foreign_keys`
* makes `ON DELETE CASCADE` drop a session's events with its row; `journal_mode * makes `ON DELETE CASCADE` drop a session's events with its row; the
* = WAL` matches the durability model the ADR records (the row shape maps 1:1 * `journal_mode` pragma is set from the plugin's `journalMode` config (`wal`
* default — the durability model the ADR records; the row shape maps 1:1
* onto `SessionEvent`; opencode runs this exact shape on SQLite/WAL). * onto `SessionEvent`; opencode runs this exact shape on SQLite/WAL).
* *
* The table-layout version is persisted in SQLite's `PRAGMA user_version` and * The table-layout version is persisted in SQLite's `PRAGMA user_version` and
@@ -66,10 +77,12 @@ export interface EventRow {
* makes the version check reject both sibling v3 databases instead of opening * makes the version check reject both sibling v3 databases instead of opening
* one against columns it does not have. * one against columns it does not have.
*/ */
export function openDatabase(path: string): DatabaseSync { export function openDatabase(path: string, journalMode: JournalMode): DatabaseSync {
const db = new DatabaseSync(path) const db = new DatabaseSync(path)
db.exec('PRAGMA foreign_keys = ON') db.exec('PRAGMA foreign_keys = ON')
db.exec('PRAGMA journal_mode = WAL') // journalMode is a closed in-code union (validated by the plugin Config), not
// user-controlled SQL — safe to interpolate (PRAGMA takes no bound params).
db.exec(`PRAGMA journal_mode = ${journalMode.toUpperCase()}`)
// `PRAGMA user_version` always returns exactly one row { user_version }. // `PRAGMA user_version` always returns exactly one row { user_version }.
const { user_version: onDisk } = db.prepare('PRAGMA user_version').get() as { user_version: number } const { user_version: onDisk } = db.prepare('PRAGMA user_version').get() as { user_version: number }
if (onDisk !== 0 && onDisk !== SCHEMA_VERSION) { if (onDisk !== 0 && onDisk !== SCHEMA_VERSION) {

View File

@@ -1,5 +1,6 @@
import { afterEach, describe, expect, it } from 'vitest' import { afterEach, describe, expect, it } from 'vitest'
import { Context } from 'cordis' import { Context } from 'cordis'
import { existsSync } from 'node:fs'
import { mkdtemp, rm } from 'node:fs/promises' import { mkdtemp, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os' import { tmpdir } from 'node:os'
import { join } from 'node:path' import { join } from 'node:path'
@@ -53,7 +54,7 @@ runCoordinatorContract('sqlite', async (): Promise<CoordinatorFixture> => {
// A row past the committed region whose `data` does not parse: scanRows // A row past the committed region whose `data` does not parse: scanRows
// bounds the preserved prefix at it and returns its seq as tornFrom, which // bounds the preserved prefix at it and returns its seq as tornFrom, which
// the backend surfaces to the coordinator as the tornMarker to delete from. // the backend surfaces to the coordinator as the tornMarker to delete from.
const db = openDatabase(path) const db = openDatabase(path, 'wal')
const next = (db.prepare('SELECT COALESCE(MAX(seq), -1) + 1 AS n FROM events WHERE session_id = ?') const next = (db.prepare('SELECT COALESCE(MAX(seq), -1) + 1 AS n FROM events WHERE session_id = ?')
.get(id) as { n: number }).n .get(id) as { n: number }).n
db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, ?, ?, ?, ?)') db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, ?, ?, ?, ?)')
@@ -192,7 +193,7 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
await b1.ctx.sessionPersistence.append(m.id, oneTurnLog()) // seqs 0..5 await b1.ctx.sessionPersistence.append(m.id, oneTurnLog()) // seqs 0..5
await b1.dispose() await b1.dispose()
// Hand-write an interrupted turn (turn/start seq 6, no turn/end). // Hand-write an interrupted turn (turn/start seq 6, no turn/end).
const db = openDatabase(path) const db = openDatabase(path, 'wal')
db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, 6, ?, 7, ?)') db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, 6, ?, 7, ?)')
.run(m.id, 'turn/start', JSON.stringify({ turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } })) .run(m.id, 'turn/start', JSON.stringify({ turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } }))
db.close() db.close()
@@ -204,7 +205,7 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
expect(loaded.events.at(-1)!.type).toBe('turn/end') expect(loaded.events.at(-1)!.type).toBe('turn/end')
// load() is mutating: the synthetic turn/end MUST be on disk so the stored log // load() is mutating: the synthetic turn/end MUST be on disk so the stored log
// is balanced and the cursor is truthful (contract: load closes, not defers). // is balanced and the cursor is truthful (contract: load closes, not defers).
const probe = openDatabase(path) const probe = openDatabase(path, 'wal')
const stored = probe.prepare('SELECT seq, type FROM events WHERE session_id = ? ORDER BY seq').all(m.id) as { seq: number; type: string }[] const stored = probe.prepare('SELECT seq, type FROM events WHERE session_id = ? ORDER BY seq').all(m.id) as { seq: number; type: string }[]
probe.close() probe.close()
expect(stored.map(r => r.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7]) expect(stored.map(r => r.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
@@ -237,21 +238,21 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
it('rejects opening a database whose schema version is not the current build (newer OR older)', async () => { it('rejects opening a database whose schema version is not the current build (newer OR older)', async () => {
const path = await freshDbPath() const path = await freshDbPath()
openDatabase(path).close() // stamp user_version = SCHEMA_VERSION openDatabase(path, 'wal').close() // stamp user_version = SCHEMA_VERSION
// Bump user_version past what this build supports. // Bump user_version past what this build supports.
const dbNewer = openDatabase(path) const dbNewer = openDatabase(path, 'wal')
dbNewer.exec(`PRAGMA user_version = ${SCHEMA_VERSION + 1}`) dbNewer.exec(`PRAGMA user_version = ${SCHEMA_VERSION + 1}`)
dbNewer.close() dbNewer.close()
expect(() => openDatabase(path)).toThrow(/incompatible with this build/) expect(() => openDatabase(path, 'wal')).toThrow(/incompatible with this build/)
// A stale OLDER version (e.g. a pre-summary-drop v1 DB) is also rejected — // A stale OLDER version (e.g. a pre-summary-drop v1 DB) is also rejected —
// we do not migrate (unreleased software, no backward-compat). // we do not migrate (unreleased software, no backward-compat).
const olderPath = await freshDbPath() const olderPath = await freshDbPath()
openDatabase(olderPath).close() openDatabase(olderPath, 'wal').close()
const dbOlder = openDatabase(olderPath) const dbOlder = openDatabase(olderPath, 'wal')
dbOlder.exec('PRAGMA user_version = 1') dbOlder.exec('PRAGMA user_version = 1')
dbOlder.close() dbOlder.close()
expect(() => openDatabase(olderPath)).toThrow(/incompatible with this build/) expect(() => openDatabase(olderPath, 'wal')).toThrow(/incompatible with this build/)
}) })
it('rejects a sibling v3 database (the merge-collided version) rather than opening it against missing columns', async () => { it('rejects a sibling v3 database (the merge-collided version) rather than opening it against missing columns', async () => {
@@ -261,11 +262,11 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
// of this build's columns, so it MUST be rejected, not opened. Stamp a v3 // of this build's columns, so it MUST be rejected, not opened. Stamp a v3
// database and confirm the version check refuses it. // database and confirm the version check refuses it.
const path = await freshDbPath() const path = await freshDbPath()
openDatabase(path).close() // creates + stamps user_version = SCHEMA_VERSION (4) openDatabase(path, 'wal').close() // creates + stamps user_version = SCHEMA_VERSION (4)
const db = openDatabase(path) const db = openDatabase(path, 'wal')
db.exec('PRAGMA user_version = 3') db.exec('PRAGMA user_version = 3')
db.close() db.close()
expect(() => openDatabase(path)).toThrow(/schema version 3, incompatible with this build/) expect(() => openDatabase(path, 'wal')).toThrow(/schema version 3, incompatible with this build/)
}) })
it('a corrupt-JSON row in the uncommitted tail is discarded on load, not unloadable', async () => { it('a corrupt-JSON row in the uncommitted tail is discarded on load, not unloadable', async () => {
@@ -281,7 +282,7 @@ describe('SessionPersistenceSqlite: durability and crash semantics', () => {
// unloadable; a torn tail must be discarded. scanRows finds the last // unloadable; a torn tail must be discarded. scanRows finds the last
// turn/end on the seq+type columns (never parsing tail `data`), so the // turn/end on the seq+type columns (never parsing tail `data`), so the
// unparsable row after it bounds the preserved prefix and is deleted by load. // unparsable row after it bounds the preserved prefix and is deleted by load.
const db = openDatabase(path) const db = openDatabase(path, 'wal')
db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, 6, ?, 7, ?)') db.prepare('INSERT INTO events (session_id, seq, type, time, data) VALUES (?, 6, ?, 7, ?)')
.run(m.id, 'turn/start', '{not valid json') .run(m.id, 'turn/start', '{not valid json')
db.close() db.close()
@@ -370,6 +371,29 @@ describe('SessionPersistenceSqlite: edge cases', () => {
await b2.dispose() await b2.dispose()
}) })
it('journalMode config reaches the database (default wal, rollback modes selectable)', async () => {
// :memory: databases always report journal_mode=memory, so probe file DBs.
const walPath = await freshDbPath()
const bWal = await backend(walPath)
await bWal.ctx.sessionPersistence.create(meta('jm-wal'))
expect((openDatabase(walPath, 'wal').prepare('PRAGMA journal_mode').get() as { journal_mode: string }).journal_mode).toBe('wal')
await bWal.dispose()
const deletePath = await freshDbPath()
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(SessionPersistenceSqlite, { path: deletePath, journalMode: 'delete' })
await ctx.sessionPersistence.create(meta('jm-delete'))
// Probe through a second connection: journal_mode=delete is a per-database
// property only insofar as no WAL files exist — assert the world, not the
// backend's self-report (no -wal sidecar after writes in delete mode).
const db = openDatabase(deletePath, 'delete')
expect((db.prepare('PRAGMA journal_mode').get() as { journal_mode: string }).journal_mode).toBe('delete')
db.close()
expect(existsSync(`${deletePath}-wal`)).toBe(false)
await fiber.dispose()
})
it('HMR: a DIFFERENT session colliding with a materialized on-disk id is rejected', async () => { it('HMR: a DIFFERENT session colliding with a materialized on-disk id is rejected', async () => {
const path = await freshDbPath() const path = await freshDbPath()
// Instance 1 materializes a session and disposes. // Instance 1 materializes a session and disposes.

View File

@@ -25,6 +25,8 @@ Unlike the in-process backends, the child does NOT share this cordis context —
| `cwd` | string | parent cwd | Working directory for the child process and its ACP session. | | `cwd` | string | parent cwd | Working directory for the child process and its ACP session. |
| `permission` | `'allow' \| 'reject'` | `reject` | How to auto-answer the child's `session/request_permission` prompts. `reject` declines every prompt (answer `cancelled`); `allow` approves via the first allow-shaped option. The first cut surfaces no prompt to a human. | | `permission` | `'allow' \| 'reject'` | `reject` | How to auto-answer the child's `session/request_permission` prompts. `reject` declines every prompt (answer `cancelled`); `allow` approves via the first allow-shaped option. The first cut surfaces no prompt to a human. |
| `env` | Record<string,string> | `{}` | Extra env vars for the child (e.g. its own `DEEPSEEK_API_KEY`). Forwarded on top of a credential-scrubbed copy of the parent env, so an explicit key reaches the child while ambient secrets do not leak implicitly. | | `env` | Record<string,string> | `{}` | Extra env vars for the child (e.g. its own `DEEPSEEK_API_KEY`). Forwarded on top of a credential-scrubbed copy of the parent env, so an explicit key reaches the child while ambient secrets do not leak implicitly. |
| `disposeEofGraceMs` | number | `6000` | Dispose ladder tier 1: how long the child gets to quiesce on its own after stdin EOF (flush persistence, tear down its nested subprocesses) before SIGTERM. |
| `disposeGraceMs` | number | `3000` | Dispose ladder tier 2: grace between SIGTERM and the SIGKILL escalation. |
```yaml ```yaml
- id: subagent-acp - id: subagent-acp

View File

@@ -21,7 +21,7 @@
import type { Context } from 'cordis' import type { Context } from 'cordis'
import z from 'schemastery' import z from 'schemastery'
import type { SubagentCapabilities, SubagentProvider, SubagentStartRequest } from '@deepseek-ai/dsh-subagent' import type { SubagentCapabilities, SubagentProvider, SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
import { type AcpRunSpec, type PermissionPolicy, startAcpRun } from './run.ts' import { type AcpRunSpec, DEFAULT_DISPOSE_EOF_GRACE_MS, DEFAULT_DISPOSE_GRACE_MS, type PermissionPolicy, startAcpRun } from './run.ts'
export const name = 'subagent-acp' export const name = 'subagent-acp'
export const inject = ['subagents'] export const inject = ['subagents']
@@ -52,6 +52,14 @@ export interface Config {
* ambient secrets do not leak implicitly. * ambient secrets do not leak implicitly.
*/ */
env: Record<string, string> env: Record<string, string>
/**
* Grace period (ms) for the child's EOF-driven quiesce on dispose — its
* window to flush persistence and tear down its own nested subprocesses
* before the parent escalates to a signal.
*/
disposeEofGraceMs?: number
/** Grace period (ms) between `SIGTERM` and the `SIGKILL` escalation on dispose. */
disposeGraceMs?: number
} }
export const Config: z<Config> = z.object({ export const Config: z<Config> = z.object({
@@ -61,8 +69,20 @@ export const Config: z<Config> = z.object({
cwd: z.string(), cwd: z.string(),
permission: z.union(['allow', 'reject'] as const).default('reject'), permission: z.union(['allow', 'reject'] as const).default('reject'),
env: z.dict(z.string()).default({}), env: z.dict(z.string()).default({}),
disposeEofGraceMs: z.number().default(DEFAULT_DISPOSE_EOF_GRACE_MS),
disposeGraceMs: z.number().default(DEFAULT_DISPOSE_GRACE_MS),
}) })
/** A dispose grace must be a positive finite number (it bounds the teardown wait). */
function assertPositiveFinite(name: string, value: number): void {
if (!Number.isFinite(value) || value <= 0) {
throw new Error(`subagent-acp: ${name} must be a positive finite number`)
}
}
/** The shape after schemastery applied the defaults (cwd has none). */
type ResolvedConfig = Required<Omit<Config, 'cwd'>> & Pick<Config, 'cwd'>
/** /**
* The ACP provider. Advertises NO start-time capabilities: an out-of-process * The ACP provider. Advertises NO start-time capabilities: an out-of-process
* child cannot honor `outputSchema`/`maxDepth`/`toolFilter` (the service rejects * child cannot honor `outputSchema`/`maxDepth`/`toolFilter` (the service rejects
@@ -71,7 +91,7 @@ export const Config: z<Config> = z.object({
class AcpProvider implements SubagentProvider { class AcpProvider implements SubagentProvider {
readonly capabilities: SubagentCapabilities = { outputSchema: false, depthLimit: false, toolFilter: false } readonly capabilities: SubagentCapabilities = { outputSchema: false, depthLimit: false, toolFilter: false }
constructor(readonly name: string, private readonly ctx: Context, private readonly config: Config) {} constructor(readonly name: string, private readonly ctx: Context, private readonly config: ResolvedConfig) {}
start(request: SubagentStartRequest) { start(request: SubagentStartRequest) {
const spec: AcpRunSpec = { const spec: AcpRunSpec = {
@@ -80,6 +100,8 @@ class AcpProvider implements SubagentProvider {
cwd: this.config.cwd ?? process.cwd(), cwd: this.config.cwd ?? process.cwd(),
permission: this.config.permission, permission: this.config.permission,
env: this.config.env, env: this.config.env,
disposeEofGraceMs: this.config.disposeEofGraceMs,
disposeGraceMs: this.config.disposeGraceMs,
onError: (error, stopReason) => { onError: (error, stopReason) => {
// The seam forbids `result` rejecting, so a child-level failure is // The seam forbids `result` rejecting, so a child-level failure is
// flattened to a stop reason — preserve it here rather than losing it. // flattened to a stop reason — preserve it here rather than losing it.
@@ -91,5 +113,9 @@ class AcpProvider implements SubagentProvider {
} }
export function apply(ctx: Context, config: Config): void { export function apply(ctx: Context, config: Config): void {
ctx.subagents.registerProvider(new AcpProvider(config.providerName, ctx, config)) // schemastery (Config) has already filled every defaulted field.
const resolved = config as ResolvedConfig
assertPositiveFinite('disposeEofGraceMs', resolved.disposeEofGraceMs)
assertPositiveFinite('disposeGraceMs', resolved.disposeGraceMs)
ctx.subagents.registerProvider(new AcpProvider(resolved.providerName, ctx, resolved))
} }

View File

@@ -73,16 +73,16 @@ export interface AcpRunSpec {
/** /**
* Grace period (ms) for the child's EOF-driven quiesce in * Grace period (ms) for the child's EOF-driven quiesce in
* {@link SubagentRun.dispose} — the window to flush persistence and tear down * {@link SubagentRun.dispose} — the window to flush persistence and tear down
* its OWN nested subprocesses before the parent escalates to a signal. Defaults * its OWN nested subprocesses before the parent escalates to a signal. The
* to {@link DEFAULT_DISPOSE_EOF_GRACE_MS}; a test injects a small value. * plugin fills this from its `disposeEofGraceMs` config.
*/ */
disposeEofGraceMs?: number disposeEofGraceMs: number
/** /**
* Grace period (ms) between `SIGTERM` and the `SIGKILL` escalation in * Grace period (ms) between `SIGTERM` and the `SIGKILL` escalation in
* {@link SubagentRun.dispose}. Defaults to {@link DEFAULT_DISPOSE_GRACE_MS}; * {@link SubagentRun.dispose}. The plugin fills this from its
* a test injects a small value to exercise the escalation without a long wait. * `disposeGraceMs` config.
*/ */
disposeGraceMs?: number disposeGraceMs: number
/** /**
* Sink for a child-level failure that the run flattened into a stop reason * Sink for a child-level failure that the run flattened into a stop reason
* (the seam contract forbids `result` rejecting). The driver calls this with * (the seam contract forbids `result` rejecting). The driver calls this with
@@ -94,19 +94,20 @@ export interface AcpRunSpec {
} }
/** /**
* Default grace for the child's EOF-driven quiesce on dispose — the window for it * Default grace for the child's EOF-driven quiesce on dispose (the
* to flush persistence and tear down its OWN nested subprocesses (which may run * `disposeEofGraceMs` config) — the window for it to flush persistence and tear
* their own `SIGTERM`→`SIGKILL` escalation) before the parent escalates to a * down its OWN nested subprocesses (which may run their own `SIGTERM`→`SIGKILL`
* signal. Deliberately LARGER than {@link DEFAULT_DISPOSE_GRACE_MS}: a cooperative * escalation) before the parent escalates to a signal. Deliberately LARGER than
* child whose teardown is itself waiting on a signal-trapping grandchild (e.g. a * {@link DEFAULT_DISPOSE_GRACE_MS}: a cooperative child whose teardown is itself
* bash subprocess in its own ~3s SIGTERM→SIGKILL grace) plus a final flush needs * waiting on a signal-trapping grandchild (e.g. a bash subprocess in its own ~3s
* MORE than a single signal-grace of headroom, or the parent's SIGTERM cuts it off * SIGTERM→SIGKILL grace) plus a final flush needs MORE than a single
* exactly as it reaches its own SIGKILL+flush. The child is an arbitrary ACP agent, * signal-grace of headroom, or the parent's SIGTERM cuts it off exactly as it
* so this is a standalone generous default, NOT derived from any child's internals. * reaches its own SIGKILL+flush. The child is an arbitrary ACP agent, so this is
* a standalone generous default, NOT derived from any child's internals.
*/ */
export const DEFAULT_DISPOSE_EOF_GRACE_MS = 6_000 export const DEFAULT_DISPOSE_EOF_GRACE_MS = 6_000
/** Default grace between SIGTERM and SIGKILL on dispose (mirrors the bash executor). */ /** Default grace between SIGTERM and SIGKILL on dispose (the `disposeGraceMs` config; mirrors the bash executor). */
export const DEFAULT_DISPOSE_GRACE_MS = 3_000 export const DEFAULT_DISPOSE_GRACE_MS = 3_000
/** /**
@@ -372,8 +373,8 @@ export function startAcpRun(request: SubagentStartRequest, spec: AcpRunSpec): Su
// Reach quiescence, not merely request it (dispose must AWAIT the child // Reach quiescence, not merely request it (dispose must AWAIT the child
// actually stopping). If the child is already gone, nothing to do. // actually stopping). If the child is already gone, nothing to do.
if (child.exitCode !== null || child.signalCode !== null) return if (child.exitCode !== null || child.signalCode !== null) return
const eofGraceMs = spec.disposeEofGraceMs ?? DEFAULT_DISPOSE_EOF_GRACE_MS const eofGraceMs = spec.disposeEofGraceMs
const graceMs = spec.disposeGraceMs ?? DEFAULT_DISPOSE_GRACE_MS const graceMs = spec.disposeGraceMs
// 1. Graceful: end the ACP request stream (stdin EOF) and let the child // 1. Graceful: end the ACP request stream (stdin EOF) and let the child
// quiesce ON ITS OWN. Our acp-agent has NO SIGTERM handler in a normal // quiesce ON ITS OWN. Our acp-agent has NO SIGTERM handler in a normal
// session — it tears down via the server bridge's connection-close path // session — it tears down via the server bridge's connection-close path

View File

@@ -8,7 +8,7 @@ import { fileURLToPath } from 'node:url'
import SubagentService from '@deepseek-ai/dsh-subagent' import SubagentService from '@deepseek-ai/dsh-subagent'
import type { Agent } from '@deepseek-ai/dsh-agent' import type { Agent } from '@deepseek-ai/dsh-agent'
import * as acp from '../src/index.ts' import * as acp from '../src/index.ts'
import { acpStopReason, acpContentText, buildChildEnv, SENSITIVE_ENV_PATTERN, startAcpRun, toAcpPrompt, type AcpRunSpec } from '../src/run.ts' import { acpStopReason, acpContentText, buildChildEnv, DEFAULT_DISPOSE_EOF_GRACE_MS, DEFAULT_DISPOSE_GRACE_MS, SENSITIVE_ENV_PATTERN, startAcpRun, toAcpPrompt, type AcpRunSpec } from '../src/run.ts'
/** /**
* Keyless integration tests for the ACP subagent backend. Each spawns a REAL * Keyless integration tests for the ACP subagent backend. Each spawns a REAL
@@ -171,7 +171,7 @@ describe('dsh-subagent-acp', () => {
const run = startAcpRun( const run = startAcpRun(
{ prompt: [{ type: 'text', text: 'p' }], parent: fakeParent, signal: controller.signal }, { prompt: [{ type: 'text', text: 'p' }], parent: fakeParent, signal: controller.signal },
// `touch <sentinel>` — runs only if the process is actually spawned. // `touch <sentinel>` — runs only if the process is actually spawned.
{ command: 'touch', args: [sentinel], cwd: tmp, permission: 'reject', env: {} }, { command: 'touch', args: [sentinel], cwd: tmp, permission: 'reject', env: {}, disposeEofGraceMs: DEFAULT_DISPOSE_EOF_GRACE_MS, disposeGraceMs: DEFAULT_DISPOSE_GRACE_MS },
) )
const result = await run.result const result = await run.result
expect(result.stopReason).toBe('aborted') expect(result.stopReason).toBe('aborted')
@@ -390,7 +390,7 @@ describe('dsh-subagent-acp', () => {
// absent-sink branch). // absent-sink branch).
const run = startAcpRun( const run = startAcpRun(
{ prompt: [{ type: 'text', text: 'p' }], parent: fakeParent }, { prompt: [{ type: 'text', text: 'p' }], parent: fakeParent },
{ command: '/nonexistent/acp-agent-binary', args: [], cwd: process.cwd(), permission: 'reject', env: {} }, { command: '/nonexistent/acp-agent-binary', args: [], cwd: process.cwd(), permission: 'reject', env: {}, disposeEofGraceMs: DEFAULT_DISPOSE_EOF_GRACE_MS, disposeGraceMs: DEFAULT_DISPOSE_GRACE_MS },
) )
const result = await run.result const result = await run.result
// The seam contract: a child-level failure resolves error, never rejects. // The seam contract: a child-level failure resolves error, never rejects.
@@ -398,6 +398,16 @@ describe('dsh-subagent-acp', () => {
await run.dispose() await run.dispose()
}) })
it('rejects a non-positive dispose grace at load', async () => {
for (const bad of [{ disposeEofGraceMs: 0 }, { disposeGraceMs: -1 }, { disposeEofGraceMs: Number.NaN }]) {
const ctx = new Context()
await ctx.plugin(SubagentService)
await expect(ctx.plugin(acp, { providerName: 'acp', command: 'true', args: [], permission: 'reject', env: {}, ...bad }))
.rejects.toThrow(/subagent-acp: dispose(?:Eof)?GraceMs must be a positive finite number/)
await ctx.fiber.dispose()
}
})
it('resolves error via the provider (real load path) when the command does not exist', async () => { it('resolves error via the provider (real load path) when the command does not exist', async () => {
const ctx = new Context() const ctx = new Context()
await ctx.plugin(SubagentService) await ctx.plugin(SubagentService)
@@ -428,6 +438,8 @@ describe('dsh-subagent-acp', () => {
cwd: process.cwd(), cwd: process.cwd(),
permission: 'reject', permission: 'reject',
env: {}, env: {},
disposeEofGraceMs: DEFAULT_DISPOSE_EOF_GRACE_MS,
disposeGraceMs: DEFAULT_DISPOSE_GRACE_MS,
onError: (error, stopReason) => { errors.push({ message: error.message, stopReason }) }, onError: (error, stopReason) => { errors.push({ message: error.message, stopReason }) },
}, },
) )

View File

@@ -17,6 +17,7 @@ Each tool is registered independently; a product that wants only one disables th
|---|---|---| |---|---|---|
| `search` | `true` | Register `web_search`. | | `search` | `true` | Register `web_search`. |
| `fetch` | `true` | Register `web_fetch`. | | `fetch` | `true` | Register `web_fetch`. |
| `searchMaxResults` | `8` | Upper bound on sources returned by one `web_search` call (the seam truncates a longer provider list and flags it). |
```yaml ```yaml
- id: tool-web - id: tool-web

View File

@@ -20,7 +20,7 @@
import type { Context } from 'cordis' import type { Context } from 'cordis'
import z from 'schemastery' import z from 'schemastery'
import type {} from '@deepseek-ai/dsh-web' import type {} from '@deepseek-ai/dsh-web'
import { applyWebSearchTool } from './search.ts' import { applyWebSearchTool, WEB_SEARCH_MAX_RESULTS } from './search.ts'
import { applyWebFetchTool } from './fetch.ts' import { applyWebFetchTool } from './fetch.ts'
export { WEB_SEARCH_MAX_RESULTS, applyWebSearchTool, formatSearchOutput, parseSearchArgs, presentSearchCall } from './search.ts' export { WEB_SEARCH_MAX_RESULTS, applyWebSearchTool, formatSearchOutput, parseSearchArgs, presentSearchCall } from './search.ts'
@@ -38,13 +38,26 @@ export interface Config {
search?: boolean search?: boolean
/** Register `web_fetch`. Defaults to true. */ /** Register `web_fetch`. Defaults to true. */
fetch?: boolean fetch?: boolean
/** Upper bound on sources returned by one `web_search` call. */
searchMaxResults?: number
} }
export const Config: z<Config> = z.object({ export const Config: z<Config> = z.object({
search: z.boolean().default(true), search: z.boolean().default(true),
fetch: z.boolean().default(true), fetch: z.boolean().default(true),
searchMaxResults: z.number().default(WEB_SEARCH_MAX_RESULTS),
}) })
/** The shape after schemastery applies its defaults to every field. */
type ResolvedConfig = Required<Config>
/** The result cap must be a positive integer (it bounds a provider's source list). */
function assertPositiveInteger(name: string, value: number): void {
if (!Number.isInteger(value) || value < 1) {
throw new Error(`tool-web: ${name} must be a positive integer`)
}
}
/** /**
* Register the enabled web tools. `search`/`fetch` default to true; a product * Register the enabled web tools. `search`/`fetch` default to true; a product
* that wants only one disables the other in config. The tools' disposers are * that wants only one disables the other in config. The tools' disposers are
@@ -52,6 +65,9 @@ export const Config: z<Config> = z.object({
* teardown is needed. * teardown is needed.
*/ */
export function apply(ctx: Context, config: Config): void { export function apply(ctx: Context, config: Config): void {
if (config.search !== false) applyWebSearchTool(ctx) // schemastery (Config) has already filled every defaulted field.
if (config.fetch !== false) applyWebFetchTool(ctx) const resolved = config as ResolvedConfig
assertPositiveInteger('searchMaxResults', resolved.searchMaxResults)
if (resolved.search) applyWebSearchTool(ctx, resolved.searchMaxResults)
if (resolved.fetch) applyWebFetchTool(ctx)
} }

View File

@@ -13,10 +13,10 @@ import type { WebSearchResult } from '@deepseek-ai/dsh-web'
import type {} from '@deepseek-ai/dsh-system-prompt' import type {} from '@deepseek-ai/dsh-system-prompt'
/** /**
* Default upper bound on returned sources. Owned by the consumer (not the * Default upper bound on returned sources (the `searchMaxResults` config).
* provider or model), mirroring `dsh-tool-fs`'s `READ_LIMIT`/`GREP_LIMIT`. The * Owned by the consumer (not the provider or model), mirroring `dsh-tool-fs`'s
* model just asks a question; the product controls how much context returns. * `READ_LIMIT`. The model just asks a question; the product controls how much
* The default `8` aligns with OpenCode's Exa default. * context returns. The default `8` aligns with OpenCode's Exa default.
*/ */
export const WEB_SEARCH_MAX_RESULTS = 8 export const WEB_SEARCH_MAX_RESULTS = 8
@@ -67,8 +67,8 @@ export function presentSearchCall(args: { query: string }): GenericCallView {
return { card: 'generic', title: args.query, kind: 'search', rawInput: args.query } return { card: 'generic', title: args.query, kind: 'search', rawInput: args.query }
} }
/** Register the `web_search` tool and its system-prompt guidance. */ /** Register the `web_search` tool and its system-prompt guidance. `maxResults` is the deployment's source cap. */
export function applyWebSearchTool(ctx: Context): void { export function applyWebSearchTool(ctx: Context, maxResults: number): void {
ctx.systemPrompt.section({ ctx.systemPrompt.section({
name: 'tool:web_search', name: 'tool:web_search',
order: 110, order: 110,
@@ -84,7 +84,7 @@ export function applyWebSearchTool(ctx: Context): void {
async execute(args, exec): Promise<ContentBlock[]> { async execute(args, exec): Promise<ContentBlock[]> {
const input = parseSearchArgs(args) const input = parseSearchArgs(args)
const result = await ctx.web.search( const result = await ctx.web.search(
{ query: input.query, maxResults: WEB_SEARCH_MAX_RESULTS }, { query: input.query, maxResults },
exec.signal ? { signal: exec.signal } : undefined, exec.signal ? { signal: exec.signal } : undefined,
) )
return [{ type: 'text', text: formatSearchOutput(result) }] return [{ type: 'text', text: formatSearchOutput(result) }]

View File

@@ -15,6 +15,7 @@ import {
presentFetchCall, presentFetchCall,
renderBody, renderBody,
htmlToMarkdown, htmlToMarkdown,
WEB_SEARCH_MAX_RESULTS,
} from '@deepseek-ai/dsh-tool-web' } from '@deepseek-ai/dsh-tool-web'
const available: WebProviderStatus = { available: true } const available: WebProviderStatus = { available: true }
@@ -279,3 +280,48 @@ describe('tool-web execution through the real registry', () => {
await fiber.dispose() await fiber.dispose()
}) })
}) })
describe('searchMaxResults is plugin config', () => {
it('forwards the default cap to the seam when unconfigured', async () => {
const seen: { maxResults?: number | undefined } = {}
const provider: WebSearchProvider = {
id: 'stub-search',
status: () => available,
search: (request) => { seen.maxResults = request.maxResults; return Promise.resolve({ providerId: 'stub-search', query: 'q', sources: [], truncated: false }) },
}
const { fiber, call } = await mountTools({ webConfig: { searchProvider: 'stub-search' }, search: provider })
await call('web_search', { query: 'q' })
expect(seen.maxResults).toBe(WEB_SEARCH_MAX_RESULTS)
await fiber.dispose()
})
it('forwards a configured cap to the seam, which enforces it', async () => {
const sources = Array.from({ length: 5 }, (_, i) => ({ url: `https://s${i}.test` }))
const provider: WebSearchProvider = {
id: 'stub-search',
status: () => available,
search: request => Promise.resolve({ providerId: 'stub-search', query: request.query, sources, truncated: false }),
}
const { fiber, call } = await mountTools({ config: { searchMaxResults: 2 }, webConfig: { searchProvider: 'stub-search' }, search: provider })
const out = await call('web_search', { query: 'q' })
expect(out.isError).toBe(false)
const body = out.content.map(b => b.text).join('')
expect(body).toContain('https://s1.test')
expect(body).not.toContain('https://s2.test')
expect(body).toContain('Showing the first 2 sources.')
await fiber.dispose()
})
it.each([
['zero', 0],
['negative', -3],
['fractional', 1.5],
])('rejects a %s searchMaxResults at load', async (_label, value) => {
const ctx = new Context()
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
await ctx.plugin(WebService, {})
await expect(ctx.plugin(ToolWeb, { searchMaxResults: value }))
.rejects.toThrow(/tool-web: searchMaxResults must be a positive integer/)
})
})

3
pnpm-lock.yaml generated
View File

@@ -323,6 +323,9 @@ importers:
diff: diff:
specifier: ^9.0.0 specifier: ^9.0.0
version: 9.0.0 version: 9.0.0
schemastery:
specifier: ^3.18.0
version: 3.18.0
devDependencies: devDependencies:
'@deepseek-ai/dsh-agent': '@deepseek-ai/dsh-agent':
specifier: workspace:^ specifier: workspace:^