Skip to content
File

Blob: src/workerd/api/queue.h

cpp483 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 
5#pragma once
6 
7#include <workerd/api/basics.h>
8#include <workerd/io/trace.h>
9#include <workerd/io/worker-interface.capnp.h>
10#include <workerd/io/worker-interface.h>
11#include <workerd/io/worker.h>
12#include <workerd/jsg/jsg.h>
13 
14#include <kj/async.h>
15#include <kj/common.h>
16 
17namespace workerd::api {
18 
19class ExecutionContext;
20 
21// Binding types
22 
23// A capability to a Worker Queue.
24class WorkerQueue: public jsg::Object {
25 public:
26 // `subrequestChannel` is what to pass to IoContext::getHttpClient() to get an HttpClient
27 // representing this queue.
28 WorkerQueue(uint subrequestChannel): subrequestChannel(subrequestChannel) {}
29 
30 // The metrics structs below (Metrics, SendMetrics, SendBatchMetrics) are deserialized from
31 // JSON responses where the upstream service uses 0 as a sentinel for "no data" on timestamp
32 // fields. Callers MUST call clearEpochSentinel() on oldestMessageTimestamp after deserialization to convert the
33 // sentinel to kj::none (JS undefined).
34 struct Metrics {
35 double backlogCount = 0;
36 double backlogBytes = 0;
37 jsg::Optional<kj::Date> oldestMessageTimestamp;
38 JSG_STRUCT(backlogCount, backlogBytes, oldestMessageTimestamp);
39 JSG_STRUCT_TS_OVERRIDE(QueueMetrics);
40 };
41 
42 struct SendMetrics {
43 double backlogCount = 0;
44 double backlogBytes = 0;
45 jsg::Optional<kj::Date> oldestMessageTimestamp;
46 JSG_STRUCT(backlogCount, backlogBytes, oldestMessageTimestamp);
47 JSG_STRUCT_TS_OVERRIDE(QueueSendMetrics);
48 };
49 
50 struct SendMetadata {
51 SendMetrics metrics;
52 JSG_STRUCT(metrics);
53 JSG_STRUCT_TS_OVERRIDE(QueueSendMetadata);
54 };
55 
56 struct SendResponse {
57 SendMetadata metadata;
58 JSG_STRUCT(metadata);
59 JSG_STRUCT_TS_OVERRIDE(QueueSendResponse);
60 };
61 
62 struct SendBatchMetrics {
63 double backlogCount = 0;
64 double backlogBytes = 0;
65 jsg::Optional<kj::Date> oldestMessageTimestamp;
66 JSG_STRUCT(backlogCount, backlogBytes, oldestMessageTimestamp);
67 JSG_STRUCT_TS_OVERRIDE(QueueSendBatchMetrics);
68 };
69 
70 struct SendBatchMetadata {
71 SendBatchMetrics metrics;
72 JSG_STRUCT(metrics);
73 JSG_STRUCT_TS_OVERRIDE(QueueSendBatchMetadata);
74 };
75 
76 struct SendBatchResponse {
77 SendBatchMetadata metadata;
78 JSG_STRUCT(metadata);
79 JSG_STRUCT_TS_OVERRIDE(QueueSendBatchResponse);
80 };
81 
82 struct SendOptions {
83 // TODO(soon): Support metadata.
84 
85 // contentType determines the serialization format of the message.
86 jsg::Optional<kj::String> contentType;
87 
88 // The number of seconds to delay the delivery of the message being sent.
89 jsg::Optional<int> delaySeconds;
90 
91 JSG_STRUCT(contentType, delaySeconds);
92 JSG_STRUCT_TS_OVERRIDE(QueueSendOptions { contentType?: QueueContentType; });
93 // NOTE: Any new fields added here should also be added to MessageSendRequest below.
94 };
95 
96 struct SendBatchOptions {
97 // The number of seconds to delay the delivery of the message being sent.
98 jsg::Optional<int> delaySeconds;
99 
100 JSG_STRUCT(delaySeconds);
101 JSG_STRUCT_TS_OVERRIDE(QueueSendBatchOptions { delaySeconds ?: number; });
102 // NOTE: Any new fields added here should also be added to MessageSendRequest below.
103 };
104 
105 struct MessageSendRequest {
106 jsg::JsRef<jsg::JsValue> body;
107 
108 // contentType determines the serialization format of the message.
109 jsg::Optional<kj::String> contentType;
110 
111 // The number of seconds to delay the delivery of the message being sent.
112 jsg::Optional<int> delaySeconds;
113 
114 JSG_STRUCT(body, contentType, delaySeconds);
115 JSG_STRUCT_TS_OVERRIDE(MessageSendRequest<Body = unknown> {
116 body: Body;
117 contentType?: QueueContentType;
118 });
119 // NOTE: Any new fields added to SendOptions must also be added here.
120 };
121 
122 jsg::Promise<SendResponse> send(jsg::Lock& js,
123 jsg::JsValue body,
124 jsg::Optional<SendOptions> options,
125 const jsg::TypeHandler<SendResponse>& responseHandler);
126 
127 jsg::Promise<SendBatchResponse> sendBatch(jsg::Lock& js,
128 jsg::Sequence<MessageSendRequest> batch,
129 jsg::Optional<SendBatchOptions> options,
130 const jsg::TypeHandler<SendBatchResponse>& responseHandler);
131 
132 jsg::Promise<Metrics> metrics(jsg::Lock& js, const jsg::TypeHandler<Metrics>& metricsHandler);
133 
134 JSG_RESOURCE_TYPE(WorkerQueue, CompatibilityFlags::Reader flags) {
135 JSG_METHOD(metrics);
136 JSG_METHOD(send);
137 JSG_METHOD(sendBatch);
138 
139 JSG_TS_ROOT();
140 JSG_TS_OVERRIDE(Queue<Body = unknown> {
141 send(message: Body, options?: QueueSendOptions): Promise<QueueSendResponse>;
142 sendBatch(messages
143 : Iterable<MessageSendRequest<Body>>, options ?: QueueSendBatchOptions)
144 : Promise<QueueSendBatchResponse>;
145 metrics(): Promise<QueueMetrics>;
146 });
147 JSG_TS_DEFINE(type QueueContentType = "text" | "bytes" | "json" | "v8");
148 }
149 
150 private:
151 uint subrequestChannel;
152};
153 
154// Event handler types
155 
156// Metadata delivered with a message batch in the queue() handler
157 
158// Same sentinel caveat as WorkerQueue::Metrics above: the capnp path uses 0 to mean "no data"
159// for oldestMessageTimestamp. As such, we must explicitly set it to kj::none (JS undefined).
160struct MessageBatchMetrics {
161 double backlogCount = 0;
162 double backlogBytes = 0;
163 jsg::Optional<kj::Date> oldestMessageTimestamp;
164 JSG_STRUCT(backlogCount, backlogBytes, oldestMessageTimestamp);
165 JSG_STRUCT_TS_OVERRIDE(MessageBatchMetrics);
166};
167 
168struct MessageBatchMetadata {
169 MessageBatchMetrics metrics;
170 JSG_STRUCT(metrics);
171 JSG_STRUCT_TS_OVERRIDE(MessageBatchMetadata);
172};
173 
174// Types for other workers passing messages into and responses out of a queue handler.
175 
176struct IncomingQueueMessage {
177 kj::String id;
178 kj::Date timestamp;
179 kj::Array<kj::byte> body;
180 kj::Maybe<kj::String> contentType;
181 uint16_t attempts;
182 JSG_STRUCT(id, timestamp, body, contentType, attempts);
183 
184 struct ContentType {
185 static constexpr kj::StringPtr TEXT = "text"_kj;
186 static constexpr kj::StringPtr BYTES = "bytes"_kj;
187 static constexpr kj::StringPtr JSON = "json"_kj;
188 static constexpr kj::StringPtr V8 = "v8"_kj;
189 };
190};
191 
192struct QueueRetryBatch {
193 bool retry;
194 jsg::Optional<int> delaySeconds;
195 JSG_STRUCT(retry, delaySeconds);
196};
197 
198struct QueueRetryMessage {
199 kj::String msgId;
200 jsg::Optional<int> delaySeconds;
201 JSG_STRUCT(msgId, delaySeconds);
202};
203 
204struct QueueResponse {
205 uint16_t outcome;
206 bool ackAll;
207 QueueRetryBatch retryBatch;
208 kj::Array<kj::String> explicitAcks;
209 kj::Array<QueueRetryMessage> retryMessages;
210 JSG_STRUCT(outcome, ackAll, retryBatch, explicitAcks, retryMessages);
211};
212 
213// Internal-only representation used to accumulate the results of a queue event.
214 
215struct QueueEventResult {
216 struct RetryOptions {
217 jsg::Optional<int> delaySeconds;
218 };
219 struct RetryBatch {
220 bool retry;
221 jsg::Optional<int> delaySeconds;
222 };
223 RetryBatch retryBatch = {.retry = false};
224 bool ackAll = false;
225 kj::HashMap<kj::String, RetryOptions> retries;
226 kj::HashSet<kj::String> explicitAcks;
227};
228 
229struct QueueRetryOptions {
230 jsg::Optional<int> delaySeconds;
231 JSG_STRUCT(delaySeconds);
232};
233 
234class QueueMessage final: public jsg::Object {
235 public:
236 QueueMessage(jsg::Lock& js, rpc::QueueMessage::Reader message, IoPtr<QueueEventResult> result);
237 QueueMessage(jsg::Lock& js, IncomingQueueMessage message, IoPtr<QueueEventResult> result);
238 
239 kj::StringPtr getId() {
240 return id;
241 }
242 kj::Date getTimestamp() {
243 return timestamp;
244 }
245 jsg::JsValue getBody(jsg::Lock& js);
246 uint16_t getAttempts() {
247 return attempts;
248 };
249 
250 void retry(jsg::Optional<QueueRetryOptions> options);
251 void ack();
252 
253 // TODO(soon): Add metadata support.
254 
255 JSG_RESOURCE_TYPE(QueueMessage) {
256 JSG_READONLY_INSTANCE_PROPERTY(id, getId);
257 JSG_READONLY_INSTANCE_PROPERTY(timestamp, getTimestamp);
258 JSG_READONLY_INSTANCE_PROPERTY(body, getBody);
259 JSG_READONLY_INSTANCE_PROPERTY(attempts, getAttempts);
260 JSG_METHOD(retry);
261 JSG_METHOD(ack);
262 
263 JSG_TS_OVERRIDE(Message<Body = unknown> {
264 readonly body: Body;
265 });
266 }
267 
268 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
269 tracker.trackField("id", id);
270 tracker.trackField("body", body);
271 tracker.trackFieldWithSize("IoPtr<QueueEventResult>", sizeof(IoPtr<QueueEventResult>));
272 }
273 
274 private:
275 kj::String id;
276 kj::Date timestamp;
277 jsg::JsRef<jsg::JsValue> body;
278 uint16_t attempts;
279 IoPtr<QueueEventResult> result;
280 
281 void visitForGc(jsg::GcVisitor& visitor) {
282 visitor.visit(body);
283 }
284};
285 
286class QueueEvent final: public ExtendableEvent {
287 public:
288 // TODO(cleanup): Should we get around the need for this alternative param type by just having the
289 // service worker caller provide us with capnp-serialized params?
290 struct Params {
291 kj::String queueName;
292 kj::Array<IncomingQueueMessage> messages;
293 MessageBatchMetadata metadata;
294 };
295 
296 explicit QueueEvent(jsg::Lock& js,
297 rpc::EventDispatcher::QueueParams::Reader params,
298 IoPtr<QueueEventResult> result);
299 explicit QueueEvent(jsg::Lock& js, Params params, IoPtr<QueueEventResult> result);
300 
301 static jsg::Ref<QueueEvent> constructor(kj::String type) = delete;
302 
303 kj::ArrayPtr<jsg::Ref<QueueMessage>> getMessages() {
304 return messages;
305 }
306 kj::StringPtr getQueueName() {
307 return queueName;
308 }
309 MessageBatchMetadata getMetadata() {
310 return metadata;
311 }
312 
313 void retryAll(jsg::Optional<QueueRetryOptions> options);
314 void ackAll();
315 
316 JSG_RESOURCE_TYPE(QueueEvent, CompatibilityFlags::Reader flags) {
317 JSG_INHERIT(ExtendableEvent);
318 
319 JSG_LAZY_READONLY_INSTANCE_PROPERTY(messages, getMessages);
320 JSG_READONLY_INSTANCE_PROPERTY(queue, getQueueName);
321 
322 JSG_READONLY_INSTANCE_PROPERTY(metadata, getMetadata);
323 
324 JSG_METHOD(retryAll);
325 JSG_METHOD(ackAll);
326 
327 JSG_TS_ROOT();
328 JSG_TS_OVERRIDE(QueueEvent<Body = unknown> {
329 readonly messages: readonly Message<Body>[];
330 readonly metadata: MessageBatchMetadata;
331 });
332 }
333 
334 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
335 for (auto& message: messages) {
336 tracker.trackField("message", message);
337 }
338 tracker.trackField("queueName", queueName);
339 tracker.trackFieldWithSize("metadata", sizeof(MessageBatchMetadata));
340 tracker.trackFieldWithSize("IoPtr<QueueEventResult>", sizeof(IoPtr<QueueEventResult>));
341 }
342 
343 struct Incomplete {};
344 struct CompletedSuccessfully {};
345 struct CompletedWithError {
346 kj::Exception error;
347 };
348 using CompletionStatus = kj::OneOf<Incomplete, CompletedSuccessfully, CompletedWithError>;
349 
350 void setCompletionStatus(CompletionStatus status) {
351 completionStatus = kj::mv(status);
352 }
353 
354 const CompletionStatus& getCompletionStatus() const {
355 return completionStatus;
356 }
357 
358 private:
359 // TODO(perf): Should we store these in a v8 array directly rather than this intermediate kj
360 // array to avoid one intermediate copy?
361 kj::Array<jsg::Ref<QueueMessage>> messages;
362 kj::String queueName;
363 MessageBatchMetadata metadata;
364 IoPtr<QueueEventResult> result;
365 CompletionStatus completionStatus = Incomplete{};
366 
367 void visitForGc(jsg::GcVisitor& visitor) {
368 visitor.visitAll(messages);
369 }
370};
371 
372// Type used when calling a module-exported queue event handler.
373class QueueController final: public jsg::Object {
374 public:
375 QueueController(jsg::Ref<QueueEvent> event): event(kj::mv(event)) {}
376 
377 kj::ArrayPtr<jsg::Ref<QueueMessage>> getMessages() {
378 return event->getMessages();
379 }
380 kj::StringPtr getQueueName() {
381 return event->getQueueName();
382 }
383 MessageBatchMetadata getMetadata() {
384 return event->getMetadata();
385 }
386 void retryAll(jsg::Optional<QueueRetryOptions> options) {
387 event->retryAll(options);
388 }
389 void ackAll() {
390 event->ackAll();
391 }
392 
393 JSG_RESOURCE_TYPE(QueueController, CompatibilityFlags::Reader flags) {
394 JSG_READONLY_INSTANCE_PROPERTY(messages, getMessages);
395 JSG_READONLY_INSTANCE_PROPERTY(queue, getQueueName);
396 
397 JSG_READONLY_INSTANCE_PROPERTY(metadata, getMetadata);
398 
399 JSG_METHOD(retryAll);
400 JSG_METHOD(ackAll);
401 
402 JSG_TS_ROOT();
403 JSG_TS_OVERRIDE(MessageBatch<Body = unknown> {
404 readonly messages: readonly Message<Body>[];
405 readonly metadata: MessageBatchMetadata;
406 });
407 }
408 
409 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
410 tracker.trackField("event", event);
411 }
412 
413 private:
414 jsg::Ref<QueueEvent> event;
415 
416 void visitForGc(jsg::GcVisitor& visitor) {
417 visitor.visit(event);
418 }
419};
420 
421// Extension of ExportedHandler covering queue handlers.
422struct QueueExportedHandler {
423 using QueueHandler = kj::Promise<void>(jsg::Ref<QueueController> controller,
424 jsg::JsRef<jsg::JsValue> env,
425 jsg::Optional<jsg::Ref<ExecutionContext>> ctx);
426 jsg::LenientOptional<jsg::Function<QueueHandler>> queue;
427 
428 JSG_STRUCT(queue);
429};
430 
431class QueueCustomEvent final: public WorkerInterface::CustomEvent, public kj::Refcounted {
432 public:
433 QueueCustomEvent(kj::OneOf<QueueEvent::Params, rpc::EventDispatcher::QueueParams::Reader> params)
434 : params(kj::mv(params)) {}
435 
436 kj::Promise<Result> run(kj::Own<IoContext_IncomingRequest> incomingRequest,
437 kj::Maybe<kj::StringPtr> entrypointName,
438 kj::Maybe<Worker::VersionInfo> versionInfo,
439 Frankenvalue props,
440 kj::TaskSet& waitUntilTasks,
441 bool isDynamicDispatch) override;
442 
443 kj::Promise<Result> sendRpc(capnp::HttpOverCapnpFactory& httpOverCapnpFactory,
444 capnp::ByteStreamFactory& byteStreamFactory,
445 rpc::EventDispatcher::Client dispatcher) override;
446 
447 static const uint16_t EVENT_TYPE = 5;
448 uint16_t getType() override {
449 return EVENT_TYPE;
450 }
451 
452 tracing::EventInfo getEventInfo() const override;
453 
454 QueueRetryBatch getRetryBatch() const {
455 return {.retry = result.retryBatch.retry, .delaySeconds = result.retryBatch.delaySeconds};
456 }
457 bool getAckAll() const {
458 return result.ackAll;
459 }
460 kj::Array<QueueRetryMessage> getRetryMessages() const;
461 kj::Array<kj::String> getExplicitAcks() const;
462 
463 kj::Promise<Result> notSupported() override {
464 KJ_UNIMPLEMENTED("queue event not supported");
465 }
466 
467 private:
468 kj::OneOf<rpc::EventDispatcher::QueueParams::Reader, QueueEvent::Params> params;
469 QueueEventResult result;
470};
471 
472#define EW_QUEUE_ISOLATE_TYPES \
473 api::WorkerQueue, api::WorkerQueue::SendMetrics, api::WorkerQueue::SendMetadata, \
474 api::WorkerQueue::SendResponse, api::WorkerQueue::SendBatchMetrics, \
475 api::WorkerQueue::SendBatchMetadata, api::WorkerQueue::SendBatchResponse, \
476 api::WorkerQueue::SendOptions, api::WorkerQueue::SendBatchOptions, \
477 api::WorkerQueue::MessageSendRequest, api::WorkerQueue::Metrics, api::MessageBatchMetrics, \
478 api::MessageBatchMetadata, api::IncomingQueueMessage, api::QueueRetryBatch, \
479 api::QueueRetryMessage, api::QueueResponse, api::QueueRetryOptions, api::QueueMessage, \
480 api::QueueEvent, api::QueueController, api::QueueExportedHandler
481 
482} // namespace workerd::api