diff --git a/cli/README.md b/cli/README.md index 606cc2c..e65f0a6 100644 --- a/cli/README.md +++ b/cli/README.md @@ -385,6 +385,7 @@ The bundle uses scrypt (N=32768, r=8, p=1) → AES-256-GCM. Salt + nonce are ran **What's included:** - Every profile's `baseUrl`, `keyPrefix`, full `mantis_live_…` API key, Cloudflare Access mode + app URL + Service-Auth client-id/secret, linked edge worker URL, and edge AES key +- All locally stored edge worker keys, including workers without a server profile. Edge-only installs can use backup/restore without logging into a server. Existing edge keys are kept on restore unless `--overwrite` is set. Older backup files remain supported. - The active-profile pointer - Plugin manifest: each plugin's `name`, `source` (GitHub `owner/repo`), pinned commit SHA, and version. **`restore` re-installs plugins via `mantis plugin add @`** — the bundle does NOT carry the plugin contents themselves, so the new machine needs network access to GitHub for the re-install. @@ -393,12 +394,17 @@ The bundle uses scrypt (N=32768, r=8, p=1) → AES-256-GCM. Salt + nonce are ran - Local-path plugins (`mantis plugin add ./some/path`) — those aren't reproducible on another machine; `backup` lists them as skipped. - `~/.cloudflared/` cached JWTs — owned by `cloudflared`, regenerated on next login. +Full backup lists the edge worker URLs it included. Independent keys require +keychain enumeration support; if your platform cannot enumerate entries, +linked profile keys are still included. Check that list before migrating an +edge-only setup and keep a separate vaulted copy of any omitted key. + **Flag reference:** | Flag | What it does | |---|---| | `mantis backup --out ` | Where to write the bundle. Default `./mantis-backup.json`. | -| `mantis backup --only ` | Back up just one profile. Default is all profiles. | +| `mantis backup --only ` | Back up just one profile and its linked edge worker. Independent edge workers are omitted; the command reports this scope. Default includes all profiles and discoverable edge worker keys. | | `mantis backup --passphrase-stdin` | Read passphrase from stdin (for scripts piping a vault into the CLI). | | `mantis backup --passphrase-env ` | Read passphrase from the named env var. | | `mantis restore ` | Decrypt + restore. By default, existing profiles on the target machine are kept; bundle entries with the same name are skipped. | diff --git a/cli/src/commands/backup.ts b/cli/src/commands/backup.ts index f24384c..0622add 100644 --- a/cli/src/commands/backup.ts +++ b/cli/src/commands/backup.ts @@ -8,7 +8,7 @@ import { setProfile, useProfile, } from "../lib/config.js"; -import { setEdgeKey } from "../lib/edge-key.js"; +import { getEdgeKey, setEdgeKey } from "../lib/edge-key.js"; import { collectBackupPayload, collectSkippedLocalPlugins, @@ -74,6 +74,12 @@ export async function backupCmd(opts: BackupCmdOpts): Promise { process.stderr.write( ` ${c.dim("plugins: ")} ${payload.plugins.length === 0 ? c.dim("(none)") : payload.plugins.map((p) => p.name).join(", ")}\n`, ); + process.stderr.write( + ` ${c.dim("edge workers:")} ${payload.edgeWorkers?.map((worker) => worker.workerUrl).join(", ") || c.dim("(none)")}\n`, + ); + if (opts.profile) { + process.stderr.write(` ${c.dim("scope:")} only profile ${opts.profile} and its linked worker; omit --only to include independent edge workers.\n`); + } if (skippedLocalPlugins.length > 0) { process.stderr.write( ` ${c.yellow("note:")} skipped ${skippedLocalPlugins.length} local-path plugin(s) (not reproducible on another machine): ${skippedLocalPlugins.join(", ")}\n`, @@ -88,6 +94,8 @@ export async function backupCmd(opts: BackupCmdOpts): Promise { profiles: payload.profiles.map((p) => p.name), plugins: payload.plugins.length, skipped_local_plugins: skippedLocalPlugins, + edge_workers: payload.edgeWorkers?.map((worker) => worker.workerUrl) ?? [], + scope: opts.profile ?? "all", }, ); } @@ -159,6 +167,29 @@ export async function restoreCmd( } } + const edgeRestored: string[] = []; + const edgeSkipped: string[] = []; + const edgeErrors: Array<{ worker: string; reason: string }> = []; + // Legacy v1 files embedded keys in profiles. Apply the same protection to + // those keys as independent worker entries, even when the profile is new. + const workerKeys = new Map(); + for (const profile of payload.profiles) { + if (profile.edgeWorkerUrl && profile.edgeKey) workerKeys.set(profile.edgeWorkerUrl, profile.edgeKey); + } + for (const { workerUrl, key } of payload.edgeWorkers ?? []) workerKeys.set(workerUrl, key); + for (const [workerUrl, key] of workerKeys) { + if (getEdgeKey(workerUrl) && !opts.overwrite) { + edgeSkipped.push(workerUrl); + continue; + } + try { + setEdgeKey(workerUrl, key); + edgeRestored.push(workerUrl); + } catch (err) { + edgeErrors.push({ worker: workerUrl, reason: err instanceof Error ? err.message : String(err) }); + } + } + // Restore active-profile pointer if the backup specified one AND we // actually restored it AND the user didn't already have a different // current profile they care about. @@ -202,6 +233,9 @@ export async function restoreCmd( `${c.red("✗")} ${e.name}: ${e.reason}\n`, ); } + if (edgeRestored.length > 0) process.stderr.write(`${c.green("✓")} restored edge keys: ${edgeRestored.join(", ")}\n`); + if (edgeSkipped.length > 0) process.stderr.write(`${c.yellow("·")} kept existing edge keys (pass --overwrite to replace): ${edgeSkipped.join(", ")}\n`); + for (const e of edgeErrors) process.stderr.write(`${c.red("✗")} edge ${e.worker}: ${e.reason}\n`); if (!opts.skipPlugins) { if (pluginsRestored.length > 0) { process.stderr.write( @@ -234,8 +268,12 @@ export async function restoreCmd( plugins_restored: pluginsRestored, plugins_failed: pluginsFailed, active_profile: payload.currentProfile, + edge_restored: edgeRestored, + edge_skipped: edgeSkipped, + edge_errors: edgeErrors, }, ); + if (errors.length || edgeErrors.length || pluginsFailed.length) process.exitCode = 1; } async function applyProfile(bp: BackupProfile): Promise { @@ -245,9 +283,6 @@ async function applyProfile(bp: BackupProfile): Promise { // about this profile" marker). setKey(entry.baseUrl, secrets.apiKey); if (secrets.cf) setCloudflareServiceAuth(entry.baseUrl, secrets.cf); - if (secrets.edgeKey && entry.edgeWorkerUrl) { - setEdgeKey(entry.edgeWorkerUrl, secrets.edgeKey); - } await setProfile(bp.name, entry); } diff --git a/cli/src/commands/device.ts b/cli/src/commands/device.ts index 72f03ad..f1cc757 100644 --- a/cli/src/commands/device.ts +++ b/cli/src/commands/device.ts @@ -112,7 +112,7 @@ export async function deviceNewCmd(opts: DeviceNewOpts): Promise { emit( () => { process.stderr.write( - `${c.green("✓")} ${c.bold(device)} armed — ${minted.length} alarm(s)\n`, + `${c.green("✓")} ${c.bold(device)} ${installed ? "armed" : "minted"} — ${minted.length} alarm(s)\n`, ); for (const m of minted) { process.stderr.write(` ${c.dim(m.memo)}\n ${c.cyan(m.url)}\n`); @@ -120,9 +120,9 @@ export async function deviceNewCmd(opts: DeviceNewOpts): Promise { if (bundlePath) { process.stderr.write(`\n${c.green("✓")} bundle → ${c.cyan(bundlePath)}\n`); } - if (!opts.install && !opts.bundle) { + if (!installed) { process.stderr.write( - `\n${c.dim("Nothing installed. Re-run with --bundle for a zip, or --install to apply here.")}\n`, + `\n${c.dim(bundlePath ? "Nothing installed. Run the install script in the exported bundle to activate these alarms." : opts.install ? "Installation canceled. Nothing installed; use the staged bundle path above to install later." : "Nothing installed. Re-run with --bundle for a zip, or --install to apply here.")}\n`, ); } }, diff --git a/cli/src/commands/edge-device.ts b/cli/src/commands/edge-device.ts index 05584ef..433507b 100644 --- a/cli/src/commands/edge-device.ts +++ b/cli/src/commands/edge-device.ts @@ -195,7 +195,7 @@ export async function edgeDeviceCmd(opts: EdgeDeviceOpts): Promise { emit( () => { process.stderr.write( - `${c.green("✓")} ${c.bold(device)} armed (edge) — ${minted.length} alarm(s)\n`, + `${c.green("✓")} ${c.bold(device)} ${installed ? "armed" : "minted"} (edge) — ${minted.length} alarm(s)\n`, ); for (const m of minted) { process.stderr.write(` ${c.dim(m.memo)}\n ${c.cyan(m.url)}\n`); @@ -213,9 +213,9 @@ export async function edgeDeviceCmd(opts: EdgeDeviceOpts): Promise { `${c.yellow("not on edge:")} ${v.slug} normally dedupes hits in a ${v.dedupeWindowSeconds}s window server-side; the stateless worker cannot remember the last hit, so expect bursts (e.g. Wi-Fi roams) to notify several times.\n`, ); } - if (!opts.install && !bundleDir) { + if (!installed) { process.stderr.write( - `\n${c.dim("Nothing installed. Re-run with --bundle for an install directory, or --install to apply here.")}\n`, + `\n${c.dim(bundleDir ? "Nothing installed. Run the install script in the exported bundle to activate these alarms." : opts.install ? "Installation canceled. Nothing installed; use the staged bundle path above to install later." : "Nothing installed. Re-run with --bundle for an install directory, or --install to apply here.")}\n`, ); } }, diff --git a/cli/src/commands/edge.ts b/cli/src/commands/edge.ts index 9b80f6b..e2037c0 100644 --- a/cli/src/commands/edge.ts +++ b/cli/src/commands/edge.ts @@ -437,7 +437,8 @@ export async function installCmd( async function renderInstaller( url: string, opts: InstallerOpts, -): Promise<{ filename: string; written: string | null }> { + emitResult = true, +): Promise<{ filename: string; written: string | null; content?: string; mime: string }> { const type = opts.type; if (!isInstallType(type)) { fail( @@ -474,7 +475,7 @@ async function renderInstaller( writtenTo = target; } - emit( + if (emitResult || !isJsonMode()) emit( () => { if (writtenTo) { process.stderr.write( @@ -497,7 +498,7 @@ async function renderInstaller( }, ); - return { filename: installer.filename, written: writtenTo }; + return { filename: installer.filename, written: writtenTo, content: writtenTo ? undefined : content, mime: installer.mime }; } // --------------------------------------------------------------------------- @@ -616,7 +617,7 @@ export async function mintCmd(opts: MintOpts): Promise { // Chained installer (--install ) runs after URL is produced. We // render it inline here so the user gets URL → test → installer in one // coherent stream, rather than spawning a second command. - let installResult: { filename: string; written: string | null } | null = null; + let installResult: Awaited> | null = null; if (opts.install) { installResult = await renderInstaller(url, { type: opts.install, @@ -624,7 +625,7 @@ export async function mintCmd(opts: MintOpts): Promise { sshOnly: opts.sshOnly, hostname: opts.hostname, memo: opts.memo, - }); + }, false); } emit( @@ -667,6 +668,9 @@ export async function mintCmd(opts: MintOpts): Promise { installer: { type: opts.install, written_to: installResult.written, + filename: installResult.filename, + mime: installResult.mime, + content: installResult.content, }, } : {}), diff --git a/cli/src/lib/backup.ts b/cli/src/lib/backup.ts index c9c774c..7c8da14 100644 --- a/cli/src/lib/backup.ts +++ b/cli/src/lib/backup.ts @@ -32,7 +32,7 @@ import { type StoredConfig, type CloudflareAccessMode, } from "./config.js"; -import { getEdgeKey } from "./edge-key.js"; +import { getEdgeKey, listEdgeKeyWorkers } from "./edge-key.js"; import { readLockfile } from "./plugins/lockfile.js"; // Manually promisify so we control the options-arg signature. Node's @@ -120,6 +120,8 @@ export type BackupPayload = { currentProfile?: string; profiles: BackupProfile[]; plugins: BackupPlugin[]; + /** Independent worker keys, including edge-only setups. Absent in older v1 bundles. */ + edgeWorkers?: Array<{ workerUrl: string; key: string }>; }; // --------------------------------------------------------------------------- @@ -200,6 +202,21 @@ export async function collectBackupPayload( ); } + // A worker does not need a server profile. Keep its key independently so + // edge-only setups and additional workers survive migration too. A scoped + // backup includes only the selected profile's linked worker. + const workerUrls = new Set( + filtered.flatMap(({ entry }) => entry.edgeWorkerUrl ? [entry.edgeWorkerUrl] : []), + ); + if (onlyProfile === undefined) { + for (const workerUrl of await listEdgeKeyWorkers()) workerUrls.add(workerUrl); + } + const edgeWorkers: NonNullable = []; + for (const workerUrl of workerUrls) { + const key = getEdgeKey(workerUrl); + if (key) edgeWorkers.push({ workerUrl, key }); + } + const lock = await readLockfile(); const plugins: BackupPlugin[] = lock.plugins // Local-path plugins aren't reproducible on another machine — skip them @@ -219,6 +236,7 @@ export async function collectBackupPayload( onlyProfile === undefined ? current ?? undefined : onlyProfile, profiles: backupProfiles, plugins, + edgeWorkers, }; } @@ -436,6 +454,14 @@ function assertPayload(v: unknown): BackupPayload { if (!Array.isArray(o.plugins)) { throw new Error("payload.plugins must be an array"); } + if (o.edgeWorkers !== undefined) { + if (!Array.isArray(o.edgeWorkers) || o.edgeWorkers.some((worker) => + typeof worker !== "object" || worker === null || + typeof worker.workerUrl !== "string" || typeof worker.key !== "string" + )) { + throw new Error("payload.edgeWorkers must contain workerUrl and key strings"); + } + } return o as unknown as BackupPayload; } diff --git a/cli/tests/backup-roundtrip.test.ts b/cli/tests/backup-roundtrip.test.ts index e2ad02d..63b3900 100644 --- a/cli/tests/backup-roundtrip.test.ts +++ b/cli/tests/backup-roundtrip.test.ts @@ -1,4 +1,4 @@ -import { mkdtemp, rm } from "node:fs/promises"; +import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { @@ -146,6 +146,76 @@ async function wipeState(): Promise { // --------------------------------------------------------------------------- describe("mantis backup → mantis restore round-trip", () => { + it("round-trips independent edge keys without any server profile and preserves existing keys", async () => { + const outPath = join(tmpHome, "edge-only.json"); + process.env.MANTIS_BACKUP_TEST_PASS = "edge-only-passphrase"; + const { setEdgeKey, getEdgeKey } = await import("../src/lib/edge-key.js"); + const worker = "https://standalone-edge.workers.dev"; + setEdgeKey(worker, "original-edge-key"); + const { backupCmd, restoreCmd } = await import("../src/commands/backup.js"); + await backupCmd({ out: outPath, passphraseEnv: "MANTIS_BACKUP_TEST_PASS" }); + await wipeState(); + await restoreCmd(outPath, { passphraseEnv: "MANTIS_BACKUP_TEST_PASS" }); + expect(getEdgeKey(worker)).toBe("original-edge-key"); + const config = await import("../src/lib/config.js"); + expect(await config.readConfig()).toBeNull(); + setEdgeKey(worker, "newer-local-key"); + await restoreCmd(outPath, { passphraseEnv: "MANTIS_BACKUP_TEST_PASS" }); + expect(getEdgeKey(worker)).toBe("newer-local-key"); + await restoreCmd(outPath, { passphraseEnv: "MANTIS_BACKUP_TEST_PASS", overwrite: true }); + expect(getEdgeKey(worker)).toBe("original-edge-key"); + delete process.env.MANTIS_BACKUP_TEST_PASS; + }, 15_000); + + it("scopes --only to the selected profile's linked worker", async () => { + await populateState(); + const { setEdgeKey } = await import("../src/lib/edge-key.js"); + setEdgeKey("https://independent.workers.dev", "independent-key"); + const { backupCmd } = await import("../src/commands/backup.js"); + const { openBundle } = await import("../src/lib/backup.js"); + const outPath = join(tmpHome, "scoped.json"); + process.env.MANTIS_BACKUP_TEST_PASS = "scope-passphrase"; + await backupCmd({ out: outPath, profile: "primary", passphraseEnv: "MANTIS_BACKUP_TEST_PASS" }); + const payload = await openBundle(JSON.parse(await readFile(outPath, "utf8")), "scope-passphrase"); + expect(payload.edgeWorkers?.map((worker) => worker.workerUrl)).toEqual(["https://primary-edge.workers.dev"]); + expect(vi.mocked(process.stderr.write).mock.calls.join(" ")).toContain("omit --only to include independent edge workers"); + delete process.env.MANTIS_BACKUP_TEST_PASS; + }); + + it("restores old v1 bundles with edge keys embedded only in profiles", async () => { + await populateState(); + const { collectBackupPayload, sealBundle } = await import("../src/lib/backup.js"); + const payload = await collectBackupPayload(undefined); + delete payload.edgeWorkers; + const outPath = join(tmpHome, "legacy.json"); + await writeFile(outPath, JSON.stringify(await sealBundle(payload, "legacy-passphrase"))); + await wipeState(); + process.env.MANTIS_BACKUP_TEST_PASS = "legacy-passphrase"; + const { restoreCmd } = await import("../src/commands/backup.js"); + await restoreCmd(outPath, { passphraseEnv: "MANTIS_BACKUP_TEST_PASS" }); + const { getEdgeKey } = await import("../src/lib/edge-key.js"); + expect(getEdgeKey("https://primary-edge.workers.dev")).toBe("MGYWRl0WT3RcVuQrMQuv4Ph9DcZakhfwHcZk0lszKnE"); + delete process.env.MANTIS_BACKUP_TEST_PASS; + }); + + it.each([false, true])("protects an existing independent worker key when restoring legacy profiles (overwrite=%s)", async (overwrite) => { + await populateState(); + const { collectBackupPayload, sealBundle } = await import("../src/lib/backup.js"); + const payload = await collectBackupPayload(undefined); + delete payload.edgeWorkers; + const outPath = join(tmpHome, "legacy-existing-worker.json"); + await writeFile(outPath, JSON.stringify(await sealBundle(payload, "legacy-passphrase"))); + await wipeState(); + const { setEdgeKey, getEdgeKey } = await import("../src/lib/edge-key.js"); + const worker = "https://primary-edge.workers.dev"; + setEdgeKey(worker, "newer-independent-key"); + process.env.MANTIS_BACKUP_TEST_PASS = "legacy-passphrase"; + const { restoreCmd } = await import("../src/commands/backup.js"); + await restoreCmd(outPath, { passphraseEnv: "MANTIS_BACKUP_TEST_PASS", overwrite }); + expect(getEdgeKey(worker)).toBe(overwrite ? "MGYWRl0WT3RcVuQrMQuv4Ph9DcZakhfwHcZk0lszKnE" : "newer-independent-key"); + delete process.env.MANTIS_BACKUP_TEST_PASS; + }); + it("restores every profile and its keychain entries on a clean machine", async () => { const outPath = join(tmpHome, "bundle.json"); process.env.MANTIS_BACKUP_TEST_PASS = "diceware-style-test-passphrase"; diff --git a/cli/tests/device-completion.test.ts b/cli/tests/device-completion.test.ts new file mode 100644 index 0000000..b0b9eaf --- /dev/null +++ b/cli/tests/device-completion.test.ts @@ -0,0 +1,40 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; + +const mocks = vi.hoisted(() => ({ + install: vi.fn(async () => false), + client: { + deviceProfiles: vi.fn(async () => ({ profiles: [{ os: "linux", defaults: ["login"], vectors: [{ slug: "login", label: "shell login", response_kind: "empty", dedupe_window_seconds: 0 }] }] })), + createKey: vi.fn(async () => ({ id: "key-one", url: "https://mantis.example/c/one" })), + deviceBundleFiles: vi.fn(async () => ({ files: {}, installScript: "install.sh", uninstallScript: "uninstall.sh" })), + }, +})); +vi.mock("../src/lib/runner.js", () => ({ withClient: async (_opts: unknown, action: (client: unknown) => unknown) => action(mocks.client) })); +vi.mock("../src/lib/device-install.js", () => ({ + applyBundleLocally: mocks.install, assertBundleInstallableHere: vi.fn(), +})); +import { deviceNewCmd } from "../src/commands/device.js"; +import { edgeDeviceCmd } from "../src/commands/edge-device.js"; + +afterEach(() => vi.restoreAllMocks()); +describe("device completion feedback", () => { + it.each(["server", "edge"])("does not call a declined installation armed (%s)", async (surface) => { + mocks.install.mockResolvedValue(false); + const errors: string[] = []; + vi.spyOn(process.stderr, "write").mockImplementation((chunk) => { errors.push(String(chunk)); return true; }); + const opts = { os: "linux", name: "audit-host", vectors: "login", install: true }; + if (surface === "server") await deviceNewCmd(opts); + else await edgeDeviceCmd({ ...opts, worker: "https://edge.example.com", webhook: "https://hooks.example.com", key: Buffer.alloc(32, 7).toString("base64url") }); + expect(errors.join("")).not.toContain(" armed"); + expect(errors.join("")).toContain("Installation canceled. Nothing installed"); + expect(errors.join("")).toContain("minted"); + }); + + it("calls a completed installation armed", async () => { + mocks.install.mockResolvedValue(true); + const errors: string[] = []; + vi.spyOn(process.stderr, "write").mockImplementation((chunk) => { errors.push(String(chunk)); return true; }); + await deviceNewCmd({ os: "linux", name: "audit-host", vectors: "login", install: true }); + expect(errors.join("")).toContain(" armed"); + expect(errors.join("")).not.toContain("Nothing installed"); + }); +}); diff --git a/cli/tests/edge-mint.test.ts b/cli/tests/edge-mint.test.ts new file mode 100644 index 0000000..d85971f --- /dev/null +++ b/cli/tests/edge-mint.test.ts @@ -0,0 +1,31 @@ +import { mkdtemp, readFile, rm } from "node:fs/promises"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { mintCmd } from "../src/commands/edge.js"; +import { setJsonMode } from "../src/lib/out.js"; + +afterEach(() => { setJsonMode(false); vi.restoreAllMocks(); }); +describe("chained edge mint JSON output", () => { + it.each([false, true])("emits one result with recoverable installer data (file=%s)", async (writeFile) => { + const dir = await mkdtemp(join(tmpdir(), "mantis-edge-mint-")); + const output: string[] = []; + vi.spyOn(process.stdout, "write").mockImplementation((chunk) => { output.push(String(chunk)); return true; }); + setJsonMode(true); + try { + await mintCmd({ + worker: "https://edge.example.com", webhook: "https://hooks.example.com/notify", + key: Buffer.alloc(32, 7).toString("base64url"), install: "shell", + ...(writeFile ? { out: join(dir, "mantis.sh") } : {}), + }); + const result = JSON.parse(output.join("")); + expect(output.join("").trim().split("\n")).toHaveLength(1); + expect(result.url).toContain("/c/"); + expect(result.installer.type).toBe("shell"); + const content = writeFile ? await readFile(result.installer.written_to, "utf8") : result.installer.content; + expect(content).toContain(result.url); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); +}); diff --git a/iot-helper/README.md b/iot-helper/README.md index ee18b50..cf4e820 100644 --- a/iot-helper/README.md +++ b/iot-helper/README.md @@ -44,6 +44,7 @@ node bin/mantis-iot-helper.js --config mantis-iot.json { "interval_seconds": 30, "cooldown_seconds": 900, + "delivery_timeout_seconds": 10, "interface": "br0", "devices": [ { @@ -77,7 +78,18 @@ that cross midnight are supported, for example `{"start":"23:00","end":"06:00"}` - Network detection is best-effort. ARP/neighbor tables only include devices recently seen by the watcher host; set `"ping": true` for devices that answer - ICMP and need active probing. + ICMP and need active probing. Failed/incomplete neighbor entries do not count + as presence. `interface` restricts matching and ping probes to that interface; + a device's own `interface` overrides the global setting. - Login detection requires a log source. Many cameras/routers can send syslog to a local host; point `log_watchers[].path` at that received log. - Run with enough permissions to read neighbor tables and log files. +- Delivery has a bounded timeout (10 seconds by default). Accepted 2xx/3xx + trigger responses start the cooldown; failures retry on the next polling + tick. Redirect responses are accepted without following the target. +- Cooldowns and log offsets are separate for each destination, so repeated + device/watcher entries can notify multiple servers independently. +- Log watching starts at the end of the file on startup. Failed log events are + retained in memory and retried even if that file rotates or disappears. This + state is not durable across helper restarts; stopping/restarting the helper + drops pending events and starts at the current end of each file again. diff --git a/iot-helper/bin/mantis-iot-helper.js b/iot-helper/bin/mantis-iot-helper.js index 37b99d4..feebd71 100644 --- a/iot-helper/bin/mantis-iot-helper.js +++ b/iot-helper/bin/mantis-iot-helper.js @@ -3,41 +3,53 @@ import { readFile, stat } from "node:fs/promises"; import { createReadStream } from "node:fs"; import { spawn } from "node:child_process"; import { platform } from "node:os"; +import { resolve } from "node:path"; +import { pathToFileURL } from "node:url"; const DEFAULT_INTERVAL_SECONDS = 30; const DEFAULT_COOLDOWN_SECONDS = 900; +const DEFAULT_DELIVERY_TIMEOUT_SECONDS = 10; -const args = parseArgs(process.argv.slice(2)); -if (args.help) { - printHelp(); - process.exit(0); -} +export async function main(argv = process.argv.slice(2)) { + const args = parseArgs(argv); + if (args.help) { + printHelp(); + return; + } + + const configPath = String(args.config ?? process.env.MANTIS_IOT_CONFIG ?? "mantis-iot.json"); + const config = await loadConfig(configPath); + const intervalMs = seconds(config.interval_seconds, DEFAULT_INTERVAL_SECONDS) * 1000; + const cooldownMs = seconds(config.cooldown_seconds, DEFAULT_COOLDOWN_SECONDS) * 1000; + const deliveryTimeoutMs = seconds(config.delivery_timeout_seconds, DEFAULT_DELIVERY_TIMEOUT_SECONDS) * 1000; + const dryRun = Boolean(args["dry-run"] ?? config.dry_run); + const once = Boolean(args.once); -const configPath = String(args.config ?? process.env.MANTIS_IOT_CONFIG ?? "mantis-iot.json"); -const config = await loadConfig(configPath); -const intervalMs = seconds(config.interval_seconds, DEFAULT_INTERVAL_SECONDS) * 1000; -const cooldownMs = seconds(config.cooldown_seconds, DEFAULT_COOLDOWN_SECONDS) * 1000; -const dryRun = Boolean(args["dry-run"] ?? config.dry_run); -const once = Boolean(args.once); + const state = createState(); -const state = { - firedAt: new Map(), - logOffsets: new Map(), -}; + console.error(`mantis-iot-helper watching ${config.devices?.length ?? 0} devices, ${config.log_watchers?.length ?? 0} logs`); -console.error(`mantis-iot-helper watching ${config.devices?.length ?? 0} devices, ${config.log_watchers?.length ?? 0} logs`); + do { + await tick(config, state, { cooldownMs, dryRun, deliveryTimeoutMs }); + if (once) break; + await sleep(intervalMs); + } while (true); +} + +export function createState() { + return { firedAt: new Map(), logOffsets: new Map(), pendingLogs: new Map() }; +} -do { - await tick(config, state, { cooldownMs, dryRun }); - if (once) break; - await sleep(intervalMs); -} while (true); +if (process.argv[1] && import.meta.url === pathToFileURL(resolve(process.argv[1])).href) { + await main(); +} -async function tick(config, state, opts) { - const neighbors = await getNeighbors(); +export async function tick(config, state, opts) { + const neighbors = opts.neighbors ?? await getNeighbors(); const now = new Date(); for (const device of config.devices ?? []) { - const present = await isDevicePresent(device, neighbors); + const networkInterface = device.interface ?? config.interface; + const present = await isDevicePresent({ ...device, interface: networkInterface }, neighbors); if (!present) continue; if (isAllowedNow(device.allowed, now)) continue; await fireWithCooldown({ @@ -45,13 +57,14 @@ async function tick(config, state, opts) { state, cooldownMs: opts.cooldownMs, dryRun: opts.dryRun, + deliveryTimeoutMs: opts.deliveryTimeoutMs, url: device.mantis_url, event: "unexpected-online", source: "iot-network", device: device.name, mac: normalizeMac(device.mac), ip: device.ip, - networkInterface: device.interface ?? config.interface, + networkInterface, payload: { device, present, at: now.toISOString() }, }); } @@ -72,69 +85,78 @@ async function getNeighbors() { return parseArp(arp); } -async function isDevicePresent(device, neighbors) { +export async function isDevicePresent(device, neighbors) { const mac = normalizeMac(device.mac); - if (mac && neighbors.byMac.has(mac)) return true; - if (device.ip && neighbors.byIp.has(device.ip)) return true; + const entries = neighbors.entries ?? [...neighbors.byMac.values(), ...neighbors.byIp.values()]; + if (entries.some((entry) => + (!device.interface || entry.dev === device.interface) && + ((mac && entry.mac === mac) || (device.ip && entry.ip === device.ip)) + )) return true; if (device.ping && device.ip) { - return ping(device.ip); + return ping(device.ip, device.interface); } return false; } -function parseIpNeighJson(raw) { - const byMac = new Map(); - const byIp = new Map(); +export function parseIpNeighJson(raw) { + const neighbors = emptyNeighbors(); try { const rows = JSON.parse(raw); for (const row of rows) { const ip = row.dst; const mac = normalizeMac(row.lladdr); - if (!ip) continue; - byIp.set(ip, { ip, mac, dev: row.dev, state: row.state }); - if (mac) byMac.set(mac, { ip, mac, dev: row.dev, state: row.state }); + if (!ip || invalidNeighborState(row.state)) continue; + addNeighbor(neighbors, { ip, mac, dev: row.dev, state: row.state }); } } catch { return emptyNeighbors(); } - return { byMac, byIp }; + return neighbors; } -function parseIpNeighText(raw) { - const byMac = new Map(); - const byIp = new Map(); +export function parseIpNeighText(raw) { + const neighbors = emptyNeighbors(); for (const line of raw.split(/\r?\n/)) { const ip = line.match(/^(\S+)/)?.[1]; const dev = line.match(/\bdev\s+(\S+)/)?.[1]; const mac = normalizeMac(line.match(/\blladdr\s+(\S+)/)?.[1]); - if (!ip) continue; - byIp.set(ip, { ip, mac, dev }); - if (mac) byMac.set(mac, { ip, mac, dev }); + if (!ip || invalidNeighborState(line)) continue; + addNeighbor(neighbors, { ip, mac, dev }); } - return { byMac, byIp }; + return neighbors; } -function parseArp(raw) { - const byMac = new Map(); - const byIp = new Map(); +export function parseArp(raw) { + const neighbors = emptyNeighbors(); for (const line of raw.split(/\r?\n/)) { const ip = line.match(/\(([^)]+)\)/)?.[1]; const mac = normalizeMac(line.match(/\bat\s+([0-9a-f:.-]+)/i)?.[1]); - if (!ip) continue; - byIp.set(ip, { ip, mac }); - if (mac) byMac.set(mac, { ip, mac }); + const dev = line.match(/\bon\s+(\S+)/)?.[1]; + if (!ip || invalidNeighborState(line)) continue; + addNeighbor(neighbors, { ip, mac, dev }); } - return { byMac, byIp }; + return neighbors; } function emptyNeighbors() { - return { byMac: new Map(), byIp: new Map() }; + return { byMac: new Map(), byIp: new Map(), entries: [] }; +} + +function invalidNeighborState(value) { + return /\b(FAILED|INCOMPLETE)\b/i.test(Array.isArray(value) ? value.join(" ") : String(value ?? "")); } -async function ping(ip) { +function addNeighbor(neighbors, entry) { + neighbors.entries.push(entry); + neighbors.byIp.set(entry.ip, entry); + if (entry.mac) neighbors.byMac.set(entry.mac, entry); +} + +async function ping(ip, networkInterface) { const args = platform() === "darwin" ? ["-c", "1", "-W", "1000", ip] : ["-c", "1", "-W", "1", ip]; + if (networkInterface) args.unshift(platform() === "darwin" ? "-b" : "-I", networkInterface); try { await run("ping", args, { timeoutMs: 2500 }); return true; @@ -143,8 +165,12 @@ async function ping(ip) { } } -async function scanLogWatcher(watcher, state, opts) { +export async function scanLogWatcher(watcher, state, opts) { if (!watcher.path || !watcher.pattern || !watcher.mantis_url) return; + // Each destination consumes its own stream, even when it watches the same + // file/name. Keep failed events outside the file so rotation cannot lose them. + const key = JSON.stringify([watcher.name, watcher.path, watcher.pattern, watcher.event, watcher.mantis_url]); + if (!await flushPendingLogs(key, state)) return; let info; try { info = await stat(watcher.path); @@ -152,34 +178,63 @@ async function scanLogWatcher(watcher, state, opts) { return; } - const previous = state.logOffsets.get(watcher.path) ?? info.size; - const start = previous > info.size ? 0 : previous; - state.logOffsets.set(watcher.path, info.size); - if (start === info.size) return; + const previous = state.logOffsets.get(key); + const fileId = `${info.dev}:${info.ino}`; + const continuing = previous && previous.fileId === fileId && previous.offset <= info.size; + const start = !previous ? info.size : continuing ? previous.offset : 0; + if (start === info.size) { + state.logOffsets.set(key, { offset: info.size, fileId, partial: continuing ? previous.partial ?? "" : "" }); + return; + } - const text = await readRange(watcher.path, start, info.size); + let text; + try { + text = await readRange(watcher.path, start, info.size); + } catch (err) { + console.error("log read failed", watcher.name, err instanceof Error ? err.message : String(err)); + return; + } + const lines = `${continuing ? previous.partial ?? "" : ""}${text}`.split(/\r?\n/); + const partial = lines.pop() ?? ""; + state.logOffsets.set(key, { offset: info.size, fileId, partial }); const re = new RegExp(watcher.pattern, "i"); - for (const line of text.split(/\r?\n/)) { + const events = []; + for (const line of lines) { if (!re.test(line)) continue; - await fireWithCooldown({ - key: `log:${watcher.name}`, + const event = { + key: `log:${key}`, state, cooldownMs: opts.cooldownMs, dryRun: opts.dryRun, + deliveryTimeoutMs: opts.deliveryTimeoutMs, url: watcher.mantis_url, event: watcher.event ?? "device-log", source: "iot-log", device: watcher.device ?? watcher.name, payload: { watcher: watcher.name, line, at: new Date().toISOString() }, - }); + }; + events.push(event); } + if (events.length) state.pendingLogs.set(key, events); + await flushPendingLogs(key, state); } -async function fireWithCooldown({ +async function flushPendingLogs(key, state) { + const events = state.pendingLogs.get(key) ?? []; + while (events.length) { + if (!await fireWithCooldown(events[0])) return false; + events.shift(); + } + state.pendingLogs.delete(key); + return true; +} + +export async function fireWithCooldown({ key, state, cooldownMs, dryRun, + deliveryTimeoutMs = DEFAULT_DELIVERY_TIMEOUT_SECONDS * 1000, url, event, source, @@ -189,11 +244,11 @@ async function fireWithCooldown({ networkInterface, payload, }) { - if (!url) return; + if (!url) return false; + const targetKey = JSON.stringify([key, url]); const now = Date.now(); - const last = state.firedAt.get(key) ?? 0; - if (now - last < cooldownMs) return; - state.firedAt.set(key, now); + const last = state.firedAt.get(targetKey); + if (last !== undefined && now - last < cooldownMs) return true; const headers = { "Content-Type": "application/json", @@ -207,7 +262,8 @@ async function fireWithCooldown({ if (dryRun) { console.error("dry-run fire", event, device, url); - return; + state.firedAt.set(targetKey, now); + return true; } try { @@ -215,11 +271,19 @@ async function fireWithCooldown({ method: "POST", headers, body: JSON.stringify(payload ?? {}), + redirect: "manual", + signal: AbortSignal.timeout(Math.max(1, Math.floor(deliveryTimeoutMs))), }); - await res.arrayBuffer().catch(() => undefined); + // The trigger has accepted the event once its response headers arrive. + // Don't wait for an arbitrary response body or follow redirect targets. + void res.body?.cancel().catch(() => {}); + if (res.status < 200 || res.status >= 400) throw new Error(`HTTP ${res.status}`); + state.firedAt.set(targetKey, Date.now()); console.error("fired", event, device ?? "-", res.status); + return true; } catch (err) { console.error("fire failed", event, device ?? "-", err instanceof Error ? err.message : String(err)); + return false; } } diff --git a/iot-helper/config.example.json b/iot-helper/config.example.json index e7fb4ba..584175e 100644 --- a/iot-helper/config.example.json +++ b/iot-helper/config.example.json @@ -1,6 +1,7 @@ { "interval_seconds": 30, "cooldown_seconds": 900, + "delivery_timeout_seconds": 10, "interface": "br0", "devices": [ { diff --git a/iot-helper/package.json b/iot-helper/package.json index 1f6d808..bec85aa 100644 --- a/iot-helper/package.json +++ b/iot-helper/package.json @@ -11,6 +11,7 @@ }, "scripts": { "check": "node --check bin/mantis-iot-helper.js", + "test": "node --test tests/*.test.js", "start": "node bin/mantis-iot-helper.js --config config.example.json" } } diff --git a/iot-helper/tests/recovery.test.js b/iot-helper/tests/recovery.test.js new file mode 100644 index 0000000..6353891 --- /dev/null +++ b/iot-helper/tests/recovery.test.js @@ -0,0 +1,182 @@ +import { afterEach, beforeEach, describe, it } from "node:test"; +import assert from "node:assert/strict"; +import { appendFile, mkdtemp, rename, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { createServer } from "node:http"; +import { + createState, fireWithCooldown, isDevicePresent, parseArp, + parseIpNeighJson, parseIpNeighText, scanLogWatcher, tick, +} from "../bin/mantis-iot-helper.js"; + +const originalFetch = globalThis.fetch; +const originalError = console.error; +const opts = { cooldownMs: 900_000, dryRun: false, deliveryTimeoutMs: 50 }; +let dir; +beforeEach(async () => { + dir = await mkdtemp(join(tmpdir(), "mantis-iot-recovery-")); + console.error = () => {}; +}); +afterEach(async () => { + globalThis.fetch = originalFetch; + console.error = originalError; + await rm(dir, { recursive: true, force: true }); +}); + +function event(state, url = "https://mantis.example/c/token") { + return { key: "device:camera", state, url, event: "unexpected-online", source: "iot-network", device: "camera", ...opts }; +} + +describe("delivery recovery", () => { + it("does not consume cooldown on failed HTTP or network delivery", async () => { + const state = createState(); + const results = [new Response(null, { status: 503 }), new Error("offline"), new Response(null, { status: 204 })]; + let calls = 0; + globalThis.fetch = async () => { + calls++; + const next = results.shift(); + if (next instanceof Error) throw next; + return next; + }; + assert.equal(await fireWithCooldown(event(state)), false); + assert.equal(state.firedAt.size, 0); + assert.equal(await fireWithCooldown(event(state)), false); + assert.equal(state.firedAt.size, 0); + assert.equal(await fireWithCooldown(event(state)), true); + assert.equal(await fireWithCooldown(event(state)), true); + assert.equal(calls, 3); + }); + + it("bounds a hanging request and permits the next attempt", async () => { + const server = createServer(() => {}); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const url = `http://127.0.0.1:${server.address().port}/c/token`; + const state = createState(); + const started = Date.now(); + try { + assert.equal(await fireWithCooldown(event(state, url)), false); + assert.ok(Date.now() - started < 1500, "delivery should finish promptly after its deadline"); + assert.equal(state.firedAt.size, 0); + globalThis.fetch = async () => new Response(null, { status: 204 }); + assert.equal(await fireWithCooldown(event(state, url)), true); + } finally { + server.closeAllConnections(); + await new Promise((resolve) => server.close(resolve)); + } + }); + + it("accepts a redirect trigger without following its destination", async () => { + const state = createState(); + globalThis.fetch = async (_url, init) => { + assert.equal(init.redirect, "manual"); + return new Response(null, { status: 302, headers: { location: "https://unreachable.example" } }); + }; + assert.equal(await fireWithCooldown(event(state)), true); + assert.equal(state.firedAt.size, 1); + }); + + it("does not wait for a streaming response body after trigger acceptance", async () => { + const server = createServer((_req, res) => { res.writeHead(200); res.write("accepted"); }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const url = `http://127.0.0.1:${server.address().port}/c/token`; + try { + const started = Date.now(); + assert.equal(await fireWithCooldown({ ...event(createState(), url), deliveryTimeoutMs: 500 }), true); + assert.ok(Date.now() - started < 1500); + } finally { + server.closeAllConnections(); + await new Promise((resolve) => server.close(resolve)); + } + }); + + it("fires the same named device independently to multiple targets", async () => { + const urls = []; + globalThis.fetch = async (url) => { urls.push(url); return new Response(null, { status: 204 }); }; + const config = { devices: ["one", "two"].map((target) => ({ + name: "camera", ip: "192.168.1.50", mantis_url: `https://${target}.example/c/token`, + allowed: [{ start: "00:00", end: "00:00", days: [] }], + })) }; + const neighbors = parseIpNeighText("192.168.1.50 dev br0 lladdr aa:bb:cc:dd:ee:ff REACHABLE"); + const state = createState(); + await tick(config, state, { ...opts, neighbors }); + await tick(config, state, { ...opts, neighbors }); + assert.deepEqual(urls, ["https://one.example/c/token", "https://two.example/c/token"]); + }); +}); + +describe("log delivery continuity", () => { + it("retries a failed log event after the source file rotates and disappears", async () => { + const path = join(dir, "camera.log"); + await writeFile(path, "old log\n"); + const watcher = { name: "camera-login", path, pattern: "auth success", mantis_url: "https://mantis.example/c/token" }; + const state = createState(); + const lines = []; + let succeed = false; + globalThis.fetch = async (_url, init) => { + lines.push(JSON.parse(init.body).line); + return new Response(null, { status: succeed ? 204 : 503 }); + }; + await scanLogWatcher(watcher, state, opts); + await appendFile(path, "auth success: admin\n"); + await scanLogWatcher(watcher, state, opts); + assert.equal(state.pendingLogs.size, 1); + assert.equal(state.firedAt.size, 0); + await rename(path, `${path}.1`); + await rm(`${path}.1`); + succeed = true; + await scanLogWatcher(watcher, state, opts); + assert.deepEqual(lines, ["auth success: admin", "auth success: admin"]); + assert.equal(state.pendingLogs.size, 0); + assert.equal(state.firedAt.size, 1); + }); + + it("gives watchers on the same file separate offsets and cooldowns per destination", async () => { + const path = join(dir, "shared.log"); + await writeFile(path, ""); + const watchers = ["one", "two"].map((target) => ({ name: "login", path, pattern: "auth success", mantis_url: `https://${target}.example/c/token` })); + const state = createState(); + const urls = []; + globalThis.fetch = async (url) => { urls.push(url); return new Response(null, { status: 204 }); }; + for (const watcher of watchers) await scanLogWatcher(watcher, state, opts); + await appendFile(path, "auth success\n"); + for (const watcher of watchers) await scanLogWatcher(watcher, state, opts); + assert.deepEqual(urls, ["https://one.example/c/token", "https://two.example/c/token"]); + }); + + it("preserves a log line written across polling ticks", async () => { + const path = join(dir, "partial.log"); + await writeFile(path, ""); + const watcher = { name: "login", path, pattern: "auth success", mantis_url: "https://mantis.example/c/token" }; + const state = createState(); + const lines = []; + globalThis.fetch = async (_url, init) => { lines.push(JSON.parse(init.body).line); return new Response(null, { status: 204 }); }; + await scanLogWatcher(watcher, state, opts); + await appendFile(path, "auth "); + await scanLogWatcher(watcher, state, opts); + await appendFile(path, "success\n"); + await scanLogWatcher(watcher, state, opts); + assert.deepEqual(lines, ["auth success"]); + }); +}); + +describe("neighbor presence", () => { + it("rejects FAILED and INCOMPLETE rows in all supported table formats", async () => { + const device = { ip: "192.168.1.50", mac: "aa:bb:cc:dd:ee:ff" }; + for (const neighbors of [ + parseIpNeighJson(JSON.stringify([{dst: device.ip, lladdr: device.mac, dev: "br0", state: ["FAILED"]}])), + parseIpNeighText(`${device.ip} dev br0 lladdr ${device.mac} INCOMPLETE`), + parseArp(`? (${device.ip}) at (incomplete) on en0 ifscope [ethernet]`), + ]) assert.equal(await isDevicePresent(device, neighbors), false); + }); + + it("respects the configured interface while retaining valid stale neighbors", async () => { + const neighbors = parseIpNeighJson(JSON.stringify([ + {dst: "192.168.1.50", lladdr: "aa:bb:cc:dd:ee:ff", dev: "br0", state: ["STALE"]}, + {dst: "192.168.1.50", lladdr: "aa:bb:cc:dd:ee:ff", dev: "eth0", state: ["REACHABLE"]}, + ])); + assert.equal(await isDevicePresent({ ip: "192.168.1.50", interface: "br0" }, neighbors), true); + assert.equal(await isDevicePresent({ ip: "192.168.1.50", interface: "wlan0" }, neighbors), false); + const arp = parseArp("? (192.168.1.50) at aa:bb:cc:dd:ee:ff on en0 ifscope [ethernet]"); + assert.equal(await isDevicePresent({ ip: "192.168.1.50", interface: "en1" }, arp), false); + }); +}); diff --git a/package.json b/package.json index b16b869..963f6a2 100644 --- a/package.json +++ b/package.json @@ -18,7 +18,7 @@ "check": "pnpm run build:core && pnpm run typecheck && pnpm --filter @mantis/cli run typecheck && pnpm --filter @mantis/edge run typecheck && pnpm --filter @mantis/iot-helper run check", "build:cli-bin": "bash cli/scripts/build-bin.sh", "test": "vitest run", - "test:all": "pnpm test && pnpm --filter @mantis/cli test && pnpm --filter @mantis/edge test", + "test:all": "pnpm test && pnpm --filter @mantis/cli test && pnpm --filter @mantis/edge test && pnpm --filter @mantis/iot-helper test", "test:watch": "vitest", "test:integration": "vitest run --config vitest.integration.config.ts", "test:integration:db": "bash scripts/test-integration.sh",