diff --git a/README.md b/README.md index c751773..dfa1da7 100644 --- a/README.md +++ b/README.md @@ -21,17 +21,19 @@ clean checkout without any database or other external dependency. src/ index.js # HTTP bootstrap: Express app + /graphql endpoint server.js # Apollo Server factory (reused by the tests) - schema.js # GraphQL type definitions - resolvers.js # Query / Mutation / field resolvers + schema.js # GraphQL type definitions and executable schema factory + resolvers.js # Query / Mutation / Subscription / field resolvers config/finance.js # Environment-driven finance connector config connectors/ # Replaceable OpenTrading, Portfolio-Watcher, tax-break adapters data/store.js # In-memory data store with seed data domain/finance.js # Canonical finance models and normalization helpers export/ # Excel (.xlsx) writer, dataset registry, and /export routes services/financeService.js # Finance aggregation, caching, and error handling + streaming/ # In-process pub/sub and the SSE streaming endpoint test/ graphql.test.js # API tests executed against the schema export.test.js # Excel export writer, dataset, and route tests + streaming.test.js # Pub/sub and Server-Sent Events streaming tests website/ src/App.jsx # Learning site: primer, tips, API Explorer src/backendSamples.js # GraphQL server samples in 10 backend languages @@ -63,6 +65,7 @@ environment variable): - Readiness check (per-upstream): - Metrics (Prometheus text): - Excel exports: +- Streaming (SSE) endpoint: Opening the GraphQL endpoint in a browser loads the Apollo Sandbox, where you can explore the schema and run the operations below. @@ -322,6 +325,8 @@ Follow-up production tasks: | `taxEstimate(taxYear, accountId, symbol, from, to, limit, offset)` | Estimate tax from tax-relevant trading activity | | `createUser(name, email)` | Create a user | | `createPost(title, content, authorId)` | Create a post for an existing user | +| `userCreated` | Subscription: streams every newly created user | +| `postCreated(authorId)` | Subscription: streams new posts, optionally for one author | ### Example queries @@ -401,6 +406,39 @@ curl http://localhost:4000/graphql \ -d '{"query":"{ users { id name posts { title } } }"}' ``` +## Streaming + +Subscriptions (and any query or mutation) can be streamed over Server-Sent +Events at `POST|GET /graphql/stream`, which takes the same +`query`/`variables`/`operationName` payload as `/graphql`. SSE keeps the +transport plain HTTP — no WebSocket upgrade or extra service is required. + +```bash +# Stream new users as they are created +curl -N -H 'Accept: text/event-stream' \ + --get http://localhost:4000/graphql/stream \ + --data-urlencode 'query=subscription { userCreated { id name email } }' + +# In another shell, trigger an event +curl -X POST http://localhost:4000/graphql \ + -H 'Content-Type: application/json' \ + -d '{"query":"mutation { createUser(name: \"Grace\", email: \"grace@example.com\") { id } }"}' +``` + +Each emitted result arrives as an `event: next` frame carrying the usual +GraphQL response body, and the stream finishes with `event: complete`. +Queries and mutations sent to the same endpoint produce exactly one `next` +frame followed by `complete`, so a single client transport covers every +operation type. Comment frames (`: ping`) act as heartbeats so idle +connections survive proxies, and a subscription is torn down as soon as the +client disconnects. + +Events are dispatched by an in-process pub/sub (`src/streaming/pubsub.js`). +It is intentionally dependency-free, which means subscribers only see events +published by their own instance; swap that module for a Redis/NATS-backed +implementation with the same `publish`/`subscribe` contract to run more than +one replica. + ## Excel export Any dataset the API serves can be downloaded as an Excel workbook (`.xlsx`) — diff --git a/package-lock.json b/package-lock.json index fb53dd2..86134ce 100644 --- a/package-lock.json +++ b/package-lock.json @@ -11,6 +11,7 @@ "dependencies": { "@apollo/server": "^5.5.1", "@as-integrations/express4": "^1.1.2", + "@graphql-tools/schema": "^10.1.1", "cors": "^2.8.5", "express": "^4.21.2", "graphql": "^16.10.0" @@ -451,12 +452,12 @@ } }, "node_modules/@graphql-tools/merge": { - "version": "9.2.3", - "resolved": "https://registry.npmjs.org/@graphql-tools/merge/-/merge-9.2.3.tgz", - "integrity": "sha512-cKRoXqJGy2zSRBLvotQpkACbXlHAb0yuHLN0l0ypKGCuL3NnF0zofYalkBJqiBxxxpK0lr3s9uowKFHCKi1/ZQ==", + "version": "9.2.4", + "resolved": "https://registry.npmjs.org/@graphql-tools/merge/-/merge-9.2.4.tgz", + "integrity": "sha512-vV+8uWNWn0+OsqT0r22lZuoT0cTe6fBqBtpLHre2rriLjI/ZTrsOHmabLPTUOyzgFnRrS2PzFSYGL3zKL1Wj4Q==", "license": "MIT", "dependencies": { - "@graphql-tools/utils": "^12.0.0", + "@graphql-tools/utils": "^12.0.1", "tslib": "^2.4.0" }, "engines": { @@ -467,13 +468,13 @@ } }, "node_modules/@graphql-tools/schema": { - "version": "10.1.0", - "resolved": "https://registry.npmjs.org/@graphql-tools/schema/-/schema-10.1.0.tgz", - "integrity": "sha512-wao48XQnfY631s3jXoNrhEHvCI8mlKXmIuWrR7F6zAdv92VuSOfHoq9P9KL2EnUMgBUnaStnByOx9Mn6RieWDg==", + "version": "10.1.1", + "resolved": "https://registry.npmjs.org/@graphql-tools/schema/-/schema-10.1.1.tgz", + "integrity": "sha512-24jJghRxEW+SG1lbJ45Zg9HEZ8ZKS923ihbFjOicspFLpryxJkWp6gZJ7D8tCVaSAhli4x1+1wDPAfhT8mLqxg==", "license": "MIT", "dependencies": { - "@graphql-tools/merge": "^9.2.3", - "@graphql-tools/utils": "^12.0.0", + "@graphql-tools/merge": "^9.2.4", + "@graphql-tools/utils": "^12.0.1", "tslib": "^2.4.0" }, "engines": { @@ -484,9 +485,9 @@ } }, "node_modules/@graphql-tools/utils": { - "version": "12.0.0", - "resolved": "https://registry.npmjs.org/@graphql-tools/utils/-/utils-12.0.0.tgz", - "integrity": "sha512-aMdIo/l+8j4lhamWAf+MWHr+lOW8zArSBKXspqtKHgt5I2JF9LS55KDLUe/L9AiJOV/O6qJJ60hgYydbnJHbxg==", + "version": "12.0.1", + "resolved": "https://registry.npmjs.org/@graphql-tools/utils/-/utils-12.0.1.tgz", + "integrity": "sha512-8YC6jn4xYDS6YTY6xDAgQAs/nGgvzRaIGHS8GeSfVItQzNK7Ms24CQtE3EJY+amvR+tBnhRSX0N83aGj+V/RIg==", "license": "MIT", "dependencies": { "@graphql-typed-document-node/core": "^3.1.1", diff --git a/package.json b/package.json index dce9e40..35a4586 100644 --- a/package.json +++ b/package.json @@ -22,6 +22,7 @@ "dependencies": { "@apollo/server": "^5.5.1", "@as-integrations/express4": "^1.1.2", + "@graphql-tools/schema": "^10.1.1", "cors": "^2.8.5", "express": "^4.21.2", "graphql": "^16.10.0" diff --git a/src/index.js b/src/index.js index a5d83c7..8d3013b 100644 --- a/src/index.js +++ b/src/index.js @@ -6,14 +6,18 @@ import { store } from './data/store.js'; import { createExportRouter } from './export/router.js'; import { logger } from './observability/logger.js'; import { metrics } from './observability/metrics.js'; +import { createSchema } from './schema.js'; import { financeService } from './services/financeService.js'; import { createApolloServer } from './server.js'; +import { pubsub } from './streaming/pubsub.js'; +import { createStreamRouter } from './streaming/sseRouter.js'; const PORT = Number(process.env.PORT) || 4000; /** Starts the HTTP server exposing the GraphQL endpoint at /graphql. */ async function main() { - const apolloServer = createApolloServer({ logger, metrics, enableObservability: true }); + const schema = createSchema(); + const apolloServer = createApolloServer({ schema, logger, metrics, enableObservability: true }); await apolloServer.start(); const app = express(); @@ -36,18 +40,33 @@ async function main() { // Spreadsheet downloads (.xlsx) for the same data the GraphQL API serves. app.use('/export', cors(), createExportRouter({ store, finance: financeService, logger })); + // Streaming endpoint (Server-Sent Events) for subscriptions and one-shot + // operations, mounted before /graphql so it keeps its own body parsing. + app.use( + '/graphql/stream', + cors(), + express.json(), + createStreamRouter({ + schema, + logger, + metrics, + contextValue: () => ({ store, finance: financeService, logger, metrics, pubsub }), + }) + ); + app.use( '/graphql', cors(), express.json(), expressMiddleware(apolloServer, { // Every request shares the same in-memory store. - context: async () => ({ store, finance: financeService, logger, metrics }), + context: async () => ({ store, finance: financeService, logger, metrics, pubsub }), }) ); await new Promise((resolve) => app.listen(PORT, resolve)); logger.info('graphql endpoint ready', { url: `http://localhost:${PORT}/graphql` }); + logger.info('graphql stream endpoint ready', { url: `http://localhost:${PORT}/graphql/stream` }); } main().catch((error) => { diff --git a/src/resolvers.js b/src/resolvers.js index a95bd85..fb18d10 100644 --- a/src/resolvers.js +++ b/src/resolvers.js @@ -1,7 +1,15 @@ import { GraphQLError } from 'graphql'; +import { pubsub as defaultPubSub, TOPICS } from './streaming/pubsub.js'; import { validateFinanceArgs } from './validation/financeArgs.js'; +/** Wraps an async iterator, forwarding only the payloads matching `predicate`. */ +async function* filterAsyncIterator(source, predicate) { + for await (const payload of source) { + if (predicate(payload)) yield payload; + } +} + async function runFinanceResolver(name, args, context, resolve, options = {}) { const issues = validateFinanceArgs(args, options); if (issues.length > 0) { @@ -46,9 +54,13 @@ export const resolvers = { }, Mutation: { - createUser: (_parent, { name, email }, { store }) => store.createUser({ name, email }), + createUser: (_parent, { name, email }, { store, pubsub = defaultPubSub }) => { + const user = store.createUser({ name, email }); + pubsub.publish(TOPICS.USER_CREATED, { userCreated: user }); + return user; + }, - createPost: (_parent, { title, content, authorId }, { store }) => { + createPost: (_parent, { title, content, authorId }, { store, pubsub = defaultPubSub }) => { // Referential integrity has to be enforced by hand with an in-memory store. if (!store.getUser(authorId)) { throw new GraphQLError(`User with id "${authorId}" does not exist.`, { @@ -56,7 +68,23 @@ export const resolvers = { }); } - return store.createPost({ title, content, authorId }); + const post = store.createPost({ title, content, authorId }); + pubsub.publish(TOPICS.POST_CREATED, { postCreated: post }); + return post; + }, + }, + + Subscription: { + userCreated: { + subscribe: (_parent, _args, { pubsub = defaultPubSub, signal }) => pubsub.subscribe(TOPICS.USER_CREATED, { signal }), + }, + + postCreated: { + subscribe: (_parent, { authorId }, { pubsub = defaultPubSub, signal }) => { + const source = pubsub.subscribe(TOPICS.POST_CREATED, { signal }); + if (!authorId) return source; + return filterAsyncIterator(source, (payload) => payload?.postCreated?.authorId === authorId); + }, }, }, diff --git a/src/schema.js b/src/schema.js index 77e198c..02b826d 100644 --- a/src/schema.js +++ b/src/schema.js @@ -1,3 +1,7 @@ +import { makeExecutableSchema } from '@graphql-tools/schema'; + +import { resolvers } from './resolvers.js'; + /** GraphQL schema definition for the demo service. */ export const typeDefs = /* GraphQL */ ` "A person who can author posts." @@ -208,4 +212,20 @@ export const typeDefs = /* GraphQL */ ` createUser(name: String!, email: String!): User! createPost(title: String!, content: String!, authorId: ID!): Post! } + + "Streaming operations delivered over Server-Sent Events at /graphql/stream." + type Subscription { + "Emitted whenever a user is created." + userCreated: User! + "Emitted whenever a post is created, optionally filtered by author." + postCreated(authorId: ID): Post! + } `; + +/** + * Builds the executable schema. Apollo Server and the SSE streaming endpoint + * share the same instance so subscriptions and queries stay in sync. + */ +export function createSchema() { + return makeExecutableSchema({ typeDefs, resolvers }); +} diff --git a/src/server.js b/src/server.js index c414601..65b7f5f 100644 --- a/src/server.js +++ b/src/server.js @@ -1,20 +1,18 @@ import { ApolloServer } from '@apollo/server'; import { createObservabilityPlugin } from './observability/apolloPlugin.js'; -import { resolvers } from './resolvers.js'; -import { typeDefs } from './schema.js'; +import { createSchema } from './schema.js'; /** * Builds an Apollo Server instance for the demo schema. * Exported separately from the HTTP bootstrap so tests can run queries * against the server without opening a port. */ -export function createApolloServer({ logger, metrics, enableObservability = false, plugins = [] } = {}) { +export function createApolloServer({ logger, metrics, enableObservability = false, plugins = [], schema = createSchema() } = {}) { const observabilityPlugins = enableObservability ? [createObservabilityPlugin({ logger, metrics })] : []; return new ApolloServer({ - typeDefs, - resolvers, + schema, plugins: [...observabilityPlugins, ...plugins], }); } diff --git a/src/streaming/pubsub.js b/src/streaming/pubsub.js new file mode 100644 index 0000000..625ae69 --- /dev/null +++ b/src/streaming/pubsub.js @@ -0,0 +1,104 @@ +/** + * Minimal in-process publish/subscribe used by GraphQL subscriptions. + * + * Keeping it in memory matches the rest of the demo service (no broker + * required). Swap this module for a Redis/NATS backed implementation to run + * more than one instance: subscribers only depend on the async-iterator + * contract returned by `subscribe`. + */ + +const DEFAULT_MAX_QUEUE = 100; + +export function createPubSub({ maxQueueSize = DEFAULT_MAX_QUEUE } = {}) { + /** @type {Map void>>} */ + const topics = new Map(); + + function publish(topic, payload) { + const listeners = topics.get(topic); + if (!listeners) return; + for (const listener of [...listeners]) listener(payload); + } + + /** + * Returns an async iterator yielding every payload published to `topic` + * after the call. The iterator stops when `signal` aborts or when the + * consumer calls `return()` (which `for await ... of` does on break). + */ + function subscribe(topic, { signal } = {}) { + const queue = []; + /** @type {((result: IteratorResult) => void) | null} */ + let pending = null; + let done = false; + + const listener = (payload) => { + if (done) return; + if (pending) { + const resolve = pending; + pending = null; + resolve({ value: payload, done: false }); + return; + } + // Drop the oldest event rather than growing without bound when a + // consumer is slower than the publisher. + if (queue.length >= maxQueueSize) queue.shift(); + queue.push(payload); + }; + + const listeners = topics.get(topic) ?? new Set(); + listeners.add(listener); + topics.set(topic, listeners); + + function stop() { + if (done) return { value: undefined, done: true }; + done = true; + // Drop buffered events so abort/return() terminate immediately instead + // of draining payloads queued before the stop. + queue.length = 0; + listeners.delete(listener); + if (listeners.size === 0) topics.delete(topic); + signal?.removeEventListener?.('abort', stop); + if (pending) { + const resolve = pending; + pending = null; + resolve({ value: undefined, done: true }); + } + return { value: undefined, done: true }; + } + + if (signal) { + if (signal.aborted) stop(); + else signal.addEventListener('abort', stop, { once: true }); + } + + return { + [Symbol.asyncIterator]() { + return this; + }, + next() { + if (queue.length > 0) return Promise.resolve({ value: queue.shift(), done: false }); + if (done) return Promise.resolve({ value: undefined, done: true }); + return new Promise((resolve) => { + pending = resolve; + }); + }, + return() { + return Promise.resolve(stop()); + }, + throw(error) { + stop(); + return Promise.reject(error); + }, + }; + } + + return { publish, subscribe }; +} + +/** Topics published by the demo resolvers. */ +export const TOPICS = { + USER_CREATED: 'USER_CREATED', + POST_CREATED: 'POST_CREATED', +}; + +/** Default pub/sub instance shared by the running server. */ +export const pubsub = createPubSub(); diff --git a/src/streaming/sseRouter.js b/src/streaming/sseRouter.js new file mode 100644 index 0000000..407e04f --- /dev/null +++ b/src/streaming/sseRouter.js @@ -0,0 +1,193 @@ +import express from 'express'; +import { GraphQLError, execute, getOperationAST, parse, subscribe, validate } from 'graphql'; + +import { createSchema } from '../schema.js'; + +const DEFAULT_HEARTBEAT_MS = 15000; +const SSE_HEADERS = { + 'content-type': 'text/event-stream', + 'cache-control': 'no-cache, no-transform', + connection: 'keep-alive', + // Disable proxy buffering so events reach the client immediately. + 'x-accel-buffering': 'no', +}; + +function readOperation(req) { + const source = req.method === 'GET' ? req.query : req.body ?? {}; + const { query, operationName } = source; + let { variables } = source; + + if (typeof variables === 'string') { + try { + variables = JSON.parse(variables); + } catch { + throw new GraphQLError('Invalid "variables" parameter: expected JSON.', { + extensions: { code: 'BAD_USER_INPUT' }, + }); + } + } + + if (typeof query !== 'string' || query.trim() === '') { + throw new GraphQLError('A GraphQL "query" is required.', { extensions: { code: 'BAD_USER_INPUT' } }); + } + + return { query, variables: variables ?? undefined, operationName: operationName ?? undefined }; +} + +function writeEvent(res, event, data) { + let flushed = res.write(`event: ${event}\n`); + if (data !== undefined) flushed = res.write(`data: ${JSON.stringify(data)}\n`) && flushed; + return res.write('\n') && flushed; +} + +/** + * Resolves once the socket has flushed its buffer, or immediately when the + * request is already aborted, so a slow client cannot make the server buffer + * payloads without bound. + */ +function waitForDrain(res, signal) { + return new Promise((resolve) => { + if (signal.aborted || res.writableEnded) { + resolve(); + return; + } + const finish = () => { + res.off('drain', finish); + res.off('close', finish); + res.off('error', finish); + signal.removeEventListener('abort', finish); + resolve(); + }; + res.once('drain', finish); + res.once('close', finish); + res.once('error', finish); + signal.addEventListener('abort', finish, { once: true }); + }); +} + +/** + * Express router streaming GraphQL results over Server-Sent Events. + * + * `GET|POST /graphql/stream` accepts the usual `query`/`variables`/ + * `operationName` payload. Subscriptions stream one `next` event per emitted + * result and finish with `complete`; queries and mutations are executed once + * and streamed as a single `next` + `complete` pair, so a single client + * transport covers every operation type. + */ +export function createStreamRouter({ + schema = createSchema(), + contextValue, + heartbeatMs = DEFAULT_HEARTBEAT_MS, + logger, + metrics, +} = {}) { + const router = express.Router(); + + async function handle(req, res) { + let operation; + try { + operation = readOperation(req); + } catch (error) { + res.status(400).json({ errors: [{ message: error.message, extensions: error.extensions }] }); + return; + } + + let document; + try { + document = parse(operation.query); + } catch (error) { + res.status(400).json({ errors: [{ message: error.message }] }); + return; + } + + const validationErrors = validate(schema, document); + if (validationErrors.length > 0) { + res.status(400).json({ errors: validationErrors.map((error) => ({ message: error.message })) }); + return; + } + + const operationType = getOperationAST(document, operation.operationName)?.operation; + + // GET requests may be prefetched or replayed, so they must stay side-effect free. + if (req.method === 'GET' && operationType === 'mutation') { + res.set('Allow', 'POST').status(405).json({ + errors: [ + { message: 'Mutations must be sent with POST.', extensions: { code: 'METHOD_NOT_ALLOWED' } }, + ], + }); + return; + } + + // Aborts in-flight subscriptions as soon as the client goes away. + const controller = new AbortController(); + res.on('close', () => controller.abort()); + + const args = { + schema, + document, + variableValues: operation.variables, + operationName: operation.operationName, + }; + + const isSubscription = operationType === 'subscription'; + + let result; + try { + // Context construction is awaited here so a throwing/rejecting context + // factory answers the request instead of becoming an unhandled rejection. + args.contextValue = { + ...(typeof contextValue === 'function' ? await contextValue({ req }) : contextValue ?? {}), + signal: controller.signal, + }; + result = await (isSubscription ? subscribe(args) : execute(args)); + } catch (error) { + res.status(500).json({ errors: [{ message: error?.message ?? String(error) }] }); + return; + } + + // Queries, mutations, and failed subscription setups produce a single + // ExecutionResult instead of an async iterator. + if (!result || typeof result[Symbol.asyncIterator] !== 'function') { + res.writeHead(200, SSE_HEADERS); + writeEvent(res, 'next', result); + writeEvent(res, 'complete'); + res.end(); + return; + } + + metrics?.increment?.('graphql_stream_subscriptions_total', { transport: 'sse' }); + res.writeHead(200, SSE_HEADERS); + // Comment frame flushes headers and keeps proxies from idling the socket. + res.write(': connected\n\n'); + + const heartbeat = heartbeatMs > 0 ? setInterval(() => res.write(': ping\n\n'), heartbeatMs) : null; + heartbeat?.unref?.(); + + try { + for await (const payload of result) { + if (controller.signal.aborted) break; + // Stop pulling from the iterator until the socket drains so a slow + // client cannot grow the response buffer at the publisher's rate. + if (writeEvent(res, 'next', payload) === false) { + await waitForDrain(res, controller.signal); + } + } + if (!controller.signal.aborted) writeEvent(res, 'complete'); + } catch (error) { + logger?.error?.('graphql stream failed', { error: error?.message ?? String(error) }); + if (!controller.signal.aborted) { + writeEvent(res, 'next', { errors: [{ message: 'Subscription failed.' }] }); + writeEvent(res, 'complete'); + } + } finally { + if (heartbeat) clearInterval(heartbeat); + await result.return?.(); + res.end(); + } + } + + router.get('/', handle); + router.post('/', handle); + + return router; +} diff --git a/test/streaming.test.js b/test/streaming.test.js new file mode 100644 index 0000000..2a2b93d --- /dev/null +++ b/test/streaming.test.js @@ -0,0 +1,219 @@ +import assert from 'node:assert/strict'; +import { createServer } from 'node:http'; +import { describe, it } from 'node:test'; + +import cors from 'cors'; +import express from 'express'; + +import { createStore } from '../src/data/store.js'; +import { createSchema } from '../src/schema.js'; +import { createFinanceService } from '../src/services/financeService.js'; +import { createPubSub, TOPICS } from '../src/streaming/pubsub.js'; +import { createStreamRouter } from '../src/streaming/sseRouter.js'; + +/** Starts the streaming endpoint on an ephemeral port and returns its base URL. */ +async function startServer() { + const store = createStore(); + const pubsub = createPubSub(); + const app = express(); + + app.use( + '/graphql/stream', + cors(), + express.json(), + createStreamRouter({ + schema: createSchema(), + heartbeatMs: 0, + contextValue: () => ({ store, finance: createFinanceService(), pubsub }), + }) + ); + + const server = createServer(app); + await new Promise((resolve) => server.listen(0, resolve)); + const { port } = server.address(); + + return { + url: `http://127.0.0.1:${port}/graphql/stream`, + store, + pubsub, + close: () => new Promise((resolve) => server.close(resolve)), + }; +} + +/** Parses an SSE byte stream into `{ event, data }` frames. */ +async function* readEvents(response, controller) { + const decoder = new TextDecoder(); + let buffer = ''; + + for await (const chunk of response.body) { + buffer += decoder.decode(chunk, { stream: true }); + + let index; + while ((index = buffer.indexOf('\n\n')) !== -1) { + const frame = buffer.slice(0, index); + buffer = buffer.slice(index + 2); + if (frame.startsWith(':')) continue; // heartbeat/comment + + const event = /^event: (.*)$/m.exec(frame)?.[1]; + const data = /^data: (.*)$/m.exec(frame)?.[1]; + yield { event, data: data ? JSON.parse(data) : undefined }; + if (event === 'complete') { + controller?.abort(); + return; + } + } + } +} + +describe('GraphQL streaming (SSE)', () => { + it('streams subscription events as they are published', async () => { + const app = await startServer(); + + try { + const controller = new AbortController(); + const response = await fetch(`${app.url}?query=${encodeURIComponent('subscription { userCreated { id name } }')}`, { + signal: controller.signal, + }); + + assert.equal(response.status, 200); + assert.match(response.headers.get('content-type'), /text\/event-stream/); + + const events = readEvents(response, controller); + // Give the subscription a tick to register before publishing. + await new Promise((resolve) => setTimeout(resolve, 20)); + + const user = app.store.createUser({ name: 'Grace Hopper', email: 'grace@example.com' }); + app.pubsub.publish(TOPICS.USER_CREATED, { userCreated: user }); + + const first = await events.next(); + assert.equal(first.value.event, 'next'); + assert.equal(first.value.data.data.userCreated.name, 'Grace Hopper'); + + controller.abort(); + } finally { + await app.close(); + } + }); + + it('filters post events by author', async () => { + const app = await startServer(); + + try { + const controller = new AbortController(); + const query = 'subscription { postCreated(authorId: "2") { id title author { name } } }'; + const response = await fetch(`${app.url}?query=${encodeURIComponent(query)}`, { signal: controller.signal }); + const events = readEvents(response, controller); + await new Promise((resolve) => setTimeout(resolve, 20)); + + app.pubsub.publish(TOPICS.POST_CREATED, { postCreated: app.store.createPost({ title: 'Ignored', content: 'x', authorId: '1' }) }); + app.pubsub.publish(TOPICS.POST_CREATED, { postCreated: app.store.createPost({ title: 'Kept', content: 'y', authorId: '2' }) }); + + const first = await events.next(); + assert.equal(first.value.data.data.postCreated.title, 'Kept'); + assert.equal(first.value.data.data.postCreated.author.name, 'Alan Turing'); + + controller.abort(); + } finally { + await app.close(); + } + }); + + it('streams a single result for queries and mutations', async () => { + const app = await startServer(); + + try { + const response = await fetch(app.url, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ query: '{ users { name } }' }), + }); + + assert.equal(response.status, 200); + const frames = []; + for await (const frame of readEvents(response)) frames.push(frame); + + assert.deepEqual(frames.map((frame) => frame.event), ['next', 'complete']); + assert.equal(frames[0].data.data.users.length, 2); + + const mutation = await fetch(app.url, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ + query: 'mutation { createUser(name: "Grace Hopper", email: "grace@example.com") { id name } }', + }), + }); + + assert.equal(mutation.status, 200); + const mutationFrames = []; + for await (const frame of readEvents(mutation)) mutationFrames.push(frame); + + assert.deepEqual(mutationFrames.map((frame) => frame.event), ['next', 'complete']); + assert.equal(mutationFrames[0].data.data.createUser.name, 'Grace Hopper'); + assert.equal(app.store.getUser(mutationFrames[0].data.data.createUser.id)?.name, 'Grace Hopper'); + } finally { + await app.close(); + } + }); + + it('rejects malformed and invalid operations with 400', async () => { + const app = await startServer(); + + try { + const missing = await fetch(app.url, { method: 'POST', headers: { 'content-type': 'application/json' }, body: '{}' }); + assert.equal(missing.status, 400); + + const syntax = await fetch(`${app.url}?query=${encodeURIComponent('{ users')}`); + assert.equal(syntax.status, 400); + + const invalid = await fetch(`${app.url}?query=${encodeURIComponent('subscription { nope }')}`); + assert.equal(invalid.status, 400); + const body = await invalid.json(); + assert.ok(body.errors[0].message.includes('nope')); + } finally { + await app.close(); + } + }); +}); + +describe('pub/sub', () => { + it('delivers buffered payloads and stops on abort', async () => { + const pubsub = createPubSub(); + const controller = new AbortController(); + const iterator = pubsub.subscribe('topic', { signal: controller.signal }); + + pubsub.publish('topic', 1); + pubsub.publish('topic', 2); + + assert.deepEqual(await iterator.next(), { value: 1, done: false }); + assert.deepEqual(await iterator.next(), { value: 2, done: false }); + + const pending = iterator.next(); + controller.abort(); + assert.deepEqual(await pending, { value: undefined, done: true }); + assert.deepEqual(await iterator.next(), { value: undefined, done: true }); + }); + + it('drops the oldest events once the queue is full', async () => { + const pubsub = createPubSub({ maxQueueSize: 2 }); + const iterator = pubsub.subscribe('topic'); + + pubsub.publish('topic', 'a'); + pubsub.publish('topic', 'b'); + pubsub.publish('topic', 'c'); + + assert.equal((await iterator.next()).value, 'b'); + assert.equal((await iterator.next()).value, 'c'); + await iterator.return(); + }); + + it('stops delivering after unsubscribe', async () => { + const pubsub = createPubSub(); + const iterator = pubsub.subscribe('topic'); + + pubsub.publish('topic', 'buffered'); + await iterator.return(); + pubsub.publish('topic', 'ignored'); + + assert.deepEqual(await iterator.next(), { value: undefined, done: true }); + }); +});