SPB Git forge
7commits 1branches 0releases
229.0 KBsize
maindefault branch
12 days agolast push
TypeScript 91.8% HTML 3.2% JavaScript 3% SQL 1.4% CSS 0.7%
3.6 KB · 129 lines typescript
Raw Blame History
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