Skip to content
File

Blob: src/workerd/io/trace-stream.c++

55.4 KB
1#include <workerd/api/global-scope.h>
2#include <workerd/io/io-context.h>
3#include <workerd/io/io-own.h>
4#include <workerd/io/trace-stream.h>
5#include <workerd/io/worker-interface.h>
6#include <workerd/jsg/jsg.h>
7#include <workerd/util/completion-membrane.h>
8#include <workerd/util/strings.h>
9#include <workerd/util/uuid.h>
10 
11#include <capnp/membrane.h>
12 
13#include <algorithm>
14 
15namespace workerd::tracing {
16namespace {
17 
18#define STRS(V) \
19 V(ALARM, "alarm") \
20 V(ATTRIBUTES, "attributes") \
21 V(BATCHSIZE, "batchSize") \
22 V(CANCELED, "canceled") \
23 V(CHANNEL, "channel") \
24 V(CFJSON, "cfJson") \
25 V(CLOSE, "close") \
26 V(CODE, "code") \
27 V(CONNECT, "connect") \
28 V(COUNT, "count") \
29 V(CPUTIME, "cpuTime") \
30 V(CRON, "cron") \
31 V(CUSTOM, "custom") \
32 V(DAEMONDOWN, "daemonDown") \
33 V(DEBUG, "debug") \
34 V(DIAGNOSTICCHANNEL, "diagnosticChannel") \
35 V(DIAGNOSTIC, "diagnostic") \
36 V(DIAGNOSTICSTYPE, "diagnosticsType") \
37 V(DISPATCHNAMESPACE, "dispatchNamespace") \
38 V(DROPPEDEVENTS, "droppedEvents") \
39 V(EMAIL, "email") \
40 V(ENTRYPOINT, "entrypoint") \
41 V(ERROR, "error") \
42 V(EVENT, "event") \
43 V(EXCEEDEDCPU, "exceededCpu") \
44 V(EXCEEDEDMEMORY, "exceededMemory") \
45 V(EXCEPTION, "exception") \
46 V(EXECUTIONMODEL, "executionModel") \
47 V(FETCH, "fetch") \
48 V(HEADERS, "headers") \
49 V(HIBERNATABLEWEBSOCKET, "hibernatableWebSocket") \
50 V(ID, "id") \
51 V(INFO, "info") \
52 V(INTERNALERROR, "internalError") \
53 V(INVOCATIONID, "invocationId") \
54 V(JSRPC, "jsrpc") \
55 V(KILLSWITCH, "killSwitch") \
56 V(LEVEL, "level") \
57 V(LOADSHED, "loadShed") \
58 V(LOG, "log") \
59 V(MAILFROM, "mailFrom") \
60 V(MESSAGE, "message") \
61 V(METHOD, "method") \
62 V(NAME, "name") \
63 V(OK, "ok") \
64 V(ONSET, "onset") \
65 V(OUTCOME, "outcome") \
66 V(PREVIEW, "preview") \
67 V(QUEUE, "queue") \
68 V(QUEUENAME, "queueName") \
69 V(RAWSIZE, "rawSize") \
70 V(RCPTTO, "rcptTo") \
71 V(RESPONSESTREAMDISCONNECTED, "responseStreamDisconnected") \
72 V(RETURN, "return") \
73 V(SCHEDULED, "scheduled") \
74 V(SCHEDULEDTIME, "scheduledTime") \
75 V(SCRIPTNAME, "scriptName") \
76 V(SCRIPTNOTFOUND, "scriptNotFound") \
77 V(SCRIPTTAGS, "scriptTags") \
78 V(SCRIPTVERSION, "scriptVersion") \
79 V(SEQUENCE, "sequence") \
80 V(SPANCLOSE, "spanClose") \
81 V(SPANCONTEXT, "spanContext") \
82 V(SPANID, "spanId") \
83 V(TRACEFLAGS, "traceFlags") \
84 V(SPANOPEN, "spanOpen") \
85 V(STACK, "stack") \
86 V(STATUSCODE, "statusCode") \
87 V(SLUG, "slug") \
88 V(STREAMDIAGEVENT, "streamDiagEvent") \
89 V(STREAMDIAGNOSTIC, "streamDiagnostic") \
90 V(TAG, "tag") \
91 V(TIMESTAMP, "timestamp") \
92 V(TRACEID, "traceId") \
93 V(TRACE, "trace") \
94 V(TRACES, "traces") \
95 V(TYPE, "type") \
96 V(UNKNOWN, "unknown") \
97 V(URL, "url") \
98 V(VALUE, "value") \
99 V(WALLTIME, "wallTime") \
100 V(WARN, "warn") \
101 V(WASCLEAN, "wasClean")
102 
103#define V(N, L) constexpr kj::LiteralStringConst N##_STR = L##_kjc;
104STRS(V)
105#undef STRS
106 
107// Utility that prevents creating duplicate JS strings while serializing a tail event.
108class StringCache final {
109 public:
110 StringCache() = default;
111 KJ_DISALLOW_COPY_AND_MOVE(StringCache);
112 
113 // Inserted string keys must live as long as the cache. For string constants (the common case),
114 // we use LiteralStringConst and avoid memory allocation. For temporary strings, we pass in a
115 // StringPtr and allocate a string. Having ConstString as the value type fits both cases.
116 jsg::JsValue get(jsg::Lock& js, kj::LiteralStringConst value) {
117 return cache
118 .findOrCreate(value, [&]() -> decltype(cache)::Entry {
119 return {value, jsg::JsRef<jsg::JsValue>(js, js.strIntern(value))};
120 }).getHandle(js);
121 }
122 jsg::JsValue get(jsg::Lock& js, kj::StringPtr value) {
123 return cache
124 .findOrCreate(value, [&]() -> decltype(cache)::Entry {
125 return {kj::ConstString(kj::str(value)), jsg::JsRef<jsg::JsValue>(js, js.strIntern(value))};
126 }).getHandle(js);
127 }
128 
129 private:
130 kj::HashMap<kj::ConstString, jsg::JsRef<jsg::JsValue>> cache;
131};
132 
133// Why ToJS(...) functions and not JSG_STRUCT? Good question. The various tracing:*
134// types are defined in the "trace" bazel target which currently does not depend on
135// jsg in any way. These also represent the internal API of these types which doesn't
136// really match exactly what we want to expose to users. In order to use JSG_STRUCT
137// we would either need to make the "trace" target depend on "jsg", which seems a bit
138// wasteful and unnecessary, or we'd need to define wrapper structs that use JSG_STRUCT
139// which also seems wasteful and unnecessary. We also don't need the type mapping for
140// these structs to be bidirectional. So, instead, let's just do the simple easy thing
141// and define a set of serializers to these types.
142 
143// Serialize attribute value
144jsg::JsValue ToJs(jsg::Lock& js, const Attribute::Value& value) {
145 KJ_SWITCH_ONEOF(value) {
146 KJ_CASE_ONEOF(str, kj::ConstString) {
147 return js.str(str);
148 }
149 KJ_CASE_ONEOF(b, bool) {
150 return js.boolean(b);
151 }
152 KJ_CASE_ONEOF(d, double) {
153 return js.num(d);
154 }
155 KJ_CASE_ONEOF(i, int64_t) {
156 return js.bigInt(i);
157 }
158 }
159 KJ_UNREACHABLE;
160}
161 
162// Serialize attribute key:value(s) pair object
163jsg::JsValue ToJs(jsg::Lock& js, const Attribute& attribute, StringCache& cache) {
164 auto obj = js.obj();
165 obj.set(js, NAME_STR, cache.get(js, attribute.name));
166 
167 if (attribute.value.size() == 1) {
168 obj.set(js, VALUE_STR, ToJs(js, attribute.value[0]));
169 } else {
170 obj.set(js, VALUE_STR, js.arr(attribute.value.asPtr(), [](jsg::Lock& js, const auto& val) {
171 return ToJs(js, val);
172 }));
173 }
174 
175 return obj;
176}
177 
178// Serialize "attributes" event
179jsg::JsValue ToJs(jsg::Lock& js, kj::ArrayPtr<const Attribute> attributes, StringCache& cache) {
180 auto obj = js.obj();
181 obj.set(js, TYPE_STR, cache.get(js, ATTRIBUTES_STR));
182 obj.set(js, INFO_STR, js.arr(attributes, [&cache](jsg::Lock& js, const auto& attr) {
183 return ToJs(js, attr, cache);
184 }));
185 return obj;
186}
187 
188jsg::JsValue ToJs(jsg::Lock& js, const FetchResponseInfo& info, StringCache& cache) {
189 auto obj = js.obj();
190 obj.set(js, TYPE_STR, cache.get(js, FETCH_STR));
191 obj.set(js, STATUSCODE_STR, js.num(info.statusCode));
192 return obj;
193}
194 
195jsg::JsValue ToJs(jsg::Lock& js, const FetchEventInfo& info, StringCache& cache) {
196 auto obj = js.obj();
197 obj.set(js, TYPE_STR, cache.get(js, FETCH_STR));
198 obj.set(js, METHOD_STR, cache.get(js, kj::str(info.method)));
199 obj.set(js, URL_STR, js.str(info.url));
200 if (info.cfJson.size() > 0) {
201 obj.set(js, CFJSON_STR, jsg::JsValue(js.parseJson(info.cfJson).getHandle(js)));
202 }
203 
204 auto ToJs = [](jsg::Lock& js, const FetchEventInfo::Header& header, StringCache& cache) {
205 auto obj = js.obj();
206 obj.set(js, NAME_STR, cache.get(js, header.name));
207 obj.set(js, VALUE_STR, js.str(header.value));
208 return obj;
209 };
210 
211 obj.set(js, HEADERS_STR,
212 js.arr(info.headers.asPtr(),
213 [&cache, &ToJs](jsg::Lock& js, const auto& header) { return ToJs(js, header, cache); }));
214 
215 return obj;
216}
217 
218jsg::JsValue ToJs(jsg::Lock& js, const JsRpcEventInfo& info, StringCache& cache) {
219 auto obj = js.obj();
220 obj.set(js, TYPE_STR, cache.get(js, JSRPC_STR));
221 return obj;
222}
223 
224jsg::JsValue ToJs(jsg::Lock& js, const ScheduledEventInfo& info, StringCache& cache) {
225 auto obj = js.obj();
226 obj.set(js, TYPE_STR, cache.get(js, SCHEDULED_STR));
227 if (isPredictableModeForTest()) {
228 obj.set(js, SCHEDULEDTIME_STR, js.date(kj::UNIX_EPOCH));
229 } else {
230 obj.set(js, SCHEDULEDTIME_STR, js.date(info.scheduledTime));
231 }
232 obj.set(js, CRON_STR, js.str(info.cron));
233 return obj;
234}
235 
236jsg::JsValue ToJs(jsg::Lock& js, const AlarmEventInfo& info, StringCache& cache) {
237 auto obj = js.obj();
238 obj.set(js, TYPE_STR, cache.get(js, ALARM_STR));
239 if (isPredictableModeForTest()) {
240 obj.set(js, SCHEDULEDTIME_STR, js.date(kj::UNIX_EPOCH));
241 } else {
242 obj.set(js, SCHEDULEDTIME_STR, js.date(info.scheduledTime));
243 }
244 return obj;
245}
246 
247jsg::JsValue ToJs(jsg::Lock& js, const QueueEventInfo& info, StringCache& cache) {
248 auto obj = js.obj();
249 obj.set(js, TYPE_STR, cache.get(js, QUEUE_STR));
250 obj.set(js, QUEUENAME_STR, js.str(info.queueName));
251 obj.set(js, BATCHSIZE_STR, js.num(info.batchSize));
252 return obj;
253}
254 
255jsg::JsValue ToJs(jsg::Lock& js, const EmailEventInfo& info, StringCache& cache) {
256 auto obj = js.obj();
257 obj.set(js, TYPE_STR, cache.get(js, EMAIL_STR));
258 obj.set(js, MAILFROM_STR, js.str(info.mailFrom));
259 obj.set(js, RCPTTO_STR, js.str(info.rcptTo));
260 obj.set(js, RAWSIZE_STR, js.num(info.rawSize));
261 return obj;
262}
263 
264jsg::JsValue ToJs(jsg::Lock& js, const TraceEventInfo& info, StringCache& cache) {
265 auto obj = js.obj();
266 obj.set(js, TYPE_STR, cache.get(js, TRACE_STR));
267 obj.set(js, TRACES_STR,
268 js.arr(info.traces.asPtr(), [](jsg::Lock& js, const auto& trace) -> jsg::JsValue {
269 KJ_IF_SOME(name, trace.scriptName) {
270 return js.str(name);
271 }
272 return js.null();
273 }));
274 return obj;
275}
276 
277jsg::JsValue ToJs(jsg::Lock& js, const HibernatableWebSocketEventInfo& info, StringCache& cache) {
278 auto obj = js.obj();
279 obj.set(js, TYPE_STR, cache.get(js, HIBERNATABLEWEBSOCKET_STR));
280 
281 KJ_SWITCH_ONEOF(info.type) {
282 KJ_CASE_ONEOF(message, HibernatableWebSocketEventInfo::Message) {
283 auto mobj = js.obj();
284 mobj.set(js, TYPE_STR, cache.get(js, MESSAGE_STR));
285 obj.set(js, INFO_STR, mobj);
286 }
287 KJ_CASE_ONEOF(error, HibernatableWebSocketEventInfo::Error) {
288 auto mobj = js.obj();
289 mobj.set(js, TYPE_STR, cache.get(js, ERROR_STR));
290 obj.set(js, INFO_STR, mobj);
291 }
292 KJ_CASE_ONEOF(close, HibernatableWebSocketEventInfo::Close) {
293 auto mobj = js.obj();
294 mobj.set(js, TYPE_STR, cache.get(js, CLOSE_STR));
295 mobj.set(js, CODE_STR, js.num(close.code));
296 mobj.set(js, WASCLEAN_STR, js.boolean(close.wasClean));
297 obj.set(js, INFO_STR, mobj);
298 }
299 }
300 
301 return obj;
302}
303 
304jsg::JsValue ToJs(jsg::Lock& js, const ConnectEventInfo& info, StringCache& cache) {
305 auto obj = js.obj();
306 obj.set(js, TYPE_STR, cache.get(js, CONNECT_STR));
307 return obj;
308}
309 
310jsg::JsValue ToJs(jsg::Lock& js, const CustomEventInfo& info, StringCache& cache) {
311 auto obj = js.obj();
312 obj.set(js, TYPE_STR, cache.get(js, CUSTOM_STR));
313 return obj;
314}
315 
316jsg::JsValue ToJs(jsg::Lock& js, const EventOutcome& outcome, StringCache& cache) {
317 switch (outcome) {
318 case EventOutcome::OK:
319 return cache.get(js, OK_STR);
320 case EventOutcome::CANCELED:
321 return cache.get(js, CANCELED_STR);
322 case EventOutcome::EXCEPTION:
323 return cache.get(js, EXCEPTION_STR);
324 case EventOutcome::KILL_SWITCH:
325 return cache.get(js, KILLSWITCH_STR);
326 case EventOutcome::DAEMON_DOWN:
327 return cache.get(js, DAEMONDOWN_STR);
328 case EventOutcome::EXCEEDED_CPU:
329 return cache.get(js, EXCEEDEDCPU_STR);
330 case EventOutcome::EXCEEDED_MEMORY:
331 return cache.get(js, EXCEEDEDMEMORY_STR);
332 case EventOutcome::LOAD_SHED:
333 return cache.get(js, LOADSHED_STR);
334 case EventOutcome::RESPONSE_STREAM_DISCONNECTED:
335 return cache.get(js, RESPONSESTREAMDISCONNECTED_STR);
336 case EventOutcome::SCRIPT_NOT_FOUND:
337 return cache.get(js, SCRIPTNOTFOUND_STR);
338 case EventOutcome::INTERNAL_ERROR:
339 return cache.get(js, INTERNALERROR_STR);
340 case EventOutcome::UNKNOWN:
341 return cache.get(js, UNKNOWN_STR);
342 }
343 KJ_UNREACHABLE;
344}
345 
346jsg::JsValue ToJs(jsg::Lock& js, const Onset& onset, StringCache& cache) {
347 auto obj = js.obj();
348 obj.set(js, TYPE_STR, cache.get(js, ONSET_STR));
349 obj.set(js, EXECUTIONMODEL_STR, cache.get(js, kj::str(onset.workerInfo.executionModel)));
350 obj.set(js, SPANID_STR, js.str(onset.spanId.toGoString()));
351 
352 KJ_IF_SOME(ns, onset.workerInfo.dispatchNamespace) {
353 obj.set(js, DISPATCHNAMESPACE_STR, js.str(ns));
354 }
355 KJ_IF_SOME(entrypoint, onset.workerInfo.entrypoint) {
356 obj.set(js, ENTRYPOINT_STR, js.str(entrypoint));
357 }
358 KJ_IF_SOME(name, onset.workerInfo.scriptName) {
359 obj.set(js, SCRIPTNAME_STR, js.str(name));
360 }
361 KJ_IF_SOME(tags, onset.workerInfo.scriptTags) {
362 obj.set(js, SCRIPTTAGS_STR,
363 js.arr(tags.asPtr(), [](jsg::Lock& js, const kj::String& tag) { return js.str(tag); }));
364 }
365 KJ_IF_SOME(version, onset.workerInfo.scriptVersion) {
366 auto vobj = js.obj();
367 auto id = version->getId();
368 KJ_IF_SOME(uuid, UUID::fromUpperLower(id.getUpper(), id.getLower())) {
369 vobj.set(js, ID_STR, js.str(uuid.toString()));
370 }
371 if (version->hasTag()) {
372 vobj.set(js, TAG_STR, js.str(version->getTag()));
373 }
374 if (version->hasMessage()) {
375 vobj.set(js, MESSAGE_STR, js.str(version->getMessage()));
376 }
377 obj.set(js, SCRIPTVERSION_STR, vobj);
378 }
379 KJ_IF_SOME(preview, onset.workerInfo.preview) {
380 auto pobj = js.obj();
381 pobj.set(js, ID_STR, js.str(preview.id));
382 pobj.set(js, SLUG_STR, js.str(preview.slug));
383 pobj.set(js, NAME_STR, js.str(preview.name));
384 obj.set(js, PREVIEW_STR, pobj);
385 }
386 
387 KJ_SWITCH_ONEOF(onset.info) {
388 KJ_CASE_ONEOF(fetch, FetchEventInfo) {
389 obj.set(js, INFO_STR, ToJs(js, fetch, cache));
390 }
391 KJ_CASE_ONEOF(jsrpc, JsRpcEventInfo) {
392 obj.set(js, INFO_STR, ToJs(js, jsrpc, cache));
393 }
394 KJ_CASE_ONEOF(scheduled, ScheduledEventInfo) {
395 obj.set(js, INFO_STR, ToJs(js, scheduled, cache));
396 }
397 KJ_CASE_ONEOF(alarm, AlarmEventInfo) {
398 obj.set(js, INFO_STR, ToJs(js, alarm, cache));
399 }
400 KJ_CASE_ONEOF(queue, QueueEventInfo) {
401 obj.set(js, INFO_STR, ToJs(js, queue, cache));
402 }
403 KJ_CASE_ONEOF(email, EmailEventInfo) {
404 obj.set(js, INFO_STR, ToJs(js, email, cache));
405 }
406 KJ_CASE_ONEOF(trace, TraceEventInfo) {
407 obj.set(js, INFO_STR, ToJs(js, trace, cache));
408 }
409 KJ_CASE_ONEOF(hws, HibernatableWebSocketEventInfo) {
410 obj.set(js, INFO_STR, ToJs(js, hws, cache));
411 }
412 KJ_CASE_ONEOF(connect, ConnectEventInfo) {
413 obj.set(js, INFO_STR, ToJs(js, connect, cache));
414 }
415 KJ_CASE_ONEOF(custom, CustomEventInfo) {
416 obj.set(js, INFO_STR, ToJs(js, custom, cache));
417 }
418 }
419 
420 if (onset.attributes.size() > 0) {
421 obj.set(js, ATTRIBUTES_STR,
422 js.arr(onset.attributes.asPtr(),
423 [&cache](jsg::Lock& js, const auto& attr) { return ToJs(js, attr, cache); }));
424 }
425 
426 return obj;
427}
428 
429jsg::JsValue ToJs(jsg::Lock& js, const Outcome& outcome, StringCache& cache) {
430 auto obj = js.obj();
431 obj.set(js, TYPE_STR, cache.get(js, OUTCOME_STR));
432 obj.set(js, OUTCOME_STR, ToJs(js, outcome.outcome, cache));
433 
434 double cpuTime = outcome.cpuTime / kj::MILLISECONDS;
435 double wallTime = outcome.wallTime / kj::MILLISECONDS;
436 
437 obj.set(js, CPUTIME_STR, js.num(cpuTime));
438 obj.set(js, WALLTIME_STR, js.num(wallTime));
439 
440 return obj;
441}
442 
443jsg::JsValue ToJs(jsg::Lock& js, const SpanOpen& spanOpen, StringCache& cache) {
444 auto obj = js.obj();
445 obj.set(js, TYPE_STR, cache.get(js, SPANOPEN_STR));
446 obj.set(js, NAME_STR, js.str(spanOpen.operationName));
447 // Export span ID as non-truncated hex value – in practice this will be a random span ID.
448 obj.set(js, SPANID_STR, js.str(spanOpen.spanId.toGoString()));
449 
450 KJ_IF_SOME(info, spanOpen.info) {
451 KJ_SWITCH_ONEOF(info) {
452 KJ_CASE_ONEOF(fetch, FetchEventInfo) {
453 obj.set(js, INFO_STR, ToJs(js, fetch, cache));
454 }
455 KJ_CASE_ONEOF(jsrpc, JsRpcEventInfo) {
456 obj.set(js, INFO_STR, ToJs(js, jsrpc, cache));
457 }
458 KJ_CASE_ONEOF(custom, CustomInfo) {
459 obj.set(js, INFO_STR, ToJs(js, custom.asPtr(), cache));
460 }
461 }
462 }
463 return obj;
464}
465 
466jsg::JsValue ToJs(jsg::Lock& js, const SpanClose& spanClose, StringCache& cache) {
467 auto obj = js.obj();
468 obj.set(js, TYPE_STR, cache.get(js, SPANCLOSE_STR));
469 obj.set(js, OUTCOME_STR, ToJs(js, spanClose.outcome, cache));
470 return obj;
471}
472 
473jsg::JsValue ToJs(jsg::Lock& js, const DiagnosticChannelEvent& dce, StringCache& cache) {
474 auto obj = js.obj();
475 obj.set(js, TYPE_STR, cache.get(js, DIAGNOSTICCHANNEL_STR));
476 obj.set(js, CHANNEL_STR, cache.get(js, dce.channel));
477 jsg::Serializer::Released released{
478 .data = kj::heapArray<kj::byte>(dce.message),
479 };
480 jsg::Deserializer deser(js, released);
481 obj.set(js, MESSAGE_STR, deser.readValue(js));
482 return obj;
483}
484 
485jsg::JsValue ToJs(jsg::Lock& js, const Exception& ex, StringCache& cache) {
486 auto obj = js.obj();
487 obj.set(js, TYPE_STR, cache.get(js, EXCEPTION_STR));
488 obj.set(js, NAME_STR, cache.get(js, ex.name));
489 obj.set(js, MESSAGE_STR, js.str(ex.message));
490 KJ_IF_SOME(stack, ex.stack) {
491 obj.set(js, STACK_STR, js.str(stack));
492 }
493 return obj;
494}
495 
496jsg::JsValue ToJs(jsg::Lock& js, const LogLevel& level, StringCache& cache) {
497 switch (level) {
498 case LogLevel::DEBUG_:
499 return cache.get(js, DEBUG_STR);
500 case LogLevel::INFO:
501 return cache.get(js, INFO_STR);
502 case LogLevel::LOG:
503 return cache.get(js, LOG_STR);
504 case LogLevel::WARN:
505 return cache.get(js, WARN_STR);
506 case LogLevel::ERROR:
507 return cache.get(js, ERROR_STR);
508 }
509 KJ_UNREACHABLE;
510}
511 
512jsg::JsValue ToJs(jsg::Lock& js, const Log& log, StringCache& cache) {
513 auto obj = js.obj();
514 obj.set(js, TYPE_STR, cache.get(js, LOG_STR));
515 obj.set(js, LEVEL_STR, ToJs(js, log.logLevel, cache));
516 // TODO(o11y): Check that we are always returning an object here
517 obj.set(js, MESSAGE_STR, jsg::JsValue(js.parseJson(log.message).getHandle(js)));
518 return obj;
519}
520 
521jsg::JsValue ToJs(jsg::Lock& js, const StreamDiagnosticsEvent& streamDiag, StringCache& cache) {
522 auto obj = js.obj();
523 obj.set(js, TYPE_STR, cache.get(js, STREAMDIAGNOSTIC_STR));
524 // At present we only support the droppedEvents type.
525 
526 // Handle droppedEvents
527 auto droppedEventsDiagnostic = js.obj();
528 droppedEventsDiagnostic.set(js, DIAGNOSTICSTYPE_STR, cache.get(js, DROPPEDEVENTS_STR));
529 droppedEventsDiagnostic.set(js, COUNT_STR, js.num(streamDiag.droppedEventsCount));
530 obj.set(js, DIAGNOSTIC_STR, kj::mv(droppedEventsDiagnostic));
531 return obj;
532}
533 
534jsg::JsValue ToJs(jsg::Lock& js, const Return& ret, StringCache& cache) {
535 auto obj = js.obj();
536 obj.set(js, TYPE_STR, cache.get(js, RETURN_STR));
537 
538 KJ_IF_SOME(info, ret.info) {
539 obj.set(js, INFO_STR, ToJs(js, info, cache));
540 }
541 
542 return obj;
543}
544 
545jsg::JsValue ToJs(jsg::Lock& js, const TailEvent& event, StringCache& cache) {
546 auto obj = js.obj();
547 
548 // Set SpanContext
549 auto sCObj = js.obj();
550 sCObj.set(js, TRACEID_STR, js.str(event.spanContext.getTraceId().toGoString()));
551 KJ_IF_SOME(spanId, event.spanContext.getSpanId()) {
552 sCObj.set(js, SPANID_STR, js.str(spanId.toGoString()));
553 }
554 KJ_IF_SOME(flags, event.spanContext.getTraceFlags()) {
555 sCObj.set(js, TRACEFLAGS_STR, js.num(flags));
556 }
557 obj.set(js, SPANCONTEXT_STR, kj::mv(sCObj));
558 
559 obj.set(js, INVOCATIONID_STR, js.str(event.invocationId.toGoString()));
560 obj.set(js, TIMESTAMP_STR, js.date(event.timestamp));
561 obj.set(js, SEQUENCE_STR, js.num(event.sequence));
562 
563 KJ_SWITCH_ONEOF(event.event) {
564 KJ_CASE_ONEOF(onset, Onset) {
565 obj.set(js, EVENT_STR, ToJs(js, onset, cache));
566 }
567 KJ_CASE_ONEOF(outcome, Outcome) {
568 obj.set(js, EVENT_STR, ToJs(js, outcome, cache));
569 }
570 KJ_CASE_ONEOF(spanOpen, SpanOpen) {
571 obj.set(js, EVENT_STR, ToJs(js, spanOpen, cache));
572 }
573 KJ_CASE_ONEOF(spanClose, SpanClose) {
574 obj.set(js, EVENT_STR, ToJs(js, spanClose, cache));
575 }
576 KJ_CASE_ONEOF(de, DiagnosticChannelEvent) {
577 obj.set(js, EVENT_STR, ToJs(js, de, cache));
578 }
579 KJ_CASE_ONEOF(ex, Exception) {
580 obj.set(js, EVENT_STR, ToJs(js, ex, cache));
581 }
582 KJ_CASE_ONEOF(log, Log) {
583 obj.set(js, EVENT_STR, ToJs(js, log, cache));
584 }
585 KJ_CASE_ONEOF(diagEvent, StreamDiagnosticsEvent) {
586 obj.set(js, EVENT_STR, ToJs(js, diagEvent, cache));
587 }
588 KJ_CASE_ONEOF(ret, Return) {
589 obj.set(js, EVENT_STR, ToJs(js, ret, cache));
590 }
591 KJ_CASE_ONEOF(attrs, CustomInfo) {
592 obj.set(js, EVENT_STR, ToJs(js, attrs, cache));
593 }
594 }
595 
596 return obj;
597}
598 
599// Returns the name of the handler function for this type of event.
600kj::Maybe<kj::StringPtr> getHandlerName(const TailEvent& event) {
601 KJ_SWITCH_ONEOF(event.event) {
602 KJ_CASE_ONEOF(_, Onset) {
603 KJ_FAIL_ASSERT("Onset event should only be provided to tailStream(), not returned handler");
604 // return ONSET_STR;
605 }
606 KJ_CASE_ONEOF(_, Outcome) {
607 return OUTCOME_STR;
608 }
609 KJ_CASE_ONEOF(_, SpanOpen) {
610 return SPANOPEN_STR;
611 }
612 KJ_CASE_ONEOF(_, SpanClose) {
613 return SPANCLOSE_STR;
614 }
615 KJ_CASE_ONEOF(_, DiagnosticChannelEvent) {
616 return DIAGNOSTICCHANNEL_STR;
617 }
618 KJ_CASE_ONEOF(_, Exception) {
619 return EXCEPTION_STR;
620 }
621 KJ_CASE_ONEOF(_, Log) {
622 return LOG_STR;
623 }
624 KJ_CASE_ONEOF(_, StreamDiagnosticsEvent) {
625 return STREAMDIAGEVENT_STR;
626 }
627 KJ_CASE_ONEOF(_, Return) {
628 return RETURN_STR;
629 }
630 KJ_CASE_ONEOF(_, CustomInfo) {
631 return ATTRIBUTES_STR;
632 }
633 }
634 return kj::none;
635}
636 
637class TailStreamTarget final: public rpc::TailStreamTarget::Server {
638 public:
639 TailStreamTarget(IoContext& ioContext,
640 kj::Maybe<kj::StringPtr> entrypointNamePtr,
641 kj::Maybe<Worker::VersionInfo> versionInfo,
642 Frankenvalue props,
643 kj::Own<kj::PromiseFulfiller<void>> doneFulfiller,
644 bool isDynamicDispatch)
645 : weakIoContext(ioContext.getWeakRef()),
646 entrypointNamePtr(kj::mv(entrypointNamePtr)),
647 versionInfo(kj::mv(versionInfo)),
648 props(kj::mv(props)),
649 doneFulfiller(kj::mv(doneFulfiller)),
650 isDynamicDispatch(isDynamicDispatch) {}
651 
652 KJ_DISALLOW_COPY_AND_MOVE(TailStreamTarget);
653 ~TailStreamTarget() {
654 if (doneFulfiller->isWaiting()) {
655 doneFulfiller->reject(KJ_EXCEPTION(DISCONNECTED, "Streaming tail session canceled."));
656 }
657 }
658 
659 kj::Promise<void> report(ReportContext reportContext) override {
660 IoContext& ioContext = KJ_REQUIRE_NONNULL(weakIoContext->tryGet(),
661 "The destination object for this tail session no longer exists.", doneReceiving);
662 
663 ioContext.getLimitEnforcer().topUpActor();
664 
665 auto ownReportContext = capnp::CallContextHook::from(reportContext).addRef();
666 // We need to be able to access the results builder from both the promise below and its
667 // exception handler.
668 auto sharedResults = kj::rc<SharedResults>(reportContext.initResults());
669 
670 auto promise = ioContext.run([this, &ioContext, sharedResults = sharedResults.addRef(),
671 reportContext, ownReportContext = ownReportContext->addRef()](
672 Worker::Lock& lock) mutable -> kj::Promise<void> {
673 auto params = reportContext.getParams();
674 KJ_ASSERT(params.hasEvents(), "Events are required.");
675 auto eventReaders = params.getEvents();
676 kj::Array<TailEvent> events = KJ_MAP(reader, eventReaders) { return TailEvent(reader); };
677 
678 // If we have not yet received the onset event, the first event in the
679 // received collection must be an Onset event and must be handled separately.
680 // We will only dispatch the remaining events if a handler is returned.
681 auto result = ([&]() -> kj::Promise<void> {
682 KJ_IF_SOME(handler, maybeHandler) {
683 KJ_IF_SOME(h, handler.tryGet()) {
684 auto handle = h.getHandle(lock);
685 return handleEvents(lock, handle, ioContext, kj::mv(events), kj::mv(sharedResults));
686 } else {
687 KJ_LOG(ERROR, "tail stream handler was destroyed while processing events");
688 JSG_FAIL_REQUIRE(Error, "Tail stream handler became invalid during event processing");
689 KJ_UNREACHABLE;
690 }
691 } else {
692 return handleOnset(lock, ioContext, kj::mv(events), kj::mv(sharedResults));
693 }
694 })();
695 
696 if (ioContext.hasOutputGate()) {
697 return result.then([weakIoContext = weakIoContext->addRef()]() mutable {
698 return KJ_REQUIRE_NONNULL(weakIoContext->tryGet()).waitForOutputLocks();
699 });
700 } else {
701 return kj::mv(result);
702 }
703 });
704 
705 auto paf = kj::newPromiseAndFulfiller<void>();
706 promise = promise.then([&fulfiller = *paf.fulfiller]() { fulfiller.fulfill(); },
707 [&, &fulfiller = *paf.fulfiller, ownReportContext = kj::mv(ownReportContext),
708 results = kj::mv(sharedResults)](kj::Exception&& e) mutable {
709 // This is the top level exception catcher for tail events being delivered. We do not want to
710 // propagate JS exceptions to the client side here, all exceptions should stay within this
711 // customEvent. Instead, we propagate the exception to the doneFulfiller, where it is used to
712 // set the right outcome code and re-thrown if appropriate. By rejecting the doneFulfiller, we
713 // also ensure that no more tail events get delivered.
714 if (jsg::isTunneledException(e.getDescription())) {
715 auto description = jsg::stripRemoteExceptionPrefix(e.getDescription());
716 if (!description.startsWith("remote.")) {
717 e.setDescription(kj::str("remote.", description));
718 }
719 }
720 // We still fulfill this fulfiller to disarm the cancellation check below
721 fulfiller.fulfill();
722 results->setStop(true);
723 doneReceiving = true;
724 doneFulfiller->reject(kj::mv(e));
725 });
726 promise = promise.attach(kj::defer([fulfiller = kj::mv(paf.fulfiller)]() mutable {
727 if (fulfiller->isWaiting()) {
728 fulfiller->reject(JSG_KJ_EXCEPTION(FAILED, Error,
729 "The destination execution context for this tail session was canceled while the "
730 "call was still running."));
731 }
732 }));
733 ioContext.addTask(kj::mv(promise));
734 
735 return kj::mv(paf.promise);
736 }
737 
738 private:
739 // Used to share the results builder (and send the stop signal) from both the main code path and
740 // the exception handler.
741 struct SharedResults: public kj::Refcounted, rpc::TailStreamTarget::TailStreamResults::Builder {
742 SharedResults(rpc::TailStreamTarget::TailStreamResults::Builder results)
743 : rpc::TailStreamTarget::TailStreamResults::Builder(kj::mv(results)) {}
744 };
745 // Handles the very first (onset) event in the tail stream. This will cause
746 // the exported tailStream handler to be called, passing the onset event
747 // as the initial argument. If the tail stream wishes to continue receiving
748 // events for this invocation, it will return a handler in the form of an
749 // object or a function. If no handler is returned, the tail session is
750 // shutdown.
751 kj::Promise<void> handleOnset(Worker::Lock& lock,
752 IoContext& ioContext,
753 kj::Array<TailEvent> events,
754 kj::Rc<SharedResults> results) {
755 // There should be only a single onset event in this batch.
756 KJ_ASSERT(
757 events.size() == 1 && events[0].event.is<Onset>(), "Expected only a single onset event");
758 auto& event = events[0];
759 
760 auto handler =
761 KJ_REQUIRE_NONNULL(lock.getExportedHandler(entrypointNamePtr, kj::mv(versionInfo),
762 kj::mv(props), ioContext.getActor(), isDynamicDispatch),
763 "Failed to get handler to worker.");
764 StringCache stringCache;
765 
766 jsg::Lock& js = lock;
767 auto target = jsg::JsObject(handler->self.getHandle(js));
768 v8::Local<v8::Value> maybeFn = target.get(js, "tailStream"_kj);
769 
770 // If there's no actual tailStream handler, or if the tailStream export is
771 // something other than a function, we will emit a warning for the user
772 // then immediately return.
773 if (!maybeFn->IsFunction()) {
774 ioContext.logWarningOnce("A worker configured to act as a streaming tail worker does "
775 "not export a tailStream() handler.");
776 results->setStop(true);
777 doneReceiving = true;
778 doneFulfiller->fulfill();
779 return kj::READY_NOW;
780 }
781 
782 // Invoke the tailStream handler function.
783 v8::Local<v8::Function> fn = maybeFn.As<v8::Function>();
784 kj::Maybe<v8::Local<v8::Object>> maybeCtx;
785 KJ_IF_SOME(hCtx, handler->getCtx()) {
786 maybeCtx = v8::Local<v8::Object>(
787 lock.getWorker().getIsolate().getApi().wrapExecutionContext(js, kj::mv(hCtx)));
788 }
789 v8::LocalVector<v8::Value> handlerArgs(js.v8Isolate, maybeCtx != kj::none ? 3 : 2);
790 handlerArgs[0] = ToJs(js, event, stringCache);
791 handlerArgs[1] = handler->env.getHandle(js);
792 KJ_IF_SOME(ctx, maybeCtx) {
793 handlerArgs[2] = ctx;
794 }
795 
796 try {
797 auto result =
798 jsg::check(fn->Call(js.v8Context(), target, handlerArgs.size(), handlerArgs.data()));
799 
800 // The handler can return a function, an object, undefined, or a promise
801 // for any of these. We will convert the result to a promise for consistent
802 // handling...
803 return ioContext.awaitJs(js,
804 js.toPromise(result).then(js,
805 ioContext.addFunctor([this, results = results.addRef(), &ioContext](
806 jsg::Lock& js, jsg::Value value) mutable {
807 // The value here can be one of a function, an object, or undefined.
808 // Any value other than these will result in a warning but will otherwise
809 // be treated like undefined.
810 
811 // If a function or object is returned, then our tail worker wishes to
812 // keep receiving events! Yay! Otherwise, we will stop the stream by
813 // setting the stop field in the results.
814 auto handle = value.getHandle(js);
815 if (handle->IsFunction() || handle->IsObject()) {
816 // Sweet! Our tail worker wants to keep receiving events. Let's store
817 // the handler and return.
818 maybeHandler = ioContext.addObjectReverse(
819 kj::heap<jsg::JsRef<jsg::JsValue>>(js, jsg::JsValue(handle)));
820 return;
821 }
822 
823 // If the handler returned any other kind of value, let's be nice and
824 // at least warn the user about it.
825 if (!handle->IsUndefined()) {
826 ioContext.logWarningOnce(
827 kj::str("tailStream() handler returned an unusable value. "
828 "The tailStream() handler is expected to return either a function, an "
829 "object, or undefined. Received ",
830 jsg::JsValue(handle).typeOf(js)));
831 }
832 // And finally, we'll stop the stream since the tail worker did not return
833 // a handler for us to continue with.
834 results->setStop(true);
835 doneReceiving = true;
836 doneFulfiller->fulfill();
837 }),
838 ioContext.addFunctor(
839 [&, results = results.addRef()](jsg::Lock& js, jsg::Value&& error) mutable {
840 // Received a JS error. Do not reject doneFulfiller yet, this will be handled when we catch
841 // the exception later.
842 results->setStop(true);
843 doneReceiving = true;
844 js.throwException(kj::mv(error));
845 })));
846 } catch (...) {
847 ioContext.logWarningOnce("A worker configured to act as a streaming tail worker did "
848 "not return a valid tailStream() handler.");
849 results->setStop(true);
850 doneReceiving = true;
851 doneFulfiller->fulfill();
852 return kj::READY_NOW;
853 }
854 KJ_UNREACHABLE;
855 }
856 
857 kj::Promise<void> handleEvents(Worker::Lock& lock,
858 const jsg::JsValue& handler,
859 IoContext& ioContext,
860 kj::Array<TailEvent> events,
861 kj::Rc<SharedResults> results) {
862 jsg::Lock& js = lock;
863 
864 // Should not ever happen but let's handle it anyway.
865 if (events.size() == 0) return kj::READY_NOW;
866 
867 // Take the received set of events and dispatch them to the correct handler.
868 
869 v8::Local<v8::Value> h = handler;
870 v8::LocalVector<v8::Value> returnValues(js.v8Isolate);
871 StringCache stringCache;
872 
873 // If any of the events delivered are an outcome event, we will signal that
874 // the stream should be stopped and will fulfill the done promise.
875 bool finishing = false;
876 
877 // When a tail worker receives its outcome event, we need to ensure that the final tail worker
878 // invocation is completed before destroying the tail worker customEvent and incomingRequest. To
879 // achieve this, we only fulfill the doneFulfiller after JS execution has completed.
880 bool doFulfill = false;
881 
882 for (auto& event: events) {
883 // If we already received an outcome event, we will stop processing any
884 // further events.
885 if (finishing) break;
886 if (event.event.is<Outcome>()) {
887 finishing = true;
888 results->setStop(true);
889 doneReceiving = true;
890 // We set doFulfill to indicate that the outcome event has been received via RPC and no more
891 // events are expected.
892 doFulfill = true;
893 };
894 
895 v8::Local<v8::Value> eventObj = ToJs(js, event, stringCache);
896 if (h->IsFunction()) {
897 // If the handler is a function, then we'll just pass all of the events to that
898 // function. If the function returns a promise and there are multiple events we
899 // will not wait for each promise to resolve before calling the next iteration.
900 // But we will wait for all promises to settle before returning the resolved
901 // kj promise.
902 auto fn = h.As<v8::Function>();
903 returnValues.push_back(jsg::check(fn->Call(js.v8Context(), h, 1, &eventObj)));
904 } else {
905 // If the handler is an object, then we need to know what kind of events
906 // we have and look for a specific handler function for each.
907 KJ_ASSERT(h->IsObject());
908 KJ_IF_SOME(name, getHandlerName(event)) {
909 jsg::JsObject obj = jsg::JsObject(h.As<v8::Object>());
910 v8::Local<v8::Value> val = obj.get(js, name);
911 // If the value is not a function, we'll ignore it entirely.
912 if (val->IsFunction()) {
913 auto fn = val.As<v8::Function>();
914 returnValues.push_back(jsg::check(fn->Call(js.v8Context(), h, 1, &eventObj)));
915 }
916 }
917 }
918 }
919 // We want the equivalent behavior to Promise.all([...]) here but v8 does not
920 // give us a C++ equivalent of Promise.all([...]) so we need to approximate it.
921 // We do so by chaining all of the promises together.
922 kj::Maybe<jsg::Promise<void>> promise;
923 for (auto& val: returnValues) {
924 KJ_IF_SOME(p, promise) {
925 promise = p.then(js,
926 [p = js.toPromise(val).whenResolved(js)](jsg::Lock& js) mutable { return kj::mv(p); });
927 } else {
928 promise = js.toPromise(val).whenResolved(js);
929 }
930 }
931 
932 KJ_IF_SOME(p, promise) {
933 // When doFulfill is set, the last promise refers to the outcome event. In that case the chain
934 // of promises provides all remaining events to the user tail handler, so we should fulfill
935 // the doneFulfiller afterwards, indicating that TailStreamTarget has received all events over
936 // the stream and has done all its work, that the stream self-evidently did not get canceled
937 // prematurely. This applies even if promises were rejected.
938 // No need to catch exceptions here: They will be handled in report() alongside exceptions
939 // from the onset event etc. JSG knows how JS exceptions look like, so we don't need an
940 // identifier for them.
941 if (doFulfill) {
942 p = p.then(js, [&](jsg::Lock& js) {
943 doneReceiving = true;
944 doneFulfiller->fulfill();
945 });
946 }
947 return ioContext.awaitJs(js, kj::mv(p));
948 } else if (doFulfill) {
949 // If we have no promises, but doFulfill is true, then none of the events we have had a
950 // handler available, but we still need to indicate that we are done since we got the outcome
951 // event – do this by calling fulfill right away.
952 doneReceiving = true;
953 doneFulfiller->fulfill();
954 }
955 return kj::READY_NOW;
956 }
957 
958 kj::Own<IoContext::WeakRef> weakIoContext;
959 kj::Maybe<kj::StringPtr> entrypointNamePtr;
960 kj::Maybe<Worker::VersionInfo> versionInfo;
961 Frankenvalue props;
962 // The done fulfiller is resolved when we receive the outcome event
963 // or rejected if the capability is dropped before receiving the outcome
964 // event.
965 kj::Own<kj::PromiseFulfiller<void>> doneFulfiller;
966 bool isDynamicDispatch;
967 
968 // The maybeHandler will be empty until we receive and process the
969 // onset event.
970 kj::Maybe<ReverseIoOwn<jsg::JsRef<jsg::JsValue>>> maybeHandler;
971 
972 // Indicates that we told (or should have told) the client that we want no further events, used
973 // to debug events arriving when the IoContext is no longer valid.
974 bool doneReceiving = false;
975};
976} // namespace
977 
978EventInfo TailStreamCustomEvent::getEventInfo() const {
979 return TraceEventInfo(kj::Array<TraceEventInfo::TraceItem>(nullptr));
980}
981 
982kj::Promise<WorkerInterface::CustomEvent::Result> TailStreamCustomEvent::run(
983 kj::Own<IoContext::IncomingRequest> incomingRequest,
984 kj::Maybe<kj::StringPtr> entrypointName,
985 kj::Maybe<Worker::VersionInfo> versionInfo,
986 Frankenvalue props,
987 kj::TaskSet& waitUntilTasks,
988 bool isDynamicDispatch) {
989 IoContext& ioContext = incomingRequest->getContext();
990 incomingRequest->delivered();
991 
992 auto [donePromise, doneFulfiller] = kj::newPromiseAndFulfiller<void>();
993 capFulfiller->fulfill(kj::heap<TailStreamTarget>(ioContext, kj::mv(entrypointName),
994 kj::mv(versionInfo), kj::mv(props), kj::mv(doneFulfiller), isDynamicDispatch));
995 
996 donePromise = donePromise.attach(ioContext.registerPendingEvent());
997 
998 KJ_DEFER({
999 // waitUntil() should allow extending execution on the server side even when the client
1000 // disconnects.
1001 waitUntilTasks.add(incomingRequest->drain().attach(kj::mv(incomingRequest)));
1002 });
1003 
1004 auto eventOutcome = co_await donePromise.exclusiveJoin(ioContext.onAbort()).then([&]() {
1005 return ioContext.waitUntilStatus();
1006 }, [&incomingRequest](kj::Exception&& e) {
1007 // If we have a JSG exception, just set the appropriate return code – this will already have
1008 // been logged and we do not need to treat it like a KJ exception. Otherwise, re-throw the
1009 // exception.
1010 if (jsg::isTunneledException(e.getDescription())) {
1011 incomingRequest->getMetrics().reportFailure(e);
1012 return EventOutcome::EXCEPTION;
1013 }
1014 kj::throwRecoverableException(kj::mv(e));
1015 KJ_UNREACHABLE;
1016 });
1017 KJ_IF_SOME(t, ioContext.getWorkerTracer()) {
1018 t.setReturn(ioContext.now());
1019 }
1020 
1021 co_return WorkerInterface::CustomEvent::Result{.outcome = eventOutcome};
1022}
1023 
1024kj::Promise<WorkerInterface::CustomEvent::Result> TailStreamCustomEvent::sendRpc(
1025 capnp::HttpOverCapnpFactory& httpOverCapnpFactory,
1026 capnp::ByteStreamFactory& byteStreamFactory,
1027 rpc::EventDispatcher::Client dispatcher) {
1028 auto revokePaf = kj::newPromiseAndFulfiller<void>();
1029 
1030 KJ_DEFER({
1031 if (revokePaf.fulfiller->isWaiting()) {
1032 revokePaf.fulfiller->reject(KJ_EXCEPTION(DISCONNECTED, "Streaming tail session canceled"));
1033 }
1034 });
1035 
1036 auto req = dispatcher.tailStreamSessionRequest();
1037 auto sent = req.send();
1038 
1039 rpc::TailStreamTarget::Client cap = sent.getTopLevel();
1040 
1041 cap = capnp::membrane(kj::mv(cap), kj::refcounted<RevokerMembrane>(kj::mv(revokePaf.promise)));
1042 
1043 auto completionPaf = kj::newPromiseAndFulfiller<void>();
1044 cap = capnp::membrane(
1045 kj::mv(cap), kj::refcounted<CompletionMembrane>(kj::mv(completionPaf.fulfiller)));
1046 
1047 capFulfiller->fulfill(kj::mv(cap));
1048 
1049 // Forked promise for completion of all capabilities associated with the cap stream. This is
1050 // expected to be resolved when the request is canceled or when the client receives the stop
1051 // signal and deallocates cap after the tail worker indicates that it has processed all events
1052 // successfully.
1053 kj::ForkedPromise<void> forked = completionPaf.promise.fork();
1054 try {
1055 EventOutcome outcome = co_await sent.then([](auto resp) {
1056 return resp.getResult();
1057 }).exclusiveJoin(forked.addBranch().then([]() { return EventOutcome::CANCELED; }));
1058 
1059 // If the sent promise returned first, we still need to wait for the parent process to drop the
1060 // capability (which should happen right after it receives the stop signal) so that no
1061 // capabilities remain in an incomplete state when we return.
1062 co_await forked.addBranch();
1063 co_return WorkerInterface::CustomEvent::Result{.outcome = outcome};
1064 } catch (...) {
1065 auto e = kj::getCaughtExceptionAsKj();
1066 if (revokePaf.fulfiller->isWaiting()) {
1067 revokePaf.fulfiller->reject(e.clone());
1068 }
1069 kj::throwFatalException(kj::mv(e));
1070 }
1071}
1072 
1073TailStreamWriter::TailStreamWriter(Pending pending, kj::TaskSet& waitUntilTasks)
1074 : inner(kj::mv(pending)),
1075 waitUntilTasks(waitUntilTasks) {}
1076 
1077bool TailStreamWriter::reportImpl(TailEvent&& event, size_t sizeHint) {
1078 // In reportImpl, our inner state must be active.
1079 auto& actives = KJ_ASSERT_NONNULL(inner.tryGet<kj::Vector<kj::Own<Active>>>());
1080 
1081 // We only care about sessions that are currently active, removing any inactive ones.
1082 auto activeEnd = std::remove_if(actives.begin(), actives.end(),
1083 [](const auto& active) { return active->capability == kj::none; });
1084 if (activeEnd == actives.begin()) {
1085 // Oh! We have no active sessions. Well, never mind then, let's
1086 // transition to a closed state and drop everything on the floor.
1087 inner = Closed{};
1088 
1089 // Since we have no more living sessions (e.g. because all tail workers failed to return a valid
1090 // handler), mark the state as closing as we can't handle future events anyway.
1091 return true;
1092 }
1093 
1094 // We have at least some active sessions. Truncate the array to get rid of any inactive ones.
1095 if (activeEnd != actives.end()) {
1096 actives.truncate(activeEnd - actives.begin());
1097 }
1098 
1099 // We do not expect any events after the outcome.
1100 bool isClosing = event.event.is<Outcome>();
1101 // Deliver the event to the queue and make sure we are processing.
1102 for (auto& active: actives) {
1103 // Only queue the event if we don't have an excessive queue size yet. Return and Outcome
1104 // events are only provided once and thus won't be dropped.
1105 if (active->queueSize < maxQueueSize || event.event.is<Outcome>() || event.event.is<Return>()) {
1106 // When we get to the outcome, no more events will be dropped. Inject an internal diagnostics
1107 // event indicating how many events were dropped if applicable.
1108 if (event.event.is<Outcome>() && active->droppedEvents > 0) {
1109 StreamDiagnosticsEvent diag(active->droppedEvents);
1110 TailEvent diagTailEvent(SpanContext::clone(event.spanContext), event.invocationId,
1111 event.timestamp, event.sequence, kj::mv(diag));
1112 active->queue.push(kj::mv(diagTailEvent));
1113 // Increment the outcome sequence number to keep things consistent.
1114 event.sequence++;
1115 }
1116 
1117 // Optimization: Elide copy for last tail worker, helpful for common case of only one STW
1118 // being present.
1119 if (&active == &actives.back()) {
1120 active->queue.push(kj::mv(event));
1121 } else {
1122 active->queue.push(event.clone());
1123 }
1124 // Adjust estimated queue size based on size hint and an arbitrary amount for serialization
1125 // overhead. As long as this estimate is reasonably accurate, we won't need to check the
1126 // size again when serializing the message.
1127 active->queueSize += tailSerializationOverhead + sizeHint;
1128 } else {
1129 active->droppedEvents++;
1130 }
1131 
1132 if (!active->pumping) {
1133 waitUntilTasks.add(pump(kj::addRef(*active)));
1134 }
1135 }
1136 
1137 return isClosing;
1138}
1139 
1140// Delivers the queued tail events to a streaming tail worker.
1141//
1142// Note: An invocation of pump() may outlive the TailStreamWriter, as it is placed in
1143// `waitUntilTasks`. Hence, it is declared `static`, and owns a strong ref to its `Active`.
1144kj::Promise<void> TailStreamWriter::pump(kj::Own<Active> current) {
1145 current->pumping = true;
1146 KJ_DEFER(current->pumping = false);
1147 
1148 try {
1149 if (!current->onsetSeen) {
1150 // Our first event... yay! Our first job here will be to dispatch
1151 // the onset event to the tail worker. If the tail worker wishes
1152 // to handle the remaining events in the stream, then it will return
1153 // a new capability to which those would be reported. This is done
1154 // via the "result.getPipeline()" API below. If hasPipeline()
1155 // returns false then that means the tail worker did not return
1156 // a handler for this stream and no further attempts to deliver
1157 // events should be made for this stream.
1158 current->onsetSeen = true;
1159 auto onsetEvent = KJ_ASSERT_NONNULL(current->queue.pop());
1160 auto builder = KJ_ASSERT_NONNULL(current->capability).reportRequest();
1161 auto eventsBuilder = builder.initEvents(1);
1162 // When sending the onset event to the tail worker, the receiving end
1163 // requires that the onset event be delivered separately, without any
1164 // other events in the bundle. So here we'll separate it out and deliver
1165 // just the one event...
1166 onsetEvent.copyTo(eventsBuilder[0]);
1167 auto result = co_await builder.send();
1168 if (result.getStop()) {
1169 // If our call to send returns a stop signal, then we'll clear
1170 // the capability and be done.
1171 current->queue.clear();
1172 current->capability = kj::none;
1173 co_return;
1174 }
1175 }
1176 
1177 // If we got this far then we have a handler for all of our events.
1178 // Deliver remaining streaming tail events in batches if possible.
1179 while (!current->queue.empty()) {
1180 auto builder = KJ_ASSERT_NONNULL(current->capability).reportRequest();
1181 auto eventsBuilder = builder.initEvents(current->queue.size());
1182 size_t n = 0;
1183 
1184 // We're synchronously draining the queue – reset its size.
1185 current->queueSize = 0;
1186 current->queue.drainTo([&](TailEvent&& event) { event.copyTo(eventsBuilder[n++]); });
1187 
1188 auto result = co_await builder.send();
1189 
1190 // Note that although we cleared the current.queue above, it is
1191 // possible/likely that additional events were added to the queue
1192 // while the above builder.send() was being awaited. If the result
1193 // comes back indicating that we should stop, then we'll stop here
1194 // without any further processing. We'll defensively clear the
1195 // queue again and drop the client stub. Otherwise, if result.getStop()
1196 // is false, we'll loop back around to send any items that have since
1197 // been added to the queue or exit this loop if there are no additional
1198 // events waiting to be sent.
1199 if (result.getStop()) {
1200 current->queue.clear();
1201 current->capability = kj::none;
1202 co_return;
1203 }
1204 }
1205 } catch (...) {
1206 // If any RPC throws an exception, we should treat it as a stop signal, as this suggests
1207 // the connection to the STW itself has been lost. (An excpetion thrown within the STW
1208 // itself would have resulted in a `stop` return value instead of an exception over RPC.)
1209 current->queue.clear();
1210 current->capability = kj::none;
1211 throw;
1212 }
1213}
1214 
1215// If we are using streaming tail workers, initialize the mechanism that will deliver events
1216// to that collection of tail workers.
1217kj::Maybe<kj::Own<TailStreamWriter>> initializeTailStreamWriter(
1218 kj::Array<kj::Own<WorkerInterface>> streamingTailWorkers, kj::TaskSet& waitUntilTasks) {
1219 if (streamingTailWorkers.size() == 0) {
1220 return kj::none;
1221 }
1222 
1223 return kj::heap<TailStreamWriter>(kj::mv(streamingTailWorkers), waitUntilTasks);
1224}
1225 
1226void TailStreamWriter::report(const InvocationSpanContext& context,
1227 TailEvent::Event&& event,
1228 kj::Date timestamp,
1229 size_t sizeHint) {
1230 // Becomes a no-op if a terminal event (close) has been reported, or if the stream closed due to
1231 // not receiving a well-formed event handler. We need to disambiguate these cases as the former
1232 // indicates an implementation error resulting in trailing events whereas the latter case is
1233 // caused by a user error and events being reported after the stream being closed are expected –
1234 // reject events following an outcome event, but otherwise just exit if the state has been closed.
1235 // This could be an assert, but just log an error in case this is prevalent in some edge case.
1236 if (outcomeSeen) {
1237 KJ_LOG(ERROR, "reported tail stream event after stream close ", event, kj::getStackTrace());
1238 }
1239 if (inner.is<Closed>()) {
1240 return;
1241 }
1242 // The onset event must be first and must only happen once.
1243 if (event.is<Onset>()) {
1244 KJ_ASSERT(!onsetSeen, "Tail stream onset already provided");
1245 onsetSeen = true;
1246 } else {
1247 KJ_ASSERT(onsetSeen, "Tail stream onset was not reported");
1248 if (event.is<Outcome>()) {
1249 outcomeSeen = true;
1250 }
1251 }
1252 
1253 // A zero spanId at the TailEvent level signifies that no spanId should be provided to the tail
1254 // worker (for Onset events). We go to great lengths to rule out getting an all-zero spanId by
1255 // chance (see SpanId::fromEntropy()), so this should be safe.
1256 TailEvent tailEvent(context.getTraceId(), context.getInvocationId(),
1257 context.getSpanId() == SpanId::nullId ? kj::none : kj::Maybe(context.getSpanId()), timestamp,
1258 sequence++, kj::mv(event), context.getTraceFlags());
1259 
1260 KJ_SWITCH_ONEOF(inner) {
1261 KJ_CASE_ONEOF(closed, Closed) {
1262 // The tail stream has already been closed because we have received an outcome event. The
1263 // writer should have failed and we actually shouldn't get here. Assert!
1264 KJ_FAIL_ASSERT("tracing::TailStreamWriter report callback invoked after close");
1265 }
1266 KJ_CASE_ONEOF(pending, Pending) {
1267 // This is our first event! It has to be an onset event as we have validated above. Start each
1268 // of our tail working sessions.
1269 
1270 // Transitions into the active state by grabbing the pending client capability.
1271 inner = kj::Vector<kj::Own<Active>>( KJ_MAP(wi, pending) {
1272 auto customEvent = kj::heap<TailStreamCustomEvent>();
1273 auto result = customEvent->getCap();
1274 auto active = kj::refcounted<Active>(kj::mv(result));
1275 
1276 // Attach the workerInterface and customEvent to the waitUntil tasks so that they stay alive
1277 // until tail worker operations including JS execution are complete, including returning the
1278 // outcome.
1279 waitUntilTasks.add(wi->customEvent(kj::mv(customEvent))
1280 .attach(kj::mv(wi), kj::addRef(*active))
1281 .ignoreResult());
1282 return active;
1283 });
1284 
1285 // At this point our writer state is "active", which means the state consists of one or more
1286 // streaming tail worker client stubs to which the event will be dispatched.
1287 }
1288 KJ_CASE_ONEOF(active, kj::Vector<kj::Own<Active>>) {
1289 // active tail stream writers have already been configured, process the event.
1290 }
1291 }
1292 
1293 // The state is determined to be closing when it receives a terminal event (tracing::Outcome),
1294 // or if there are no active tail workers left, we can close the internal state at that point.
1295 if (reportImpl(kj::mv(tailEvent), sizeHint)) {
1296 inner = Closed{};
1297 }
1298}
1299 
1300} // namespace workerd::tracing