feat(notifications): ntfy/Slack/Discord/generic webhook for the push events

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
Devvyn
2026-10-02 22:06:29 +08:00
co-authored by Claude Sonnet 5.5
parent 9240493c43
commit 0ff57ce304
11 changed files with 1257 additions and 22 deletions
+47
View File
@@ -2601,6 +2601,53 @@
</div>
</div>
</div>
<div class="set-group" id="webhookGroup" style="display:none">
<div class="set-group-head"><h4>Webhook (ntfy, Slack, Discord)</h4><span class="set-scope">server</span></div>
<div class="set-group-body">
<div class="set-row" data-search="webhook ntfy slack discord notification phone headless">
<div class="set-row-text">
<span class="set-row-label">Send alerts to a webhook</span>
<span class="set-row-desc">Posts the same events as push notifications (permission prompts, questions, errors, idle) to ntfy, Slack, Discord or any URL, so a server with no browser open can still reach your phone. The URL is a secret: it is stored on the server only and is never shown again once saved.</span>
</div>
<label class="switch switch-sm"><input type="checkbox" id="webhookEnabled"><span class="slider"></span></label>
</div>
<div class="set-row has-field">
<div class="set-row-text"><span class="set-row-label">Service</span></div>
<select id="webhookKind" class="set-select">
<option value="ntfy">ntfy</option>
<option value="slack">Slack</option>
<option value="discord">Discord</option>
<option value="generic">Generic JSON</option>
</select>
</div>
<div class="set-row has-field">
<div class="set-row-text">
<span class="set-row-label">Webhook URL</span>
<span class="set-row-desc" id="webhookUrlHint">Nothing saved yet.</span>
</div>
<input type="password" id="webhookUrl" class="set-select" autocomplete="off" spellcheck="false" placeholder="https://ntfy.sh/your-topic">
</div>
<div class="set-row has-field">
<div class="set-row-text">
<span class="set-row-label">Which events</span>
<span class="set-row-desc">"Needs attention" skips the routine "response complete" message.</span>
</div>
<select id="webhookScope" class="set-select">
<option value="attention">Needs attention</option>
<option value="all">Everything</option>
</select>
</div>
<div class="set-row">
<div class="set-row-text"><span class="set-row-label">Save and test</span></div>
<span>
<button class="btn-toolbar btn-sm btn-primary" id="webhookSaveBtn" onclick="app.saveWebhook()">Save</button>
<button class="btn-toolbar btn-sm" id="webhookTestBtn" onclick="app.testWebhook()">Send test</button>
</span>
</div>
<div id="webhookResult" class="set-note" style="display:none" data-i18n-skip></div>
</div>
</div>
</section>
<!-- ══ Voice ════════════════════════════════════════════════════ -->
+83
View File
@@ -416,6 +416,7 @@ Object.assign(CodemanApp.prototype, {
// .checked fires no onchange, so the list's visibility (and lazy load)
// needs an explicit sync on every open, not just a save.
this.applyCliManagementVisibility();
this.loadWebhook();
// Read My Mind: synced, default OFF (opt-in; capture + prediction cost real tokens).
document.getElementById('appSettingsReadMyMind').checked = settings.readMyMindEnabled === true;
document.getElementById('appSettingsUltracodeFloatingWindows').checked =
@@ -1114,6 +1115,88 @@ Object.assign(CodemanApp.prototype, {
this._updateCheck = null;
},
/**
* Webhook notifications (Settings → Notifications). Server-side config behind /api/webhook, not a
* settings-payload field: the URL is a secret, so it never round-trips through settings.json or
* this page. The URL box is write-only; the status line shows scheme + host only.
*/
_webhookSay(text, bad = false) {
const out = document.getElementById('webhookResult');
if (!out) return;
out.textContent = text;
out.style.display = text ? 'block' : 'none';
out.style.color = bad ? 'var(--danger, #e5534b)' : '';
},
async loadWebhook() {
const group = document.getElementById('webhookGroup');
if (!group) return;
const res = await this._api('/api/webhook');
if (!res || !res.ok) {
group.style.display = 'none'; // not an admin in multi-user mode, or the server predates the route
return;
}
let body = null;
try { body = await res.json(); } catch { /* leave hidden */ }
if (!body || body.success === false) { group.style.display = 'none'; return; }
const d = body.data;
group.style.display = '';
document.getElementById('webhookEnabled').checked = d.enabled === true;
document.getElementById('webhookKind').value = d.kind;
document.getElementById('webhookScope').value = d.scope;
const url = document.getElementById('webhookUrl');
url.value = '';
url.placeholder = d.hasUrl ? 'Saved. Paste a new URL to replace it' : 'https://ntfy.sh/your-topic';
document.getElementById('webhookUrlHint').textContent = d.hasUrl ? `Saved: ${d.urlMasked}` : 'Nothing saved yet.';
if (d.lastResult) {
const when = new Date(d.lastResult.at).toLocaleString();
this._webhookSay(
d.lastResult.ok ? `Last delivery succeeded (${when}).` : `Last delivery failed (${when}): ${d.lastResult.error}`,
!d.lastResult.ok
);
} else {
this._webhookSay('');
}
},
async saveWebhook() {
const payload = {
enabled: document.getElementById('webhookEnabled').checked,
kind: document.getElementById('webhookKind').value,
scope: document.getElementById('webhookScope').value,
};
const url = document.getElementById('webhookUrl').value.trim();
if (url) payload.url = url; // blank = keep the saved one
const res = await this._api('/api/webhook', { method: 'PUT', body: payload });
let body = null;
try { body = res ? await res.json() : null; } catch { /* fall through */ }
if (!res || !res.ok || !body || body.success === false) {
this._webhookSay(body?.error || 'Could not save the webhook.', true);
return;
}
await this.loadWebhook();
this._webhookSay('Saved.');
},
async testWebhook() {
const btn = document.getElementById('webhookTestBtn');
if (btn) btn.disabled = true;
this._webhookSay('Sending…');
try {
const res = await this._apiPost('/api/webhook/test', {});
let body = null;
try { body = res ? await res.json() : null; } catch { /* fall through */ }
if (!res || !res.ok || !body || body.success === false) {
this._webhookSay(body?.error || 'Could not send the test.', true);
return;
}
const r = body.data;
this._webhookSay(r.ok ? 'Test sent. Check your phone or channel.' : `Delivery failed: ${r.error}`, !r.ok);
} finally {
if (btn) btn.disabled = false;
}
},
_setUpdateResult(html) {
const el = this.$('updateResult');
if (el) { el.style.display = 'block'; el.innerHTML = html; }
+1
View File
@@ -28,6 +28,7 @@ export { registerWsRoutes } from './ws-routes.js';
export { registerVoiceRoutes } from './voice-routes.js';
export { registerWebviewRoutes, tryWebviewRefererFallback } from './webview-routes.js';
export { registerTabLayoutRoutes } from './tab-layout-routes.js';
export { registerWebhookRoutes } from './webhook-routes.js';
export {
registerCustomModelRoutes,
refreshAllCustomModelHosts,
+115
View File
@@ -0,0 +1,115 @@
/**
* @fileoverview Webhook notification settings (src/webhook-notify.ts).
*
* GET /api/webhook — the config WITHOUT its URL (scheme + host only), and the last delivery result
* PUT /api/webhook — change enabled / kind / url / scope; an empty `url` clears it
* POST /api/webhook/test — send one test message with the saved config
*
* The URL is a bearer secret (anyone holding a Slack/Discord webhook URL can post as it), so it is
* stored in its own 0600 file and never returned. In multi-user mode all three routes are admin only:
* the channel receives every session's events, the same reach an admin's own Web Push has.
*/
import type { FastifyInstance, FastifyReply, FastifyRequest } from 'fastify';
import { ApiErrorCode, createErrorResponse, getErrorMessage, type ApiResponse } from '../../types.js';
import { isAdmin, parseBody } from '../route-helpers.js';
import { isMultiUserMode } from '../../config/multiuser.js';
import { WebhookUpdateSchema } from '../schemas.js';
import {
maskWebhookUrl,
readWebhookConfig,
webhookUrlProblem,
writeWebhookConfig,
type WebhookKind,
type WebhookNotifier,
type WebhookResult,
type WebhookScope,
} from '../../webhook-notify.js';
export interface WebhookStatus {
enabled: boolean;
kind: WebhookKind;
scope: WebhookScope;
hasUrl: boolean;
/** Scheme + host only; the path and query are the secret. */
urlMasked: string;
lastResult: WebhookResult | null;
}
export interface WebhookRouteDeps {
notifier: WebhookNotifier;
configDir: string;
/** The instance's window title, so a test message says which machine sent it. */
hostTitle: () => string;
}
export function registerWebhookRoutes(app: FastifyInstance, deps: WebhookRouteDeps): void {
const denied = (req: FastifyRequest, reply: FastifyReply): ApiResponse<never> | null => {
if (isMultiUserMode() && !isAdmin(req)) {
reply.code(403);
return createErrorResponse(ApiErrorCode.FORBIDDEN, 'Admin only in multi-user mode');
}
return null;
};
const status = async (): Promise<WebhookStatus> => {
const cfg = await readWebhookConfig(deps.configDir);
return {
enabled: cfg.enabled,
kind: cfg.kind,
scope: cfg.scope,
hasUrl: cfg.url !== '',
urlMasked: maskWebhookUrl(cfg.url),
lastResult: deps.notifier.lastResult,
};
};
app.get('/api/webhook', async (req, reply): Promise<ApiResponse<WebhookStatus>> => {
const no = denied(req, reply);
if (no) return no;
return { success: true, data: await status() };
});
app.put('/api/webhook', async (req, reply): Promise<ApiResponse<WebhookStatus>> => {
const no = denied(req, reply);
if (no) return no;
const patch = parseBody(WebhookUpdateSchema, req.body, 'Invalid webhook settings');
const current = await readWebhookConfig(deps.configDir);
const next = {
enabled: patch.enabled ?? current.enabled,
kind: patch.kind ?? current.kind,
scope: patch.scope ?? current.scope,
url: patch.url !== undefined ? patch.url.trim() : current.url,
};
if (next.url) {
const problem = webhookUrlProblem(next.url);
if (problem) {
reply.code(400);
return createErrorResponse(ApiErrorCode.INVALID_INPUT, problem);
}
}
if (next.enabled && !next.url) {
reply.code(400);
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Add a webhook URL before enabling notifications');
}
try {
await writeWebhookConfig(deps.configDir, next);
} catch (err) {
reply.code(500);
return createErrorResponse(ApiErrorCode.OPERATION_FAILED, getErrorMessage(err));
}
return { success: true, data: await status() };
});
app.post('/api/webhook/test', async (req, reply): Promise<ApiResponse<WebhookResult>> => {
const no = denied(req, reply);
if (no) return no;
const cfg = await readWebhookConfig(deps.configDir);
if (!cfg.url) {
reply.code(400);
return createErrorResponse(ApiErrorCode.INVALID_INPUT, 'Save a webhook URL first');
}
// 200 even when delivery failed: the request to Codeman worked, `data.ok` says whether the webhook did.
return { success: true, data: await deps.notifier.sendTest(cfg, deps.hostTitle()) };
});
}
+13
View File
@@ -1827,6 +1827,19 @@ export const RespawnEnableSchema = z.object({
// ========== Web Push ==========
/** POST /api/push/subscribe */
/**
* PUT /api/webhook. `.strict()` like every settings-shaped schema; `url` is optional so a change of
* kind or scope never needs the secret re-sent, and an empty string clears it.
*/
export const WebhookUpdateSchema = z
.object({
enabled: z.boolean().optional(),
kind: z.enum(['ntfy', 'slack', 'discord', 'generic']).optional(),
scope: z.enum(['attention', 'all']).optional(),
url: z.string().max(2048).optional(),
})
.strict();
export const PushSubscribeSchema = z.object({
endpoint: z
.string()
+51 -22
View File
@@ -44,6 +44,8 @@ import { hostname as getHostname, uptime as osUptime } from 'node:os';
import { looksLikeHostReboot, newestPersistedActivity, planRebootRestore } from '../reboot-restore.js';
import { rebootRestoreRegistry } from './reboot-restore-registry.js';
import { dataPath, getDataDir, CODEMAN_INSTANCE } from '../config/instance.js';
import { WebhookNotifier, readWebhookConfig, type WebhookUrgency } from '../webhook-notify.js';
import { webviewFetch } from './webview-egress.js';
import { readRemoteHosts, rehydrateRemoteHostFields } from '../remote-hosts.js';
import type { RemoteWakeRegistry } from '../remote-wake.js';
import { normalizeBasePath, stripBasePath, joinBasePath } from '../config/base-path.js';
@@ -197,6 +199,7 @@ import {
registerVoiceRoutes,
registerWebviewRoutes,
registerTabLayoutRoutes,
registerWebhookRoutes,
registerCustomModelRoutes,
refreshAllCustomModelHosts,
readCustomModelEndpointsEnabled,
@@ -368,6 +371,8 @@ export class WebServer extends EventEmitter {
private hookSecretFailures: StaleExpirationMap<string, number> | null = null;
private userFailures: StaleExpirationMap<string, number> | null = null;
private pushStore: PushSubscriptionStore = new PushSubscriptionStore();
/** ntfy / Slack / Discord / generic webhook for the push events; config in webhook.json (0600). */
private webhookNotifier = new WebhookNotifier(() => readWebhookConfig(getDataDir()), webviewFetch);
private teamWatcher: TeamWatcher = new TeamWatcher();
private _orchestratorLoop: import('../orchestrator-loop.js').OrchestratorLoop | null = null;
private readonly titleHostname: string;
@@ -1130,6 +1135,11 @@ export class WebServer extends EventEmitter {
registerOrchestratorRoutes(this.app, ctx);
registerWebviewRoutes(this.app, ctx, this.basePath);
registerTabLayoutRoutes(this.app, ctx);
registerWebhookRoutes(this.app, {
notifier: this.webhookNotifier,
configDir: getDataDir(),
hostTitle: () => this.windowTitle,
});
registerCustomModelRoutes(this.app);
registerCliRegistryRoutes(this.app);
@@ -2652,6 +2662,25 @@ export class WebServer extends EventEmitter {
const template = WebServer.PUSH_EVENT_MAP[event];
if (!template) return;
const sessionName = (data.sessionName as string) || '';
const sessionId = (data.sessionId as string) || '';
const body = WebServer.pushBodyText(event, data, sessionName);
// Webhook channel (ntfy / Slack / Discord / generic): independent of Web Push, so it runs BEFORE
// the "no subscriptions" return below, which is exactly the headless-server case it exists for.
// Fire-and-forget; WebhookNotifier dedupes, caps what is in flight and never throws.
void this.webhookNotifier
.notify({
event,
title: template.title,
body,
urgency: template.urgency as WebhookUrgency,
sessionId: sessionId || undefined,
sessionName: sessionName || undefined,
host: this.windowTitle,
})
.catch(() => undefined);
const subscriptions = this.pushStore.getAll();
if (subscriptions.length === 0) return;
@@ -2670,9 +2699,6 @@ export class WebServer extends EventEmitter {
const vapidKeys = this.pushStore.getVapidKeys();
webpush.setVapidDetails('mailto:codeman@localhost', vapidKeys.publicKey, vapidKeys.privateKey);
const sessionName = (data.sessionName as string) || '';
const sessionId = (data.sessionId as string) || '';
// Multi-user: a session-scoped push (all PUSH_EVENT_MAP events carry a sessionId)
// must reach only the owner's devices (+ admins) — the body embeds the session
// name + activity, so cross-user delivery would leak it. Resolved once here; the
@@ -2680,25 +2706,6 @@ export class WebServer extends EventEmitter {
const multiUserPush = isMultiUserMode();
const pushSessionOwner = sessionId ? this.sessions.get(sessionId)?.owner : undefined;
// Build body text from event data
let body = sessionName ? `[${sessionName}]` : '';
if (event === SseEvent.SessionError && data.error) {
body += body ? ' ' : '';
body += String(data.error).slice(0, 200);
} else if (event === SseEvent.RespawnBlocked && data.reason) {
body += body ? ' ' : '';
body += String(data.reason);
} else if (event === SseEvent.SessionRalphCompletionDetected && data.phrase) {
body += body ? ' ' : '';
body += String(data.phrase);
} else if (event === SseEvent.SessionRespawnBreakerTripped && data.count) {
body += body ? ' ' : '';
body += `Stopped after ${Number(data.count)} rapid crashes — restart the session to retry`;
} else if (event === SseEvent.HookPermissionPrompt && data.tool_name) {
body += body ? ' ' : '';
body += `Tool: ${String(data.tool_name)}`;
}
const payload = JSON.stringify({
title: template.title,
// Hostname-aware prefix so OS-level notifications from multiple Codeman
@@ -2752,6 +2759,28 @@ export class WebServer extends EventEmitter {
}
}
/** The notification body for an event (shared by Web Push and the webhook channel). */
private static pushBodyText(event: string, data: Record<string, unknown>, sessionName: string): string {
let body = sessionName ? `[${sessionName}]` : '';
if (event === SseEvent.SessionError && data.error) {
body += body ? ' ' : '';
body += String(data.error).slice(0, 200);
} else if (event === SseEvent.RespawnBlocked && data.reason) {
body += body ? ' ' : '';
body += String(data.reason);
} else if (event === SseEvent.SessionRalphCompletionDetected && data.phrase) {
body += body ? ' ' : '';
body += String(data.phrase);
} else if (event === SseEvent.SessionRespawnBreakerTripped && data.count) {
body += body ? ' ' : '';
body += `Stopped after ${Number(data.count)} rapid crashes — restart the session to retry`;
} else if (event === SseEvent.HookPermissionPrompt && data.tool_name) {
body += body ? ' ' : '';
body += `Tool: ${String(data.tool_name)}`;
}
return body;
}
private cleanupDeadSSEClients(): void {
this.sse.cleanupDeadClients();
}
+328
View File
@@ -0,0 +1,328 @@
/**
* @fileoverview Webhook notifications (ntfy, Slack, Discord, generic JSON) for the events that
* already trigger Web Push, so a headless server can reach a phone without a browser tab or a
* push subscription.
*
* Split in three, so the parts that matter are testable without a network:
* - pure: `webhookUrlProblem`, `maskWebhookUrl`, `shouldSendWebhook`, `buildWebhookRequest`
* - store: `~/.codeman/webhook.json`, written 0600 via tmp+rename (the URL is a bearer secret:
* anyone holding a Slack/Discord webhook URL can post as it)
* - IO: `sendWebhook` (injected fetch) and `WebhookNotifier` (dedupe, in-flight cap, last result)
*
* Rules the code keeps and the tests pin:
* - The URL is configured only through the admin-only `/api/webhook` routes and kept OUT of
* `settings.json`, which every logged-in user can read through `GET /api/settings`.
* - Delivery goes through `webviewFetch`: link-local and cloud-metadata targets are refused on
* the RESOLVED address at connect time, redirects are not followed, and the call is bounded
* by a timeout. Loopback and LAN stay allowed on purpose (a local ntfy is the feature).
* - The URL never appears in a log line, a result, or an error message.
* - Session names and error text are user/agent-controlled, so they cannot ping a channel:
* Discord gets `allowed_mentions: { parse: [] }` and Slack control characters are escaped.
*
* @module webhook-notify
*/
import { existsSync, mkdirSync } from 'node:fs';
import fs from 'node:fs/promises';
import { join } from 'node:path';
import { blockedWebviewHostReason } from './web/webview-egress-policy.js';
const WEBHOOK_FILE = 'webhook.json';
const MAX_URL_LENGTH = 2048;
const SEND_TIMEOUT_MS = 5000;
const MAX_BODY_CHARS = 500;
/** Same event + session within this window is sent once: a flapping prompt must not flood a channel. */
const DEDUPE_WINDOW_MS = 3000;
const MAX_IN_FLIGHT = 5;
export const WEBHOOK_KINDS = ['ntfy', 'slack', 'discord', 'generic'] as const;
export type WebhookKind = (typeof WEBHOOK_KINDS)[number];
/** `attention`: only events that need a human (critical / warning). `all`: also "response complete". */
export const WEBHOOK_SCOPES = ['attention', 'all'] as const;
export type WebhookScope = (typeof WEBHOOK_SCOPES)[number];
export type WebhookUrgency = 'critical' | 'warning' | 'info';
export interface WebhookConfig {
enabled: boolean;
kind: WebhookKind;
url: string;
scope: WebhookScope;
}
export const DEFAULT_WEBHOOK_CONFIG: WebhookConfig = { enabled: false, kind: 'ntfy', url: '', scope: 'attention' };
export interface WebhookMessage {
event: string;
title: string;
body: string;
urgency: WebhookUrgency;
sessionId?: string;
sessionName?: string;
/** The Codeman instance's window title, so several machines are told apart. */
host?: string;
}
export interface WebhookResult {
ok: boolean;
status?: number;
error?: string;
at: number;
}
// ---------------------------------------------------------------------------
// Pure
// ---------------------------------------------------------------------------
/** Why `raw` cannot be a webhook URL, or null. Used at save time; delivery re-checks the resolved address. */
export function webhookUrlProblem(raw: string): string | null {
if (raw.length > MAX_URL_LENGTH) return 'URL is too long';
let url: URL;
try {
url = new URL(raw);
} catch {
return 'Not a valid URL';
}
if (url.protocol !== 'https:' && url.protocol !== 'http:') return 'Only http and https URLs are allowed';
if (url.username || url.password) return 'Put credentials in the path or a header-less token, not user:password@';
const blocked = blockedWebviewHostReason(url.hostname);
if (blocked) return `Refused: ${blocked}`;
return null;
}
/** Scheme + host only: the path and query of a webhook URL are the secret. */
export function maskWebhookUrl(raw: string): string {
try {
const url = new URL(raw);
return `${url.protocol}//${url.host}/•••`;
} catch {
return '';
}
}
export function shouldSendWebhook(cfg: WebhookConfig, urgency: WebhookUrgency): boolean {
if (!cfg.enabled || !cfg.url) return false;
return cfg.scope === 'all' || urgency !== 'info';
}
const clip = (s: string, n: number): string => (s.length > n ? `${s.slice(0, n - 1)}…` : s);
/** A header value must be single-line printable ASCII; anything else goes out RFC 2047 encoded. */
function headerSafe(value: string): string {
const oneLine = value.replace(/[\r\n]+/g, ' ').trim();
return /^[\x20-\x7e]*$/.test(oneLine) ? oneLine : `=?UTF-8?B?${Buffer.from(oneLine, 'utf8').toString('base64')}?=`;
}
/** Slack parses `<!channel>`, `<@U123>` and `<url|text>`; escaping the three control characters turns them to text. */
const slackEscape = (s: string): string => s.replace(/&/g, '&amp;').replace(/</g, '&lt;').replace(/>/g, '&gt;');
const NTFY_PRIORITY: Record<WebhookUrgency, string> = { critical: '5', warning: '4', info: '3' };
const NTFY_TAGS: Record<WebhookUrgency, string> = {
critical: 'rotating_light',
warning: 'bell',
info: 'white_check_mark',
};
export interface WebhookRequest {
method: 'POST';
headers: Record<string, string>;
body: string;
}
export function buildWebhookRequest(kind: WebhookKind, msg: WebhookMessage, now: Date = new Date()): WebhookRequest {
const title = clip(msg.title, 120);
const body = clip(msg.body, MAX_BODY_CHARS);
const prefix = msg.host ? `${clip(msg.host, 60)}: ` : '';
switch (kind) {
case 'ntfy':
return {
method: 'POST',
headers: {
'Content-Type': 'text/plain; charset=utf-8',
Title: headerSafe(`${prefix}${title}`),
Priority: NTFY_PRIORITY[msg.urgency],
Tags: NTFY_TAGS[msg.urgency],
},
body: body || title,
};
case 'slack':
return {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ text: `*${slackEscape(`${prefix}${title}`)}*${body ? `\n${slackEscape(body)}` : ''}` }),
};
case 'discord':
return {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
content: clip(`**${prefix}${title}**${body ? `\n${body}` : ''}`, 1900),
// Agent output and session names are not trusted to @everyone a channel.
allowed_mentions: { parse: [] },
}),
};
case 'generic':
return {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
event: msg.event,
title,
body,
urgency: msg.urgency,
sessionId: msg.sessionId ?? null,
sessionName: msg.sessionName ?? null,
host: msg.host ?? null,
at: now.toISOString(),
}),
};
}
}
// ---------------------------------------------------------------------------
// Store
// ---------------------------------------------------------------------------
export function webhookConfigPath(configDir: string): string {
return join(configDir, WEBHOOK_FILE);
}
function coerce(raw: unknown): WebhookConfig {
const r = (typeof raw === 'object' && raw !== null ? raw : {}) as Record<string, unknown>;
return {
enabled: r.enabled === true,
kind: (WEBHOOK_KINDS as readonly unknown[]).includes(r.kind)
? (r.kind as WebhookKind)
: DEFAULT_WEBHOOK_CONFIG.kind,
url: typeof r.url === 'string' ? r.url : '',
scope: (WEBHOOK_SCOPES as readonly unknown[]).includes(r.scope)
? (r.scope as WebhookScope)
: DEFAULT_WEBHOOK_CONFIG.scope,
};
}
export async function readWebhookConfig(configDir: string): Promise<WebhookConfig> {
try {
return coerce(JSON.parse(await fs.readFile(webhookConfigPath(configDir), 'utf-8')));
} catch {
return { ...DEFAULT_WEBHOOK_CONFIG };
}
}
/** 0600 via tmp+rename: `mode` on writeFile only applies to a file being created. */
export async function writeWebhookConfig(configDir: string, cfg: WebhookConfig): Promise<void> {
if (!existsSync(configDir)) mkdirSync(configDir, { recursive: true });
const target = webhookConfigPath(configDir);
const tmp = `${target}.${process.pid}.tmp`;
try {
await fs.writeFile(tmp, JSON.stringify(coerce(cfg), null, 2), { mode: 0o600 });
await fs.rename(tmp, target);
} catch (err) {
await fs.unlink(tmp).catch(() => undefined);
throw err;
}
}
// ---------------------------------------------------------------------------
// IO
// ---------------------------------------------------------------------------
export type WebhookFetch = (target: URL, init: RequestInit) => Promise<Response>;
/** What went wrong, without the URL: blocked / timed out / refused / an HTTP status. */
function describeError(err: unknown): string {
const e = err as { name?: string; message?: string; cause?: { code?: string; message?: string } };
if (e?.name === 'TimeoutError' || e?.name === 'AbortError') return 'Timed out';
const text = `${e?.message ?? ''} ${e?.cause?.message ?? ''}`;
if (/link-local|cloud-metadata|EGRESS/i.test(text))
return 'Refused: target is a link-local or cloud-metadata address';
if (e?.cause?.code === 'ENOTFOUND') return 'Host not found';
if (e?.cause?.code === 'ECONNREFUSED') return 'Connection refused';
return 'Network error';
}
export async function sendWebhook(
cfg: Pick<WebhookConfig, 'kind' | 'url'>,
msg: WebhookMessage,
fetchImpl: WebhookFetch
): Promise<WebhookResult> {
const at = Date.now();
const problem = webhookUrlProblem(cfg.url);
if (problem) return { ok: false, error: problem, at };
const req = buildWebhookRequest(cfg.kind, msg);
try {
const res = await fetchImpl(new URL(cfg.url), {
method: req.method,
headers: req.headers,
body: req.body,
redirect: 'manual',
signal: AbortSignal.timeout(SEND_TIMEOUT_MS),
});
void res.body?.cancel().catch(() => undefined);
if (res.status >= 300 && res.status < 400) {
return { ok: false, status: res.status, error: 'The URL redirects; use the final URL', at };
}
return res.ok
? { ok: true, status: res.status, at }
: { ok: false, status: res.status, error: `HTTP ${res.status}`, at };
} catch (err) {
return { ok: false, error: describeError(err), at };
}
}
/**
* Sends the notifications the server decides on. Fire-and-forget by design (a slow webhook must
* never delay Web Push or a request), so it dedupes, caps what is in flight, and remembers only
* the last result for the Settings status line.
*/
export class WebhookNotifier {
private lastSent = new Map<string, number>();
private inFlight = 0;
private last: WebhookResult | null = null;
constructor(
private readonly load: () => Promise<WebhookConfig>,
private readonly fetchImpl: WebhookFetch,
private readonly now: () => number = Date.now
) {}
get lastResult(): WebhookResult | null {
return this.last;
}
async notify(msg: WebhookMessage): Promise<void> {
const cfg = await this.load();
if (!shouldSendWebhook(cfg, msg.urgency)) return;
const key = `${msg.event}:${msg.sessionId ?? ''}`;
const t = this.now();
const prev = this.lastSent.get(key);
if (prev !== undefined && t - prev < DEDUPE_WINDOW_MS) return;
if (this.inFlight >= MAX_IN_FLIGHT) return;
this.lastSent.set(key, t);
if (this.lastSent.size > 256) {
for (const [k, v] of this.lastSent) if (t - v > DEDUPE_WINDOW_MS) this.lastSent.delete(k);
}
this.inFlight++;
try {
this.last = await sendWebhook(cfg, msg, this.fetchImpl);
} finally {
this.inFlight--;
}
}
/** A deliberate test send: bypasses `enabled`, scope and dedupe, and records the result. */
async sendTest(cfg: Pick<WebhookConfig, 'kind' | 'url'>, host?: string): Promise<WebhookResult> {
const result = await sendWebhook(
cfg,
{
event: 'webhook:test',
title: 'Codeman test notification',
body: 'If you can read this, webhook notifications are working.',
urgency: 'info',
host,
},
this.fetchImpl
);
this.last = result;
return result;
}
}