Add E2B PTY, LSP, and code runtime providers
This commit is contained in:
@@ -2,5 +2,5 @@
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write packages/pty/README.md
|
||||
README.md: a4f743056b4a524be9623b0f700f37e0534b463f
|
||||
README.zh.md: c84ad3f1b59afcdbbd111f1b82c57c56aa24fdcf
|
||||
README.md: 239c751f30f8e5c543e6737e65a67f9bb22f79ad
|
||||
README.zh.md: edaae7f948d957e2fce2b3f331cc2af61cfabd36
|
||||
|
||||
@@ -7,7 +7,8 @@ English | [中文](README.zh.md)
|
||||
| Package | Role | ctx key |
|
||||
|---|---|---|
|
||||
| [`pty`](pty/README.md) (`@deepseek-ai/dsh-pty`) | Backend registry, branded ids, exact-Agent ownership, session operations, and awaited cleanup | `ctx.pty` |
|
||||
| `pty-local` (`@deepseek-ai/dsh-pty-local`) | Shell backend over `ctx.subprocess.spawnTerminal`: readiness detection, bounded terminal state, sandbox policy, and session operations | registers on `ctx.pty` |
|
||||
| [`pty-local`](pty-local/README.md) (`@deepseek-ai/dsh-pty-local`) | Local `node-pty` backend, readiness detection, bounded terminal state, sandboxing, and process-session supervision | registers on `ctx.pty` |
|
||||
| [`pty-e2b`](pty-e2b/README.md) (`@deepseek-ai/dsh-pty-e2b`) | E2B byte-PTY backend, remote foreground signaling, bounded terminal state, and awaited remote cleanup | registers on `ctx.pty` |
|
||||
| `tool-pty` (`@deepseek-ai/dsh-tool-pty`) | Six model-facing tools and generic task integration for background sends | registers on `ctx.tools` |
|
||||
|
||||
The design and deferred boundaries live in the [persistent PTY Agent Note](../../.agents/notes/implemented/feature/2026-07-16-persistent-pty-sessions.md).
|
||||
The core design lives in the [persistent PTY Agent Note](../../.agents/notes/implemented/feature/2026-07-16-persistent-pty-sessions.md); the remote ownership boundary lives in the [E2B extension note](../../.agents/notes/implemented/feature/2026-07-28-e2b-interactive-semantic-code-runtime-poc.md).
|
||||
|
||||
@@ -7,7 +7,8 @@
|
||||
| 包 | 职责 | ctx 键 |
|
||||
|---|---|---|
|
||||
| [`pty`](pty/README.md)(`@deepseek-ai/dsh-pty`) | 后端注册表、品牌化 id、精确的 Agent 所有权、会话操作与等待完成的清理 | `ctx.pty` |
|
||||
| `pty-local`(`@deepseek-ai/dsh-pty-local`) | `ctx.subprocess.spawnTerminal` 之上的 shell 后端:就绪检测、有界终端状态、沙箱策略与会话操作 | 注册到 `ctx.pty` |
|
||||
| [`pty-local`](pty-local/README.md)(`@deepseek-ai/dsh-pty-local`) | 本地 `node-pty` 后端、就绪检测、有界终端状态、沙箱与进程会话监管 | 注册到 `ctx.pty` |
|
||||
| [`pty-e2b`](pty-e2b/README.md)(`@deepseek-ai/dsh-pty-e2b`) | E2B 字节 PTY 后端、远程前台信号传递、有界终端状态与等待完成的远程清理 | 注册到 `ctx.pty` |
|
||||
| `tool-pty`(`@deepseek-ai/dsh-tool-pty`) | 6 个面向模型的工具,并为后台发送集成通用任务 | 注册到 `ctx.tools` |
|
||||
|
||||
设计与暂缓边界记录在[持久 PTY Agent Note](../../.agents/notes/implemented/feature/2026-07-16-persistent-pty-sessions.md) 中。
|
||||
核心设计记录在[持久 PTY Agent Note](../../.agents/notes/implemented/feature/2026-07-16-persistent-pty-sessions.md) 中;远程所有权边界记录在 [E2B 扩展 Agent Note](../../.agents/notes/implemented/feature/2026-07-28-e2b-interactive-semantic-code-runtime-poc.md) 中。
|
||||
|
||||
6
packages/pty/pty-e2b/README.i18n.yaml
Normal file
6
packages/pty/pty-e2b/README.i18n.yaml
Normal file
@@ -0,0 +1,6 @@
|
||||
# Bilingual-pair consistency record (docs/i18n/README.md): the git blob hash of each
|
||||
# side as of the last confirmed-consistent state. Both languages carry equal authority;
|
||||
# after editing either side, bring the other along and re-record with:
|
||||
# pnpm run verify-translation-pairing --write packages/pty/pty-e2b/README.md
|
||||
README.md: 7a8ddc9e08f4d28d69b97825232b00e2d8f394b7
|
||||
README.zh.md: 2151ca1f2d91b01a7b0106153f708ace2c3d2726
|
||||
54
packages/pty/pty-e2b/README.md
Normal file
54
packages/pty/pty-e2b/README.md
Normal file
@@ -0,0 +1,54 @@
|
||||
# @deepseek-ai/dsh-pty-e2b
|
||||
|
||||
English | [中文](README.zh.md)
|
||||
|
||||
E2B byte-PTY backend for [`ctx.pty`](../pty/README.md). It creates persistent interactive shells inside the shared `ctx.e2b` sandbox while the PTY registry keeps session identity, exact-Agent ownership, and cleanup policy on the host.
|
||||
|
||||
## Plugin and configuration
|
||||
|
||||
The `pty-e2b` plugin injects `e2b` and `pty`, then registers one backend under `backendType`.
|
||||
|
||||
| Key | Default | Meaning |
|
||||
|---|---|---|
|
||||
| `backendType` | `shell` | Registry type selected by `terminal_open`. |
|
||||
| `rows` / `cols` | `40` / `160` | Initial remote PTY size. |
|
||||
| `scrollbackLines` | `10000` | Maximum retained logical lines. |
|
||||
| `scrollbackMaxBytes` | `4194304` | Maximum retained UTF-8 scrollback bytes. |
|
||||
| `maxReadBytes` | `262144` | Maximum bytes returned by one read or settled send. |
|
||||
| `pollIntervalMs` | `50` | Host readiness-poll interval. |
|
||||
| `idleSilenceMs` | `3000` | Output silence that yields `inferred_idle`. |
|
||||
| `timeoutMs` | `30000` | Absolute startup and send wait bound. |
|
||||
| `disposeGraceMs` | `3000` | TERM-to-KILL cleanup grace. |
|
||||
|
||||
Numeric values are positive safe integers, `backendType` is non-empty, and `maxReadBytes` cannot exceed `scrollbackMaxBytes`. A relative spawn cwd resolves against `ctx.e2b.cwd`; an absolute remote path remains absolute.
|
||||
|
||||
## Runtime contract
|
||||
|
||||
The backend uses E2B's byte-oriented PTY callback with a streaming fatal UTF-8 decoder, then the backend-neutral line sanitizer and bounded buffers from `dsh-pty`. It installs a controlled Bash prompt marker and waits for printable prompt text; when that marker is unavailable, observed output plus the configured silence bound yields `inferred_idle`. Startup with no output reaches the absolute timeout and fails instead of publishing an empty session.
|
||||
|
||||
Each send writes UTF-8 bytes and an optional carriage-return submit sequence. Cancellation and explicit signals resolve the remote terminal's foreground process group through `ps`, then signal that group; `SIGKILL` refuses to target the shell itself. Close sends `SIGTERM` to the PTY process group, waits, escalates through E2B's PTY kill, and does not resolve until the SDK handle reports exit. A startup failure closes the unpublished PTY, and `PtyBackendCleanupError` preserves a concurrent cleanup failure.
|
||||
|
||||
The remote PTY process and its child processes live in E2B. Prompt/readiness state, scrollback, operation handles, owner authority, and SDK event delivery remain in host memory.
|
||||
|
||||
## Model Experience
|
||||
|
||||
### Indirect consumer
|
||||
|
||||
#### What the model sees
|
||||
|
||||
Nothing directly. Through `@deepseek-ai/dsh-tool-pty`, the model may receive bounded MOTD, send deltas, scrollback pages, readiness reasons, signal results, and cleanup failures.
|
||||
|
||||
#### Token effect
|
||||
|
||||
None until a consumer returns bounded backend output. Retained host PTY scrollback is not placed in model history by this package.
|
||||
|
||||
#### KV Cache effect
|
||||
|
||||
No direct invalidation; the consumer owns prompts, schemas, and appended results.
|
||||
|
||||
## Known Limitations and Deferred Work
|
||||
|
||||
- **Line-oriented terminal model** — CSI/OSC control sequences are removed; alternate-screen and full terminal emulation remain unsupported.
|
||||
- **Readiness is marker-or-silence based** — E2B exposes foreground process groups but not the local backend's Linux syscall inspection, so `inferred_idle` is deliberately possible.
|
||||
- **UTF-8 only** — invalid byte sequences fail the session instead of returning lossy text.
|
||||
- **No reconnectable terminal handles** — retaining an E2B sandbox preserves remote files, not host ownership, buffers, callbacks, or live PTY sessions.
|
||||
54
packages/pty/pty-e2b/README.zh.md
Normal file
54
packages/pty/pty-e2b/README.zh.md
Normal file
@@ -0,0 +1,54 @@
|
||||
# @deepseek-ai/dsh-pty-e2b
|
||||
|
||||
[English](README.md) | 中文
|
||||
|
||||
用于 [`ctx.pty`](../pty/README.md) 的 E2B 字节 PTY 后端。它在共享的 `ctx.e2b` 沙箱内创建持久交互式 shell;PTY 注册表则在宿主侧维护会话身份、精确的 Agent 所有权和清理策略。
|
||||
|
||||
## 插件与配置
|
||||
|
||||
`pty-e2b` 插件注入 `e2b` 和 `pty`,然后以 `backendType` 注册一个后端。
|
||||
|
||||
| 配置键 | 默认值 | 含义 |
|
||||
|---|---|---|
|
||||
| `backendType` | `shell` | `terminal_open` 选择的注册表类型。 |
|
||||
| `rows` / `cols` | `40` / `160` | 远程 PTY 的初始尺寸。 |
|
||||
| `scrollbackLines` | `10000` | 保留的逻辑行数上限。 |
|
||||
| `scrollbackMaxBytes` | `4194304` | 保留的 UTF-8 scrollback 字节数上限。 |
|
||||
| `maxReadBytes` | `262144` | 单次读取或发送结算时返回的字节数上限。 |
|
||||
| `pollIntervalMs` | `50` | 宿主就绪轮询间隔。 |
|
||||
| `idleSilenceMs` | `3000` | 触发 `inferred_idle` 的输出静默时长。 |
|
||||
| `timeoutMs` | `30000` | 启动与发送等待的绝对上限。 |
|
||||
| `disposeGraceMs` | `3000` | TERM 到 KILL 的清理宽限期。 |
|
||||
|
||||
数值必须是正的安全整数,`backendType` 必须非空,且 `maxReadBytes` 不得超过 `scrollbackMaxBytes`。相对的 spawn cwd 以 `ctx.e2b.cwd` 为基准解析;绝对远程路径保持不变。
|
||||
|
||||
## 运行时契约
|
||||
|
||||
该后端为 E2B 面向字节的 PTY 回调配备流式、遇到无效序列即失败的 UTF-8 解码器,随后使用 `dsh-pty` 提供的后端无关行清理器与有界缓冲区。它会安装受控的 Bash 提示符标记,并等待可打印的提示符文本;若该标记不可用,系统会在已经观察到输出且达到已配置的静默上限时得出 `inferred_idle`。零输出的启动过程会达到绝对超时并失败,不会发布空会话。
|
||||
|
||||
每次发送都会写入 UTF-8 字节,并可选写入回车提交序列。取消与显式信号会通过 `ps` 确定远程终端的前台进程组,再向该组发送信号;发送 `SIGKILL` 时拒绝以 shell 本身为目标。关闭操作向 PTY 进程组发送 `SIGTERM`,等待后通过 E2B 的 PTY kill 操作升级,并且直到 SDK 句柄报告退出才结算。如果启动失败,系统会关闭尚未发布的 PTY;若清理同时失败,`PtyBackendCleanupError` 会保留这项失败。
|
||||
|
||||
远程 PTY 进程及其子进程位于 E2B。提示符/就绪状态、scrollback、操作句柄、所有者权限和 SDK 事件交付仍保留在宿主内存中。
|
||||
|
||||
## 模型体验
|
||||
|
||||
### 间接消费方
|
||||
|
||||
#### 模型看到的内容
|
||||
|
||||
没有直接可见内容。模型通过 `@deepseek-ai/dsh-tool-pty` 可能收到有界的 MOTD、发送增量、scrollback 页、就绪原因、信号结果和清理失败。
|
||||
|
||||
#### Token 影响
|
||||
|
||||
消费方返回有界的后端输出前没有影响。本包不会把宿主保留的 PTY scrollback 放入模型历史。
|
||||
|
||||
#### KV Cache 影响
|
||||
|
||||
不会直接失效;提示词、schema 和追加结果由消费方负责。
|
||||
|
||||
## 已知限制与暂缓工作
|
||||
|
||||
- **面向行的终端模型**:CSI/OSC 控制序列会被移除;备用屏幕与完整终端仿真仍不受支持。
|
||||
- **就绪判断基于标记或静默**:E2B 会公开前台进程组,但不提供本地后端使用的 Linux syscall 检查,因此系统有意保留返回 `inferred_idle` 的可能性。
|
||||
- **仅支持 UTF-8**:无效字节序列会使会话失败,而不是返回有损文本。
|
||||
- **没有可重连的终端句柄**:保留 E2B 沙箱会保留远程文件,但不会保留宿主所有权、缓冲区、回调或实时 PTY 会话。
|
||||
46
packages/pty/pty-e2b/package.json
Normal file
46
packages/pty/pty-e2b/package.json
Normal file
@@ -0,0 +1,46 @@
|
||||
{
|
||||
"name": "@deepseek-ai/dsh-pty-e2b",
|
||||
"description": "E2B PTY provider for DeepSeek Harness",
|
||||
"version": "0.0.1",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "lib/index.js",
|
||||
"types": "lib/types/index.d.ts",
|
||||
"exports": {
|
||||
".": {
|
||||
"types": "./lib/types/index.d.ts",
|
||||
"default": "./lib/index.js"
|
||||
},
|
||||
"./invariant": {
|
||||
"types": "./lib/types/invariant.d.ts",
|
||||
"default": "./lib/invariant.js"
|
||||
},
|
||||
"./src/*": "./src/*",
|
||||
"./package.json": "./package.json"
|
||||
},
|
||||
"files": [
|
||||
"lib/index.js",
|
||||
"lib/invariant.js",
|
||||
"lib/types/**/*.d.ts",
|
||||
"lib/types/**/*.d.ts.map",
|
||||
"src"
|
||||
],
|
||||
"license": "BSD-3-Clause",
|
||||
"peerDependencies": {
|
||||
"@deepseek-ai/dsh-e2b": "^0.0.1",
|
||||
"@deepseek-ai/dsh-invariants": "^0.0.1",
|
||||
"@deepseek-ai/dsh-pty": "^0.0.1",
|
||||
"cordis": "^4.0.0-rc.7"
|
||||
},
|
||||
"dependencies": {
|
||||
"schemastery": "^3.18.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@deepseek-ai/dsh-agent": "workspace:^",
|
||||
"@deepseek-ai/dsh-e2b": "workspace:^",
|
||||
"@deepseek-ai/dsh-invariants": "workspace:^",
|
||||
"@deepseek-ai/dsh-pty": "workspace:^",
|
||||
"@deepseek-ai/dsh-session": "workspace:^",
|
||||
"cordis": "^4.0.0-rc.7"
|
||||
}
|
||||
}
|
||||
64
packages/pty/pty-e2b/src/config.ts
Normal file
64
packages/pty/pty-e2b/src/config.ts
Normal file
@@ -0,0 +1,64 @@
|
||||
/** Validated configuration for the E2B PTY backend. */
|
||||
|
||||
import z from 'schemastery'
|
||||
|
||||
/** Public plugin configuration. */
|
||||
export interface Config {
|
||||
/** Backend registry type. */
|
||||
backendType?: string
|
||||
/** Initial terminal rows. */
|
||||
rows?: number
|
||||
/** Initial terminal columns. */
|
||||
cols?: number
|
||||
/** Maximum retained logical lines. */
|
||||
scrollbackLines?: number
|
||||
/** Maximum retained UTF-8 bytes. */
|
||||
scrollbackMaxBytes?: number
|
||||
/** Maximum bytes returned by one read or settled viewport. */
|
||||
maxReadBytes?: number
|
||||
/** Readiness polling interval. */
|
||||
pollIntervalMs?: number
|
||||
/** Output silence duration that yields `inferred_idle`. */
|
||||
idleSilenceMs?: number
|
||||
/** Absolute send and startup wait bound. */
|
||||
timeoutMs?: number
|
||||
/** Grace before PTY teardown escalates from TERM to KILL. */
|
||||
disposeGraceMs?: number
|
||||
}
|
||||
|
||||
/** Configuration after Schemastery defaults. */
|
||||
export type ResolvedConfig = Required<Config>
|
||||
|
||||
/* jscpd:ignore-start -- Loader requires a backend-local schema and load-time diagnostics. */
|
||||
/** Schemastery config exposed by the plugin. */
|
||||
export const Config: z<Config> = z.object({
|
||||
backendType: z.string().default('shell'),
|
||||
rows: z.number().default(40),
|
||||
cols: z.number().default(160),
|
||||
scrollbackLines: z.number().default(10_000),
|
||||
scrollbackMaxBytes: z.number().default(4 * 1024 * 1024),
|
||||
maxReadBytes: z.number().default(256 * 1024),
|
||||
pollIntervalMs: z.number().default(50),
|
||||
idleSilenceMs: z.number().default(3_000),
|
||||
timeoutMs: z.number().default(30_000),
|
||||
disposeGraceMs: z.number().default(3_000),
|
||||
})
|
||||
|
||||
/**
|
||||
* Validate the resolved configuration before publishing the backend.
|
||||
* @param config - Schemastery-resolved plugin configuration.
|
||||
* @returns Nothing; success narrows every optional field to its resolved value.
|
||||
*/
|
||||
export function validateConfig(config: Config): asserts config is ResolvedConfig {
|
||||
const resolved = config as ResolvedConfig
|
||||
if (resolved.backendType.length === 0) throw new Error('pty-e2b: backendType must be non-empty')
|
||||
for (const [name, value] of Object.entries(resolved)) {
|
||||
if (typeof value === 'number' && (!Number.isSafeInteger(value) || value <= 0)) {
|
||||
throw new Error(`pty-e2b: ${name} must be a positive safe integer`)
|
||||
}
|
||||
}
|
||||
if (resolved.maxReadBytes > resolved.scrollbackMaxBytes) {
|
||||
throw new Error('pty-e2b: maxReadBytes must not exceed scrollbackMaxBytes')
|
||||
}
|
||||
}
|
||||
/* jscpd:ignore-end */
|
||||
93
packages/pty/pty-e2b/src/index.ts
Normal file
93
packages/pty/pty-e2b/src/index.ts
Normal file
@@ -0,0 +1,93 @@
|
||||
/** E2B byte-PTY backend for persistent interactive terminal sessions. */
|
||||
|
||||
import { posix } from 'node:path'
|
||||
import type { Context } from 'cordis'
|
||||
import type { CommandHandle, Sandbox } from '@deepseek-ai/dsh-e2b'
|
||||
import { PtyBackendCleanupError } from '@deepseek-ai/dsh-pty'
|
||||
import type { PtyBackend, PtyBackendSpawnSpec } from '@deepseek-ai/dsh-pty'
|
||||
import { type Config, type ResolvedConfig, validateConfig } from './config.ts'
|
||||
import { E2BPtySession } from './session.ts'
|
||||
|
||||
export { Config } from './config.ts'
|
||||
export type { Config as PtyE2BConfig } from './config.ts'
|
||||
export { E2BPtySession } from './session.ts'
|
||||
|
||||
/** Cordis plugin name. */
|
||||
export const name = 'pty-e2b'
|
||||
/** Required shared sandbox owner and PTY registry. */
|
||||
export const inject = ['e2b', 'pty']
|
||||
|
||||
function terminalEnvironment(spec: PtyBackendSpawnSpec): Record<string, string> {
|
||||
return {
|
||||
TERM: 'dumb',
|
||||
PAGER: 'cat',
|
||||
GIT_PAGER: 'cat',
|
||||
PS1: 'dsh> ',
|
||||
PROMPT_COMMAND: 'printf "\\033]133;D;%s\\007" "$?"',
|
||||
BASH_SILENCE_DEPRECATION_WARNING: '1',
|
||||
DSH_SHELL: '1',
|
||||
DSH_SESSION_ID: spec.owner.id,
|
||||
DSH_PTY_SESSION_ID: spec.sessionId,
|
||||
}
|
||||
}
|
||||
|
||||
/** E2B backend registered under the configured terminal type. */
|
||||
export class E2BPtyBackend implements PtyBackend {
|
||||
readonly type: string
|
||||
|
||||
constructor(
|
||||
private readonly ctx: Context,
|
||||
private readonly config: ResolvedConfig,
|
||||
private readonly createPty: (
|
||||
sandbox: Sandbox,
|
||||
options: Parameters<Sandbox['pty']['create']>[0],
|
||||
) => Promise<CommandHandle> = (sandbox, options) => sandbox.pty.create(options),
|
||||
) {
|
||||
this.type = config.backendType
|
||||
}
|
||||
|
||||
/** Create, initialize, and publish one remote PTY session. */
|
||||
async spawn(spec: PtyBackendSpawnSpec): Promise<E2BPtySession> {
|
||||
spec.signal?.throwIfAborted()
|
||||
const sandbox = await this.ctx.e2b.getSandbox()
|
||||
spec.signal?.throwIfAborted()
|
||||
const pending: Uint8Array[] = []
|
||||
const created: { session?: E2BPtySession } = {}
|
||||
const handle = await this.createPty(sandbox, {
|
||||
rows: this.config.rows,
|
||||
cols: this.config.cols,
|
||||
cwd: posix.resolve(this.ctx.e2b.cwd, spec.cwd ?? this.ctx.e2b.cwd),
|
||||
envs: terminalEnvironment(spec),
|
||||
timeoutMs: 0,
|
||||
...spec.signal === undefined ? {} : { signal: spec.signal },
|
||||
onData: (data) => {
|
||||
if (created.session === undefined) pending.push(Uint8Array.from(data))
|
||||
else created.session.onData(data)
|
||||
},
|
||||
})
|
||||
if (!Number.isSafeInteger(handle.pid) || handle.pid <= 0) {
|
||||
await handle.kill().catch(() => false)
|
||||
throw new Error(`pty-e2b: E2B returned invalid PTY pid ${handle.pid}`)
|
||||
}
|
||||
const session = new E2BPtySession(sandbox, handle, this.config)
|
||||
created.session = session
|
||||
for (const data of pending) session.onData(data)
|
||||
try {
|
||||
await session.initialize(spec.signal)
|
||||
return session
|
||||
} catch (error: unknown) {
|
||||
try {
|
||||
await session.close('E2B PTY startup failed')
|
||||
} catch (cleanupError: unknown) {
|
||||
throw new PtyBackendCleanupError(error, cleanupError)
|
||||
}
|
||||
throw error
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Register the E2B PTY backend. */
|
||||
export function apply(ctx: Context, config: Config): void {
|
||||
validateConfig(config)
|
||||
ctx.pty.registerBackend(new E2BPtyBackend(ctx, config))
|
||||
}
|
||||
20
packages/pty/pty-e2b/src/invariant.ts
Normal file
20
packages/pty/pty-e2b/src/invariant.ts
Normal file
@@ -0,0 +1,20 @@
|
||||
/** Package-owned invariant companion for `@deepseek-ai/dsh-pty-e2b`. */
|
||||
|
||||
/* jscpd:ignore-start */
|
||||
import type { Context } from 'cordis'
|
||||
import type { InvariantInstaller } from '@deepseek-ai/dsh-invariants'
|
||||
|
||||
const PACKAGE_NAME = '@deepseek-ai/dsh-pty-e2b'
|
||||
|
||||
/** Cordis companion plugin name. */
|
||||
export const name = 'pty-e2b-invariant'
|
||||
/** Service required before the companion can reserve package ownership. */
|
||||
export const inject = ['invariants']
|
||||
|
||||
/** No runtime invariant: the PTY registry owns publication and cleanup. */
|
||||
const install: InvariantInstaller = () => {}
|
||||
|
||||
/** Register this package's invariant companion. */
|
||||
export const apply = (ctx: Context): Promise<() => void> =>
|
||||
Promise.resolve(ctx.invariants.register(PACKAGE_NAME, install))
|
||||
/* jscpd:ignore-end */
|
||||
367
packages/pty/pty-e2b/src/session.ts
Normal file
367
packages/pty/pty-e2b/src/session.ts
Normal file
@@ -0,0 +1,367 @@
|
||||
/** One byte-oriented E2B PTY session projected onto the harness PTY seam. */
|
||||
|
||||
import { Buffer } from 'node:buffer'
|
||||
import type { CommandHandle, Sandbox } from '@deepseek-ai/dsh-e2b'
|
||||
import { CommandExitError } from '@deepseek-ai/dsh-e2b'
|
||||
import {
|
||||
PtyTerminalSanitizer,
|
||||
PtyTextBuffer,
|
||||
ptySignalName,
|
||||
ptyUtf8Tail,
|
||||
} from '@deepseek-ai/dsh-pty'
|
||||
import type {
|
||||
PtyBackendSession,
|
||||
PtyReadRequest,
|
||||
PtyReadResult,
|
||||
PtySendOperation,
|
||||
PtySendRead,
|
||||
PtySendRequest,
|
||||
PtySendResult,
|
||||
PtySessionStatus,
|
||||
PtySignal,
|
||||
PtySignalResult,
|
||||
PtyWaitReason,
|
||||
} from '@deepseek-ai/dsh-pty'
|
||||
import type { ResolvedConfig } from './config.ts'
|
||||
|
||||
function delay(ms: number): Promise<void> {
|
||||
return new Promise(resolve => setTimeout(resolve, ms))
|
||||
}
|
||||
|
||||
/* jscpd:ignore-start -- Operation state stays backend-local because process readiness and cleanup identities diverge. */
|
||||
class E2BSendOperation implements PtySendOperation {
|
||||
private readonly output: PtyTextBuffer
|
||||
private readonly result = Promise.withResolvers<PtySendResult>()
|
||||
private finished = false
|
||||
|
||||
constructor(
|
||||
maxBytes: number,
|
||||
readonly startedAt: number,
|
||||
private readonly onCancel: () => void,
|
||||
) {
|
||||
this.output = new PtyTextBuffer(maxBytes)
|
||||
}
|
||||
|
||||
get done(): Promise<PtySendResult> {
|
||||
return this.result.promise
|
||||
}
|
||||
|
||||
append(text: string): void {
|
||||
if (!this.finished) this.output.append(text)
|
||||
}
|
||||
|
||||
settle(waitReason: PtyWaitReason, sessionStatus: PtySessionStatus, inheritedTruncation: boolean): void {
|
||||
if (this.finished) return
|
||||
this.finished = true
|
||||
const read = this.output.snapshot()
|
||||
this.result.resolve({
|
||||
viewport: read.text,
|
||||
waitReason,
|
||||
sessionStatus,
|
||||
truncated: read.truncated || inheritedTruncation,
|
||||
})
|
||||
}
|
||||
|
||||
fail(error: unknown): void {
|
||||
if (this.finished) return
|
||||
this.finished = true
|
||||
this.result.reject(error)
|
||||
}
|
||||
|
||||
readOutput(): PtySendRead {
|
||||
return this.output.consume()
|
||||
}
|
||||
|
||||
cancel(): boolean {
|
||||
if (this.finished) return false
|
||||
this.onCancel()
|
||||
return true
|
||||
}
|
||||
}
|
||||
/* jscpd:ignore-end */
|
||||
|
||||
/** Live session around one E2B SDK PTY handle. */
|
||||
export class E2BPtySession implements PtyBackendSession {
|
||||
motd = ''
|
||||
readonly pid: number
|
||||
private readonly decoder = new TextDecoder('utf-8', { fatal: true })
|
||||
private readonly sanitizer: PtyTerminalSanitizer
|
||||
private readonly scrollback: PtyTextBuffer
|
||||
private readonly exited = Promise.withResolvers<void>()
|
||||
private statusValue: PtySessionStatus = { kind: 'running' }
|
||||
private active: E2BSendOperation | undefined
|
||||
private activeTimer: NodeJS.Timeout | undefined
|
||||
private activeAbort: (() => void) | undefined
|
||||
private promptSeen = false
|
||||
private promptTextSeen = false
|
||||
private initializing = false
|
||||
private lastOutputAt = Date.now()
|
||||
private closing = false
|
||||
private closePromise: Promise<void> | undefined
|
||||
private closeSignal: NodeJS.Signals | null = null
|
||||
private transportFailure: Error | undefined
|
||||
private remoteExited = false
|
||||
|
||||
constructor(
|
||||
private readonly sandbox: Sandbox,
|
||||
private readonly handle: CommandHandle,
|
||||
private readonly config: ResolvedConfig,
|
||||
) {
|
||||
this.pid = handle.pid
|
||||
this.sanitizer = new PtyTerminalSanitizer(config.maxReadBytes)
|
||||
this.scrollback = new PtyTextBuffer(config.scrollbackMaxBytes, config.scrollbackLines)
|
||||
const completion = handle.wait()
|
||||
void completion.then(
|
||||
(result) => { this.onExit(result.exitCode) },
|
||||
(error: unknown) => {
|
||||
if (error instanceof CommandExitError) this.onExit(error.exitCode)
|
||||
else this.onTransportFailure(error)
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* Consume bytes received by the SDK's PTY callback.
|
||||
* @param data - Exact callback bytes in delivery order.
|
||||
*/
|
||||
onData(data: Uint8Array): void {
|
||||
let decoded: string
|
||||
try {
|
||||
decoded = this.decoder.decode(data, { stream: true })
|
||||
} catch (error: unknown) {
|
||||
this.onTransportFailure(new Error('pty-e2b: PTY emitted invalid UTF-8', { cause: error }))
|
||||
return
|
||||
}
|
||||
const sanitized = this.sanitizer.push(decoded)
|
||||
this.appendOutput(sanitized.text)
|
||||
if (sanitized.prompt) {
|
||||
this.promptSeen = true
|
||||
this.promptTextSeen = sanitized.promptText === true
|
||||
this.lastOutputAt = Date.now()
|
||||
} else if (this.promptSeen && sanitized.promptText === true) {
|
||||
this.promptTextSeen = true
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Await the first prompt or bounded startup fallback.
|
||||
* @param signal - Optional startup cancellation signal.
|
||||
*/
|
||||
async initialize(signal?: AbortSignal): Promise<void> {
|
||||
this.initializing = true
|
||||
try {
|
||||
const operation = this.startSend({ text: '', submit: false, ...signal === undefined ? {} : { signal } })
|
||||
const result = await operation.done
|
||||
if (result.waitReason === 'session_exit') throw new Error('E2B PTY shell exited during startup')
|
||||
if (result.waitReason === 'timeout') throw new Error('E2B PTY shell did not reach readiness before startup timeout')
|
||||
this.motd = result.viewport
|
||||
} catch (error: unknown) {
|
||||
signal?.throwIfAborted()
|
||||
throw error
|
||||
} finally {
|
||||
this.initializing = false
|
||||
}
|
||||
}
|
||||
|
||||
/* jscpd:ignore-start -- PTY backends share request admission while owning distinct input and readiness transports. */
|
||||
startSend(request: PtySendRequest): PtySendOperation {
|
||||
if (this.closing) throw new Error('E2B PTY session is closing')
|
||||
if (this.statusValue.kind === 'exited') throw new Error('E2B PTY session has exited')
|
||||
if (this.active !== undefined) throw new Error('E2B PTY session already has an active send')
|
||||
if (request.signal?.aborted === true) throw new Error('E2B PTY send aborted before write')
|
||||
|
||||
const operation = new E2BSendOperation(
|
||||
this.config.maxReadBytes,
|
||||
Date.now(),
|
||||
() => { this.interrupt(operation) },
|
||||
)
|
||||
this.active = operation
|
||||
this.lastOutputAt = Date.now()
|
||||
this.promptSeen = false
|
||||
this.promptTextSeen = false
|
||||
if (request.signal !== undefined) {
|
||||
const onAbort = (): void => { operation.cancel() }
|
||||
request.signal.addEventListener('abort', onAbort, { once: true })
|
||||
this.activeAbort = () => request.signal?.removeEventListener('abort', onAbort)
|
||||
}
|
||||
|
||||
const input = `${request.text}${request.submit ? '\r' : ''}`
|
||||
if (input.length > 0) {
|
||||
void this.sandbox.pty.sendInput(this.pid, Buffer.from(input)).catch((error: unknown) => {
|
||||
if (this.active === operation) this.failActive(error)
|
||||
})
|
||||
}
|
||||
this.activeTimer = setInterval(() => { this.pollReadiness(operation) }, this.config.pollIntervalMs)
|
||||
return operation
|
||||
}
|
||||
/* jscpd:ignore-end */
|
||||
|
||||
/* jscpd:ignore-start -- The seam requires identical bounded-read coordinates across backend buffers. */
|
||||
read(request: PtyReadRequest): PtyReadResult {
|
||||
const snapshot = this.scrollback.snapshot()
|
||||
const lines = snapshot.text.split('\n')
|
||||
const totalLines = snapshot.text.length === 0 ? 0 : lines.length
|
||||
const offset = request.offset ?? 0
|
||||
const count = request.count ?? 500
|
||||
if (!Number.isSafeInteger(offset) || offset < 0) throw new Error('PTY read offset must be a non-negative safe integer')
|
||||
if (!Number.isSafeInteger(count) || count <= 0) throw new Error('PTY read count must be a positive safe integer')
|
||||
if (offset >= totalLines) {
|
||||
return { text: '', totalLines, lineBegin: offset, lineEnd: offset, truncated: snapshot.truncated }
|
||||
}
|
||||
const end = totalLines - offset
|
||||
const start = Math.max(0, end - count)
|
||||
const bounded = ptyUtf8Tail(lines.slice(start, end).join('\n'), this.config.maxReadBytes)
|
||||
const returnedLines = bounded.text.length === 0 ? 0 : bounded.text.split('\n').length
|
||||
return {
|
||||
text: bounded.text,
|
||||
totalLines,
|
||||
lineBegin: offset,
|
||||
lineEnd: offset + returnedLines,
|
||||
truncated: snapshot.truncated || bounded.truncated,
|
||||
}
|
||||
}
|
||||
/* jscpd:ignore-end */
|
||||
|
||||
/* jscpd:ignore-start -- Signal, status, and close methods preserve the seam shape around remote identities. */
|
||||
async signal(signal: PtySignal): Promise<PtySignalResult> {
|
||||
const pgid = await this.foregroundPgid()
|
||||
if (signal === 'SIGKILL' && pgid === this.pid) {
|
||||
throw new Error('refusing to SIGKILL the E2B PTY shell; use terminal_close')
|
||||
}
|
||||
await this.sandbox.commands.run(`kill -${signal.slice(3)} -- -${pgid}`)
|
||||
return { delivered: true, targetPgid: pgid }
|
||||
}
|
||||
|
||||
status(): PtySessionStatus {
|
||||
return this.statusValue
|
||||
}
|
||||
|
||||
close(reason: string): Promise<void> {
|
||||
this.closing = true
|
||||
if (this.closePromise !== undefined) return this.closePromise
|
||||
const closing = this.closeOnce(reason).catch((error: unknown) => {
|
||||
this.closePromise = undefined
|
||||
this.failActive(error)
|
||||
throw error
|
||||
})
|
||||
this.closePromise = closing
|
||||
return closing
|
||||
}
|
||||
/* jscpd:ignore-end */
|
||||
|
||||
private appendOutput(text: string): void {
|
||||
if (text.length === 0) return
|
||||
this.lastOutputAt = Date.now()
|
||||
this.scrollback.append(text)
|
||||
this.active?.append(text)
|
||||
}
|
||||
|
||||
private pollReadiness(operation: E2BSendOperation): void {
|
||||
if (this.active !== operation) return
|
||||
if (this.statusValue.kind === 'exited') {
|
||||
this.settleActive('session_exit')
|
||||
return
|
||||
}
|
||||
const elapsed = Date.now() - operation.startedAt
|
||||
const idleFor = Date.now() - this.lastOutputAt
|
||||
if (this.promptSeen && this.promptTextSeen && idleFor >= this.config.pollIntervalMs) {
|
||||
this.settleActive('stdin_read')
|
||||
return
|
||||
}
|
||||
const startupHasOutput = !this.initializing || this.scrollback.snapshot().text.length > 0
|
||||
if (startupHasOutput && idleFor >= this.config.idleSilenceMs) {
|
||||
this.settleActive('inferred_idle')
|
||||
return
|
||||
}
|
||||
if (elapsed >= this.config.timeoutMs) this.settleActive('timeout')
|
||||
}
|
||||
|
||||
private settleActive(waitReason: PtyWaitReason): void {
|
||||
const operation = this.active
|
||||
if (operation === undefined) return
|
||||
const inherited = this.scrollback.snapshot().truncated
|
||||
this.clearActive()
|
||||
operation.settle(waitReason, this.statusValue, inherited)
|
||||
}
|
||||
|
||||
private clearActive(): void {
|
||||
if (this.activeTimer !== undefined) clearInterval(this.activeTimer)
|
||||
this.activeTimer = undefined
|
||||
this.activeAbort?.()
|
||||
this.activeAbort = undefined
|
||||
this.active = undefined
|
||||
}
|
||||
|
||||
private failActive(error: unknown): void {
|
||||
const operation = this.active
|
||||
if (operation === undefined) return
|
||||
this.clearActive()
|
||||
operation.fail(error)
|
||||
}
|
||||
|
||||
private interrupt(operation: E2BSendOperation): void {
|
||||
if (this.active !== operation) return
|
||||
void this.signal('SIGINT').catch((error: unknown) => { this.failActive(error) })
|
||||
}
|
||||
|
||||
private async foregroundPgid(): Promise<number> {
|
||||
const result = await this.sandbox.commands.run(`ps -o tpgid= -p ${this.pid}`)
|
||||
const raw = result.stdout.trim()
|
||||
const pgid = Number(raw)
|
||||
if (!/^[1-9][0-9]*$/.test(raw) || !Number.isSafeInteger(pgid)) {
|
||||
throw new Error(`cannot resolve foreground process group for E2B PTY ${this.pid}`)
|
||||
}
|
||||
return pgid
|
||||
}
|
||||
|
||||
private onExit(exitCode: number): void {
|
||||
this.remoteExited = true
|
||||
let tail = ''
|
||||
try {
|
||||
tail = this.decoder.decode()
|
||||
} catch (error: unknown) {
|
||||
this.transportFailure ??= new Error('pty-e2b: PTY ended with invalid UTF-8', { cause: error })
|
||||
}
|
||||
this.appendOutput(this.sanitizer.push(tail).text)
|
||||
this.appendOutput(this.sanitizer.flush())
|
||||
const inferredSignal = this.closeSignal ?? (exitCode > 128 ? ptySignalName(exitCode - 128) : null)
|
||||
this.statusValue = {
|
||||
kind: 'exited',
|
||||
exitCode: inferredSignal === null ? exitCode : null,
|
||||
signal: inferredSignal,
|
||||
}
|
||||
if (this.transportFailure === undefined) this.settleActive('session_exit')
|
||||
else this.failActive(this.transportFailure)
|
||||
this.exited.resolve()
|
||||
}
|
||||
|
||||
private onTransportFailure(error: unknown): void {
|
||||
const failure = error instanceof Error ? error : new Error(String(error))
|
||||
this.transportFailure ??= failure
|
||||
this.statusValue = { kind: 'exited', exitCode: null, signal: null }
|
||||
this.failActive(failure)
|
||||
}
|
||||
|
||||
private async closeOnce(reason: string): Promise<void> {
|
||||
if (!this.remoteExited) {
|
||||
this.closeSignal = 'SIGTERM'
|
||||
try {
|
||||
await this.sandbox.commands.run(`kill -TERM -- -${this.pid}`)
|
||||
} catch (error: unknown) {
|
||||
if (!(error instanceof CommandExitError)) throw error
|
||||
}
|
||||
await Promise.race([this.exited.promise, delay(this.config.disposeGraceMs)])
|
||||
}
|
||||
if (!this.remoteExited) {
|
||||
this.closeSignal = 'SIGKILL'
|
||||
await this.sandbox.pty.kill(this.pid)
|
||||
await Promise.race([this.exited.promise, delay(this.config.disposeGraceMs)])
|
||||
}
|
||||
if (!this.remoteExited) {
|
||||
throw new Error(`E2B PTY cleanup failed (${reason}); surviving pid: ${this.pid}`)
|
||||
}
|
||||
this.settleActive('session_exit')
|
||||
await this.handle.disconnect().catch(() => {})
|
||||
if (this.transportFailure !== undefined) throw this.transportFailure
|
||||
}
|
||||
}
|
||||
195
packages/pty/pty-e2b/tests/index.spec.ts
Normal file
195
packages/pty/pty-e2b/tests/index.spec.ts
Normal file
@@ -0,0 +1,195 @@
|
||||
import { Context } from 'cordis'
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { CommandHandle, Sandbox } from '@deepseek-ai/dsh-e2b'
|
||||
import type E2BSandboxService from '@deepseek-ai/dsh-e2b'
|
||||
import PtyService, { PtyBackendCleanupError, PtySessionId } from '@deepseek-ai/dsh-pty'
|
||||
import { E2BPtyBackend, apply } from '@deepseek-ai/dsh-pty-e2b'
|
||||
import { validateConfig } from '@deepseek-ai/dsh-pty-e2b/src/config.ts'
|
||||
import * as E2BPtyInvariant from '../src/invariant.ts'
|
||||
import InvariantService from '@deepseek-ai/dsh-invariants'
|
||||
import { AgentMessageId, type Agent } from '@deepseek-ai/dsh-agent'
|
||||
import { Session, SessionId } from '@deepseek-ai/dsh-session'
|
||||
|
||||
function config() {
|
||||
return {
|
||||
backendType: 'shell', rows: 24, cols: 80,
|
||||
scrollbackLines: 10, scrollbackMaxBytes: 128, maxReadBytes: 64,
|
||||
pollIntervalMs: 1, idleSilenceMs: 2, timeoutMs: 5, disposeGraceMs: 1,
|
||||
}
|
||||
}
|
||||
|
||||
function owner(ctx: Context): Agent {
|
||||
const id = SessionId('owner')
|
||||
return {
|
||||
id, options: {}, session: new Session(id), status: 'idle', acceptsNextStep: false, ctx,
|
||||
followup: () => AgentMessageId('unused'), steer: () => AgentMessageId('unused'),
|
||||
inject: () => AgentMessageId('unused'), send: () => AgentMessageId('unused'),
|
||||
cancel() {}, whenIdle: () => Promise.resolve(),
|
||||
}
|
||||
}
|
||||
|
||||
function handle(pid = 123, kill = vi.fn().mockResolvedValue(true)): CommandHandle {
|
||||
const result = Promise.withResolvers<{ exitCode: number; stdout: string; stderr: string }>()
|
||||
return {
|
||||
pid,
|
||||
wait: () => result.promise,
|
||||
kill,
|
||||
disconnect: vi.fn().mockResolvedValue(undefined),
|
||||
} as unknown as CommandHandle
|
||||
}
|
||||
|
||||
describe('E2BPtyBackend and plugin', () => {
|
||||
it('creates a remote PTY with isolated environment and initializes the session', async () => {
|
||||
vi.useFakeTimers()
|
||||
const ctx = new Context()
|
||||
const sandbox = {} as Sandbox
|
||||
ctx.provide('e2b', {
|
||||
cwd: '/workspace',
|
||||
getSandbox: async () => sandbox,
|
||||
} as E2BSandboxService)
|
||||
const created = handle()
|
||||
let options: Parameters<Sandbox['pty']['create']>[0] | undefined
|
||||
const backend = new E2BPtyBackend(ctx, config(), async (_sandbox, received) => {
|
||||
options = received
|
||||
void received.onData(Buffer.from('banner\n'))
|
||||
setTimeout(() => { void received.onData(Buffer.from('\x1b]133;D;0\x07dsh> ')) }, 0)
|
||||
return created
|
||||
})
|
||||
const pending = backend.spawn({
|
||||
sessionId: PtySessionId('pty-1'), owner: owner(ctx), type: 'shell', cwd: 'project',
|
||||
signal: new AbortController().signal,
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(2)
|
||||
const session = await pending
|
||||
|
||||
expect(session.motd).toBe('dsh> ')
|
||||
expect(options).toMatchObject({ rows: 24, cols: 80, cwd: '/workspace/project', timeoutMs: 0 })
|
||||
expect(options?.envs).toMatchObject({
|
||||
TERM: 'dumb', PAGER: 'cat', GIT_PAGER: 'cat', PS1: 'dsh> ',
|
||||
DSH_SHELL: '1', DSH_SESSION_ID: 'owner', DSH_PTY_SESSION_ID: 'pty-1',
|
||||
})
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('uses the SDK PTY create method and the shared cwd by default', async () => {
|
||||
vi.useFakeTimers()
|
||||
const ctx = new Context()
|
||||
const created = handle()
|
||||
const create = vi.fn(async (received: Parameters<Sandbox['pty']['create']>[0]) => {
|
||||
setTimeout(() => { void received.onData(Buffer.from('\x1b]133;D;0\x07dsh> ')) }, 0)
|
||||
return created
|
||||
})
|
||||
const sandbox = { pty: { create } } as unknown as Sandbox
|
||||
ctx.provide('e2b', { cwd: '/workspace', getSandbox: async () => sandbox } as unknown as E2BSandboxService)
|
||||
const backend = new E2BPtyBackend(ctx, config())
|
||||
const pending = backend.spawn({ sessionId: PtySessionId('default'), owner: owner(ctx), type: 'shell' })
|
||||
await vi.advanceTimersByTimeAsync(2)
|
||||
await pending
|
||||
expect(create).toHaveBeenCalledWith(expect.objectContaining({ cwd: '/workspace' }))
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('rejects aborts and invalid pids, killing a malformed SDK handle', async () => {
|
||||
const ctx = new Context()
|
||||
const sandbox = {} as Sandbox
|
||||
ctx.provide('e2b', { cwd: '/workspace', getSandbox: async () => sandbox } as E2BSandboxService)
|
||||
const create = vi.fn().mockResolvedValue(handle(0))
|
||||
const backend = new E2BPtyBackend(ctx, config(), create)
|
||||
const aborted = AbortSignal.abort(new Error('stop'))
|
||||
await expect(backend.spawn({ sessionId: PtySessionId('one'), owner: owner(ctx), type: 'shell', signal: aborted })).rejects.toThrow('stop')
|
||||
expect(create).not.toHaveBeenCalled()
|
||||
|
||||
const malformedKill = vi.fn().mockResolvedValue(true)
|
||||
const malformed = handle(0, malformedKill)
|
||||
const invalid = new E2BPtyBackend(ctx, config(), async () => malformed)
|
||||
await expect(invalid.spawn({ sessionId: PtySessionId('two'), owner: owner(ctx), type: 'shell' })).rejects.toThrow('invalid PTY pid')
|
||||
expect(malformedKill).toHaveBeenCalledOnce()
|
||||
|
||||
const killFailureKill = vi.fn().mockRejectedValue(new Error('already gone'))
|
||||
const killFailure = handle(0, killFailureKill)
|
||||
const raced = new E2BPtyBackend(ctx, config(), async () => killFailure)
|
||||
await expect(raced.spawn({ sessionId: PtySessionId('three'), owner: owner(ctx), type: 'shell' })).rejects.toThrow('invalid PTY pid')
|
||||
})
|
||||
|
||||
it('cleans failed startup and aggregates a cleanup failure', async () => {
|
||||
vi.useFakeTimers()
|
||||
const ctx = new Context()
|
||||
const sandbox = {
|
||||
commands: { run: vi.fn().mockResolvedValue({ exitCode: 0, stdout: '', stderr: '' }) },
|
||||
pty: { kill: vi.fn().mockRejectedValue(new Error('cleanup failed')) },
|
||||
} as unknown as Sandbox
|
||||
ctx.provide('e2b', { cwd: '/workspace', getSandbox: async () => sandbox } as E2BSandboxService)
|
||||
const failedHandle = handle()
|
||||
const backend = new E2BPtyBackend(ctx, config(), async () => failedHandle)
|
||||
const pending = backend.spawn({ sessionId: PtySessionId('failed'), owner: owner(ctx), type: 'shell' })
|
||||
const rejected = expect(pending).rejects.toMatchObject({
|
||||
name: 'PtyBackendCleanupError',
|
||||
cleanupError: expect.objectContaining({ message: 'cleanup failed' }),
|
||||
} satisfies Partial<PtyBackendCleanupError>)
|
||||
await vi.advanceTimersByTimeAsync(6)
|
||||
await vi.advanceTimersByTimeAsync(2)
|
||||
await rejected
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('preserves startup failure when cleanup succeeds', async () => {
|
||||
vi.useFakeTimers()
|
||||
const ctx = new Context()
|
||||
const completion = Promise.withResolvers<{ exitCode: number; stdout: string; stderr: string }>()
|
||||
const created = {
|
||||
pid: 123,
|
||||
wait: () => completion.promise,
|
||||
disconnect: vi.fn().mockResolvedValue(undefined),
|
||||
} as unknown as CommandHandle
|
||||
const sandbox = {
|
||||
commands: {
|
||||
run: vi.fn(async (command: string) => {
|
||||
if (command.startsWith('kill -TERM')) completion.resolve({ exitCode: 143, stdout: '', stderr: '' })
|
||||
return { exitCode: 0, stdout: '', stderr: '' }
|
||||
}),
|
||||
},
|
||||
pty: { kill: vi.fn().mockResolvedValue(true) },
|
||||
} as unknown as Sandbox
|
||||
ctx.provide('e2b', { cwd: '/workspace', getSandbox: async () => sandbox } as unknown as E2BSandboxService)
|
||||
const backend = new E2BPtyBackend(ctx, config(), async () => created)
|
||||
const rejected = expect(backend.spawn({ sessionId: PtySessionId('failed-clean'), owner: owner(ctx), type: 'shell' }))
|
||||
.rejects.toThrow('startup timeout')
|
||||
await vi.advanceTimersByTimeAsync(6)
|
||||
await rejected
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
it('validates configuration and registers the selected backend type', async () => {
|
||||
const valid = config()
|
||||
expect(() => { validateConfig(valid) }).not.toThrow()
|
||||
for (const invalid of [
|
||||
{ ...valid, backendType: '' },
|
||||
{ ...valid, rows: 0 },
|
||||
{ ...valid, rows: 1.5 },
|
||||
{ ...valid, maxReadBytes: 129 },
|
||||
]) {
|
||||
expect(() => { validateConfig(invalid) }).toThrow()
|
||||
}
|
||||
|
||||
const registerBackend = vi.fn()
|
||||
apply({ pty: { registerBackend } } as unknown as Context, valid)
|
||||
expect(registerBackend).toHaveBeenCalledWith(expect.objectContaining({ type: 'shell' }))
|
||||
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(PtyService)
|
||||
ctx.provide('e2b', { cwd: '/workspace', getSandbox: async () => ({}) } as never)
|
||||
const fiber = await ctx.plugin({
|
||||
inject: ['pty', 'e2b'],
|
||||
apply: (pluginCtx: Context) => { apply(pluginCtx, valid) },
|
||||
})
|
||||
expect(ctx.pty.listBackends()).toEqual(['shell'])
|
||||
await fiber.dispose()
|
||||
})
|
||||
|
||||
it('registers the package-owned invariant companion', async () => {
|
||||
const ctx = new Context()
|
||||
await ctx.plugin(InvariantService, { enabled: true })
|
||||
const fiber = await ctx.plugin(E2BPtyInvariant).await()
|
||||
await fiber.dispose()
|
||||
})
|
||||
})
|
||||
398
packages/pty/pty-e2b/tests/session.spec.ts
Normal file
398
packages/pty/pty-e2b/tests/session.spec.ts
Normal file
@@ -0,0 +1,398 @@
|
||||
import { Buffer } from 'node:buffer'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
CommandExitError,
|
||||
type CommandHandle,
|
||||
type CommandResult,
|
||||
type Sandbox,
|
||||
} from '@deepseek-ai/dsh-e2b'
|
||||
import type { PtySendOperation, PtySessionStatus } from '@deepseek-ai/dsh-pty'
|
||||
import { E2BPtySession } from '@deepseek-ai/dsh-pty-e2b'
|
||||
import type { ResolvedConfig } from '@deepseek-ai/dsh-pty-e2b/src/config.ts'
|
||||
|
||||
function commandError(exitCode: number): CommandExitError {
|
||||
return new CommandExitError({ exitCode, stdout: '', stderr: '', error: `exit ${exitCode}` })
|
||||
}
|
||||
|
||||
class FakePtyHandle {
|
||||
pid = 123
|
||||
readonly result = Promise.withResolvers<CommandResult>()
|
||||
disconnects = 0
|
||||
kills = 0
|
||||
disconnectError: unknown
|
||||
private settled = false
|
||||
|
||||
wait(): Promise<CommandResult> {
|
||||
return this.result.promise
|
||||
}
|
||||
|
||||
async disconnect(): Promise<void> {
|
||||
this.disconnects += 1
|
||||
if (this.disconnectError !== undefined) throw this.disconnectError
|
||||
}
|
||||
|
||||
async kill(): Promise<boolean> {
|
||||
this.kills += 1
|
||||
return true
|
||||
}
|
||||
|
||||
exit(exitCode = 0): void {
|
||||
if (this.settled) return
|
||||
this.settled = true
|
||||
this.result.resolve({ exitCode, stdout: '', stderr: '' })
|
||||
}
|
||||
|
||||
failExit(exitCode: number): void {
|
||||
if (this.settled) return
|
||||
this.settled = true
|
||||
this.result.reject(commandError(exitCode))
|
||||
}
|
||||
|
||||
crash(error: unknown): void {
|
||||
if (this.settled) return
|
||||
this.settled = true
|
||||
this.result.reject(error)
|
||||
}
|
||||
|
||||
asHandle(): CommandHandle {
|
||||
return this as unknown as CommandHandle
|
||||
}
|
||||
}
|
||||
|
||||
class FakeSandbox {
|
||||
readonly sent: Array<{ pid: number; data: Buffer }> = []
|
||||
readonly commands: string[] = []
|
||||
readonly killed: number[] = []
|
||||
pgid = '456\n'
|
||||
sendError: unknown
|
||||
commandError: unknown
|
||||
killError: unknown
|
||||
onTerm: (() => void) | undefined
|
||||
onKill: (() => void) | undefined
|
||||
|
||||
readonly sandbox = {
|
||||
pty: {
|
||||
sendInput: async (pid: number, data: Uint8Array): Promise<void> => {
|
||||
this.sent.push({ pid, data: Buffer.from(data) })
|
||||
if (this.sendError !== undefined) throw this.sendError
|
||||
},
|
||||
kill: async (pid: number): Promise<boolean> => {
|
||||
this.killed.push(pid)
|
||||
if (this.killError !== undefined) throw this.killError
|
||||
this.onKill?.()
|
||||
return true
|
||||
},
|
||||
},
|
||||
commands: {
|
||||
run: async (command: string): Promise<CommandResult> => {
|
||||
this.commands.push(command)
|
||||
if (this.commandError !== undefined) {
|
||||
const error = this.commandError
|
||||
this.commandError = undefined
|
||||
throw error
|
||||
}
|
||||
if (command.startsWith('ps ')) return { exitCode: 0, stdout: this.pgid, stderr: '' }
|
||||
if (command.startsWith('kill -TERM')) this.onTerm?.()
|
||||
return { exitCode: 0, stdout: '', stderr: '' }
|
||||
},
|
||||
},
|
||||
} as unknown as Sandbox
|
||||
}
|
||||
|
||||
function config(overrides: Partial<ResolvedConfig> = {}): ResolvedConfig {
|
||||
return {
|
||||
backendType: 'shell', rows: 24, cols: 80,
|
||||
scrollbackLines: 10, scrollbackMaxBytes: 128, maxReadBytes: 64,
|
||||
pollIntervalMs: 10, idleSilenceMs: 40, timeoutMs: 100, disposeGraceMs: 20,
|
||||
...overrides,
|
||||
}
|
||||
}
|
||||
|
||||
async function initialize(session: E2BPtySession): Promise<void> {
|
||||
const pending = session.initialize()
|
||||
session.onData(Buffer.from('\x1b]133;D;0\x07dsh> '))
|
||||
await vi.advanceTimersByTimeAsync(10)
|
||||
await pending
|
||||
}
|
||||
|
||||
afterEach(() => { vi.useRealTimers() })
|
||||
|
||||
describe('E2BPtySession readiness, output, and signals', () => {
|
||||
it('initializes, sends UTF-8 input, settles at a prompt, and reads bounded scrollback', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fake = new FakeSandbox()
|
||||
const handle = new FakePtyHandle()
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config({ maxReadBytes: 12 }))
|
||||
expect(session.read({})).toMatchObject({ text: '', totalLines: 0 })
|
||||
await initialize(session)
|
||||
expect(session.motd).toBe('dsh> ')
|
||||
|
||||
const operation = session.startSend({ text: 'printf 你好', submit: true })
|
||||
expect(fake.sent).toEqual([{ pid: 123, data: Buffer.from('printf 你好\r') }])
|
||||
session.onData(Buffer.from('一\n二\n三\x1b]133;D;0\x07dsh> '))
|
||||
const bounded = operation.readOutput()
|
||||
expect(bounded.delta).toContain('三')
|
||||
expect(bounded.truncated).toBe(true)
|
||||
await vi.advanceTimersByTimeAsync(10)
|
||||
expect(await operation.done).toMatchObject({ waitReason: 'stdin_read', sessionStatus: { kind: 'running' } })
|
||||
expect(operation.cancel()).toBe(false)
|
||||
expect(session.read({ count: 2 }).text).toContain('dsh>')
|
||||
expect(session.read({ offset: 99 })).toMatchObject({ text: '', lineBegin: 99, lineEnd: 99 })
|
||||
expect(() => session.read({ offset: -1 })).toThrow('non-negative safe integer')
|
||||
expect(() => session.read({ offset: 1.5 })).toThrow('non-negative safe integer')
|
||||
expect(() => session.read({ count: 0 })).toThrow('positive safe integer')
|
||||
expect(() => session.read({ count: 1.5 })).toThrow('positive safe integer')
|
||||
|
||||
await expect(session.signal('SIGTERM')).resolves.toEqual({ delivered: true, targetPgid: 456 })
|
||||
expect(fake.commands).toContain('kill -TERM -- -456')
|
||||
expect(session.status()).toEqual({ kind: 'running' })
|
||||
})
|
||||
|
||||
it('distinguishes inferred idle, timeout, session exit, and no-output startup timeout', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fake = new FakeSandbox()
|
||||
const handle = new FakePtyHandle()
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config())
|
||||
await initialize(session)
|
||||
|
||||
const inferred = session.startSend({ text: '', submit: false })
|
||||
await vi.advanceTimersByTimeAsync(40)
|
||||
expect((await inferred.done).waitReason).toBe('inferred_idle')
|
||||
|
||||
const timeout = session.startSend({ text: '', submit: false })
|
||||
for (let index = 0; index < 3; index += 1) {
|
||||
await vi.advanceTimersByTimeAsync(30)
|
||||
session.onData(Buffer.from('.'))
|
||||
}
|
||||
await vi.advanceTimersByTimeAsync(10)
|
||||
expect((await timeout.done).waitReason).toBe('timeout')
|
||||
|
||||
const exiting = session.startSend({ text: '', submit: false })
|
||||
handle.failExit(143)
|
||||
expect(await exiting.done).toMatchObject({
|
||||
waitReason: 'session_exit',
|
||||
sessionStatus: { kind: 'exited', exitCode: null, signal: 'SIGTERM' },
|
||||
})
|
||||
expect(() => session.startSend({ text: '', submit: false })).toThrow('has exited')
|
||||
|
||||
const startupHandle = new FakePtyHandle()
|
||||
const startup = new E2BPtySession(fake.sandbox, startupHandle.asHandle(), config())
|
||||
const timedOut = expect(startup.initialize()).rejects.toThrow('startup timeout')
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
await timedOut
|
||||
})
|
||||
|
||||
it('handles split prompt text, stale operations, and explicit cancellation', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fake = new FakeSandbox()
|
||||
const handle = new FakePtyHandle()
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config())
|
||||
const initializing = session.initialize()
|
||||
session.onData(Buffer.from('\x1b]133;D;0\x07'))
|
||||
await vi.advanceTimersByTimeAsync(20)
|
||||
session.onData(Buffer.from('dsh> '))
|
||||
await vi.advanceTimersByTimeAsync(10)
|
||||
await initializing
|
||||
|
||||
const operation = session.startSend({ text: 'sleep', submit: true })
|
||||
const internal = session as unknown as {
|
||||
pollReadiness(operation: PtySendOperation): void
|
||||
interrupt(operation: PtySendOperation): void
|
||||
settleActive(reason: 'timeout'): void
|
||||
failActive(error: unknown): void
|
||||
appendOutput(text: string): void
|
||||
statusValue: PtySessionStatus
|
||||
}
|
||||
internal.pollReadiness({} as PtySendOperation)
|
||||
internal.interrupt({} as PtySendOperation)
|
||||
internal.appendOutput('')
|
||||
fake.pgid = '789\n'
|
||||
expect(operation.cancel()).toBe(true)
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
expect(fake.commands).toContain('kill -INT -- -789')
|
||||
session.onData(Buffer.from('\x1b]133;D;130\x07dsh> '))
|
||||
await vi.advanceTimersByTimeAsync(10)
|
||||
await operation.done
|
||||
|
||||
internal.settleActive('timeout')
|
||||
internal.failActive(new Error('ignored'))
|
||||
const operationInternal = operation as unknown as {
|
||||
append(text: string): void
|
||||
settle(reason: 'timeout', status: PtySessionStatus, inherited: boolean): void
|
||||
fail(error: unknown): void
|
||||
}
|
||||
operationInternal.append('ignored')
|
||||
operationInternal.settle('timeout', { kind: 'running' }, false)
|
||||
operationInternal.fail(new Error('ignored'))
|
||||
})
|
||||
|
||||
it('observes AbortSignal and contains send or foreground lookup failures', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fake = new FakeSandbox()
|
||||
const handle = new FakePtyHandle()
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config())
|
||||
await initialize(session)
|
||||
|
||||
const controller = new AbortController()
|
||||
const aborting = session.startSend({ text: '', submit: false, signal: controller.signal })
|
||||
expect(() => session.startSend({ text: '', submit: false })).toThrow('active send')
|
||||
fake.pgid = 'not-a-pgid\n'
|
||||
controller.abort()
|
||||
await expect(aborting.done).rejects.toThrow('cannot resolve foreground process group')
|
||||
|
||||
const already = new AbortController()
|
||||
already.abort()
|
||||
expect(() => session.startSend({ text: '', submit: false, signal: already.signal })).toThrow('aborted before write')
|
||||
|
||||
fake.sendError = new Error('send failed')
|
||||
const failed = session.startSend({ text: 'x', submit: false })
|
||||
await expect(failed.done).rejects.toThrow('send failed')
|
||||
|
||||
fake.pgid = '123\n'
|
||||
await expect(session.signal('SIGKILL')).rejects.toThrow('refusing to SIGKILL')
|
||||
fake.pgid = '0\n'
|
||||
await expect(session.signal('SIGINT')).rejects.toThrow('cannot resolve')
|
||||
|
||||
const deferred = Promise.withResolvers<undefined>()
|
||||
fake.sendError = undefined
|
||||
const sendInput = vi.spyOn(fake.sandbox.pty, 'sendInput').mockReturnValueOnce(deferred.promise)
|
||||
const late = session.startSend({ text: 'late', submit: false })
|
||||
session.onData(Buffer.from('\x1b]133;D;0\x07dsh> '))
|
||||
await vi.advanceTimersByTimeAsync(10)
|
||||
await late.done
|
||||
deferred.reject(new Error('late failure'))
|
||||
await vi.advanceTimersByTimeAsync(0)
|
||||
expect(sendInput).toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('preserves startup abort reasons and classifies invalid UTF-8 transport failures', async () => {
|
||||
const fake = new FakeSandbox()
|
||||
const abortHandle = new FakePtyHandle()
|
||||
const abortSession = new E2BPtySession(fake.sandbox, abortHandle.asHandle(), config())
|
||||
const controller = new AbortController()
|
||||
const reason = new Error('startup cancelled')
|
||||
const initializing = abortSession.initialize(controller.signal)
|
||||
const rejected = expect(initializing).rejects.toBe(reason)
|
||||
controller.abort(reason)
|
||||
await rejected
|
||||
|
||||
const invalidHandle = new FakePtyHandle()
|
||||
const invalid = new E2BPtySession(fake.sandbox, invalidHandle.asHandle(), config())
|
||||
const pending = invalid.startSend({ text: '', submit: false })
|
||||
invalid.onData(Uint8Array.from([0xff]))
|
||||
await expect(pending.done).rejects.toThrow('invalid UTF-8')
|
||||
expect(invalid.status()).toEqual({ kind: 'exited', exitCode: null, signal: null })
|
||||
|
||||
const crashHandle = new FakePtyHandle()
|
||||
const crashed = new E2BPtySession(fake.sandbox, crashHandle.asHandle(), config())
|
||||
const active = crashed.startSend({ text: '', submit: false })
|
||||
crashHandle.crash('transport gone')
|
||||
await expect(active.done).rejects.toEqual(new Error('transport gone'))
|
||||
|
||||
const startupExitHandle = new FakePtyHandle()
|
||||
const startupExit = new E2BPtySession(fake.sandbox, startupExitHandle.asHandle(), config())
|
||||
const exitedDuringStartup = expect(startupExit.initialize()).rejects.toThrow('exited during startup')
|
||||
startupExitHandle.exit(7)
|
||||
await exitedDuringStartup
|
||||
})
|
||||
|
||||
it('covers empty bounded reads and polling an exited active session', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fake = new FakeSandbox()
|
||||
const tinyHandle = new FakePtyHandle()
|
||||
const tiny = new E2BPtySession(fake.sandbox, tinyHandle.asHandle(), config({ maxReadBytes: 1 }))
|
||||
tiny.onData(Buffer.from('你'))
|
||||
expect(tiny.read({ count: 1 })).toMatchObject({ text: '', lineEnd: 0 })
|
||||
|
||||
const handle = new FakePtyHandle()
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config())
|
||||
const operation = session.startSend({ text: '', submit: false })
|
||||
const internal = session as unknown as {
|
||||
pollReadiness(operation: PtySendOperation): void
|
||||
clearActive(): void
|
||||
statusValue: PtySessionStatus
|
||||
}
|
||||
internal.statusValue = { kind: 'exited', exitCode: 7, signal: null }
|
||||
internal.pollReadiness(operation)
|
||||
expect((await operation.done).waitReason).toBe('session_exit')
|
||||
internal.clearActive()
|
||||
})
|
||||
})
|
||||
|
||||
describe('E2BPtySession teardown', () => {
|
||||
it('terminates the process group once, awaits exit, and disconnects', async () => {
|
||||
vi.useFakeTimers()
|
||||
const fake = new FakeSandbox()
|
||||
const handle = new FakePtyHandle()
|
||||
fake.onTerm = () => { handle.failExit(143) }
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config())
|
||||
const first = session.close('done')
|
||||
expect(session.close('again')).toBe(first)
|
||||
await first
|
||||
expect(session.status()).toEqual({ kind: 'exited', exitCode: null, signal: 'SIGTERM' })
|
||||
expect(handle.disconnects).toBe(1)
|
||||
expect(() => session.startSend({ text: '', submit: false })).toThrow('closing')
|
||||
})
|
||||
|
||||
it('contains an already-gone TERM, escalates to KILL, and reports a survivor', async () => {
|
||||
vi.useFakeTimers()
|
||||
const gone = new FakeSandbox()
|
||||
const goneHandle = new FakePtyHandle()
|
||||
gone.commandError = commandError(1)
|
||||
gone.onKill = () => { goneHandle.failExit(137) }
|
||||
const goneSession = new E2BPtySession(gone.sandbox, goneHandle.asHandle(), config())
|
||||
const closingGone = goneSession.close('gone')
|
||||
await vi.advanceTimersByTimeAsync(20)
|
||||
await closingGone
|
||||
expect(gone.killed).toEqual([123])
|
||||
expect(goneSession.status()).toEqual({ kind: 'exited', exitCode: null, signal: 'SIGKILL' })
|
||||
|
||||
const survivor = new FakeSandbox()
|
||||
const survivorHandle = new FakePtyHandle()
|
||||
const survivorSession = new E2BPtySession(survivor.sandbox, survivorHandle.asHandle(), config())
|
||||
const failed = expect(survivorSession.close('still alive')).rejects.toThrow('surviving pid: 123')
|
||||
await vi.advanceTimersByTimeAsync(40)
|
||||
await failed
|
||||
survivorHandle.exit()
|
||||
await expect(survivorSession.close('retry')).resolves.toBeUndefined()
|
||||
})
|
||||
|
||||
it('propagates cleanup transport failures and lets close retry', async () => {
|
||||
const fake = new FakeSandbox()
|
||||
const handle = new FakePtyHandle()
|
||||
fake.commandError = new Error('TERM transport failed')
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config())
|
||||
await expect(session.close('failure')).rejects.toThrow('TERM transport failed')
|
||||
handle.exit()
|
||||
await expect(session.close('retry')).resolves.toBeUndefined()
|
||||
|
||||
const invalidTailHandle = new FakePtyHandle()
|
||||
const invalidTail = new E2BPtySession(fake.sandbox, invalidTailHandle.asHandle(), config())
|
||||
invalidTail.onData(Uint8Array.from([0xe2]))
|
||||
invalidTailHandle.exit()
|
||||
await expect(invalidTail.close('invalid tail')).rejects.toThrow('invalid UTF-8')
|
||||
|
||||
const normalHandle = new FakePtyHandle()
|
||||
normalHandle.disconnectError = new Error('disconnect raced')
|
||||
const normal = new E2BPtySession(fake.sandbox, normalHandle.asHandle(), config())
|
||||
normalHandle.exit(7)
|
||||
await Promise.resolve()
|
||||
expect(normal.status()).toEqual({ kind: 'exited', exitCode: 7, signal: null })
|
||||
await expect(normal.close('already exited')).resolves.toBeUndefined()
|
||||
})
|
||||
|
||||
it('kills a remotely live PTY after its host transport fails', async () => {
|
||||
const fake = new FakeSandbox()
|
||||
const handle = new FakePtyHandle()
|
||||
const session = new E2BPtySession(fake.sandbox, handle.asHandle(), config())
|
||||
const active = session.startSend({ text: '', submit: false })
|
||||
session.onData(Uint8Array.from([0xff]))
|
||||
await expect(active.done).rejects.toThrow('invalid UTF-8')
|
||||
expect(session.status()).toEqual({ kind: 'exited', exitCode: null, signal: null })
|
||||
|
||||
fake.onTerm = () => { handle.failExit(143) }
|
||||
await expect(session.close('transport failed')).rejects.toThrow('invalid UTF-8')
|
||||
expect(fake.commands).toContain('kill -TERM -- -123')
|
||||
expect(handle.disconnects).toBe(1)
|
||||
})
|
||||
})
|
||||
16
packages/pty/pty-e2b/tsconfig.json
Normal file
16
packages/pty/pty-e2b/tsconfig.json
Normal file
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"extends": "../../../tsconfig.base.json",
|
||||
"compilerOptions": {
|
||||
"rootDir": "src",
|
||||
"outDir": "lib/types"
|
||||
},
|
||||
"include": ["src"],
|
||||
"references": [
|
||||
{ "path": "../../../vendor/cosmokit" },
|
||||
{ "path": "../../../vendor/cordis" },
|
||||
{ "path": "../../../vendor/schemastery" },
|
||||
{ "path": "../../e2b/e2b" },
|
||||
{ "path": "../pty" },
|
||||
{ "path": "../../support/invariants" }
|
||||
]
|
||||
}
|
||||
@@ -1,188 +0,0 @@
|
||||
/** Streaming terminal-control sanitizer for the line-oriented first release. */
|
||||
|
||||
import { Buffer } from 'node:buffer'
|
||||
|
||||
/** OSC marker emitted by the controlled bash before each prompt. */
|
||||
export const PROMPT_MARKER_PREFIX = '133;D;'
|
||||
|
||||
/** Exact printable prompt emitted after the private marker. */
|
||||
export const CONTROLLED_PROMPT = 'dsh> '
|
||||
|
||||
/** One sanitized chunk plus whether it contained the owned prompt marker. */
|
||||
export interface SanitizedChunk {
|
||||
text: string
|
||||
prompt: boolean
|
||||
/** Printable text after the latest owned marker in this chunk. */
|
||||
promptTail?: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Remove CSI/OSC/short escape sequences while preserving split-sequence carry.
|
||||
* Full terminal emulation is deliberately deferred; ordinary line output and
|
||||
* the private prompt marker are the supported contract.
|
||||
*/
|
||||
export class TerminalSanitizer {
|
||||
private pending = ''
|
||||
private discardMode: 'osc' | 'csi' | undefined
|
||||
private discardOscEscape = false
|
||||
private trailingCarriageReturn = false
|
||||
private trackingPromptTail = false
|
||||
|
||||
constructor(private readonly maxPendingBytes: number) {}
|
||||
|
||||
/**
|
||||
* Consume one decoded `node-pty` data chunk.
|
||||
* @param chunk - decoded terminal data.
|
||||
* @returns Printable text and whether the private prompt marker completed.
|
||||
*/
|
||||
push(chunk: string): SanitizedChunk {
|
||||
this.pending += this.discardPrefix(chunk)
|
||||
let text = ''
|
||||
let prompt = false
|
||||
let includePromptTail = this.trackingPromptTail
|
||||
let promptTail = ''
|
||||
let index = 0
|
||||
const appendText = (value: string): void => {
|
||||
text += value
|
||||
if (this.trackingPromptTail) promptTail += value
|
||||
}
|
||||
while (index < this.pending.length) {
|
||||
const escape = this.pending.indexOf('\x1b', index)
|
||||
if (escape < 0) {
|
||||
appendText(this.pending.slice(index))
|
||||
index = this.pending.length
|
||||
break
|
||||
}
|
||||
appendText(this.pending.slice(index, escape))
|
||||
if (escape + 1 >= this.pending.length) {
|
||||
index = escape
|
||||
break
|
||||
}
|
||||
const kind = this.pending[escape + 1]
|
||||
if (kind === ']') {
|
||||
const bel = this.pending.indexOf('\x07', escape + 2)
|
||||
const stringTerminator = this.pending.indexOf('\x1b\\', escape + 2)
|
||||
let end = -1
|
||||
if (bel >= 0 && stringTerminator >= 0) end = Math.min(bel + 1, stringTerminator + 2)
|
||||
else if (bel >= 0) end = bel + 1
|
||||
else if (stringTerminator >= 0) end = stringTerminator + 2
|
||||
if (end < 0) {
|
||||
index = escape
|
||||
break
|
||||
}
|
||||
const terminatorBytes = this.pending[end - 1] === '\x07' ? 1 : 2
|
||||
const content = this.pending.slice(escape + 2, end - terminatorBytes)
|
||||
if (content.startsWith(PROMPT_MARKER_PREFIX)) {
|
||||
prompt = true
|
||||
this.trackingPromptTail = true
|
||||
includePromptTail = true
|
||||
promptTail = ''
|
||||
}
|
||||
index = end
|
||||
continue
|
||||
}
|
||||
if (kind === '[') {
|
||||
let end = escape + 2
|
||||
while (end < this.pending.length) {
|
||||
const code = this.pending.charCodeAt(end)
|
||||
if (code >= 0x40 && code <= 0x7e) break
|
||||
end += 1
|
||||
}
|
||||
if (end >= this.pending.length) {
|
||||
index = escape
|
||||
break
|
||||
}
|
||||
index = end + 1
|
||||
continue
|
||||
}
|
||||
// Two-byte escape family (save/restore cursor and similar).
|
||||
index = escape + 2
|
||||
}
|
||||
this.pending = this.pending.slice(index)
|
||||
this.enforcePendingBound()
|
||||
return {
|
||||
text: this.normalizeText(text),
|
||||
prompt,
|
||||
...includePromptTail ? { promptTail } : {},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Flush a trailing printable fragment when the PTY exits.
|
||||
* @returns Remaining printable text; incomplete escapes are discarded.
|
||||
*/
|
||||
flush(): string {
|
||||
const text = this.pending.startsWith('\x1b') ? '' : this.pending
|
||||
this.pending = ''
|
||||
this.discardMode = undefined
|
||||
this.discardOscEscape = false
|
||||
this.trackingPromptTail = false
|
||||
const normalized = this.normalizeText(text)
|
||||
if (!this.trailingCarriageReturn) return normalized
|
||||
this.trailingCarriageReturn = false
|
||||
return `${normalized}\n`
|
||||
}
|
||||
|
||||
private normalizeText(text: string): string {
|
||||
let complete = this.trailingCarriageReturn ? `\r${text}` : text
|
||||
this.trailingCarriageReturn = false
|
||||
if (complete.endsWith('\r')) {
|
||||
complete = complete.slice(0, -1)
|
||||
this.trailingCarriageReturn = true
|
||||
}
|
||||
return normalizeTerminalText(complete)
|
||||
}
|
||||
|
||||
private enforcePendingBound(): void {
|
||||
if (Buffer.byteLength(this.pending) <= this.maxPendingBytes) return
|
||||
this.discardMode = this.pending[1] === ']' ? 'osc' : 'csi'
|
||||
this.pending = ''
|
||||
}
|
||||
|
||||
private discardPrefix(chunk: string): string {
|
||||
if (this.discardMode === undefined) return chunk
|
||||
if (this.discardMode === 'csi') {
|
||||
for (let index = 0; index < chunk.length; index += 1) {
|
||||
const code = chunk.charCodeAt(index)
|
||||
if (code >= 0x40 && code <= 0x7e) {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(index + 1)
|
||||
}
|
||||
}
|
||||
return ''
|
||||
}
|
||||
|
||||
let index = 0
|
||||
if (this.discardOscEscape) {
|
||||
this.discardOscEscape = false
|
||||
if (chunk.startsWith('\\')) {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(1)
|
||||
}
|
||||
}
|
||||
while (index < chunk.length) {
|
||||
if (chunk[index] === '\x07') {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(index + 1)
|
||||
}
|
||||
if (chunk[index] === '\x1b') {
|
||||
if (chunk[index + 1] === '\\') {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(index + 2)
|
||||
}
|
||||
if (index + 1 === chunk.length) this.discardOscEscape = true
|
||||
}
|
||||
index += 1
|
||||
}
|
||||
return ''
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalize CRLF and standalone carriage returns for line-oriented rendering.
|
||||
* @param text - sanitized terminal text.
|
||||
* @returns Line-normalized text with BEL removed.
|
||||
*/
|
||||
export function normalizeTerminalText(text: string): string {
|
||||
return text.replaceAll('\r\n', '\n').replaceAll('\r', '\n').replaceAll('\x07', '')
|
||||
}
|
||||
@@ -1,12 +1,7 @@
|
||||
/** Persistent PTY session over the subprocess seam's terminal primitive. */
|
||||
/** Local `node-pty` session: bounded output, readiness, signals, and teardown. */
|
||||
|
||||
import { Buffer } from 'node:buffer'
|
||||
import type {
|
||||
SubprocessOutcome,
|
||||
SubprocessTerminalForeground,
|
||||
SubprocessTerminalHandle,
|
||||
} from '@deepseek-ai/dsh-subprocess'
|
||||
import { PtyError } from '@deepseek-ai/dsh-pty'
|
||||
import type { IDisposable, IPty } from 'node-pty'
|
||||
import { PtyTerminalSanitizer, PtyTextBuffer, ptySignalName, ptyUtf8Tail } from '@deepseek-ai/dsh-pty'
|
||||
import type {
|
||||
PtyBackendSession,
|
||||
PtyReadRequest,
|
||||
@@ -21,89 +16,30 @@ import type {
|
||||
PtyWaitReason,
|
||||
} from '@deepseek-ai/dsh-pty'
|
||||
import type { ResolvedConfig } from './config.ts'
|
||||
import { CONTROLLED_PROMPT, TerminalSanitizer } from './sanitize.ts'
|
||||
import type { ProcessIdentity, ProcessInspector } from './process-inspector.ts'
|
||||
|
||||
function utf8Tail(text: string, maxBytes: number): { text: string; truncated: boolean } {
|
||||
if (Buffer.byteLength(text) <= maxBytes) return { text, truncated: false }
|
||||
const chars = Array.from(text)
|
||||
let bytes = 0
|
||||
let start = chars.length
|
||||
while (start > 0) {
|
||||
const next = Buffer.byteLength(chars[start - 1] as string)
|
||||
if (bytes + next > maxBytes) break
|
||||
bytes += next
|
||||
start -= 1
|
||||
}
|
||||
return { text: chars.slice(start).join(''), truncated: true }
|
||||
}
|
||||
|
||||
class BoundedTextBuffer {
|
||||
private value = ''
|
||||
private dropped = false
|
||||
|
||||
constructor(
|
||||
private readonly maxBytes: number,
|
||||
private readonly maxLines?: number,
|
||||
) {}
|
||||
|
||||
append(text: string): void {
|
||||
if (text.length === 0) return
|
||||
this.value += text
|
||||
if (this.maxLines !== undefined) {
|
||||
const lines = this.value.split('\n')
|
||||
if (lines.length > this.maxLines) {
|
||||
this.value = lines.slice(lines.length - this.maxLines).join('\n')
|
||||
this.dropped = true
|
||||
}
|
||||
}
|
||||
const tail = utf8Tail(this.value, this.maxBytes)
|
||||
this.value = tail.text
|
||||
this.dropped ||= tail.truncated
|
||||
}
|
||||
|
||||
consume(): PtySendRead {
|
||||
const delta = this.value
|
||||
const truncated = this.dropped
|
||||
this.value = ''
|
||||
this.dropped = false
|
||||
return { delta, truncated }
|
||||
}
|
||||
|
||||
snapshot(): { text: string; truncated: boolean } {
|
||||
return { text: this.value, truncated: this.dropped }
|
||||
}
|
||||
function delay(ms: number): Promise<void> {
|
||||
return new Promise(resolve => setTimeout(resolve, ms))
|
||||
}
|
||||
|
||||
class LocalSendOperation implements PtySendOperation {
|
||||
private readonly output: BoundedTextBuffer
|
||||
private readonly output: PtyTextBuffer
|
||||
private readonly promise: PromiseWithResolvers<PtySendResult>
|
||||
private finished = false
|
||||
private cancellationRequested = false
|
||||
private initialForegroundLeftWait: boolean
|
||||
private initialForegroundPgid: number | undefined
|
||||
|
||||
constructor(
|
||||
maxBytes: number,
|
||||
readonly startedAt: number,
|
||||
private readonly onCancel: () => void,
|
||||
) {
|
||||
this.output = new BoundedTextBuffer(maxBytes)
|
||||
this.output = new PtyTextBuffer(maxBytes)
|
||||
this.promise = Promise.withResolvers<PtySendResult>()
|
||||
this.initialForegroundLeftWait = true
|
||||
}
|
||||
|
||||
get done(): Promise<PtySendResult> {
|
||||
return this.promise.promise
|
||||
}
|
||||
|
||||
get settled(): boolean {
|
||||
return this.finished
|
||||
}
|
||||
|
||||
get cancelRequested(): boolean {
|
||||
return this.cancellationRequested
|
||||
}
|
||||
|
||||
append(text: string): void {
|
||||
if (!this.finished) this.output.append(text)
|
||||
}
|
||||
@@ -130,75 +66,50 @@ class LocalSendOperation implements PtySendOperation {
|
||||
return this.output.consume()
|
||||
}
|
||||
|
||||
setInitialForeground(foreground: SubprocessTerminalForeground | undefined): void {
|
||||
this.initialForegroundPgid = foreground?.processGroupId
|
||||
this.initialForegroundLeftWait = foreground?.inputWaiting !== true
|
||||
}
|
||||
|
||||
acceptsStdinWait(pgid: number, waiting: boolean): boolean {
|
||||
// The same group may still expose the wait that existed before terminal.write.
|
||||
// Observe every poll so a departure before the exact-settlement threshold
|
||||
// still makes a later return to that wait post-write evidence.
|
||||
if (pgid !== this.initialForegroundPgid) return waiting
|
||||
if (!waiting) this.initialForegroundLeftWait = true
|
||||
return waiting && this.initialForegroundLeftWait
|
||||
}
|
||||
|
||||
cancel(): boolean {
|
||||
if (this.finished) return false
|
||||
this.cancellationRequested = true
|
||||
this.onCancel()
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
/** Backend session wrapping one provider-owned terminal process. */
|
||||
/** Backend session wrapping one `node-pty` process and its captured process tree. */
|
||||
export class LocalPtySession implements PtyBackendSession {
|
||||
motd = ''
|
||||
readonly pid: number
|
||||
private readonly decoder = new TextDecoder()
|
||||
private readonly sanitizer: TerminalSanitizer
|
||||
private readonly scrollback: BoundedTextBuffer
|
||||
private readonly outputEnded = Promise.withResolvers<void>()
|
||||
private readonly completion: Promise<void>
|
||||
private readonly sanitizer: PtyTerminalSanitizer
|
||||
private readonly scrollback: PtyTextBuffer
|
||||
private readonly exitPromise: PromiseWithResolvers<void> = Promise.withResolvers<void>()
|
||||
private readonly dataDisposable: IDisposable
|
||||
private readonly exitDisposable: IDisposable
|
||||
private statusValue: PtySessionStatus = { kind: 'running' }
|
||||
// TODO(pty-send-state-consolidation): Fold the per-send fields below
|
||||
// (active/activeTimer/activeDeadlineTimer/activeAbort/interrupting/
|
||||
// activeWrite/pollingReady/polling) into one send-lifecycle owner; the
|
||||
// cancellation/readiness interplay now has enough pinned tests to carry
|
||||
// that refactor safely.
|
||||
private active: LocalSendOperation | undefined
|
||||
private activeTimer: NodeJS.Timeout | undefined
|
||||
private activeDeadlineTimer: NodeJS.Timeout | undefined
|
||||
private activeAbort: (() => void) | undefined
|
||||
private interrupting: LocalSendOperation | undefined
|
||||
private activeWrite: Promise<boolean> | undefined
|
||||
private pollingReady: LocalSendOperation | undefined
|
||||
private polling = false
|
||||
private promptSeen = false
|
||||
private promptTextSeen = false
|
||||
private promptTail = ''
|
||||
private shellPgid: number | undefined
|
||||
private initializing = false
|
||||
private lastOutputAt = Date.now()
|
||||
private closing = false
|
||||
private closePromise: Promise<void> | undefined
|
||||
private transportFailure: Error | undefined
|
||||
|
||||
constructor(
|
||||
private readonly terminal: SubprocessTerminalHandle,
|
||||
private readonly terminal: IPty,
|
||||
private readonly inspector: ProcessInspector,
|
||||
private readonly config: ResolvedConfig,
|
||||
) {
|
||||
this.pid = terminal.pid
|
||||
this.sanitizer = new TerminalSanitizer(config.maxReadBytes)
|
||||
this.scrollback = new BoundedTextBuffer(config.scrollbackMaxBytes, config.scrollbackLines)
|
||||
terminal.output.on('data', this.onTerminalData)
|
||||
terminal.output.once('end', this.onTerminalEnd)
|
||||
terminal.output.once('error', this.onTerminalError)
|
||||
this.completion = terminal.done.then(
|
||||
outcome => this.onExit(outcome),
|
||||
(error: unknown) => { this.onTransportFailure(error) },
|
||||
)
|
||||
this.sanitizer = new PtyTerminalSanitizer(config.maxReadBytes)
|
||||
this.scrollback = new PtyTextBuffer(config.scrollbackMaxBytes, config.scrollbackLines)
|
||||
this.dataDisposable = terminal.onData((data) => { this.onData(data) })
|
||||
this.exitDisposable = terminal.onExit(({ exitCode, signal }) => {
|
||||
const tail = this.sanitizer.flush()
|
||||
this.appendOutput(tail)
|
||||
this.statusValue = { kind: 'exited', exitCode, signal: ptySignalName(signal) }
|
||||
this.settleActive('session_exit')
|
||||
this.exitPromise.resolve()
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -225,14 +136,7 @@ export class LocalPtySession implements PtyBackendSession {
|
||||
startSend(request: PtySendRequest): PtySendOperation {
|
||||
if (this.closing) throw new Error('PTY session is closing')
|
||||
if (this.statusValue.kind === 'exited') throw new Error('PTY session has exited')
|
||||
if (this.active !== undefined) {
|
||||
const draining = this.activeWrite !== undefined
|
||||
? ' or draining provider write'
|
||||
: this.interrupting !== undefined
|
||||
? ' or draining foreground interrupt'
|
||||
: ''
|
||||
throw new PtyError(`PTY session already has an active send${draining}`, 'SEND_ACTIVE')
|
||||
}
|
||||
if (this.active !== undefined) throw new Error('PTY session already has an active send')
|
||||
if (request.signal?.aborted === true) throw new Error('PTY send aborted before write')
|
||||
|
||||
const operation = new LocalSendOperation(
|
||||
@@ -241,79 +145,29 @@ export class LocalPtySession implements PtyBackendSession {
|
||||
() => { this.interrupt(operation) },
|
||||
)
|
||||
this.active = operation
|
||||
this.resetReadinessEvidence()
|
||||
this.lastOutputAt = Date.now()
|
||||
this.promptSeen = false
|
||||
this.promptTextSeen = false
|
||||
|
||||
if (request.signal !== undefined) {
|
||||
const onAbort = (): void => { operation.cancel() }
|
||||
request.signal.addEventListener('abort', onAbort, { once: true })
|
||||
this.activeAbort = () => request.signal?.removeEventListener('abort', onAbort)
|
||||
}
|
||||
this.activeDeadlineTimer = setTimeout(() => {
|
||||
if (this.active === operation) {
|
||||
this.settleActive('timeout', this.activeWrite !== undefined || this.interrupting === operation)
|
||||
}
|
||||
}, this.config.timeoutMs)
|
||||
void this.beginSend(operation, request)
|
||||
|
||||
try {
|
||||
if (request.text.length > 0) this.terminal.write(request.text)
|
||||
if (request.submit) this.terminal.write('\r')
|
||||
} catch (error: unknown) {
|
||||
this.clearActive()
|
||||
operation.fail(error)
|
||||
return operation
|
||||
}
|
||||
|
||||
this.activeTimer = setInterval(() => { this.pollReadiness(operation) }, this.config.pollIntervalMs)
|
||||
return operation
|
||||
}
|
||||
|
||||
private async beginSend(operation: LocalSendOperation, request: PtySendRequest): Promise<void> {
|
||||
let foreground: SubprocessTerminalForeground | undefined
|
||||
try {
|
||||
foreground = await this.terminal.inspectForeground()
|
||||
} catch (error: unknown) {
|
||||
// A pre-write inspection failure while cancellation owns the slot must not
|
||||
// release it: interruptOnce's in-flight foreground signal could land on a
|
||||
// successor's foreground group. The interrupt path's post-signal tail
|
||||
// resumes polling, whose guarded catch propagates a persistent failure.
|
||||
// A retained settled operation implies that same in-flight interrupt, so
|
||||
// this guard admits only an unsettled active send.
|
||||
if (this.active === operation && !this.closing && this.interrupting !== operation) {
|
||||
this.failActive(error)
|
||||
}
|
||||
return
|
||||
}
|
||||
try {
|
||||
if (this.active !== operation || this.closing || this.interrupting === operation) return
|
||||
operation.setInitialForeground(foreground)
|
||||
const input = `${request.text}${request.submit ? '\r' : ''}`
|
||||
if (input.length > 0 && !operation.cancelRequested) {
|
||||
this.resetReadinessEvidence()
|
||||
const write = this.terminal.write(input)
|
||||
this.activeWrite = write.then(() => true, () => false)
|
||||
try {
|
||||
await write
|
||||
} finally {
|
||||
this.activeWrite = undefined
|
||||
}
|
||||
}
|
||||
// Cancellation owns post-write signalling and reservation release.
|
||||
if (operation.cancelRequested) return
|
||||
if (this.active === operation && operation.settled) {
|
||||
this.clearActive()
|
||||
return
|
||||
}
|
||||
// Closing can race the awaited provider write even though static analysis sees only local assignments.
|
||||
// oxlint-disable-next-line typescript/no-unnecessary-condition -- awaited provider writes can close the session.
|
||||
if (this.active === operation && !this.closing) {
|
||||
this.pollingReady = operation
|
||||
this.schedulePoll(operation)
|
||||
}
|
||||
} catch (error: unknown) {
|
||||
if (this.active === operation && !this.closing) {
|
||||
if (operation.settled) this.clearActive()
|
||||
else this.failActive(error)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private resetReadinessEvidence(): void {
|
||||
this.lastOutputAt = Date.now()
|
||||
this.promptSeen = false
|
||||
this.promptTextSeen = false
|
||||
this.promptTail = ''
|
||||
}
|
||||
|
||||
read(request: PtyReadRequest): PtyReadResult {
|
||||
const snapshot = this.scrollback.snapshot()
|
||||
const lines = snapshot.text.split('\n')
|
||||
@@ -328,7 +182,7 @@ export class LocalPtySession implements PtyBackendSession {
|
||||
const end = totalLines - offset
|
||||
const start = Math.max(0, end - count)
|
||||
const requested = lines.slice(start, end).join('\n')
|
||||
const bounded = utf8Tail(requested, this.config.maxReadBytes)
|
||||
const bounded = ptyUtf8Tail(requested, this.config.maxReadBytes)
|
||||
const returnedLines = bounded.text.length === 0 ? 0 : bounded.text.split('\n').length
|
||||
return {
|
||||
text: bounded.text,
|
||||
@@ -339,10 +193,16 @@ export class LocalPtySession implements PtyBackendSession {
|
||||
}
|
||||
}
|
||||
|
||||
async signal(signal: PtySignal): Promise<PtySignalResult> {
|
||||
if (this.closing) throw new Error('PTY session is closing')
|
||||
const targetPgid = await this.terminal.signalForeground(signal)
|
||||
return { delivered: true, targetPgid }
|
||||
signal(signal: PtySignal): Promise<PtySignalResult> {
|
||||
return Promise.resolve().then(() => {
|
||||
const pgid = this.inspector.foregroundPgid(this.pid)
|
||||
if (pgid === undefined) throw new Error(`cannot resolve foreground process group for PTY ${this.pid}`)
|
||||
if (signal === 'SIGKILL' && pgid === this.pid) {
|
||||
throw new Error('refusing to SIGKILL the PTY shell; use terminal_close')
|
||||
}
|
||||
this.inspector.signalGroup(pgid, signal)
|
||||
return { delivered: true, targetPgid: pgid }
|
||||
})
|
||||
}
|
||||
|
||||
status(): PtySessionStatus {
|
||||
@@ -361,56 +221,21 @@ export class LocalPtySession implements PtyBackendSession {
|
||||
return closing
|
||||
}
|
||||
|
||||
private readonly onTerminalData = (chunk: Buffer | Uint8Array | string): void => {
|
||||
const bytes = typeof chunk === 'string' ? Buffer.from(chunk, 'utf8') : chunk
|
||||
this.onData(this.decoder.decode(bytes, { stream: true }))
|
||||
}
|
||||
|
||||
private readonly onTerminalEnd = (): void => {
|
||||
this.onData(this.decoder.decode())
|
||||
this.appendOutput(this.sanitizer.flush())
|
||||
this.outputEnded.resolve()
|
||||
}
|
||||
|
||||
private readonly onTerminalError = (error: Error): void => {
|
||||
this.onTransportFailure(error)
|
||||
this.outputEnded.resolve()
|
||||
}
|
||||
|
||||
private onData(data: string): void {
|
||||
const sanitized = this.sanitizer.push(data)
|
||||
this.appendOutput(sanitized.text)
|
||||
if (sanitized.prompt) {
|
||||
// TODO(pty-delayed-signal-prompt): With a reproducer, define a marker-generation boundary
|
||||
// before attributing a signal-delayed prompt to a later send.
|
||||
const foregroundPgid = this.inspector.foregroundPgid(this.pid)
|
||||
if (this.shellPgid === undefined) this.shellPgid = foregroundPgid
|
||||
// Bash can print PROMPT_COMMAND before the kernel publishes its return
|
||||
// to the foreground process group. Retain the marker; polling below is
|
||||
// the authority that accepts it only after bash owns the foreground.
|
||||
this.promptSeen = true
|
||||
this.promptTail = ''
|
||||
this.promptTextSeen = sanitized.promptText === true
|
||||
this.lastOutputAt = Date.now()
|
||||
} else if (this.promptSeen && sanitized.promptText === true) {
|
||||
this.promptTextSeen = true
|
||||
}
|
||||
if (this.promptSeen && sanitized.promptTail !== undefined) {
|
||||
const remaining = Math.max(0, CONTROLLED_PROMPT.length + 1 - this.promptTail.length)
|
||||
this.promptTail += sanitized.promptTail.slice(0, remaining)
|
||||
if (sanitized.promptTail.length > remaining) this.promptTail = `${CONTROLLED_PROMPT}\0`
|
||||
this.promptTextSeen = this.promptTail === CONTROLLED_PROMPT
|
||||
}
|
||||
}
|
||||
|
||||
private async onExit(outcome: SubprocessOutcome): Promise<void> {
|
||||
await this.outputEnded.promise
|
||||
if (this.transportFailure !== undefined) return
|
||||
this.statusValue = { kind: 'exited', exitCode: outcome.exitCode, signal: outcome.signal }
|
||||
this.settleActive('session_exit')
|
||||
}
|
||||
|
||||
private onTransportFailure(error: unknown): void {
|
||||
const failure = error instanceof Error ? error : new Error(String(error))
|
||||
this.transportFailure ??= failure
|
||||
this.statusValue = { kind: 'exited', exitCode: null, signal: null }
|
||||
this.failActive(failure)
|
||||
void this.terminal.terminate().catch(() => {})
|
||||
}
|
||||
|
||||
private appendOutput(text: string): void {
|
||||
@@ -420,94 +245,60 @@ export class LocalPtySession implements PtyBackendSession {
|
||||
this.active?.append(text)
|
||||
}
|
||||
|
||||
private schedulePoll(operation: LocalSendOperation, delayMs = this.config.pollIntervalMs): void {
|
||||
if (this.active !== operation || this.interrupting === operation || this.polling) return
|
||||
if (this.activeTimer !== undefined) clearTimeout(this.activeTimer)
|
||||
this.activeTimer = setTimeout(() => {
|
||||
this.activeTimer = undefined
|
||||
void this.pollReadiness(operation)
|
||||
}, delayMs)
|
||||
}
|
||||
|
||||
private async pollReadiness(operation: LocalSendOperation): Promise<void> {
|
||||
if (this.active !== operation || this.polling) return
|
||||
this.polling = true
|
||||
try {
|
||||
if (this.statusValue.kind === 'exited') {
|
||||
this.settleActive('session_exit')
|
||||
return
|
||||
}
|
||||
const foreground = await this.terminal.inspectForeground()
|
||||
if (this.active !== operation || this.closing || this.interrupting === operation) return
|
||||
const idleFor = Date.now() - this.lastOutputAt
|
||||
if (this.promptSeen && foreground !== undefined && this.shellPgid === undefined) {
|
||||
this.shellPgid = foreground.processGroupId
|
||||
}
|
||||
if (this.promptSeen && this.promptTextSeen && idleFor >= this.config.pollIntervalMs
|
||||
&& foreground?.processGroupId === this.shellPgid) {
|
||||
this.settleActive('stdin_read')
|
||||
return
|
||||
}
|
||||
const elapsed = Date.now() - operation.startedAt
|
||||
const startupHasOutput = !this.initializing || this.scrollback.snapshot().text.length > 0
|
||||
const acceptsStdinWait = startupHasOutput && foreground !== undefined
|
||||
&& operation.acceptsStdinWait(foreground.processGroupId, foreground.inputWaiting)
|
||||
if (elapsed >= this.config.exactProbeAfterMs && acceptsStdinWait) {
|
||||
this.settleActive('stdin_read')
|
||||
return
|
||||
}
|
||||
// A prompt candidate can race bash's foreground handoff, but an interactive
|
||||
// child also inherits PROMPT_COMMAND. Silence therefore remains the bound
|
||||
// on waiting for shell ownership instead of letting a child marker suppress
|
||||
// readiness until the absolute timeout.
|
||||
const handoffGrace = this.promptSeen ? this.config.handoffGraceMs : 0
|
||||
if (startupHasOutput && idleFor >= this.config.idleSilenceMs + handoffGrace) {
|
||||
this.settleActive('inferred_idle')
|
||||
}
|
||||
} catch (error: unknown) {
|
||||
if (this.active === operation && !this.closing && this.interrupting !== operation) this.failActive(error)
|
||||
} finally {
|
||||
this.polling = false
|
||||
const active = this.active
|
||||
// Awaited provider inspection can clear or replace the active send despite static analysis.
|
||||
// oxlint-disable-next-line typescript/no-unnecessary-condition -- awaited inspection can replace the active send.
|
||||
if (active !== undefined && this.pollingReady === active) this.schedulePoll(active)
|
||||
private pollReadiness(operation: LocalSendOperation): void {
|
||||
if (this.active !== operation) return
|
||||
if (this.statusValue.kind === 'exited') {
|
||||
this.settleActive('session_exit')
|
||||
return
|
||||
}
|
||||
if (this.promptSeen && this.promptTextSeen && Date.now() - this.lastOutputAt >= this.config.pollIntervalMs) {
|
||||
const pgid = this.inspector.foregroundPgid(this.pid)
|
||||
if (this.shellPgid !== undefined && pgid === this.shellPgid) {
|
||||
this.settleActive('stdin_read')
|
||||
return
|
||||
}
|
||||
}
|
||||
const elapsed = Date.now() - operation.startedAt
|
||||
const startupHasOutput = !this.initializing || this.scrollback.snapshot().text.length > 0
|
||||
if (startupHasOutput && elapsed >= this.config.exactProbeAfterMs) {
|
||||
const pgid = this.inspector.foregroundPgid(this.pid)
|
||||
if (pgid !== undefined && this.inspector.isStdinWaiting(pgid)) {
|
||||
this.settleActive('stdin_read')
|
||||
return
|
||||
}
|
||||
}
|
||||
// A prompt candidate can race bash's foreground handoff, but an interactive
|
||||
// child also inherits PROMPT_COMMAND. Silence therefore remains the bound
|
||||
// on waiting for shell ownership instead of letting a child marker suppress
|
||||
// readiness until the absolute timeout. When a prompt marker was seen, the
|
||||
// configured grace holds the fallback past the silence bound so polls in
|
||||
// that window can observe the foreground handoff and settle as stdin_read.
|
||||
const idleFor = Date.now() - this.lastOutputAt
|
||||
const handoffGrace = this.promptSeen ? this.config.handoffGraceMs : 0
|
||||
if (startupHasOutput && idleFor >= this.config.idleSilenceMs + handoffGrace) {
|
||||
this.settleActive('inferred_idle')
|
||||
return
|
||||
}
|
||||
if (elapsed >= this.config.timeoutMs) this.settleActive('timeout')
|
||||
}
|
||||
|
||||
private settleActive(waitReason: PtyWaitReason, retainOwnership = false): void {
|
||||
private settleActive(waitReason: PtyWaitReason): void {
|
||||
const operation = this.active
|
||||
if (operation === undefined) return
|
||||
const scrollbackTruncated = this.scrollback.snapshot().truncated
|
||||
if (retainOwnership) {
|
||||
this.stopPolling()
|
||||
this.activeAbort?.()
|
||||
this.activeAbort = undefined
|
||||
} else {
|
||||
this.clearActive()
|
||||
}
|
||||
this.clearActive()
|
||||
operation.settle(waitReason, this.statusValue, scrollbackTruncated)
|
||||
}
|
||||
|
||||
private stopPolling(): void {
|
||||
this.stopReadinessPolling()
|
||||
if (this.activeDeadlineTimer !== undefined) clearTimeout(this.activeDeadlineTimer)
|
||||
this.activeDeadlineTimer = undefined
|
||||
}
|
||||
|
||||
private stopReadinessPolling(): void {
|
||||
if (this.activeTimer !== undefined) clearTimeout(this.activeTimer)
|
||||
if (this.activeTimer !== undefined) clearInterval(this.activeTimer)
|
||||
this.activeTimer = undefined
|
||||
this.pollingReady = undefined
|
||||
}
|
||||
|
||||
private clearActive(): void {
|
||||
const operation = this.active
|
||||
this.stopPolling()
|
||||
this.activeAbort?.()
|
||||
this.activeAbort = undefined
|
||||
if (this.interrupting === operation) this.interrupting = undefined
|
||||
this.pollingReady = undefined
|
||||
this.active = undefined
|
||||
}
|
||||
|
||||
@@ -520,46 +311,104 @@ export class LocalPtySession implements PtyBackendSession {
|
||||
|
||||
private interrupt(operation: LocalSendOperation): void {
|
||||
if (this.active !== operation) return
|
||||
this.interrupting = operation
|
||||
this.stopReadinessPolling()
|
||||
void this.interruptOnce(operation)
|
||||
try {
|
||||
const pgid = this.inspector.foregroundPgid(this.pid)
|
||||
if (pgid === undefined) throw new Error(`cannot resolve foreground process group for PTY ${this.pid}`)
|
||||
this.inspector.signalGroup(pgid, 'SIGINT')
|
||||
} catch (error: unknown) {
|
||||
this.failActive(error)
|
||||
}
|
||||
}
|
||||
|
||||
private async interruptOnce(operation: LocalSendOperation): Promise<void> {
|
||||
try {
|
||||
const activeWrite = this.activeWrite
|
||||
if (activeWrite !== undefined && !await activeWrite) return
|
||||
await this.terminal.signalForeground('SIGINT')
|
||||
} catch (error: unknown) {
|
||||
if (this.active === operation && !this.closing) this.onTransportFailure(error)
|
||||
return
|
||||
} finally {
|
||||
if (this.interrupting === operation) this.interrupting = undefined
|
||||
private survivors(members: ProcessIdentity[]): ProcessIdentity[] {
|
||||
return members.filter(member => this.inspector.isAlive(member))
|
||||
}
|
||||
|
||||
private descendants(): ProcessIdentity[] {
|
||||
return this.inspector.processTree(this.pid).filter(member => member.pid !== this.pid)
|
||||
}
|
||||
|
||||
private async waitForExit(members: ProcessIdentity[]): Promise<ProcessIdentity[]> {
|
||||
const deadline = Date.now() + this.config.disposeGraceMs
|
||||
let survivors = this.survivors(members)
|
||||
while (survivors.length > 0 && Date.now() < deadline) {
|
||||
await delay(Math.min(25, Math.max(1, deadline - Date.now())))
|
||||
survivors = this.survivors(members)
|
||||
}
|
||||
if (this.active === operation && operation.settled) {
|
||||
this.clearActive()
|
||||
} else if (this.active === operation && !this.closing) {
|
||||
this.pollingReady = operation
|
||||
this.schedulePoll(operation, 0)
|
||||
return survivors
|
||||
}
|
||||
|
||||
private signalMembers(members: ProcessIdentity[], signal: 'SIGTERM' | 'SIGKILL'): void {
|
||||
for (const member of members) {
|
||||
try {
|
||||
this.inspector.signalProcess(member, signal)
|
||||
} catch (_alreadyExitedDuringSignal) {
|
||||
// Identity is rechecked by the inspector; a same-tick exit is success.
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private unionMembers(...groups: ProcessIdentity[][]): ProcessIdentity[] {
|
||||
const members: ProcessIdentity[] = []
|
||||
const seen = new Set<string>()
|
||||
for (const group of groups) {
|
||||
for (const member of group) {
|
||||
const key = JSON.stringify([member.pid, member.started])
|
||||
if (seen.has(key)) continue
|
||||
seen.add(key)
|
||||
members.push(member)
|
||||
}
|
||||
}
|
||||
return members
|
||||
}
|
||||
|
||||
private async stopDescendants(): Promise<ProcessIdentity[]> {
|
||||
const captured = this.descendants()
|
||||
this.signalMembers(captured, 'SIGTERM')
|
||||
const capturedSurvivors = await this.waitForExit(captured)
|
||||
// A TERM-handling descendant may have forked while winding down. Rescan
|
||||
// while the shell can still reap every member, then kill both the fresh
|
||||
// tree and captured survivors that were reparented out of that tree.
|
||||
const members = this.unionMembers(capturedSurvivors, this.descendants())
|
||||
this.signalMembers(members, 'SIGKILL')
|
||||
const survivors = await this.waitForExit(members)
|
||||
return this.survivors(this.unionMembers(survivors, this.descendants()))
|
||||
}
|
||||
|
||||
private async stopShell(): Promise<void> {
|
||||
try {
|
||||
this.terminal.kill('SIGTERM')
|
||||
} catch (_topLevelAlreadyExitedDuringTerm) {
|
||||
// The exit notification remains authoritative.
|
||||
}
|
||||
if (this.statusValue.kind === 'running') {
|
||||
await Promise.race([this.exitPromise.promise, delay(this.config.disposeGraceMs)])
|
||||
}
|
||||
if (this.statusValue.kind === 'running') {
|
||||
try {
|
||||
this.terminal.kill('SIGKILL')
|
||||
} catch (_topLevelAlreadyExitedDuringKill) {
|
||||
// The exit notification remains authoritative.
|
||||
}
|
||||
await Promise.race([this.exitPromise.promise, delay(this.config.disposeGraceMs)])
|
||||
}
|
||||
if (this.statusValue.kind === 'running') {
|
||||
throw new Error(`PTY cleanup failed; surviving pids: ${this.pid}`)
|
||||
}
|
||||
}
|
||||
|
||||
private async closeOnce(reason: string): Promise<void> {
|
||||
this.dataDisposable.dispose()
|
||||
// Stop readiness polling but retain the active operation: teardown settles
|
||||
// it as session_exit below, so an in-flight send is never mis-settled as
|
||||
// stdin_read/inferred_idle/timeout during the grace period.
|
||||
this.stopPolling()
|
||||
try {
|
||||
await this.terminal.terminate()
|
||||
} catch (error: unknown) {
|
||||
throw new Error(`PTY cleanup failed (${reason})`, { cause: error })
|
||||
const survivors = await this.stopDescendants()
|
||||
if (survivors.length > 0) {
|
||||
throw new Error(`PTY cleanup failed (${reason}); surviving pids: ${survivors.map(member => member.pid).join(', ')}`)
|
||||
}
|
||||
// Quiescence is the active send's terminal outcome.
|
||||
await this.stopShell()
|
||||
this.settleActive('session_exit')
|
||||
await this.completion
|
||||
this.terminal.output.off('data', this.onTerminalData)
|
||||
this.terminal.output.off('end', this.onTerminalEnd)
|
||||
this.terminal.output.off('error', this.onTerminalError)
|
||||
if (this.transportFailure !== undefined) throw this.transportFailure
|
||||
this.exitDisposable.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,17 +1,17 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { normalizeTerminalText, TerminalSanitizer } from '@deepseek-ai/dsh-pty-local/src/sanitize.ts'
|
||||
import { normalizePtyTerminalText, PtyTerminalSanitizer } from '@deepseek-ai/dsh-pty'
|
||||
|
||||
describe('TerminalSanitizer', () => {
|
||||
describe('PtyTerminalSanitizer', () => {
|
||||
it('removes split CSI and owned OSC prompt markers', () => {
|
||||
const sanitizer = new TerminalSanitizer(64)
|
||||
const sanitizer = new PtyTerminalSanitizer(64)
|
||||
expect(sanitizer.push('red\x1b[3')).toEqual({ text: 'red', prompt: false })
|
||||
expect(sanitizer.push('1m text\x1b[0m\r\n')).toEqual({ text: ' text\n', prompt: false })
|
||||
expect(sanitizer.push('\x1b]133;')).toEqual({ text: '', prompt: false })
|
||||
expect(sanitizer.push('D;0\x07dsh> ')).toEqual({ text: 'dsh> ', prompt: true, promptTail: 'dsh> ' })
|
||||
expect(sanitizer.push('D;0\x07dsh> ')).toEqual({ text: 'dsh> ', prompt: true, promptText: true })
|
||||
})
|
||||
|
||||
it('drops unrelated OSC, short escapes, BEL, and incomplete trailing escape', () => {
|
||||
const sanitizer = new TerminalSanitizer(64)
|
||||
const sanitizer = new PtyTerminalSanitizer(64)
|
||||
expect(sanitizer.push('a\x1b]0;title\x1b\\b\x1b7c\x07')).toEqual({ text: 'abc', prompt: false })
|
||||
expect(sanitizer.push('tail\x1b')).toEqual({ text: 'tail', prompt: false })
|
||||
expect(sanitizer.flush()).toBe('')
|
||||
@@ -22,11 +22,11 @@ describe('TerminalSanitizer', () => {
|
||||
})
|
||||
|
||||
it('normalizes CRLF and standalone carriage returns', () => {
|
||||
expect(normalizeTerminalText('a\r\nb\rc\x07')).toBe('a\nb\nc')
|
||||
expect(normalizePtyTerminalText('a\r\nb\rc\x07')).toBe('a\nb\nc')
|
||||
})
|
||||
|
||||
it('carries a trailing carriage return across data chunks and flushes standalone CR', () => {
|
||||
const sanitizer = new TerminalSanitizer(64)
|
||||
const sanitizer = new PtyTerminalSanitizer(64)
|
||||
expect(sanitizer.push('a\r')).toEqual({ text: 'a', prompt: false })
|
||||
expect(sanitizer.push('\nb')).toEqual({ text: '\nb', prompt: false })
|
||||
expect(sanitizer.push('\r')).toEqual({ text: '', prompt: false })
|
||||
@@ -34,41 +34,41 @@ describe('TerminalSanitizer', () => {
|
||||
})
|
||||
|
||||
it('reports printable prompt text that follows a marker in a later chunk', () => {
|
||||
const sanitizer = new TerminalSanitizer(64)
|
||||
expect(sanitizer.push('\x1b]133;D;0\x07')).toEqual({ text: '', prompt: true, promptTail: '' })
|
||||
expect(sanitizer.push('dsh> ')).toEqual({ text: 'dsh> ', prompt: false, promptTail: 'dsh> ' })
|
||||
const sanitizer = new PtyTerminalSanitizer(64)
|
||||
expect(sanitizer.push('\x1b]133;D;0\x07')).toEqual({ text: '', prompt: true })
|
||||
expect(sanitizer.push('dsh> ')).toEqual({ text: 'dsh> ', prompt: false, promptText: true })
|
||||
})
|
||||
|
||||
it('bounds and discards unterminated control sequences through their terminators', () => {
|
||||
const oscBel = new TerminalSanitizer(8)
|
||||
const oscBel = new PtyTerminalSanitizer(8)
|
||||
expect(oscBel.push(`\x1b]0;${'x'.repeat(16)}`)).toEqual({ text: '', prompt: false })
|
||||
expect(oscBel.push('more\x07tail')).toEqual({ text: 'tail', prompt: false })
|
||||
|
||||
const oscSt = new TerminalSanitizer(8)
|
||||
const oscSt = new PtyTerminalSanitizer(8)
|
||||
oscSt.push(`\x1b]0;${'x'.repeat(16)}`)
|
||||
expect(oscSt.push('more\x1b')).toEqual({ text: '', prompt: false })
|
||||
expect(oscSt.push('\\tail')).toEqual({ text: 'tail', prompt: false })
|
||||
|
||||
const oscDirectSt = new TerminalSanitizer(8)
|
||||
const oscDirectSt = new PtyTerminalSanitizer(8)
|
||||
oscDirectSt.push(`\x1b]0;${'x'.repeat(16)}`)
|
||||
expect(oscDirectSt.push('more\x1b\\tail')).toEqual({ text: 'tail', prompt: false })
|
||||
|
||||
const oscFalseSt = new TerminalSanitizer(8)
|
||||
const oscFalseSt = new PtyTerminalSanitizer(8)
|
||||
oscFalseSt.push(`\x1b]0;${'x'.repeat(16)}`)
|
||||
oscFalseSt.push('\x1b')
|
||||
expect(oscFalseSt.push('more')).toEqual({ text: '', prompt: false })
|
||||
expect(oscFalseSt.push('\x07tail')).toEqual({ text: 'tail', prompt: false })
|
||||
|
||||
const oscNonTerminatingEscape = new TerminalSanitizer(8)
|
||||
const oscNonTerminatingEscape = new PtyTerminalSanitizer(8)
|
||||
oscNonTerminatingEscape.push(`\x1b]0;${'x'.repeat(16)}`)
|
||||
expect(oscNonTerminatingEscape.push('more\x1bxmore\x07tail')).toEqual({ text: 'tail', prompt: false })
|
||||
|
||||
const csi = new TerminalSanitizer(8)
|
||||
const csi = new PtyTerminalSanitizer(8)
|
||||
expect(csi.push(`\x1b[${'1'.repeat(16)}`)).toEqual({ text: '', prompt: false })
|
||||
expect(csi.push('123')).toEqual({ text: '', prompt: false })
|
||||
expect(csi.push('mtext')).toEqual({ text: 'text', prompt: false })
|
||||
|
||||
const flushed = new TerminalSanitizer(8)
|
||||
const flushed = new PtyTerminalSanitizer(8)
|
||||
flushed.push(`\x1b]0;${'x'.repeat(16)}`)
|
||||
expect(flushed.flush()).toBe('')
|
||||
expect(flushed.push('text')).toEqual({ text: 'text', prompt: false })
|
||||
|
||||
@@ -41,6 +41,15 @@ export type {
|
||||
PtyWaitReason,
|
||||
} from './types.ts'
|
||||
export { PtyBackendCleanupError } from './types.ts'
|
||||
export {
|
||||
normalizePtyTerminalText,
|
||||
PTY_PROMPT_MARKER_PREFIX,
|
||||
PtyTerminalSanitizer,
|
||||
PtyTextBuffer,
|
||||
ptySignalName,
|
||||
ptyUtf8Tail,
|
||||
} from './terminal.ts'
|
||||
export type { PtySanitizedChunk } from './terminal.ts'
|
||||
|
||||
/** Opaque identity minted by {@link PtyService} for one live PTY session. */
|
||||
export type PtySessionId = PtySessionIdValue
|
||||
|
||||
273
packages/pty/pty/src/terminal.ts
Normal file
273
packages/pty/pty/src/terminal.ts
Normal file
@@ -0,0 +1,273 @@
|
||||
/** Backend-neutral line-oriented terminal buffering and control-sequence sanitization. */
|
||||
|
||||
import { Buffer } from 'node:buffer'
|
||||
import { constants } from 'node:os'
|
||||
import type { PtySendRead } from './types.ts'
|
||||
|
||||
/** OSC marker emitted by a controlled bash before each prompt. */
|
||||
export const PTY_PROMPT_MARKER_PREFIX = '133;D;'
|
||||
|
||||
/** One sanitized chunk plus whether it contained the controlled prompt marker. */
|
||||
export interface PtySanitizedChunk {
|
||||
/** Printable, line-normalized terminal text. */
|
||||
text: string
|
||||
/** Whether the chunk completed the controlled prompt marker. */
|
||||
prompt: boolean
|
||||
/** Present when printable text followed the latest controlled prompt marker. */
|
||||
promptText?: true
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the largest code-point-aligned UTF-8 tail within a byte cap.
|
||||
* @param text - Candidate terminal text.
|
||||
* @param maxBytes - Maximum retained UTF-8 bytes.
|
||||
* @returns The retained tail and whether its head was dropped.
|
||||
*/
|
||||
export function ptyUtf8Tail(text: string, maxBytes: number): { text: string; truncated: boolean } {
|
||||
if (Buffer.byteLength(text) <= maxBytes) return { text, truncated: false }
|
||||
const chars = Array.from(text)
|
||||
let bytes = 0
|
||||
let start = chars.length
|
||||
while (start > 0) {
|
||||
const next = Buffer.byteLength(chars[start - 1] as string)
|
||||
if (bytes + next > maxBytes) break
|
||||
bytes += next
|
||||
start -= 1
|
||||
}
|
||||
return { text: chars.slice(start).join(''), truncated: true }
|
||||
}
|
||||
|
||||
/** UTF-8 and optionally line-bounded terminal text buffer. */
|
||||
export class PtyTextBuffer {
|
||||
private value = ''
|
||||
private dropped = false
|
||||
|
||||
/**
|
||||
* @param maxBytes - Maximum retained UTF-8 bytes.
|
||||
* @param maxLines - Optional maximum retained logical lines.
|
||||
*/
|
||||
constructor(
|
||||
private readonly maxBytes: number,
|
||||
private readonly maxLines?: number,
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Append terminal text and drop the oldest excess.
|
||||
* @param text - Decoded and sanitized terminal text.
|
||||
*/
|
||||
append(text: string): void {
|
||||
if (text.length === 0) return
|
||||
this.value += text
|
||||
if (this.maxLines !== undefined) {
|
||||
const lines = this.value.split('\n')
|
||||
if (lines.length > this.maxLines) {
|
||||
this.value = lines.slice(lines.length - this.maxLines).join('\n')
|
||||
this.dropped = true
|
||||
}
|
||||
}
|
||||
const tail = ptyUtf8Tail(this.value, this.maxBytes)
|
||||
this.value = tail.text
|
||||
this.dropped ||= tail.truncated
|
||||
}
|
||||
|
||||
/**
|
||||
* Consume all currently retained operation text.
|
||||
* @returns The delta and whether older text was dropped.
|
||||
*/
|
||||
consume(): PtySendRead {
|
||||
const delta = this.value
|
||||
const truncated = this.dropped
|
||||
this.value = ''
|
||||
this.dropped = false
|
||||
return { delta, truncated }
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the retained text without consuming it.
|
||||
* @returns The retained text and whether its head was dropped.
|
||||
*/
|
||||
snapshot(): { text: string; truncated: boolean } {
|
||||
return { text: this.value, truncated: this.dropped }
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Streaming terminal-control sanitizer for line-oriented PTY backends.
|
||||
* Full terminal emulation is deliberately outside the PTY seam.
|
||||
*/
|
||||
export class PtyTerminalSanitizer {
|
||||
private pending = ''
|
||||
private discardMode: 'osc' | 'csi' | undefined
|
||||
private discardOscEscape = false
|
||||
private trailingCarriageReturn = false
|
||||
private awaitingPromptText = false
|
||||
|
||||
/** @param maxPendingBytes - Bound for an incomplete terminal-control sequence. */
|
||||
constructor(private readonly maxPendingBytes: number) {}
|
||||
|
||||
/**
|
||||
* Consume one decoded PTY data chunk.
|
||||
* @param chunk - Decoded terminal data.
|
||||
* @returns Printable text and prompt-marker facts.
|
||||
*/
|
||||
push(chunk: string): PtySanitizedChunk {
|
||||
this.pending += this.discardPrefix(chunk)
|
||||
let text = ''
|
||||
let prompt = false
|
||||
let promptText = false
|
||||
let index = 0
|
||||
const appendText = (value: string): boolean => {
|
||||
text += value
|
||||
if (this.awaitingPromptText && value.replace(/[\r\n\x07]/g, '').length > 0) {
|
||||
this.awaitingPromptText = false
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
while (index < this.pending.length) {
|
||||
const escape = this.pending.indexOf('\x1b', index)
|
||||
if (escape < 0) {
|
||||
promptText = appendText(this.pending.slice(index)) || promptText
|
||||
index = this.pending.length
|
||||
break
|
||||
}
|
||||
promptText = appendText(this.pending.slice(index, escape)) || promptText
|
||||
if (escape + 1 >= this.pending.length) {
|
||||
index = escape
|
||||
break
|
||||
}
|
||||
const kind = this.pending[escape + 1]
|
||||
if (kind === ']') {
|
||||
const bel = this.pending.indexOf('\x07', escape + 2)
|
||||
const stringTerminator = this.pending.indexOf('\x1b\\', escape + 2)
|
||||
let end = -1
|
||||
if (bel >= 0 && stringTerminator >= 0) end = Math.min(bel + 1, stringTerminator + 2)
|
||||
else if (bel >= 0) end = bel + 1
|
||||
else if (stringTerminator >= 0) end = stringTerminator + 2
|
||||
if (end < 0) {
|
||||
index = escape
|
||||
break
|
||||
}
|
||||
const terminatorBytes = this.pending[end - 1] === '\x07' ? 1 : 2
|
||||
const content = this.pending.slice(escape + 2, end - terminatorBytes)
|
||||
if (content.startsWith(PTY_PROMPT_MARKER_PREFIX)) {
|
||||
prompt = true
|
||||
promptText = false
|
||||
this.awaitingPromptText = true
|
||||
}
|
||||
index = end
|
||||
continue
|
||||
}
|
||||
if (kind === '[') {
|
||||
let end = escape + 2
|
||||
while (end < this.pending.length) {
|
||||
const code = this.pending.charCodeAt(end)
|
||||
if (code >= 0x40 && code <= 0x7e) break
|
||||
end += 1
|
||||
}
|
||||
if (end >= this.pending.length) {
|
||||
index = escape
|
||||
break
|
||||
}
|
||||
index = end + 1
|
||||
continue
|
||||
}
|
||||
index = escape + 2
|
||||
}
|
||||
this.pending = this.pending.slice(index)
|
||||
this.enforcePendingBound()
|
||||
return { text: this.normalizeText(text), prompt, ...promptText ? { promptText: true } : {} }
|
||||
}
|
||||
|
||||
/**
|
||||
* Flush printable trailing data and discard incomplete controls.
|
||||
* @returns Remaining normalized printable text.
|
||||
*/
|
||||
flush(): string {
|
||||
const text = this.pending.startsWith('\x1b') ? '' : this.pending
|
||||
this.pending = ''
|
||||
this.discardMode = undefined
|
||||
this.discardOscEscape = false
|
||||
this.awaitingPromptText = false
|
||||
const normalized = this.normalizeText(text)
|
||||
if (!this.trailingCarriageReturn) return normalized
|
||||
this.trailingCarriageReturn = false
|
||||
return `${normalized}\n`
|
||||
}
|
||||
|
||||
private normalizeText(text: string): string {
|
||||
let complete = this.trailingCarriageReturn ? `\r${text}` : text
|
||||
this.trailingCarriageReturn = false
|
||||
if (complete.endsWith('\r')) {
|
||||
complete = complete.slice(0, -1)
|
||||
this.trailingCarriageReturn = true
|
||||
}
|
||||
return normalizePtyTerminalText(complete)
|
||||
}
|
||||
|
||||
private enforcePendingBound(): void {
|
||||
if (Buffer.byteLength(this.pending) <= this.maxPendingBytes) return
|
||||
this.discardMode = this.pending[1] === ']' ? 'osc' : 'csi'
|
||||
this.pending = ''
|
||||
}
|
||||
|
||||
private discardPrefix(chunk: string): string {
|
||||
if (this.discardMode === undefined) return chunk
|
||||
if (this.discardMode === 'csi') {
|
||||
for (let index = 0; index < chunk.length; index += 1) {
|
||||
const code = chunk.charCodeAt(index)
|
||||
if (code >= 0x40 && code <= 0x7e) {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(index + 1)
|
||||
}
|
||||
}
|
||||
return ''
|
||||
}
|
||||
|
||||
let index = 0
|
||||
if (this.discardOscEscape) {
|
||||
this.discardOscEscape = false
|
||||
if (chunk.startsWith('\\')) {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(1)
|
||||
}
|
||||
}
|
||||
while (index < chunk.length) {
|
||||
if (chunk[index] === '\x07') {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(index + 1)
|
||||
}
|
||||
if (chunk[index] === '\x1b') {
|
||||
if (chunk[index + 1] === '\\') {
|
||||
this.discardMode = undefined
|
||||
return chunk.slice(index + 2)
|
||||
}
|
||||
if (index + 1 === chunk.length) this.discardOscEscape = true
|
||||
}
|
||||
index += 1
|
||||
}
|
||||
return ''
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalize CRLF and standalone carriage returns for line-oriented rendering.
|
||||
* @param text - Sanitized terminal text.
|
||||
* @returns Line-normalized text with BEL removed.
|
||||
*/
|
||||
export function normalizePtyTerminalText(text: string): string {
|
||||
return text.replaceAll('\r\n', '\n').replaceAll('\r', '\n').replaceAll('\x07', '')
|
||||
}
|
||||
|
||||
/**
|
||||
* Convert a platform signal number into the seam's signal-name vocabulary.
|
||||
* @param number - Platform signal number, zero, or an absent signal.
|
||||
* @returns The matching Node signal name, or `null` when unknown or absent.
|
||||
*/
|
||||
export function ptySignalName(number: number | undefined): NodeJS.Signals | null {
|
||||
if (number === undefined || number === 0) return null
|
||||
for (const [name, value] of Object.entries(constants.signals)) {
|
||||
if (value === number) return name as NodeJS.Signals
|
||||
}
|
||||
return null
|
||||
}
|
||||
Reference in New Issue
Block a user