theseus/bundled-addons/pithos/core/daemon.js

186 lines
7.1 KiB
JavaScript

// Owns the s3d process: start/stop, a log ring buffer, and one-shot CLI runs
// (users/keys/version) against the same config file the daemon uses.
import { spawn, execFile } from 'node:child_process';
import { EventEmitter } from 'node:events';
import path from 'node:path';
import { sendBreak } from './winctrl.js';
const LOG_LINES = 2000;
const ANSI = /\x1b\[[0-9;]*m/g;
export class Daemon extends EventEmitter {
constructor({ binary, configFile }) {
super();
this.binary = binary;
this.configFile = configFile;
this.child = null;
this.state = 'stopped'; // stopped | starting | running | stopping | crashed
this.exitCode = null;
this.startedAt = null;
this.logs = [];
this.seq = 0;
this.busy = 0; // flushes Pithos started and is waiting on
this.lastUploadAt = 0; // s3d's last "uploading object group"
}
env() {
// S3D_CONFIG_FILE pins every invocation to the config this GUI edits; the
// cwd is the config's folder because s3d also searches ./s3d.yml first.
return { ...process.env, S3D_CONFIG_FILE: this.configFile };
}
cwd() { return path.dirname(this.configFile); }
status() {
return {
state: this.state,
pid: this.child?.pid ?? null,
exitCode: this.exitCode,
startedAt: this.startedAt,
binary: this.binary,
needsLogin: !!this.needsLogin,
configFile: this.configFile,
};
}
setState(state) {
this.state = state;
this.emit('status', this.status());
}
pushLog(stream, chunk) {
for (const raw of chunk.toString('utf8').split(/\r?\n/)) {
const line = raw.replace(ANSI, '');
if (!line) continue;
const entry = { seq: ++this.seq, t: Date.now(), stream, line };
// s3d logs "server started" once every listener is up.
if (this.state === 'starting' && /server started/.test(line)) this.setState('running');
if (/No app key found/.test(line)) this.needsLogin = true;
if (/uploading object group/.test(line)) this.lastUploadAt = Date.now();
this.logs.push(entry);
if (this.logs.length > LOG_LINES) this.logs.shift();
this.emit('log', entry);
}
}
logsSince(seq = 0) { return this.logs.filter((l) => l.seq > seq); }
// s3d migrates its database when it starts, and so does every one-shot
// command (users list, keys…) that opens it. Two of them at once failed
// with "table users already exists", so a start waits for running
// commands to finish, and commands wait while s3d is starting.
start() {
if (!this.binary) throw new Error('s3d binary not found');
if (this.child || this.pendingStart) return this.status();
this.exitCode = null;
this.needsLogin = false;
this.setState('starting');
if (this.cliRuns > 0) {
this.pendingStart = true;
this.once('cli-idle', () => { this.pendingStart = false; if (this.state === 'starting' && !this.child) this.spawnServer(); });
} else this.spawnServer();
return this.status();
}
spawnServer() {
const child = spawn(this.binary, [], {
cwd: this.cwd(), env: this.env(), stdio: ['ignore', 'pipe', 'pipe'], windowsHide: true,
});
this.child = child;
this.startedAt = Date.now();
child.stdout.on('data', (d) => this.pushLog('out', d));
child.stderr.on('data', (d) => this.pushLog('err', d));
child.on('error', (err) => {
this.pushLog('err', `failed to start s3d: ${err.message}`);
// A spawn failure never emits 'exit'.
if (child.pid === undefined) { this.child = null; this.setState('crashed'); }
});
child.on('exit', (code, signal) => {
this.child = null;
this.exitCode = code;
this.pushLog('sys', `s3d exited (${signal || `code ${code}`})`);
this.setState(this.state === 'stopping' ? 'stopped' : 'crashed');
});
return this.status();
}
// Work s3d must not be stopped in the middle of (Upload now, autoflush).
async whileBusy(fn) {
this.busy++;
try { return await fn(); } finally { this.busy--; }
}
// True once no flush of ours is running and s3d has not started an upload
// batch for a while (it logs the start of one, not the end).
quiet() { return this.busy === 0 && Date.now() - this.lastUploadAt > 120_000; }
async waitForQuiet(maxMs) {
const until = Date.now() + maxMs;
let told = false;
while (!this.quiet() && Date.now() < until && this.child) {
if (!told) { this.pushLog('sys', 'waiting for the upload in progress to finish before stopping s3d'); told = true; }
await new Promise((r) => setTimeout(r, 1000));
}
return this.quiet();
}
// An s3d stopped in the middle of an upload has lost buffered files whose
// rows still said "waiting", so a stop first lets uploads finish, then
// asks s3d to shut down cleanly (Ctrl+Break on Windows, SIGTERM elsewhere)
// and only ends it by force if it is still there a minute later.
// waitMs: how long to wait for uploads (short when the host is quitting).
async stop({ waitMs = 5 * 60_000 } = {}) {
const child = this.child;
if (!child && this.pendingStart) { this.pendingStart = false; this.setState('stopped'); return this.status(); }
if (!child) return this.status();
this.setState('stopping');
const exited = new Promise((r) => { if (this.child !== child) r(); else child.once('exit', r); });
if (!(await this.waitForQuiet(waitMs)) && this.child === child) this.pushLog('sys', 'an upload is still running; stopping s3d anyway');
if (this.child === child) {
let asked = false;
if (process.platform === 'win32') asked = await sendBreak(child.pid);
else { try { asked = child.kill('SIGTERM'); } catch {} }
if (!asked) { try { child.kill(); } catch {} }
}
const timer = setTimeout(() => {
if (this.child !== child) return;
this.pushLog('sys', 's3d did not stop within a minute; ending it');
try { child.kill('SIGKILL'); } catch {}
}, 60_000);
await exited;
clearTimeout(timer);
return this.status();
}
// One-shot CLI run. Resolves { code, stdout, stderr }; never rejects on a
// non-zero exit so callers can surface s3d's own error text.
async run(args, { timeout = 30_000 } = {}) {
if (!this.binary) throw new Error('s3d binary not found');
if (this.state === 'starting') await this.settled(20_000);
this.cliRuns = (this.cliRuns || 0) + 1;
try {
return await new Promise((resolve, reject) => {
execFile(this.binary, args, {
cwd: this.cwd(), env: this.env(), timeout, windowsHide: true, maxBuffer: 16 << 20,
}, (err, stdout, stderr) => {
if (err && typeof err.code !== 'number') return reject(err);
resolve({ code: err ? err.code : 0, stdout: stdout.replace(ANSI, ''), stderr: stderr.replace(ANSI, '') });
});
});
} finally {
if (--this.cliRuns === 0) this.emit('cli-idle');
}
}
// Resolves once s3d is no longer starting (or after ms).
settled(ms) {
return new Promise((resolve) => {
if (this.state !== 'starting') return resolve();
const done = () => { clearTimeout(t); this.off('status', on); resolve(); };
const on = (st) => { if (st.state !== 'starting') done(); };
const t = setTimeout(done, ms);
this.on('status', on);
});
}
}