spb/social-runtime-crawler
Public
TypeScript 91.8%
HTML 3.2%
JavaScript 3%
SQL 1.4%
CSS 0.7%
1import fs from "node:fs";2import path from "node:path";3import { newId, nowIso, type Platform, type Provenance } from "@src/shared";45/** Universal Social Event Format (§33, §34). */6export type SocialEventType =7 | "SESSION_STARTED"8 | "SESSION_ENDED"9 | "PAGE_OPENED"10 | "PAGE_CLASSIFIED"11 | "ENTITY_DISCOVERED"12 | "ENTITY_OBSERVED"13 | "ENTITY_UPDATED"14 | "POST_DISCOVERED"15 | "COMMENT_DISCOVERED"16 | "PROFILE_DISCOVERED"17 | "MEDIA_DISCOVERED"18 | "VIDEO_DISCOVERED"19 | "IMAGE_DISCOVERED"20 | "FEED_ITEM_OBSERVED"21 | "NETWORK_RESPONSE_OBSERVED"22 | "NETWORK_SCHEMA_DISCOVERED"23 | "WEBSOCKET_FRAME_OBSERVED"24 | "DOM_CHANGED"25 | "ACTION_PLANNED"26 | "ACTION_EXECUTED"27 | "ACTION_FAILED"28 | "NAVIGATION_COMPLETED"29 | "LOOP_DETECTED"30 | "BUDGET_EXHAUSTED"31 | "AUTH_REQUIRED"32 | "CONNECTOR_PATTERN_LEARNED"33 | "CONNECTOR_DEGRADED"34 | "CONNECTOR_REPAIRED"35 | "WORKER_ERROR";3637export interface SocialEvent<T = Record<string, unknown>> {38 event_id: string;39 event_type: SocialEventType;40 platform: Platform;41 session_id: string;42 step?: number;43 timestamp: string;44 payload: T;45 provenance?: Provenance[];46 discovered_via?: { action?: string; action_id?: string; url?: string };47}4849export type EventHandler = (ev: SocialEvent) => void | Promise<void>;5051/**52 * Typed in-process event bus. Observers publish, storage / learner / dashboard subscribe.53 * Handlers are awaited sequentially so persistence stays ordered.54 */55export class EventBus {56 private handlers: { type: SocialEventType | "*"; fn: EventHandler }[] = [];57 private queue: Promise<void> = Promise.resolve();58 public count = 0;5960 on(type: SocialEventType | "*", fn: EventHandler): () => void {61 const h = { type, fn };62 this.handlers.push(h);63 return () => {64 this.handlers = this.handlers.filter((x) => x !== h);65 };66 }6768 emit<T extends Record<string, unknown>>(69 partial: Omit<SocialEvent<T>, "event_id" | "timestamp"> & { event_id?: string; timestamp?: string },70 ): SocialEvent<T> {71 const ev: SocialEvent<T> = {72 event_id: partial.event_id ?? newId("ev"),73 timestamp: partial.timestamp ?? nowIso(),74 ...partial,75 };76 this.count++;77 const targets = this.handlers.filter((h) => h.type === "*" || h.type === ev.event_type);78 this.queue = this.queue.then(async () => {79 for (const h of targets) {80 try {81 await h.fn(ev as SocialEvent);82 } catch (err) {83 process.stderr.write(`[events] handler error on ${ev.event_type}: ${(err as Error).message}\n`);84 }85 }86 });87 return ev;88 }8990 /** Wait for all queued handlers (call before shutdown). */91 flush(): Promise<void> {92 return this.queue;93 }94}9596/** Append-only JSONL session log — the raw record every replay (§54) is built from. */97export class JsonlEventLog {98 private stream: fs.WriteStream;99 readonly file: string;100101 constructor(sessionDir: string) {102 fs.mkdirSync(sessionDir, { recursive: true });103 this.file = path.join(sessionDir, "events.jsonl");104 this.stream = fs.createWriteStream(this.file, { flags: "a" });105 }106107 write(ev: SocialEvent): void {108 this.stream.write(JSON.stringify(ev) + "\n");109 }110111 attach(bus: EventBus): () => void {112 return bus.on("*", (ev) => this.write(ev));113 }114115 close(): Promise<void> {116 return new Promise((res) => this.stream.end(res));117 }118119 static read(sessionDir: string): SocialEvent[] {120 const file = path.join(sessionDir, "events.jsonl");121 if (!fs.existsSync(file)) return [];122 return fs123 .readFileSync(file, "utf8")124 .split("\n")125 .filter(Boolean)126 .map((l) => JSON.parse(l) as SocialEvent);127 }128}129