/**
* DC-055: Host journald reader
*
* Wraps the host's `journalctl` binary so the API can stream host service
* logs (caddy, dashcaddy-api, docker, ...) without exposing the binary
* directly to the web layer. The CLI is invoked with --directory pointed at
* the bind-mounted /var/log/journal from start.sh so we don't need the
* systemd-journal remote protocol or a privileged socket.
*
* Security contract:
* - `unit` MUST be in the allow-list `ALLOWED_UNITS`. We never accept a
* raw unit name from the caller and pass it to the shell, even with
* shell:false — because an attacker who can set unit=caddy.service;
* touch /tmp/x could use the CLI itself as a confused-deputy vector.
* - All journalctl invocations use `spawn` (not `exec`) and pass arguments
* as an array (`shell:false`). No shell metacharacters can be smuggled
* in through any field — the unit, since/until, search, tail numbers
* are validated separately before being added to argv.
* - Streams (SSE) cap to MAX_STREAM_BYTES and kill the child on overflow
* so a `tail=999999999999` request can't OOM the process.
*
* Failure modes that surface to the route layer:
* - journalctl missing in the container (DN container, dev container):
* every call throws Error('journalctl unavailable'). Route 503s.
* - unit not in allow-list: throws ValidationError. Route 400s.
* - non-zero exit code: child stderr is captured and surfaced verbatim
* up to LOG_PREVIEW_BYTES so the operator can see "Failed to open
* directory" instead of a generic 500.
*/
const { spawn } = require('child_process');
const path = require('path');
const fs = require('fs');
const JOURNAL_DIR = '/var/log/journal';
const ALLOWED_UNITS = Object.freeze([
// Core reverse proxy + DNS host services
'caddy',
'dashcaddy-api',
'docker',
'systemd-journald',
'networkd-dispatcher',
'tailscaled',
'ssh',
// Permit the unit with and without the .service suffix. The CLI accepts
// both; we store the bare name and append nothing — journalctl treats
// "caddy" and "caddy.service" identically.
]);
// Cap how much a single request can read — prevents `tail=999999999` from
// piping half the journal into memory. The dashboard doesn't have a UI for
// "load 100MB of logs" and journalctl itself caps at 2GB anyway.
const MAX_TAIL_LINES = 5000;
// Streaming cap: how many journal entries we hand to the SSE consumer
// before killing the child. The dashboard shouldn't accumulate more than
// this in memory — pair with MAX_OUTPUT_BUFFER for a defense-in-depth
// bound on what the route layer will hold.
const MAX_STREAM_LINES = 5000;
const MAX_OUTPUT_BUFFER = 2 * 1024 * 1024; // 2MB hard cap on total stdout
const LOG_PREVIEW_BYTES = 4096;
const UNIT_PATTERN = /^[a-zA-Z0-9_.@-]+$/;
const ISO_PATTERN = /^\d{4}-\d{2}-\d{2}(?:[T ]\d{2}:\d{2}(?::\d{2}(?:\.\d+)?)?(?:Z|[+-]\d{2}:?\d{2})?)?$/;
/**
* Validate a unit name against the allow-list. Returns the canonical name
* or throws ValidationError.
*/
function assertUnitAllowed(unit) {
if (typeof unit !== 'string' || !unit) {
const err = new Error('unit is required');
err.name = 'ValidationError';
throw err;
}
// Strip the .service suffix defensively so callers don't have to remember
// which form journalctl prefers for a given unit.
const normalised = unit.endsWith('.service') ? unit.slice(0, -8) : unit;
if (!UNIT_PATTERN.test(normalised)) {
const err = new Error(`unit contains invalid characters: ${unit}`);
err.name = 'ValidationError';
throw err;
}
if (!ALLOWED_UNITS.includes(normalised)) {
const err = new Error(`unit not in allow-list: ${normalised}`);
err.name = 'ValidationError';
throw err;
}
return normalised;
}
/**
* Parse tail to a bounded positive integer.
*/
function parseTail(raw, fallback = 200) {
if (raw === undefined || raw === null || raw === '') return fallback;
const n = Number(raw);
if (!Number.isFinite(n) || !Number.isInteger(n) || n <= 0) {
const err = new Error(`tail must be a positive integer (got ${raw})`);
err.name = 'ValidationError';
throw err;
}
return Math.min(n, MAX_TAIL_LINES);
}
/**
* Parse since/until — accept either an ISO timestamp, a unix epoch in ms, or
* journalctl's relative syntax ("30 min ago", "today", "yesterday"). The
* dashboard uses ISO timestamps from ``; the
* relative syntax is for power users typing into the search bar.
*/
function parseTimestamp(raw, fieldName) {
if (raw === undefined || raw === null || raw === '') return null;
if (typeof raw !== 'string') {
const err = new Error(`${fieldName} must be a string`);
err.name = 'ValidationError';
throw err;
}
// ISO 8601
if (ISO_PATTERN.test(raw)) {
const ms = Date.parse(raw);
if (!Number.isFinite(ms)) {
const err = new Error(`${fieldName} is not a valid ISO timestamp: ${raw}`);
err.name = 'ValidationError';
throw err;
}
return new Date(ms).toISOString();
}
// Numeric (unix epoch seconds OR ms — journalctl accepts seconds)
if (/^-?\d+$/.test(raw)) {
const n = Number(raw);
const ms = n > 1e12 ? n : n * 1000;
if (!Number.isFinite(ms)) {
const err = new Error(`${fieldName} is not a valid epoch: ${raw}`);
err.name = 'ValidationError';
throw err;
}
return new Date(ms).toISOString();
}
// Relative syntax: pass through to journalctl, but cap to 1024 chars and
// disallow shell metacharacters.
if (raw.length > 1024 || /[`$;&|><\\\n\r]/.test(raw)) {
const err = new Error(`${fieldName} contains forbidden characters: ${raw}`);
err.name = 'ValidationError';
throw err;
}
return raw;
}
/**
* Detect whether journalctl is reachable. Cheap probe (no-op flag) so we
* don't shell out on every request when the binary is missing (dev
* container, Windows host, etc.).
*/
function isAvailable({ journalDir = JOURNAL_DIR, exec = spawn } = {}) {
if (!fs.existsSync(journalDir)) return false;
return new Promise((resolve) => {
const child = exec('journalctl', ['--no-pager', '--version'], { stdio: 'ignore' });
child.on('error', () => resolve(false));
child.on('exit', (code) => resolve(code === 0));
});
}
/**
* Build argv for journalctl. Exposed so tests can assert exactly what we
* shell out — never build the arg array inline anywhere else.
*/
function buildArgv({ unit, since, until, tail, search, follow = false }) {
const argv = [
'--directory', JOURNAL_DIR,
'--no-pager',
'--output=short',
'-u', unit,
];
if (since) argv.push('--since', since);
if (until) argv.push('--until', until);
if (typeof tail === 'number') argv.push('-n', String(tail));
if (search) {
// journalctl -S matches the searchable text fields (MESSAGE + others).
// Quote-enforcing isn't needed because spawn argv doesn't touch a shell.
argv.push('-S', search);
}
if (follow) argv.push('--follow');
return argv;
}
/**
* Read a bounded tail of journal entries for a unit. Resolves to an array
* of {timestamp, text} lines, oldest first. Throws ValidationError on bad
* input, Error('journalctl unavailable') if the binary or journal dir is
* missing, and Error('journalctl exited N: ') for CLI failures.
*/
/**
* Spawn journalctl with the given argv and collect stdout/stderr up to
* the configured caps. Resolves to a Buffer of stdout on success, rejects
* with Error('journalctl unavailable') on ENOENT or
* Error('journalctl exited N: ') on non-zero exit. Exceeding the
* output cap rejects with an explicit overflow message.
*
* Kept as a free function (not inside `readEntries`) so the same plumbing
* can be reused for streaming without code duplication.
*/
function runJournalctl({ exec, argv }) {
return new Promise((resolve, reject) => {
const child = exec('journalctl', argv, { stdio: ['ignore', 'pipe', 'pipe'] });
let stdout = Buffer.alloc(0);
let stderr = '';
child.stdout.on('data', (chunk) => {
if (stdout.length + chunk.length > MAX_OUTPUT_BUFFER) {
child.kill('SIGKILL');
reject(new Error(`output exceeded ${MAX_OUTPUT_BUFFER} bytes`));
return;
}
stdout = Buffer.concat([stdout, chunk]);
});
child.stderr.on('data', (chunk) => {
if (stderr.length < LOG_PREVIEW_BYTES) {
stderr += chunk.toString('utf8');
if (stderr.length > LOG_PREVIEW_BYTES) {
stderr = stderr.slice(0, LOG_PREVIEW_BYTES) + '…';
}
}
});
child.on('error', (err) => {
if (err.code === 'ENOENT') {
reject(new Error('journalctl unavailable'));
} else {
reject(err);
}
});
child.on('exit', (code, signal) => {
if (signal === 'SIGKILL' && stdout.length >= MAX_OUTPUT_BUFFER) return; // already rejected
if (code !== 0) {
reject(new Error(`journalctl exited ${code}${stderr ? ': ' + stderr.trim() : ''}`));
return;
}
resolve({ stdout, stderr });
});
});
}
/**
* Parse a journalctl --output=short line into a structured entry.
* Lines look like: "Aug 18 00:42:46 vmi3080415 caddy[3620580]: {...}"
*/
function parseShortLine(line, fallbackUnit) {
const tsMatch = line.match(/^([A-Z][a-z]{2} \d{2} \d{2}:\d{2}:\d{2}) (\S+) (.+?)\[\d+\]: (.*)$/);
if (tsMatch) {
return {
timestamp: tsMatch[1],
hostname: tsMatch[2],
unit: tsMatch[3],
text: tsMatch[4],
};
}
return { timestamp: null, hostname: null, unit: fallbackUnit, text: line };
}
function readEntries(opts, { exec = spawn } = {}) {
return Promise.resolve().then(async () => {
const unit = assertUnitAllowed(opts.unit);
const tail = parseTail(opts.tail);
const since = parseTimestamp(opts.since, 'since');
const until = parseTimestamp(opts.until, 'until');
const search = typeof opts.search === 'string' && opts.search.length > 0
? opts.search.slice(0, 1024)
: null;
const argv = buildArgv({ unit, tail, since, until, search, follow: false });
const { stdout } = await runJournalctl({ exec, argv });
const lines = stdout.toString('utf8').split('\n').filter(Boolean);
return lines.map((line) => parseShortLine(line, unit));
});
}
/**
* Stream journal entries as they arrive. Returns { child, onData, onError,
* kill } — the route wires `onData`/`onError` to the SSE socket and calls
* `kill()` on disconnect.
*
* The child is spawned with --follow and we cap total bytes received; on
* overflow we kill the child and emit a synthetic 'overflow' message so the
* client knows to reconnect with a narrower window.
*/
function streamEntries(opts, { exec = spawn, onData, onError } = {}) {
const unit = assertUnitAllowed(opts.unit);
const since = parseTimestamp(opts.since, 'since');
const search = typeof opts.search === 'string' && opts.search.length > 0
? opts.search.slice(0, 1024)
: null;
const argv = buildArgv({ unit, since, search, follow: true });
let child;
try {
child = exec('journalctl', argv, { stdio: ['ignore', 'pipe', 'pipe'] });
} catch (err) {
if (err.code === 'ENOENT') {
const e = new Error('journalctl unavailable');
onError && onError(e);
return { kill: () => {}, child: null };
}
throw err;
}
// Closure-scoped stream bookkeeping: the previous version attached a
// counter to the onData function itself, which made the 5000-line cap
// unreachable (a function has its own properties — the count was never
// incremented). Closure scope is the right place.
let stdout = Buffer.alloc(0);
let lineCount = 0;
let overflowEmitted = false;
child.stdout.on('data', (chunk) => {
if (stdout.length + chunk.length > MAX_OUTPUT_BUFFER) {
child.kill('SIGKILL');
onError && onError(new Error(`stream exceeded ${MAX_OUTPUT_BUFFER} bytes`));
return;
}
stdout = Buffer.concat([stdout, chunk]);
if (onData) {
const text = stdout.toString('utf8');
const lines = text.split('\n');
// Hold back the last partial line; flush on the next chunk or exit.
stdout = Buffer.from(lines.pop(), 'utf8');
for (const line of lines) {
if (!line) continue;
lineCount++;
if (lineCount > MAX_STREAM_LINES && !overflowEmitted) {
overflowEmitted = true;
child.kill('SIGKILL');
onError && onError(new Error(`stream exceeded ${MAX_STREAM_LINES} lines`));
return;
}
const tsMatch = line.match(/^([A-Z][a-z]{2} \d{2} \d{2}:\d{2}:\d{2}) (\S+) (.+?)\[\d+\]: (.*)$/);
onData({
timestamp: tsMatch ? tsMatch[1] : null,
hostname: tsMatch ? tsMatch[2] : null,
unit: tsMatch ? tsMatch[3] : unit,
text: tsMatch ? tsMatch[4] : line,
});
}
}
});
let stderr = '';
child.stderr.on('data', (chunk) => {
if (stderr.length < LOG_PREVIEW_BYTES) {
stderr += chunk.toString('utf8');
if (stderr.length > LOG_PREVIEW_BYTES) {
stderr = stderr.slice(0, LOG_PREVIEW_BYTES) + '…';
}
}
});
child.on('error', (err) => {
onError && onError(err);
});
child.on('exit', (code) => {
if (code !== 0 && stderr) {
onError && onError(new Error(`journalctl exited ${code}: ${stderr.trim()}`));
}
});
return {
child,
kill() {
try { child.kill('SIGTERM'); } catch (_) { /* already dead */ }
},
};
}
/**
* List units that currently have journal entries (for the dashboard
* dropdown). Walks the allow-list and asks journalctl for the most recent
* entry per unit. Units with no entries are omitted.
*/
async function listUnits({ exec = spawn } = {}) {
if (!fs.existsSync(JOURNAL_DIR)) return [];
const out = [];
for (const unit of ALLOWED_UNITS) {
const lines = await new Promise((resolve) => {
const child = exec('journalctl', [
'--directory', JOURNAL_DIR,
'--no-pager', '-q',
'-u', unit,
'-n', '1',
'--output=short',
], { stdio: ['ignore', 'pipe', 'ignore'] });
let buf = '';
child.stdout.on('data', (c) => { buf += c.toString('utf8'); });
child.on('error', () => resolve(''));
child.on('exit', () => resolve(buf));
});
if (lines.trim()) {
out.push({ unit, hasEntries: true });
}
}
return out;
}
module.exports = {
ALLOWED_UNITS,
MAX_TAIL_LINES,
MAX_OUTPUT_BUFFER,
isAvailable,
readEntries,
streamEntries,
listUnits,
assertUnitAllowed,
parseTail,
parseTimestamp,
parseShortLine,
buildArgv,
};