mirror of
https://github.com/n8n-io/n8n.git
synced 2026-09-28 13:13:26 +08:00
282 lines
9.3 KiB
JavaScript
282 lines
9.3 KiB
JavaScript
import { spawn, spawnSync } from 'node:child_process';
|
||
import { readFileSync } from 'node:fs';
|
||
import { createServer } from 'node:net';
|
||
import { setTimeout as delay } from 'node:timers/promises';
|
||
import { parseArgs } from 'node:util';
|
||
|
||
import { basicAuth, openCodeProxy } from './cloud-session-opencode-proxy.mjs';
|
||
|
||
const SIGNALS = ['SIGINT', 'SIGTERM', 'SIGHUP'];
|
||
|
||
export function parseOpenCodeArgs(args) {
|
||
let parsed;
|
||
try {
|
||
parsed = parseArgs({
|
||
args,
|
||
allowPositionals: true,
|
||
options: { web: { type: 'boolean' }, new: { type: 'boolean' }, port: { type: 'string' } },
|
||
});
|
||
} catch (error) {
|
||
throw new Error(`${error.message}. Use --help. Use --legacy for remote OpenCode CLI flags.`);
|
||
}
|
||
const { positionals, values } = parsed;
|
||
if (positionals.length > 1 || (positionals[0] && !/^\w[\w-]*$/.test(positionals[0])))
|
||
throw new Error('Use one workspace name with letters, digits, - or _.');
|
||
let port = values.web ? 4096 : 0;
|
||
if (values.port !== undefined) {
|
||
if (!/^\d+$/.test(values.port) || +values.port < 1 || +values.port > 65535)
|
||
throw new Error('--port must be 1–65535.');
|
||
port = +values.port;
|
||
}
|
||
return { name: positionals[0] ?? 'agent', web: !!values.web, fresh: !!values.new, port };
|
||
}
|
||
|
||
function localVersion() {
|
||
const result = spawnSync('opencode', ['--version'], { encoding: 'utf8' });
|
||
if (result.error || result.status !== 0)
|
||
throw new Error('Install OpenCode locally, or use --web to connect with a browser.');
|
||
const version = result.stdout.trim();
|
||
if (!/^\d+\.\d+\.\d+(?:[-+][\w.-]+)?$/.test(version))
|
||
throw new Error('Cannot read the local OpenCode version. Use --web or --legacy.');
|
||
return version;
|
||
}
|
||
|
||
async function freePort() {
|
||
const server = createServer();
|
||
await new Promise((resolve, reject) => {
|
||
server.once('error', reject);
|
||
server.listen(0, '127.0.0.1', resolve);
|
||
});
|
||
const { port } = server.address();
|
||
await new Promise((resolve) => server.close(resolve));
|
||
return port;
|
||
}
|
||
|
||
function startChild(command, args, options = {}) {
|
||
const child = spawn(command, args, options);
|
||
let finished = false;
|
||
const done = new Promise((resolve) => {
|
||
child.once('error', (error) => {
|
||
finished = true;
|
||
resolve({ code: 1, error });
|
||
});
|
||
// Wait for stdout to close before parsing the SSH response.
|
||
child.once('close', (code, signal) => {
|
||
finished = true;
|
||
resolve({ code: code ?? (signal ? 1 : 0) });
|
||
});
|
||
});
|
||
const kill = (signal) => {
|
||
if (!child.pid) return;
|
||
try {
|
||
if (options.detached) process.kill(-child.pid, signal);
|
||
else child.kill(signal);
|
||
} catch (error) {
|
||
if (error.code !== 'ESRCH') throw error;
|
||
}
|
||
};
|
||
return {
|
||
child,
|
||
done,
|
||
get finished() {
|
||
return finished;
|
||
},
|
||
async stop() {
|
||
if (finished) return;
|
||
// gh starts ssh as a child. Terminate the process group to close the tunnel too.
|
||
kill('SIGTERM');
|
||
await Promise.race([done, delay(2000, undefined, { ref: false })]);
|
||
if (!finished) kill('SIGKILL');
|
||
},
|
||
};
|
||
}
|
||
|
||
async function bootstrap(codespace, options, signal) {
|
||
signal.throwIfAborted();
|
||
console.log(`Preparing OpenCode workspace '${options.name}' on ${codespace}…`);
|
||
const command = `umask 077; mkdir -p /workspaces/.n8n-opencode && flock --close -w 600 /workspaces/.n8n-opencode/launch.lock node --input-type=module - ${options.name} ${options.fresh} ${options.web}`;
|
||
const remote = startChild('gh', ['codespace', 'ssh', '-c', codespace, '--', command], {
|
||
stdio: ['pipe', 'pipe', 'inherit'],
|
||
detached: true,
|
||
});
|
||
let stdout = '';
|
||
remote.child.stdout.setEncoding('utf8');
|
||
remote.child.stdout.on('data', (chunk) => {
|
||
stdout += chunk;
|
||
});
|
||
remote.child.stdin.on('error', () => {});
|
||
const abort = () => {
|
||
void remote.stop();
|
||
};
|
||
signal.addEventListener('abort', abort, { once: true });
|
||
try {
|
||
remote.child.stdin.end(
|
||
readFileSync(new URL('../.devcontainer/codespaces/opencode-server.mjs', import.meta.url)),
|
||
);
|
||
const result = await remote.done;
|
||
signal.throwIfAborted();
|
||
if (result.code !== 0)
|
||
throw new Error('Could not prepare OpenCode. Check the SSH output above, then retry.');
|
||
let state;
|
||
try {
|
||
state = JSON.parse(stdout.trim().split('\n').at(-1));
|
||
} catch {
|
||
throw new Error('Invalid OpenCode server response.');
|
||
}
|
||
// The server owns the workspace layout and the credential format. Check the shape only.
|
||
const text = (value) => typeof value === 'string' && value.length > 0;
|
||
if (
|
||
!state ||
|
||
!Number.isInteger(state.port) ||
|
||
!text(state.password) ||
|
||
!text(state.directory) ||
|
||
(options.web && !text(state.sessionID))
|
||
) {
|
||
throw new Error('Invalid OpenCode server response.');
|
||
}
|
||
return state;
|
||
} finally {
|
||
signal.removeEventListener('abort', abort);
|
||
await remote.stop();
|
||
}
|
||
}
|
||
|
||
async function waitForTunnel(url, password, tunnel, signal) {
|
||
for (let attempt = 0; attempt < 60; attempt++) {
|
||
signal.throwIfAborted();
|
||
if (tunnel.finished)
|
||
throw new Error('The SSH tunnel closed. Check the SSH output above, then reconnect.');
|
||
let response;
|
||
try {
|
||
response = await fetch(`${url}/global/health`, {
|
||
headers: { authorization: basicAuth(password) },
|
||
signal: AbortSignal.any([signal, AbortSignal.timeout(1000)]),
|
||
});
|
||
} catch {
|
||
/* The listener is not ready yet. */
|
||
}
|
||
if (response) {
|
||
if (!response.ok)
|
||
throw new Error(
|
||
`OpenCode health check failed (${response.status}). Reconnect after checking the remote server.`,
|
||
);
|
||
const health = await response.json();
|
||
if (health.healthy !== true || typeof health.version !== 'string')
|
||
throw new Error('Invalid OpenCode health response.');
|
||
return health;
|
||
}
|
||
await delay(500, undefined, { signal });
|
||
}
|
||
throw new Error('Timed out while connecting to OpenCode. Run the session command again.');
|
||
}
|
||
|
||
function openBrowser(url) {
|
||
const command = process.platform === 'darwin' ? 'open' : 'xdg-open';
|
||
const browser = spawn(command, [url], { stdio: 'ignore', detached: true });
|
||
browser.on('error', () => console.log('Open the URL above in your browser.'));
|
||
browser.on('exit', (code) => {
|
||
if (code) console.log('Open the URL above in your browser.');
|
||
});
|
||
browser.unref();
|
||
}
|
||
|
||
export async function connectOpenCode(options, ensureCodespace) {
|
||
if (process.platform === 'win32')
|
||
throw new Error('Run this command in WSL. Native Windows is not supported.');
|
||
const version = options.web ? undefined : localVersion();
|
||
const controller = new AbortController();
|
||
const interrupt = () => controller.abort();
|
||
for (const signal of SIGNALS) process.on(signal, interrupt);
|
||
let tunnel;
|
||
let client;
|
||
let proxy;
|
||
try {
|
||
const codespace = ensureCodespace();
|
||
const state = await bootstrap(codespace, options, controller.signal);
|
||
// Bind the browser port first so the tunnel cannot receive the same port.
|
||
if (options.web) proxy = await openCodeProxy({ password: state.password, port: options.port });
|
||
const port = !options.web && options.port ? options.port : await freePort();
|
||
const url = `http://127.0.0.1:${port}`;
|
||
tunnel = startChild(
|
||
'gh',
|
||
[
|
||
'codespace',
|
||
'ssh',
|
||
'-c',
|
||
codespace,
|
||
'--',
|
||
'-N',
|
||
'-o',
|
||
'ExitOnForwardFailure=yes',
|
||
'-o',
|
||
'ServerAliveInterval=15',
|
||
'-o',
|
||
'ServerAliveCountMax=3',
|
||
'-L',
|
||
`127.0.0.1:${port}:127.0.0.1:${state.port}`,
|
||
],
|
||
{ stdio: ['ignore', 'ignore', 'inherit'], detached: true },
|
||
);
|
||
const health = await waitForTunnel(url, state.password, tunnel, controller.signal);
|
||
if (version && health.version !== version) {
|
||
throw new Error(
|
||
`OpenCode versions differ: local ${version}, server ${health.version}. Install the matching version with \`npm install -g opencode-ai@${health.version}\`, or use --web or --legacy. The remote server is still running.`,
|
||
);
|
||
}
|
||
const stopped = new Promise((resolve) =>
|
||
controller.signal.addEventListener('abort', () => resolve('interrupted'), { once: true }),
|
||
);
|
||
controller.signal.throwIfAborted();
|
||
if (proxy) {
|
||
proxy.targetPort = port;
|
||
const webUrl = `${proxy.origin}/${Buffer.from(state.directory).toString('base64url')}/session/${state.sessionID}`;
|
||
console.log(
|
||
`OpenCode: ${webUrl}\nKeep this command running. Press Ctrl-C to disconnect. The remote server stays running.`,
|
||
);
|
||
openBrowser(webUrl);
|
||
} else {
|
||
console.log(`Connecting to OpenCode ${health.version} in ${state.directory}…`);
|
||
// OpenCode's own resume spans all worktrees of the repository. The
|
||
// server resolved the workspace's latest conversation instead, so
|
||
// target it directly. A fresh TUI conversation needs no target.
|
||
client = startChild(
|
||
'opencode',
|
||
[
|
||
'attach',
|
||
url,
|
||
'--dir',
|
||
state.directory,
|
||
...(state.sessionID ? ['--session', state.sessionID] : []),
|
||
],
|
||
{
|
||
stdio: 'inherit',
|
||
env: {
|
||
...process.env,
|
||
OPENCODE_SERVER_PASSWORD: state.password,
|
||
OPENCODE_SERVER_USERNAME: 'opencode',
|
||
},
|
||
},
|
||
);
|
||
}
|
||
const result = await Promise.race([
|
||
stopped,
|
||
tunnel.done.then(() => 'tunnel'),
|
||
...(client ? [client.done] : []),
|
||
]);
|
||
if (typeof result === 'object' && result.code !== 0)
|
||
throw new Error(
|
||
'The local OpenCode client exited with an error. The remote session is saved.',
|
||
);
|
||
if (result === 'tunnel')
|
||
throw new Error(
|
||
'The SSH connection closed. Run the same command to reconnect to the saved session.',
|
||
);
|
||
} catch (error) {
|
||
if (!controller.signal.aborted) throw error;
|
||
} finally {
|
||
proxy?.close();
|
||
await Promise.all([client?.stop(), tunnel?.stop()]);
|
||
for (const signal of SIGNALS) process.removeListener(signal, interrupt);
|
||
}
|
||
}
|