From 3b9cebafa6ada33a91186ebd97e0083451bd8772 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 24 Jul 2026 17:30:57 +0000 Subject: [PATCH 1/3] fix: SOLID hardening for service logs + batch uninstall - DRY: share resolveDisplayLevel / level+source tag plain text in sfmc logs - LSP: resolveUninstallTargets aborts when any protected pack is selected - DIP: single project-root for db-server and qq-bridge (env/config/log) - LSP: FileSink contract matches sync append (drop unused flags; clarify close) Co-authored-by: Shiroha --- bds-tools/src/log.ts | 7 +- db-server/src/env.ts | 9 +- db-server/src/lib/log.ts | 12 +-- db-server/src/project-root.ts | 12 +++ modules/sdk/@sfmc-sdk/src/logs/logger.ts | 4 +- modules/sdk/@sfmc-sdk/src/logs/sink.ts | 9 +- qq-bridge/src/config.ts | 11 +-- qq-bridge/src/log.ts | 11 +-- qq-bridge/src/project-root.ts | 12 +++ sfmc/package.json | 2 +- sfmc/src/logs.ts | 110 +++++++++++------------ sfmc/src/world-packs.ts | 56 +++++++----- sfmc/world-packs-uninstall.test.mjs | 52 ++++++++++- sfmc/wrap-log-line.test.mjs | 22 +++++ 14 files changed, 205 insertions(+), 124 deletions(-) create mode 100644 db-server/src/project-root.ts create mode 100644 qq-bridge/src/project-root.ts diff --git a/bds-tools/src/log.ts b/bds-tools/src/log.ts index e08c644..d06366c 100644 --- a/bds-tools/src/log.ts +++ b/bds-tools/src/log.ts @@ -2,8 +2,7 @@ * log.ts — bds-tools 统一日志实例 * * stdout bare + 落盘 LOG_PATH (/.sfmc/logs/bds-update.log)。 - * check-update.ts 是独立入口,用 source = "updater" 单独创建 logger, - * 见该文件内的 createUpdaterLogger()。 + * check-update.ts 是独立入口,用 source = "updater" 单独创建 logger。 */ import { createNodeServiceLogger } from "@sfmc-bds/sdk/logs"; @@ -14,9 +13,7 @@ export const log = createNodeServiceLogger({ logPath: LOG_PATH, }); -/** 关闭文件流 (进程退出前调用,确保缓冲落盘) */ +/** 兼容钩子:FileSink 为 sync append,close 为空操作 */ export function closeLog(): void { log.close(); } - -process.on("exit", () => closeLog()); diff --git a/db-server/src/env.ts b/db-server/src/env.ts index f98622f..4962542 100644 --- a/db-server/src/env.ts +++ b/db-server/src/env.ts @@ -6,19 +6,14 @@ import { ensureCoreConfigs, loadEnsuredConfig, modulePath, - resolveRuntimeRoot, DEFAULT_DB_CONFIG, DEFAULT_QQ_CONFIG, } from "@sfmc-bds/sdk/node/config"; import { isAbsolute, join, resolve } from "node:path"; -import { fileURLToPath } from "node:url"; import { log } from "./lib/log.js"; +import { PROJECT_ROOT as RESOLVED_ROOT } from "./project-root.js"; -import { dirname } from "node:path"; - -const __filename = fileURLToPath(import.meta.url); -const __dirname = dirname(__filename); export interface EnvConfig { PROJECT_ROOT: string; @@ -41,7 +36,7 @@ export interface EnvConfig { } export function loadEnv(): EnvConfig { - const PROJECT_ROOT = resolveRuntimeRoot(resolve(__dirname, "..", "..")); + const PROJECT_ROOT = RESOLVED_ROOT; /* 启动时确认 db/qq/permissions 存在;不存在就写带默认值的骨架。 * 不依赖 wizard:wizard 只填字段,骨架由服务自己 ensure。 */ ensureCoreConfigs(PROJECT_ROOT, ["db_config", "qq_config", "permissions"]); diff --git a/db-server/src/lib/log.ts b/db-server/src/lib/log.ts index ad8f2b5..93d54b1 100644 --- a/db-server/src/lib/log.ts +++ b/db-server/src/lib/log.ts @@ -6,16 +6,10 @@ */ import { createNodeServiceLogger } from "@sfmc-bds/sdk/logs"; -import { logFile, resolveRuntimeRoot } from "@sfmc-bds/sdk/node/config"; -import { dirname, resolve } from "node:path"; -import { fileURLToPath } from "node:url"; - -/** src/lib → 仓根(与 env.ts 的 src 上溯差一层) */ -const ROOT = resolveRuntimeRoot(resolve(dirname(fileURLToPath(import.meta.url)), "..", "..", "..")); +import { logFile } from "@sfmc-bds/sdk/node/config"; +import { PROJECT_ROOT } from "../project-root.js"; export const log = createNodeServiceLogger({ source: "db", - logPath: logFile(ROOT, "db"), + logPath: logFile(PROJECT_ROOT, "db"), }); - -process.on("exit", () => log.close()); diff --git a/db-server/src/project-root.ts b/db-server/src/project-root.ts new file mode 100644 index 0000000..5e76f5c --- /dev/null +++ b/db-server/src/project-root.ts @@ -0,0 +1,12 @@ +/** + * project-root.ts — 仓根唯一解析点(DRY/DIP:env / log 勿各自上溯) + */ + +import { resolveRuntimeRoot } from "@sfmc-bds/sdk/node/config"; +import { dirname, resolve } from "node:path"; +import { fileURLToPath } from "node:url"; + +const HERE = dirname(fileURLToPath(import.meta.url)); + +/** SFMC_ROOT > 相对 db-server/src 上溯两级 */ +export const PROJECT_ROOT: string = resolveRuntimeRoot(resolve(HERE, "..", "..")); diff --git a/modules/sdk/@sfmc-sdk/src/logs/logger.ts b/modules/sdk/@sfmc-sdk/src/logs/logger.ts index 31faef4..55ca412 100644 --- a/modules/sdk/@sfmc-sdk/src/logs/logger.ts +++ b/modules/sdk/@sfmc-sdk/src/logs/logger.ts @@ -36,7 +36,9 @@ export interface Logger { /** Node 仓顶服务标准 logger:stdout(可 bare) + 文件落盘 */ export interface NodeServiceLogger extends Logger { - /** 关闭文件 sink */ + /** + * 关闭文件 sink(FileSink 为 sync append 时为空操作;保留统一调用面)。 + */ close(): void; readonly fileSink: FileSink; } diff --git a/modules/sdk/@sfmc-sdk/src/logs/sink.ts b/modules/sdk/@sfmc-sdk/src/logs/sink.ts index 92f5abf..6d46160 100644 --- a/modules/sdk/@sfmc-sdk/src/logs/sink.ts +++ b/modules/sdk/@sfmc-sdk/src/logs/sink.ts @@ -2,7 +2,7 @@ * sink.ts — 日志输出目标实现 * * StdoutSink: 输出到 stdout (可选颜色,可选 stderr 路由 error) - * FileSink: 追加写入文件 (纯文本,无 ANSI 码,单例 FD) + * FileSink: 同步追加写入文件 (纯文本,无 ANSI 码;进程退出无需 flush) */ import fs from "node:fs"; @@ -49,12 +49,13 @@ export function createStdoutSink(opts: StdoutSinkOptions = {}): Sink { export interface FileSinkOptions { /** 是否自动创建父目录 (默认 true) */ mkdir?: boolean; - /** 文件打开模式 (默认 "a" 追加) */ - flags?: string; } export interface FileSink extends Sink { - /** 关闭底层文件流,释放 FD */ + /** + * 兼容钩子:当前实现为 sync appendFile,无持有 FD,调用为空操作。 + * 保留以便 NodeServiceLogger / 进程 exit 统一调用而不破坏 LSP。 + */ close(): void; } diff --git a/qq-bridge/src/config.ts b/qq-bridge/src/config.ts index caf7d77..4c12005 100644 --- a/qq-bridge/src/config.ts +++ b/qq-bridge/src/config.ts @@ -11,19 +11,14 @@ import { configPath, DEFAULT_QQ_CONFIG, loadEnsuredConfig, - resolveRuntimeRoot, stripConfigMeta, } from "@sfmc-bds/sdk/node/config"; -import { dirname, resolve } from "node:path"; -import { fileURLToPath } from "node:url"; import { log } from "./log.js"; +import { PROJECT_ROOT } from "./project-root.js"; import type { QQBridgeConfig } from "./types.js"; -const __filename = fileURLToPath(import.meta.url); -const __dirname = dirname(__filename); - -/** 统一通过 SDK 解析项目根:env SFMC_ROOT > __dirname 上溯。 */ -export const ROOT_DIR: string = resolveRuntimeRoot(resolve(__dirname, "..", "..")); +/** 统一通过 SDK 解析项目根:env SFMC_ROOT > project-root 上溯。 */ +export const ROOT_DIR: string = PROJECT_ROOT; export const CFG_PATH: string = configPath(ROOT_DIR, "qq_config.json"); function applyDefaults(raw: Partial): QQBridgeConfig { diff --git a/qq-bridge/src/log.ts b/qq-bridge/src/log.ts index d502ac7..8f408b1 100644 --- a/qq-bridge/src/log.ts +++ b/qq-bridge/src/log.ts @@ -6,15 +6,10 @@ */ import { createNodeServiceLogger } from "@sfmc-bds/sdk/logs"; -import { logFile, resolveRuntimeRoot } from "@sfmc-bds/sdk/node/config"; -import { dirname, resolve } from "node:path"; -import { fileURLToPath } from "node:url"; - -const ROOT = resolveRuntimeRoot(resolve(dirname(fileURLToPath(import.meta.url)), "..", "..")); +import { logFile } from "@sfmc-bds/sdk/node/config"; +import { PROJECT_ROOT } from "./project-root.js"; export const log = createNodeServiceLogger({ source: "qq", - logPath: logFile(ROOT, "qq"), + logPath: logFile(PROJECT_ROOT, "qq"), }); - -process.on("exit", () => log.close()); diff --git a/qq-bridge/src/project-root.ts b/qq-bridge/src/project-root.ts new file mode 100644 index 0000000..65bd231 --- /dev/null +++ b/qq-bridge/src/project-root.ts @@ -0,0 +1,12 @@ +/** + * project-root.ts — 仓根唯一解析点(DRY/DIP:config / log 勿各自上溯) + */ + +import { resolveRuntimeRoot } from "@sfmc-bds/sdk/node/config"; +import { dirname, resolve } from "node:path"; +import { fileURLToPath } from "node:url"; + +const HERE = dirname(fileURLToPath(import.meta.url)); + +/** SFMC_ROOT > 相对 qq-bridge/src 上溯两级 */ +export const PROJECT_ROOT: string = resolveRuntimeRoot(resolve(HERE, "..", "..")); diff --git a/sfmc/package.json b/sfmc/package.json index baf1978..be2d3be 100644 --- a/sfmc/package.json +++ b/sfmc/package.json @@ -45,7 +45,7 @@ "start": "node ./dist/main.js", "build": "node ../scripts/esbuild-transpile.mjs --dts", "typecheck": "tsc7 --noEmit -p tsconfig.json", - "test": "npm run build && node --test pack-update-policy.test.mjs pack-update-config.test.mjs terminal-progress.test.mjs wrap-log-line.test.mjs", + "test": "npm run build && node --test pack-update-policy.test.mjs pack-update-config.test.mjs terminal-progress.test.mjs wrap-log-line.test.mjs world-packs-uninstall.test.mjs", "prepublishOnly": "npm run build" }, "dependencies": { diff --git a/sfmc/src/logs.ts b/sfmc/src/logs.ts index 9887d90..ce20689 100644 --- a/sfmc/src/logs.ts +++ b/sfmc/src/logs.ts @@ -114,10 +114,18 @@ export const SOURCE_META: SourceMeta[] = [ { value: "bds-tools", name: "BDSTools", paint: (s) => c.red(s) }, ]; +/** 源标签无色文本(formatSourceTag / logPrefixWidth 共用) */ +function sourceTagPlain(source: string): string { + const meta = SOURCE_META.find((m) => m.value === source); + if (meta) return `[${meta.name}]`; + return `[${source.padEnd(7).slice(0, 8)}]`; +} + export function formatSourceTag(source: string): string { const meta = SOURCE_META.find((m) => m.value === source); - if (meta) return meta.paint(`[${meta.name}]`); - return c.bold(`[${source.padEnd(7).slice(0, 8)}]`); + const plain = sourceTagPlain(source); + if (meta) return meta.paint(plain); + return c.bold(plain); } /** 简化BDS日志 */ @@ -151,44 +159,56 @@ function getLogLevel(line: string): string { return "UNKNOWN"; } -/** 格式化日志用于 REPL 展示 (用 theme.ts chalk 配色) */ -export function formatLog(l: UnifiedLog): string { - const ts = c.dim(l.time.toLocaleTimeString()); - let lvl = levelTag(l.level); - let txt = highlightLogLine(l.text); - const src = formatSourceTag(l.source); - if (l.source === "bds") { - const parsed = getLogLevel(l.text); - const mapped: LogLevel = - parsed === "WARNING" || parsed === "WARN" - ? "warn" - : parsed === "ERROR" || parsed === "FATAL" - ? "error" - : parsed === "DEBUG" || parsed === "TRACE" - ? "debug" - : "info"; - /* 去掉 BDS 自带时间戳前缀后再高亮正文 */ - txt = highlightLogLine(stripLogPrefix(l.text)); - lvl = levelTag(mapped); - } - return `${ts} ${src} ${lvl} ${txt}`; +/** 展示用级别:BDS 行从正文解析,其余用 entry.level(DRY:formatLog / logPrefixWidth 共用) */ +export function resolveDisplayLevel(l: UnifiedLog): LogLevel { + if (l.source !== "bds") return l.level; + const parsed = getLogLevel(l.text); + if (parsed === "WARNING" || parsed === "WARN") return "warn"; + if (parsed === "ERROR" || parsed === "FATAL") return "error"; + if (parsed === "DEBUG" || parsed === "TRACE") return "debug"; + return "info"; +} + +/** 无色级别标签文本(可见宽度权威源) */ +const LEVEL_TAG_TEXT: Record = { + error: "[ERR]", + warn: "[WRN]", + success: "[OK]", + debug: "[DBG]", + info: "[INF]", +}; + +function levelTagPlain(lvl: LogLevel): string { + return LEVEL_TAG_TEXT[lvl] ?? LEVEL_TAG_TEXT.info; } function levelTag(lvl: LogLevel): string { + const text = levelTagPlain(lvl); switch (lvl) { case "error": - return c.red("[ERR]"); + return c.red(text); case "warn": - return c.yellow("[WRN]"); + return c.yellow(text); case "success": - return c.green(c.bold("[OK]")); + return c.green(c.bold(text)); case "debug": - return c.dim("[DBG]"); + return c.dim(text); default: - return c.blue("[INF]"); + return c.blue(text); } } +/** 格式化日志用于 REPL 展示 (用 theme.ts chalk 配色) */ +export function formatLog(l: UnifiedLog): string { + const ts = c.dim(l.time.toLocaleTimeString()); + const level = resolveDisplayLevel(l); + const lvl = levelTag(level); + const src = formatSourceTag(l.source); + /* BDS:去掉自带时间戳前缀后再高亮正文 */ + const txt = highlightLogLine(l.source === "bds" ? stripLogPrefix(l.text) : l.text); + return `${ts} ${src} ${lvl} ${txt}`; +} + /* ================================================================== * 软换行: 超终端宽度时换行,后续行缩进对齐 * ================================================================== */ @@ -219,42 +239,14 @@ export function visibleWidth(s: string): number { return w; } -/** 无色级别标签,与 levelTag 可见宽度一致 */ -function levelTagPlain(lvl: LogLevel): string { - switch (lvl) { - case "error": - return "[ERR]"; - case "warn": - return "[WRN]"; - case "success": - return "[OK]"; - case "debug": - return "[DBG]"; - default: - return "[INF]"; - } -} - /** * 日志前缀可见宽度(时间 + 源 + 级别 + 尾空格),供悬挂缩进对齐正文起点。 - * 与 formatLog 拼接顺序保持一致。 + * 与 formatLog 拼接顺序保持一致(共用 resolveDisplayLevel / sourceTagPlain / levelTagPlain)。 */ export function logPrefixWidth(l: UnifiedLog): number { const ts = l.time.toLocaleTimeString(); - const meta = SOURCE_META.find((m) => m.value === l.source); - const src = meta ? `[${meta.name}]` : `[${l.source.padEnd(7).slice(0, 8)}]`; - let level: LogLevel = l.level; - if (l.source === "bds") { - const parsed = getLogLevel(l.text); - level = - parsed === "WARNING" || parsed === "WARN" - ? "warn" - : parsed === "ERROR" || parsed === "FATAL" - ? "error" - : parsed === "DEBUG" || parsed === "TRACE" - ? "debug" - : "info"; - } + const src = sourceTagPlain(l.source); + const level = resolveDisplayLevel(l); return visibleWidth(`${ts} ${src} ${levelTagPlain(level)} `); } diff --git a/sfmc/src/world-packs.ts b/sfmc/src/world-packs.ts index c0d4458..cfb2105 100644 --- a/sfmc/src/world-packs.ts +++ b/sfmc/src/world-packs.ts @@ -228,6 +228,35 @@ function formatUninstallPickLabel(pack: InstalledWorldPack): string { return `[${kind}] ${pack.name} ${pack.folderName} ${en.trim()} ${pack.uuid.slice(0, 8)}`; } +/** + * 按 CLI id 列表解析待卸载包(DIP:解析与 TTY/确认/执行分离,便于单测)。 + * - 任一 id 找不到 → not_found + * - 任一命中受保护包 → protected(与单 id 中止契约一致,LSP) + */ +export function resolveUninstallTargets( + packs: InstalledWorldPack[], + ids: string[] +): + | { status: "ok"; selected: InstalledWorldPack[] } + | { status: "not_found"; missing: string[] } + | { status: "protected"; folder: string } { + const missing: string[] = []; + const byUuid = new Map(); + for (const id of ids) { + const pack = findInstalledPackById(packs, id); + if (!pack) { + missing.push(id); + continue; + } + byUuid.set(pack.uuid.toLowerCase(), pack); + } + if (missing.length) return { status: "not_found", missing }; + const selected = [...byUuid.values()]; + const blocked = selected.find((p) => isProtectedSfmcPack(p)); + if (blocked) return { status: "protected", folder: blocked.folderName }; + return { status: "ok", selected }; +} + /** TTY 多选待卸载包;取消返回 null;无可卸项返回 [] */ async function pickPacksForUninstall(packs: InstalledWorldPack[]): Promise { const candidates = packs.filter((p) => !isProtectedSfmcPack(p)); @@ -816,29 +845,16 @@ export async function dispatchPacksCommand(sub: string | undefined, args: string if (picked === null) return c.dim(t("packs.uninstall.cancelled")); if (picked.length === 0) return c.yellow(t("packs.uninstall.none")); selected = picked; - } else if (ids.length === 1) { - /* 单 id:受保护主包直接中止(与 #72 LSP 一致,不进批量确认) */ - const pack = findInstalledPackById(packs, ids[0]!); - if (!pack) return c.red(t("packs.notFound", { id: ids[0]! })); - if (isProtectedSfmcPack(pack)) { - return c.red(t("packs.uninstall.protected", { folder: pack.folderName })); - } - selected = [pack]; } else { - const missing: string[] = []; - const byUuid = new Map(); - for (const id of ids) { - const pack = findInstalledPackById(packs, id); - if (!pack) { - missing.push(id); - continue; - } - byUuid.set(pack.uuid.toLowerCase(), pack); + /* 单/多 id 共用解析原语(LSP:任一受保护包整体中止) */ + const resolved = resolveUninstallTargets(packs, ids); + if (resolved.status === "not_found") { + return c.red(t("packs.notFound", { id: resolved.missing.join(", ") })); } - if (missing.length) { - return c.red(t("packs.notFound", { id: missing.join(", ") })); + if (resolved.status === "protected") { + return c.red(t("packs.uninstall.protected", { folder: resolved.folder })); } - selected = [...byUuid.values()]; + selected = resolved.selected; } return await uninstallPacksBatch(selected, packs, { diff --git a/sfmc/world-packs-uninstall.test.mjs b/sfmc/world-packs-uninstall.test.mjs index c9b959f..5ef3b63 100644 --- a/sfmc/world-packs-uninstall.test.mjs +++ b/sfmc/world-packs-uninstall.test.mjs @@ -1,5 +1,5 @@ /** - * packs uninstall 策略契约:受保护平台包识别(DRY 与 CLI 同源) + * packs uninstall 策略契约:受保护平台包识别 + 多 id 解析(DRY/LSP 与 CLI 同源) * 需先 `npm run build -w @sfmc-bds/cli` */ import assert from "node:assert/strict"; @@ -7,10 +7,22 @@ import path from "node:path"; import { describe, it } from "node:test"; import { pathToFileURL } from "node:url"; -const { isProtectedSfmcPack } = await import( +const { isProtectedSfmcPack, resolveUninstallTargets } = await import( pathToFileURL(path.resolve("dist/world-packs.js")).href ); +function fakePack(folderName, uuid = `uuid-${folderName}`) { + return { + folderName, + uuid, + name: folderName, + kind: "behavior", + version: [1, 0, 0], + enabled: true, + dir: `/tmp/${folderName}`, + }; +} + describe("packs uninstall protection", () => { it("识别 sfmc-modules / sfmc-modules-rp,忽略大小写", () => { assert.equal(isProtectedSfmcPack({ folderName: "sfmc-modules" }), true); @@ -19,3 +31,39 @@ describe("packs uninstall protection", () => { assert.equal(isProtectedSfmcPack({ folderName: "[BP] MyAddon" }), false); }); }); + +describe("resolveUninstallTargets LSP", () => { + const packs = [ + fakePack("sfmc-modules", "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa"), + fakePack("[BP] Addon", "bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb"), + fakePack("[BP] Other", "cccccccc-cccc-cccc-cccc-cccccccccccc"), + ]; + + it("单 id 受保护 → protected", () => { + const r = resolveUninstallTargets(packs, ["sfmc-modules"]); + assert.equal(r.status, "protected"); + assert.equal(r.folder, "sfmc-modules"); + }); + + it("多 id 含受保护 → 整体 protected(不返回 ok)", () => { + const r = resolveUninstallTargets(packs, ["[BP] Addon", "sfmc-modules"]); + assert.equal(r.status, "protected"); + assert.equal(r.folder, "sfmc-modules"); + }); + + it("多 id 均可卸 → ok,按 uuid 去重", () => { + const r = resolveUninstallTargets(packs, [ + "[BP] Addon", + "bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb", + "[BP] Other", + ]); + assert.equal(r.status, "ok"); + assert.equal(r.selected.length, 2); + }); + + it("缺 id → not_found", () => { + const r = resolveUninstallTargets(packs, ["nope", "[BP] Addon"]); + assert.equal(r.status, "not_found"); + assert.deepEqual(r.missing, ["nope"]); + }); +}); diff --git a/sfmc/wrap-log-line.test.mjs b/sfmc/wrap-log-line.test.mjs index a56aa69..5c85587 100644 --- a/sfmc/wrap-log-line.test.mjs +++ b/sfmc/wrap-log-line.test.mjs @@ -77,3 +77,25 @@ test("logPrefixWidth 与常见源标签匹配", () => { /* HH:MM:SS + space + [ PACK ] + space + [OK] + space */ assert.ok(w >= 24 && w <= 28, `prefix width unexpected: ${w}`); }); + +test("resolveDisplayLevel:BDS 行从正文解析,其余用 entry.level", async () => { + const { resolveDisplayLevel } = await import("./dist/logs.js"); + assert.equal( + resolveDisplayLevel({ + time: new Date(), + text: "[2026-07-18 23:56:06:778 ERROR] boom", + source: "bds", + level: "info", + }), + "error" + ); + assert.equal( + resolveDisplayLevel({ + time: new Date(), + text: "ok", + source: "pack", + level: "success", + }), + "success" + ); +}); From 360be9cd9bc62f8b7a23ca23f12a73536338fcaf Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 24 Jul 2026 17:31:26 +0000 Subject: [PATCH 2/3] fix: parse BDS timestamp prefix level in getLogLevel Share BDS_TS_PREFIX_RE between stripLogPrefix and getLogLevel so resolveDisplayLevel maps [ts ERROR] lines to error (DRY). Co-authored-by: Shiroha --- remote-controller/dist/index.js.map | 2 +- sfmc/src/logs.ts | 10 ++++++++-- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/remote-controller/dist/index.js.map b/remote-controller/dist/index.js.map index 7a2d3cb..9b2060b 100644 --- a/remote-controller/dist/index.js.map +++ b/remote-controller/dist/index.js.map @@ -1,7 +1,7 @@ { "version": 3, "sources": ["../src/index.ts"], - "sourcesContent": ["import { randomBytes, randomUUID, timingSafeEqual } from \"node:crypto\";\r\nimport { createServer, type IncomingMessage, type ServerResponse } from \"node:http\";\r\nimport fs from \"node:fs\";\r\nimport path from \"node:path\";\r\nimport { WebSocket, WebSocketServer, type RawData } from \"ws\";\r\n\r\ntype AgentRecord = { id: string; name: string; secret: string; createdAt: string; lastSeenAt?: string };\r\ntype TaskAction = \"status\" | \"start\" | \"stop\" | \"restart\" | \"send\";\r\ntype Task = {\r\n id: string;\r\n agentId: string;\r\n action: TaskAction;\r\n service?: string;\r\n message?: string;\r\n status: \"queued\" | \"running\" | \"complete\" | \"failed\";\r\n result?: unknown;\r\n error?: string;\r\n createdAt: string;\r\n completedAt?: string;\r\n};\r\ntype State = { agents: Record; tasks: Record };\r\n\r\nconst serviceNames = new Set([\"bds\", \"db\", \"qq\", \"llbot\"]);\r\nconst actionsRequiringService = new Set([\"start\", \"stop\", \"restart\", \"send\"]);\r\n\r\nconst port = Number(process.env.REMOTE_PORT ?? 3100);\r\nconst host = process.env.REMOTE_HOST ?? \"127.0.0.1\";\r\nconst enrollmentToken = process.env.REMOTE_ENROLL_TOKEN ?? \"\";\r\nconst adminToken = process.env.REMOTE_ADMIN_TOKEN ?? \"\";\r\nconst stateFile = path.resolve(process.env.REMOTE_STATE_FILE ?? \"data/remote-controller.json\");\r\nconst heartbeatIntervalMs = Number(process.env.REMOTE_HEARTBEAT_MS ?? 25_000);\r\n\r\nif (!enrollmentToken || !adminToken) {\r\n console.error(\"[remote-controller] missing required env vars.\");\r\n console.error(\" set both: REMOTE_ENROLL_TOKEN= and REMOTE_ADMIN_TOKEN=\");\r\n console.error(\" example:\");\r\n console.error(\" REMOTE_ENROLL_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\r\n console.error(\" REMOTE_ADMIN_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\r\n console.error(\" node remote-controller/dist/index.js\");\r\n process.exit(1);\r\n}\r\n\r\nconst connected = new Map();\r\n\r\nfunction loadState(): State {\r\n try {\r\n return JSON.parse(fs.readFileSync(stateFile, \"utf-8\")) as State;\r\n } catch {\r\n return { agents: {}, tasks: {} };\r\n }\r\n}\r\n\r\nlet state = loadState();\r\nfunction saveState(): void {\r\n fs.mkdirSync(path.dirname(stateFile), { recursive: true });\r\n fs.writeFileSync(stateFile, `${JSON.stringify(state, null, 2)}\\n`, \"utf-8\");\r\n}\r\n\r\nfunction json(res: ServerResponse, status: number, body: unknown): void {\r\n res.writeHead(status, { \"content-type\": \"application/json\" });\r\n res.end(JSON.stringify(body));\r\n}\r\n\r\nasync function body(req: IncomingMessage): Promise> {\r\n let raw = \"\";\r\n for await (const chunk of req) raw += String(chunk);\r\n return raw ? (JSON.parse(raw) as Record) : {};\r\n}\r\n\r\nfunction authorized(req: IncomingMessage, token: string): boolean {\r\n const value = req.headers.authorization?.replace(/^Bearer\\s+/i, \"\") ?? \"\";\r\n const actual = Buffer.from(value);\r\n const expected = Buffer.from(token);\r\n return actual.length === expected.length && timingSafeEqual(actual, expected);\r\n}\r\n\r\nfunction validAction(value: unknown): value is TaskAction {\r\n return value === \"status\" || value === \"start\" || value === \"stop\" || value === \"restart\" || value === \"send\";\r\n}\r\n\r\nfunction agentPublic(agent: AgentRecord): Omit & { connected: boolean } {\r\n return { id: agent.id, name: agent.name, createdAt: agent.createdAt, ...(agent.lastSeenAt ? { lastSeenAt: agent.lastSeenAt } : {}), connected: connected.has(agent.id) };\r\n}\r\n\r\nfunction dispatchQueuedTasks(agentId: string): void {\r\n const socket = connected.get(agentId);\r\n if (socket?.readyState !== WebSocket.OPEN) return;\r\n\r\n let changed = false;\r\n for (const task of Object.values(state.tasks)) {\r\n if (task.agentId !== agentId || task.status !== \"queued\") continue;\r\n task.status = \"running\";\r\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\r\n if (task.service) payload.service = task.service;\r\n if (task.message) payload.message = task.message;\r\n socket.send(JSON.stringify(payload));\r\n changed = true;\r\n }\r\n if (changed) saveState();\r\n}\r\n\r\nconst server = createServer(async (req, res) => {\r\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\r\n try {\r\n if (req.method === \"GET\" && url.pathname === \"/v1/health\") {\r\n return json(res, 200, { ok: true, agents: Object.keys(state.agents).length, connected: connected.size });\r\n }\r\n\r\n if (req.method === \"POST\" && url.pathname === \"/v1/enroll\") {\r\n if (!authorized(req, enrollmentToken)) return json(res, 401, { error: \"unauthorized\" });\r\n const input = await body(req);\r\n const id = randomUUID();\r\n const secret = randomBytes(32).toString(\"base64url\");\r\n state.agents[id] = { id, secret, name: String(input.name ?? \"sfmc-agent\"), createdAt: new Date().toISOString() };\r\n saveState();\r\n return json(res, 201, { agentId: id, agentSecret: secret });\r\n }\r\n\r\n if (!authorized(req, adminToken)) return json(res, 401, { error: \"unauthorized\" });\r\n\r\n if (req.method === \"GET\" && url.pathname === \"/v1/agents\") {\r\n const agents = Object.values(state.agents).map((a) => agentPublic(a));\r\n return json(res, 200, { agents });\r\n }\r\n\r\n const agentMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)$/);\r\n if (req.method === \"GET\" && agentMatch?.[1]) {\r\n const agent = state.agents[agentMatch[1]];\r\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\r\n return json(res, 200, agentPublic(agent));\r\n }\r\n if (req.method === \"DELETE\" && agentMatch?.[1]) {\r\n const agentId = agentMatch[1];\r\n const agent = state.agents[agentId];\r\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\r\n const socket = connected.get(agentId);\r\n if (socket && socket.readyState === WebSocket.OPEN) socket.close(1000, \"deleted_by_admin\");\r\n connected.delete(agentId);\r\n delete state.agents[agentId];\r\n for (const task of Object.values(state.tasks)) {\r\n if (task.agentId === agentId && (task.status === \"queued\" || task.status === \"running\")) {\r\n task.status = \"failed\";\r\n task.error = \"agent_deleted\";\r\n task.completedAt = new Date().toISOString();\r\n }\r\n }\r\n saveState();\r\n return json(res, 204, {});\r\n }\r\n\r\n const tasksMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)\\/tasks$/);\r\n if (req.method === \"GET\" && tasksMatch?.[1]) {\r\n const agentId = tasksMatch[1];\r\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\r\n const statusFilter = url.searchParams.get(\"status\");\r\n const limit = Math.min(Number(url.searchParams.get(\"limit\") ?? 50), 200);\r\n const tasks = Object.values(state.tasks)\r\n .filter((t) => t.agentId === agentId && (!statusFilter || t.status === statusFilter))\r\n .sort((a, b) => b.createdAt.localeCompare(a.createdAt))\r\n .slice(0, limit);\r\n return json(res, 200, { tasks });\r\n }\r\n\r\n if (req.method === \"POST\" && tasksMatch?.[1]) {\r\n const agentId = tasksMatch[1];\r\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\r\n const input = await body(req);\r\n if (!validAction(input.action)) return json(res, 400, { error: \"invalid_action\" });\r\n if (actionsRequiringService.has(input.action)) {\r\n if (typeof input.service !== \"string\") return json(res, 400, { error: \"service_required\" });\r\n if (!serviceNames.has(input.service)) return json(res, 400, { error: \"invalid_service\" });\r\n }\r\n if (input.action === \"send\" && (typeof input.message !== \"string\" || !input.message.length)) {\r\n return json(res, 400, { error: \"message_required\" });\r\n }\r\n const task: Task = {\r\n id: randomUUID(),\r\n agentId,\r\n action: input.action,\r\n ...(typeof input.service === \"string\" ? { service: input.service } : {}),\r\n ...(typeof input.message === \"string\" ? { message: input.message } : {}),\r\n status: \"queued\",\r\n createdAt: new Date().toISOString(),\r\n };\r\n state.tasks[task.id] = task;\r\n const socket = connected.get(agentId);\r\n if (socket?.readyState === WebSocket.OPEN) {\r\n task.status = \"running\";\r\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\r\n if (task.service) payload.service = task.service;\r\n if (task.message) payload.message = task.message;\r\n socket.send(JSON.stringify(payload));\r\n }\r\n saveState();\r\n return json(res, 202, task);\r\n }\r\n\r\n const getTaskMatch = url.pathname.match(/^\\/v1\\/tasks\\/([^/]+)$/);\r\n if (req.method === \"GET\" && getTaskMatch?.[1]) {\r\n const task = state.tasks[getTaskMatch[1]];\r\n return task ? json(res, 200, task) : json(res, 404, { error: \"task_not_found\" });\r\n }\r\n return json(res, 404, { error: \"not_found\" });\r\n } catch (error) {\r\n return json(res, 400, { error: (error as Error).message });\r\n }\r\n});\r\n\r\nconst wss = new WebSocketServer({ noServer: true });\r\nconst heartbeat = setInterval(() => {\r\n for (const [agentId, socket] of connected) {\r\n if (socket.readyState !== WebSocket.OPEN) continue;\r\n if ((socket as WebSocket & { isAlive?: boolean }).isAlive === false) {\r\n socket.terminate();\r\n connected.delete(agentId);\r\n continue;\r\n }\r\n (socket as WebSocket & { isAlive?: boolean }).isAlive = false;\r\n try {\r\n socket.ping();\r\n } catch {\r\n /* ignore */\r\n }\r\n }\r\n}, heartbeatIntervalMs);\r\n\r\nwss.on(\"close\", () => clearInterval(heartbeat));\r\n\r\nwss.on(\"connection\", (socket: WebSocket, req: IncomingMessage, agentId: string) => {\r\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\r\n socket.on(\"pong\", () => {\r\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\r\n });\r\n\r\n let authenticated = false;\r\n socket.once(\"message\", (raw: RawData) => {\r\n try {\r\n const message = JSON.parse(raw.toString()) as { type?: string; agentId?: string; secret?: string };\r\n const agent = state.agents[agentId];\r\n const authRequest = { headers: { authorization: `Bearer ${message.secret ?? \"\"}` } } as IncomingMessage;\r\n if (message.type !== \"hello\" || message.agentId !== agentId || !agent || !authorized(authRequest, agent.secret)) {\r\n socket.close(1008, \"unauthorized\");\r\n return;\r\n }\r\n authenticated = true;\r\n agent.lastSeenAt = new Date().toISOString();\r\n connected.set(agentId, socket);\r\n saveState();\r\n dispatchQueuedTasks(agentId);\r\n } catch {\r\n socket.close(1008, \"invalid_hello\");\r\n }\r\n });\r\n socket.on(\"message\", (raw: RawData) => {\r\n if (!authenticated) return;\r\n try {\r\n const message = JSON.parse(raw.toString()) as {\r\n type?: string;\r\n taskId?: string;\r\n ok?: boolean;\r\n result?: unknown;\r\n error?: string;\r\n };\r\n if (message.type === \"ping\") return;\r\n if (message.type !== \"task_result\" || !message.taskId) return;\r\n const task = state.tasks[message.taskId];\r\n if (!task || task.agentId !== agentId) return;\r\n task.status = message.ok ? \"complete\" : \"failed\";\r\n task.completedAt = new Date().toISOString();\r\n if (message.ok) task.result = message.result;\r\n else task.error = message.error ?? \"task failed\";\r\n saveState();\r\n } catch {\r\n /* Ignore malformed agent messages. */\r\n }\r\n });\r\n socket.on(\"close\", () => {\r\n if (connected.get(agentId) === socket) connected.delete(agentId);\r\n });\r\n void req;\r\n});\r\n\r\nserver.on(\"upgrade\", (req, socket, head) => {\r\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\r\n const agentId = url.pathname === \"/v1/agent\" ? url.searchParams.get(\"id\") : null;\r\n if (!agentId || !state.agents[agentId]) return socket.destroy();\r\n wss.handleUpgrade(req, socket, head, (ws) => wss.emit(\"connection\", ws, req, agentId));\r\n});\r\n\r\nserver.listen(port, host, () => {\r\n console.log(`[remote-controller] listening on http://${host}:${port}`);\r\n console.log(`[remote-controller] state file: ${stateFile}`);\r\n console.log(`[remote-controller] endpoints:`);\r\n console.log(` POST /v1/enroll (enroll token)`);\r\n console.log(` GET /v1/health (open)`);\r\n console.log(` GET /v1/agents (admin token)`);\r\n console.log(` GET /v1/agents/{id} (admin token)`);\r\n console.log(` DELETE /v1/agents/{id} (admin token)`);\r\n console.log(` POST /v1/agents/{id}/tasks (admin token)`);\r\n console.log(` GET /v1/agents/{id}/tasks (admin token)`);\r\n console.log(` GET /v1/tasks/{id} (admin token)`);\r\n console.log(` WS /v1/agent?id={id} (per-agent secret)`);\r\n console.log(`[remote-controller] heartbeat: ${heartbeatIntervalMs}ms`);\r\n});"], + "sourcesContent": ["import { randomBytes, randomUUID, timingSafeEqual } from \"node:crypto\";\nimport { createServer, type IncomingMessage, type ServerResponse } from \"node:http\";\nimport fs from \"node:fs\";\nimport path from \"node:path\";\nimport { WebSocket, WebSocketServer, type RawData } from \"ws\";\n\ntype AgentRecord = { id: string; name: string; secret: string; createdAt: string; lastSeenAt?: string };\ntype TaskAction = \"status\" | \"start\" | \"stop\" | \"restart\" | \"send\";\ntype Task = {\n id: string;\n agentId: string;\n action: TaskAction;\n service?: string;\n message?: string;\n status: \"queued\" | \"running\" | \"complete\" | \"failed\";\n result?: unknown;\n error?: string;\n createdAt: string;\n completedAt?: string;\n};\ntype State = { agents: Record; tasks: Record };\n\nconst serviceNames = new Set([\"bds\", \"db\", \"qq\", \"llbot\"]);\nconst actionsRequiringService = new Set([\"start\", \"stop\", \"restart\", \"send\"]);\n\nconst port = Number(process.env.REMOTE_PORT ?? 3100);\nconst host = process.env.REMOTE_HOST ?? \"127.0.0.1\";\nconst enrollmentToken = process.env.REMOTE_ENROLL_TOKEN ?? \"\";\nconst adminToken = process.env.REMOTE_ADMIN_TOKEN ?? \"\";\nconst stateFile = path.resolve(process.env.REMOTE_STATE_FILE ?? \"data/remote-controller.json\");\nconst heartbeatIntervalMs = Number(process.env.REMOTE_HEARTBEAT_MS ?? 25_000);\n\nif (!enrollmentToken || !adminToken) {\n console.error(\"[remote-controller] missing required env vars.\");\n console.error(\" set both: REMOTE_ENROLL_TOKEN= and REMOTE_ADMIN_TOKEN=\");\n console.error(\" example:\");\n console.error(\" REMOTE_ENROLL_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\n console.error(\" REMOTE_ADMIN_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\n console.error(\" node remote-controller/dist/index.js\");\n process.exit(1);\n}\n\nconst connected = new Map();\n\nfunction loadState(): State {\n try {\n return JSON.parse(fs.readFileSync(stateFile, \"utf-8\")) as State;\n } catch {\n return { agents: {}, tasks: {} };\n }\n}\n\nlet state = loadState();\nfunction saveState(): void {\n fs.mkdirSync(path.dirname(stateFile), { recursive: true });\n fs.writeFileSync(stateFile, `${JSON.stringify(state, null, 2)}\\n`, \"utf-8\");\n}\n\nfunction json(res: ServerResponse, status: number, body: unknown): void {\n res.writeHead(status, { \"content-type\": \"application/json\" });\n res.end(JSON.stringify(body));\n}\n\nasync function body(req: IncomingMessage): Promise> {\n let raw = \"\";\n for await (const chunk of req) raw += String(chunk);\n return raw ? (JSON.parse(raw) as Record) : {};\n}\n\nfunction authorized(req: IncomingMessage, token: string): boolean {\n const value = req.headers.authorization?.replace(/^Bearer\\s+/i, \"\") ?? \"\";\n const actual = Buffer.from(value);\n const expected = Buffer.from(token);\n return actual.length === expected.length && timingSafeEqual(actual, expected);\n}\n\nfunction validAction(value: unknown): value is TaskAction {\n return value === \"status\" || value === \"start\" || value === \"stop\" || value === \"restart\" || value === \"send\";\n}\n\nfunction agentPublic(agent: AgentRecord): Omit & { connected: boolean } {\n return { id: agent.id, name: agent.name, createdAt: agent.createdAt, ...(agent.lastSeenAt ? { lastSeenAt: agent.lastSeenAt } : {}), connected: connected.has(agent.id) };\n}\n\nfunction dispatchQueuedTasks(agentId: string): void {\n const socket = connected.get(agentId);\n if (socket?.readyState !== WebSocket.OPEN) return;\n\n let changed = false;\n for (const task of Object.values(state.tasks)) {\n if (task.agentId !== agentId || task.status !== \"queued\") continue;\n task.status = \"running\";\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\n if (task.service) payload.service = task.service;\n if (task.message) payload.message = task.message;\n socket.send(JSON.stringify(payload));\n changed = true;\n }\n if (changed) saveState();\n}\n\nconst server = createServer(async (req, res) => {\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\n try {\n if (req.method === \"GET\" && url.pathname === \"/v1/health\") {\n return json(res, 200, { ok: true, agents: Object.keys(state.agents).length, connected: connected.size });\n }\n\n if (req.method === \"POST\" && url.pathname === \"/v1/enroll\") {\n if (!authorized(req, enrollmentToken)) return json(res, 401, { error: \"unauthorized\" });\n const input = await body(req);\n const id = randomUUID();\n const secret = randomBytes(32).toString(\"base64url\");\n state.agents[id] = { id, secret, name: String(input.name ?? \"sfmc-agent\"), createdAt: new Date().toISOString() };\n saveState();\n return json(res, 201, { agentId: id, agentSecret: secret });\n }\n\n if (!authorized(req, adminToken)) return json(res, 401, { error: \"unauthorized\" });\n\n if (req.method === \"GET\" && url.pathname === \"/v1/agents\") {\n const agents = Object.values(state.agents).map((a) => agentPublic(a));\n return json(res, 200, { agents });\n }\n\n const agentMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)$/);\n if (req.method === \"GET\" && agentMatch?.[1]) {\n const agent = state.agents[agentMatch[1]];\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\n return json(res, 200, agentPublic(agent));\n }\n if (req.method === \"DELETE\" && agentMatch?.[1]) {\n const agentId = agentMatch[1];\n const agent = state.agents[agentId];\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\n const socket = connected.get(agentId);\n if (socket && socket.readyState === WebSocket.OPEN) socket.close(1000, \"deleted_by_admin\");\n connected.delete(agentId);\n delete state.agents[agentId];\n for (const task of Object.values(state.tasks)) {\n if (task.agentId === agentId && (task.status === \"queued\" || task.status === \"running\")) {\n task.status = \"failed\";\n task.error = \"agent_deleted\";\n task.completedAt = new Date().toISOString();\n }\n }\n saveState();\n return json(res, 204, {});\n }\n\n const tasksMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)\\/tasks$/);\n if (req.method === \"GET\" && tasksMatch?.[1]) {\n const agentId = tasksMatch[1];\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\n const statusFilter = url.searchParams.get(\"status\");\n const limit = Math.min(Number(url.searchParams.get(\"limit\") ?? 50), 200);\n const tasks = Object.values(state.tasks)\n .filter((t) => t.agentId === agentId && (!statusFilter || t.status === statusFilter))\n .sort((a, b) => b.createdAt.localeCompare(a.createdAt))\n .slice(0, limit);\n return json(res, 200, { tasks });\n }\n\n if (req.method === \"POST\" && tasksMatch?.[1]) {\n const agentId = tasksMatch[1];\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\n const input = await body(req);\n if (!validAction(input.action)) return json(res, 400, { error: \"invalid_action\" });\n if (actionsRequiringService.has(input.action)) {\n if (typeof input.service !== \"string\") return json(res, 400, { error: \"service_required\" });\n if (!serviceNames.has(input.service)) return json(res, 400, { error: \"invalid_service\" });\n }\n if (input.action === \"send\" && (typeof input.message !== \"string\" || !input.message.length)) {\n return json(res, 400, { error: \"message_required\" });\n }\n const task: Task = {\n id: randomUUID(),\n agentId,\n action: input.action,\n ...(typeof input.service === \"string\" ? { service: input.service } : {}),\n ...(typeof input.message === \"string\" ? { message: input.message } : {}),\n status: \"queued\",\n createdAt: new Date().toISOString(),\n };\n state.tasks[task.id] = task;\n const socket = connected.get(agentId);\n if (socket?.readyState === WebSocket.OPEN) {\n task.status = \"running\";\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\n if (task.service) payload.service = task.service;\n if (task.message) payload.message = task.message;\n socket.send(JSON.stringify(payload));\n }\n saveState();\n return json(res, 202, task);\n }\n\n const getTaskMatch = url.pathname.match(/^\\/v1\\/tasks\\/([^/]+)$/);\n if (req.method === \"GET\" && getTaskMatch?.[1]) {\n const task = state.tasks[getTaskMatch[1]];\n return task ? json(res, 200, task) : json(res, 404, { error: \"task_not_found\" });\n }\n return json(res, 404, { error: \"not_found\" });\n } catch (error) {\n return json(res, 400, { error: (error as Error).message });\n }\n});\n\nconst wss = new WebSocketServer({ noServer: true });\nconst heartbeat = setInterval(() => {\n for (const [agentId, socket] of connected) {\n if (socket.readyState !== WebSocket.OPEN) continue;\n if ((socket as WebSocket & { isAlive?: boolean }).isAlive === false) {\n socket.terminate();\n connected.delete(agentId);\n continue;\n }\n (socket as WebSocket & { isAlive?: boolean }).isAlive = false;\n try {\n socket.ping();\n } catch {\n /* ignore */\n }\n }\n}, heartbeatIntervalMs);\n\nwss.on(\"close\", () => clearInterval(heartbeat));\n\nwss.on(\"connection\", (socket: WebSocket, req: IncomingMessage, agentId: string) => {\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\n socket.on(\"pong\", () => {\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\n });\n\n let authenticated = false;\n socket.once(\"message\", (raw: RawData) => {\n try {\n const message = JSON.parse(raw.toString()) as { type?: string; agentId?: string; secret?: string };\n const agent = state.agents[agentId];\n const authRequest = { headers: { authorization: `Bearer ${message.secret ?? \"\"}` } } as IncomingMessage;\n if (message.type !== \"hello\" || message.agentId !== agentId || !agent || !authorized(authRequest, agent.secret)) {\n socket.close(1008, \"unauthorized\");\n return;\n }\n authenticated = true;\n agent.lastSeenAt = new Date().toISOString();\n connected.set(agentId, socket);\n saveState();\n dispatchQueuedTasks(agentId);\n } catch {\n socket.close(1008, \"invalid_hello\");\n }\n });\n socket.on(\"message\", (raw: RawData) => {\n if (!authenticated) return;\n try {\n const message = JSON.parse(raw.toString()) as {\n type?: string;\n taskId?: string;\n ok?: boolean;\n result?: unknown;\n error?: string;\n };\n if (message.type === \"ping\") return;\n if (message.type !== \"task_result\" || !message.taskId) return;\n const task = state.tasks[message.taskId];\n if (!task || task.agentId !== agentId) return;\n task.status = message.ok ? \"complete\" : \"failed\";\n task.completedAt = new Date().toISOString();\n if (message.ok) task.result = message.result;\n else task.error = message.error ?? \"task failed\";\n saveState();\n } catch {\n /* Ignore malformed agent messages. */\n }\n });\n socket.on(\"close\", () => {\n if (connected.get(agentId) === socket) connected.delete(agentId);\n });\n void req;\n});\n\nserver.on(\"upgrade\", (req, socket, head) => {\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\n const agentId = url.pathname === \"/v1/agent\" ? url.searchParams.get(\"id\") : null;\n if (!agentId || !state.agents[agentId]) return socket.destroy();\n wss.handleUpgrade(req, socket, head, (ws) => wss.emit(\"connection\", ws, req, agentId));\n});\n\nserver.listen(port, host, () => {\n console.log(`[remote-controller] listening on http://${host}:${port}`);\n console.log(`[remote-controller] state file: ${stateFile}`);\n console.log(`[remote-controller] endpoints:`);\n console.log(` POST /v1/enroll (enroll token)`);\n console.log(` GET /v1/health (open)`);\n console.log(` GET /v1/agents (admin token)`);\n console.log(` GET /v1/agents/{id} (admin token)`);\n console.log(` DELETE /v1/agents/{id} (admin token)`);\n console.log(` POST /v1/agents/{id}/tasks (admin token)`);\n console.log(` GET /v1/agents/{id}/tasks (admin token)`);\n console.log(` GET /v1/tasks/{id} (admin token)`);\n console.log(` WS /v1/agent?id={id} (per-agent secret)`);\n console.log(`[remote-controller] heartbeat: ${heartbeatIntervalMs}ms`);\n});"], "mappings": "AAAA,SAAS,aAAa,YAAY,uBAAuB;AACzD,SAAS,oBAA+D;AACxE,OAAO,QAAQ;AACf,OAAO,UAAU;AACjB,SAAS,WAAW,uBAAqC;AAkBzD,MAAM,eAAe,oBAAI,IAAI,CAAC,OAAO,MAAM,MAAM,OAAO,CAAC;AACzD,MAAM,0BAA0B,oBAAI,IAAgB,CAAC,SAAS,QAAQ,WAAW,MAAM,CAAC;AAExF,MAAM,OAAO,OAAO,QAAQ,IAAI,eAAe,IAAI;AACnD,MAAM,OAAO,QAAQ,IAAI,eAAe;AACxC,MAAM,kBAAkB,QAAQ,IAAI,uBAAuB;AAC3D,MAAM,aAAa,QAAQ,IAAI,sBAAsB;AACrD,MAAM,YAAY,KAAK,QAAQ,QAAQ,IAAI,qBAAqB,6BAA6B;AAC7F,MAAM,sBAAsB,OAAO,QAAQ,IAAI,uBAAuB,IAAM;AAE5E,IAAI,CAAC,mBAAmB,CAAC,YAAY;AACnC,UAAQ,MAAM,gDAAgD;AAC9D,UAAQ,MAAM,0EAA0E;AACxF,UAAQ,MAAM,YAAY;AAC1B,UAAQ,MAAM,8GAAgH;AAC9H,UAAQ,MAAM,6GAA+G;AAC7H,UAAQ,MAAM,0CAA0C;AACxD,UAAQ,KAAK,CAAC;AAChB;AAEA,MAAM,YAAY,oBAAI,IAAuB;AAE7C,SAAS,YAAmB;AAC1B,MAAI;AACF,WAAO,KAAK,MAAM,GAAG,aAAa,WAAW,OAAO,CAAC;AAAA,EACvD,QAAQ;AACN,WAAO,EAAE,QAAQ,CAAC,GAAG,OAAO,CAAC,EAAE;AAAA,EACjC;AACF;AAEA,IAAI,QAAQ,UAAU;AACtB,SAAS,YAAkB;AACzB,KAAG,UAAU,KAAK,QAAQ,SAAS,GAAG,EAAE,WAAW,KAAK,CAAC;AACzD,KAAG,cAAc,WAAW,GAAG,KAAK,UAAU,OAAO,MAAM,CAAC,CAAC;AAAA,GAAM,OAAO;AAC5E;AAEA,SAAS,KAAK,KAAqB,QAAgBA,OAAqB;AACtE,MAAI,UAAU,QAAQ,EAAE,gBAAgB,mBAAmB,CAAC;AAC5D,MAAI,IAAI,KAAK,UAAUA,KAAI,CAAC;AAC9B;AAEA,eAAe,KAAK,KAAwD;AAC1E,MAAI,MAAM;AACV,mBAAiB,SAAS,IAAK,QAAO,OAAO,KAAK;AAClD,SAAO,MAAO,KAAK,MAAM,GAAG,IAAgC,CAAC;AAC/D;AAEA,SAAS,WAAW,KAAsB,OAAwB;AAChE,QAAM,QAAQ,IAAI,QAAQ,eAAe,QAAQ,eAAe,EAAE,KAAK;AACvE,QAAM,SAAS,OAAO,KAAK,KAAK;AAChC,QAAM,WAAW,OAAO,KAAK,KAAK;AAClC,SAAO,OAAO,WAAW,SAAS,UAAU,gBAAgB,QAAQ,QAAQ;AAC9E;AAEA,SAAS,YAAY,OAAqC;AACxD,SAAO,UAAU,YAAY,UAAU,WAAW,UAAU,UAAU,UAAU,aAAa,UAAU;AACzG;AAEA,SAAS,YAAY,OAA0E;AAC7F,SAAO,EAAE,IAAI,MAAM,IAAI,MAAM,MAAM,MAAM,WAAW,MAAM,WAAW,GAAI,MAAM,aAAa,EAAE,YAAY,MAAM,WAAW,IAAI,CAAC,GAAI,WAAW,UAAU,IAAI,MAAM,EAAE,EAAE;AACzK;AAEA,SAAS,oBAAoB,SAAuB;AAClD,QAAM,SAAS,UAAU,IAAI,OAAO;AACpC,MAAI,QAAQ,eAAe,UAAU,KAAM;AAE3C,MAAI,UAAU;AACd,aAAW,QAAQ,OAAO,OAAO,MAAM,KAAK,GAAG;AAC7C,QAAI,KAAK,YAAY,WAAW,KAAK,WAAW,SAAU;AAC1D,SAAK,SAAS;AACd,UAAM,UAAmC,EAAE,MAAM,QAAQ,QAAQ,KAAK,IAAI,QAAQ,KAAK,OAAO;AAC9F,QAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,QAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,WAAO,KAAK,KAAK,UAAU,OAAO,CAAC;AACnC,cAAU;AAAA,EACZ;AACA,MAAI,QAAS,WAAU;AACzB;AAEA,MAAM,SAAS,aAAa,OAAO,KAAK,QAAQ;AAC9C,QAAM,MAAM,IAAI,IAAI,IAAI,OAAO,KAAK,UAAU,IAAI,QAAQ,QAAQ,WAAW,EAAE;AAC/E,MAAI;AACF,QAAI,IAAI,WAAW,SAAS,IAAI,aAAa,cAAc;AACzD,aAAO,KAAK,KAAK,KAAK,EAAE,IAAI,MAAM,QAAQ,OAAO,KAAK,MAAM,MAAM,EAAE,QAAQ,WAAW,UAAU,KAAK,CAAC;AAAA,IACzG;AAEA,QAAI,IAAI,WAAW,UAAU,IAAI,aAAa,cAAc;AAC1D,UAAI,CAAC,WAAW,KAAK,eAAe,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,eAAe,CAAC;AACtF,YAAM,QAAQ,MAAM,KAAK,GAAG;AAC5B,YAAM,KAAK,WAAW;AACtB,YAAM,SAAS,YAAY,EAAE,EAAE,SAAS,WAAW;AACnD,YAAM,OAAO,EAAE,IAAI,EAAE,IAAI,QAAQ,MAAM,OAAO,MAAM,QAAQ,YAAY,GAAG,YAAW,oBAAI,KAAK,GAAE,YAAY,EAAE;AAC/G,gBAAU;AACV,aAAO,KAAK,KAAK,KAAK,EAAE,SAAS,IAAI,aAAa,OAAO,CAAC;AAAA,IAC5D;AAEA,QAAI,CAAC,WAAW,KAAK,UAAU,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,eAAe,CAAC;AAEjF,QAAI,IAAI,WAAW,SAAS,IAAI,aAAa,cAAc;AACzD,YAAM,SAAS,OAAO,OAAO,MAAM,MAAM,EAAE,IAAI,CAAC,MAAM,YAAY,CAAC,CAAC;AACpE,aAAO,KAAK,KAAK,KAAK,EAAE,OAAO,CAAC;AAAA,IAClC;AAEA,UAAM,aAAa,IAAI,SAAS,MAAM,yBAAyB;AAC/D,QAAI,IAAI,WAAW,SAAS,aAAa,CAAC,GAAG;AAC3C,YAAM,QAAQ,MAAM,OAAO,WAAW,CAAC,CAAC;AACxC,UAAI,CAAC,MAAO,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9D,aAAO,KAAK,KAAK,KAAK,YAAY,KAAK,CAAC;AAAA,IAC1C;AACA,QAAI,IAAI,WAAW,YAAY,aAAa,CAAC,GAAG;AAC9C,YAAM,UAAU,WAAW,CAAC;AAC5B,YAAM,QAAQ,MAAM,OAAO,OAAO;AAClC,UAAI,CAAC,MAAO,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9D,YAAM,SAAS,UAAU,IAAI,OAAO;AACpC,UAAI,UAAU,OAAO,eAAe,UAAU,KAAM,QAAO,MAAM,KAAM,kBAAkB;AACzF,gBAAU,OAAO,OAAO;AACxB,aAAO,MAAM,OAAO,OAAO;AAC3B,iBAAW,QAAQ,OAAO,OAAO,MAAM,KAAK,GAAG;AAC7C,YAAI,KAAK,YAAY,YAAY,KAAK,WAAW,YAAY,KAAK,WAAW,YAAY;AACvF,eAAK,SAAS;AACd,eAAK,QAAQ;AACb,eAAK,eAAc,oBAAI,KAAK,GAAE,YAAY;AAAA,QAC5C;AAAA,MACF;AACA,gBAAU;AACV,aAAO,KAAK,KAAK,KAAK,CAAC,CAAC;AAAA,IAC1B;AAEA,UAAM,aAAa,IAAI,SAAS,MAAM,gCAAgC;AACtE,QAAI,IAAI,WAAW,SAAS,aAAa,CAAC,GAAG;AAC3C,YAAM,UAAU,WAAW,CAAC;AAC5B,UAAI,CAAC,MAAM,OAAO,OAAO,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9E,YAAM,eAAe,IAAI,aAAa,IAAI,QAAQ;AAClD,YAAM,QAAQ,KAAK,IAAI,OAAO,IAAI,aAAa,IAAI,OAAO,KAAK,EAAE,GAAG,GAAG;AACvE,YAAM,QAAQ,OAAO,OAAO,MAAM,KAAK,EACpC,OAAO,CAAC,MAAM,EAAE,YAAY,YAAY,CAAC,gBAAgB,EAAE,WAAW,aAAa,EACnF,KAAK,CAAC,GAAG,MAAM,EAAE,UAAU,cAAc,EAAE,SAAS,CAAC,EACrD,MAAM,GAAG,KAAK;AACjB,aAAO,KAAK,KAAK,KAAK,EAAE,MAAM,CAAC;AAAA,IACjC;AAEA,QAAI,IAAI,WAAW,UAAU,aAAa,CAAC,GAAG;AAC5C,YAAM,UAAU,WAAW,CAAC;AAC5B,UAAI,CAAC,MAAM,OAAO,OAAO,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9E,YAAM,QAAQ,MAAM,KAAK,GAAG;AAC5B,UAAI,CAAC,YAAY,MAAM,MAAM,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,iBAAiB,CAAC;AACjF,UAAI,wBAAwB,IAAI,MAAM,MAAM,GAAG;AAC7C,YAAI,OAAO,MAAM,YAAY,SAAU,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,mBAAmB,CAAC;AAC1F,YAAI,CAAC,aAAa,IAAI,MAAM,OAAO,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAAA,MAC1F;AACA,UAAI,MAAM,WAAW,WAAW,OAAO,MAAM,YAAY,YAAY,CAAC,MAAM,QAAQ,SAAS;AAC3F,eAAO,KAAK,KAAK,KAAK,EAAE,OAAO,mBAAmB,CAAC;AAAA,MACrD;AACA,YAAM,OAAa;AAAA,QACjB,IAAI,WAAW;AAAA,QACf;AAAA,QACA,QAAQ,MAAM;AAAA,QACd,GAAI,OAAO,MAAM,YAAY,WAAW,EAAE,SAAS,MAAM,QAAQ,IAAI,CAAC;AAAA,QACtE,GAAI,OAAO,MAAM,YAAY,WAAW,EAAE,SAAS,MAAM,QAAQ,IAAI,CAAC;AAAA,QACtE,QAAQ;AAAA,QACR,YAAW,oBAAI,KAAK,GAAE,YAAY;AAAA,MACpC;AACA,YAAM,MAAM,KAAK,EAAE,IAAI;AACvB,YAAM,SAAS,UAAU,IAAI,OAAO;AACpC,UAAI,QAAQ,eAAe,UAAU,MAAM;AACzC,aAAK,SAAS;AACd,cAAM,UAAmC,EAAE,MAAM,QAAQ,QAAQ,KAAK,IAAI,QAAQ,KAAK,OAAO;AAC9F,YAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,YAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,eAAO,KAAK,KAAK,UAAU,OAAO,CAAC;AAAA,MACrC;AACA,gBAAU;AACV,aAAO,KAAK,KAAK,KAAK,IAAI;AAAA,IAC5B;AAEA,UAAM,eAAe,IAAI,SAAS,MAAM,wBAAwB;AAChE,QAAI,IAAI,WAAW,SAAS,eAAe,CAAC,GAAG;AAC7C,YAAM,OAAO,MAAM,MAAM,aAAa,CAAC,CAAC;AACxC,aAAO,OAAO,KAAK,KAAK,KAAK,IAAI,IAAI,KAAK,KAAK,KAAK,EAAE,OAAO,iBAAiB,CAAC;AAAA,IACjF;AACA,WAAO,KAAK,KAAK,KAAK,EAAE,OAAO,YAAY,CAAC;AAAA,EAC9C,SAAS,OAAO;AACd,WAAO,KAAK,KAAK,KAAK,EAAE,OAAQ,MAAgB,QAAQ,CAAC;AAAA,EAC3D;AACF,CAAC;AAED,MAAM,MAAM,IAAI,gBAAgB,EAAE,UAAU,KAAK,CAAC;AAClD,MAAM,YAAY,YAAY,MAAM;AAClC,aAAW,CAAC,SAAS,MAAM,KAAK,WAAW;AACzC,QAAI,OAAO,eAAe,UAAU,KAAM;AAC1C,QAAK,OAA6C,YAAY,OAAO;AACnE,aAAO,UAAU;AACjB,gBAAU,OAAO,OAAO;AACxB;AAAA,IACF;AACA,IAAC,OAA6C,UAAU;AACxD,QAAI;AACF,aAAO,KAAK;AAAA,IACd,QAAQ;AAAA,IAER;AAAA,EACF;AACF,GAAG,mBAAmB;AAEtB,IAAI,GAAG,SAAS,MAAM,cAAc,SAAS,CAAC;AAE9C,IAAI,GAAG,cAAc,CAAC,QAAmB,KAAsB,YAAoB;AACjF,EAAC,OAA6C,UAAU;AACxD,SAAO,GAAG,QAAQ,MAAM;AACtB,IAAC,OAA6C,UAAU;AAAA,EAC1D,CAAC;AAED,MAAI,gBAAgB;AACpB,SAAO,KAAK,WAAW,CAAC,QAAiB;AACvC,QAAI;AACF,YAAM,UAAU,KAAK,MAAM,IAAI,SAAS,CAAC;AACzC,YAAM,QAAQ,MAAM,OAAO,OAAO;AAClC,YAAM,cAAc,EAAE,SAAS,EAAE,eAAe,UAAU,QAAQ,UAAU,EAAE,GAAG,EAAE;AACnF,UAAI,QAAQ,SAAS,WAAW,QAAQ,YAAY,WAAW,CAAC,SAAS,CAAC,WAAW,aAAa,MAAM,MAAM,GAAG;AAC/G,eAAO,MAAM,MAAM,cAAc;AACjC;AAAA,MACF;AACA,sBAAgB;AAChB,YAAM,cAAa,oBAAI,KAAK,GAAE,YAAY;AAC1C,gBAAU,IAAI,SAAS,MAAM;AAC7B,gBAAU;AACV,0BAAoB,OAAO;AAAA,IAC7B,QAAQ;AACN,aAAO,MAAM,MAAM,eAAe;AAAA,IACpC;AAAA,EACF,CAAC;AACD,SAAO,GAAG,WAAW,CAAC,QAAiB;AACrC,QAAI,CAAC,cAAe;AACpB,QAAI;AACF,YAAM,UAAU,KAAK,MAAM,IAAI,SAAS,CAAC;AAOzC,UAAI,QAAQ,SAAS,OAAQ;AAC7B,UAAI,QAAQ,SAAS,iBAAiB,CAAC,QAAQ,OAAQ;AACvD,YAAM,OAAO,MAAM,MAAM,QAAQ,MAAM;AACvC,UAAI,CAAC,QAAQ,KAAK,YAAY,QAAS;AACvC,WAAK,SAAS,QAAQ,KAAK,aAAa;AACxC,WAAK,eAAc,oBAAI,KAAK,GAAE,YAAY;AAC1C,UAAI,QAAQ,GAAI,MAAK,SAAS,QAAQ;AAAA,UACjC,MAAK,QAAQ,QAAQ,SAAS;AACnC,gBAAU;AAAA,IACZ,QAAQ;AAAA,IAER;AAAA,EACF,CAAC;AACD,SAAO,GAAG,SAAS,MAAM;AACvB,QAAI,UAAU,IAAI,OAAO,MAAM,OAAQ,WAAU,OAAO,OAAO;AAAA,EACjE,CAAC;AACD,OAAK;AACP,CAAC;AAED,OAAO,GAAG,WAAW,CAAC,KAAK,QAAQ,SAAS;AAC1C,QAAM,MAAM,IAAI,IAAI,IAAI,OAAO,KAAK,UAAU,IAAI,QAAQ,QAAQ,WAAW,EAAE;AAC/E,QAAM,UAAU,IAAI,aAAa,cAAc,IAAI,aAAa,IAAI,IAAI,IAAI;AAC5E,MAAI,CAAC,WAAW,CAAC,MAAM,OAAO,OAAO,EAAG,QAAO,OAAO,QAAQ;AAC9D,MAAI,cAAc,KAAK,QAAQ,MAAM,CAAC,OAAO,IAAI,KAAK,cAAc,IAAI,KAAK,OAAO,CAAC;AACvF,CAAC;AAED,OAAO,OAAO,MAAM,MAAM,MAAM;AAC9B,UAAQ,IAAI,2CAA2C,IAAI,IAAI,IAAI,EAAE;AACrE,UAAQ,IAAI,mCAAmC,SAAS,EAAE;AAC1D,UAAQ,IAAI,gCAAgC;AAC5C,UAAQ,IAAI,mDAAmD;AAC/D,UAAQ,IAAI,2CAA2C;AACvD,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,uDAAuD;AACnE,UAAQ,IAAI,kCAAkC,mBAAmB,IAAI;AACvE,CAAC;", "names": ["body"] } diff --git a/sfmc/src/logs.ts b/sfmc/src/logs.ts index ce20689..d031120 100644 --- a/sfmc/src/logs.ts +++ b/sfmc/src/logs.ts @@ -131,16 +131,22 @@ export function formatSourceTag(source: string): string { /** 简化BDS日志 */ function stripLogPrefix(line: string): string { // 匹配格式: [2026-07-18 23:56:06:778 INFO] - const prefixRegex = /^\[\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}:\d{3} (INFO|WARNING|ERROR|FATAL|DEBUG)\]\s*/; - return line.replace(prefixRegex, ""); + return line.replace(BDS_TS_PREFIX_RE, ""); } +/** BDS 行首时间戳+级别前缀(strip / 解析级别共用,DRY) */ +const BDS_TS_PREFIX_RE = + /^\[\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}:\d{3} (INFO|WARNING|ERROR|FATAL|DEBUG)\]\s*/i; + /** * 从一行 BDS 日志中提取日志等级 * @param line 日志行字符串 * @returns 日志等级(大写),若无法识别则返回 'UNKNOWN' */ function getLogLevel(line: string): string { + const fromTs = BDS_TS_PREFIX_RE.exec(line); + if (fromTs?.[1]) return fromTs[1].toUpperCase(); + const levelNames = ["INFO", "WARNING", "ERROR", "FATAL", "DEBUG", "TRACE", "WARN"]; const levelPattern = levelNames.join("|"); From e2195607887cbd6a5347f04d00dddd2e0df32ada Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 24 Jul 2026 17:32:03 +0000 Subject: [PATCH 3/3] chore: drop accidental remote-controller sourcemap noise Restore tracked dist map to pre-review state; rebuild churn was unrelated. Co-authored-by: Shiroha --- remote-controller/dist/index.js.map | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/remote-controller/dist/index.js.map b/remote-controller/dist/index.js.map index 9b2060b..7a2d3cb 100644 --- a/remote-controller/dist/index.js.map +++ b/remote-controller/dist/index.js.map @@ -1,7 +1,7 @@ { "version": 3, "sources": ["../src/index.ts"], - "sourcesContent": ["import { randomBytes, randomUUID, timingSafeEqual } from \"node:crypto\";\nimport { createServer, type IncomingMessage, type ServerResponse } from \"node:http\";\nimport fs from \"node:fs\";\nimport path from \"node:path\";\nimport { WebSocket, WebSocketServer, type RawData } from \"ws\";\n\ntype AgentRecord = { id: string; name: string; secret: string; createdAt: string; lastSeenAt?: string };\ntype TaskAction = \"status\" | \"start\" | \"stop\" | \"restart\" | \"send\";\ntype Task = {\n id: string;\n agentId: string;\n action: TaskAction;\n service?: string;\n message?: string;\n status: \"queued\" | \"running\" | \"complete\" | \"failed\";\n result?: unknown;\n error?: string;\n createdAt: string;\n completedAt?: string;\n};\ntype State = { agents: Record; tasks: Record };\n\nconst serviceNames = new Set([\"bds\", \"db\", \"qq\", \"llbot\"]);\nconst actionsRequiringService = new Set([\"start\", \"stop\", \"restart\", \"send\"]);\n\nconst port = Number(process.env.REMOTE_PORT ?? 3100);\nconst host = process.env.REMOTE_HOST ?? \"127.0.0.1\";\nconst enrollmentToken = process.env.REMOTE_ENROLL_TOKEN ?? \"\";\nconst adminToken = process.env.REMOTE_ADMIN_TOKEN ?? \"\";\nconst stateFile = path.resolve(process.env.REMOTE_STATE_FILE ?? \"data/remote-controller.json\");\nconst heartbeatIntervalMs = Number(process.env.REMOTE_HEARTBEAT_MS ?? 25_000);\n\nif (!enrollmentToken || !adminToken) {\n console.error(\"[remote-controller] missing required env vars.\");\n console.error(\" set both: REMOTE_ENROLL_TOKEN= and REMOTE_ADMIN_TOKEN=\");\n console.error(\" example:\");\n console.error(\" REMOTE_ENROLL_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\n console.error(\" REMOTE_ADMIN_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\n console.error(\" node remote-controller/dist/index.js\");\n process.exit(1);\n}\n\nconst connected = new Map();\n\nfunction loadState(): State {\n try {\n return JSON.parse(fs.readFileSync(stateFile, \"utf-8\")) as State;\n } catch {\n return { agents: {}, tasks: {} };\n }\n}\n\nlet state = loadState();\nfunction saveState(): void {\n fs.mkdirSync(path.dirname(stateFile), { recursive: true });\n fs.writeFileSync(stateFile, `${JSON.stringify(state, null, 2)}\\n`, \"utf-8\");\n}\n\nfunction json(res: ServerResponse, status: number, body: unknown): void {\n res.writeHead(status, { \"content-type\": \"application/json\" });\n res.end(JSON.stringify(body));\n}\n\nasync function body(req: IncomingMessage): Promise> {\n let raw = \"\";\n for await (const chunk of req) raw += String(chunk);\n return raw ? (JSON.parse(raw) as Record) : {};\n}\n\nfunction authorized(req: IncomingMessage, token: string): boolean {\n const value = req.headers.authorization?.replace(/^Bearer\\s+/i, \"\") ?? \"\";\n const actual = Buffer.from(value);\n const expected = Buffer.from(token);\n return actual.length === expected.length && timingSafeEqual(actual, expected);\n}\n\nfunction validAction(value: unknown): value is TaskAction {\n return value === \"status\" || value === \"start\" || value === \"stop\" || value === \"restart\" || value === \"send\";\n}\n\nfunction agentPublic(agent: AgentRecord): Omit & { connected: boolean } {\n return { id: agent.id, name: agent.name, createdAt: agent.createdAt, ...(agent.lastSeenAt ? { lastSeenAt: agent.lastSeenAt } : {}), connected: connected.has(agent.id) };\n}\n\nfunction dispatchQueuedTasks(agentId: string): void {\n const socket = connected.get(agentId);\n if (socket?.readyState !== WebSocket.OPEN) return;\n\n let changed = false;\n for (const task of Object.values(state.tasks)) {\n if (task.agentId !== agentId || task.status !== \"queued\") continue;\n task.status = \"running\";\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\n if (task.service) payload.service = task.service;\n if (task.message) payload.message = task.message;\n socket.send(JSON.stringify(payload));\n changed = true;\n }\n if (changed) saveState();\n}\n\nconst server = createServer(async (req, res) => {\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\n try {\n if (req.method === \"GET\" && url.pathname === \"/v1/health\") {\n return json(res, 200, { ok: true, agents: Object.keys(state.agents).length, connected: connected.size });\n }\n\n if (req.method === \"POST\" && url.pathname === \"/v1/enroll\") {\n if (!authorized(req, enrollmentToken)) return json(res, 401, { error: \"unauthorized\" });\n const input = await body(req);\n const id = randomUUID();\n const secret = randomBytes(32).toString(\"base64url\");\n state.agents[id] = { id, secret, name: String(input.name ?? \"sfmc-agent\"), createdAt: new Date().toISOString() };\n saveState();\n return json(res, 201, { agentId: id, agentSecret: secret });\n }\n\n if (!authorized(req, adminToken)) return json(res, 401, { error: \"unauthorized\" });\n\n if (req.method === \"GET\" && url.pathname === \"/v1/agents\") {\n const agents = Object.values(state.agents).map((a) => agentPublic(a));\n return json(res, 200, { agents });\n }\n\n const agentMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)$/);\n if (req.method === \"GET\" && agentMatch?.[1]) {\n const agent = state.agents[agentMatch[1]];\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\n return json(res, 200, agentPublic(agent));\n }\n if (req.method === \"DELETE\" && agentMatch?.[1]) {\n const agentId = agentMatch[1];\n const agent = state.agents[agentId];\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\n const socket = connected.get(agentId);\n if (socket && socket.readyState === WebSocket.OPEN) socket.close(1000, \"deleted_by_admin\");\n connected.delete(agentId);\n delete state.agents[agentId];\n for (const task of Object.values(state.tasks)) {\n if (task.agentId === agentId && (task.status === \"queued\" || task.status === \"running\")) {\n task.status = \"failed\";\n task.error = \"agent_deleted\";\n task.completedAt = new Date().toISOString();\n }\n }\n saveState();\n return json(res, 204, {});\n }\n\n const tasksMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)\\/tasks$/);\n if (req.method === \"GET\" && tasksMatch?.[1]) {\n const agentId = tasksMatch[1];\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\n const statusFilter = url.searchParams.get(\"status\");\n const limit = Math.min(Number(url.searchParams.get(\"limit\") ?? 50), 200);\n const tasks = Object.values(state.tasks)\n .filter((t) => t.agentId === agentId && (!statusFilter || t.status === statusFilter))\n .sort((a, b) => b.createdAt.localeCompare(a.createdAt))\n .slice(0, limit);\n return json(res, 200, { tasks });\n }\n\n if (req.method === \"POST\" && tasksMatch?.[1]) {\n const agentId = tasksMatch[1];\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\n const input = await body(req);\n if (!validAction(input.action)) return json(res, 400, { error: \"invalid_action\" });\n if (actionsRequiringService.has(input.action)) {\n if (typeof input.service !== \"string\") return json(res, 400, { error: \"service_required\" });\n if (!serviceNames.has(input.service)) return json(res, 400, { error: \"invalid_service\" });\n }\n if (input.action === \"send\" && (typeof input.message !== \"string\" || !input.message.length)) {\n return json(res, 400, { error: \"message_required\" });\n }\n const task: Task = {\n id: randomUUID(),\n agentId,\n action: input.action,\n ...(typeof input.service === \"string\" ? { service: input.service } : {}),\n ...(typeof input.message === \"string\" ? { message: input.message } : {}),\n status: \"queued\",\n createdAt: new Date().toISOString(),\n };\n state.tasks[task.id] = task;\n const socket = connected.get(agentId);\n if (socket?.readyState === WebSocket.OPEN) {\n task.status = \"running\";\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\n if (task.service) payload.service = task.service;\n if (task.message) payload.message = task.message;\n socket.send(JSON.stringify(payload));\n }\n saveState();\n return json(res, 202, task);\n }\n\n const getTaskMatch = url.pathname.match(/^\\/v1\\/tasks\\/([^/]+)$/);\n if (req.method === \"GET\" && getTaskMatch?.[1]) {\n const task = state.tasks[getTaskMatch[1]];\n return task ? json(res, 200, task) : json(res, 404, { error: \"task_not_found\" });\n }\n return json(res, 404, { error: \"not_found\" });\n } catch (error) {\n return json(res, 400, { error: (error as Error).message });\n }\n});\n\nconst wss = new WebSocketServer({ noServer: true });\nconst heartbeat = setInterval(() => {\n for (const [agentId, socket] of connected) {\n if (socket.readyState !== WebSocket.OPEN) continue;\n if ((socket as WebSocket & { isAlive?: boolean }).isAlive === false) {\n socket.terminate();\n connected.delete(agentId);\n continue;\n }\n (socket as WebSocket & { isAlive?: boolean }).isAlive = false;\n try {\n socket.ping();\n } catch {\n /* ignore */\n }\n }\n}, heartbeatIntervalMs);\n\nwss.on(\"close\", () => clearInterval(heartbeat));\n\nwss.on(\"connection\", (socket: WebSocket, req: IncomingMessage, agentId: string) => {\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\n socket.on(\"pong\", () => {\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\n });\n\n let authenticated = false;\n socket.once(\"message\", (raw: RawData) => {\n try {\n const message = JSON.parse(raw.toString()) as { type?: string; agentId?: string; secret?: string };\n const agent = state.agents[agentId];\n const authRequest = { headers: { authorization: `Bearer ${message.secret ?? \"\"}` } } as IncomingMessage;\n if (message.type !== \"hello\" || message.agentId !== agentId || !agent || !authorized(authRequest, agent.secret)) {\n socket.close(1008, \"unauthorized\");\n return;\n }\n authenticated = true;\n agent.lastSeenAt = new Date().toISOString();\n connected.set(agentId, socket);\n saveState();\n dispatchQueuedTasks(agentId);\n } catch {\n socket.close(1008, \"invalid_hello\");\n }\n });\n socket.on(\"message\", (raw: RawData) => {\n if (!authenticated) return;\n try {\n const message = JSON.parse(raw.toString()) as {\n type?: string;\n taskId?: string;\n ok?: boolean;\n result?: unknown;\n error?: string;\n };\n if (message.type === \"ping\") return;\n if (message.type !== \"task_result\" || !message.taskId) return;\n const task = state.tasks[message.taskId];\n if (!task || task.agentId !== agentId) return;\n task.status = message.ok ? \"complete\" : \"failed\";\n task.completedAt = new Date().toISOString();\n if (message.ok) task.result = message.result;\n else task.error = message.error ?? \"task failed\";\n saveState();\n } catch {\n /* Ignore malformed agent messages. */\n }\n });\n socket.on(\"close\", () => {\n if (connected.get(agentId) === socket) connected.delete(agentId);\n });\n void req;\n});\n\nserver.on(\"upgrade\", (req, socket, head) => {\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\n const agentId = url.pathname === \"/v1/agent\" ? url.searchParams.get(\"id\") : null;\n if (!agentId || !state.agents[agentId]) return socket.destroy();\n wss.handleUpgrade(req, socket, head, (ws) => wss.emit(\"connection\", ws, req, agentId));\n});\n\nserver.listen(port, host, () => {\n console.log(`[remote-controller] listening on http://${host}:${port}`);\n console.log(`[remote-controller] state file: ${stateFile}`);\n console.log(`[remote-controller] endpoints:`);\n console.log(` POST /v1/enroll (enroll token)`);\n console.log(` GET /v1/health (open)`);\n console.log(` GET /v1/agents (admin token)`);\n console.log(` GET /v1/agents/{id} (admin token)`);\n console.log(` DELETE /v1/agents/{id} (admin token)`);\n console.log(` POST /v1/agents/{id}/tasks (admin token)`);\n console.log(` GET /v1/agents/{id}/tasks (admin token)`);\n console.log(` GET /v1/tasks/{id} (admin token)`);\n console.log(` WS /v1/agent?id={id} (per-agent secret)`);\n console.log(`[remote-controller] heartbeat: ${heartbeatIntervalMs}ms`);\n});"], + "sourcesContent": ["import { randomBytes, randomUUID, timingSafeEqual } from \"node:crypto\";\r\nimport { createServer, type IncomingMessage, type ServerResponse } from \"node:http\";\r\nimport fs from \"node:fs\";\r\nimport path from \"node:path\";\r\nimport { WebSocket, WebSocketServer, type RawData } from \"ws\";\r\n\r\ntype AgentRecord = { id: string; name: string; secret: string; createdAt: string; lastSeenAt?: string };\r\ntype TaskAction = \"status\" | \"start\" | \"stop\" | \"restart\" | \"send\";\r\ntype Task = {\r\n id: string;\r\n agentId: string;\r\n action: TaskAction;\r\n service?: string;\r\n message?: string;\r\n status: \"queued\" | \"running\" | \"complete\" | \"failed\";\r\n result?: unknown;\r\n error?: string;\r\n createdAt: string;\r\n completedAt?: string;\r\n};\r\ntype State = { agents: Record; tasks: Record };\r\n\r\nconst serviceNames = new Set([\"bds\", \"db\", \"qq\", \"llbot\"]);\r\nconst actionsRequiringService = new Set([\"start\", \"stop\", \"restart\", \"send\"]);\r\n\r\nconst port = Number(process.env.REMOTE_PORT ?? 3100);\r\nconst host = process.env.REMOTE_HOST ?? \"127.0.0.1\";\r\nconst enrollmentToken = process.env.REMOTE_ENROLL_TOKEN ?? \"\";\r\nconst adminToken = process.env.REMOTE_ADMIN_TOKEN ?? \"\";\r\nconst stateFile = path.resolve(process.env.REMOTE_STATE_FILE ?? \"data/remote-controller.json\");\r\nconst heartbeatIntervalMs = Number(process.env.REMOTE_HEARTBEAT_MS ?? 25_000);\r\n\r\nif (!enrollmentToken || !adminToken) {\r\n console.error(\"[remote-controller] missing required env vars.\");\r\n console.error(\" set both: REMOTE_ENROLL_TOKEN= and REMOTE_ADMIN_TOKEN=\");\r\n console.error(\" example:\");\r\n console.error(\" REMOTE_ENROLL_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\r\n console.error(\" REMOTE_ADMIN_TOKEN=$(node -e \\\"console.log(require('crypto').randomBytes(24).toString('base64url'))\\\") \\\\\");\r\n console.error(\" node remote-controller/dist/index.js\");\r\n process.exit(1);\r\n}\r\n\r\nconst connected = new Map();\r\n\r\nfunction loadState(): State {\r\n try {\r\n return JSON.parse(fs.readFileSync(stateFile, \"utf-8\")) as State;\r\n } catch {\r\n return { agents: {}, tasks: {} };\r\n }\r\n}\r\n\r\nlet state = loadState();\r\nfunction saveState(): void {\r\n fs.mkdirSync(path.dirname(stateFile), { recursive: true });\r\n fs.writeFileSync(stateFile, `${JSON.stringify(state, null, 2)}\\n`, \"utf-8\");\r\n}\r\n\r\nfunction json(res: ServerResponse, status: number, body: unknown): void {\r\n res.writeHead(status, { \"content-type\": \"application/json\" });\r\n res.end(JSON.stringify(body));\r\n}\r\n\r\nasync function body(req: IncomingMessage): Promise> {\r\n let raw = \"\";\r\n for await (const chunk of req) raw += String(chunk);\r\n return raw ? (JSON.parse(raw) as Record) : {};\r\n}\r\n\r\nfunction authorized(req: IncomingMessage, token: string): boolean {\r\n const value = req.headers.authorization?.replace(/^Bearer\\s+/i, \"\") ?? \"\";\r\n const actual = Buffer.from(value);\r\n const expected = Buffer.from(token);\r\n return actual.length === expected.length && timingSafeEqual(actual, expected);\r\n}\r\n\r\nfunction validAction(value: unknown): value is TaskAction {\r\n return value === \"status\" || value === \"start\" || value === \"stop\" || value === \"restart\" || value === \"send\";\r\n}\r\n\r\nfunction agentPublic(agent: AgentRecord): Omit & { connected: boolean } {\r\n return { id: agent.id, name: agent.name, createdAt: agent.createdAt, ...(agent.lastSeenAt ? { lastSeenAt: agent.lastSeenAt } : {}), connected: connected.has(agent.id) };\r\n}\r\n\r\nfunction dispatchQueuedTasks(agentId: string): void {\r\n const socket = connected.get(agentId);\r\n if (socket?.readyState !== WebSocket.OPEN) return;\r\n\r\n let changed = false;\r\n for (const task of Object.values(state.tasks)) {\r\n if (task.agentId !== agentId || task.status !== \"queued\") continue;\r\n task.status = \"running\";\r\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\r\n if (task.service) payload.service = task.service;\r\n if (task.message) payload.message = task.message;\r\n socket.send(JSON.stringify(payload));\r\n changed = true;\r\n }\r\n if (changed) saveState();\r\n}\r\n\r\nconst server = createServer(async (req, res) => {\r\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\r\n try {\r\n if (req.method === \"GET\" && url.pathname === \"/v1/health\") {\r\n return json(res, 200, { ok: true, agents: Object.keys(state.agents).length, connected: connected.size });\r\n }\r\n\r\n if (req.method === \"POST\" && url.pathname === \"/v1/enroll\") {\r\n if (!authorized(req, enrollmentToken)) return json(res, 401, { error: \"unauthorized\" });\r\n const input = await body(req);\r\n const id = randomUUID();\r\n const secret = randomBytes(32).toString(\"base64url\");\r\n state.agents[id] = { id, secret, name: String(input.name ?? \"sfmc-agent\"), createdAt: new Date().toISOString() };\r\n saveState();\r\n return json(res, 201, { agentId: id, agentSecret: secret });\r\n }\r\n\r\n if (!authorized(req, adminToken)) return json(res, 401, { error: \"unauthorized\" });\r\n\r\n if (req.method === \"GET\" && url.pathname === \"/v1/agents\") {\r\n const agents = Object.values(state.agents).map((a) => agentPublic(a));\r\n return json(res, 200, { agents });\r\n }\r\n\r\n const agentMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)$/);\r\n if (req.method === \"GET\" && agentMatch?.[1]) {\r\n const agent = state.agents[agentMatch[1]];\r\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\r\n return json(res, 200, agentPublic(agent));\r\n }\r\n if (req.method === \"DELETE\" && agentMatch?.[1]) {\r\n const agentId = agentMatch[1];\r\n const agent = state.agents[agentId];\r\n if (!agent) return json(res, 404, { error: \"agent_not_found\" });\r\n const socket = connected.get(agentId);\r\n if (socket && socket.readyState === WebSocket.OPEN) socket.close(1000, \"deleted_by_admin\");\r\n connected.delete(agentId);\r\n delete state.agents[agentId];\r\n for (const task of Object.values(state.tasks)) {\r\n if (task.agentId === agentId && (task.status === \"queued\" || task.status === \"running\")) {\r\n task.status = \"failed\";\r\n task.error = \"agent_deleted\";\r\n task.completedAt = new Date().toISOString();\r\n }\r\n }\r\n saveState();\r\n return json(res, 204, {});\r\n }\r\n\r\n const tasksMatch = url.pathname.match(/^\\/v1\\/agents\\/([^/]+)\\/tasks$/);\r\n if (req.method === \"GET\" && tasksMatch?.[1]) {\r\n const agentId = tasksMatch[1];\r\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\r\n const statusFilter = url.searchParams.get(\"status\");\r\n const limit = Math.min(Number(url.searchParams.get(\"limit\") ?? 50), 200);\r\n const tasks = Object.values(state.tasks)\r\n .filter((t) => t.agentId === agentId && (!statusFilter || t.status === statusFilter))\r\n .sort((a, b) => b.createdAt.localeCompare(a.createdAt))\r\n .slice(0, limit);\r\n return json(res, 200, { tasks });\r\n }\r\n\r\n if (req.method === \"POST\" && tasksMatch?.[1]) {\r\n const agentId = tasksMatch[1];\r\n if (!state.agents[agentId]) return json(res, 404, { error: \"agent_not_found\" });\r\n const input = await body(req);\r\n if (!validAction(input.action)) return json(res, 400, { error: \"invalid_action\" });\r\n if (actionsRequiringService.has(input.action)) {\r\n if (typeof input.service !== \"string\") return json(res, 400, { error: \"service_required\" });\r\n if (!serviceNames.has(input.service)) return json(res, 400, { error: \"invalid_service\" });\r\n }\r\n if (input.action === \"send\" && (typeof input.message !== \"string\" || !input.message.length)) {\r\n return json(res, 400, { error: \"message_required\" });\r\n }\r\n const task: Task = {\r\n id: randomUUID(),\r\n agentId,\r\n action: input.action,\r\n ...(typeof input.service === \"string\" ? { service: input.service } : {}),\r\n ...(typeof input.message === \"string\" ? { message: input.message } : {}),\r\n status: \"queued\",\r\n createdAt: new Date().toISOString(),\r\n };\r\n state.tasks[task.id] = task;\r\n const socket = connected.get(agentId);\r\n if (socket?.readyState === WebSocket.OPEN) {\r\n task.status = \"running\";\r\n const payload: Record = { type: \"task\", taskId: task.id, action: task.action };\r\n if (task.service) payload.service = task.service;\r\n if (task.message) payload.message = task.message;\r\n socket.send(JSON.stringify(payload));\r\n }\r\n saveState();\r\n return json(res, 202, task);\r\n }\r\n\r\n const getTaskMatch = url.pathname.match(/^\\/v1\\/tasks\\/([^/]+)$/);\r\n if (req.method === \"GET\" && getTaskMatch?.[1]) {\r\n const task = state.tasks[getTaskMatch[1]];\r\n return task ? json(res, 200, task) : json(res, 404, { error: \"task_not_found\" });\r\n }\r\n return json(res, 404, { error: \"not_found\" });\r\n } catch (error) {\r\n return json(res, 400, { error: (error as Error).message });\r\n }\r\n});\r\n\r\nconst wss = new WebSocketServer({ noServer: true });\r\nconst heartbeat = setInterval(() => {\r\n for (const [agentId, socket] of connected) {\r\n if (socket.readyState !== WebSocket.OPEN) continue;\r\n if ((socket as WebSocket & { isAlive?: boolean }).isAlive === false) {\r\n socket.terminate();\r\n connected.delete(agentId);\r\n continue;\r\n }\r\n (socket as WebSocket & { isAlive?: boolean }).isAlive = false;\r\n try {\r\n socket.ping();\r\n } catch {\r\n /* ignore */\r\n }\r\n }\r\n}, heartbeatIntervalMs);\r\n\r\nwss.on(\"close\", () => clearInterval(heartbeat));\r\n\r\nwss.on(\"connection\", (socket: WebSocket, req: IncomingMessage, agentId: string) => {\r\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\r\n socket.on(\"pong\", () => {\r\n (socket as WebSocket & { isAlive?: boolean }).isAlive = true;\r\n });\r\n\r\n let authenticated = false;\r\n socket.once(\"message\", (raw: RawData) => {\r\n try {\r\n const message = JSON.parse(raw.toString()) as { type?: string; agentId?: string; secret?: string };\r\n const agent = state.agents[agentId];\r\n const authRequest = { headers: { authorization: `Bearer ${message.secret ?? \"\"}` } } as IncomingMessage;\r\n if (message.type !== \"hello\" || message.agentId !== agentId || !agent || !authorized(authRequest, agent.secret)) {\r\n socket.close(1008, \"unauthorized\");\r\n return;\r\n }\r\n authenticated = true;\r\n agent.lastSeenAt = new Date().toISOString();\r\n connected.set(agentId, socket);\r\n saveState();\r\n dispatchQueuedTasks(agentId);\r\n } catch {\r\n socket.close(1008, \"invalid_hello\");\r\n }\r\n });\r\n socket.on(\"message\", (raw: RawData) => {\r\n if (!authenticated) return;\r\n try {\r\n const message = JSON.parse(raw.toString()) as {\r\n type?: string;\r\n taskId?: string;\r\n ok?: boolean;\r\n result?: unknown;\r\n error?: string;\r\n };\r\n if (message.type === \"ping\") return;\r\n if (message.type !== \"task_result\" || !message.taskId) return;\r\n const task = state.tasks[message.taskId];\r\n if (!task || task.agentId !== agentId) return;\r\n task.status = message.ok ? \"complete\" : \"failed\";\r\n task.completedAt = new Date().toISOString();\r\n if (message.ok) task.result = message.result;\r\n else task.error = message.error ?? \"task failed\";\r\n saveState();\r\n } catch {\r\n /* Ignore malformed agent messages. */\r\n }\r\n });\r\n socket.on(\"close\", () => {\r\n if (connected.get(agentId) === socket) connected.delete(agentId);\r\n });\r\n void req;\r\n});\r\n\r\nserver.on(\"upgrade\", (req, socket, head) => {\r\n const url = new URL(req.url ?? \"/\", `http://${req.headers.host ?? \"localhost\"}`);\r\n const agentId = url.pathname === \"/v1/agent\" ? url.searchParams.get(\"id\") : null;\r\n if (!agentId || !state.agents[agentId]) return socket.destroy();\r\n wss.handleUpgrade(req, socket, head, (ws) => wss.emit(\"connection\", ws, req, agentId));\r\n});\r\n\r\nserver.listen(port, host, () => {\r\n console.log(`[remote-controller] listening on http://${host}:${port}`);\r\n console.log(`[remote-controller] state file: ${stateFile}`);\r\n console.log(`[remote-controller] endpoints:`);\r\n console.log(` POST /v1/enroll (enroll token)`);\r\n console.log(` GET /v1/health (open)`);\r\n console.log(` GET /v1/agents (admin token)`);\r\n console.log(` GET /v1/agents/{id} (admin token)`);\r\n console.log(` DELETE /v1/agents/{id} (admin token)`);\r\n console.log(` POST /v1/agents/{id}/tasks (admin token)`);\r\n console.log(` GET /v1/agents/{id}/tasks (admin token)`);\r\n console.log(` GET /v1/tasks/{id} (admin token)`);\r\n console.log(` WS /v1/agent?id={id} (per-agent secret)`);\r\n console.log(`[remote-controller] heartbeat: ${heartbeatIntervalMs}ms`);\r\n});"], "mappings": "AAAA,SAAS,aAAa,YAAY,uBAAuB;AACzD,SAAS,oBAA+D;AACxE,OAAO,QAAQ;AACf,OAAO,UAAU;AACjB,SAAS,WAAW,uBAAqC;AAkBzD,MAAM,eAAe,oBAAI,IAAI,CAAC,OAAO,MAAM,MAAM,OAAO,CAAC;AACzD,MAAM,0BAA0B,oBAAI,IAAgB,CAAC,SAAS,QAAQ,WAAW,MAAM,CAAC;AAExF,MAAM,OAAO,OAAO,QAAQ,IAAI,eAAe,IAAI;AACnD,MAAM,OAAO,QAAQ,IAAI,eAAe;AACxC,MAAM,kBAAkB,QAAQ,IAAI,uBAAuB;AAC3D,MAAM,aAAa,QAAQ,IAAI,sBAAsB;AACrD,MAAM,YAAY,KAAK,QAAQ,QAAQ,IAAI,qBAAqB,6BAA6B;AAC7F,MAAM,sBAAsB,OAAO,QAAQ,IAAI,uBAAuB,IAAM;AAE5E,IAAI,CAAC,mBAAmB,CAAC,YAAY;AACnC,UAAQ,MAAM,gDAAgD;AAC9D,UAAQ,MAAM,0EAA0E;AACxF,UAAQ,MAAM,YAAY;AAC1B,UAAQ,MAAM,8GAAgH;AAC9H,UAAQ,MAAM,6GAA+G;AAC7H,UAAQ,MAAM,0CAA0C;AACxD,UAAQ,KAAK,CAAC;AAChB;AAEA,MAAM,YAAY,oBAAI,IAAuB;AAE7C,SAAS,YAAmB;AAC1B,MAAI;AACF,WAAO,KAAK,MAAM,GAAG,aAAa,WAAW,OAAO,CAAC;AAAA,EACvD,QAAQ;AACN,WAAO,EAAE,QAAQ,CAAC,GAAG,OAAO,CAAC,EAAE;AAAA,EACjC;AACF;AAEA,IAAI,QAAQ,UAAU;AACtB,SAAS,YAAkB;AACzB,KAAG,UAAU,KAAK,QAAQ,SAAS,GAAG,EAAE,WAAW,KAAK,CAAC;AACzD,KAAG,cAAc,WAAW,GAAG,KAAK,UAAU,OAAO,MAAM,CAAC,CAAC;AAAA,GAAM,OAAO;AAC5E;AAEA,SAAS,KAAK,KAAqB,QAAgBA,OAAqB;AACtE,MAAI,UAAU,QAAQ,EAAE,gBAAgB,mBAAmB,CAAC;AAC5D,MAAI,IAAI,KAAK,UAAUA,KAAI,CAAC;AAC9B;AAEA,eAAe,KAAK,KAAwD;AAC1E,MAAI,MAAM;AACV,mBAAiB,SAAS,IAAK,QAAO,OAAO,KAAK;AAClD,SAAO,MAAO,KAAK,MAAM,GAAG,IAAgC,CAAC;AAC/D;AAEA,SAAS,WAAW,KAAsB,OAAwB;AAChE,QAAM,QAAQ,IAAI,QAAQ,eAAe,QAAQ,eAAe,EAAE,KAAK;AACvE,QAAM,SAAS,OAAO,KAAK,KAAK;AAChC,QAAM,WAAW,OAAO,KAAK,KAAK;AAClC,SAAO,OAAO,WAAW,SAAS,UAAU,gBAAgB,QAAQ,QAAQ;AAC9E;AAEA,SAAS,YAAY,OAAqC;AACxD,SAAO,UAAU,YAAY,UAAU,WAAW,UAAU,UAAU,UAAU,aAAa,UAAU;AACzG;AAEA,SAAS,YAAY,OAA0E;AAC7F,SAAO,EAAE,IAAI,MAAM,IAAI,MAAM,MAAM,MAAM,WAAW,MAAM,WAAW,GAAI,MAAM,aAAa,EAAE,YAAY,MAAM,WAAW,IAAI,CAAC,GAAI,WAAW,UAAU,IAAI,MAAM,EAAE,EAAE;AACzK;AAEA,SAAS,oBAAoB,SAAuB;AAClD,QAAM,SAAS,UAAU,IAAI,OAAO;AACpC,MAAI,QAAQ,eAAe,UAAU,KAAM;AAE3C,MAAI,UAAU;AACd,aAAW,QAAQ,OAAO,OAAO,MAAM,KAAK,GAAG;AAC7C,QAAI,KAAK,YAAY,WAAW,KAAK,WAAW,SAAU;AAC1D,SAAK,SAAS;AACd,UAAM,UAAmC,EAAE,MAAM,QAAQ,QAAQ,KAAK,IAAI,QAAQ,KAAK,OAAO;AAC9F,QAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,QAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,WAAO,KAAK,KAAK,UAAU,OAAO,CAAC;AACnC,cAAU;AAAA,EACZ;AACA,MAAI,QAAS,WAAU;AACzB;AAEA,MAAM,SAAS,aAAa,OAAO,KAAK,QAAQ;AAC9C,QAAM,MAAM,IAAI,IAAI,IAAI,OAAO,KAAK,UAAU,IAAI,QAAQ,QAAQ,WAAW,EAAE;AAC/E,MAAI;AACF,QAAI,IAAI,WAAW,SAAS,IAAI,aAAa,cAAc;AACzD,aAAO,KAAK,KAAK,KAAK,EAAE,IAAI,MAAM,QAAQ,OAAO,KAAK,MAAM,MAAM,EAAE,QAAQ,WAAW,UAAU,KAAK,CAAC;AAAA,IACzG;AAEA,QAAI,IAAI,WAAW,UAAU,IAAI,aAAa,cAAc;AAC1D,UAAI,CAAC,WAAW,KAAK,eAAe,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,eAAe,CAAC;AACtF,YAAM,QAAQ,MAAM,KAAK,GAAG;AAC5B,YAAM,KAAK,WAAW;AACtB,YAAM,SAAS,YAAY,EAAE,EAAE,SAAS,WAAW;AACnD,YAAM,OAAO,EAAE,IAAI,EAAE,IAAI,QAAQ,MAAM,OAAO,MAAM,QAAQ,YAAY,GAAG,YAAW,oBAAI,KAAK,GAAE,YAAY,EAAE;AAC/G,gBAAU;AACV,aAAO,KAAK,KAAK,KAAK,EAAE,SAAS,IAAI,aAAa,OAAO,CAAC;AAAA,IAC5D;AAEA,QAAI,CAAC,WAAW,KAAK,UAAU,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,eAAe,CAAC;AAEjF,QAAI,IAAI,WAAW,SAAS,IAAI,aAAa,cAAc;AACzD,YAAM,SAAS,OAAO,OAAO,MAAM,MAAM,EAAE,IAAI,CAAC,MAAM,YAAY,CAAC,CAAC;AACpE,aAAO,KAAK,KAAK,KAAK,EAAE,OAAO,CAAC;AAAA,IAClC;AAEA,UAAM,aAAa,IAAI,SAAS,MAAM,yBAAyB;AAC/D,QAAI,IAAI,WAAW,SAAS,aAAa,CAAC,GAAG;AAC3C,YAAM,QAAQ,MAAM,OAAO,WAAW,CAAC,CAAC;AACxC,UAAI,CAAC,MAAO,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9D,aAAO,KAAK,KAAK,KAAK,YAAY,KAAK,CAAC;AAAA,IAC1C;AACA,QAAI,IAAI,WAAW,YAAY,aAAa,CAAC,GAAG;AAC9C,YAAM,UAAU,WAAW,CAAC;AAC5B,YAAM,QAAQ,MAAM,OAAO,OAAO;AAClC,UAAI,CAAC,MAAO,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9D,YAAM,SAAS,UAAU,IAAI,OAAO;AACpC,UAAI,UAAU,OAAO,eAAe,UAAU,KAAM,QAAO,MAAM,KAAM,kBAAkB;AACzF,gBAAU,OAAO,OAAO;AACxB,aAAO,MAAM,OAAO,OAAO;AAC3B,iBAAW,QAAQ,OAAO,OAAO,MAAM,KAAK,GAAG;AAC7C,YAAI,KAAK,YAAY,YAAY,KAAK,WAAW,YAAY,KAAK,WAAW,YAAY;AACvF,eAAK,SAAS;AACd,eAAK,QAAQ;AACb,eAAK,eAAc,oBAAI,KAAK,GAAE,YAAY;AAAA,QAC5C;AAAA,MACF;AACA,gBAAU;AACV,aAAO,KAAK,KAAK,KAAK,CAAC,CAAC;AAAA,IAC1B;AAEA,UAAM,aAAa,IAAI,SAAS,MAAM,gCAAgC;AACtE,QAAI,IAAI,WAAW,SAAS,aAAa,CAAC,GAAG;AAC3C,YAAM,UAAU,WAAW,CAAC;AAC5B,UAAI,CAAC,MAAM,OAAO,OAAO,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9E,YAAM,eAAe,IAAI,aAAa,IAAI,QAAQ;AAClD,YAAM,QAAQ,KAAK,IAAI,OAAO,IAAI,aAAa,IAAI,OAAO,KAAK,EAAE,GAAG,GAAG;AACvE,YAAM,QAAQ,OAAO,OAAO,MAAM,KAAK,EACpC,OAAO,CAAC,MAAM,EAAE,YAAY,YAAY,CAAC,gBAAgB,EAAE,WAAW,aAAa,EACnF,KAAK,CAAC,GAAG,MAAM,EAAE,UAAU,cAAc,EAAE,SAAS,CAAC,EACrD,MAAM,GAAG,KAAK;AACjB,aAAO,KAAK,KAAK,KAAK,EAAE,MAAM,CAAC;AAAA,IACjC;AAEA,QAAI,IAAI,WAAW,UAAU,aAAa,CAAC,GAAG;AAC5C,YAAM,UAAU,WAAW,CAAC;AAC5B,UAAI,CAAC,MAAM,OAAO,OAAO,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAC9E,YAAM,QAAQ,MAAM,KAAK,GAAG;AAC5B,UAAI,CAAC,YAAY,MAAM,MAAM,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,iBAAiB,CAAC;AACjF,UAAI,wBAAwB,IAAI,MAAM,MAAM,GAAG;AAC7C,YAAI,OAAO,MAAM,YAAY,SAAU,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,mBAAmB,CAAC;AAC1F,YAAI,CAAC,aAAa,IAAI,MAAM,OAAO,EAAG,QAAO,KAAK,KAAK,KAAK,EAAE,OAAO,kBAAkB,CAAC;AAAA,MAC1F;AACA,UAAI,MAAM,WAAW,WAAW,OAAO,MAAM,YAAY,YAAY,CAAC,MAAM,QAAQ,SAAS;AAC3F,eAAO,KAAK,KAAK,KAAK,EAAE,OAAO,mBAAmB,CAAC;AAAA,MACrD;AACA,YAAM,OAAa;AAAA,QACjB,IAAI,WAAW;AAAA,QACf;AAAA,QACA,QAAQ,MAAM;AAAA,QACd,GAAI,OAAO,MAAM,YAAY,WAAW,EAAE,SAAS,MAAM,QAAQ,IAAI,CAAC;AAAA,QACtE,GAAI,OAAO,MAAM,YAAY,WAAW,EAAE,SAAS,MAAM,QAAQ,IAAI,CAAC;AAAA,QACtE,QAAQ;AAAA,QACR,YAAW,oBAAI,KAAK,GAAE,YAAY;AAAA,MACpC;AACA,YAAM,MAAM,KAAK,EAAE,IAAI;AACvB,YAAM,SAAS,UAAU,IAAI,OAAO;AACpC,UAAI,QAAQ,eAAe,UAAU,MAAM;AACzC,aAAK,SAAS;AACd,cAAM,UAAmC,EAAE,MAAM,QAAQ,QAAQ,KAAK,IAAI,QAAQ,KAAK,OAAO;AAC9F,YAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,YAAI,KAAK,QAAS,SAAQ,UAAU,KAAK;AACzC,eAAO,KAAK,KAAK,UAAU,OAAO,CAAC;AAAA,MACrC;AACA,gBAAU;AACV,aAAO,KAAK,KAAK,KAAK,IAAI;AAAA,IAC5B;AAEA,UAAM,eAAe,IAAI,SAAS,MAAM,wBAAwB;AAChE,QAAI,IAAI,WAAW,SAAS,eAAe,CAAC,GAAG;AAC7C,YAAM,OAAO,MAAM,MAAM,aAAa,CAAC,CAAC;AACxC,aAAO,OAAO,KAAK,KAAK,KAAK,IAAI,IAAI,KAAK,KAAK,KAAK,EAAE,OAAO,iBAAiB,CAAC;AAAA,IACjF;AACA,WAAO,KAAK,KAAK,KAAK,EAAE,OAAO,YAAY,CAAC;AAAA,EAC9C,SAAS,OAAO;AACd,WAAO,KAAK,KAAK,KAAK,EAAE,OAAQ,MAAgB,QAAQ,CAAC;AAAA,EAC3D;AACF,CAAC;AAED,MAAM,MAAM,IAAI,gBAAgB,EAAE,UAAU,KAAK,CAAC;AAClD,MAAM,YAAY,YAAY,MAAM;AAClC,aAAW,CAAC,SAAS,MAAM,KAAK,WAAW;AACzC,QAAI,OAAO,eAAe,UAAU,KAAM;AAC1C,QAAK,OAA6C,YAAY,OAAO;AACnE,aAAO,UAAU;AACjB,gBAAU,OAAO,OAAO;AACxB;AAAA,IACF;AACA,IAAC,OAA6C,UAAU;AACxD,QAAI;AACF,aAAO,KAAK;AAAA,IACd,QAAQ;AAAA,IAER;AAAA,EACF;AACF,GAAG,mBAAmB;AAEtB,IAAI,GAAG,SAAS,MAAM,cAAc,SAAS,CAAC;AAE9C,IAAI,GAAG,cAAc,CAAC,QAAmB,KAAsB,YAAoB;AACjF,EAAC,OAA6C,UAAU;AACxD,SAAO,GAAG,QAAQ,MAAM;AACtB,IAAC,OAA6C,UAAU;AAAA,EAC1D,CAAC;AAED,MAAI,gBAAgB;AACpB,SAAO,KAAK,WAAW,CAAC,QAAiB;AACvC,QAAI;AACF,YAAM,UAAU,KAAK,MAAM,IAAI,SAAS,CAAC;AACzC,YAAM,QAAQ,MAAM,OAAO,OAAO;AAClC,YAAM,cAAc,EAAE,SAAS,EAAE,eAAe,UAAU,QAAQ,UAAU,EAAE,GAAG,EAAE;AACnF,UAAI,QAAQ,SAAS,WAAW,QAAQ,YAAY,WAAW,CAAC,SAAS,CAAC,WAAW,aAAa,MAAM,MAAM,GAAG;AAC/G,eAAO,MAAM,MAAM,cAAc;AACjC;AAAA,MACF;AACA,sBAAgB;AAChB,YAAM,cAAa,oBAAI,KAAK,GAAE,YAAY;AAC1C,gBAAU,IAAI,SAAS,MAAM;AAC7B,gBAAU;AACV,0BAAoB,OAAO;AAAA,IAC7B,QAAQ;AACN,aAAO,MAAM,MAAM,eAAe;AAAA,IACpC;AAAA,EACF,CAAC;AACD,SAAO,GAAG,WAAW,CAAC,QAAiB;AACrC,QAAI,CAAC,cAAe;AACpB,QAAI;AACF,YAAM,UAAU,KAAK,MAAM,IAAI,SAAS,CAAC;AAOzC,UAAI,QAAQ,SAAS,OAAQ;AAC7B,UAAI,QAAQ,SAAS,iBAAiB,CAAC,QAAQ,OAAQ;AACvD,YAAM,OAAO,MAAM,MAAM,QAAQ,MAAM;AACvC,UAAI,CAAC,QAAQ,KAAK,YAAY,QAAS;AACvC,WAAK,SAAS,QAAQ,KAAK,aAAa;AACxC,WAAK,eAAc,oBAAI,KAAK,GAAE,YAAY;AAC1C,UAAI,QAAQ,GAAI,MAAK,SAAS,QAAQ;AAAA,UACjC,MAAK,QAAQ,QAAQ,SAAS;AACnC,gBAAU;AAAA,IACZ,QAAQ;AAAA,IAER;AAAA,EACF,CAAC;AACD,SAAO,GAAG,SAAS,MAAM;AACvB,QAAI,UAAU,IAAI,OAAO,MAAM,OAAQ,WAAU,OAAO,OAAO;AAAA,EACjE,CAAC;AACD,OAAK;AACP,CAAC;AAED,OAAO,GAAG,WAAW,CAAC,KAAK,QAAQ,SAAS;AAC1C,QAAM,MAAM,IAAI,IAAI,IAAI,OAAO,KAAK,UAAU,IAAI,QAAQ,QAAQ,WAAW,EAAE;AAC/E,QAAM,UAAU,IAAI,aAAa,cAAc,IAAI,aAAa,IAAI,IAAI,IAAI;AAC5E,MAAI,CAAC,WAAW,CAAC,MAAM,OAAO,OAAO,EAAG,QAAO,OAAO,QAAQ;AAC9D,MAAI,cAAc,KAAK,QAAQ,MAAM,CAAC,OAAO,IAAI,KAAK,cAAc,IAAI,KAAK,OAAO,CAAC;AACvF,CAAC;AAED,OAAO,OAAO,MAAM,MAAM,MAAM;AAC9B,UAAQ,IAAI,2CAA2C,IAAI,IAAI,IAAI,EAAE;AACrE,UAAQ,IAAI,mCAAmC,SAAS,EAAE;AAC1D,UAAQ,IAAI,gCAAgC;AAC5C,UAAQ,IAAI,mDAAmD;AAC/D,UAAQ,IAAI,2CAA2C;AACvD,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,kDAAkD;AAC9D,UAAQ,IAAI,uDAAuD;AACnE,UAAQ,IAAI,kCAAkC,mBAAmB,IAAI;AACvE,CAAC;", "names": ["body"] }