Skip to content
File

Blob: src/cloudflare/internal/pipeline-transform.ts

typescript200 lines
1// Copyright (c) 2024 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5import entrypoints from 'cloudflare-internal:workers';
6 
7/**
8 * Reads a stream line by line, yielding each line as it becomes available.
9 *
10 * This function consumes a ReadableStream of Uint8Array chunks (binary data),
11 * converts it to text using TextDecoderStream, and yields each line
12 * encountered. Lines are delimited by newline characters ('\n'). The final
13 * line is yielded even if it doesn't end with a newline.
14 *
15 * @param stream - A ReadableStream containing binary data to be decoded as text
16 * @returns An AsyncGenerator that yields each line from the stream
17 */
18async function* readLines(
19 stream: ReadableStream<Uint8Array>
20): AsyncGenerator<string> {
21 // @ts-expect-error TS2345 TODO(soon): Fix this.
22 const textStream = stream.pipeThrough(new TextDecoderStream());
23 const reader = textStream.getReader();
24 
25 let buffer = '';
26 
27 try {
28 // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition
29 while (true) {
30 const { done, value } = await reader.read();
31 
32 // Add any new content to the buffer
33 if (value) buffer += value;
34 
35 // Process complete lines
36 let lineEndIndex = buffer.indexOf('\n');
37 while (lineEndIndex >= 0) {
38 yield buffer.substring(0, lineEndIndex);
39 buffer = buffer.substring(lineEndIndex + 1);
40 lineEndIndex = buffer.indexOf('\n');
41 }
42 
43 // If we're done and have processed all complete lines,
44 // yield any remaining content and exit
45 if (done) {
46 if (buffer.length > 0) {
47 yield buffer;
48 }
49 break;
50 }
51 }
52 } finally {
53 reader.releaseLock();
54 }
55}
56 
57type Batch = {
58 id: string; // unique identifier for the batch
59 shard: string; // assigned shard
60 ts: number; // creation timestamp of the batch
61 
62 format: FormatType;
63 size: {
64 bytes: number;
65 rows: number;
66 };
67 data: unknown;
68};
69 
70const Format = {
71 JSON_STREAM: 'json_stream' as const, // jsonl
72};
73type FormatType = (typeof Format)[keyof typeof Format];
74 
75type JsonStream = Batch & {
76 format: typeof Format.JSON_STREAM;
77 data: ReadableStream<Uint8Array>;
78};
79 
80type PipelineBatchMetadata = {
81 pipelineId: string;
82 pipelineName: string;
83};
84 
85type PipelineRecord = Record<string, unknown>;
86 
87export class PipelineTransformImpl<
88 I extends PipelineRecord,
89 O extends PipelineRecord,
90>
91 extends entrypoints.WorkerEntrypoint
92{
93 #batch?: Batch;
94 #initalized: boolean = false;
95 
96 // stub overridden on the subclass
97 // eslint-disable-next-line @typescript-eslint/require-await
98 async run(_records: I[], _metadata: PipelineBatchMetadata): Promise<O[]> {
99 throw new Error('should be implemented by parent');
100 }
101 
102 // called by the dispatcher to validate that run is properly implemented by the subclass
103 // @ts-expect-error This is OK. We use this method in tests.
104 // eslint-disable-next-line no-restricted-syntax
105 private async _ping(): Promise<void> {
106 // making sure the function was overridden by an implementing subclass
107 if (this.run !== PipelineTransformImpl.prototype.run) {
108 return Promise.resolve();
109 } else {
110 return Promise.reject(
111 new Error(
112 'the run method must be overridden by the PipelineTransformationEntrypoint subclass'
113 )
114 );
115 }
116 }
117 
118 // called by the dispatcher which then calls the subclass methods
119 // the reason this is typescript private and not javascript private is that this must be
120 // able to be called by the dispatcher but should not be called by the class implementer
121 // @ts-expect-error This is OK. We use this method in tests.
122 // eslint-disable-next-line no-restricted-syntax
123 private async _run(
124 batch: Batch,
125 metadata: PipelineBatchMetadata
126 ): Promise<JsonStream> {
127 if (this.#initalized) {
128 throw new Error('pipeline entrypoint has already been initialized');
129 }
130 
131 this.#batch = batch;
132 this.#initalized = true;
133 
134 // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition
135 if (this.#batch.format === Format.JSON_STREAM) {
136 const records: I[] = await this.#readJsonStream();
137 const transformed = await this.run(records, metadata);
138 return this.#sendJson(transformed);
139 } else {
140 throw new Error(
141 'PipelineTransformationEntrypoint run supports only the JSON_STREAM batch format'
142 );
143 }
144 }
145 
146 async #readJsonStream(): Promise<I[]> {
147 if (this.#batch?.format !== Format.JSON_STREAM) {
148 throw new Error(`expected JSON_STREAM not ${this.#batch?.format}`);
149 }
150 
151 const batch = this.#batch.data as ReadableStream<Uint8Array>;
152 
153 const data: I[] = [];
154 for await (const line of readLines(batch)) {
155 if (line.trim().length > 0) {
156 // guard against empty lines
157 data.push(JSON.parse(line) as I);
158 }
159 }
160 
161 return data;
162 }
163 
164 #sendJson(records: O[]): JsonStream {
165 if (!(records instanceof Array)) {
166 throw new Error('transformations must return an array of PipelineRecord');
167 }
168 
169 let written = 0;
170 const encoder = new TextEncoder();
171 const readable = new ReadableStream<Uint8Array>({
172 start(controller): void {
173 for (const record of records) {
174 const encoded = encoder.encode(`${JSON.stringify(record)}\n`);
175 written += encoded.length;
176 controller.enqueue(encoded);
177 }
178 
179 controller.close();
180 },
181 });
182 
183 if (!this.#batch) {
184 throw new Error('Batch should have been defined. Assertion error.');
185 }
186 
187 return {
188 id: this.#batch.id,
189 shard: this.#batch.shard,
190 ts: this.#batch.ts,
191 format: Format.JSON_STREAM,
192 size: {
193 bytes: written,
194 rows: records.length,
195 },
196 data: readable,
197 };
198 }
199}