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
42 changes: 40 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -63,6 +65,7 @@ environment variable):
- Readiness check (per-upstream): <http://localhost:4000/ready>
- Metrics (Prometheus text): <http://localhost:4000/metrics>
- Excel exports: <http://localhost:4000/export>
- Streaming (SSE) endpoint: <http://localhost:4000/graphql/stream>

Opening the GraphQL endpoint in a browser loads the Apollo Sandbox, where you can
explore the schema and run the operations below.
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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.
Comment thread
charles2ke marked this conversation as resolved.

```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`) —
Expand Down
25 changes: 13 additions & 12 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
23 changes: 21 additions & 2 deletions src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand All @@ -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) => {
Expand Down
34 changes: 31 additions & 3 deletions src/resolvers.js
Original file line number Diff line number Diff line change
@@ -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) {
Expand Down Expand Up @@ -46,17 +54,37 @@ 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.`, {
extensions: { code: 'BAD_USER_INPUT' },
});
}

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);
},
},
},

Expand Down
20 changes: 20 additions & 0 deletions src/schema.js
Original file line number Diff line number Diff line change
@@ -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."
Expand Down Expand Up @@ -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 });
}
8 changes: 3 additions & 5 deletions src/server.js
Original file line number Diff line number Diff line change
@@ -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],
});
}
Loading
Loading