Files
gitea-macos-vm-orchestrator/Sources/RunnerCore/SSHExec.swift
T

590 lines
23 KiB
Swift

import Foundation
import NIOCore
import NIOPosix
import NIOSSH
/// The outcome of a command run inside a guest.
public struct SSHCommandResult: Sendable, Equatable {
/// The remote process's exit status. `0` on success.
public let exitCode: Int32
/// Captured stdout, UTF-8 decoded with lossy replacement.
public let stdout: String
/// Captured stderr, UTF-8 decoded with lossy replacement.
public let stderr: String
public init(exitCode: Int32, stdout: String, stderr: String) {
self.exitCode = exitCode
self.stdout = stdout
self.stderr = stderr
}
/// Whether the command exited zero.
public var succeeded: Bool { exitCode == 0 }
/// Throws ``CoreError/sshFailed(_:)`` unless the command exited zero.
///
/// - Parameter command: Echoed into the error message for context.
public func throwIfFailed(command: String) throws {
guard exitCode != 0 else { return }
let detail = stderr.isEmpty ? stdout : stderr
let trimmed = detail.trimmingCharacters(in: .whitespacesAndNewlines)
throw CoreError.sshFailed(
"command failed (exit \(exitCode)): \(command)" + (trimmed.isEmpty ? "" : "\n\(trimmed)")
)
}
}
/// The seam for talking to a guest.
///
/// The orchestrator and the provisioner are written against this rather than
/// against ``SSHExecutor`` so that provisioning logic can be unit-tested with a
/// recording fake, and so a future vsock-based transport could be dropped in
/// without touching callers.
public protocol GuestExecutor: Sendable {
/// Runs a shell command in the guest and waits for it to exit.
///
/// - Parameters:
/// - command: A `/bin/sh`-compatible command line.
/// - timeout: Wall-clock ceiling; exceeding it throws
/// ``CoreError/timeout(_:)`` and closes the channel.
/// - Returns: Exit status and captured output.
func run(_ command: String, timeout: Duration) async throws -> SSHCommandResult
/// Copies a local file into the guest.
///
/// - Parameters:
/// - localPath: Source path on the host.
/// - remotePath: Destination path in the guest.
func upload(localPath: String, remotePath: String) async throws
/// Writes bytes to a guest file with an explicit mode.
///
/// Used for secrets — notably the registration token, which is written with
/// mode `0600` and deleted immediately after `gitea-runner register` reads
/// it, so it never appears in a process argument list.
///
/// - Parameters:
/// - data: File contents.
/// - remotePath: Destination path in the guest.
/// - mode: Octal mode string, e.g. `"0600"`.
func uploadData(_ data: Data, remotePath: String, mode: String) async throws
}
extension GuestExecutor {
/// ``run(_:timeout:)`` with a two-minute default ceiling.
public func run(_ command: String) async throws -> SSHCommandResult {
try await run(command, timeout: .seconds(120))
}
/// Runs a command and throws unless it exits zero.
///
/// - Returns: The successful result.
@discardableResult
public func runChecked(_ command: String, timeout: Duration = .seconds(120)) async throws -> SSHCommandResult {
let result = try await run(command, timeout: timeout)
try result.throwIfFailed(command: command)
return result
}
}
/// SSH client over swift-nio-ssh using password authentication.
///
/// Password auth (rather than keys) is deliberate: the guest is a throwaway VM
/// on a host-private NAT network whose credentials come from the same config
/// that created it, and injecting a key would mean another provisioning step
/// during the window before SSH is up.
///
/// - Important: Host keys are **not** verified. The peer is a VM this process
/// just booted, on a link no other host shares; there is no trust-on-first-use
/// story that would add security here, and pinning would break on every clone.
public final class SSHExecutor: GuestExecutor, @unchecked Sendable {
/// Guest IP, as learned from ``DHCPLeaseParser``.
public let host: String
/// SSH port; `22` for a stock guest with Remote Login enabled.
public let port: Int
/// Guest account name.
public let username: String
/// Guest account password.
public let password: String
/// Creates an executor. No connection is made until the first command.
///
/// - Parameters:
/// - host: Guest IP address.
/// - port: SSH port. Defaults to `22`.
/// - username: Guest account.
/// - password: Guest password.
public init(host: String, port: Int = 22, username: String, password: String) {
self.host = host
self.port = port
self.username = username
self.password = password
}
public func run(_ command: String, timeout: Duration) async throws -> SSHCommandResult {
do {
return try await execute(command, stdin: nil, timeout: timeout)
} catch let error as SSHTransportError {
throw error.asCoreError
}
}
public func upload(localPath: String, remotePath: String) async throws {
let url = URL(fileURLWithPath: localPath)
guard let data = try? Data(contentsOf: url) else {
throw CoreError.notFound("local file for upload: \(localPath)")
}
try await uploadData(data, remotePath: remotePath, mode: "0644")
}
public func uploadData(_ data: Data, remotePath: String, mode: String) async throws {
// Deliberately not SFTP or SCP: a stock macOS guest runs an sshd whose
// subsystem set we do not control at this point in provisioning, and an
// exec channel with the payload as stdin needs nothing beyond what we
// already use for every other command.
let quotedPath = Self.shellQuote(remotePath)
guard mode.allSatisfy(\.isNumber), !mode.isEmpty else {
throw CoreError.sshFailed("invalid file mode \(mode.debugDescription) for \(remotePath)")
}
let command = """
mkdir -p "$(dirname \(quotedPath))" && cat > \(quotedPath) && chmod \(mode) \(quotedPath)
"""
let result: SSHCommandResult
do {
result = try await execute(command, stdin: data, timeout: .seconds(300))
} catch let error as SSHTransportError {
throw error.asCoreError
}
try result.throwIfFailed(command: "upload to \(remotePath)")
}
/// Releases any pooled connection and event loop resources.
///
/// Connections are not pooled — each command opens and closes its own — and
/// the event loop group is NIO's process-wide singleton, so there is nothing
/// to release. Kept so callers can be written against a lifecycle that a
/// future pooling or vsock transport may need.
public func close() async {}
// MARK: - Transport
/// Opens a connection, runs one exec channel, and tears both down.
///
/// - Parameters:
/// - command: The `/bin/sh` command line to exec.
/// - stdin: Bytes to stream as the command's standard input. Standard
/// input is closed (channel EOF) either way, so a command that would
/// otherwise read from the terminal exits instead of hanging.
/// - timeout: Wall-clock ceiling on the whole exchange.
/// - Throws: ``SSHTransportError`` for connect/auth problems (which
/// ``waitForSSH(host:port:username:password:timeout:pollInterval:)`` needs
/// to tell apart), or ``CoreError/timeout(_:)`` when the ceiling elapses.
func execute(_ command: String, stdin: Data?, timeout: Duration) async throws -> SSHCommandResult {
let group = MultiThreadedEventLoopGroup.singleton
let auth = SSHAuthOutcome()
let username = self.username
let password = self.password
let host = self.host
let port = self.port
let bootstrap = ClientBootstrap(group: group)
.channelOption(ChannelOptions.socketOption(.tcp_nodelay), value: 1)
.channelInitializer { channel in
channel.eventLoop.makeCompletedFuture {
let configuration = SSHClientConfiguration(
userAuthDelegate: PasswordOnlyAuthDelegate(
username: username,
password: password,
outcome: auth
),
serverAuthDelegate: AcceptAnyHostKeyDelegate()
)
try channel.pipeline.syncOperations.addHandler(
NIOSSHHandler(
role: .client(configuration),
allocator: channel.allocator,
inboundChildChannelInitializer: nil
)
)
}
}
let channel: Channel
do {
channel = try await bootstrap.connect(host: host, port: port).get()
} catch {
// No TCP connection at all: sshd is not listening yet (or the guest
// is unreachable). Recoverable — this is what waitForSSH retries on.
throw SSHTransportError.connectFailed(host: host, port: port, underlying: error)
}
let loop = channel.eventLoop
let resultPromise = loop.makePromise(of: SSHCommandResult.self)
let stdinBuffer = stdin.map { ByteBuffer(bytes: $0) }
let description = command
let timeoutTask = loop.scheduleTask(in: .nanoseconds(Self.nanoseconds(timeout))) {
resultPromise.fail(CoreError.timeout("ssh command on \(host): \(description)"))
channel.close(promise: nil)
}
// Completing an already-completed NIO promise is a no-op, so these
// racing completions are safe: whichever fires first wins.
channel.closeFuture.whenComplete { _ in
if auth.wasRejected {
resultPromise.fail(
SSHTransportError.authenticationFailed(host: host, username: username)
)
} else {
resultPromise.fail(
CoreError.sshFailed("ssh connection to \(host):\(port) closed before the command finished")
)
}
}
channel.pipeline.handler(type: NIOSSHHandler.self).flatMap { sshHandler -> EventLoopFuture<Channel> in
let childPromise = loop.makePromise(of: Channel.self)
sshHandler.createChannel(childPromise, channelType: .session) { child, channelType in
guard channelType == .session else {
return child.eventLoop.makeFailedFuture(
CoreError.sshFailed("unexpected SSH channel type \(channelType)")
)
}
return child.eventLoop.makeCompletedFuture {
try child.pipeline.syncOperations.addHandler(
ExecChannelHandler(
command: description,
stdin: stdinBuffer,
promise: resultPromise
)
)
}
// Without this the guest's EOF would close the channel before the
// exit-status request arrives.
.flatMap { child.setOption(ChannelOptions.allowRemoteHalfClosure, value: true) }
}
return childPromise.futureResult
}.whenFailure { error in
resultPromise.fail(error)
channel.close(promise: nil)
}
do {
let result = try await resultPromise.futureResult.get()
timeoutTask.cancel()
channel.close(promise: nil)
return result
} catch {
timeoutTask.cancel()
channel.close(promise: nil)
throw error
}
}
/// Wraps a path (or any argument) so `/bin/sh` sees it literally.
static func shellQuote(_ value: String) -> String {
"'" + value.replacingOccurrences(of: "'", with: "'\\''") + "'"
}
private static func nanoseconds(_ duration: Duration) -> Int64 {
let components = duration.components
let seconds = components.seconds.multipliedReportingOverflow(by: 1_000_000_000)
guard !seconds.overflow else { return .max }
let sum = seconds.partialValue.addingReportingOverflow(
Int64(components.attoseconds / 1_000_000_000)
)
return sum.overflow ? .max : sum.partialValue
}
}
// MARK: - Transport failures
/// Connection-level failures, kept distinct from ``CoreError`` so that
/// ``waitForSSH(host:port:username:password:timeout:pollInterval:)`` can tell
/// "sshd is not up yet" (retry) from "the password is wrong" (give up now).
enum SSHTransportError: Error {
/// No TCP connection could be established.
case connectFailed(host: String, port: Int, underlying: Error)
/// The server rejected our credentials.
case authenticationFailed(host: String, username: String)
var asCoreError: CoreError {
switch self {
case .connectFailed(let host, let port, let underlying):
return .sshFailed("cannot connect to \(host):\(port): \(underlying)")
case .authenticationFailed(let host, let username):
return .sshFailed("authentication failed for \(username)@\(host)")
}
}
}
/// Shared, thread-safe record of whether the server rejected our password.
///
/// The auth delegate runs on the event loop; the value is read from the async
/// caller, hence the lock.
final class SSHAuthOutcome: @unchecked Sendable {
private let lock = NSLock()
private var rejected = false
var wasRejected: Bool {
lock.lock()
defer { lock.unlock() }
return rejected
}
func markRejected() {
lock.lock()
defer { lock.unlock() }
rejected = true
}
}
// MARK: - Delegates
/// Accepts every host key.
///
/// The peer is a VM this process booted seconds ago, on a host-private NAT link
/// no other machine shares, from an image that is destroyed after one job. Its
/// host key is freshly generated per clone, so there is nothing to pin:
/// trust-on-first-use would accept whatever the first connection presented —
/// exactly what this does — while a pinned key would reject every legitimate
/// guest. See docs/DESIGN.md §6.
final class AcceptAnyHostKeyDelegate: NIOSSHClientServerAuthenticationDelegate {
func validateHostKey(hostKey: NIOSSHPublicKey, validationCompletePromise: EventLoopPromise<Void>) {
validationCompletePromise.succeed(())
}
}
/// Offers the configured password once, then reports rejection.
///
/// NIOSSH asks again after a failed attempt; a second ask means the server
/// refused the first, which is a provisioning bug rather than a transient
/// condition, so it is recorded for the caller to fail fast on.
final class PasswordOnlyAuthDelegate: NIOSSHClientUserAuthenticationDelegate {
private let username: String
private let password: String
private let outcome: SSHAuthOutcome
private var offered = false
init(username: String, password: String, outcome: SSHAuthOutcome) {
self.username = username
self.password = password
self.outcome = outcome
}
func nextAuthenticationType(
availableMethods: NIOSSHAvailableUserAuthenticationMethods,
nextChallengePromise: EventLoopPromise<NIOSSHUserAuthenticationOffer?>
) {
guard !offered, availableMethods.contains(.password) else {
// Either the server refused our password, or it never offered
// password auth at all. Both mean this guest will not let us in.
outcome.markRejected()
nextChallengePromise.succeed(nil)
return
}
offered = true
nextChallengePromise.succeed(
NIOSSHUserAuthenticationOffer(
username: username,
serviceName: "",
offer: .password(.init(password: password))
)
)
}
}
// MARK: - Exec channel
/// Drives one exec channel: sends the request, streams stdin, splits stdout from
/// stderr, and captures the exit status.
final class ExecChannelHandler: ChannelInboundHandler {
typealias InboundIn = SSHChannelData
typealias OutboundOut = SSHChannelData
/// SSH channel data is framed into packets; keep writes comfortably under
/// the 128 KiB default maximum packet size.
private static let chunkSize = 32 * 1024
private let command: String
private var stdin: ByteBuffer?
private var promise: EventLoopPromise<SSHCommandResult>?
private var stdout = ByteBufferAllocator().buffer(capacity: 0)
private var stderr = ByteBufferAllocator().buffer(capacity: 0)
private var exitCode: Int32?
init(command: String, stdin: ByteBuffer?, promise: EventLoopPromise<SSHCommandResult>) {
self.command = command
self.stdin = stdin
self.promise = promise
}
func channelActive(context: ChannelHandlerContext) {
let request = SSHChannelRequestEvent.ExecRequest(command: command, wantReply: true)
let sent = context.eventLoop.makePromise(of: Void.self)
// Capture only Sendable values: the handler and its context must not
// escape onto another thread.
let resultPromise = promise
let channel = context.channel
sent.futureResult.whenFailure { error in
resultPromise?.fail(error)
channel.close(promise: nil)
}
context.triggerUserOutboundEvent(request, promise: sent)
context.fireChannelActive()
}
func userInboundEventTriggered(context: ChannelHandlerContext, event: Any) {
switch event {
case is ChannelSuccessEvent:
sendStandardInput(context: context)
case is ChannelFailureEvent:
fail(context: context, error: CoreError.sshFailed("guest refused to exec: \(command)"))
case let status as SSHChannelRequestEvent.ExitStatus:
exitCode = Int32(truncatingIfNeeded: status.exitStatus)
case let signal as SSHChannelRequestEvent.ExitSignal:
// A signalled process has no exit status; report it the way a shell
// would, and keep the signal name in stderr so it is not lost.
exitCode = 128
var note = ByteBuffer(string: "\nterminated by SIG\(signal.signalName): \(signal.errorMessage)\n")
stderr.writeBuffer(&note)
default:
context.fireUserInboundEventTriggered(event)
}
}
func channelRead(context: ChannelHandlerContext, data: NIOAny) {
let channelData = unwrapInboundIn(data)
guard case .byteBuffer(var bytes) = channelData.data else { return }
switch channelData.type {
case .channel: stdout.writeBuffer(&bytes)
case .stdErr: stderr.writeBuffer(&bytes)
default: break // An extended data type we did not ask for.
}
}
func channelInactive(context: ChannelHandlerContext) {
complete()
context.fireChannelInactive()
}
func handlerRemoved(context: ChannelHandlerContext) {
complete()
}
func errorCaught(context: ChannelHandlerContext, error: Error) {
fail(context: context, error: error)
}
private func sendStandardInput(context: ChannelHandlerContext) {
if var payload = stdin {
stdin = nil
while payload.readableBytes > 0 {
let slice = payload.readSlice(length: min(Self.chunkSize, payload.readableBytes))!
context.write(
wrapOutboundOut(SSHChannelData(type: .channel, data: .byteBuffer(slice))),
promise: nil
)
}
context.flush()
}
// EOF either way: `cat > file` needs it to finish, and a command that
// would otherwise block reading stdin gets an immediate end of input.
context.close(mode: .output, promise: nil)
}
private func complete() {
guard let promise else { return }
self.promise = nil
if let exitCode {
promise.succeed(
SSHCommandResult(
exitCode: exitCode,
stdout: String(buffer: stdout),
stderr: String(buffer: stderr)
)
)
} else {
promise.fail(
CoreError.sshFailed("guest closed the channel without an exit status: \(command)")
)
}
}
private func fail(context: ChannelHandlerContext, error: Error) {
if let promise {
self.promise = nil
promise.fail(error)
}
context.close(promise: nil)
}
}
/// Blocks until a guest accepts an authenticated SSH session, or the deadline
/// passes.
///
/// Called after a DHCP lease appears but before any provisioning: a fresh guest
/// answers on port 22 only once `launchd` has started `sshd`, which lags the
/// lease by tens of seconds.
///
/// - Parameters:
/// - host: Guest IP.
/// - port: SSH port. Defaults to `22`.
/// - username: Guest account.
/// - password: Guest password.
/// - timeout: Overall ceiling.
/// - pollInterval: Delay between attempts. Defaults to 2 s.
/// - Throws: ``CoreError/timeout(_:)`` if the guest never answers.
public func waitForSSH(
host: String,
port: Int = 22,
username: String,
password: String,
timeout: Duration,
pollInterval: Duration = .seconds(2)
) async throws {
let executor = SSHExecutor(host: host, port: port, username: username, password: password)
let started = ContinuousClock.now
var lastError: Error?
while true {
do {
// A real authenticated session running a trivial command, not a bare
// TCP probe: sshd binds the port before it is ready to authenticate,
// so a connect that succeeds proves very little.
_ = try await executor.execute("true", stdin: nil, timeout: .seconds(20))
return
} catch let error as SSHTransportError {
if case .authenticationFailed = error {
// Wrong credentials will not become right by waiting: the guest
// was provisioned with a different account or password, which is
// a build failure, not a boot delay.
throw error.asCoreError
}
lastError = error
} catch {
// Timeouts and mid-handshake closures are what a guest that is still
// starting `sshd` looks like. Keep waiting.
lastError = error
}
guard ContinuousClock.now - started < timeout else { break }
try await Task.sleep(for: pollInterval)
guard ContinuousClock.now - started < timeout else { break }
}
let detail = lastError.map { "; last error: \($0)" } ?? ""
throw CoreError.timeout("ssh on \(host):\(port)\(detail)")
}