// ============================================================================= // Socket.io client that connects to the ACP backend and dispatches events. // ============================================================================= import { io, type Socket } from "socket.io-client"; import { SocketEvent, type AcpJobEventData } from "./types.js"; const PAGERDUTY_EVENTS_API_URL = "https://events.pagerduty.com/v2/enqueue"; const DEFAULT_HEARTBEAT_INTERVAL_MS = 5 * 60 * 1000; const DEFAULT_DISCONNECT_ALERT_THRESHOLD_MS = 2 * 60 * 1000; const DEFAULT_DISCONNECT_MONITOR_INTERVAL_MS = 30 * 1000; const DEFAULT_MANUAL_RECONNECT_INTERVAL_MS = 5 * 1000; const DEFAULT_FAILED_RECONNECTS_BEFORE_ALERT = 3; type PagerDutyAction = "trigger" | "resolve"; type SocketConnectOptions = { auth: { walletAddress: string }; transports: ["websocket"]; }; export interface SocketLike { connected: boolean; on(event: string, handler: (...args: any[]) => void): unknown; connect(): void; disconnect(): void; } export interface AcpSocketDeps { createSocket?: (acpUrl: string, options: SocketConnectOptions) => SocketLike; fetchFn?: typeof fetch; setIntervalFn?: typeof setInterval; clearIntervalFn?: typeof clearInterval; nowFn?: () => number; heartbeatIntervalMs?: number; disconnectAlertThresholdMs?: number; disconnectMonitorIntervalMs?: number; manualReconnectIntervalMs?: number; failedReconnectsBeforeAlert?: number; pagerDutyRetryDelayMs?: number; } interface PagerDutyOptions { routingKey: string; dedupKey: string; source: string; action: PagerDutyAction; summary: string; severity?: "critical" | "error" | "warning" | "info"; details?: Record; fetchFn: typeof fetch; retryDelayMs?: number; shouldRetry?: () => boolean; } const PAGERDUTY_RETRY_DELAY_MS = 2_000; async function sendPagerDutyEvent(opts: PagerDutyOptions): Promise { const { routingKey, dedupKey, source, action, summary, severity = "critical", details, fetchFn, retryDelayMs = PAGERDUTY_RETRY_DELAY_MS, shouldRetry, } = opts; if (!routingKey) { console.warn( "[socket] PAGERDUTY_ROUTING_KEY not set; skipping PagerDuty event", ); return; } const payload = { routing_key: routingKey, event_action: action, dedup_key: dedupKey, payload: { summary, source, severity, component: "acp-socket", group: "seller-runtime", class: "socket-connectivity", custom_details: details ?? {}, }, }; const attempt = async (): Promise<{ ok: boolean; retryable: boolean }> => { try { const response = await fetchFn(PAGERDUTY_EVENTS_API_URL, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify(payload), }); if (response.ok) { console.log(`[socket] PagerDuty ${action} sent`); return { ok: true, retryable: false }; } const body = await response.text().catch(() => ""); console.error( `[socket] PagerDuty ${action} failed: ${response.status} ${body}`, ); // Only retry on server errors (5xx); 4xx are client errors and won't benefit from retry return { ok: false, retryable: response.status >= 500 }; } catch (err) { const message = err instanceof Error ? err.message : String(err); console.error(`[socket] PagerDuty ${action} error: ${message}`); // Network failures are retryable return { ok: false, retryable: true }; } }; const first = await attempt(); if (first.ok || !first.retryable) return; if (shouldRetry && !shouldRetry()) { console.log(`[socket] PagerDuty ${action} retry skipped (state changed)`); return; } console.log( `[socket] PagerDuty ${action} — retrying in ${retryDelayMs}ms...`, ); await new Promise((resolve) => setTimeout(resolve, retryDelayMs)); if (shouldRetry && !shouldRetry()) { console.log(`[socket] PagerDuty ${action} retry aborted (state changed)`); return; } await attempt(); } export interface AcpSocketCallbacks { onNewTask: (data: AcpJobEventData) => void; onEvaluate?: (data: AcpJobEventData) => void; } export interface AcpSocketOptions { acpUrl: string; walletAddress: string; callbacks: AcpSocketCallbacks; } /** * Connect to the ACP socket and start listening for seller events. * Returns a cleanup function that disconnects the socket. */ export function connectAcpSocket( opts: AcpSocketOptions, deps: AcpSocketDeps = {}, ): () => void { const { acpUrl, walletAddress, callbacks } = opts; const createSocket = deps.createSocket ?? ((url: string, options: SocketConnectOptions): SocketLike => io(url, options) as Socket); const fetchFn = deps.fetchFn ?? fetch; const setIntervalFn = deps.setIntervalFn ?? setInterval; const clearIntervalFn = deps.clearIntervalFn ?? clearInterval; const now = deps.nowFn ?? Date.now; const heartbeatIntervalMs = deps.heartbeatIntervalMs ?? DEFAULT_HEARTBEAT_INTERVAL_MS; const disconnectAlertThresholdMs = deps.disconnectAlertThresholdMs ?? DEFAULT_DISCONNECT_ALERT_THRESHOLD_MS; const disconnectMonitorIntervalMs = deps.disconnectMonitorIntervalMs ?? DEFAULT_DISCONNECT_MONITOR_INTERVAL_MS; const manualReconnectIntervalMs = deps.manualReconnectIntervalMs ?? DEFAULT_MANUAL_RECONNECT_INTERVAL_MS; const failedReconnectsBeforeAlert = deps.failedReconnectsBeforeAlert ?? DEFAULT_FAILED_RECONNECTS_BEFORE_ALERT; const pagerDutyRetryDelayMs = deps.pagerDutyRetryDelayMs ?? PAGERDUTY_RETRY_DELAY_MS; const pdRoutingKey = process.env.PAGERDUTY_ROUTING_KEY ?? ""; const pdDedupKey = `acp-socket-${walletAddress.toLowerCase()}`; const pdSource = `openclaw-acp-seller:${walletAddress}`; let disconnectedAt: number | null = null; let failedReconnectAttempts = 0; let pdIncidentOpen = false; let reconnectInterval: ReturnType | null = null; const socket = createSocket(acpUrl, { auth: { walletAddress }, transports: ["websocket"], }); const triggerPagerDuty = ( summary: string, details: Record, ): void => { if (pdIncidentOpen) return; pdIncidentOpen = true; void sendPagerDutyEvent({ routingKey: pdRoutingKey, dedupKey: pdDedupKey, source: pdSource, action: "trigger", summary, severity: "critical", details, fetchFn, retryDelayMs: pagerDutyRetryDelayMs, shouldRetry: () => pdIncidentOpen, }); }; const resolvePagerDuty = ( summary: string, details: Record, ): void => { if (!pdIncidentOpen) return; pdIncidentOpen = false; void sendPagerDutyEvent({ routingKey: pdRoutingKey, dedupKey: pdDedupKey, source: pdSource, action: "resolve", summary, severity: "info", details, fetchFn, retryDelayMs: pagerDutyRetryDelayMs, }); }; const stopReconnectLoop = (): void => { if (reconnectInterval) { clearIntervalFn(reconnectInterval); reconnectInterval = null; } }; const startReconnectLoop = (): void => { if (reconnectInterval) return; console.log( `[socket] Server-initiated disconnect — starting manual reconnect loop (${manualReconnectIntervalMs}ms)`, ); reconnectInterval = setIntervalFn(() => { if (socket.connected) { stopReconnectLoop(); return; } failedReconnectAttempts += 1; console.log( `[socket] Manual reconnect attempt #${failedReconnectAttempts}`, ); socket.connect(); if (failedReconnectAttempts >= failedReconnectsBeforeAlert) { triggerPagerDuty( `ACP socket failed to reconnect after ${failedReconnectsBeforeAlert} attempts`, { walletAddress, failedReconnectAttempts, acpUrl, }, ); } }, manualReconnectIntervalMs); }; socket.on( SocketEvent.ROOM_JOINED, (_data: unknown, callback?: (ack: boolean) => void) => { console.log("[socket] Joined ACP room"); if (typeof callback === "function") callback(true); }, ); socket.on( SocketEvent.ON_NEW_TASK, (data: AcpJobEventData, callback?: (ack: boolean) => void) => { if (typeof callback === "function") callback(true); console.log(`[socket] onNewTask jobId=${data.id} phase=${data.phase}`); try { callbacks.onNewTask(data); } catch (err) { console.error( `[socket] CRITICAL: onNewTask callback threw synchronously for job ${data.id}:`, err instanceof Error ? err.message : String(err), ); } }, ); socket.on( SocketEvent.ON_EVALUATE, (data: AcpJobEventData, callback?: (ack: boolean) => void) => { if (typeof callback === "function") callback(true); console.log(`[socket] onEvaluate jobId=${data.id} phase=${data.phase}`); try { if (callbacks.onEvaluate) { callbacks.onEvaluate(data); } } catch (err) { console.error( `[socket] CRITICAL: onEvaluate callback threw synchronously for job ${data.id}:`, err instanceof Error ? err.message : String(err), ); } }, ); socket.on("connect", () => { const wasDisconnected = disconnectedAt !== null; let disconnectedMs = 0; if (disconnectedAt !== null) { disconnectedMs = now() - disconnectedAt; } disconnectedAt = null; failedReconnectAttempts = 0; stopReconnectLoop(); console.log("[socket] Connected to ACP"); if (wasDisconnected) { resolvePagerDuty("ACP socket reconnected successfully", { walletAddress, disconnectedMs, acpUrl, }); } }); socket.on("disconnect", (reason) => { console.log(`[socket] Disconnected: ${reason}`); if (disconnectedAt === null) { disconnectedAt = now(); } // "io server disconnect" = server forcibly closed the connection. // Socket.io will NOT auto-reconnect — we must explicitly reconnect. if (reason === "io server disconnect") { startReconnectLoop(); } }); socket.on("connect_error", (err: Error) => { console.error(`[socket] Connection error: ${err.message}`); }); const heartbeatInterval = setIntervalFn(() => { if (socket.connected) { console.log("[socket] Heartbeat: connected to ACP"); return; } const downForMs = disconnectedAt !== null ? now() - disconnectedAt : 0; console.log( `[socket] Heartbeat: disconnected for ${Math.floor(downForMs / 1000)}s`, ); }, heartbeatIntervalMs); const disconnectMonitor = setIntervalFn(() => { if (disconnectedAt === null) return; const downForMs = now() - disconnectedAt; if (downForMs > disconnectAlertThresholdMs) { triggerPagerDuty("ACP socket disconnected for over 2 minutes", { walletAddress, disconnectedForSeconds: Math.floor(downForMs / 1000), acpUrl, }); } }, disconnectMonitorIntervalMs); let isDisconnected = false; const disconnect = () => { if (isDisconnected) return; isDisconnected = true; stopReconnectLoop(); clearIntervalFn(heartbeatInterval); clearIntervalFn(disconnectMonitor); socket.disconnect(); process.off("SIGINT", onSigInt); process.off("SIGTERM", onSigTerm); }; const onSigInt = () => { disconnect(); process.exit(0); }; const onSigTerm = () => { disconnect(); process.exit(0); }; process.on("SIGINT", onSigInt); process.on("SIGTERM", onSigTerm); return disconnect; }