Skip to content
File

Blob: src/workerd/api/tests/queue-test.js

javascript160 lines
1// Copyright (c) 2023 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 assert from 'node:assert';
6import { Buffer } from 'node:buffer';
7 
8let serializedBody;
9 
10export default {
11 // Producer receiver (from `env.QUEUE`)
12 async fetch(request, env, ctx) {
13 const { pathname } = new URL(request.url);
14 if (pathname === '/message') {
15 assert.strictEqual(request.method, 'POST');
16 const format = request.headers.get('X-Msg-Fmt') ?? 'v8';
17 if (format === 'text') {
18 assert.strictEqual(request.headers.get('X-Msg-Delay-Secs'), '2');
19 assert.strictEqual(await request.text(), 'abc');
20 } else if (format === 'bytes') {
21 const array = new Uint16Array(await request.arrayBuffer());
22 assert.deepStrictEqual(array, new Uint16Array([1, 2, 3]));
23 } else if (format === 'json') {
24 assert.deepStrictEqual(await request.json(), { a: 1 });
25 } else if (format === 'v8') {
26 // workerd doesn't provide V8 deserialization APIs, so just look for expected strings
27 const buffer = Buffer.from(await request.arrayBuffer());
28 assert(buffer.includes('key'));
29 assert(buffer.includes('value'));
30 serializedBody = buffer;
31 } else {
32 assert.fail(`Unexpected format: ${JSON.stringify(format)}`);
33 }
34 } else if (pathname === '/batch') {
35 assert.strictEqual(request.method, 'POST');
36 assert.strictEqual(request.headers.get('X-Msg-Delay-Secs'), '2');
37 
38 const body = await request.json();
39 
40 assert.strictEqual(typeof body, 'object');
41 assert(Array.isArray(body?.messages));
42 assert.strictEqual(body.messages.length, 4);
43 
44 assert.strictEqual(body.messages[0].contentType, 'text');
45 assert.strictEqual(
46 Buffer.from(body.messages[0].body, 'base64').toString(),
47 'def'
48 );
49 
50 assert.strictEqual(body.messages[1].contentType, 'bytes');
51 assert.deepStrictEqual(
52 Buffer.from(body.messages[1].body, 'base64'),
53 Buffer.from([4, 5, 6])
54 );
55 
56 assert.strictEqual(body.messages[2].contentType, 'json');
57 assert.deepStrictEqual(
58 JSON.parse(Buffer.from(body.messages[2].body, 'base64')),
59 [7, 8, { b: 9 }]
60 );
61 
62 assert.strictEqual(body.messages[3].contentType, 'v8');
63 assert(Buffer.from(body.messages[3].body, 'base64').includes('value'));
64 assert.strictEqual(body.messages[3].delaySecs, 1);
65 } else {
66 assert.fail(`Unexpected pathname: ${JSON.stringify(pathname)}`);
67 }
68 return Response.json({
69 metadata: {
70 metrics: {
71 backlogCount: 0,
72 backlogBytes: 0,
73 oldestMessageTimestamp: 0,
74 },
75 },
76 });
77 },
78 
79 // Consumer receiver (from `env.SERVICE`)
80 async queue(batch, env, ctx) {
81 assert.strictEqual(batch.queue, 'test-queue');
82 assert.strictEqual(batch.messages.length, 5);
83 
84 assert.strictEqual(batch.messages[0].id, '#0');
85 assert.strictEqual(batch.messages[0].body, 'ghi');
86 assert.strictEqual(batch.messages[0].attempts, 1);
87 
88 assert.strictEqual(batch.messages[1].id, '#1');
89 assert.deepStrictEqual(batch.messages[1].body, new Uint8Array([7, 8, 9]));
90 assert.strictEqual(batch.messages[1].attempts, 2);
91 
92 assert.strictEqual(batch.messages[2].id, '#2');
93 assert.deepStrictEqual(batch.messages[2].body, { c: { d: 10 } });
94 assert.strictEqual(batch.messages[2].attempts, 3);
95 batch.messages[2].retry();
96 
97 assert.strictEqual(batch.messages[3].id, '#3');
98 assert.deepStrictEqual(batch.messages[3].body, batch.messages[3].timestamp);
99 assert.strictEqual(batch.messages[3].attempts, 4);
100 batch.messages[3].retry({ delaySeconds: 2 });
101 
102 assert.strictEqual(batch.messages[4].id, '#4');
103 assert.deepStrictEqual(batch.messages[4].body, new Map([['key', 'value']]));
104 assert.strictEqual(batch.messages[4].attempts, 5);
105 
106 batch.ackAll();
107 },
108 
109 async test(ctrl, env, ctx) {
110 await env.QUEUE.send('abc', { contentType: 'text', delaySeconds: 2 });
111 await env.QUEUE.send(new Uint16Array([1, 2, 3]), { contentType: 'bytes' });
112 await env.QUEUE.send({ a: 1 }, { contentType: 'json' });
113 await env.QUEUE.send(new Map([['key', 'value']]), { contentType: 'v8' });
114 
115 await env.QUEUE.sendBatch(
116 [
117 { body: 'def', contentType: 'text' },
118 { body: new Uint8Array([4, 5, 6]), contentType: 'bytes' },
119 { body: [7, 8, { b: 9 }], contentType: 'json' },
120 { body: new Set(['value']), contentType: 'v8', delaySeconds: 1 },
121 ],
122 { delaySeconds: 2 }
123 );
124 
125 const timestamp = new Date();
126 const response = await env.SERVICE.queue('test-queue', [
127 { id: '#0', timestamp, body: 'ghi', attempts: 1 },
128 { id: '#1', timestamp, body: new Uint8Array([7, 8, 9]), attempts: 2 },
129 { id: '#2', timestamp, body: { c: { d: 10 } }, attempts: 3 },
130 { id: '#3', timestamp, body: timestamp, attempts: 4 },
131 { id: '#4', timestamp, serializedBody, attempts: 5 },
132 ]);
133 assert.strictEqual(response.outcome, 'ok');
134 assert(!response.retryBatch.retry);
135 assert(response.ackAll);
136 assert.deepStrictEqual(response.retryMessages, [
137 { msgId: '#2' },
138 { msgId: '#3', delaySeconds: 2 },
139 ]);
140 assert.deepStrictEqual(response.explicitAcks, []);
141 
142 await assert.rejects(
143 env.SERVICE.queue('test-queue', [{ id: '#0', timestamp, attempts: 1 }]),
144 {
145 name: 'TypeError',
146 message: 'Expected one of body or serializedBody for each message',
147 }
148 );
149 await assert.rejects(
150 env.SERVICE.queue('test-queue', [
151 { id: '#0', timestamp, body: '', serializedBody, attempts: 1 },
152 ]),
153 {
154 name: 'TypeError',
155 message: 'Expected one of body or serializedBody for each message',
156 }
157 );
158 },
159};