aboutsummaryrefslogtreecommitdiff
path: root/source/library/fetcher/sse.ts
blob: c28436b94700a2a020e832830a7f609255ebde89 (plain) (blame)
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 ------------------------------------------- ***/

import { createClient, type Client, type ClientOptions } from "graphql-sse";

/*** UTILITY ------------------------------------------ ***/

import type { Fetcher, FetcherOptions, FetcherResult } from "./types.ts";

/*** EXPORT ------------------------------------------- ***/

export type SseFetcherOptions = FetcherOptions & {
  client?: Partial<Omit<ClientOptions, "url">>;
};

export function createSseFetcher(options: SseFetcherOptions): Fetcher {
  const client: Client = createClient({
    fetchFn: options.fetch,
    headers: options.headers,
    url: options.url,
    ...options.client
  });

  return (req) => {
    return {
      [Symbol.asyncIterator](): AsyncIterator<FetcherResult> {
        const queue: FetcherResult[] = [];
        const resolvers: ((step: IteratorResult<FetcherResult>) => void)[] = [];
        let done = false;
        let error: unknown = null;

        const dispose = client.subscribe<FetcherResult, Record<string, unknown>>(
          {
            extensions: req.headers as Record<string, unknown> | undefined,
            operationName: req.operationName ?? undefined,
            query: req.query,
            variables: req.variables
          },
          {
            complete() {
              done = true;
              flush();
            },
            error(err) {
              error = err;
              done = true;
              flush();
            },
            next(value) {
              queue.push(value as FetcherResult);
              flush();
            }
          }
        );

        function flush() {
          while (resolvers.length > 0 && (queue.length > 0 || done)) {
            const resolve = resolvers.shift();

            if (!resolve)
              break;

            if (queue.length > 0)
              resolve({ done: false, value: queue.shift() as FetcherResult });
            else if (error)
              resolve({ done: true, value: undefined });
            else
              resolve({ done: true, value: undefined });
          }
        }

        return {
          next(): Promise<IteratorResult<FetcherResult>> {
            if (error) {
              const err = error;
              error = null;
              return Promise.reject(err);
            }

            if (queue.length > 0)
              return Promise.resolve({ done: false, value: queue.shift() as FetcherResult });

            if (done)
              return Promise.resolve({ done: true, value: undefined });

            return new Promise<IteratorResult<FetcherResult>>((resolve) => {
              resolvers.push(resolve);
            });
          },
          return(): Promise<IteratorResult<FetcherResult>> {
            dispose();
            done = true;
            flush();
            return Promise.resolve({ done: true, value: undefined });
          }
        };
      }
    };
  };
}