Skip to content
File

Blob: src/cloudflare/internal/test/pipeline-transform/transform-test.js

javascript126 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// eslint-disable-next-line @typescript-eslint/ban-ts-comment
5// @ts-nocheck
6 
7import assert from 'node:assert';
8import { PipelineTransformationEntrypoint } from 'cloudflare:pipelines';
9 
10// this is how "Pipeline" would be implemented by the user
11const customTransform = class MyEntrypoint extends PipelineTransformationEntrypoint {
12 /**
13 * @param {any} batch
14 * @override
15 */
16 async run(records, _) {
17 for (const record of records) {
18 record.dispatcher = 'was here!';
19 await new Promise((resolve) => setTimeout(resolve, 1));
20 record.wait = 'happened!';
21 }
22 
23 return records;
24 }
25};
26 
27const lines = [];
28const jimmy = `${JSON.stringify({ name: 'jimmy', age: '42', bio: 'hello\nworld\n\n' })}\n`;
29const jonny = `${JSON.stringify({ name: 'jonny', age: '9' })}\n`;
30const joey = `${JSON.stringify({ name: 'joey', age: '108' })}\n\n`;
31 
32for (let i = 0; i < 1000; i++) {
33 lines.push(jimmy, jonny, joey);
34}
35 
36function newBatch() {
37 return {
38 id: 'test',
39 shard: '0',
40 ts: Date.now(),
41 format: 'json_stream',
42 data: new ReadableStream({
43 start(controller) {
44 const encoder = new TextEncoder();
45 for (const line of lines) {
46 controller.enqueue(encoder.encode(line));
47 }
48 controller.close();
49 },
50 }),
51 };
52}
53 
54// bazel test //src/cloudflare/internal/test/pipeline-transform:transform --test_output=errors --sandbox_debug
55export const tests = {
56 async test(ctr, env, ctx) {
57 {
58 // should fail dispatcher test call when PipelineTransform class not extended
59 const transformer = new PipelineTransformationEntrypoint(ctx, env);
60 await assert.rejects(transformer._ping(), (err) => {
61 assert.strictEqual(
62 err.message,
63 'the run method must be overridden by the PipelineTransformationEntrypoint subclass'
64 );
65 return true;
66 });
67 }
68 
69 {
70 // should correctly handle dispatcher test call
71 const transform = new customTransform(ctx, env);
72 await assert.doesNotReject(transform._ping());
73 }
74 
75 {
76 // should return mutated batch
77 const transformer = new customTransform(ctx, env);
78 const batch = newBatch();
79 
80 const result = await transformer._run(batch, {
81 id: 'abc',
82 name: 'mypipeline',
83 });
84 assert.equal(true, result.data instanceof ReadableStream);
85 
86 const reader = result.data
87 .pipeThrough(new TextDecoderStream())
88 .getReader();
89 
90 let data = '';
91 while (true) {
92 const { done, value } = await reader.read();
93 if (done) {
94 break;
95 } else {
96 data += value;
97 }
98 }
99 reader.releaseLock();
100 
101 assert.notEqual(data.length, 0);
102 
103 const objects = [];
104 const resultLines = data.split('\n');
105 for (const line of resultLines) {
106 if (line.trim().length > 0) {
107 // Guard against empty lines
108 objects.push(JSON.parse(line));
109 }
110 }
111 assert.equal(objects.length, 3000);
112 
113 let index = 0;
114 for (const obj of objects) {
115 assert.equal(obj.dispatcher, 'was here!');
116 delete obj.dispatcher;
117 assert.equal(obj.wait, 'happened!');
118 delete obj.wait;
119 
120 assert.equal(`${JSON.stringify(obj)}`, lines[index].trim());
121 index++;
122 }
123 }
124 },
125};