mirror of
https://github.com/stablyai/orca.git
synced 2026-10-08 08:02:32 +00:00
* Revert "feat(orcad): source-side dormant export of a relay-hosted SSH target (#16741 T6-8) (#24519)" This reverts commit783101b304. * Revert "feat(ssh): update, roll back, recover and stop a managed orcad server (#16741 T6-5 follow-up) (#24463)" This reverts commit38c2d1dcb9. * Revert "feat(ssh): deploy and pair an empty managed orcad server over SSH (#16741 T6-5) (#24453)" This reverts commit8b76683b40. * Revert "fix(ssh): orcad GC honors the activation journal; readiness requires proven daemon coverage (#16741 T6 follow-up) (#24451)" This reverts commitd3f8c5063b. * Revert "feat(ssh): remote orcad stop by request file and journaled decommission (#16741 T6-4) (#24449)" This reverts commit43d9b43d3f. * Revert "feat(orcad): supervisable server: stop requests, managed stop receipts and a lifetime that keeps its lock on failed teardown (#16741 T6-3) (#24433)" This reverts commitb093d3ab20. * Revert "feat(ssh): crash-safe orcad activation, rollback and recovery (#16741 T6-2) (#24423)" This reverts commit1a9ac0e955. * Revert "feat(runtime): SSH access links for paired servers in a downgrade-safe sidecar (#16741 T5-1+T5-2) (#24420)" This reverts commit99db2bfae4. * Revert "feat(relay): capability-gated owner reset with a durable preparation journal (#16741 T3 R1) (#24418)" This reverts commit34a582bd39. * Revert "feat(ssh): track connection-manager drains, test probes and provider continuations (#16741 T2 P3+P8a) (#24407)" This reverts commitd53063d2b1. * Revert "feat(daemon): idle retirement, session census and recovery-only provider (#16741 T2 P4b) (#24409)" This reverts commitff212dbbef. * Revert "feat(ssh): add pty.resumeClient and split SSH PTY process listing (#16741 T2 P5+P6) (#24414)" This reverts commit92cb71765e. * Revert "feat(relay): await owned watcher and agent children on shutdown (#16741 T2 P1) (#24400)" This reverts commit6b36e4f85b. * Revert "feat(session): retry failed renderer session writes and verify local folder PTYs (#16741 T2 P9) (#24406)" This reverts commitd23ecef301. * Revert "feat(ssh): remote orcad primitives on the pinned Node runtime (#16741 T6-1) (#24419)" This reverts commitdd87ae578d. * Revert "fix(runtime): fence runtime-environment subscriptions and status probes by identity (#16741 T5-3) (#24421)" This reverts commitece9e4d2e3. * Revert "feat(orcad): migration manifest and dormant-state contracts (#16741 T6-7) (#24422)" This reverts commit3fbdaba262. * Revert "feat(ssh): wire SshConnection through the work and transport close ledgers (#16741 T2 P2) (#24401)" This reverts commit4e8edc8872. * Revert "feat(profiles): carry markdown frontmatter visibility in project transfers (#16741 T2 P7) (#24405)" This reverts commit60c93263cc. * Revert "fix(runtime): project the PTY incarnation onto mobile session tabs (#24413)" This reverts commit99e0303572. * Revert "feat(daemon): tag daemon stream data with the PTY incarnation id (#16741 T2 P4a) (#24402)" This reverts commit817af768b0. * Revert "feat(ssh): port the SSH connection work ledger and transport close ledger (#16741 T2) (#24210)" This reverts commitc9918931c8. * Revert "feat(relay): fence and drain file and git response streams on shutdown (#24185)" This reverts commitdc08ffeba9. * Revert "refactor(runtime-rpc): extract the Node WebSocket lifecycle; opt-in pinned port (#24186)" This reverts commita789233bbb. * Revert "feat(relay): route relay handlers through work admission; producer publication drain (#24181)" This reverts commit0b812bd698. * Revert "feat(relay): land the #16741 T1 seam (work drain, publication drain, release gate) (#24156)" This reverts commit3aa2d3af7c. --------- Co-authored-by: m4air <m4air@Mac.localdomain>
219 lines
7.9 KiB
TypeScript
219 lines
7.9 KiB
TypeScript
import type { Readable } from 'node:stream'
|
|
import type { ClientChannel } from 'ssh2'
|
|
import type { SshConnection } from './ssh-connection'
|
|
import { createSshOperationAbortError, type SshExecOptions } from './ssh-connection-utils'
|
|
import {
|
|
redactRelayInstallMarkerError,
|
|
redactRelayInstallMarkerTokens
|
|
} from './ssh-relay-install-marker'
|
|
import type { SystemSshCommandChannel } from './system-ssh-command'
|
|
|
|
const EXEC_TIMEOUT_MS = 30_000
|
|
const COMMAND_CLOSE_GRACE_MS = 5_000
|
|
const MAX_EXEC_OUTPUT_CHARS = 1024 * 1024
|
|
|
|
type ExecCommandOptions = SshExecOptions & {
|
|
timeoutMs?: number
|
|
// Why: a zero-exit command resolves with stdout alone, so the reason a wrapped-in-`|| echo`
|
|
// probe failed is discarded. Callers that need that diagnostic opt in here rather than
|
|
// folding stderr into stdout, where it would match the probe's own token strings.
|
|
// On the system-ssh transport this stream also carries local OpenSSH noise; log-only.
|
|
onStderr?: (stderr: string) => void
|
|
/** Streamed into the command's stdin, then EOF. A read error terminates the command. */
|
|
stdin?: Readable
|
|
}
|
|
|
|
type SshCommandTerminationError = Error & {
|
|
sshChannelCloseConfirmed: boolean
|
|
}
|
|
|
|
// Why: callers must tell "the host answered no" from "the host never answered". Matching the
|
|
// message text is what let an unanswered probe be read as a definitive negative.
|
|
export const SSH_EXEC_TIMEOUT_CODE = 'SSH_EXEC_TIMEOUT'
|
|
|
|
export function isSshExecTimeout(error: unknown): boolean {
|
|
return (
|
|
error instanceof Error && (error as Partial<{ code: string }>).code === SSH_EXEC_TIMEOUT_CODE
|
|
)
|
|
}
|
|
|
|
export function isUnconfirmedSshCommandTermination(
|
|
error: unknown
|
|
): error is SshCommandTerminationError {
|
|
return (
|
|
error instanceof Error &&
|
|
(error as Partial<SshCommandTerminationError>).sshChannelCloseConfirmed === false
|
|
)
|
|
}
|
|
|
|
export async function execCommand(
|
|
conn: SshConnection,
|
|
command: string,
|
|
options?: ExecCommandOptions
|
|
): Promise<string> {
|
|
const { timeoutMs = EXEC_TIMEOUT_MS, onStderr, stdin, ...execOptions } = options ?? {}
|
|
const signal = options?.signal
|
|
if (signal?.aborted) {
|
|
throw createSshOperationAbortError()
|
|
}
|
|
// Why: reconnect/disconnect can flip the connection back to ssh2 before a
|
|
// killed local OpenSSH child emits close; the channel's transport is immutable.
|
|
const openedWithSystemSsh = conn.usesSystemSshTransport?.() === true
|
|
let channel: ClientChannel
|
|
try {
|
|
channel = await conn.exec(command, execOptions)
|
|
} catch (error) {
|
|
// Preserve identity/classifier fields while removing install-owner tokens.
|
|
redactRelayInstallMarkerError(error)
|
|
throw error
|
|
}
|
|
return new Promise((resolve, reject) => {
|
|
let stdout = ''
|
|
let stderr = ''
|
|
let settled = false
|
|
let terminationError: SshCommandTerminationError | null = null
|
|
let closeGraceTimer: ReturnType<typeof setTimeout> | null = null
|
|
|
|
const cleanup = (): void => {
|
|
clearTimeout(timeout)
|
|
if (closeGraceTimer) {
|
|
clearTimeout(closeGraceTimer)
|
|
}
|
|
signal?.removeEventListener('abort', onAbort)
|
|
channel.off('error', fail)
|
|
channel.stderr.off('error', fail)
|
|
channel.off('data', onStdoutData)
|
|
channel.stderr.off('data', onStderrData)
|
|
channel.off('close', onClose)
|
|
if (stdin) {
|
|
stdin.off('error', fail)
|
|
stdin.unpipe(channel.stdin)
|
|
stdin.destroy()
|
|
}
|
|
}
|
|
const settle = (fn: typeof resolve | typeof reject, val: string | Error): void => {
|
|
if (settled) {
|
|
return
|
|
}
|
|
settled = true
|
|
cleanup()
|
|
fn(val as never)
|
|
}
|
|
const guardUnconfirmedTeardown = (): void => {
|
|
const swallowLateError = (): void => {}
|
|
const cleanupGuards = (): void => {
|
|
channel.off('error', swallowLateError)
|
|
channel.stderr.off('error', swallowLateError)
|
|
channel.off('close', cleanupGuards)
|
|
}
|
|
channel.on('error', swallowLateError)
|
|
channel.stderr.on('error', swallowLateError)
|
|
channel.once('close', cleanupGuards)
|
|
// Why: an unconfirmed close can arrive after the caller's bounded wait;
|
|
// keep draining discarded streams so ssh2 can finish CHANNEL_CLOSE.
|
|
channel.resume()
|
|
channel.stderr.resume()
|
|
}
|
|
// Why: sshd counts the session against MaxSessions until CHANNEL_CLOSE
|
|
// completes. Settling on abort before the channel actually closes lets the
|
|
// concurrent-bootstrap sequential fallback reissue an exec while the slot
|
|
// is still held, so it gets refused again. Close and settle from onClose.
|
|
const requestTermination = (error: Error): void => {
|
|
if (terminationError) {
|
|
return
|
|
}
|
|
terminationError = Object.assign(error, { sshChannelCloseConfirmed: false })
|
|
clearTimeout(timeout)
|
|
// Why: callers must not release an install lock while its remote npm
|
|
// process can still mutate node_modules. Prefer confirmed channel close,
|
|
// but bound a broken transport's teardown wait.
|
|
closeGraceTimer = setTimeout(() => {
|
|
guardUnconfirmedTeardown()
|
|
settle(reject, error)
|
|
}, COMMAND_CLOSE_GRACE_MS)
|
|
channel.close()
|
|
}
|
|
const fail = (err: Error): void => {
|
|
redactRelayInstallMarkerError(err)
|
|
requestTermination(err)
|
|
}
|
|
const onAbort = (): void => requestTermination(createSshOperationAbortError())
|
|
const onStdoutData = (data: Buffer): void => {
|
|
stdout = appendExecOutputTail(stdout, data.toString('utf-8'))
|
|
}
|
|
const onStderrData = (data: Buffer): void => {
|
|
stderr = appendExecOutputTail(stderr, data.toString('utf-8'))
|
|
}
|
|
const onClose = (code: number): void => {
|
|
if (
|
|
!terminationError &&
|
|
openedWithSystemSsh &&
|
|
(channel as SystemSshCommandChannel)._closeRequested
|
|
) {
|
|
terminationError = Object.assign(createSshOperationAbortError(), {
|
|
sshChannelCloseConfirmed: false
|
|
})
|
|
}
|
|
if (terminationError) {
|
|
// Why: a system-SSH channel closes when the local OpenSSH child exits;
|
|
// that does not prove the remote command stopped, especially with a ControlMaster.
|
|
if (!openedWithSystemSsh) {
|
|
terminationError.sshChannelCloseConfirmed = true
|
|
}
|
|
settle(reject, terminationError)
|
|
} else if (code !== 0) {
|
|
// Why: on the system-ssh transport channel.stderr carries local OpenSSH
|
|
// client noise; preferring it masks the real failure in stdout (2>&1).
|
|
const output = redactRelayInstallMarkerTokens(
|
|
[stderr.trim(), stdout.trim()].filter(Boolean).join('\n')
|
|
)
|
|
settle(
|
|
reject,
|
|
new Error(
|
|
`Command "${redactRelayInstallMarkerTokens(command)}" failed (exit ${code}): ${output}`
|
|
)
|
|
)
|
|
} else {
|
|
if (stderr && onStderr) {
|
|
onStderr(redactRelayInstallMarkerTokens(stderr))
|
|
}
|
|
settle(resolve, stdout)
|
|
}
|
|
}
|
|
const timeout = setTimeout(() => {
|
|
requestTermination(
|
|
Object.assign(
|
|
new Error(
|
|
`Command "${redactRelayInstallMarkerTokens(command)}" timed out after ${
|
|
timeoutMs / 1000
|
|
}s`
|
|
),
|
|
{ code: SSH_EXEC_TIMEOUT_CODE }
|
|
)
|
|
)
|
|
}, timeoutMs)
|
|
|
|
// Why: remote reboot tears down exec channels with stream errors. Without
|
|
// scoped listeners, Node treats those as uncaught exceptions.
|
|
signal?.addEventListener('abort', onAbort, { once: true })
|
|
channel.on('error', fail)
|
|
channel.stderr.on('error', fail)
|
|
channel.on('data', onStdoutData)
|
|
channel.stderr.on('data', onStderrData)
|
|
channel.on('close', onClose)
|
|
if (stdin) {
|
|
stdin.on('error', fail)
|
|
// Why `.stdin`: ssh2 aliases it to the channel, and a system-ssh channel only ends it there.
|
|
stdin.pipe(channel.stdin)
|
|
}
|
|
if (signal?.aborted) {
|
|
onAbort()
|
|
}
|
|
})
|
|
}
|
|
|
|
function appendExecOutputTail(existing: string, chunk: string): string {
|
|
const combined = existing + chunk
|
|
return combined.length > MAX_EXEC_OUTPUT_CHARS ? combined.slice(-MAX_EXEC_OUTPUT_CHARS) : combined
|
|
}
|