diff --git a/docs/deepseek-integration.md b/docs/deepseek-integration.md index e0a44467..de142648 100644 --- a/docs/deepseek-integration.md +++ b/docs/deepseek-integration.md @@ -223,6 +223,12 @@ reader (`src/deepseek-transcript.ts`): which is returned as `Turn error: …` (and an early stop such as `max-tokens` as `Turn ended: …`) instead of an empty string that reads as "still thinking". +The transcript reader applies to **local** dsh sessions only. A Docker case's +harness writes its transcript inside the container's own `~/.dsh` (the workspace +bind mount does not cover it), and a remote-SSH case's lives on the remote host, +so the local reader could never find those files — such sessions keep the pane +segmenter, coarse but real. The splash caveat above applies to them accordingly. + ### As an agent worker Because dsh has both halves — a real end-of-turn signal and a real transcript — an diff --git a/src/deepseek-transcript.ts b/src/deepseek-transcript.ts index 6a10d032..dc374f10 100644 --- a/src/deepseek-transcript.ts +++ b/src/deepseek-transcript.ts @@ -193,8 +193,11 @@ export function zstdFrameRanges(buf: Buffer): Array<[number, number]> { * a zstd magic is passed through unchanged, which is what lets the same reader * open a plain `session.jsonl` (dsh writes one when compression is off). * - * A frame that fails to decompress is skipped rather than fatal: a half-written - * tail frame must not cost the caller the whole conversation. + * A frame that fails to decompress truncates the decode THERE rather than + * failing it: everything decoded before it is kept, so a half-written tail + * frame does not cost the caller the whole conversation. (Not "skipped" — a + * frame after a corrupt one is never reached, which is the safe reading: dsh + * appends, so a bad frame means everything after it is suspect too.) */ export function decodeZstdFrames(buf: Buffer): string { if (buf.length < 4) return buf.toString('utf8'); @@ -246,12 +249,16 @@ function stripReasoningPrefix(text: string): string { return close === -1 ? text : text.slice(close + ''.length); } -function textOfContent(content: unknown): string { +/** `stripReasoning` is for ASSISTANT content only: a user prompt containing a + * literal `` (someone pasting a transcript, say) must render whole. */ +function textOfContent(content: unknown, stripReasoning = true): string { const parts: string[] = []; for (const entry of asArray(content)) { const block = asRecord(entry); if (!block) continue; - if (block.type === 'text' && typeof block.text === 'string') parts.push(stripReasoningPrefix(block.text)); + if (block.type === 'text' && typeof block.text === 'string') { + parts.push(stripReasoning ? stripReasoningPrefix(block.text) : block.text); + } } return parts.join('').trim(); } @@ -374,7 +381,7 @@ export function parseDeepSeekTranscript(raw: string, options: { blocks?: boolean // snapshot dsh injects every turn (sandbox policy, approvals, cwd). if (asRecord(data.source)?.kind !== 'user') break; if (!wantBlocks) break; - const text = textOfContent(data.content); + const text = textOfContent(data.content, false); if (text) blocks.push({ kind: 'prompt', label: 'Prompt', role: 'user', text }); break; } @@ -508,13 +515,19 @@ const TRANSCRIPT_FILES = ['session.jsonl.zstd', 'session.jsonl']; * 1. a transcript whose header `createdAt` sits within `PAIRING_WINDOW_MS` of * this session's start — that is this pane's own boot, and it stays right * even when a sibling session is running in the same case directory; - * 2. otherwise the newest transcript created after this session started — - * which is what `/new` inside a live session produces; + * 2. otherwise the newest transcript created after this session started; * 3. otherwise nothing. * - * ⚠️ Step 2 cannot tell a `/new` from a sibling that started later in the same - * workspace. Agent fleets do not hit it (each worker gets its own case, and the - * skill refuses a duplicate name); two humans sharing one case directory can. + * ⚠️ The boot transcript wins for as long as it exists on disk — deliberately, + * and even over a LATER transcript in the same workspace. Step 2 cannot tell a + * `/new` from a sibling session that started later in the same directory, so + * preferring newest-eligible would hand a worker its busier sibling's reply + * (the exact bug the hard rule above was measured against, one seat over). + * The cost of that choice: after an interactive `/new` in a dsh tab, this + * reader keeps serving the pre-`/new` conversation (the same session's own + * earlier turns — stale, never foreign); step 2 is reached only when no + * boot-window transcript exists. Worker fleets never `/new`, so they only + * ever see step 1. */ export async function findDeepSeekTranscript(options: { dshHome: string; @@ -611,6 +624,24 @@ async function readTranscriptHeader(path: string): Promise<{ cwd?: string; id?: * file past this is pathological and is not worth a synchronous decode. */ const MAX_TRANSCRIPT_BYTES = 64 * 1024 * 1024; +/** + * Memo of the last few decoded transcripts, keyed on (path, mtime, size, + * blocks). The skill's `last_text` polls once per second, and each poll used + * to zstdDecompressSync + reparse the WHOLE file on the event loop even when + * nothing had been appended — a multi-MB transcript made that a repeated + * ~100ms-class stall on the single-threaded server. A poll that finds the + * file unchanged now costs one stat. Insertion-order eviction; tiny, because + * an entry only earns its keep while a session is being actively polled. + */ +const parseMemo = new Map(); +const PARSE_MEMO_MAX = 16; + +/** Test seam: a fixture that rewrites one path in place inside a single mtime + * tick would otherwise read its predecessor back out of the memo. */ +export function resetDeepSeekTranscriptMemoForTest(): void { + parseMemo.clear(); +} + /** * Read one dsh session's last answer. * @@ -645,11 +676,21 @@ export async function readDeepSeekLastResponse( const stat = await fs.stat(path).catch(() => null); if (!stat || stat.size > MAX_TRANSCRIPT_BYTES) return empty; + const memoKey = `${path}|${stat.mtimeMs}|${stat.size}|${options.blocks ? 1 : 0}`; + const memoized = parseMemo.get(memoKey); + if (memoized) return memoized; + let buf: Buffer; try { buf = await fs.readFile(path); } catch { return empty; } - return parseDeepSeekTranscript(decodeZstdFrames(buf), options); + const result = parseDeepSeekTranscript(decodeZstdFrames(buf), options); + if (parseMemo.size >= PARSE_MEMO_MAX) { + const oldest = parseMemo.keys().next().value; + if (oldest !== undefined) parseMemo.delete(oldest); + } + parseMemo.set(memoKey, result); + return result; } diff --git a/src/web/routes/session-routes.ts b/src/web/routes/session-routes.ts index e7957d5f..efcbf7e1 100644 --- a/src/web/routes/session-routes.ts +++ b/src/web/routes/session-routes.ts @@ -2051,7 +2051,14 @@ export function registerSessionRoutes( // reply reads as a reply. So an EMPTY transcript result still wins over the // pane — "nothing said yet" is the honest answer. Only `null`, meaning a // Node too old to decode zstd, falls through to the segmenter below. - if (session.mode === 'deepseek') { + // ⚠️ Local sessions only: a docker case's harness writes its transcript + // inside the CONTAINER's ~/.dsh (the workspace bind-mount does not cover + // it) and a remote-SSH case's lives on the remote host, so the local + // reader would scan a $DSH_HOME that can never hold this session's file + // and return "nothing said yet" forever — an agent polling that worker + // would starve on an answer that exists. Those configurations keep the + // pane segmenter below: coarse, but the real conversation. + if (session.mode === 'deepseek' && !session.docker && !session.remote) { const deepSeekQuery = req.query as { context?: string }; const full = deepSeekQuery.context === 'full'; const transcript = await readDeepSeekLastResponse(session, { blocks: full }); diff --git a/test/deepseek-mode.test.ts b/test/deepseek-mode.test.ts index 1cdb7b71..8f954a81 100644 --- a/test/deepseek-mode.test.ts +++ b/test/deepseek-mode.test.ts @@ -249,6 +249,18 @@ describe('DeepSeek status bridge', () => { expect(server).toContain("if (!session || session.mode !== 'claude') return;"); }); + it('keeps the transcript reader off docker and remote-SSH sessions', () => { + // A docker case's harness writes its transcript inside the CONTAINER's + // ~/.dsh and a remote-SSH case's lives on the remote host, so the local + // reader would scan a $DSH_HOME that can never hold the file and return + // "nothing said yet" forever — starving an agent that polls the worker. + // Those sessions must keep the pane segmenter. Static, because standing up + // a docker/remote session in the unit harness is exactly what the tmux + // test-mode mocks exist to avoid. + const routes = readFileSync(join(process.cwd(), 'src/web/routes/session-routes.ts'), 'utf-8'); + expect(routes).toMatch(/session\.mode === 'deepseek' && !session\.docker && !session\.remote/); + }); + it('maps the harness lifecycle states onto real hook events', () => { expect(DEEPSEEK_STATE_TO_HOOK_EVENT.idle).toBe('stop'); expect(DEEPSEEK_STATE_TO_HOOK_EVENT.blocked).toBe('permission_prompt'); diff --git a/test/deepseek-transcript.test.ts b/test/deepseek-transcript.test.ts index 0a1afaeb..e0837070 100644 --- a/test/deepseek-transcript.test.ts +++ b/test/deepseek-transcript.test.ts @@ -23,6 +23,7 @@ import { findDeepSeekTranscript, parseDeepSeekTranscript, readDeepSeekLastResponse, + resetDeepSeekTranscriptMemoForTest, resolveDeepSeekHome, zstdFrameRanges, zstdSupported, @@ -358,4 +359,37 @@ describe.skipIf(!zstdSupported())('session-to-transcript pairing', () => { }); expect(result).toEqual({ text: '', timestamp: '', blocks: [] }); }); + + it('serves an unchanged transcript from the memo and re-reads when the stat moves', async () => { + const home = await mkdtemp(join(tmpdir(), 'dsh-home-')); + try { + const dir = join(home, 'sessions', '--home-tester-cases-worker-1--', 'memo'); + await mkdir(dir, { recursive: true }); + const path = join(dir, 'session.jsonl'); + const stamp = new Date(sessionStart + 1_000); + const read = () => + readDeepSeekLastResponse({ workingDir: workspace, createdAt: sessionStart, deepSeekHomeOverride: home }); + const body = (answer: string) => + `${[sessionHeader(workspace, sessionStart + 1_000), assistantMessage(answer), turnEnd(1, { kind: 'completed' })].join('\n')}\n`; + + resetDeepSeekTranscriptMemoForTest(); + await writeFile(path, body('AAAA')); + await utimes(path, stamp, stamp); + expect((await read())?.text).toBe('AAAA'); + + // Same byte length, same forced mtime: indistinguishable from unchanged + // by stat, and deliberately served from the memo — the 1s/poll skill loop + // must not decode an unchanged file, and dsh only ever APPENDS, so a + // same-stat rewrite does not exist outside a test. + await writeFile(path, body('BBBB')); + await utimes(path, stamp, stamp); + expect((await read())?.text).toBe('AAAA'); + + // An append moves mtime (and normally size), which is the invalidation. + await utimes(path, new Date(sessionStart + 2_000), new Date(sessionStart + 2_000)); + expect((await read())?.text).toBe('BBBB'); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); });