Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
115 changes: 115 additions & 0 deletions src/events.ts
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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<T>(
emitter: EventEmitter,
eventName: string,
maxQueue = 100
): AsyncIterableIterator<T> {
const pullQueue: Array<(result: IteratorResult<T>) => 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<T> = {
next(): Promise<IteratorResult<T>> {
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<IteratorResult<T>> {
stop();
return Promise.resolve({ value: undefined as unknown as T, done: true });
},
throw(err): Promise<IteratorResult<T>> {
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<T>(
iterator: AsyncIterableIterator<T>,
predicate: (value: T) => boolean
): AsyncIterableIterator<T> {
const filtered: AsyncIterableIterator<T> = {
async next(): Promise<IteratorResult<T>> {
while (true) {
const result = await iterator.next();
if (result.done || predicate(result.value)) {
return result;
}
}
},
return(value?: unknown): Promise<IteratorResult<T>> {
return iterator.return
? (iterator.return(value) as Promise<IteratorResult<T>>)
: Promise.resolve({ value: undefined as unknown as T, done: true });
},
throw(err): Promise<IteratorResult<T>> {
return iterator.throw
? iterator.throw(err)
: Promise.reject(err);
},
[Symbol.asyncIterator]() {
return filtered;
},
};

return filtered;
}
230 changes: 230 additions & 0 deletions src/graphql/__tests__/subscriptions.test.ts
Original file line number Diff line number Diff line change
@@ -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> = {}): 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> = {}): 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<void> }> {
const server = http.createServer();
attachGraphQLSubscriptions(server);

await new Promise<void>((resolve) => server.listen(0, resolve));
const { port } = server.address() as { port: number };

return {
url: `ws://localhost:${port}${SUBSCRIPTIONS_PATH}`,
close: () =>
new Promise<void>((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<WebSocket> {
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<string, unknown>
): void {
ws.send(JSON.stringify({ id, type: "subscribe", payload: { query, variables } }));
}

function collectNext(ws: WebSocket, id: string, n: number): Promise<Record<string, unknown>[]> {
return new Promise((resolve, reject) => {
const results: Record<string, unknown>[] = [];
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<void>;

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<string, unknown>;
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<string, unknown>).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<string, unknown>).eventId).toBe("hfl-a");
expect((msg.hostFnLogAdded as Record<string, unknown>).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();
});
});
6 changes: 3 additions & 3 deletions src/graphql/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import {
import { costLimitPlugin } from "./costLimit";
import { persistedQueryPlugin } from "./persisted";

const typeDefs = `#graphql
export const typeDefs = `#graphql
enum TransferDirection {
INCOMING
OUTGOING
Expand Down Expand Up @@ -65,7 +65,7 @@ const typeDefs = `#graphql

type TransferDirection = "INCOMING" | "OUTGOING" | "ALL";

function formatTransfer(row: Record<string, unknown>) {
export function formatTransfer(row: Record<string, unknown>) {
return {
...row,
ledgerClosedAt:
Expand All @@ -75,7 +75,7 @@ function formatTransfer(row: Record<string, unknown>) {
};
}

const resolvers = {
export const resolvers = {
Query: {
health: () => ({ ok: true, version: process.env.npm_package_version ?? "1.0.0" }),

Expand Down
Loading
Loading