mirror of
https://github.com/Ark0N/Codeman.git
synced 2026-09-30 12:39:42 +02:00
fix(deepseek): review fixes for the transcript reader — docker/remote gate, poll memo, honest pairing docs
Four review findings on the worker-transcript feature: - Docker and remote-SSH dsh sessions now keep the pane segmenter: their transcripts live in the container's / remote host's own ~/.dsh, which the local reader can never see, so the transcript path returned 'nothing said yet' forever and an agent polling such a worker starved on an answer that existed. Gated on !session.docker && !session.remote (statically pinned) and documented in the integration guide. - last-response reads are memoized on (path, mtime, size, blocks): the skill's last_text polls once per second, and each poll decompressed and reparsed the whole file on the event loop even when nothing had been appended. An unchanged poll now costs one stat. - The pairing ladder's comment claimed /new is served by step 2; in truth the boot-window transcript wins for as long as it exists (deliberately: preferring newest-eligible would hand a worker its busier sibling's reply). The comment now states the real tradeoff instead of the aspirational one. Same for decodeZstdFrames' 'skipped' wording — a corrupt frame truncates the decode there, which is the safe behavior. - stripReasoningPrefix no longer runs on user prompt text, so a prompt containing a literal </think> renders whole in blocks view. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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
|
||||
|
||||
+52
-11
@@ -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 + '</think>'.length);
|
||||
}
|
||||
|
||||
function textOfContent(content: unknown): string {
|
||||
/** `stripReasoning` is for ASSISTANT content only: a user prompt containing a
|
||||
* literal `</think>` (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<string, DeepSeekTranscriptResult>();
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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');
|
||||
|
||||
@@ -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 });
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user