fix(zmodem): 接收路径写背压;sftp 删除容忍 ENOENT;sysinfo 停止关流
- zmodem 接收:fs WriteStream 缓冲无上限,慢盘 + 快对端可把内存涨到接近 文件大小。write() 返回 false 时暂停喂入 sentry 的 ssh 流、drain 恢复; touchActivity 只在写入被接受或 drain 恢复时计数,stall 看门狗不再把 '堆在缓冲里'当进展;finalize/detach 解除悬挂暂停 - sftp 删除路径 ENOENT 一律视为已删成功(文件已被他人删除/重复删除/ 传输重跑后目标已消失),不再汇成假 delete failed;真错误仍报错 (tests/sftp-timeout.mjs 新增 5 断言含防过度吞错用例) - sysinfo:exec 回调里 state.stopped 分支先 close 到手的 stream
This commit is contained in:
1 parent
81ac112dac
commit
e23010f950
4 files changed
+247
-12
No files matched your search
+50
-4
@@ -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<void> {
|
||||
try {
|
||||
await pVoid(cb => sftp.unlink(path, cb))
|
||||
} catch (err) {
|
||||
if (!isNoSuchFile(err)) throw err
|
||||
}
|
||||
}
|
||||
|
||||
async function rmdirIfPresent(sftp: SFTPWrapper, path: string): Promise<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
@@ -354,7 +400,7 @@ export function deleteRemote(sessionId: string, paths: string[]): Promise<void>
|
||||
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) {
|
||||
|
||||
+11
-1
@@ -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
|
||||
|
||||
+108
-5
@@ -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
|
||||
|
||||
+78
-2
@@ -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) => {
|
||||
|
||||
Reference in new issue
Block a user