Files
orca/src/main/native-chat/transcript-stream-lines.ts
JinjingandOrca 2c28b6c92c Gate WSL transcript filesystem I/O to prevent stalls (STA-4049) (#14203)
* fix(ai-vault): gate post-resolution WSL transcript I/O (STA-4049)

PR #14090 admitted only path *resolution* through the WSL transcript
filesystem gate. Every byte read afterwards from the resulting
\\wsl.localhost\... UNC path ran raw, so a distro that answers the first
access() and then stalls hung Native Chat at "loading" and AI Vault at
"scanning" with no timeout and no error.

Route that I/O through a new wsl-transcript-fs-access accessor, which is
a verbatim node:fs passthrough off UNC and an admitted, deadlined task on
it. open/positional-read opt out of coalescing (dedupe: false): joiners
would share one FileHandle or one caller's buffer.

Refusals now surface as the existing retryable message rather than
notFound, per-root scan failures are contained to an AiVaultScanIssue,
and the memoized Codex/Kimi indexes evict on refusal so a stall cannot
pin "no titles"/"no cwd" until the index changes.

* fix(ai-vault): stop caching WSL gate refusals as results (STA-4049)

Code review 1 P1 fixes on top of the transcript gate:

- transcript-read-cache: never store a gate refusal. The refusal leaves the
  file's mtime untouched, so the cached error would have been served to every
  later call until the transcript itself changed.
- kimi/grok/opencode parsers: rethrow WslTranscriptFsError instead of folding it
  into "no session"/"no transcript", so the session parse cache cannot store a
  null or partial answer under an unchanged mtime. Ordinary missing/half-written
  files stay contained.
- opencode-usage scanner: gate the data-directory readdir and the absolute
  OPENCODE_DB stat. The AI Vault's primary OpenCode source reaches them
  transitively, which is why the direct-import guard never saw them.
- gated stat/lstat: accept an AbortSignal, matching gated open/read, so a
  cancelled watch install or title probe detaches immediately instead of holding
  a waiter to its deadline.
- gated open: close a FileHandle whose syscall lands after the last waiter gave
  up, and close handles off UNC verbatim (awaited, failures surfaced).

Co-authored-by: Orca <help@stably.ai>

* fix(native-chat): decode gated chunks incrementally and cancel drain I/O (STA-4049)

Addresses the CR2 blockers.

UTF-8 chunk-boundary corruption: the UNC branch yielded raw 1 MiB Buffer
slices that `decodeTranscriptStream` decoded independently, so any multibyte
codepoint straddling a boundary became U+FFFD on both sides — corrupting the
JSONL line and shifting `consumedBytes` (which seeds fallback message ids).
`gatedChunks` now holds a StringDecoder when `encoding` is set, and
`decodeTranscriptStream` holds one for the Buffer path, matching what
`createReadStream`'s decoder already did off UNC.

Watcher teardown: `installTranscriptWatcher` owns an AbortController that
`unsubscribe()` aborts, threaded through every gated call on the drain path.
Waiters now detach at teardown instead of holding to the 30s deadline, and
the gate's aborted-signal pre-check stops an in-flight drain from admitting
new tasks after close.

Rovo `session_context.json`: `readJsonObjectIfExists` rethrows
WslTranscriptFsError so `parseSessionCandidate` records a scan issue, instead
of caching an un-enriched session under an unchanged mtime that never re-reads.

Primary OpenCode source: `listOpenCodeDatabases` takes an optional refusal
reporter so a refused `OPENCODE_DB`/`XDG_DATA_HOME` surfaces an
AiVaultScanIssue, matching `listOpenCodeDatabasesInDirectory`.

`boundaryFingerprint` moved to its own module to keep the watcher engine
under the max-lines cap.

* refactor(native-chat): consolidate transcript I/O and remove fallback te

- Move boundaryFingerprint from its own module to transcript-file-version.ts
- Extract runPathOperation helper to eliminate duplicate UNC path routing
- Remove tests for fallback behaviors when transcripts are unavailable or incomplete
- Clean up implementation comments and verbose test documentation

* consolidate scan issues and gate session scanner I/O (STA-4049)

Both local and remote session scans hit stalled WSL distros identically:
one failed probe per discovered path. Unifying issue recording and gate
refusal handling prevents duplication and ensures consistent behavior.

- Gate all file operations (stat, readdir, read, open) through WSL
  stall detection instead of scattered or missing gates
- Serve cached transcripts when stat stalls; distinguish gate refusals
  from missing files
- Serialize UNC close operations to prevent thread pool exhaustion
- Incremental chunk decoding in streams handles codepoint boundaries
  correctly

---------

Co-authored-by: Orca <help@stably.ai>
2026-08-13 00:54:19 -07:00

53 lines
1.8 KiB
TypeScript

import type { Readable } from 'node:stream'
import { StringDecoder } from 'node:string_decoder'
import type { NativeChatMessage } from '../../shared/native-chat-types'
import { transcriptFallbackId } from './transcript-fallback-id'
type TranscriptDecoder = (line: string, fallbackId: string) => NativeChatMessage | null
export async function decodeTranscriptStream(
stream: Readable,
filePath: string,
start: number,
decode: TranscriptDecoder,
includeTrailingLine: boolean
): Promise<{ messages: NativeChatMessage[]; consumedBytes: number }> {
const messages: NativeChatMessage[] = []
// Why: a Buffer chunk can end mid-codepoint, and decoding it standalone would
// both corrupt the line and shift `consumedBytes` (which seeds fallback ids).
const decoder = new StringDecoder('utf8')
let pending = ''
let consumedBytes = 0
for await (const chunk of stream) {
pending += typeof chunk === 'string' ? chunk : decoder.write(Buffer.from(chunk))
let newlineIndex = pending.indexOf('\n')
while (newlineIndex !== -1) {
const segment = pending.slice(0, newlineIndex + 1)
decodeLine(segment.slice(0, -1), consumedBytes)
consumedBytes += Buffer.byteLength(segment, 'utf8')
pending = pending.slice(newlineIndex + 1)
newlineIndex = pending.indexOf('\n')
}
}
pending += decoder.end()
if (includeTrailingLine && pending.length > 0) {
decodeLine(pending, consumedBytes)
consumedBytes += Buffer.byteLength(pending, 'utf8')
}
return { messages, consumedBytes }
function decodeLine(rawLine: string, relativeOffset: number): void {
const line = rawLine.endsWith('\r') ? rawLine.slice(0, -1) : rawLine
if (!line) {
return
}
const message = decode(line, transcriptFallbackId(filePath, start + relativeOffset))
if (message) {
messages.push(message)
}
}
}