diff --git a/src/main/sftp.ts b/src/main/sftp.ts index b4c17b4..1c319ad 100644 --- a/src/main/sftp.ts +++ b/src/main/sftp.ts @@ -325,12 +325,58 @@ export function renameRemote(sessionId: string, from: string, to: string): Promi return withSftp(sessionId, sftp => pVoid(cb => sftp.rename(from, to, cb))) } +/** + * "The file is already gone" — the one server-side failure a delete may treat + * as success. ssh2 reports the SFTP status both ways: `code` carries the numeric + * status (2 = SSH_FX_NO_SUCH_FILE) and `message` the library's own English text, + * so both are checked (builds fill one or the other). Our own translated errors + * never reach this check — they arrive as OperationError, which the caller + * passes through untouched. + */ +function isNoSuchFile(err: unknown): boolean { + if ((err as { code?: unknown } | undefined)?.code === 2) return true + return /no such file/i.test((err as Error | undefined)?.message ?? '') +} + +/** + * unlink/rmdir where "already gone" counts as success. + * + * Deleting is not idempotent by nature: the entry may have been removed by + * someone else between the listing and the click, or the same delete may be run + * again after an earlier pass already removed it. Reporting a path that is + * verifiably gone as "delete failed" is a fake error — the end state the user + * asked for already holds. + */ +async function unlinkIfPresent(sftp: SFTPWrapper, path: string): Promise { + try { + await pVoid(cb => sftp.unlink(path, cb)) + } catch (err) { + if (!isNoSuchFile(err)) throw err + } +} + +async function rmdirIfPresent(sftp: SFTPWrapper, path: string): Promise { + try { + await pVoid(cb => sftp.rmdir(path, cb)) + } catch (err) { + if (!isNoSuchFile(err)) throw err + } +} + async function deleteRecursive(sftp: SFTPWrapper, path: string, isDir: boolean): Promise { if (!isDir) { - await pVoid(cb => sftp.unlink(path, cb)) + await unlinkIfPresent(sftp, path) return } - const names = await readdir(sftp, path) + let names: string[] + try { + names = await readdir(sftp, path) + } catch (err) { + // The directory is gone (removed by someone else, or by an earlier pass of + // this same delete): everything inside went with it, so the intent holds. + if (isNoSuchFile(err)) return + throw err + } const base = path.replace(/\/+$/, '') || '/' for (const name of names) { const child = `${base}/${name}` @@ -342,7 +388,7 @@ async function deleteRecursive(sftp: SFTPWrapper, path: string, isDir: boolean): } await deleteRecursive(sftp, child, childIsDir) } - await pVoid(cb => sftp.rmdir(path, cb)) + await rmdirIfPresent(sftp, path) } export function deleteRemote(sessionId: string, paths: string[]): Promise { @@ -354,7 +400,7 @@ export function deleteRemote(sessionId: string, paths: string[]): Promise try { isDir = (await lstat(sftp, path)).isDirectory() } catch { - // stat failed — fall through to unlink + // stat failed — fall through to unlink, which tolerates "already gone" } await deleteRecursive(sftp, path, isDir) } catch (err) { diff --git a/src/main/sysinfo.ts b/src/main/sysinfo.ts index b5351dd..623bc3c 100644 --- a/src/main/sysinfo.ts +++ b/src/main/sysinfo.ts @@ -157,7 +157,17 @@ function pollOnce(id: string, state: PollState): void { try { client.exec(COLLECT_CMD, (err: Error | undefined, stream) => { - if (state.stopped) return + if (state.stopped) { + // Polling stopped while this exec was in flight: the channel is already + // open, so simply dropping the stream would leak it for the rest of the + // ssh session. Close it (best effort — it may be dead already) and go. + try { + stream?.close() + } catch { + // best effort + } + return + } if (err || !stream) { handleError(id, state, err?.message || t('main.sysinfo.execFailed')) return diff --git a/src/main/zmodem.ts b/src/main/zmodem.ts index 36b4b7d..4c60b14 100644 --- a/src/main/zmodem.ts +++ b/src/main/zmodem.ts @@ -65,6 +65,19 @@ export interface ZmodemDeps { toTerminal(sessionId: string, data: Buffer): void /** Write bytes out to the ssh stream (zmodem frames + CAN abort sequence). */ writeStream(sessionId: string, data: Buffer): void + /** + * Backpressure over the ssh stream that FEEDS this session's sentry. pty.ts + * owns that channel and injects these two, mirroring writeStream: the engine + * has no handle on it, and pausing the source is the only thing that bounds + * how much a peer can push into a file sink that is over its highWaterMark. + * + * Optional so a harness with no real stream still runs (and so an injector + * that predates this hook keeps working): without it a backed-up sink still + * stops counting progress, which turns a wedged disk into a stall-watchdog + * failure instead of an unbounded queue — but the queue itself stays. + */ + pauseSource?(sessionId: string): void + resumeSource?(sessionId: string): void } interface Engine { @@ -94,7 +107,18 @@ interface Engine { offerTimer: NodeJS.Timeout | null stallTimer: NodeJS.Timeout | null receiveStream: Writable | null + /** + * Octets the current file's sink ACCEPTED. A write that only landed in the + * WriteStream's own queue (write() returned false) is not counted here until + * the queue drains — see pendingBytes and the receive on_input. + */ bytes: number + /** Octets handed to a backed-up sink; folded into `bytes` on drain/teardown. */ + pendingBytes: number + /** True while the source is paused because the sink is over its highWaterMark. */ + sourcePaused: boolean + /** One-shot: the injector supplied no pauseSource, so the pause is a no-op. */ + warnedNoSource: boolean totalBytes: number } @@ -176,6 +200,9 @@ function finalize( // mistaken for a clean finish. engine.ending = true cancelTimers(engine) + // A pause must never outlive its transfer: the 'drain' that would lift it may + // still be queued, or may never come at all once the sink is ending here. + releaseSource(engine) const stream = engine.receiveStream engine.receiveStream = null // Clear the slot *before* end(): an errored/destroyed sink must not be ended @@ -235,6 +262,58 @@ function touchActivity(engine: Engine): void { engine.stallTimer = timer } +/** + * Stop the ssh stream that feeds this engine (receive path only: the send path + * pushes into the sink instead of pulling from the peer). + * + * Idempotent — every subpacket of a blocked window arrives through the same + * on_input, and pausing an already-paused channel is pointless work on the hot + * path. The flag is set even when the injector supplied no hook, so the + * accounting below stays consistent either way. + */ +function pauseSource(engine: Engine): void { + if (engine.sourcePaused) return + engine.sourcePaused = true + if (!engine.deps.pauseSource) { + // Say it once per transfer instead of silently doing nothing: the hook is + // the whole mechanism, so an injector that forgot it still has an unbounded + // sink queue (only the progress accounting below improves). + if (!engine.warnedNoSource) { + engine.warnedNoSource = true + console.warn('[zmodem] receive sink is over its highWaterMark but no pauseSource hook was injected') + } + return + } + try { + engine.deps.pauseSource(engine.id) + } catch { + // the stream may already be gone + } +} + +/** + * Undo a backpressure pause, folding the octets that only ever reached the + * sink's queue into the received count. + * + * Called from the sink's 'drain' AND from every teardown path: a pause that + * outlived its transfer would leave the session's ssh stream stopped for good + * (terminal output and the peer both look frozen, with no way back), and + * waiting for a 'drain' that never comes — dead disk, destroyed stream — is + * exactly how that happens. Returns whether a pause was actually outstanding. + */ +function releaseSource(engine: Engine): boolean { + if (!engine.sourcePaused) return false + engine.sourcePaused = false + engine.bytes += engine.pendingBytes + engine.pendingBytes = 0 + try { + engine.deps.resumeSource?.(engine.id) + } catch { + // the stream may already be gone + } + return true +} + function wireSession(engine: Engine): void { const session = engine.session if (!session) return @@ -383,23 +462,40 @@ function handleOffer(engine: Engine, offer: Zmodem.Offer): void { finalize(engine, false, err instanceof Error ? err.message : t('main.zmodem.receiveFailed')) }) + // Backpressure release. The sink's queue is empty again, so the octets that + // were parked there are really written, the peer may push more, and this is + // the ONLY progress signal while blocked — it must therefore also reset the + // stall watchdog (a slow but working disk is not a stalled transfer). + stream.on('drain', () => { + if (engine.receiveStream !== stream) return + if (releaseSource(engine)) touchActivity(engine) + }) + offer .accept({ on_input: (payload: Uint8Array | number[]) => { // zmodem.js delivers the payload as a plain octet Array (not a // Uint8Array); coerce defensively to a Buffer before writing. const buf = Array.isArray(payload) ? Buffer.from(payload) : Buffer.from(payload.buffer, payload.byteOffset, payload.byteLength) - if (!stream.destroyed && !stream.closed) { - stream.write(buf) + const writable = !stream.destroyed && !stream.closed + if (writable && stream.write(buf)) { + engine.bytes += buf.byteLength + // Not throttled: the stall watchdog measures byte progress, not UI events. + touchActivity(engine) + } else if (writable) { + // write() === false: the sink is over its highWaterMark, so these + // octets only sit in Node's queue — for a 100MB file into a slow disk + // the whole file would end up there. Pausing the ssh stream makes the + // peer wait instead; counting the bytes or touching the watchdog now + // would report progress for data that has not reached the disk yet. + engine.pendingBytes += buf.byteLength + pauseSource(engine) } - engine.bytes += buf.byteLength emitProgressThrottled(engine, { file: name, bytes: engine.bytes, totalBytes: engine.totalBytes }) - // Not throttled: the stall watchdog measures byte progress, not UI events. - touchActivity(engine) } }) .then(() => { @@ -521,6 +617,9 @@ export function attachZmodem(sessionId: string, deps: ZmodemDeps): void { stallTimer: null, receiveStream: null, bytes: 0, + pendingBytes: 0, + sourcePaused: false, + warnedNoSource: false, totalBytes: 0 } engines.set(sessionId, engine) @@ -683,6 +782,10 @@ export function detachZmodem(sessionId: string): void { if (engine.active) { finalize(engine, false, t('main.zmodem.sessionClosed'), 'cancelled') } + // The engine is about to be dropped: a pause still outstanding here would + // never be lifted (finalize already released the active one; this covers the + // path where it was not active). + releaseSource(engine) engine.active = false cancelTimers(engine) const stream = engine.receiveStream diff --git a/tests/sftp-timeout.mjs b/tests/sftp-timeout.mjs index 5b66d9f..f0b6249 100644 --- a/tests/sftp-timeout.mjs +++ b/tests/sftp-timeout.mjs @@ -51,8 +51,9 @@ sftp.setLocalPathPolicy({ /** * A fake ssh2 SFTPWrapper whose per-call behavior is scripted. * `behaviour(op, args)` returns `'silent'` (never calls back — the dead-channel - * case), `'error'` (calls back with an error), or `'ok'` (calls back with - * `result`). Every call is recorded so the test can assert on what was opened. + * case), `'enoent'` (calls back with the server's "no such file"), `'error'` + * (calls back with a generic error), or `'ok'` (calls back with `result`). Every + * call is recorded so the test can assert on what was opened. */ const defaultStat = () => ({ isDirectory: () => false, @@ -63,6 +64,9 @@ const defaultStat = () => ({ gid: 1000 }) +/** SSH_FX_NO_SUCH_FILE as ssh2 surfaces it: a numeric status on an Error. */ +const noSuchFile = () => Object.assign(new Error('No such file'), { code: 2 }) + let calls let behaviour let openedChannels @@ -101,6 +105,7 @@ const makeWrapper = () => { calls.push({ op, args }) const verdict = behaviour(op, args) if (verdict === 'silent') return + if (verdict === 'enoent') return cb(noSuchFile()) if (verdict === 'error') return cb(new Error(`${op} failed`)) cb(null, voidResult ? undefined : result) } @@ -433,6 +438,77 @@ console.log('default budgets are sane relative to each other') ok(elapsed < 3000, `the check itself was quick (${elapsed}ms)`) } +// ---- 7. a delete is idempotent: "already gone" is success ------------------ +console.log('deleting an entry that is already gone succeeds instead of reporting a failure') +{ + // The file may have been removed by someone else between the listing and the + // click, or the same delete may be run again after an earlier pass already + // removed it. "No such file" means the end state the user asked for already + // holds, so it must not surface as "delete failed". + reset((op) => (op === 'lstat' || op === 'unlink' ? 'enoent' : 'ok')) + sftp.closeSftp('del-gone') + const missing = await expectReject(sftp.deleteRemote('del-gone', ['/remote/gone.txt']), 'deleteRemote') + ok(!missing.rejected, `a missing file is not a delete failure (${missing.message})`) + + // The tolerance is for that one status only: a real unlink failure still has + // to reach the user. + reset((op) => (op === 'unlink' ? 'error' : 'ok')) + sftp.closeSftp('del-denied') + const denied = await expectReject(sftp.deleteRemote('del-denied', ['/remote/x.txt']), 'deleteRemote') + ok(denied.rejected, 'a genuine unlink failure is still reported') + + // The other shape of the same status: no numeric code, only the library's + // English message. + reset(() => 'ok') + sftp.registerSftpClientProvider(() => ({ + sftp: (cb) => { + const w = makeWrapper() + w.unlink = (p, cb) => { + calls.push({ op: 'unlink', args: [p] }) + cb(new Error('No such file')) + } + openedChannels.push(w) + cb(null, w) + } + })) + sftp.closeSftp('del-gone-msg') + const byMessage = await expectReject(sftp.deleteRemote('del-gone-msg', ['/remote/gone.txt']), 'deleteRemote') + ok(!byMessage.rejected, `the message-only form is tolerated too (${byMessage.message})`) + + // The recursive shapes: the directory itself is gone, and a tree whose entries + // were already removed reports the same status from unlink and rmdir. + reset((op, args) => (op === 'readdir' && args[0] === '/remote/dir' ? 'enoent' : 'ok')) + sftp.registerSftpClientProvider(() => ({ + sftp: (cb) => { + const w = makeWrapper() + // Only the top-level path is a directory, or the children would recurse. + w.lstat = (p, cb) => { + calls.push({ op: 'lstat', args: [p] }) + cb(null, { ...defaultStat(), isDirectory: () => p === '/remote/dir' }) + } + openedChannels.push(w) + cb(null, w) + } + })) + sftp.closeSftp('del-dir-gone') + const dirGone = await expectReject(sftp.deleteRemote('del-dir-gone', ['/remote/dir']), 'deleteRemote') + ok(!dirGone.rejected, `a directory that no longer exists is not a failure (${dirGone.message})`) + + reset((op) => (op === 'unlink' || op === 'rmdir' ? 'enoent' : 'ok')) + sftp.closeSftp('del-dir-partial') + const dirPartial = await expectReject(sftp.deleteRemote('del-dir-partial', ['/remote/dir']), 'deleteRemote') + ok(!dirPartial.rejected, `a tree whose entries were already removed is not a failure (${dirPartial.message})`) + + // Restore the standard provider for any later section. + sftp.registerSftpClientProvider(() => ({ + sftp: (cb) => { + const w = makeWrapper() + openedChannels.push(w) + cb(null, w) + } + })) +} + // ---- helpers ---------------------------------------------------------------- function waitFor(pred, timeoutMs) { return new Promise((resolve) => {