|
| 1 | +import { Client } from "npm:pg@^8.13.0" |
| 2 | +import { FlowcoreDataPumpCluster } from "../../src/data-pump/data-pump-cluster.ts" |
| 3 | +import { PostgresCoordinator } from "./postgres-coordinator.ts" |
| 4 | +import { PostgresStateManager } from "./postgres-state-manager.ts" |
| 5 | +import { FakeDataSource } from "./fake-data-source.ts" |
| 6 | + |
| 7 | +const DATABASE_URL = Deno.env.get("DATABASE_URL") ?? "postgres://postgres:postgres@localhost:5432/datapump_test" |
| 8 | +const POD_NAME = Deno.env.get("POD_NAME") ?? "local-pod" |
| 9 | +const POD_IP = Deno.env.get("POD_IP") ?? "127.0.0.1" |
| 10 | +const TOTAL_EVENTS = parseInt(Deno.env.get("TOTAL_EVENTS") ?? "100", 10) |
| 11 | +const WS_PORT = parseInt(Deno.env.get("WS_PORT") ?? "8080", 10) |
| 12 | + |
| 13 | +const log = { |
| 14 | + debug: (msg: string, meta?: Record<string, unknown>) => console.log(`[DEBUG] [${POD_NAME}] ${msg}`, meta ?? ""), |
| 15 | + info: (msg: string, meta?: Record<string, unknown>) => console.log(`[INFO] [${POD_NAME}] ${msg}`, meta ?? ""), |
| 16 | + warn: (msg: string, meta?: Record<string, unknown>) => console.warn(`[WARN] [${POD_NAME}] ${msg}`, meta ?? ""), |
| 17 | + error: (msg: string | Error, meta?: Record<string, unknown>) => |
| 18 | + console.error(`[ERROR] [${POD_NAME}] ${msg}`, meta ?? ""), |
| 19 | +} |
| 20 | + |
| 21 | +// connect to PG |
| 22 | +const db = new Client({ connectionString: DATABASE_URL }) |
| 23 | +await db.connect() |
| 24 | +log.info("Connected to PostgreSQL") |
| 25 | + |
| 26 | +// create tables |
| 27 | +await db.query(` |
| 28 | + CREATE TABLE IF NOT EXISTS flowcore_pump_leases ( |
| 29 | + key TEXT PRIMARY KEY, |
| 30 | + holder TEXT NOT NULL, |
| 31 | + expires_at TIMESTAMPTZ NOT NULL |
| 32 | + ); |
| 33 | +
|
| 34 | + CREATE TABLE IF NOT EXISTS flowcore_pump_instances ( |
| 35 | + instance_id TEXT PRIMARY KEY, |
| 36 | + address TEXT NOT NULL, |
| 37 | + last_heartbeat TIMESTAMPTZ NOT NULL DEFAULT NOW() |
| 38 | + ); |
| 39 | +
|
| 40 | + CREATE TABLE IF NOT EXISTS pump_state ( |
| 41 | + id TEXT PRIMARY KEY, |
| 42 | + time_bucket TEXT NOT NULL, |
| 43 | + event_id TEXT, |
| 44 | + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() |
| 45 | + ); |
| 46 | +
|
| 47 | + CREATE TABLE IF NOT EXISTS processed_events ( |
| 48 | + id SERIAL PRIMARY KEY, |
| 49 | + pod_name TEXT NOT NULL, |
| 50 | + event_id TEXT NOT NULL, |
| 51 | + processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() |
| 52 | + ); |
| 53 | +`) |
| 54 | +log.info("Database tables ready") |
| 55 | + |
| 56 | +// create components |
| 57 | +const coordinator = new PostgresCoordinator(db) |
| 58 | +const stateManager = new PostgresStateManager(db) |
| 59 | +const fakeDataSource = new FakeDataSource(TOTAL_EVENTS) |
| 60 | + |
| 61 | +// create cluster |
| 62 | +const cluster = new FlowcoreDataPumpCluster({ |
| 63 | + auth: { getBearerToken: () => Promise.resolve("fake") }, |
| 64 | + dataSource: { |
| 65 | + tenant: "integration-test", |
| 66 | + dataCore: "test-data-core", |
| 67 | + flowType: "test-flow-type", |
| 68 | + eventTypes: ["test-event"], |
| 69 | + }, |
| 70 | + stateManager, |
| 71 | + coordinator, |
| 72 | + dataSourceOverride: fakeDataSource, |
| 73 | + advertisedAddress: `ws://${POD_IP}:${WS_PORT}`, |
| 74 | + noTranslation: true, |
| 75 | + processor: { |
| 76 | + concurrency: 5, |
| 77 | + handler: async (events) => { |
| 78 | + for (const event of events) { |
| 79 | + await db.query(`INSERT INTO processed_events (pod_name, event_id) VALUES ($1, $2)`, [ |
| 80 | + POD_NAME, |
| 81 | + event.eventId, |
| 82 | + ]) |
| 83 | + } |
| 84 | + log.info(`Processed ${events.length} events`) |
| 85 | + }, |
| 86 | + }, |
| 87 | + notifier: { type: "poller", intervalMs: 2000 }, |
| 88 | + leaseTtlMs: 15000, |
| 89 | + leaseRenewIntervalMs: 5000, |
| 90 | + heartbeatIntervalMs: 3000, |
| 91 | + workerConcurrency: 5, |
| 92 | + logger: log, |
| 93 | +}) |
| 94 | + |
| 95 | +// start HTTP + WS server |
| 96 | +Deno.serve({ port: WS_PORT }, (req) => { |
| 97 | + const url = new URL(req.url) |
| 98 | + |
| 99 | + if (url.pathname === "/health") { |
| 100 | + return new Response( |
| 101 | + JSON.stringify({ |
| 102 | + instanceId: cluster.id, |
| 103 | + isLeader: cluster.isLeaderInstance, |
| 104 | + workerCount: cluster.activeWorkerCount, |
| 105 | + isRunning: cluster.isRunning, |
| 106 | + podName: POD_NAME, |
| 107 | + }), |
| 108 | + { headers: { "content-type": "application/json" } }, |
| 109 | + ) |
| 110 | + } |
| 111 | + |
| 112 | + if (req.headers.get("upgrade")?.toLowerCase() === "websocket") { |
| 113 | + const { socket, response } = Deno.upgradeWebSocket(req) |
| 114 | + cluster.handleConnection(socket) |
| 115 | + return response |
| 116 | + } |
| 117 | + |
| 118 | + return new Response("Not found", { status: 404 }) |
| 119 | +}) |
| 120 | + |
| 121 | +log.info(`HTTP/WS server listening on :${WS_PORT}`) |
| 122 | + |
| 123 | +// start cluster |
| 124 | +await cluster.start() |
| 125 | +log.info("Cluster started") |
| 126 | + |
| 127 | +// graceful shutdown |
| 128 | +Deno.addSignalListener("SIGTERM", async () => { |
| 129 | + log.info("SIGTERM received, shutting down...") |
| 130 | + await cluster.stop() |
| 131 | + await db.end() |
| 132 | + Deno.exit(0) |
| 133 | +}) |
0 commit comments