/
MemoryEventProcessor.ts
102 lines (91 loc) · 3.7 KB
/
MemoryEventProcessor.ts
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
import { ExecutionResult, GraphQLSchema } from 'graphql';
import { getAsyncIterator, isAsyncIterable } from 'iterall';
import { ArrayPubSub } from './ArrayPubSub';
import { formatMessage } from './formatMessage';
import { execute, ExecuteOptions } from './execute';
import {
IConnectionManager,
ISubscriptionEvent,
ISubscriptionManager,
IEventProcessor,
} from './types';
import { SERVER_EVENT_TYPES } from './protocol';
import { Server } from './Server';
// polyfill Symbol.asyncIterator
if (Symbol.asyncIterator === undefined) {
(Symbol as any).asyncIterator = Symbol.for('asyncIterator');
}
interface MemoryEventProcessorOptions {
connectionManager: IConnectionManager;
context?: ExecuteOptions['context'];
schema: GraphQLSchema;
subscriptionManager: ISubscriptionManager;
}
export type EventProcessorFn = (
events: ISubscriptionEvent[],
lambdaContext?: any,
) => Promise<void>;
export class MemoryEventProcessor<TServer extends Server = Server>
implements IEventProcessor<TServer, EventProcessorFn> {
public createHandler(server: TServer): EventProcessorFn {
return async function processEvents(events, lambdaContext = {}) {
const options = await server.createGraphQLServerOptions(
events as any,
lambdaContext,
);
const { connectionManager, subscriptionManager } = options.$$internal;
for (const event of events) {
// iterate over subscribers that listen to this event
// and for each connection:
// - create a schema (so we have subscribers registered in PubSub)
// - execute operation from event againt schema
// - if iterator returns a result, send it to client
// - clean up subscriptions and follow with next page of subscriptions
// - if the are no more subscriptions, process next event
// make sure that you won't throw any errors otherwise dynamo will call
// handler with same events again
for await (const subscribers of subscriptionManager.subscribersByEventName(
event.event,
)) {
const promises = subscribers
.map(async subscriber => {
// create PubSub for this subscriber
const pubSub = new ArrayPubSub([event]);
// execute operation by executing it and then publishing the event
const iterable = await execute({
connectionManager,
subscriptionManager,
schema: options.schema,
event: {} as any, // we don't have api gateway event here
lambdaContext: lambdaContext as any, // we don't have a lambda's context here
context: options.context,
connection: subscriber.connection,
operation: subscriber.operation,
pubSub,
registerSubscriptions: false,
});
if (!isAsyncIterable(iterable)) {
// something went wrong, probably there is an error
return Promise.resolve();
}
const iterator = getAsyncIterator(iterable);
const result: IteratorResult<ExecutionResult> = await iterator.next();
if (result.value != null) {
return connectionManager.sendToConnection(
subscriber.connection,
formatMessage({
id: subscriber.operationId,
payload: result.value,
type: SERVER_EVENT_TYPES.GQL_DATA,
}),
);
}
return Promise.resolve();
})
.map(promise => promise.catch(e => console.log(e)));
await Promise.all(promises);
}
}
};
}
}