fix(zmodem): finalize on sink errors, throttle progress, emit terminal events on detach

This commit is contained in:
Bill committed 2026-10-07 20:44:11 +08:00
1 parent 655c660607
commit 46b76ebca7
2 files changed
+454 -15

No files matched your search

+308
View File
@@ -8,6 +8,16 @@
* a temp dir; assert content matches + done(ok) + progress.
* 场景 2 (upload): our engine uploads a local temp file -> remote `rz`
* receives it; assert remote content matches + done(ok).
* 场景 3 (throttle): slow peer, 2KB every 3ms -> the running progress
* broadcasts must stay under the 150ms floor.
* 场景 4 (sink err): save dir does not exist -> createWriteStream errors, the
* engine must finalize at once (not at the 90s stall) and
* release the terminal data plane.
* 场景 5 (detach): session killed with an offer pending / a transfer in
* flight -> detachZmodem must emit ZMODEM_DONE + a terminal
* progress event, exactly once. Also asserts that cancelling
* a confirmed session reports ok=false, not a fake success
* (abort() fires 'session_end' synchronously).
*
* Pre-req (the bundle is built automatically on first run):
* npx esbuild src/main/zmodem.ts --bundle --platform=node --format=cjs \
@@ -102,6 +112,71 @@ async function driveRemoteSend(session, filePath) {
await session.close()
}
/**
* Same as driveRemoteSend, but paced: `chunkSize` bytes every `delayMs`. Keeps
* the transfer running long enough for the progress throttle to be measurable
* (zmodem.js caps a subpacket at 8192 bytes, so an unpaced transfer is over in
* a few milliseconds and produces almost no events either way).
*/
async function driveRemoteSendPaced(session, filePath, chunkSize, delayMs) {
const data = readFileSync(filePath)
const name = basename(filePath)
const xfer = await session.send_offer({ name, size: data.length, mtime: new Date() })
if (!xfer) throw new Error('offer was skipped')
let off = 0
while (off < data.length) {
const end = Math.min(off + chunkSize, data.length)
const slice = new Uint8Array(data.buffer, data.byteOffset + off, end - off)
if (end < data.length) {
xfer.send(slice)
await sleep(delayMs)
} else {
await xfer.end(slice)
}
off = end
}
await session.close()
}
/** Wait for the first event on `channel` (ms budget); returns its payload. */
async function waitForEvent(events, channel, budgetMs = 5000) {
const deadline = Date.now() + budgetMs
for (;;) {
const hit = events.find((e) => e.channel === channel)
if (hit) return hit.payload
if (Date.now() >= deadline) return undefined
await sleep(10)
}
}
/**
* Feed octets to the peer sentry. The peer's zsession throws "peer_aborted"
* once our engine writes its abort CAN, which is the expected outcome in the
* scenarios where the engine deliberately kills the transfer.
*/
function safeConsume(sentry, data) {
try {
sentry.consume(data)
} catch {
// peer aborted by our own CAN
}
}
/** Attach our engine to a fake ssh session, collecting every broadcast. */
function attachCollecting(sessionId, writeStream) {
const events = []
engine.attachZmodem(sessionId, {
broadcast: (channel, payload) => events.push({ channel, payload }),
toTerminal: () => {},
writeStream
})
return events
}
const progressOf = (events) =>
events.filter((e) => e.channel === 'transfer:progress').map((e) => e.payload)
const doneOf = (events) => events.filter((e) => e.channel === 'zmodem:done').map((e) => e.payload)
/** Read all bytes a remote receive side has drained so far. */
function collectRemote(self) {
return Buffer.concat(self._incoming)
@@ -256,11 +331,244 @@ async function scenarioUpload() {
console.log(`[upload] progress events: ${prog.length}, done(ok)=true`)
}
// ---- SCENARIO 3: progress throttle -------------------------------------------
async function scenarioThrottle() {
console.log('\n=== 场景 3: 进度节流 (download, 慢速对端) ===')
const sessionId = `sess-${randomBytes(3).toString('hex')}`
const srcDir = mkTemp('.zm-src-thr-')
const srcFile = join(srcDir, 'paced.bin')
const payload = randomBytes(300_000)
writeFileSync(srcFile, payload)
const outDir = mkTemp('.zm-out-thr-')
const remote = makeRemoteSentry(
(detection) => {
const s = detection.confirm()
if (s.type === 'send') {
// 2KB per tick => ~150 on_input calls: an unthrottled engine turns
// every one of them into a broadcast.
driveRemoteSendPaced(s, srcFile, 2048, 3).catch((e) => fail(`remote send error: ${e}`))
}
},
(buf) => setTimeout(() => engine.feedZmodem(sessionId, buf), 0)
)
const events = attachCollecting(sessionId, (id, data) =>
setTimeout(() => safeConsume(remote.sentry, data), 0)
)
engine.feedZmodem(sessionId, Buffer.from(zmodem.Header.build('ZRQINIT').to_hex()))
const offer = await waitForEvent(events, 'zmodem:offer')
if (!offer) fail('throttle: no ZMODEM_OFFER received')
const t0 = Date.now()
engine.respondZmodem({ id: offer.id, cancelled: false, dir: outDir })
const done = await waitForEvent(events, 'zmodem:done')
const elapsed = Date.now() - t0
if (!done) fail('throttle: no ZMODEM_DONE received')
if (!done.ok) fail(`throttle: expected ok, got ${JSON.stringify(done)}`)
const got = readFileSync(join(outDir, 'paced.bin'))
if (got.length !== payload.length || !got.equals(payload)) {
fail(`throttle: content mismatch (len ${got.length} vs ${payload.length})`)
}
const prog = progressOf(events)
// 150ms floor => at most one event per window, plus the deliberately forced
// file-boundary/final events (4) and slack for timer jitter. Without the
// throttle this is one event per subpacket (~150), so the gap is wide.
const bound = Math.ceil(elapsed / 150) + 8
if (prog.length > bound) {
fail(`throttle: ${prog.length} progress events in ${elapsed}ms exceeds bound ${bound}`)
}
if (prog.length < 2) fail(`throttle: only ${prog.length} progress events, throttle too aggressive`)
if (prog[prog.length - 1].state !== 'done') {
fail(`throttle: last progress state ${prog[prog.length - 1].state}, expected done`)
}
console.log(`[throttle] ${prog.length} progress events in ${elapsed}ms (bound ${bound}), content matches`)
}
// ---- SCENARIO 4: dead write sink finalizes at once ---------------------------
async function scenarioWriteError() {
console.log('\n=== 场景 4: 接收落盘失败 (目录不存在) -> 立即终结 ===')
const sessionId = `sess-${randomBytes(3).toString('hex')}`
const srcDir = mkTemp('.zm-src-err-')
const srcFile = join(srcDir, 'fail.bin')
writeFileSync(srcFile, randomBytes(120_000))
const outDir = mkTemp('.zm-out-err-')
// Never created: createWriteStream() emits ENOENT, the path a full disk or a
// revoked permission would take.
const badDir = join(outDir, 'no-such-dir')
const remote = makeRemoteSentry(
(detection) => {
const s = detection.confirm()
if (s.type === 'send') {
// Paced on purpose: an unpaced in-process transfer finishes before the
// fs.open ENOENT even surfaces, and the sink would only fail after the
// session had already ended successfully. The engine aborts as soon as
// the sink fails, so the peer dying here is the expected outcome, not a
// test failure.
driveRemoteSendPaced(s, srcFile, 4096, 5).catch(() => {})
}
},
(buf) => setTimeout(() => engine.feedZmodem(sessionId, buf), 0)
)
const events = attachCollecting(sessionId, (id, data) =>
setTimeout(() => safeConsume(remote.sentry, data), 0)
)
engine.feedZmodem(sessionId, Buffer.from(zmodem.Header.build('ZRQINIT').to_hex()))
const offer = await waitForEvent(events, 'zmodem:offer')
if (!offer) fail('write-error: no ZMODEM_OFFER received')
const t0 = Date.now()
engine.respondZmodem({ id: offer.id, cancelled: false, dir: badDir })
// 2s budget: the point is that this does NOT wait for the 90s stall watchdog.
const done = await waitForEvent(events, 'zmodem:done', 2000)
if (!done) fail('write-error: no ZMODEM_DONE within 2s (engine stuck active?)')
if (done.ok) {
fail(`write-error: expected ok=false, got ${JSON.stringify(done)} — sink error lost the race`)
}
if (engine.isZmodemActive(sessionId)) {
fail('write-error: engine still active — terminal data plane stays suppressed')
}
const prog = progressOf(events)
const last = prog[prog.length - 1]
if (last.state !== 'error') fail(`write-error: last progress state ${last.state}, expected error`)
console.log(`[write-error] done(ok=false) after ${Date.now() - t0}ms, state=error, engine inactive`)
}
// ---- SCENARIO 5: detach (session killed) emits the terminal events ------------
async function scenarioDetach() {
console.log('\n=== 场景 5: 会话被 kill -> detach 补发终结事件 ===')
// (a) offer still pending: the renderer's modal is up and nothing will answer
{
const sessionId = `sess-${randomBytes(3).toString('hex')}`
const events = attachCollecting(sessionId, () => {})
engine.feedZmodem(sessionId, Buffer.from(zmodem.Header.build('ZRQINIT').to_hex()))
const offer = await waitForEvent(events, 'zmodem:offer')
if (!offer) fail('detach(pending): no ZMODEM_OFFER received')
engine.detachZmodem(sessionId)
const done = doneOf(events)
if (done.length !== 1) fail(`detach(pending): expected 1 ZMODEM_DONE, got ${done.length}`)
if (done[0].ok) fail('detach(pending): expected ok=false')
if (engine.isZmodemActive(sessionId)) fail('detach(pending): engine still active')
const prog = progressOf(events)
if (prog[prog.length - 1].state !== 'cancelled') {
fail(`detach(pending): last progress state ${prog[prog.length - 1].state}, expected cancelled`)
}
// Idempotent: the engine is already gone, a second detach emits nothing.
engine.detachZmodem(sessionId)
if (doneOf(events).length !== 1) fail('detach(pending): ZMODEM_DONE emitted twice')
console.log('[detach] pending offer -> done(ok=false) + cancelled progress, idempotent')
}
// (b) transfer in flight
{
const sessionId = `sess-${randomBytes(3).toString('hex')}`
const srcDir = mkTemp('.zm-src-det-')
const srcFile = join(srcDir, 'killed.bin')
writeFileSync(srcFile, randomBytes(200_000))
const outDir = mkTemp('.zm-out-det-')
let detached = false
const remote = makeRemoteSentry(
(detection) => {
const s = detection.confirm()
if (s.type === 'send') {
// Killed mid-flight on purpose: a dying peer is expected here.
driveRemoteSendPaced(s, srcFile, 2048, 3).catch(() => {})
}
},
(buf) => {
if (detached) return
setTimeout(() => engine.feedZmodem(sessionId, buf), 0)
}
)
const events = attachCollecting(sessionId, (id, data) => {
if (detached) return
setTimeout(() => remote.sentry.consume(data), 0)
})
engine.feedZmodem(sessionId, Buffer.from(zmodem.Header.build('ZRQINIT').to_hex()))
const offer = await waitForEvent(events, 'zmodem:offer')
if (!offer) fail('detach(live): no ZMODEM_OFFER received')
engine.respondZmodem({ id: offer.id, cancelled: false, dir: outDir })
if (!engine.isZmodemActive(sessionId)) fail('detach(live): engine not active after respond')
const mark = events.length
detached = true
engine.detachZmodem(sessionId)
await sleep(300)
const done = doneOf(events)
if (done.length !== 1) fail(`detach(live): expected 1 ZMODEM_DONE, got ${done.length}`)
if (done[0].ok) fail('detach(live): expected ok=false')
if (engine.isZmodemActive(sessionId)) fail('detach(live): engine still active')
const after = progressOf(events.slice(mark))
if (after.length === 0) fail('detach(live): no terminal progress emitted')
if (after[after.length - 1].state !== 'cancelled') {
fail(`detach(live): last progress state ${after[after.length - 1].state}, expected cancelled`)
}
if (after.some((p) => p.state === 'running')) {
fail('detach(live): progress still running after detach')
}
engine.detachZmodem(sessionId)
if (doneOf(events).length !== 1) fail('detach(live): ZMODEM_DONE emitted twice')
console.log('[detach] live transfer -> done(ok=false) + cancelled progress, exactly once')
}
// (c) user cancels a *confirmed* session (no files chosen). abortSession()
// aborts before it finalizes, and the abort fires 'session_end' synchronously:
// without the engine.ending guard that re-entered finalize(ok=true) and the
// cancel was reported as a successful transfer.
{
const sessionId = `sess-${randomBytes(3).toString('hex')}`
const remote = makeRemoteSentry(
(detection) => {
const s = detection.confirm()
if (s.type === 'receive') s.start().catch(() => {})
},
(buf) => setTimeout(() => engine.feedZmodem(sessionId, buf), 0)
)
const events = attachCollecting(sessionId, (id, data) =>
setTimeout(() => safeConsume(remote.sentry, data), 0)
)
// Remote rz initiates: it detects "receive" from ZRQINIT, confirms and
// start()s, which sends ZRINIT toward our engine.
remote.sentry.consume(Buffer.from(zmodem.Header.build('ZRQINIT').to_hex()))
const offer = await waitForEvent(events, 'zmodem:offer')
if (!offer) fail('detach(cancel): no ZMODEM_OFFER received')
if (offer.mode !== 'send') fail(`detach(cancel): expected send offer, got ${offer.mode}`)
engine.respondZmodem({ id: offer.id, cancelled: false, paths: [] })
const done = doneOf(events)
if (done.length !== 1) fail(`detach(cancel): expected 1 ZMODEM_DONE, got ${done.length}`)
if (done[0].ok) fail('detach(cancel): a cancelled transfer was reported as success')
if (engine.isZmodemActive(sessionId)) fail('detach(cancel): engine still active')
const prog = progressOf(events)
if (prog[prog.length - 1].state !== 'cancelled') {
fail(`detach(cancel): last progress state ${prog[prog.length - 1].state}, expected cancelled`)
}
console.log('[detach] confirmed session cancelled -> done(ok=false), no fake success')
}
}
// ---- run ----------------------------------------------------------------------
const timer = setTimeout(() => fail('zmodem-e2e timed out'), 30000)
try {
await scenarioDownload()
await scenarioUpload()
await scenarioThrottle()
await scenarioWriteError()
await scenarioDetach()
} finally {
clearTimeout(timer)
cleanupTempDirs()