mirror of
https://github.com/anomalyco/opencode.git
synced 2026-10-06 15:36:34 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dda732362b |
No files matched your search
@@ -1,6 +1,6 @@
|
||||
import { Effect, Option, Schema } from "effect"
|
||||
import { HttpClient } from "effect/unstable/http"
|
||||
import { parseTarget, quote, runSsh, sshArgs, SshFailure } from "./command"
|
||||
import { SshFailure } from "./command"
|
||||
import { RemoteCli } from "./remote-cli"
|
||||
|
||||
// Use commands supported by released V2 CLIs. The registration is the service's
|
||||
@@ -74,22 +74,14 @@ function connectionAddress(address: string, password: string) {
|
||||
}
|
||||
|
||||
export const bootstrap = Effect.fn("Ssh.bootstrap")(function* (input: {
|
||||
target: ReturnType<typeof parseTarget>
|
||||
run: (script: string, stdin?: Uint8Array) => Effect.Effect<string, SshFailure>
|
||||
version: string
|
||||
development?: boolean
|
||||
env: NodeJS.ProcessEnv
|
||||
replace?: boolean
|
||||
stage: (stage: "checking" | "downloading" | "uploading" | "starting") => Effect.Effect<void>
|
||||
}) {
|
||||
const run = (script: string) =>
|
||||
runSsh({
|
||||
args: [...sshArgs(input.target), input.target.host, "sh -l -s"],
|
||||
env: input.env,
|
||||
stdin: script,
|
||||
})
|
||||
|
||||
yield* input.stage("checking")
|
||||
const registered = parseRegistration(yield* run(discoverScript))
|
||||
const registered = parseRegistration(yield* input.run(discoverScript))
|
||||
|
||||
if (registered && (input.development || registered.version === input.version)) {
|
||||
yield* input.stage("starting")
|
||||
@@ -102,7 +94,7 @@ export const bootstrap = Effect.fn("Ssh.bootstrap")(function* (input: {
|
||||
|
||||
if (registered && !input.replace) return yield* Effect.fail(new SshFailure("version", registered.version))
|
||||
const destination = yield* Effect.try({ try: () => binaryPath(input.version), catch: SshFailure.from })
|
||||
const existing = yield* run(RemoteCli.versionScript(`"${destination}"`))
|
||||
const existing = yield* input.run(RemoteCli.versionScript(`"${destination}"`))
|
||||
const staged = RemoteCli.parseVersion(existing) === input.version
|
||||
|
||||
// Source worktree versions are unpublished. Use the installer's beta channel
|
||||
@@ -113,7 +105,7 @@ export const bootstrap = Effect.fn("Ssh.bootstrap")(function* (input: {
|
||||
const setup = { version, directory: `.opencode/desktop-ssh/${version}` }
|
||||
|
||||
if (!staged) {
|
||||
const output = yield* run(RemoteCli.probeScript).pipe(Effect.mapError(() => new SshFailure("platform")))
|
||||
const output = yield* input.run(RemoteCli.probeScript).pipe(Effect.mapError(() => new SshFailure("platform")))
|
||||
|
||||
const target = output
|
||||
.split(/\r?\n/)
|
||||
@@ -122,7 +114,7 @@ export const bootstrap = Effect.fn("Ssh.bootstrap")(function* (input: {
|
||||
|
||||
const url = yield* Effect.try({ try: () => RemoteCli.archiveUrl(target ?? "", version), catch: SshFailure.from })
|
||||
yield* input.stage("downloading")
|
||||
yield* run(RemoteCli.installScript({ ...setup, source: { type: "download", url } })).pipe(
|
||||
yield* input.run(RemoteCli.installScript({ ...setup, source: { type: "download", url } })).pipe(
|
||||
Effect.catch(
|
||||
Effect.fnUntraced(function* (error) {
|
||||
yield* input.stage("uploading")
|
||||
@@ -138,23 +130,16 @@ export const bootstrap = Effect.fn("Ssh.bootstrap")(function* (input: {
|
||||
)
|
||||
const archive = new Uint8Array(yield* response.arrayBuffer.pipe(Effect.mapError(SshFailure.from)))
|
||||
|
||||
// The upload uses stdin; the script itself must be the remote command.
|
||||
return yield* runSsh({
|
||||
args: [
|
||||
...sshArgs(input.target),
|
||||
input.target.host,
|
||||
`sh -c ${quote(RemoteCli.installScript({ ...setup, source: { type: "archive" } }))}`,
|
||||
],
|
||||
env: input.env,
|
||||
stdin: archive,
|
||||
}).pipe(Effect.mapError(() => new SshFailure("install", error.message)))
|
||||
return yield* input
|
||||
.run(RemoteCli.installScript({ ...setup, source: { type: "archive" } }), archive)
|
||||
.pipe(Effect.mapError(() => new SshFailure("install", error.message)))
|
||||
}),
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
yield* input.stage("starting")
|
||||
const registration = parseRegistration(yield* run(startScript(version, input.replace)))
|
||||
const registration = parseRegistration(yield* input.run(startScript(version, input.replace)))
|
||||
|
||||
if (!registration) return yield* Effect.fail(new SshFailure("service"))
|
||||
|
||||
|
||||
@@ -1,10 +1,12 @@
|
||||
import { describe, expect, test } from "bun:test"
|
||||
import { PlatformError } from "effect"
|
||||
import { parseTarget, quote, sshArgs, tunnelArgs, commandFailureDetail, SshFailure } from "./command"
|
||||
import { parseTarget, quote, sshArgs, commandFailureDetail, SshFailure } from "./command"
|
||||
|
||||
describe("SSH connection commands", () => {
|
||||
test("classifies a missing SSH executable from the platform error", () => {
|
||||
// SAFETY: `systemError` takes the reason's tag as data; it has no constructor per tag.
|
||||
const error = PlatformError.systemError({
|
||||
// oxlint-disable-next-line anti-slop-effect/no-manual-tagged-construction -- see SAFETY above
|
||||
_tag: "NotFound",
|
||||
module: "ChildProcessSpawner",
|
||||
method: "spawn",
|
||||
@@ -13,20 +15,6 @@ describe("SSH connection commands", () => {
|
||||
|
||||
expect(SshFailure.from(error).code).toBe("ssh-missing")
|
||||
})
|
||||
test("forwarding overrides bootstrap persistence before reusing the control socket", () => {
|
||||
const args = tunnelArgs(
|
||||
{
|
||||
host: "devbox",
|
||||
args: ["-o", "ControlMaster=auto", "-o", "ControlPersist=60", "-o", "ControlPath=/test/socket"],
|
||||
},
|
||||
1234,
|
||||
{ host: "127.0.0.1", port: 5678 },
|
||||
)
|
||||
|
||||
expect(args.slice(0, 4)).toEqual(["-o", "ControlMaster=no", "-o", "ControlPersist=no"])
|
||||
expect(args).toContain("ControlPath=/test/socket")
|
||||
expect(args.slice(-4)).toEqual(["-L", "127.0.0.1:1234:127.0.0.1:5678", "devbox", "sh -c 'exec cat >/dev/null'"])
|
||||
})
|
||||
test("retains CLI stdout failures without exposing private connection details", () => {
|
||||
expect(
|
||||
commandFailureDetail(1, {
|
||||
|
||||
@@ -154,28 +154,6 @@ export function sshArgs(target: ReturnType<typeof parseTarget>) {
|
||||
]
|
||||
}
|
||||
|
||||
export function tunnelArgs(
|
||||
target: ReturnType<typeof parseTarget>,
|
||||
localPort: number,
|
||||
remote: { host: string; port: number },
|
||||
) {
|
||||
// A multiplexed `ssh -N` may exit after handing forwarding to its master.
|
||||
// Keep a session open on stdin instead; the scoped process owns that pipe.
|
||||
return [
|
||||
"-o",
|
||||
"ControlMaster=no",
|
||||
"-o",
|
||||
"ControlPersist=no",
|
||||
...sshArgs(target),
|
||||
"-o",
|
||||
"ExitOnForwardFailure=yes",
|
||||
"-L",
|
||||
`127.0.0.1:${localPort}:${remote.host}:${remote.port}`,
|
||||
target.host,
|
||||
"sh -c 'exec cat >/dev/null'",
|
||||
]
|
||||
}
|
||||
|
||||
export const runSsh = Effect.fn("Ssh.run")(function* (input: {
|
||||
args: string[]
|
||||
env?: NodeJS.ProcessEnv
|
||||
|
||||
@@ -10,17 +10,17 @@ import {
|
||||
Path,
|
||||
Predicate,
|
||||
PubSub,
|
||||
Ref,
|
||||
Schedule,
|
||||
Scope,
|
||||
Stream,
|
||||
} from "effect"
|
||||
import { HttpClient } from "effect/unstable/http"
|
||||
import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"
|
||||
import { ChildProcessSpawner } from "effect/unstable/process"
|
||||
import type { SshConfig, SshHttp, SshItem, SshStart, SshState } from "./contract"
|
||||
import { createAskpass } from "./askpass"
|
||||
import { bootstrap } from "./bootstrap"
|
||||
import { parseTarget, quote, runSsh, sshArgs, sshExecutable, tunnelArgs, SshFailure } from "./command"
|
||||
import { parseTarget, quote, runSsh, sshArgs, SshFailure } from "./command"
|
||||
import { forward, openSession } from "./session"
|
||||
|
||||
type Connection = {
|
||||
owner?: number
|
||||
@@ -92,7 +92,6 @@ export const createSshController = Effect.fn("Ssh.controller")(function* (input:
|
||||
const connect = Effect.fn("Ssh.connect")(function* (config: SshConfig, connection: Connection, replace = false) {
|
||||
const target = yield* Effect.try({ try: () => parseTarget(config.target), catch: SshFailure.from })
|
||||
const directory = yield* fs.makeTempDirectoryScoped({ prefix: "oc-ssh-" })
|
||||
const control = path.join(directory, "s")
|
||||
|
||||
const helper =
|
||||
input.command && input.command.length > 1 && process.platform !== "win32"
|
||||
@@ -104,15 +103,6 @@ export const createSshController = Effect.fn("Ssh.controller")(function* (input:
|
||||
mode: 0o700,
|
||||
})
|
||||
|
||||
if (process.platform !== "win32") {
|
||||
target.args.unshift("-o", "ControlMaster=auto", "-o", "ControlPersist=60", "-o", `ControlPath=${control}`)
|
||||
// Close only our local SSH master. The remote OpenCode service owns its
|
||||
// own lifetime and must survive disconnect, failure, and app shutdown.
|
||||
yield* Effect.addFinalizer(() =>
|
||||
run({ args: ["-o", `ControlPath=${control}`, "-O", "exit", target.host], timeout: 2000 }).pipe(Effect.ignore),
|
||||
)
|
||||
}
|
||||
|
||||
const authentication = yield* Deferred.make<void>()
|
||||
|
||||
const askpass = yield* createAskpass({
|
||||
@@ -154,56 +144,28 @@ export const createSshController = Effect.fn("Ssh.controller")(function* (input:
|
||||
destination: `${fields.get("user") ?? ""}@${fields.get("hostname")}:${fields.get("port") ?? "22"}`,
|
||||
})
|
||||
|
||||
const socks = yield* freePort
|
||||
|
||||
const session = yield* openSession({ target, env: askpass.env, socks }).pipe(
|
||||
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
|
||||
)
|
||||
|
||||
const remote = yield* bootstrap({
|
||||
target,
|
||||
run: session.run,
|
||||
version: input.version,
|
||||
development: input.development,
|
||||
env: askpass.env,
|
||||
replace,
|
||||
stage: (stage) => update(config.id, { stage, prompt: undefined }),
|
||||
}).pipe(
|
||||
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
|
||||
Effect.provideService(HttpClient.HttpClient, httpClient),
|
||||
)
|
||||
}).pipe(Effect.provideService(HttpClient.HttpClient, httpClient))
|
||||
|
||||
const port = yield* freePort
|
||||
const http = { url: `http://127.0.0.1:${port}`, password: remote.password }
|
||||
|
||||
const tunnel = yield* spawner.spawn(
|
||||
ChildProcess.make(sshExecutable(), tunnelArgs(target, port, remote), {
|
||||
env: askpass.env,
|
||||
extendEnv: true,
|
||||
windowsHide: true,
|
||||
stdin: "pipe",
|
||||
stdout: "ignore",
|
||||
killSignal: "SIGTERM",
|
||||
forceKillAfter: "2 seconds",
|
||||
}),
|
||||
)
|
||||
|
||||
const detail = yield* Ref.make("")
|
||||
|
||||
const stderr = yield* tunnel.stderr.pipe(
|
||||
Stream.decodeText(),
|
||||
Stream.runForEach((text) => Ref.update(detail, (tail) => (tail + text).slice(-8192))),
|
||||
Effect.forkScoped,
|
||||
)
|
||||
|
||||
const closed = Effect.gen(function* () {
|
||||
const exitCode = yield* tunnel.exitCode
|
||||
yield* Fiber.join(stderr)
|
||||
|
||||
return yield* Effect.fail(
|
||||
new SshFailure("connection", (yield* Ref.get(detail)) || JSON.stringify({ exitCode })),
|
||||
)
|
||||
})
|
||||
const http = { url: `http://127.0.0.1:${yield* forward(socks, remote)}`, password: remote.password }
|
||||
|
||||
yield* waitReady(http, () => items.get(config.id)?.stage === "authentication").pipe(
|
||||
Effect.provideService(HttpClient.HttpClient, httpClient),
|
||||
Effect.catch(() =>
|
||||
Ref.get(detail).pipe(Effect.flatMap((detail) => Effect.fail(new SshFailure("service", detail)))),
|
||||
session.detail.pipe(Effect.flatMap((detail) => Effect.fail(new SshFailure("service", detail)))),
|
||||
),
|
||||
Effect.raceFirst(closed),
|
||||
Effect.raceFirst(session.closed),
|
||||
)
|
||||
const saved = new Map(configs).set(config.id, config)
|
||||
yield* input.save([...saved.values()])
|
||||
@@ -211,7 +173,7 @@ export const createSshController = Effect.fn("Ssh.controller")(function* (input:
|
||||
failures.delete(config.id)
|
||||
yield* update(config.id, { http, stage: "ready", saved: true, detail: "", prompt: undefined, error: undefined })
|
||||
yield* Deferred.succeed(connection.ready, http)
|
||||
yield* closed
|
||||
yield* session.closed
|
||||
}).pipe(
|
||||
Effect.raceFirst(askpass.closed),
|
||||
Effect.raceFirst(Deferred.await(authentication).pipe(Effect.andThen(Effect.interrupt))),
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
import { expect } from "bun:test"
|
||||
import { NodeServices } from "@effect/platform-node"
|
||||
import { Effect, Exit, Predicate } from "effect"
|
||||
import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"
|
||||
import { createServer, connect, type AddressInfo, type Server } from "node:net"
|
||||
import { testEffect } from "../../../core/test/lib/effect"
|
||||
import { SshFailure } from "./command"
|
||||
import { forward, openSession } from "./session"
|
||||
|
||||
const it = testEffect(NodeServices.layer)
|
||||
|
||||
const listen = (server: Server) =>
|
||||
Effect.promise(
|
||||
() =>
|
||||
new Promise<number>((resolve) =>
|
||||
// SAFETY: a TCP server listening on a host and port reports an AddressInfo.
|
||||
server.listen(0, "127.0.0.1", () => resolve((server.address() as AddressInfo).port)),
|
||||
),
|
||||
)
|
||||
|
||||
// Stands in for `ssh host sh -l -s`: the session's stdin reaches a local shell.
|
||||
const open = Effect.fn("test.session.open")(function* (remote: string) {
|
||||
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner
|
||||
|
||||
return yield* openSession({ target: { host: "devbox", args: [] }, env: {}, socks: 1 }).pipe(
|
||||
Effect.provideService(
|
||||
ChildProcessSpawner.ChildProcessSpawner,
|
||||
ChildProcessSpawner.make((command) =>
|
||||
spawner.spawn(
|
||||
Predicate.isTagged(command, "StandardCommand")
|
||||
? ChildProcess.make("sh", ["-c", remote], command.options)
|
||||
: command,
|
||||
),
|
||||
),
|
||||
),
|
||||
)
|
||||
})
|
||||
|
||||
it.live(
|
||||
"runs every bootstrap step over one connection after login output",
|
||||
Effect.gen(function* () {
|
||||
const session = yield* open("echo login banner; printf 'partial'; exec sh -s")
|
||||
|
||||
expect(yield* session.run("set -eu\necho out\necho err >&2\nprintf 'no newline'")).toBe("out\nno newline")
|
||||
|
||||
const failed = yield* session.run("set -eu\necho before >&2\nexit 3\necho never").pipe(Effect.flip)
|
||||
expect(failed).toBeInstanceOf(SshFailure)
|
||||
expect(failed.message).toBe("before")
|
||||
|
||||
// Larger than one read, so the driver must assemble it from several pipe reads.
|
||||
const archive = new TextEncoder().encode("0123456789abcdef\n".repeat(200_000))
|
||||
expect(yield* session.run("cat", archive)).toBe(new TextDecoder().decode(archive))
|
||||
expect(yield* session.run('printf "%s" "$(cat)"', new TextEncoder().encode("stdin is per step"))).toBe(
|
||||
"stdin is per step",
|
||||
)
|
||||
expect(yield* session.run("cat")).toBe("")
|
||||
}),
|
||||
)
|
||||
|
||||
it.live(
|
||||
"fails every step with ssh diagnostics when the connection closes",
|
||||
Effect.gen(function* () {
|
||||
const session = yield* open("echo 'user@devbox: Permission denied (password).' >&2; exit 255")
|
||||
const failed = yield* session.run("echo never").pipe(Effect.flip)
|
||||
expect(failed.message).toBe("user@devbox: Permission denied (password).")
|
||||
expect((yield* session.closed.pipe(Effect.flip)).message).toBe(failed.message)
|
||||
}),
|
||||
)
|
||||
|
||||
it.live(
|
||||
"forwards local connections to the remote service through SOCKS",
|
||||
Effect.gen(function* () {
|
||||
const requests: string[] = []
|
||||
|
||||
const service = createServer((socket) =>
|
||||
socket.on("data", (data) => {
|
||||
requests.push(data.toString())
|
||||
socket.end("pong")
|
||||
}),
|
||||
)
|
||||
|
||||
// A minimal SOCKS5 server, as `ssh -D` provides: no authentication, CONNECT only.
|
||||
const socks = createServer((client) => {
|
||||
client.once("data", () => {
|
||||
client.write(Buffer.from([5, 0]))
|
||||
client.once("data", (request) => {
|
||||
const host = request.subarray(5, 5 + (request[4] ?? 0)).toString()
|
||||
const port = request.readUInt16BE(5 + (request[4] ?? 0))
|
||||
|
||||
const upstream = connect(port, host, () => {
|
||||
// Coalesce the reply with service bytes to cover data after the handshake.
|
||||
client.write(Buffer.from([5, 0, 0, 1, 0, 0, 0, 0, 0, 0]))
|
||||
client.pipe(upstream).pipe(client)
|
||||
})
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
yield* Effect.addFinalizer(() => Effect.sync(() => (service.close(), socks.close())))
|
||||
const port = yield* forward(yield* listen(socks), { host: "127.0.0.1", port: yield* listen(service) })
|
||||
|
||||
const reply = yield* Effect.promise(
|
||||
() =>
|
||||
new Promise<string>((resolve) => {
|
||||
const socket = connect(port, "127.0.0.1", () => socket.write("ping"))
|
||||
const chunks: Buffer[] = []
|
||||
socket.on("data", (data) => chunks.push(data)).on("close", () => resolve(Buffer.concat(chunks).toString()))
|
||||
}),
|
||||
)
|
||||
|
||||
expect(reply).toBe("pong")
|
||||
expect(requests).toEqual(["ping"])
|
||||
}),
|
||||
)
|
||||
|
||||
it.live(
|
||||
"closes the local forward with the attempt",
|
||||
Effect.gen(function* () {
|
||||
const port = yield* Effect.scoped(forward(1, { host: "127.0.0.1", port: 1 }))
|
||||
|
||||
const exit = yield* Effect.exit(
|
||||
Effect.tryPromise(
|
||||
() =>
|
||||
new Promise<void>((resolve, reject) => {
|
||||
const socket = connect(port, "127.0.0.1", () => (socket.destroy(), resolve()))
|
||||
socket.on("error", reject)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
expect(Exit.isFailure(exit)).toBe(true)
|
||||
}),
|
||||
)
|
||||
@@ -0,0 +1,272 @@
|
||||
import { Cause, Deferred, Effect, Fiber, Queue, Semaphore, Stream } from "effect"
|
||||
import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"
|
||||
import { randomUUID } from "node:crypto"
|
||||
import { createServer, connect, type AddressInfo, type Socket } from "node:net"
|
||||
import { commandFailureDetail, parseTarget, sshArgs, sshExecutable, SshFailure } from "./command"
|
||||
|
||||
// Runs every bootstrap step and the HTTP forward over one SSH connection, so the
|
||||
// user authenticates once per attempt on every platform, including Windows
|
||||
// OpenSSH, which has no connection multiplexing.
|
||||
//
|
||||
// `sh -l -s` parses this group before running it, and the client sends no frame
|
||||
// until the ready marker, so the shell has nothing buffered past the group. Each
|
||||
// frame is `<id> <script bytes> <stdin bytes>\n` followed by both payloads; `dd`
|
||||
// with `count=1` never reads past the frame. The reply is a result header and
|
||||
// the step's stdout and stderr.
|
||||
export function driver(nonce: string) {
|
||||
return `{
|
||||
dir=$(mktemp -d) || exit 1
|
||||
trap 'rm -rf "$dir"' EXIT
|
||||
trap 'exit 1' HUP INT TERM PIPE
|
||||
take() {
|
||||
: > "$2"
|
||||
size=0
|
||||
while [ "$size" -lt "$1" ]; do
|
||||
chunk=$(($1 - size))
|
||||
if [ "$chunk" -gt 65536 ]; then chunk=65536; fi
|
||||
dd bs="$chunk" count=1 2>/dev/null >> "$2" || return 1
|
||||
next=$(($(wc -c < "$2")))
|
||||
if [ "$next" -eq "$size" ]; then return 1; fi
|
||||
size=$next
|
||||
done
|
||||
}
|
||||
printf '\\nOPENCODE_SSH_READY_${nonce}\\n'
|
||||
while read -r id script input; do
|
||||
take "$script" "$dir/script" && take "$input" "$dir/input" || exit 1
|
||||
( . "$dir/script" ) < "$dir/input" > "$dir/stdout" 2> "$dir/stderr"
|
||||
printf 'OPENCODE_SSH_RESULT %s %s %s %s\\n' "$id" "$?" "$(($(wc -c < "$dir/stdout")))" "$(($(wc -c < "$dir/stderr")))"
|
||||
cat "$dir/stdout" "$dir/stderr"
|
||||
rm -f "$dir/script" "$dir/input" "$dir/stdout" "$dir/stderr"
|
||||
done
|
||||
exit 0
|
||||
}
|
||||
`
|
||||
}
|
||||
|
||||
export function sessionArgs(target: ReturnType<typeof parseTarget>, socks: number) {
|
||||
// Our own process must hold the forward: a ControlMaster from the user's
|
||||
// config could keep it alive in the background after this attempt ends.
|
||||
return [
|
||||
"-o",
|
||||
"ControlMaster=no",
|
||||
"-o",
|
||||
"ControlPath=none",
|
||||
...sshArgs(target),
|
||||
"-o",
|
||||
"ExitOnForwardFailure=yes",
|
||||
"-D",
|
||||
`127.0.0.1:${socks}`,
|
||||
target.host,
|
||||
"sh -l -s",
|
||||
]
|
||||
}
|
||||
|
||||
type Result = { id: number; code: number; stdout: string; stderr: string }
|
||||
|
||||
export const openSession = Effect.fn("Ssh.session")(function* (input: {
|
||||
target: ReturnType<typeof parseTarget>
|
||||
env: NodeJS.ProcessEnv
|
||||
socks: number
|
||||
}) {
|
||||
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner
|
||||
const nonce = randomUUID().replaceAll("-", "")
|
||||
const marker = Buffer.from(`\nOPENCODE_SSH_READY_${nonce}\n`)
|
||||
const stdin = yield* Queue.unbounded<Uint8Array, Cause.Done>()
|
||||
const results = yield* Queue.unbounded<Result, SshFailure>()
|
||||
const ready = yield* Deferred.make<void, SshFailure>()
|
||||
const exited = yield* Deferred.make<never, SshFailure>()
|
||||
const steps = yield* Semaphore.make(1)
|
||||
const state = { buffer: Buffer.alloc(0), ready: false, id: 0, stderr: "" }
|
||||
yield* Effect.addFinalizer(() => Queue.end(stdin))
|
||||
yield* Queue.offer(stdin, Buffer.from(driver(nonce)))
|
||||
|
||||
const child = yield* spawner.spawn(
|
||||
ChildProcess.make(sshExecutable(), sessionArgs(input.target, input.socks), {
|
||||
env: input.env,
|
||||
extendEnv: true,
|
||||
windowsHide: true,
|
||||
killSignal: "SIGTERM",
|
||||
forceKillAfter: "2 seconds",
|
||||
stdin: { stream: Stream.fromQueue(stdin), endOnDone: true },
|
||||
}),
|
||||
)
|
||||
|
||||
const deliver = Effect.fnUntraced(function* (chunk: Uint8Array) {
|
||||
state.buffer = Buffer.concat([state.buffer, chunk])
|
||||
|
||||
if (!state.ready) {
|
||||
const at = state.buffer.indexOf(marker)
|
||||
|
||||
// Login shells may print before the driver starts. Keep only a possible
|
||||
// partial marker.
|
||||
if (at === -1) {
|
||||
state.buffer = state.buffer.subarray(-marker.length)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
state.buffer = state.buffer.subarray(at + marker.length)
|
||||
state.ready = true
|
||||
yield* Deferred.succeed(ready, undefined)
|
||||
}
|
||||
|
||||
while (true) {
|
||||
const end = state.buffer.indexOf(10)
|
||||
|
||||
if (end === -1) return
|
||||
const header = /^OPENCODE_SSH_RESULT (\d+) (\d+) (\d+) (\d+)$/.exec(state.buffer.subarray(0, end).toString())
|
||||
|
||||
if (!header) return yield* Effect.fail(new SshFailure("connection", "Unexpected remote shell output"))
|
||||
const stdout = end + 1 + Number(header[3])
|
||||
const stderr = stdout + Number(header[4])
|
||||
|
||||
if (state.buffer.length < stderr) return
|
||||
yield* Queue.offer(results, {
|
||||
id: Number(header[1]),
|
||||
code: Number(header[2]),
|
||||
stdout: state.buffer.subarray(end + 1, stdout).toString(),
|
||||
stderr: state.buffer.subarray(stdout, stderr).toString(),
|
||||
})
|
||||
state.buffer = state.buffer.subarray(stderr)
|
||||
}
|
||||
})
|
||||
|
||||
const stderr = yield* child.stderr.pipe(
|
||||
Stream.decodeText(),
|
||||
Stream.runForEach((text) => Effect.sync(() => (state.stderr = (state.stderr + text).slice(-16_384)))),
|
||||
Effect.forkScoped,
|
||||
)
|
||||
|
||||
// Every step fails once the connection closes or its output stops following
|
||||
// the protocol; an SSH failure carries ssh's own diagnostics.
|
||||
yield* child.stdout.pipe(
|
||||
Stream.runForEach(deliver),
|
||||
Effect.matchEffect({
|
||||
onFailure: (error) => Effect.succeed(SshFailure.from(error)),
|
||||
onSuccess: () =>
|
||||
Effect.gen(function* () {
|
||||
const code = yield* child.exitCode.pipe(Effect.orElseSucceed(() => null))
|
||||
yield* Fiber.join(stderr).pipe(Effect.ignore)
|
||||
|
||||
return new SshFailure("connection", commandFailureDetail(code, { stdout: "", stderr: state.stderr }))
|
||||
}),
|
||||
}),
|
||||
Effect.flatMap((failure) =>
|
||||
Effect.all([Deferred.fail(ready, failure), Queue.fail(results, failure), Deferred.fail(exited, failure)]),
|
||||
),
|
||||
Effect.forkScoped,
|
||||
)
|
||||
|
||||
return {
|
||||
/** Fails with the SSH diagnostic once the connection closes. */
|
||||
closed: Deferred.await(exited),
|
||||
/** The connection's latest SSH diagnostics, for failures it did not cause. */
|
||||
detail: Effect.sync(() => state.stderr),
|
||||
/** Runs one remote shell script and returns its stdout, like `ssh host sh -l -s`. */
|
||||
run: (script: string, data?: Uint8Array) =>
|
||||
steps
|
||||
.withPermit(
|
||||
Effect.gen(function* () {
|
||||
yield* Deferred.await(ready)
|
||||
const id = ++state.id
|
||||
const body = Buffer.from(script)
|
||||
yield* Queue.offer(
|
||||
stdin,
|
||||
Buffer.concat([
|
||||
Buffer.from(`${id} ${body.length} ${data?.length ?? 0}\n`),
|
||||
body,
|
||||
data ?? Buffer.alloc(0),
|
||||
]),
|
||||
)
|
||||
const result = yield* Queue.take(results)
|
||||
|
||||
if (result.id !== id)
|
||||
return yield* Effect.fail(new SshFailure("connection", "Unexpected remote shell output"))
|
||||
|
||||
if (result.code !== 0)
|
||||
return yield* Effect.fail(new SshFailure("connection", commandFailureDetail(result.code, result)))
|
||||
|
||||
return result.stdout
|
||||
}),
|
||||
)
|
||||
.pipe(
|
||||
Effect.timeoutOrElse({
|
||||
duration: "10 minutes",
|
||||
orElse: () => Effect.fail(new SshFailure("connection", "Remote command timed out")),
|
||||
}),
|
||||
),
|
||||
}
|
||||
})
|
||||
|
||||
// Serves the remote OpenCode service on a local port through the session's
|
||||
// SOCKS forward, so the forward needs no port before bootstrap finds the service.
|
||||
export const forward = Effect.fn("Ssh.forward")(function* (socks: number, remote: { host: string; port: number }) {
|
||||
const host = Buffer.from(remote.host.replace(/^\[(.*)\]$/, "$1"))
|
||||
|
||||
const request = Buffer.concat([
|
||||
Buffer.from([5, 1, 0, 3, host.length]),
|
||||
host,
|
||||
Buffer.from([remote.port >> 8, remote.port & 255]),
|
||||
])
|
||||
|
||||
const sockets = new Set<Socket>()
|
||||
|
||||
const server = createServer((client) => {
|
||||
const upstream = connect(socks, "127.0.0.1")
|
||||
const state = { pending: Buffer.alloc(0), method: false }
|
||||
|
||||
const close = () => {
|
||||
client.destroy()
|
||||
upstream.destroy()
|
||||
sockets.delete(client)
|
||||
sockets.delete(upstream)
|
||||
}
|
||||
|
||||
sockets.add(client).add(upstream)
|
||||
client.pause()
|
||||
client.on("error", close).on("close", close)
|
||||
upstream.on("error", close).on("close", close)
|
||||
upstream.on("connect", () => upstream.write(Buffer.from([5, 1, 0])))
|
||||
|
||||
const handshake = (chunk: Buffer) => {
|
||||
state.pending = Buffer.concat([state.pending, chunk])
|
||||
|
||||
if (!state.method) {
|
||||
if (state.pending.length < 2) return
|
||||
|
||||
if (state.pending[0] !== 5 || state.pending[1] !== 0) return close()
|
||||
state.method = true
|
||||
state.pending = state.pending.subarray(2)
|
||||
upstream.write(request)
|
||||
}
|
||||
|
||||
// OpenSSH always answers CONNECT with an IPv4 address: 10 bytes.
|
||||
if (state.pending.length < 10) return
|
||||
|
||||
if (state.pending[0] !== 5 || state.pending[1] !== 0 || state.pending[3] !== 1) return close()
|
||||
upstream.off("data", handshake)
|
||||
|
||||
if (state.pending.length > 10) client.write(state.pending.subarray(10))
|
||||
client.pipe(upstream)
|
||||
upstream.pipe(client)
|
||||
client.resume()
|
||||
}
|
||||
|
||||
upstream.on("data", handshake)
|
||||
})
|
||||
|
||||
yield* Effect.acquireRelease(
|
||||
Effect.callback<void, SshFailure>((resume) => {
|
||||
server.once("error", (error) => resume(Effect.fail(SshFailure.from(error))))
|
||||
server.listen(0, "127.0.0.1", () => resume(Effect.void))
|
||||
}),
|
||||
() =>
|
||||
Effect.sync(() => {
|
||||
server.close()
|
||||
sockets.forEach((socket) => socket.destroy())
|
||||
}),
|
||||
)
|
||||
|
||||
// SAFETY: a TCP server listening on a host and port reports an AddressInfo.
|
||||
return (server.address() as AddressInfo).port
|
||||
})
|
||||
Reference in new issue
Block a user