From 97a1948866ebad73a90c14a9de47f2744ef8cca1 Mon Sep 17 00:00:00 2001 From: Ezedike-egwom Collins Date: Thu, 27 Aug 2026 07:09:50 +0100 Subject: [PATCH 1/3] feat: emit host-fn log events for real-time subscribers --- src/events.ts | 115 +++++++++++++++++++++++++++++++++++++++++++++++++ src/indexer.ts | 3 +- 2 files changed, 117 insertions(+), 1 deletion(-) diff --git a/src/events.ts b/src/events.ts index d5ffcaa1..22683db5 100644 --- a/src/events.ts +++ b/src/events.ts @@ -1,5 +1,6 @@ import { EventEmitter } from "events"; import type { TransferRecord } from "./db"; +import type { HostFnRecord } from "./indexer/host-fn-log"; // Singleton emitter for real-time transfer notifications. // setMaxListeners(0) removes the default 10-listener cap — one listener per @@ -12,3 +13,117 @@ export type TransferEvent = TransferRecord; export function emitTransfer(transfer: TransferEvent): void { transferEmitter.emit("transfer:new", transfer); } + +// Singleton emitter for real-time host-fn log notifications (GraphQL +// subscriptions, #99). Mirrors transferEmitter above. +export const hostFnLogEmitter = new EventEmitter(); +hostFnLogEmitter.setMaxListeners(0); + +export type HostFnLogEvent = HostFnRecord; + +export function emitHostFnLog(log: HostFnLogEvent): void { + hostFnLogEmitter.emit("hostfnlog:new", log); +} + +/** + * Turn an EventEmitter channel into a pull-based AsyncIterableIterator, the + * shape GraphQL subscription resolvers require. + * + * Bounds memory for a slow consumer: once more than `maxQueue` events have + * piled up waiting to be pulled, the oldest queued event is dropped so a + * stalled subscriber can never grow the queue unboundedly. + */ +export function eventsToAsyncIterator( + emitter: EventEmitter, + eventName: string, + maxQueue = 100 +): AsyncIterableIterator { + const pullQueue: Array<(result: IteratorResult) => void> = []; + const pushQueue: T[] = []; + let listening = true; + + const pushValue = (value: T) => { + if (pullQueue.length > 0) { + const resolve = pullQueue.shift()!; + resolve({ value, done: false }); + return; + } + + pushQueue.push(value); + if (pushQueue.length > maxQueue) { + pushQueue.shift(); + } + }; + + emitter.on(eventName, pushValue); + + const stop = () => { + if (!listening) return; + listening = false; + emitter.off(eventName, pushValue); + while (pullQueue.length > 0) { + pullQueue.shift()!({ value: undefined as unknown as T, done: true }); + } + }; + + const iterator: AsyncIterableIterator = { + next(): Promise> { + if (pushQueue.length > 0) { + return Promise.resolve({ value: pushQueue.shift()!, done: false }); + } + if (!listening) { + return Promise.resolve({ value: undefined as unknown as T, done: true }); + } + return new Promise((resolve) => pullQueue.push(resolve)); + }, + return(): Promise> { + stop(); + return Promise.resolve({ value: undefined as unknown as T, done: true }); + }, + throw(err): Promise> { + stop(); + return Promise.reject(err); + }, + [Symbol.asyncIterator]() { + return iterator; + }, + }; + + return iterator; +} + +/** + * Wrap an async iterator so only values matching `predicate` are yielded, + * used to apply subscription arguments (e.g. contractId) without adding a + * dedicated listener per filter combination. + */ +export function filterAsyncIterator( + iterator: AsyncIterableIterator, + predicate: (value: T) => boolean +): AsyncIterableIterator { + const filtered: AsyncIterableIterator = { + async next(): Promise> { + while (true) { + const result = await iterator.next(); + if (result.done || predicate(result.value)) { + return result; + } + } + }, + return(value?: unknown): Promise> { + return iterator.return + ? (iterator.return(value) as Promise>) + : Promise.resolve({ value: undefined as unknown as T, done: true }); + }, + throw(err): Promise> { + return iterator.throw + ? iterator.throw(err) + : Promise.reject(err); + }, + [Symbol.asyncIterator]() { + return filtered; + }, + }; + + return filtered; +} diff --git a/src/indexer.ts b/src/indexer.ts index 3cdb051c..d21b9da9 100644 --- a/src/indexer.ts +++ b/src/indexer.ts @@ -11,7 +11,7 @@ import { setLastIndexedLedger, pruneOldTransfers, } from "./db"; -import { emitTransfer } from "./events"; +import { emitTransfer, emitHostFnLog } from "./events"; import { parseHostFnEvent, upsertHostFnLogs, type HostFnRecord } from "./indexer/host-fn-log"; import { tagSacTransfers } from "./indexer/sac-detect"; import { pollParallel } from "./indexer/parallel"; @@ -168,6 +168,7 @@ async function pollOnce( await upsertHostFnLogs(hostFnRecords).catch(err => console.error("[indexer] host-fn log error:", err), ); + hostFnRecords.forEach(emitHostFnLog); } // ── NFT path ───────────────────────────────────────────────────────────────── From 5949c11a862bae82bfef593c57cdc2af4f2aec35 Mon Sep 17 00:00:00 2001 From: Ezedike-egwom Collins Date: Thu, 27 Aug 2026 07:10:01 +0100 Subject: [PATCH 2/3] feat: add GraphQL subscriptions over WebSocket --- src/graphql/server.ts | 6 +- src/graphql/subscriptions.ts | 261 +++++++++++++++++++++++++++++++++++ src/index.ts | 5 + src/ws.ts | 3 + 4 files changed, 272 insertions(+), 3 deletions(-) create mode 100644 src/graphql/subscriptions.ts diff --git a/src/graphql/server.ts b/src/graphql/server.ts index 39daa32b..d329ac34 100644 --- a/src/graphql/server.ts +++ b/src/graphql/server.ts @@ -9,7 +9,7 @@ import { import { costLimitPlugin } from "./costLimit"; import { persistedQueryPlugin } from "./persisted"; -const typeDefs = `#graphql +export const typeDefs = `#graphql enum TransferDirection { INCOMING OUTGOING @@ -65,7 +65,7 @@ const typeDefs = `#graphql type TransferDirection = "INCOMING" | "OUTGOING" | "ALL"; -function formatTransfer(row: Record) { +export function formatTransfer(row: Record) { return { ...row, ledgerClosedAt: @@ -75,7 +75,7 @@ function formatTransfer(row: Record) { }; } -const resolvers = { +export const resolvers = { Query: { health: () => ({ ok: true, version: process.env.npm_package_version ?? "1.0.0" }), diff --git a/src/graphql/subscriptions.ts b/src/graphql/subscriptions.ts new file mode 100644 index 00000000..b9d9cff2 --- /dev/null +++ b/src/graphql/subscriptions.ts @@ -0,0 +1,261 @@ +/** + * GraphQL subscriptions over WebSocket (#99). + * + * Streams newly-ingested TokenTransfer / HostFnLog rows to subscribed + * clients, with per-client contract filters and bounded per-subscriber + * queues so a slow consumer can't grow server memory unboundedly. + * + * Built directly on `graphql` + `ws` (both already dependencies) rather + * than the graphql-ws / subscriptions-transport-ws packages, using a small + * JSON message protocol modelled on graphql-ws: + * + * → { id, type: "subscribe", payload: { query, variables? } } + * ← { id, type: "next", payload: } (repeated) + * ← { id, type: "error", payload: [] } + * ← { id, type: "complete" } + * → { id, type: "complete" } (client unsubscribes) + */ +import { + GraphQLObjectType, + GraphQLSchema, + buildSchema, + parse, + subscribe, + validate, + type ExecutionResult, + type GraphQLFieldResolver, +} from "graphql"; +import { WebSocketServer, WebSocket } from "ws"; +import type { IncomingMessage, Server } from "http"; +import { typeDefs as baseTypeDefs, resolvers as baseResolvers, formatTransfer } from "./server"; +import { toDisplayAmount } from "../api"; +import { + transferEmitter, + hostFnLogEmitter, + eventsToAsyncIterator, + filterAsyncIterator, + type TransferEvent, + type HostFnLogEvent, +} from "../events"; + +export const SUBSCRIPTIONS_PATH = "/graphql/subscriptions"; + +// Refuse to enqueue more data on an already-saturated socket — the client +// simply misses the update rather than the server buffering it forever. +const MAX_BUFFERED_BYTES = 1_000_000; + +const subscriptionTypeDefs = `#graphql + scalar JSON + + type HostFnLog { + contractId: String! + functionName: String! + args: JSON + result: JSON + ledger: Int! + ledgerClosedAt: String! + txHash: String! + eventId: String! + } + + type Subscription { + transferAdded(contractId: String): Transfer! + hostFnLogAdded(contractId: String): HostFnLog! + } +`; + +function formatHostFnLog(log: HostFnLogEvent) { + return { + ...log, + ledgerClosedAt: + log.ledgerClosedAt instanceof Date + ? log.ledgerClosedAt.toISOString() + : String(log.ledgerClosedAt), + }; +} + +type FieldResolverMap = Record< + string, + | GraphQLFieldResolver + | { subscribe?: GraphQLFieldResolver; resolve?: GraphQLFieldResolver } +>; + +/** + * graphql-js's `buildSchema` produces a schema with default (identity) + * resolvers only. Attach the real Query resolvers plus our Subscription + * resolvers onto the built schema — the same technique makeExecutableSchema + * uses under the hood, without needing that package as a dependency. + */ +function attachResolvers(schema: GraphQLSchema, resolverMap: Record): void { + for (const [typeName, fields] of Object.entries(resolverMap)) { + const type = schema.getType(typeName); + if (!(type instanceof GraphQLObjectType)) continue; + + const typeFields = type.getFields(); + for (const [fieldName, fieldResolver] of Object.entries(fields)) { + const field = typeFields[fieldName]; + if (!field) continue; + + if (typeof fieldResolver === "function") { + field.resolve = fieldResolver; + } else { + if (fieldResolver.subscribe) field.subscribe = fieldResolver.subscribe; + if (fieldResolver.resolve) field.resolve = fieldResolver.resolve; + } + } + } +} + +function buildSubscriptionSchema(): GraphQLSchema { + const schema = buildSchema(baseTypeDefs + subscriptionTypeDefs); + + attachResolvers(schema, baseResolvers); + attachResolvers(schema, { + Subscription: { + transferAdded: { + subscribe: (_parent, args: { contractId?: string }) => + filterAsyncIterator( + eventsToAsyncIterator(transferEmitter, "transfer:new"), + (t) => !args.contractId || t.contractId === args.contractId + ), + resolve: (payload: unknown) => { + const transfer = payload as TransferEvent; + return { + ...formatTransfer(transfer as unknown as Record), + displayAmount: toDisplayAmount(transfer.amount), + }; + }, + }, + hostFnLogAdded: { + subscribe: (_parent, args: { contractId?: string }) => + filterAsyncIterator( + eventsToAsyncIterator(hostFnLogEmitter, "hostfnlog:new"), + (l) => !args.contractId || l.contractId === args.contractId + ), + resolve: (payload: unknown) => formatHostFnLog(payload as HostFnLogEvent), + }, + }, + }); + + return schema; +} + +interface ClientMessage { + id?: string; + type?: "subscribe" | "complete"; + payload?: { query?: string; variables?: Record }; +} + +/** + * Attach a WebSocket server implementing GraphQL subscriptions. + * + * Clients connect to: ws://host/graphql/subscriptions + */ +export function attachGraphQLSubscriptions(server: Server): void { + const schema = buildSubscriptionSchema(); + const wss = new WebSocketServer({ noServer: true }); + + server.on("upgrade", (req: IncomingMessage, socket, head) => { + const url = req.url ?? ""; + // Owns the whole /graphql/* upgrade namespace (ws.ts explicitly defers + // it here) so any unmatched sub-path still gets a clean 404 instead of + // hanging with no response. + if (!url.startsWith("/graphql/")) return; + if (!url.startsWith(SUBSCRIPTIONS_PATH)) { + socket.write("HTTP/1.1 404 Not Found\r\n\r\n"); + socket.destroy(); + return; + } + + wss.handleUpgrade(req, socket, head, (ws) => { + wss.emit("connection", ws, req); + }); + }); + + wss.on("connection", (ws: WebSocket) => { + // One entry per active subscription id on this connection, so a + // "complete" message or socket close can release its async iterator. + const active = new Map>(); + + const send = (msg: Record) => { + if (ws.readyState !== WebSocket.OPEN) return; + if (ws.bufferedAmount > MAX_BUFFERED_BYTES) return; + ws.send(JSON.stringify(msg)); + }; + + const stop = (id: string) => { + const iterator = active.get(id); + active.delete(id); + iterator?.return?.(undefined); + }; + + ws.on("message", (data: Buffer | string) => { + let msg: ClientMessage; + try { + msg = JSON.parse(data.toString()); + } catch { + send({ type: "error", payload: [{ message: "Invalid JSON" }] }); + return; + } + + if (msg.type === "complete") { + if (msg.id) stop(msg.id); + return; + } + + if (msg.type !== "subscribe" || !msg.id || !msg.payload?.query) return; + const { id } = msg; + const { query, variables } = msg.payload; + + void (async () => { + let document; + try { + document = parse(query); + } catch (err) { + send({ id, type: "error", payload: [{ message: (err as Error).message }] }); + return; + } + + const validationErrors = validate(schema, document); + if (validationErrors.length > 0) { + send({ id, type: "error", payload: validationErrors.map((e) => ({ message: e.message })) }); + return; + } + + const result = await subscribe({ schema, document, variableValues: variables }); + + if (!(Symbol.asyncIterator in result)) { + send({ + id, + type: "error", + payload: (result as ExecutionResult).errors ?? [{ message: "Subscription failed" }], + }); + return; + } + + const iterator = result as AsyncIterableIterator; + active.set(id, iterator); + + try { + for await (const event of iterator) { + if (!active.has(id)) break; // client unsubscribed mid-stream + send({ id, type: "next", payload: event }); + } + if (active.has(id)) { + active.delete(id); + send({ id, type: "complete" }); + } + } catch (err) { + active.delete(id); + send({ id, type: "error", payload: [{ message: (err as Error).message }] }); + } + })(); + }); + + const cleanup = () => { + for (const id of [...active.keys()]) stop(id); + }; + ws.on("close", cleanup); + ws.on("error", cleanup); + }); +} diff --git a/src/index.ts b/src/index.ts index c8c9f703..06e63c97 100644 --- a/src/index.ts +++ b/src/index.ts @@ -5,6 +5,7 @@ import { createApp } from "./api"; import { startIndexer } from "./indexer"; import { prisma } from "./db"; import { attachWebSocketServer } from "./ws"; +import { attachGraphQLSubscriptions, SUBSCRIPTIONS_PATH } from "./graphql/subscriptions"; import { startWebhookWorker } from "./workers/webhooks"; import { startPartitionRetentionJob } from "./jobs/retention"; @@ -33,9 +34,13 @@ async function main() { // Attach WebSocket upgrade handler — clients connect to /subscribe/:address attachWebSocketServer(server); + // Attach GraphQL subscriptions — clients connect to /graphql/subscriptions + attachGraphQLSubscriptions(server); + server.listen(PORT, () => { console.log(`[wraith] API listening on http://localhost:${PORT}`); console.log(`[wraith] WebSocket subscriptions available at ws://localhost:${PORT}/subscribe/:address`); + console.log(`[wraith] GraphQL subscriptions available at ws://localhost:${PORT}${SUBSCRIPTIONS_PATH}`); }); // ── Start webhook worker ─────────────────────────────────────────────────── diff --git a/src/ws.ts b/src/ws.ts index c811fd96..4f84bc15 100644 --- a/src/ws.ts +++ b/src/ws.ts @@ -27,6 +27,9 @@ export function attachWebSocketServer(server: Server): void { server.on("upgrade", (req: IncomingMessage, socket, head) => { const url = req.url ?? ""; + // GraphQL subscriptions own /graphql/* upgrades — leave those to + // attachGraphQLSubscriptions() rather than 404'ing them here. + if (url.startsWith("/graphql/")) return; if (!SUBSCRIBE_RE.test(url)) { socket.write("HTTP/1.1 404 Not Found\r\n\r\n"); socket.destroy(); From a39e802d9fdb6113a5856673b2da4e77f6b7586c Mon Sep 17 00:00:00 2001 From: Ezedike-egwom Collins Date: Thu, 27 Aug 2026 07:10:05 +0100 Subject: [PATCH 3/3] test: cover GraphQL subscriptions over WebSocket --- src/graphql/__tests__/subscriptions.test.ts | 230 ++++++++++++++++++++ 1 file changed, 230 insertions(+) create mode 100644 src/graphql/__tests__/subscriptions.test.ts diff --git a/src/graphql/__tests__/subscriptions.test.ts b/src/graphql/__tests__/subscriptions.test.ts new file mode 100644 index 00000000..a71ae65b --- /dev/null +++ b/src/graphql/__tests__/subscriptions.test.ts @@ -0,0 +1,230 @@ +/** + * GraphQL subscription tests — issue #99 + * + * Spins up a real HTTP + WS server in-process (no Docker) and exercises: + * 1. connect / subscribe / receive / unsubscribe cycle for transferAdded + * 2. contractId filter on transferAdded + * 3. hostFnLogAdded delivers newly-logged host-fn invocations + * 4. slow consumer doesn't crash the server (bounded queue) + */ +import http from "http"; +import { WebSocket } from "ws"; +import { attachGraphQLSubscriptions, SUBSCRIPTIONS_PATH } from "../subscriptions"; +import { emitTransfer, emitHostFnLog } from "../../events"; +import type { TransferEvent, HostFnLogEvent } from "../../events"; + +jest.mock("../../db", () => ({ + queryAllTransfers: jest.fn().mockResolvedValue({ total: 0, transfers: [], nextCursor: null }), + queryByTxHash: jest.fn().mockResolvedValue([]), + querySummary: jest.fn().mockResolvedValue([]), + queryTransfers: jest.fn().mockResolvedValue({ total: 0, transfers: [], nextCursor: null }), +})); + +function makeTransfer(overrides: Partial = {}): TransferEvent { + return { + contractId: "CTOKEN", + eventType: "transfer", + fromAddress: "GSENDER", + toAddress: "GRECV", + amount: "10000000", + ledger: 100, + ledgerClosedAt: new Date("2025-01-01T00:00:00Z"), + txHash: "txhash", + eventId: "ev-1", + ...overrides, + } as TransferEvent; +} + +function makeHostFnLog(overrides: Partial = {}): HostFnLogEvent { + return { + contractId: "CTOKEN", + functionName: "swap", + args: ["a", "b"], + result: { ok: true }, + gasUsed: null, + ledger: 100, + ledgerClosedAt: new Date("2025-01-01T00:00:00Z"), + txHash: "txhash", + eventId: "hfl-1", + ...overrides, + }; +} + +async function startServer(): Promise<{ url: string; close: () => Promise }> { + const server = http.createServer(); + attachGraphQLSubscriptions(server); + + await new Promise((resolve) => server.listen(0, resolve)); + const { port } = server.address() as { port: number }; + + return { + url: `ws://localhost:${port}${SUBSCRIPTIONS_PATH}`, + close: () => + new Promise((resolve, reject) => + server.close((err) => (err ? reject(err) : resolve())) + ), + }; +} + +// Tracked so afterAll can force-close any socket a failed assertion left +// open — otherwise http.Server#close() hangs waiting for it and blows the +// hook timeout instead of reporting the real test failure. +const openSockets: WebSocket[] = []; + +function connect(url: string): Promise { + return new Promise((resolve, reject) => { + const ws = new WebSocket(url); + openSockets.push(ws); + ws.once("open", () => resolve(ws)); + ws.once("error", reject); + }); +} + +function subscribeOp( + ws: WebSocket, + id: string, + query: string, + variables?: Record +): void { + ws.send(JSON.stringify({ id, type: "subscribe", payload: { query, variables } })); +} + +function collectNext(ws: WebSocket, id: string, n: number): Promise[]> { + return new Promise((resolve, reject) => { + const results: Record[] = []; + const handler = (data: Buffer) => { + const msg = JSON.parse(data.toString()); + if (msg.id !== id) return; + if (msg.type === "next") { + results.push(msg.payload.data); + if (results.length >= n) { + ws.off("message", handler); + resolve(results); + } + } else if (msg.type === "error") { + ws.off("message", handler); + reject(new Error(JSON.stringify(msg.payload))); + } + }; + ws.on("message", handler); + }); +} + +describe("GraphQL subscriptions /graphql/subscriptions", () => { + let serverUrl: string; + let closeServer: () => Promise; + + beforeAll(async () => { + const srv = await startServer(); + serverUrl = srv.url; + closeServer = srv.close; + }); + + afterAll(async () => { + for (const ws of openSockets.splice(0)) { + if (ws.readyState === WebSocket.OPEN || ws.readyState === WebSocket.CONNECTING) { + ws.terminate(); + } + } + await closeServer(); + }); + + it("rejects upgrade on an unknown /graphql sub-path with 404", async () => { + const other = serverUrl.replace(SUBSCRIPTIONS_PATH, "/graphql/nope"); + await expect(connect(other)).rejects.toThrow(); + }); + + it("streams a new transfer to a transferAdded subscriber", async () => { + const ws = await connect(serverUrl); + subscribeOp(ws, "1", "subscription { transferAdded { contractId eventId displayAmount } }"); + + // Give the server a tick to register the subscription's async iterator + // before emitting, mirroring the existing WS subscription test. + await new Promise((r) => setTimeout(r, 20)); + + const pending = collectNext(ws, "1", 1); + emitTransfer(makeTransfer({ eventId: "ev-a" })); + + const [msg] = await pending; + const transfer = msg.transferAdded as Record; + expect(transfer.eventId).toBe("ev-a"); + expect(transfer.displayAmount).toBe("1.0000000"); + + ws.close(); + }); + + it("filters transferAdded by contractId", async () => { + const ws = await connect(serverUrl); + subscribeOp( + ws, + "2", + "subscription($contractId: String) { transferAdded(contractId: $contractId) { contractId eventId } }", + { contractId: "CWANTED" } + ); + await new Promise((r) => setTimeout(r, 20)); + + const pending = collectNext(ws, "2", 1); + + emitTransfer(makeTransfer({ contractId: "COTHER", eventId: "ev-skip" })); + emitTransfer(makeTransfer({ contractId: "CWANTED", eventId: "ev-match" })); + + const [msg] = await pending; + expect((msg.transferAdded as Record).eventId).toBe("ev-match"); + + ws.close(); + }); + + it("streams hostFnLogAdded events", async () => { + const ws = await connect(serverUrl); + subscribeOp(ws, "3", "subscription { hostFnLogAdded { contractId functionName eventId } }"); + await new Promise((r) => setTimeout(r, 20)); + + const pending = collectNext(ws, "3", 1); + emitHostFnLog(makeHostFnLog({ eventId: "hfl-a" })); + + const [msg] = await pending; + expect((msg.hostFnLogAdded as Record).eventId).toBe("hfl-a"); + expect((msg.hostFnLogAdded as Record).functionName).toBe("swap"); + + ws.close(); + }); + + it("stops delivering after the client sends complete", async () => { + const ws = await connect(serverUrl); + subscribeOp(ws, "4", "subscription { transferAdded { eventId } }"); + await new Promise((r) => setTimeout(r, 20)); + + let received = 0; + ws.on("message", (data) => { + const msg = JSON.parse(data.toString()); + if (msg.id === "4" && msg.type === "next") received++; + }); + + ws.send(JSON.stringify({ id: "4", type: "complete" })); + await new Promise((r) => setTimeout(r, 20)); + + emitTransfer(makeTransfer({ eventId: "ev-after-complete" })); + await new Promise((r) => setTimeout(r, 30)); + + expect(received).toBe(0); + ws.close(); + }); + + it("backpressure — buffers many rapid transfers without crashing the server", async () => { + const ws = await connect(serverUrl); + subscribeOp(ws, "5", "subscription { transferAdded { eventId } }"); + await new Promise((r) => setTimeout(r, 20)); + + const COUNT = 50; + const pending = collectNext(ws, "5", COUNT); + + for (let i = 0; i < COUNT; i++) { + emitTransfer(makeTransfer({ eventId: `bp-${i}` })); + } + + const msgs = await pending; + expect(msgs).toHaveLength(COUNT); + + ws.close(); + }); +});