diff --git a/deno.json b/deno.json index f17b44e..2194924 100644 --- a/deno.json +++ b/deno.json @@ -1,6 +1,6 @@ { "name": "@moonlight-protocol/network-dashboard-platform", - "version": "0.1.3", + "version": "0.1.4", "license": "MIT", "exports": "./src/main.ts", "tasks": { diff --git a/src/config/logger.ts b/src/config/logger.ts index 18fdc59..428abf2 100644 --- a/src/config/logger.ts +++ b/src/config/logger.ts @@ -1,59 +1,10 @@ -import chalk from "chalk"; -import { LOG_LEVEL } from "@/config/env.ts"; - -export enum LogLevel { - FATAL = 0, - ERROR = 1, - WARN = 2, - INFO = 3, - DEBUG = 4, - TRACE = 5, -} - -class Logger { - constructor(private level: LogLevel) {} - - private format(args: unknown[]): string { - return args - .map((arg) => { - if (typeof arg === "string") return arg; - try { - return chalk.cyan(JSON.stringify(arg)); - } catch (err) { - return chalk.cyan(`[Unstringifiable: ${(err as Error).message}]`); - } - }) - .join(" "); - } - - private write(level: LogLevel, color: typeof chalk.blue, ...args: unknown[]) { - if (this.level < level) return; - const ts = new Date().toISOString(); - const prefix = chalk.gray(`[${ts}::${LogLevel[level]}]`); - console.log(`${prefix} ${color(this.format(args))}`); - } - - trace(...args: unknown[]) { - this.write(LogLevel.TRACE, chalk.white, ...args); - } - debug(...args: unknown[]) { - this.write(LogLevel.DEBUG, chalk.green, ...args); - } - info(...args: unknown[]) { - this.write(LogLevel.INFO, chalk.blue, ...args); - } - warn(...args: unknown[]) { - this.write(LogLevel.WARN, chalk.yellow, ...args); - } - error(...args: unknown[]) { - this.write(LogLevel.ERROR, chalk.red, ...args); - } - fatal(...args: unknown[]) { - this.write(LogLevel.FATAL, chalk.bgRed.white, ...args); - } +import { type Logger, newLogger, parseLevel } from "@/utils/logger/index.ts"; + +/** + * Creates the root logger from `LOG_LEVEL` env var. Called once in main.ts; + * the returned logger is threaded through to every service and free function + * via dependency injection. There is no module-level singleton. + */ +export function createLogger(): Logger { + return newLogger(parseLevel(Deno.env.get("LOG_LEVEL"))); } - -const resolvedLevel = LogLevel[LOG_LEVEL as keyof typeof LogLevel] ?? - LogLevel.INFO; - -export const LOG = new Logger(resolvedLevel); diff --git a/src/core/events/bus.ts b/src/core/events/bus.ts index 7d44193..7f12adf 100644 --- a/src/core/events/bus.ts +++ b/src/core/events/bus.ts @@ -1,3 +1,4 @@ +import type { Logger } from "@/utils/logger/index.ts"; import type { NetworkEvent } from "./types.ts"; type Listener = (event: NetworkEvent) => void; @@ -11,24 +12,35 @@ type Listener = (event: NetworkEvent) => void; * replaced. * * A misbehaving listener must never break the publish loop — every - * delivery is wrapped in a try/catch. + * delivery is wrapped in a try/catch and reported via the injected logger. */ -class NetworkEventBus { +export class NetworkEventBus { private listeners = new Set(); + private log: Logger; + + constructor(deps: { log: Logger }) { + this.log = deps.log.scope("NetworkEventBus"); + } subscribe(listener: Listener): () => void { + this.log.info("subscribe"); this.listeners.add(listener); + this.log.debug("listenerCount", this.listeners.size); return () => { + this.log.info("unsubscribe"); this.listeners.delete(listener); }; } publish(event: NetworkEvent): void { + this.log.info("publish"); + this.log.debug("eventKind", event.kind); + this.log.debug("listenerCount", this.listeners.size); for (const listener of this.listeners) { try { listener(event); } catch (err) { - console.warn("[network-event-bus] listener threw:", err); + this.log.error(err, "listener threw during publish"); } } } @@ -37,5 +49,3 @@ class NetworkEventBus { return this.listeners.size; } } - -export const networkEventBus = new NetworkEventBus(); diff --git a/src/core/events/bus_test.ts b/src/core/events/bus_test.ts index 2646700..5760d4a 100644 --- a/src/core/events/bus_test.ts +++ b/src/core/events/bus_test.ts @@ -1,5 +1,6 @@ import { assertEquals } from "@std/assert"; -import { networkEventBus } from "./bus.ts"; +import { NetworkEventBus } from "./bus.ts"; +import { newNoop } from "@/utils/logger/index.ts"; import type { NetworkEvent } from "./types.ts"; function ev(id: string): NetworkEvent { @@ -14,33 +15,40 @@ function ev(id: string): NetworkEvent { }; } +function newBus(): NetworkEventBus { + return new NetworkEventBus({ log: newNoop() }); +} + Deno.test("subscribe delivers, unsubscribe stops delivery", () => { + const bus = newBus(); const received: string[] = []; - const unsub = networkEventBus.subscribe((e) => received.push(e.id)); - networkEventBus.publish(ev("a")); - networkEventBus.publish(ev("b")); + const unsub = bus.subscribe((e) => received.push(e.id)); + bus.publish(ev("a")); + bus.publish(ev("b")); unsub(); - networkEventBus.publish(ev("c")); + bus.publish(ev("c")); assertEquals(received, ["a", "b"]); }); Deno.test("publish survives a throwing listener", () => { + const bus = newBus(); const received: string[] = []; - const u1 = networkEventBus.subscribe(() => { + const u1 = bus.subscribe(() => { throw new Error("boom"); }); - const u2 = networkEventBus.subscribe((e) => received.push(e.id)); - networkEventBus.publish(ev("x")); + const u2 = bus.subscribe((e) => received.push(e.id)); + bus.publish(ev("x")); u1(); u2(); assertEquals(received, ["x"]); }); Deno.test("listenerCount reflects subscriptions", () => { - const u1 = networkEventBus.subscribe(() => {}); - const u2 = networkEventBus.subscribe(() => {}); - assertEquals(networkEventBus.listenerCount(), 2); + const bus = newBus(); + const u1 = bus.subscribe(() => {}); + const u2 = bus.subscribe(() => {}); + assertEquals(bus.listenerCount(), 2); u1(); u2(); - assertEquals(networkEventBus.listenerCount(), 0); + assertEquals(bus.listenerCount(), 0); }); diff --git a/src/core/sync/contract-init-listener.ts b/src/core/sync/contract-init-listener.ts index 9771650..3e496b4 100644 --- a/src/core/sync/contract-init-listener.ts +++ b/src/core/sync/contract-init-listener.ts @@ -1,9 +1,10 @@ import { Address, xdr } from "stellar-sdk"; import { Server } from "stellar-sdk/rpc"; -import { LOG } from "@/config/logger.ts"; +import type { Logger } from "@/utils/logger/index.ts"; import { STELLAR_RPC_URL } from "@/config/env.ts"; import { networkState } from "@/core/state/store.ts"; import { refreshTopology } from "./topology-refresh.ts"; +import type { NetworkEventBus } from "@/core/events/bus.ts"; import { isKnownChannelAuthHash, isReady as isWasmRegistryReady, @@ -74,29 +75,30 @@ export function isContractInitListenerEnabled(): boolean { */ export async function evaluateUnknownContract( contractId: string, + deps: { log: Logger; bus: NetworkEventBus }, ): Promise { + const log = deps.log.scope("evaluateUnknownContract"); + log.info("evaluateUnknownContract"); + log.debug("contractId", contractId); + if (!isContractInitListenerEnabled()) return; if (!contractId) return; if (pendingAdoption.has(contractId)) return; if (notMoonlight.has(contractId)) return; if (networkState.hasCouncil(contractId)) return; - const wasmHash = await fetchWasmHash(contractId); + const wasmHash = await fetchWasmHash(contractId, { log }); if (wasmHash === null) return; if (!isKnownChannelAuthHash(wasmHash)) { notMoonlight.add(contractId); - LOG.debug("Ignored contract_initialized from non-Moonlight contract", { - contractId, - wasmHash, - }); + log.debug("wasmHash", wasmHash); + log.event("ignored contract_initialized from non-Moonlight contract"); return; } pendingAdoption.add(contractId); - LOG.info("Detected new Channel Auth deploy via contract_initialized", { - contractId, - wasmHash, - }); - await drainPendingAdoptions(); + log.debug("wasmHash", wasmHash); + log.event("detected new Channel Auth deploy via contract_initialized"); + await drainPendingAdoptions(deps); } /** @@ -104,15 +106,20 @@ export async function evaluateUnknownContract( * WASM hash matched but which council-platform hasn't yet registered. * Called once per pollTick — no-op when nothing is pending. */ -export async function drainPendingAdoptions(): Promise { +export async function drainPendingAdoptions( + deps: { log: Logger; bus: NetworkEventBus }, +): Promise { if (pendingAdoption.size === 0) return; - await refreshTopology(`pending=${pendingAdoption.size}`); + const log = deps.log.scope("drainPendingAdoptions"); + log.info("drainPendingAdoptions"); + log.debug("pendingCount", pendingAdoption.size); + + await refreshTopology(`pending=${pendingAdoption.size}`, deps); for (const cid of [...pendingAdoption]) { if (networkState.hasCouncil(cid)) { pendingAdoption.delete(cid); - LOG.info("Pending Channel Auth contract adopted into topology", { - contractId: cid, - }); + log.debug("contractId", cid); + log.event("pending Channel Auth contract adopted into topology"); } } } @@ -123,7 +130,11 @@ export function __resetForTests(): void { pendingAdoption.clear(); } -async function fetchWasmHash(contractId: string): Promise { +async function fetchWasmHash( + contractId: string, + deps: { log: Logger }, +): Promise { + const log = deps.log.scope("fetchWasmHash"); try { const server = getServer(); const key = xdr.LedgerKey.contractData( @@ -144,10 +155,8 @@ async function fetchWasmHash(contractId: string): Promise { .map((b) => b.toString(16).padStart(2, "0")) .join(""); } catch (err) { - LOG.debug("fetchWasmHash failed", { - contractId, - error: err instanceof Error ? err.message : String(err), - }); + log.debug("contractId", contractId); + log.error(err, "fetchWasmHash failed"); return null; } } diff --git a/src/core/sync/council-fetch.ts b/src/core/sync/council-fetch.ts index 0e2e580..4203337 100644 --- a/src/core/sync/council-fetch.ts +++ b/src/core/sync/council-fetch.ts @@ -1,4 +1,4 @@ -import { LOG } from "@/config/logger.ts"; +import type { Logger } from "@/utils/logger/index.ts"; import { COUNCIL_PLATFORM_URL } from "@/config/env.ts"; import type { CouncilTopologyEntry } from "@/core/events/types.ts"; @@ -20,17 +20,25 @@ type PublicCouncil = { type PublicCouncilsResponse = { data?: PublicCouncil[] }; -export async function fetchCouncilTopology(): Promise { +export async function fetchCouncilTopology( + deps: { log: Logger }, +): Promise { + const log = deps.log.scope("fetchCouncilTopology"); + log.info("fetchCouncilTopology"); + const base = COUNCIL_PLATFORM_URL.replace(/\/+$/, ""); const url = `${base}/api/v1/public/councils`; + log.debug("url", url); const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), 15_000); try { + log.event("requesting council list"); const res = await fetch(url, { signal: controller.signal }); if (!res.ok) { throw new Error(`council-platform returned HTTP ${res.status}`); } + log.event("council list received"); const body = (await res.json()) as PublicCouncilsResponse; const entries: CouncilTopologyEntry[] = []; for (const c of body.data ?? []) { @@ -60,7 +68,8 @@ export async function fetchCouncilTopology(): Promise { .filter((code): code is string => !!code), }); } - LOG.info("Fetched council-platform topology", { count: entries.length }); + log.debug("count", entries.length); + log.event("council-platform topology built"); return entries; } finally { clearTimeout(timer); diff --git a/src/core/sync/scheduler.ts b/src/core/sync/scheduler.ts index bd08b5f..76131f2 100644 --- a/src/core/sync/scheduler.ts +++ b/src/core/sync/scheduler.ts @@ -1,4 +1,5 @@ -import { LOG } from "@/config/logger.ts"; +import type { Logger } from "@/utils/logger/index.ts"; +import type { NetworkEventBus } from "@/core/events/bus.ts"; import { networkState } from "@/core/state/store.ts"; import { refreshTopology } from "./topology-refresh.ts"; @@ -9,38 +10,37 @@ let hourlyTimer: number | null = null; let minuteTimer: number | null = null; let running = false; -function hourlyResync(): Promise { - return refreshTopology("hourly resync"); -} +export function startScheduler( + deps: { log: Logger; bus: NetworkEventBus }, +): void { + if (running) return; + running = true; + const log = deps.log.scope("scheduler"); -function minuteSweep(): void { - const purged = networkState.sweepWindow(); - if (purged > 0) { - LOG.debug("Minute sweep dropped stale window entries", { purged }); + function minuteSweep(): void { + const purged = networkState.sweepWindow(); + if (purged > 0) { + log.debug("purged", purged); + log.event("minute sweep dropped stale window entries"); + } } -} -export function startScheduler(): void { - if (running) return; - running = true; hourlyTimer = setInterval( () => { - hourlyResync().catch((err) => { - LOG.error("Hourly re-sync threw", { - error: err instanceof Error ? err.message : String(err), - }); + refreshTopology("hourly resync", deps).catch((err) => { + log.error(err, "hourly re-sync threw"); }); }, HOURLY_RESYNC_MS, ) as unknown as number; minuteTimer = setInterval(minuteSweep, MINUTE_SWEEP_MS) as unknown as number; - LOG.info("Scheduler started", { - hourlyResyncMs: HOURLY_RESYNC_MS, - minuteSweepMs: MINUTE_SWEEP_MS, - }); + + log.debug("hourlyResyncMs", HOURLY_RESYNC_MS); + log.debug("minuteSweepMs", MINUTE_SWEEP_MS); + log.event("scheduler started"); } -export function stopScheduler(): void { +export function stopScheduler(deps: { log: Logger }): void { running = false; if (hourlyTimer !== null) { clearInterval(hourlyTimer); @@ -50,5 +50,5 @@ export function stopScheduler(): void { clearInterval(minuteTimer); minuteTimer = null; } - LOG.info("Scheduler stopped"); + deps.log.scope("scheduler").event("scheduler stopped"); } diff --git a/src/core/sync/soroban-watcher.ts b/src/core/sync/soroban-watcher.ts index 97089ef..69b7dc5 100644 --- a/src/core/sync/soroban-watcher.ts +++ b/src/core/sync/soroban-watcher.ts @@ -1,8 +1,8 @@ import { Server } from "stellar-sdk/rpc"; -import { LOG } from "@/config/logger.ts"; +import type { Logger } from "@/utils/logger/index.ts"; import { STELLAR_RPC_URL } from "@/config/env.ts"; import { networkState } from "@/core/state/store.ts"; -import { networkEventBus } from "@/core/events/bus.ts"; +import type { NetworkEventBus } from "@/core/events/bus.ts"; import { mapChainEvent, type RawChainEvent } from "./event-mapper.ts"; import { CONTRACT_INITIALIZED_TOPIC_PATTERN, @@ -59,17 +59,25 @@ function watchedContractIds(): string[] { return Array.from(ids); } -function publish( +function publishMappedEvent( event: ReturnType, ledgerClosedAtMs: number | null, + bus: NetworkEventBus, + log: Logger, ): void { + log.info("publishMappedEvent"); if (!event) return; + log.debug("kind", event.kind); const latencyMs = ledgerClosedAtMs === null ? null : Math.max(0, Date.now() - ledgerClosedAtMs); const wasNew = networkState.recordEvent(event, latencyMs); - if (!wasNew) return; - networkEventBus.publish(event); + if (!wasNew) { + log.event("event already seen, skipping publish"); + return; + } + log.event("publishing event to bus"); + bus.publish(event); } /** @@ -88,7 +96,12 @@ type ProcessedEvent = { ledgerClosedAtMs: number | null; }; -function processRawEventBatch(raws: RawChainEvent[]): ProcessedEvent[] { +function processRawEventBatch( + raws: RawChainEvent[], + log: Logger, +): ProcessedEvent[] { + log.info("processRawEventBatch"); + log.debug("rawCount", raws.length); const byTx = new Map(); const txOrder: string[] = []; for (const raw of raws) { @@ -126,20 +139,6 @@ function processRawEventBatch(raws: RawChainEvent[]): ProcessedEvent[] { return out; } -/** - * Cold-start scan: walk trailing 24h on the current contractId set, - * map events, and seed the rolling window + ring buffer in chronological - * order. Sets the forward cursor to one past the latest ledger seen. - */ -function describeErr(err: unknown): string { - if (err instanceof Error) return err.message; - try { - return JSON.stringify(err); - } catch { - return String(err); - } -} - /** * Parse Soroban's `ledgerClosedAt` (ISO string) to ms-since-epoch. Older * SDK responses may omit it; in that case latency stays null for the event. @@ -156,12 +155,22 @@ function parseLedgerClosedAt(raw: unknown): number | null { * null if the error doesn't match that pattern. */ function parseValidRangeFloor(err: unknown): number | null { - const msg = describeErr(err); + const msg = err instanceof Error ? err.message : String(err); const match = msg.match(/ledger range:\s*(\d+)\s*-\s*\d+/); return match ? Number(match[1]) : null; } -export async function coldStartScan(): Promise { +/** + * Cold-start scan: walk trailing 24h on the current contractId set, + * map events, and seed the rolling window + ring buffer in chronological + * order. Sets the forward cursor to one past the latest ledger seen. + */ +export async function coldStartScan( + deps: { log: Logger; bus: NetworkEventBus }, +): Promise { + const log = deps.log.scope("coldStartScan"); + log.info("coldStartScan"); + const server = getServer(); const latest = await server.getLatestLedger(); // Soroban's *event* retention is much shorter than its ledger retention @@ -187,28 +196,25 @@ export async function coldStartScan(): Promise { if (probe.events.length > 0) { initialStart = tryStart; if (back !== LOOKBACK_LEDGERS_24H) { - LOG.info("Cold-start scan clamped to events retention floor", { - desiredStart, - workingStart: tryStart, - lookbackLedgers: back, - }); + log.debug("desiredStart", desiredStart); + log.debug("workingStart", tryStart); + log.debug("lookbackLedgers", back); + log.event("cold-start scan clamped to events retention floor"); } break; } } catch (err) { - LOG.debug("Cold-start probe failed at lookback", { - back, - error: err instanceof Error ? err.message : String(err), - }); + log.debug("back", back); + log.error(err, "cold-start probe failed at lookback"); } } const contractIds = watchedContractIds(); if (contractIds.length === 0) { lastLedgerSeen = latest.sequence; - LOG.info( - "Cold-start scan skipped — no contracts to watch yet (no councils registered).", - { latestLedger: latest.sequence }, + log.debug("latestLedger", latest.sequence); + log.event( + "cold-start scan skipped — no contracts to watch yet (no councils registered)", ); return; } @@ -217,11 +223,10 @@ export async function coldStartScan(): Promise { // the forward poller with a null cursor. lastLedgerSeen = latest.sequence; - LOG.info("Cold-start scan starting", { - startLedger: initialStart, - latestLedger: latest.sequence, - contractCount: contractIds.length, - }); + log.debug("startLedger", initialStart); + log.debug("latestLedger", latest.sequence); + log.debug("contractCount", contractIds.length); + log.event("cold-start scan starting"); const rawBatch: RawChainEvent[] = []; @@ -246,10 +251,11 @@ export async function coldStartScan(): Promise { // First-page out-of-range failure: retry once at the RPC's valid floor. const floor = page === 0 ? parseValidRangeFloor(err) : null; if (floor !== null && floor > nextLedger) { - LOG.warn("Cold-start scan startLedger below retention; retrying", { - requestedStartLedger: nextLedger, - retentionFloor: floor, - }); + log.debug("requestedStartLedger", nextLedger); + log.debug("retentionFloor", floor); + log.event( + "cold-start scan startLedger below retention; retrying at floor", + ); nextLedger = floor; try { res = await server.getEvents({ @@ -258,19 +264,15 @@ export async function coldStartScan(): Promise { limit: PAGE_LIMIT, }); } catch (err2) { - LOG.warn("Cold-start scan retry at retention floor failed", { - startLedger: nextLedger, - error: describeErr(err2), - }); + log.debug("startLedger", nextLedger); + log.error(err2, "cold-start scan retry at retention floor failed"); break; } } else { - LOG.warn("Cold-start scan page failed (stopping chunk)", { - page, - startLedger: nextLedger, - chunkSize: chunk.length, - error: describeErr(err), - }); + log.debug("page", page); + log.debug("startLedger", nextLedger); + log.debug("chunkSize", chunk.length); + log.error(err, "cold-start scan page failed (stopping chunk)"); break; } } @@ -292,10 +294,9 @@ export async function coldStartScan(): Promise { const lastLedgerInPage = res.events[res.events.length - 1].ledger; nextLedger = lastLedgerInPage + 1; if (page >= 50) { - LOG.warn("Cold-start scan hit page cap (50) for chunk; stopping", { - eventsSoFar: rawBatch.length, - chunkSize: chunk.length, - }); + log.debug("eventsSoFar", rawBatch.length); + log.debug("chunkSize", chunk.length); + log.event("cold-start scan hit page cap (50) for chunk; stopping"); break; } } @@ -308,20 +309,24 @@ export async function coldStartScan(): Promise { // events grouped per chunk, so re-sort by ledger here to keep the // ring-buffer chronological. rawBatch.sort((a, b) => a.ledger - b.ledger); - const chronological = processRawEventBatch(rawBatch).map((p) => p.event); + const chronological = processRawEventBatch(rawBatch, log).map((p) => p.event); networkState.seedWindow(chronological); // Recent ring buffer: keep newest at index 0. const newestFirst = [...chronological].reverse(); networkState.seedRecent(newestFirst); - LOG.info("Cold-start scan complete", { - pagesWalked: page, - eventsSeeded: chronological.length, - lastLedgerSeen, - }); + + log.debug("pagesWalked", page); + log.debug("eventsSeeded", chronological.length); + log.debug("lastLedgerSeen", lastLedgerSeen); + log.event("cold-start scan complete"); } -async function pollTick(): Promise { +async function pollTick( + deps: { log: Logger; bus: NetworkEventBus }, +): Promise { if (!running || lastLedgerSeen === null) return; + const log = deps.log.scope("pollTick"); + const contractIds = watchedContractIds(); // Soroban's getEvents intersects multiple filter entries within a single @@ -361,10 +366,8 @@ async function pollTick(): Promise { }); } } catch (err) { - LOG.warn("Soroban poll (known contracts) failed", { - chunkSize: chunk.length, - error: describeErr(err), - }); + log.debug("chunkSize", chunk.length); + log.error(err, "Soroban poll (known contracts) failed"); } } @@ -388,61 +391,64 @@ async function pollTick(): Promise { unknownCandidates.add(cid); } } catch (err) { - LOG.warn("Soroban poll (contract_initialized) failed", { - error: err instanceof Error ? err.message : String(err), - }); + log.error(err, "Soroban poll (contract_initialized) failed"); } } - for (const processed of processRawEventBatch(rawBatch)) { - publish(processed.event, processed.ledgerClosedAtMs); + for (const processed of processRawEventBatch(rawBatch, log)) { + publishMappedEvent( + processed.event, + processed.ledgerClosedAtMs, + deps.bus, + log, + ); } for (const cid of unknownCandidates) { - evaluateUnknownContract(cid).catch((err) => { - LOG.warn("evaluateUnknownContract failed", { - contractId: cid, - error: err instanceof Error ? err.message : String(err), - }); + evaluateUnknownContract(cid, deps).catch((err) => { + log.debug("contractId", cid); + log.error(err, "evaluateUnknownContract failed"); }); } // Retry any matched-but-not-yet-registered contracts. No-op when nothing // is pending, so this is cheap when council-platform is caught up. - drainPendingAdoptions().catch((err) => { - LOG.warn("drainPendingAdoptions failed", { - error: err instanceof Error ? err.message : String(err), - }); + drainPendingAdoptions(deps).catch((err) => { + log.error(err, "drainPendingAdoptions failed"); }); lastLedgerSeen = nextLastLedger; } -function scheduleNext(): void { +function scheduleNext(deps: { log: Logger; bus: NetworkEventBus }): void { if (!running) return; pollTimer = setTimeout(async () => { - await pollTick(); - scheduleNext(); + await pollTick(deps); + scheduleNext(deps); }, POLL_INTERVAL_MS) as unknown as number; } -export function startSorobanWatcher(): void { +export function startSorobanWatcher( + deps: { log: Logger; bus: NetworkEventBus }, +): void { if (running) return; running = true; - LOG.info("Soroban watcher started", { - intervalMs: POLL_INTERVAL_MS, - lastLedgerSeen, - }); - scheduleNext(); + const log = deps.log.scope("sorobanWatcher"); + log.debug("intervalMs", POLL_INTERVAL_MS); + log.debug("lastLedgerSeen", lastLedgerSeen); + log.event("soroban watcher started"); + scheduleNext(deps); } -export function stopSorobanWatcher(): void { +export function stopSorobanWatcher(deps: { log: Logger }): void { running = false; if (pollTimer !== null) { clearTimeout(pollTimer); pollTimer = null; } - LOG.info("Soroban watcher stopped"); + deps.log.scope("sorobanWatcher").event("soroban watcher stopped"); } /** Re-anchor the rolling 24h counter window after the hourly re-sync. */ -export async function rescanRollingWindow(): Promise { - await coldStartScan(); +export async function rescanRollingWindow( + deps: { log: Logger; bus: NetworkEventBus }, +): Promise { + await coldStartScan(deps); } diff --git a/src/core/sync/topology-refresh.ts b/src/core/sync/topology-refresh.ts index df08b7b..ac0d1ac 100644 --- a/src/core/sync/topology-refresh.ts +++ b/src/core/sync/topology-refresh.ts @@ -1,5 +1,6 @@ -import { LOG } from "@/config/logger.ts"; +import type { Logger } from "@/utils/logger/index.ts"; import { networkState } from "@/core/state/store.ts"; +import type { NetworkEventBus } from "@/core/events/bus.ts"; import { fetchCouncilTopology } from "./council-fetch.ts"; import { rescanRollingWindow } from "./soroban-watcher.ts"; import { refreshWasmRegistry } from "./wasm-registry.ts"; @@ -17,34 +18,42 @@ import { refreshWasmRegistry } from "./wasm-registry.ts"; let inFlight: Promise | null = null; -export function refreshTopology(reason: string): Promise { +export function refreshTopology( + reason: string, + deps: { log: Logger; bus: NetworkEventBus }, +): Promise { if (inFlight) return inFlight; - inFlight = run(reason).finally(() => { + inFlight = run(reason, deps).finally(() => { inFlight = null; }); return inFlight; } -async function run(reason: string): Promise { +async function run( + reason: string, + deps: { log: Logger; bus: NetworkEventBus }, +): Promise { + const log = deps.log.scope("topologyRefresh"); + log.info("refreshTopology"); + log.debug("reason", reason); + try { // Re-fetch the wasm-hash registry alongside the topology so any new // soroban-core release becomes recognised within an hour without a // dashboard-backend restart. Errors are absorbed by the registry // itself — they don't block the topology refresh. - await refreshWasmRegistry(); - const topology = await fetchCouncilTopology(); + log.event("refreshing WASM registry"); + await refreshWasmRegistry({ log }); + log.event("fetching council topology"); + const topology = await fetchCouncilTopology({ log }); networkState.replaceTopology(topology); - await rescanRollingWindow(); - LOG.info("Topology refreshed", { - reason, - councils: networkState.getCouncilIds().length, - providers: networkState.countActiveProviders(), - assets: networkState.countAssetsRegistered(), - }); + log.event("topology replaced in network state"); + await rescanRollingWindow({ log, bus: deps.bus }); + log.debug("councils", networkState.getCouncilIds().length); + log.debug("providers", networkState.countActiveProviders()); + log.debug("assets", networkState.countAssetsRegistered()); + log.event("topology refreshed"); } catch (err) { - LOG.warn("Topology refresh failed", { - reason, - error: err instanceof Error ? err.message : String(err), - }); + log.error(err, "topology refresh failed"); } } diff --git a/src/core/sync/wasm-registry.ts b/src/core/sync/wasm-registry.ts index 2bb661b..4ef1f28 100644 --- a/src/core/sync/wasm-registry.ts +++ b/src/core/sync/wasm-registry.ts @@ -1,4 +1,4 @@ -import { LOG } from "@/config/logger.ts"; +import type { Logger } from "@/utils/logger/index.ts"; /** * Registry of accepted Channel Auth WASM hashes. @@ -57,7 +57,13 @@ export function __resetForTests(): void { * refreshes are kept so deploys against older council versions remain * recognised even if GitHub temporarily omits an old release. */ -export async function refreshWasmRegistry(): Promise { +export async function refreshWasmRegistry( + deps: { log: Logger }, +): Promise { + const log = deps.log.scope("refreshWasmRegistry"); + log.info("refreshWasmRegistry"); + log.debug("repo", RELEASES_REPO); + let releases: Array< { tag_name: string; @@ -65,22 +71,23 @@ export async function refreshWasmRegistry(): Promise { } >; try { + log.event("fetching release list"); const res = await fetch(RELEASES_URL, { headers: { "Accept": "application/vnd.github+json" }, }); if (!res.ok) { - LOG.warn("Channel Auth WASM registry fetch failed", { - repo: RELEASES_REPO, - status: res.status, - }); + log.debug("status", res.status); + log.error( + new Error(`HTTP ${res.status}`), + "Channel Auth WASM registry fetch failed", + ); return; } releases = await res.json(); + log.event("release list fetched"); + log.debug("releaseCount", releases.length); } catch (err) { - LOG.warn("Channel Auth WASM registry fetch threw", { - repo: RELEASES_REPO, - error: err instanceof Error ? err.message : String(err), - }); + log.error(err, "Channel Auth WASM registry fetch threw"); return; } @@ -91,37 +98,34 @@ export async function refreshWasmRegistry(): Promise { try { const dl = await fetch(asset.browser_download_url); if (!dl.ok) { - LOG.warn("Channel Auth WASM download failed", { - tag: r.tag_name, - status: dl.status, - }); + log.debug("tag", r.tag_name); + log.debug("status", dl.status); + log.error( + new Error(`HTTP ${dl.status}`), + "Channel Auth WASM download failed", + ); continue; } const buf = await dl.arrayBuffer(); const hash = await sha256Hex(buf); if (!validHashes.has(hash)) { validHashes.add(hash); - LOG.info("Channel Auth WASM registered", { - tag: r.tag_name, - wasmHash: hash, - bytes: buf.byteLength, - }); + log.debug("tag", r.tag_name); + log.debug("wasmHash", hash); + log.debug("bytes", buf.byteLength); + log.event("Channel Auth WASM registered"); } } catch (err) { - LOG.warn("Channel Auth WASM probe threw", { - tag: r.tag_name, - error: err instanceof Error ? err.message : String(err), - }); + log.debug("tag", r.tag_name); + log.error(err, "Channel Auth WASM probe threw"); } } lastRefreshOk = validHashes.size > 0; - LOG.info("Channel Auth WASM registry refreshed", { - repo: RELEASES_REPO, - totalReleases: releases.length, - knownHashes: validHashes.size, - addedSinceLast: validHashes.size - beforeCount, - }); + log.debug("totalReleases", releases.length); + log.debug("knownHashes", validHashes.size); + log.debug("addedSinceLast", validHashes.size - beforeCount); + log.event("Channel Auth WASM registry refreshed"); } async function sha256Hex(buf: ArrayBuffer): Promise { diff --git a/src/http/v1/network-ws.ts b/src/http/v1/network-ws.ts index ae4a50a..a0a5f8c 100644 --- a/src/http/v1/network-ws.ts +++ b/src/http/v1/network-ws.ts @@ -1,7 +1,7 @@ import type { Context } from "@oak/oak"; import { Router } from "@oak/oak"; -import { LOG } from "@/config/logger.ts"; -import { networkEventBus } from "@/core/events/bus.ts"; +import type { Logger } from "@/utils/logger/index.ts"; +import type { NetworkEventBus } from "@/core/events/bus.ts"; import { buildSnapshotFrame } from "@/core/state/snapshot.ts"; import { NETWORK_WS_SUBPROTOCOL, @@ -25,68 +25,80 @@ const IDLE_TIMEOUT_SECONDS = 30; * * No client → server frames. Clients reconnect rather than keep-alive. */ -export function networkWsHandler(ctx: Context): void { - if (!ctx.isUpgradable) { - ctx.response.status = 426; - ctx.response.body = { error: "WebSocket upgrade required" }; - return; - } +export function handleNetworkWs( + deps: { log: Logger; bus: NetworkEventBus }, +): (ctx: Context) => void { + const log = deps.log.scope("networkWs"); - const socket = ctx.upgrade({ - protocol: NETWORK_WS_SUBPROTOCOL, - idleTimeout: IDLE_TIMEOUT_SECONDS, - }); + return (ctx: Context) => { + log.info("handleNetworkWs"); - let unsubscribe: (() => void) | null = null; - let closed = false; - - const cleanup = () => { - if (closed) return; - closed = true; - if (unsubscribe) { - unsubscribe(); - unsubscribe = null; + if (!ctx.isUpgradable) { + ctx.response.status = 426; + ctx.response.body = { error: "WebSocket upgrade required" }; + return; } - }; - const sendFrame = (frame: ServerFrame): void => { - if (socket.readyState !== WebSocket.OPEN) return; - try { - socket.send(JSON.stringify(frame)); - } catch (err) { - LOG.warn("Failed to send network WS frame", { - type: frame.type, - error: err instanceof Error ? err.message : String(err), - }); - } - }; + const socket = ctx.upgrade({ + protocol: NETWORK_WS_SUBPROTOCOL, + idleTimeout: IDLE_TIMEOUT_SECONDS, + }); + + let unsubscribe: (() => void) | null = null; + let closed = false; + + const cleanup = () => { + if (closed) return; + closed = true; + if (unsubscribe) { + unsubscribe(); + unsubscribe = null; + } + }; - socket.onopen = () => { - sendFrame(buildSnapshotFrame()); - unsubscribe = networkEventBus.subscribe((event) => { - sendFrame({ - type: "event", - event, - counters: buildSnapshotFrame().counters, + const sendFrame = (frame: ServerFrame): void => { + if (socket.readyState !== WebSocket.OPEN) return; + try { + socket.send(JSON.stringify(frame)); + } catch (err) { + log.debug("type", frame.type); + log.error(err, "failed to send network WS frame"); + } + }; + + socket.onopen = () => { + sendFrame(buildSnapshotFrame()); + unsubscribe = deps.bus.subscribe((event) => { + sendFrame({ + type: "event", + event, + counters: buildSnapshotFrame().counters, + }); }); - }); - LOG.info("Network WS opened", { - subscribers: networkEventBus.listenerCount(), - }); - }; + log.debug("subscribers", deps.bus.listenerCount()); + log.event("network WS opened"); + }; - socket.onclose = () => { - cleanup(); - LOG.info("Network WS closed"); - }; + socket.onclose = () => { + cleanup(); + log.event("network WS closed"); + }; - socket.onerror = (event) => { - LOG.warn("Network WS error", { - message: event instanceof ErrorEvent ? event.message : "unknown", - }); - cleanup(); + socket.onerror = (event) => { + log.debug( + "message", + event instanceof ErrorEvent ? event.message : "unknown", + ); + log.error(event, "network WS error"); + cleanup(); + }; }; } -export const networkWsRouter = new Router(); -networkWsRouter.get("/network/ws", networkWsHandler); +export function buildNetworkWsRouter( + deps: { log: Logger; bus: NetworkEventBus }, +): Router { + const router = new Router(); + router.get("/network/ws", handleNetworkWs(deps)); + return router; +} diff --git a/src/http/v1/v1.routes.ts b/src/http/v1/v1.routes.ts index 7f70f80..a9d6365 100644 --- a/src/http/v1/v1.routes.ts +++ b/src/http/v1/v1.routes.ts @@ -1,14 +1,25 @@ import { Router } from "@oak/oak"; +import type { Logger } from "@/utils/logger/index.ts"; +import type { NetworkEventBus } from "@/core/events/bus.ts"; import { healthRouter } from "./health.ts"; -import { networkWsRouter } from "./network-ws.ts"; +import { buildNetworkWsRouter } from "./network-ws.ts"; -const apiRouter = new Router(); +export function buildApiRouter( + deps: { log: Logger; bus: NetworkEventBus }, +): Router { + const apiRouter = new Router(); + const networkWsRouter = buildNetworkWsRouter(deps); -apiRouter.use("/api/v1", healthRouter.routes(), healthRouter.allowedMethods()); -apiRouter.use( - "/api/v1", - networkWsRouter.routes(), - networkWsRouter.allowedMethods(), -); + apiRouter.use( + "/api/v1", + healthRouter.routes(), + healthRouter.allowedMethods(), + ); + apiRouter.use( + "/api/v1", + networkWsRouter.routes(), + networkWsRouter.allowedMethods(), + ); -export default apiRouter; + return apiRouter; +} diff --git a/src/main.ts b/src/main.ts index b172deb..99e19a4 100644 --- a/src/main.ts +++ b/src/main.ts @@ -1,9 +1,10 @@ import { Application } from "@oak/oak"; -import apiV1 from "@/http/v1/v1.routes.ts"; +import { buildApiRouter } from "@/http/v1/v1.routes.ts"; import { corsMiddleware } from "@/http/middleware/cors.ts"; import { PORT } from "@/config/env.ts"; -import { LOG } from "@/config/logger.ts"; +import { createLogger } from "@/config/logger.ts"; +import { NetworkEventBus } from "@/core/events/bus.ts"; import { networkState } from "@/core/state/store.ts"; import { fetchCouncilTopology } from "@/core/sync/council-fetch.ts"; import { @@ -34,46 +35,55 @@ import { refreshWasmRegistry } from "@/core/sync/wasm-registry.ts"; * a boot-time outage self-heals. */ async function bootstrap() { + const rootLog = createLogger(); + const log = rootLog.scope("bootstrap"); + log.info("bootstrap"); + + const bus = new NetworkEventBus({ log: rootLog }); + const deps = { log: rootLog, bus }; + try { try { - await refreshWasmRegistry(); + await refreshWasmRegistry({ log: rootLog }); } catch (err) { - LOG.error("Initial WASM registry fetch failed (continuing degraded)", { - error: err instanceof Error ? err.message : String(err), - }); + log.error( + err, + "initial WASM registry fetch failed (continuing degraded)", + ); } try { - const topology = await fetchCouncilTopology(); + const topology = await fetchCouncilTopology({ log: rootLog }); networkState.replaceTopology(topology); } catch (err) { - LOG.error("Initial council-platform fetch failed (continuing degraded)", { - error: err instanceof Error ? err.message : String(err), - }); + log.error( + err, + "initial council-platform fetch failed (continuing degraded)", + ); } try { - await coldStartScan(); + await coldStartScan(deps); } catch (err) { - LOG.error("Cold-start scan failed (continuing degraded)", { - error: err instanceof Error ? err.message : String(err), - }); + log.error(err, "cold-start scan failed (continuing degraded)"); } - startSorobanWatcher(); - startScheduler(); + startSorobanWatcher(deps); + startScheduler(deps); const app = new Application(); app.use(corsMiddleware); + const apiV1 = buildApiRouter(deps); app.use(apiV1.routes()); app.use(apiV1.allowedMethods()); - LOG.info(`network-dashboard-platform running on http://localhost:${PORT}`); + log.debug("port", PORT); + log.event(`network-dashboard-platform running on http://localhost:${PORT}`); const shutdown = () => { - LOG.info("Shutting down..."); - stopSorobanWatcher(); - stopScheduler(); + log.event("shutting down"); + stopSorobanWatcher({ log: rootLog }); + stopScheduler({ log: rootLog }); Deno.exit(0); }; Deno.addSignalListener("SIGINT", shutdown); @@ -81,11 +91,9 @@ async function bootstrap() { await app.listen({ port: PORT }); } catch (err) { - LOG.fatal("Failed to start", { - error: err instanceof Error ? err.message : String(err), - }); - stopSorobanWatcher(); - stopScheduler(); + log.error(err, "failed to start"); + stopSorobanWatcher({ log: rootLog }); + stopScheduler({ log: rootLog }); Deno.exit(1); } } diff --git a/src/utils/logger/index.ts b/src/utils/logger/index.ts new file mode 100644 index 0000000..5623024 --- /dev/null +++ b/src/utils/logger/index.ts @@ -0,0 +1,231 @@ +// Draft Logger module — TypeScript port of github.com/AquiGorka/go-logger. +// Lives at src/utils/logger/index.ts in each backend repo. + +import chalk from "chalk"; + +export enum Level { + Debug = 0, + Info = 1, + Event = 2, + Disabled = 3, +} + +export interface Logger { + info(msg: string): void; + event(msg: string): void; + debug(key: string, value: unknown): void; + error(err: unknown, msg: string): void; + scope(name: string): Logger; +} + +export interface Writer { + write(line: string): void; +} + +export interface LoggerOptions { + /** Custom stdout writer. Replaces console.log. Useful for tests. */ + writer?: Writer; + /** Opt-in file path for JSON-formatted records. Created if missing. */ + file?: string; +} + +interface Record { + ts: string; + level: "debug" | "info" | "event" | "error"; + scope: string; + msg?: string; + key?: string; + value?: unknown; + error?: string; +} + +type Format = (r: Record) => string; + +interface Sink { + writer: Writer; + format: Format; +} + +export function parseLevel(s: string | undefined): Level { + switch ((s ?? "").toLowerCase()) { + case "debug": + return Level.Debug; + case "info": + return Level.Info; + case "event": + return Level.Event; + default: + return Level.Disabled; + } +} + +const stdoutWriter: Writer = { + write: (line) => console.log(line), +}; + +class FileWriter implements Writer { + private file: Deno.FsFile; + constructor(path: string) { + const dir = path.substring(0, path.lastIndexOf("/")); + if (dir) Deno.mkdirSync(dir, { recursive: true }); + this.file = Deno.openSync(path, { + append: true, + create: true, + write: true, + }); + } + write(line: string): void { + this.file.writeSync(new TextEncoder().encode(line + "\n")); + } +} + +function stringify(v: unknown): string { + if (typeof v === "string") return v; + if (v instanceof Error) return v.message; + try { + return JSON.stringify(v); + } catch (err) { + return `[Unstringifiable: ${(err as Error).message}]`; + } +} + +function humanFormat(colored: boolean): Format { + const grayLb = colored ? chalk.gray : (s: string) => s; + const greenLb = colored ? chalk.green : (s: string) => s; + const whiteLb = colored ? chalk.white : (s: string) => s; + const cyanLb = colored ? chalk.cyan : (s: string) => s; + const redLb = colored ? chalk.red : (s: string) => s; + + return (r) => { + const ts = grayLb(`[${r.ts}]`); + switch (r.level) { + case "info": + return `${ts} ${greenLb("INF")} [${r.scope}] ${r.msg}`; + case "event": + return `${ts} ${whiteLb("EVT")} -${r.msg} (${r.scope})`; + case "debug": + return `${ts} ${cyanLb("DBG")} ${r.key}: ${ + stringify(r.value) + } (${r.scope})`; + case "error": + return `${ts} ${redLb("ERR")} [${r.scope}] ${r.msg} error="${r.error}"`; + } + }; +} + +const jsonFormat: Format = (r) => { + // Stable schema. Only the fields relevant to each level are emitted. + const out: Record = { + ts: r.ts, + level: r.level, + scope: r.scope, + }; + if (r.msg !== undefined) out.msg = r.msg; + if (r.key !== undefined) out.key = r.key; + if (r.value !== undefined) out.value = safeJsonValue(r.value); + if (r.error !== undefined) out.error = r.error; + try { + return JSON.stringify(out); + } catch (err) { + return JSON.stringify({ + ts: r.ts, + level: r.level, + scope: r.scope, + msg: `[unserializable record: ${(err as Error).message}]`, + }); + } +}; + +function safeJsonValue(v: unknown): unknown { + // BigInt + circular refs would break JSON.stringify. Coerce to string when + // we can't keep the structured value. + if (typeof v === "bigint") return v.toString(); + if (v instanceof Error) return { message: v.message, name: v.name }; + if (v === null || typeof v !== "object") return v; + try { + JSON.stringify(v); + return v; + } catch { + return stringify(v); + } +} + +class LoggerImpl implements Logger { + constructor( + private readonly level: Level, + private readonly sinks: Sink[], + private readonly scopePath: string, + ) {} + + info(msg: string): void { + if (this.level > Level.Info) return; + this.emit({ ts: now(), level: "info", scope: this.scopePath, msg }); + } + + event(msg: string): void { + if (this.level > Level.Event) return; + this.emit({ ts: now(), level: "event", scope: this.scopePath, msg }); + } + + debug(key: string, value: unknown): void { + if (this.level > Level.Debug) return; + this.emit({ ts: now(), level: "debug", scope: this.scopePath, key, value }); + } + + error(err: unknown, msg: string): void { + // ERR always emits regardless of level (matches go-logger / zerolog). + const detail = err instanceof Error ? err.message : String(err); + this.emit({ + ts: now(), + level: "error", + scope: this.scopePath, + msg, + error: detail, + }); + } + + scope(name: string): Logger { + return new LoggerImpl(this.level, this.sinks, `${this.scopePath}.${name}`); + } + + private emit(r: Record): void { + for (const sink of this.sinks) { + sink.writer.write(sink.format(r)); + } + } +} + +function now(): string { + return new Date().toISOString(); +} + +export function newLogger(level: Level, opts: LoggerOptions = {}): Logger { + const sinks: Sink[] = []; + + // stdout sink — always present. Human format. Colored when TTY (and only + // when caller did not pass a custom writer; tests get plain output). + const consoleWriter = opts.writer ?? stdoutWriter; + const colored = opts.writer === undefined && Deno.stdout.isTerminal(); + sinks.push({ writer: consoleWriter, format: humanFormat(colored) }); + + // file sink — opt-in. JSON format. + if (opts.file !== undefined) { + sinks.push({ writer: new FileWriter(opts.file), format: jsonFormat }); + } + + return new LoggerImpl(level, sinks, "main"); +} + +class NoopLogger implements Logger { + info(): void {} + event(): void {} + debug(): void {} + error(): void {} + scope(): Logger { + return this; + } +} + +export function newNoop(): Logger { + return new NoopLogger(); +}