Skip to content
File

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

59.0 KB
1// Copyright (c) 2017-2022 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#include <workerd/io/trace.h>
6#include <workerd/util/entropy.h>
7#include <workerd/util/thread-scopes.h>
8 
9#include <capnp/message.h>
10#include <capnp/schema.h>
11#include <kj/compat/http.h>
12#include <kj/debug.h>
13#include <kj/time.h>
14 
15#include <atomic>
16#include <cstdlib>
17 
18namespace workerd {
19 
20namespace tracing {
21namespace {
22kj::Maybe<kj::uint> tryFromHexDigit(char c) {
23 if ('0' <= c && c <= '9') {
24 return c - '0';
25 } else if ('a' <= c && c <= 'f') {
26 return c - ('a' - 10);
27 } else if ('A' <= c && c <= 'F') {
28 return c - ('A' - 10);
29 } else {
30 return kj::none;
31 }
32}
33 
34kj::Maybe<uint64_t> hexToUint64(kj::ArrayPtr<const char> s) {
35 KJ_ASSERT(s.size() <= 16);
36 uint64_t value = 0;
37 for (auto ch: s) {
38 KJ_IF_SOME(d, tryFromHexDigit(ch)) {
39 value = (value << 4) + d;
40 } else {
41 return kj::none;
42 }
43 }
44 return value;
45}
46 
47void addHex(kj::Vector<char>& out, uint64_t v) {
48 constexpr char HEX_DIGITS[] = "0123456789abcdef";
49 for (int i = 0; i < 16; ++i) {
50 out.add(HEX_DIGITS[v >> (64 - 4)]);
51 v = v << 4;
52 }
53};
54 
55void addBigEndianBytes(kj::Vector<byte>& out, uint64_t v) {
56 for (int i = 0; i < 8; ++i) {
57 out.add(v >> (64 - 8));
58 v = v << 8;
59 }
60};
61} // namespace
62 
63// Reference: https://github.com/jaegertracing/jaeger/blob/e46f8737/model/ids.go#L58
64kj::Maybe<TraceId> TraceId::fromGoString(kj::ArrayPtr<const char> s) {
65 auto n = s.size();
66 if (n > 32) {
67 return kj::none;
68 } else if (n <= 16) {
69 KJ_IF_SOME(low, hexToUint64(s)) {
70 return TraceId(low, 0);
71 }
72 } else {
73 KJ_IF_SOME(high, hexToUint64(s.slice(0, n - 16))) {
74 KJ_IF_SOME(low, hexToUint64(s.slice(n - 16, n))) {
75 return TraceId(low, high);
76 }
77 }
78 }
79 return kj::none;
80}
81 
82// Reference: https://github.com/jaegertracing/jaeger/blob/e46f8737/model/ids.go#L50
83kj::String TraceId::toGoString() const {
84 if (high == 0) {
85 kj::Vector<char> s(17);
86 addHex(s, low);
87 s.add('\0');
88 return kj::String(s.releaseAsArray());
89 }
90 kj::Vector<char> s(33);
91 addHex(s, high);
92 addHex(s, low);
93 s.add('\0');
94 return kj::String(s.releaseAsArray());
95}
96 
97// Reference: https://github.com/jaegertracing/jaeger/blob/e46f8737/model/ids.go#L111
98kj::Maybe<TraceId> TraceId::fromProtobuf(kj::ArrayPtr<const byte> buf) {
99 if (buf.size() != 16) {
100 return kj::none;
101 }
102 uint64_t high = 0;
103 for (auto i: kj::zeroTo(8)) {
104 high = (high << 8) + buf[i];
105 }
106 uint64_t low = 0;
107 for (auto i: kj::zeroTo(8)) {
108 low = (low << 8) + buf[i + 8];
109 }
110 return TraceId(low, high);
111}
112 
113// Reference: https://github.com/jaegertracing/jaeger/blob/e46f8737/model/ids.go#L81
114kj::Array<byte> TraceId::toProtobuf() const {
115 kj::Vector<byte> s(16);
116 addBigEndianBytes(s, high);
117 addBigEndianBytes(s, low);
118 return s.releaseAsArray();
119}
120 
121// Reference https://www.w3.org/TR/trace-context/#trace-id
122kj::String TraceId::toW3C() const {
123 kj::Vector<char> s(32);
124 addHex(s, high);
125 addHex(s, low);
126 return kj::str(s.releaseAsArray());
127}
128 
129namespace {
130uint64_t getRandom64Bit(const kj::Maybe<kj::EntropySource&>& entropySource) {
131 uint64_t ret = 0;
132 uint8_t tries = 0;
133 
134 do {
135 tries++;
136 KJ_IF_SOME(entropy, entropySource) {
137 entropy.generate(kj::asBytes(ret));
138 } else {
139 getEntropy(kj::asBytes(ret));
140 }
141 // On the extreme off chance that we ended with with zeroes
142 // let's try again, but only up to three times.
143 } while (ret == 0 && tries < 3);
144 
145 return ret;
146}
147} // namespace
148 
149TraceId TraceId::fromEntropy(kj::Maybe<kj::EntropySource&> entropySource) {
150 if (isPredictableModeForTest()) {
151 // Produce deterministic but distinct IDs per call so that traceIds of
152 // independent (untriggered) invocations don't collide -- collisions confuse
153 // test fixtures that key spans by traceId or invocationId. Triggered
154 // invocations still inherit their caller's traceId via newForInvocation, so
155 // the propagation chain is preserved.
156 static std::atomic<uint64_t> counter{0};
157 uint64_t n = counter.fetch_add(1, std::memory_order_relaxed);
158 return TraceId(staticSpanId ^ n, staticSpanId ^ n);
159 }
160 
161 return TraceId(getRandom64Bit(entropySource), getRandom64Bit(entropySource));
162}
163 
164kj::String SpanId::toGoString() const {
165 kj::Vector<char> s(16);
166 addHex(s, id);
167 s.add('\0');
168 return kj::String(s.releaseAsArray());
169}
170 
171SpanId SpanId::fromEntropy(kj::Maybe<kj::EntropySource&> entropySource) {
172 return SpanId(getRandom64Bit(entropySource));
173}
174 
175kj::String KJ_STRINGIFY(const SpanId& id) {
176 return id;
177}
178 
179kj::String KJ_STRINGIFY(const TraceId& id) {
180 return id;
181}
182 
183InvocationSpanContext::InvocationSpanContext(kj::Badge<InvocationSpanContext>,
184 kj::Maybe<kj::EntropySource&> entropySource,
185 TraceId traceId,
186 TraceId invocationId,
187 SpanId spanId,
188 kj::Maybe<const InvocationSpanContext&> parentSpanContext,
189 kj::Maybe<TraceFlags> traceFlags)
190 : entropySource(entropySource),
191 traceId(kj::mv(traceId)),
192 invocationId(kj::mv(invocationId)),
193 spanId(kj::mv(spanId)),
194 parentSpanContext(parentSpanContext.map([](const InvocationSpanContext& ctx) {
195 return kj::heap<InvocationSpanContext>(ctx.clone());
196 })),
197 traceFlags(kj::mv(traceFlags)) {}
198 
199InvocationSpanContext InvocationSpanContext::newChild() const {
200 KJ_ASSERT(!isTrigger(), "unable to create child spans on this context");
201 kj::Maybe<kj::EntropySource&> otherEntropySource = entropySource.map(
202 [](auto& es) -> kj::EntropySource& { return const_cast<kj::EntropySource&>(es); });
203 return InvocationSpanContext(kj::Badge<InvocationSpanContext>(), otherEntropySource, traceId,
204 invocationId, SpanId::fromEntropy(otherEntropySource), *this, traceFlags);
205}
206 
207InvocationSpanContext InvocationSpanContext::newForInvocation(
208 kj::Maybe<const InvocationSpanContext&> triggerContext,
209 kj::Maybe<kj::EntropySource&> entropySource) {
210 kj::Maybe<const InvocationSpanContext&> parent;
211 kj::Maybe<TraceFlags> flags;
212 auto traceId = triggerContext
213 .map([&](auto& ctx) mutable {
214 parent = ctx;
215 flags = ctx.traceFlags;
216 return ctx.traceId;
217 }).orDefault([&] { return TraceId::fromEntropy(entropySource); });
218 return InvocationSpanContext(kj::Badge<InvocationSpanContext>(), entropySource, kj::mv(traceId),
219 TraceId::fromEntropy(entropySource), SpanId::fromEntropy(entropySource), kj::mv(parent),
220 kj::mv(flags));
221}
222 
223TraceId TraceId::fromCapnp(rpc::TraceId::Reader reader) {
224 return TraceId(reader.getLow(), reader.getHigh());
225}
226 
227void TraceId::toCapnp(rpc::TraceId::Builder writer) const {
228 writer.setLow(low);
229 writer.setHigh(high);
230}
231 
232kj::Maybe<InvocationSpanContext> InvocationSpanContext::fromCapnp(
233 rpc::InvocationSpanContext::Reader reader) {
234 if (!reader.hasTraceId() || !reader.hasInvocationId()) {
235 // If the reader does not have a traceId or invocationId field then it is
236 // invalid and we will just ignore it.
237 return kj::none;
238 }
239 
240 kj::Maybe<TraceFlags> flags;
241 if (reader.hasTraceFlags() && reader.getTraceFlags().getValue().isSet()) {
242 flags = TraceFlags(reader.getTraceFlags().getValue().getSet());
243 }
244 
245 auto sc = InvocationSpanContext(kj::Badge<InvocationSpanContext>(), kj::none,
246 TraceId::fromCapnp(reader.getTraceId()), TraceId::fromCapnp(reader.getInvocationId()),
247 reader.getSpanId(), kj::none, flags);
248 // If the traceId or invocationId are invalid, then we'll ignore them.
249 if (!sc.getTraceId() || !sc.getInvocationId()) return kj::none;
250 return kj::mv(sc);
251}
252 
253void InvocationSpanContext::toCapnp(rpc::InvocationSpanContext::Builder writer) const {
254 traceId.toCapnp(writer.initTraceId());
255 invocationId.toCapnp(writer.initInvocationId());
256 writer.setSpanId(spanId);
257 KJ_IF_SOME(flags, traceFlags) {
258 writer.initTraceFlags().getValue().setSet(flags);
259 }
260}
261 
262InvocationSpanContext InvocationSpanContext::clone() const {
263 kj::Maybe<kj::EntropySource&> otherEntropySource = entropySource.map(
264 [](auto& es) -> kj::EntropySource& { return const_cast<kj::EntropySource&>(es); });
265 return InvocationSpanContext(kj::Badge<InvocationSpanContext>(), otherEntropySource, traceId,
266 invocationId, spanId,
267 parentSpanContext.map([](auto& ctx) -> const InvocationSpanContext& { return *ctx.get(); }),
268 traceFlags);
269}
270 
271kj::String KJ_STRINGIFY(const InvocationSpanContext& context) {
272 return kj::str(context.getTraceId(), "-", context.getInvocationId(), "-", context.getSpanId());
273}
274 
275kj::String KJ_STRINGIFY(const TailEvent::Event& event) {
276 KJ_SWITCH_ONEOF(event) {
277 KJ_CASE_ONEOF(onset, Onset) {
278 return kj::str("Onset");
279 }
280 KJ_CASE_ONEOF(outcome, Outcome) {
281 return kj::str("Outcome");
282 }
283 KJ_CASE_ONEOF(spanOpen, SpanOpen) {
284 return spanOpen.toString();
285 }
286 KJ_CASE_ONEOF(spanClose, SpanClose) {
287 return spanClose.toString();
288 }
289 KJ_CASE_ONEOF(diagnosticChannelEvent, DiagnosticChannelEvent) {
290 return kj::str("diagnosticChannelEvent");
291 }
292 KJ_CASE_ONEOF(exception, Exception) {
293 return kj::str("Exception");
294 }
295 KJ_CASE_ONEOF(log, Log) {
296 return kj::str("Log");
297 }
298 KJ_CASE_ONEOF(streamDiag, StreamDiagnosticsEvent) {
299 return kj::str("StreamDiagnosticsEvent(droppedEvents: ", streamDiag.droppedEventsCount, ")");
300 }
301 KJ_CASE_ONEOF(ret, Return) {
302 return kj::str("Return");
303 }
304 KJ_CASE_ONEOF(customInfo, CustomInfo) {
305 return kj::str(customInfo);
306 }
307 }
308 KJ_UNREACHABLE
309}
310 
311kj::String KJ_STRINGIFY(const CustomInfo& customInfo) {
312 return kj::str(
313 "CustomInfo: ", kj::strArray(KJ_MAP(attr, customInfo) { return kj::str(attr); }, ", "));
314}
315 
316SpanContext SpanContext::fromCapnp(rpc::SpanContext::Reader reader) {
317 auto info = reader.getInfo();
318 kj::Maybe<SpanId> spanId;
319 if (info.isSpanId()) {
320 spanId = info.getSpanId();
321 }
322 
323 kj::Maybe<TraceFlags> flags;
324 if (reader.hasTraceFlags() && reader.getTraceFlags().getValue().isSet()) {
325 flags = TraceFlags(reader.getTraceFlags().getValue().getSet());
326 }
327 
328 return SpanContext(TraceId::fromCapnp(reader.getTraceId()), spanId, flags);
329}
330 
331void SpanContext::toCapnp(rpc::SpanContext::Builder writer) const {
332 traceId.toCapnp(writer.initTraceId());
333 auto info = writer.initInfo();
334 KJ_IF_SOME(s, spanId) {
335 info.setSpanId(s);
336 }
337 KJ_IF_SOME(flags, traceFlags) {
338 writer.initTraceFlags().getValue().setSet(flags);
339 }
340}
341 
342kj::Maybe<SpanContext> SpanContext::tryFromTraceparent(kj::StringPtr tp) {
343 // The W3C Trace Context traceparent header has a fixed-length format:
344 // {version:2}-{trace-id:32}-{parent-id:16}-{flags:2}
345 //
346 // See: https://www.w3.org/TR/trace-context/#traceparent-header
347 //
348 // The spec mandates fixed-width hex fields with fixed positions
349 constexpr size_t kVersionStart = 0;
350 constexpr size_t kVersionEnd = 2;
351 constexpr size_t kTraceIdStart = 3;
352 constexpr size_t kTraceIdMid = 19;
353 constexpr size_t kTraceIdEnd = 35;
354 constexpr size_t kParentIdStart = 36;
355 constexpr size_t kParentIdEnd = 52;
356 constexpr size_t kFlagsStart = 53;
357 constexpr size_t kFlagsEnd = 55;
358 
359 if (tp.size() != kFlagsEnd) return kj::none;
360 if (tp[kVersionEnd] != '-' || tp[kTraceIdEnd] != '-' || tp[kParentIdEnd] != '-') {
361 return kj::none;
362 }
363 
364 uint8_t version =
365 KJ_UNWRAP_OR_RETURN(hexToUint64(tp.slice(kVersionStart, kVersionEnd)), kj::none);
366 uint64_t traceHigh =
367 KJ_UNWRAP_OR_RETURN(hexToUint64(tp.slice(kTraceIdStart, kTraceIdMid)), kj::none);
368 uint64_t traceLow =
369 KJ_UNWRAP_OR_RETURN(hexToUint64(tp.slice(kTraceIdMid, kTraceIdEnd)), kj::none);
370 uint64_t parentId =
371 KJ_UNWRAP_OR_RETURN(hexToUint64(tp.slice(kParentIdStart, kParentIdEnd)), kj::none);
372 uint8_t flags = KJ_UNWRAP_OR_RETURN(hexToUint64(tp.slice(kFlagsStart, kFlagsEnd)), kj::none);
373 
374 if (version != 0 || (traceHigh == 0 && traceLow == 0) || parentId == 0) return kj::none;
375 
376 return SpanContext(TraceId(traceLow, traceHigh), SpanId(parentId), TraceFlags(flags));
377}
378 
379kj::String KJ_STRINGIFY(const SpanContext& context) {
380 return kj::str(context.getTraceId(), "-", context.getSpanId());
381}
382 
383namespace {
384 
385static kj::HttpMethod validateMethod(capnp::HttpMethod method) {
386 KJ_REQUIRE(method <= capnp::HttpMethod::BAN, "unknown method", method);
387 return static_cast<kj::HttpMethod>(method);
388}
389 
390} // namespace
391 
392ConnectEventInfo::ConnectEventInfo() {}
393 
394ConnectEventInfo::ConnectEventInfo(rpc::Trace::ConnectEventInfo::Reader reader) {}
395 
396void ConnectEventInfo::copyTo(rpc::Trace::ConnectEventInfo::Builder builder) const {}
397 
398ConnectEventInfo ConnectEventInfo::clone() const {
399 return ConnectEventInfo();
400}
401 
402FetchEventInfo::FetchEventInfo(
403 kj::HttpMethod method, kj::String url, kj::String cfJson, kj::Array<Header> headers)
404 : method(method),
405 url(kj::mv(url)),
406 cfJson(kj::mv(cfJson)),
407 headers(kj::mv(headers)) {}
408 
409FetchEventInfo::FetchEventInfo(rpc::Trace::FetchEventInfo::Reader reader)
410 : method(validateMethod(reader.getMethod())),
411 url(kj::str(reader.getUrl())),
412 cfJson(kj::str(reader.getCfJson())) {
413 kj::Vector<Header> v;
414 v.addAll(reader.getHeaders());
415 headers = v.releaseAsArray();
416}
417 
418void FetchEventInfo::copyTo(rpc::Trace::FetchEventInfo::Builder builder) const {
419 builder.setMethod(static_cast<capnp::HttpMethod>(method));
420 builder.setUrl(url);
421 builder.setCfJson(cfJson);
422 
423 auto list = builder.initHeaders(headers.size());
424 for (auto i: kj::indices(headers)) {
425 headers[i].copyTo(list[i]);
426 }
427}
428 
429FetchEventInfo FetchEventInfo::clone() const {
430 return FetchEventInfo(
431 method, kj::str(url), kj::str(cfJson), KJ_MAP(h, headers) { return h.clone(); });
432}
433 
434kj::String FetchEventInfo::toString() const {
435 return kj::str("FetchEventInfo: ",
436 kj::delimited(
437 kj::arr(kj::str(method), kj::str(url), kj::str(cfJson), kj::str(headers)), ", "_kjc));
438}
439 
440FetchEventInfo::Header::Header(kj::String name, kj::String value)
441 : name(kj::mv(name)),
442 value(kj::mv(value)) {}
443 
444FetchEventInfo::Header::Header(rpc::Trace::FetchEventInfo::Header::Reader reader)
445 : name(kj::str(reader.getName())),
446 value(kj::str(reader.getValue())) {}
447 
448void FetchEventInfo::Header::copyTo(rpc::Trace::FetchEventInfo::Header::Builder builder) const {
449 builder.setName(name);
450 builder.setValue(value);
451}
452 
453FetchEventInfo::Header FetchEventInfo::Header::clone() const {
454 return Header(kj::str(name), kj::str(value));
455}
456 
457kj::String FetchEventInfo::Header::toString() const {
458 return kj::str("FetchEventInfo::Header: ", name, ", ", value);
459}
460 
461JsRpcEventInfo::JsRpcEventInfo(kj::String methodName): methodName(kj::mv(methodName)) {}
462 
463JsRpcEventInfo::JsRpcEventInfo(rpc::Trace::JsRpcEventInfo::Reader reader)
464 : methodName(kj::str(reader.getMethodName())) {}
465 
466void JsRpcEventInfo::copyTo(rpc::Trace::JsRpcEventInfo::Builder builder) const {
467 builder.setMethodName(methodName);
468}
469 
470JsRpcEventInfo JsRpcEventInfo::clone() const {
471 return JsRpcEventInfo(kj::str(methodName));
472}
473 
474kj::String JsRpcEventInfo::toString() const {
475 return kj::str("JsRpcEventInfo: ", methodName);
476}
477 
478ScheduledEventInfo::ScheduledEventInfo(double scheduledTime, kj::String cron)
479 : scheduledTime(scheduledTime),
480 cron(kj::mv(cron)) {}
481 
482ScheduledEventInfo::ScheduledEventInfo(rpc::Trace::ScheduledEventInfo::Reader reader)
483 : scheduledTime(reader.getScheduledTime()),
484 cron(kj::str(reader.getCron())) {}
485 
486void ScheduledEventInfo::copyTo(rpc::Trace::ScheduledEventInfo::Builder builder) const {
487 builder.setScheduledTime(scheduledTime);
488 builder.setCron(cron);
489}
490 
491ScheduledEventInfo ScheduledEventInfo::clone() const {
492 return ScheduledEventInfo(scheduledTime, kj::str(cron));
493}
494 
495AlarmEventInfo::AlarmEventInfo(kj::Date scheduledTime): scheduledTime(scheduledTime) {}
496 
497AlarmEventInfo::AlarmEventInfo(rpc::Trace::AlarmEventInfo::Reader reader)
498 : scheduledTime(reader.getScheduledTimeMs() * kj::MILLISECONDS + kj::UNIX_EPOCH) {}
499 
500void AlarmEventInfo::copyTo(rpc::Trace::AlarmEventInfo::Builder builder) const {
501 builder.setScheduledTimeMs((scheduledTime - kj::UNIX_EPOCH) / kj::MILLISECONDS);
502}
503 
504AlarmEventInfo AlarmEventInfo::clone() const {
505 return AlarmEventInfo(scheduledTime);
506}
507 
508QueueEventInfo::QueueEventInfo(kj::String queueName, uint32_t batchSize)
509 : queueName(kj::mv(queueName)),
510 batchSize(batchSize) {}
511 
512QueueEventInfo::QueueEventInfo(rpc::Trace::QueueEventInfo::Reader reader)
513 : queueName(kj::heapString(reader.getQueueName())),
514 batchSize(reader.getBatchSize()) {}
515 
516void QueueEventInfo::copyTo(rpc::Trace::QueueEventInfo::Builder builder) const {
517 builder.setQueueName(queueName);
518 builder.setBatchSize(batchSize);
519}
520 
521QueueEventInfo QueueEventInfo::clone() const {
522 return QueueEventInfo(kj::str(queueName), batchSize);
523}
524 
525EmailEventInfo::EmailEventInfo(kj::String mailFrom, kj::String rcptTo, uint32_t rawSize)
526 : mailFrom(kj::mv(mailFrom)),
527 rcptTo(kj::mv(rcptTo)),
528 rawSize(rawSize) {}
529 
530EmailEventInfo::EmailEventInfo(rpc::Trace::EmailEventInfo::Reader reader)
531 : mailFrom(kj::heapString(reader.getMailFrom())),
532 rcptTo(kj::heapString(reader.getRcptTo())),
533 rawSize(reader.getRawSize()) {}
534 
535void EmailEventInfo::copyTo(rpc::Trace::EmailEventInfo::Builder builder) const {
536 builder.setMailFrom(mailFrom);
537 builder.setRcptTo(rcptTo);
538 builder.setRawSize(rawSize);
539}
540 
541EmailEventInfo EmailEventInfo::clone() const {
542 return EmailEventInfo(kj::str(mailFrom), kj::str(rcptTo), rawSize);
543}
544 
545namespace {
546kj::Vector<TraceEventInfo::TraceItem> getTraceItemsFromTraces(
547 kj::ArrayPtr<const kj::Own<Trace>> traces) {
548 return KJ_MAP(t, traces) { return TraceEventInfo::TraceItem(mapCopyString(t->scriptName)); };
549}
550 
551kj::Vector<TraceEventInfo::TraceItem> getTraceItemsFromReader(
552 rpc::Trace::TraceEventInfo::Reader reader) {
553 return KJ_MAP(r, reader.getTraces()) { return TraceEventInfo::TraceItem(r); };
554}
555} // namespace
556 
557TraceEventInfo::TraceEventInfo(kj::ArrayPtr<const kj::Own<Trace>> traces)
558 : traces(getTraceItemsFromTraces(traces)) {}
559 
560TraceEventInfo::TraceEventInfo(rpc::Trace::TraceEventInfo::Reader reader)
561 : traces(getTraceItemsFromReader(reader)) {}
562 
563void TraceEventInfo::copyTo(rpc::Trace::TraceEventInfo::Builder builder) const {
564 auto list = builder.initTraces(traces.size());
565 for (auto i: kj::indices(traces)) {
566 traces[i].copyTo(list[i]);
567 }
568}
569 
570TraceEventInfo TraceEventInfo::clone() const {
571 return TraceEventInfo(KJ_MAP(item, traces) { return item.clone(); });
572}
573 
574TracePreview::TracePreview(kj::String id, kj::String slug, kj::String name)
575 : id(kj::mv(id)),
576 slug(kj::mv(slug)),
577 name(kj::mv(name)) {}
578 
579TracePreview::TracePreview(rpc::Trace::TracePreviewInfo::Reader reader)
580 : id(kj::str(reader.getId())),
581 slug(kj::str(reader.getSlug())),
582 name(kj::str(reader.getName())) {}
583 
584void TracePreview::copyTo(rpc::Trace::TracePreviewInfo::Builder builder) const {
585 builder.setId(id);
586 builder.setSlug(slug);
587 builder.setName(name);
588}
589 
590TracePreview TracePreview::clone() const {
591 return TracePreview(kj::str(id), kj::str(slug), kj::str(name));
592}
593 
594TraceEventInfo::TraceItem::TraceItem(kj::Maybe<kj::String> scriptName)
595 : scriptName(kj::mv(scriptName)) {}
596 
597TraceEventInfo::TraceItem::TraceItem(rpc::Trace::TraceEventInfo::TraceItem::Reader reader)
598 : scriptName(kj::str(reader.getScriptName())) {}
599 
600void TraceEventInfo::TraceItem::copyTo(
601 rpc::Trace::TraceEventInfo::TraceItem::Builder builder) const {
602 KJ_IF_SOME(name, scriptName) {
603 builder.setScriptName(name);
604 }
605}
606 
607TraceEventInfo::TraceItem TraceEventInfo::TraceItem::clone() const {
608 return TraceItem(mapCopyString(scriptName));
609}
610 
611DiagnosticChannelEvent::DiagnosticChannelEvent(
612 kj::Date timestamp, kj::String channel, kj::Array<kj::byte> message)
613 : timestamp(timestamp),
614 channel(kj::mv(channel)),
615 message(kj::mv(message)) {}
616 
617DiagnosticChannelEvent::DiagnosticChannelEvent(rpc::Trace::DiagnosticChannelEvent::Reader reader)
618 : timestamp(kj::UNIX_EPOCH + reader.getTimestampNs() * kj::NANOSECONDS),
619 channel(kj::heapString(reader.getChannel())),
620 message(kj::heapArray<kj::byte>(reader.getMessage())) {}
621 
622void DiagnosticChannelEvent::copyTo(rpc::Trace::DiagnosticChannelEvent::Builder builder) const {
623 builder.setTimestampNs((timestamp - kj::UNIX_EPOCH) / kj::NANOSECONDS);
624 builder.setChannel(channel);
625 builder.setMessage(message);
626}
627 
628DiagnosticChannelEvent DiagnosticChannelEvent::clone() const {
629 return DiagnosticChannelEvent(timestamp, kj::str(channel), kj::heapArray<kj::byte>(message));
630}
631 
632StreamDiagnosticsEvent::StreamDiagnosticsEvent(uint32_t droppedEventsCount)
633 : droppedEventsCount(droppedEventsCount) {}
634 
635StreamDiagnosticsEvent::StreamDiagnosticsEvent(rpc::Trace::StreamDiagnosticsEvent::Reader reader) {
636 auto diagnosticReader = reader.getDiagnostic();
637 switch (diagnosticReader.which()) {
638 case rpc::Trace::StreamDiagnosticsEvent::Diagnostic::UNDEFINED:
639 KJ_FAIL_ASSERT("received invalid diagnostics event");
640 break;
641 case rpc::Trace::StreamDiagnosticsEvent::Diagnostic::DROPPED_EVENTS:
642 auto droppedEvents = diagnosticReader.getDroppedEvents();
643 droppedEventsCount = droppedEvents.getCount();
644 KJ_DASSERT(droppedEventsCount > 0);
645 break;
646 }
647}
648 
649void StreamDiagnosticsEvent::copyTo(rpc::Trace::StreamDiagnosticsEvent::Builder builder) const {
650 KJ_DASSERT(droppedEventsCount > 0);
651 auto diagnosticBuilder = builder.initDiagnostic();
652 auto droppedEventsBuilder = diagnosticBuilder.initDroppedEvents();
653 droppedEventsBuilder.setCount(droppedEventsCount);
654}
655 
656StreamDiagnosticsEvent StreamDiagnosticsEvent::clone() const {
657 return StreamDiagnosticsEvent(droppedEventsCount);
658}
659 
660HibernatableWebSocketEventInfo::HibernatableWebSocketEventInfo(Type type): type(type) {}
661 
662HibernatableWebSocketEventInfo::HibernatableWebSocketEventInfo(
663 rpc::Trace::HibernatableWebSocketEventInfo::Reader reader)
664 : type(readFrom(reader)) {}
665 
666void HibernatableWebSocketEventInfo::copyTo(
667 rpc::Trace::HibernatableWebSocketEventInfo::Builder builder) const {
668 auto typeBuilder = builder.initType();
669 KJ_SWITCH_ONEOF(type) {
670 KJ_CASE_ONEOF(_, Message) {
671 typeBuilder.setMessage();
672 }
673 KJ_CASE_ONEOF(close, Close) {
674 auto closeBuilder = typeBuilder.initClose();
675 closeBuilder.setCode(close.code);
676 closeBuilder.setWasClean(close.wasClean);
677 }
678 KJ_CASE_ONEOF(_, Error) {
679 typeBuilder.setError();
680 }
681 }
682}
683 
684HibernatableWebSocketEventInfo HibernatableWebSocketEventInfo::clone() const {
685 KJ_SWITCH_ONEOF(type) {
686 KJ_CASE_ONEOF(_, Message) {
687 return HibernatableWebSocketEventInfo(Message{});
688 }
689 KJ_CASE_ONEOF(_, Error) {
690 return HibernatableWebSocketEventInfo(Error{});
691 }
692 KJ_CASE_ONEOF(close, Close) {
693 return HibernatableWebSocketEventInfo(Close{
694 .code = close.code,
695 .wasClean = close.wasClean,
696 });
697 }
698 }
699 KJ_UNREACHABLE;
700}
701 
702HibernatableWebSocketEventInfo::Type HibernatableWebSocketEventInfo::readFrom(
703 rpc::Trace::HibernatableWebSocketEventInfo::Reader reader) {
704 auto type = reader.getType();
705 switch (type.which()) {
706 case rpc::Trace::HibernatableWebSocketEventInfo::Type::MESSAGE: {
707 return Message{};
708 }
709 case rpc::Trace::HibernatableWebSocketEventInfo::Type::CLOSE: {
710 auto close = type.getClose();
711 return Close{
712 .code = close.getCode(),
713 .wasClean = close.getWasClean(),
714 };
715 }
716 case rpc::Trace::HibernatableWebSocketEventInfo::Type::ERROR: {
717 return Error{};
718 }
719 }
720}
721 
722FetchResponseInfo::FetchResponseInfo(uint16_t statusCode): statusCode(statusCode) {}
723 
724FetchResponseInfo::FetchResponseInfo(rpc::Trace::FetchResponseInfo::Reader reader)
725 : statusCode(reader.getStatusCode()) {}
726 
727void FetchResponseInfo::copyTo(rpc::Trace::FetchResponseInfo::Builder builder) const {
728 builder.setStatusCode(statusCode);
729}
730 
731FetchResponseInfo FetchResponseInfo::clone() const {
732 return FetchResponseInfo(statusCode);
733}
734 
735Log::Log(kj::Date timestamp, LogLevel logLevel, kj::String message)
736 : timestamp(timestamp),
737 logLevel(logLevel),
738 message(kj::mv(message)) {}
739 
740void Log::copyTo(rpc::Trace::Log::Builder builder) const {
741 builder.setTimestampNs((timestamp - kj::UNIX_EPOCH) / kj::NANOSECONDS);
742 builder.setLogLevel(logLevel);
743 builder.setMessage(message);
744}
745 
746Log Log::clone() const {
747 return Log(timestamp, logLevel, kj::str(message));
748}
749 
750Exception::Exception(
751 kj::Date timestamp, kj::String name, kj::String message, kj::Maybe<kj::String> stack)
752 : timestamp(timestamp),
753 name(kj::mv(name)),
754 message(kj::mv(message)),
755 stack(kj::mv(stack)) {}
756 
757Log::Log(rpc::Trace::Log::Reader reader)
758 : timestamp(kj::UNIX_EPOCH + reader.getTimestampNs() * kj::NANOSECONDS),
759 logLevel(reader.getLogLevel()),
760 message(kj::str(reader.getMessage())) {}
761 
762Exception::Exception(rpc::Trace::Exception::Reader reader)
763 : timestamp(kj::UNIX_EPOCH + reader.getTimestampNs() * kj::NANOSECONDS),
764 name(kj::str(reader.getName())),
765 message(kj::str(reader.getMessage())) {
766 if (reader.hasStack()) {
767 stack = kj::str(reader.getStack());
768 }
769}
770 
771void Exception::copyTo(rpc::Trace::Exception::Builder builder) const {
772 builder.setTimestampNs((timestamp - kj::UNIX_EPOCH) / kj::NANOSECONDS);
773 builder.setName(name);
774 builder.setMessage(message);
775 KJ_IF_SOME(s, stack) {
776 builder.setStack(s);
777 }
778}
779 
780Exception Exception::clone() const {
781 return Exception(timestamp, kj::str(name), kj::str(message), mapCopyString(stack));
782}
783} // namespace tracing
784 
785Trace::Trace(kj::Maybe<kj::String> stableId,
786 kj::Maybe<kj::String> scriptName,
787 kj::Maybe<kj::Own<ScriptVersion::Reader>> scriptVersion,
788 kj::Maybe<kj::String> dispatchNamespace,
789 kj::Maybe<kj::String> scriptId,
790 kj::Array<kj::String> scriptTags,
791 kj::Maybe<kj::String> entrypoint,
792 ExecutionModel executionModel,
793 kj::Maybe<kj::String> durableObjectId,
794 kj::Maybe<tracing::TracePreview> preview)
795 : stableId(kj::mv(stableId)),
796 scriptName(kj::mv(scriptName)),
797 scriptVersion(kj::mv(scriptVersion)),
798 dispatchNamespace(kj::mv(dispatchNamespace)),
799 scriptId(kj::mv(scriptId)),
800 scriptTags(kj::mv(scriptTags)),
801 entrypoint(kj::mv(entrypoint)),
802 preview(kj::mv(preview)),
803 durableObjectId(kj::mv(durableObjectId)),
804 executionModel(executionModel) {}
805Trace::Trace(rpc::Trace::Reader reader) {
806 mergeFrom(reader, PipelineLogLevel::FULL);
807}
808 
809Trace::~Trace() noexcept(false) {}
810 
811void Trace::copyTo(rpc::Trace::Builder builder) const {
812 {
813 auto list = builder.initLogs(logs.size());
814 for (auto i: kj::indices(logs)) {
815 logs[i].copyTo(list[i]);
816 }
817 }
818 
819 {
820 auto list = builder.initExceptions(exceptions.size());
821 for (auto i: kj::indices(exceptions)) {
822 exceptions[i].copyTo(list[i]);
823 }
824 }
825 
826 builder.setTruncated(truncated);
827 builder.setOutcome(outcome);
828 builder.setCpuTime(cpuTime / kj::MILLISECONDS);
829 builder.setWallTime(wallTime / kj::MILLISECONDS);
830 KJ_IF_SOME(name, scriptName) {
831 builder.setScriptName(name);
832 }
833 KJ_IF_SOME(version, scriptVersion) {
834 builder.setScriptVersion(*version);
835 }
836 KJ_IF_SOME(id, scriptId) {
837 builder.setScriptId(id);
838 }
839 KJ_IF_SOME(ns, dispatchNamespace) {
840 builder.setDispatchNamespace(ns);
841 }
842 builder.setExecutionModel(executionModel);
843 
844 {
845 auto list = builder.initScriptTags(scriptTags.size());
846 for (auto i: kj::indices(scriptTags)) {
847 list.set(i, scriptTags[i]);
848 }
849 }
850 
851 KJ_IF_SOME(tags, tailAttributes) {
852 auto list = builder.initTailAttributes(tags.size());
853 for (auto i: kj::indices(tags)) {
854 tags[i].copyTo(list[i]);
855 }
856 }
857 
858 KJ_IF_SOME(e, entrypoint) {
859 builder.setEntrypoint(e);
860 }
861 
862 KJ_IF_SOME(p, preview) {
863 p.copyTo(builder.initPreview());
864 }
865 
866 KJ_IF_SOME(id, durableObjectId) {
867 builder.setDurableObjectId(id);
868 }
869 
870 builder.setEventTimestampNs((eventTimestamp - kj::UNIX_EPOCH) / kj::NANOSECONDS);
871 
872 auto eventInfoBuilder = builder.initEventInfo();
873 KJ_IF_SOME(e, eventInfo) {
874 KJ_SWITCH_ONEOF(e) {
875 KJ_CASE_ONEOF(fetch, tracing::FetchEventInfo) {
876 auto fetchBuilder = eventInfoBuilder.initFetch();
877 fetch.copyTo(fetchBuilder);
878 }
879 KJ_CASE_ONEOF(jsRpc, tracing::JsRpcEventInfo) {
880 auto jsRpcBuilder = eventInfoBuilder.initJsRpc();
881 jsRpc.copyTo(jsRpcBuilder);
882 }
883 KJ_CASE_ONEOF(connect, tracing::ConnectEventInfo) {
884 auto connectBuilder = eventInfoBuilder.initConnect();
885 connect.copyTo(connectBuilder);
886 }
887 KJ_CASE_ONEOF(scheduled, tracing::ScheduledEventInfo) {
888 auto scheduledBuilder = eventInfoBuilder.initScheduled();
889 scheduled.copyTo(scheduledBuilder);
890 }
891 KJ_CASE_ONEOF(alarm, tracing::AlarmEventInfo) {
892 auto alarmBuilder = eventInfoBuilder.initAlarm();
893 alarm.copyTo(alarmBuilder);
894 }
895 KJ_CASE_ONEOF(queue, tracing::QueueEventInfo) {
896 auto queueBuilder = eventInfoBuilder.initQueue();
897 queue.copyTo(queueBuilder);
898 }
899 KJ_CASE_ONEOF(email, tracing::EmailEventInfo) {
900 auto emailBuilder = eventInfoBuilder.initEmail();
901 email.copyTo(emailBuilder);
902 }
903 KJ_CASE_ONEOF(trace, tracing::TraceEventInfo) {
904 auto traceBuilder = eventInfoBuilder.initTrace();
905 trace.copyTo(traceBuilder);
906 }
907 KJ_CASE_ONEOF(hibWs, tracing::HibernatableWebSocketEventInfo) {
908 auto hibWsBuilder = eventInfoBuilder.initHibernatableWebSocket();
909 hibWs.copyTo(hibWsBuilder);
910 }
911 KJ_CASE_ONEOF(custom, tracing::CustomEventInfo) {
912 eventInfoBuilder.initCustom();
913 }
914 }
915 } else {
916 eventInfoBuilder.setNone();
917 }
918 
919 KJ_IF_SOME(fetchResponseInfo, this->fetchResponseInfo) {
920 auto fetchResponseInfoBuilder = builder.initResponse();
921 fetchResponseInfo.copyTo(fetchResponseInfoBuilder);
922 }
923 
924 {
925 auto list = builder.initDiagnosticChannelEvents(diagnosticChannelEvents.size());
926 for (auto i: kj::indices(diagnosticChannelEvents)) {
927 diagnosticChannelEvents[i].copyTo(list[i]);
928 }
929 }
930}
931 
932void Trace::mergeFrom(rpc::Trace::Reader reader, PipelineLogLevel pipelineLogLevel) {
933 // Sandboxed workers currently record their traces as if the pipeline log level were set to
934 // "full", so we may need to filter out the extra data after receiving the traces back.
935 if (pipelineLogLevel != PipelineLogLevel::NONE) {
936 logs.addAll(reader.getLogs());
937 exceptions.addAll(reader.getExceptions());
938 diagnosticChannelEvents.addAll(reader.getDiagnosticChannelEvents());
939 }
940 
941 truncated = reader.getTruncated();
942 outcome = reader.getOutcome();
943 cpuTime = reader.getCpuTime() * kj::MILLISECONDS;
944 wallTime = reader.getWallTime() * kj::MILLISECONDS;
945 
946 // mergeFrom() is called both when deserializing traces from a sandboxed
947 // worker and when deserializing traces sent to a sandboxed trace worker. In
948 // the former case, the trace's scriptName (and other fields like
949 // scriptVersion) are already set and the deserialized value is missing, so
950 // we need to be careful not to overwrite the set value.
951 if (reader.hasScriptName()) {
952 scriptName = kj::str(reader.getScriptName());
953 }
954 
955 if (reader.hasScriptVersion()) {
956 scriptVersion = capnp::clone(reader.getScriptVersion());
957 }
958 
959 if (reader.hasScriptId()) {
960 scriptId = kj::str(reader.getScriptId());
961 }
962 
963 if (reader.hasDispatchNamespace()) {
964 dispatchNamespace = kj::str(reader.getDispatchNamespace());
965 }
966 executionModel = reader.getExecutionModel();
967 
968 if (auto tags = reader.getScriptTags(); tags.size() > 0) {
969 scriptTags = KJ_MAP(tag, tags) { return kj::str(tag); };
970 }
971 
972 if (auto tags = reader.getTailAttributes(); tags.size() > 0) {
973 tailAttributes = KJ_MAP(tag, tags) { return tracing::Attribute(tag); };
974 }
975 
976 if (reader.hasEntrypoint()) {
977 entrypoint = kj::str(reader.getEntrypoint());
978 }
979 
980 if (reader.hasPreview()) {
981 preview = tracing::TracePreview(reader.getPreview());
982 }
983 
984 if (reader.hasDurableObjectId()) {
985 durableObjectId = kj::str(reader.getDurableObjectId());
986 }
987 
988 eventTimestamp = kj::UNIX_EPOCH + reader.getEventTimestampNs() * kj::NANOSECONDS;
989 
990 if (pipelineLogLevel == PipelineLogLevel::NONE) {
991 eventInfo = kj::none;
992 } else {
993 auto e = reader.getEventInfo();
994 switch (e.which()) {
995 case rpc::Trace::EventInfo::Which::FETCH:
996 eventInfo = tracing::FetchEventInfo(e.getFetch());
997 break;
998 case rpc::Trace::EventInfo::Which::JS_RPC:
999 eventInfo = tracing::JsRpcEventInfo(e.getJsRpc());
1000 break;
1001 case rpc::Trace::EventInfo::Which::CONNECT:
1002 eventInfo = tracing::ConnectEventInfo(e.getConnect());
1003 break;
1004 case rpc::Trace::EventInfo::Which::SCHEDULED:
1005 eventInfo = tracing::ScheduledEventInfo(e.getScheduled());
1006 break;
1007 case rpc::Trace::EventInfo::Which::ALARM:
1008 eventInfo = tracing::AlarmEventInfo(e.getAlarm());
1009 break;
1010 case rpc::Trace::EventInfo::Which::QUEUE:
1011 eventInfo = tracing::QueueEventInfo(e.getQueue());
1012 break;
1013 case rpc::Trace::EventInfo::Which::EMAIL:
1014 eventInfo = tracing::EmailEventInfo(e.getEmail());
1015 break;
1016 case rpc::Trace::EventInfo::Which::TRACE:
1017 eventInfo = tracing::TraceEventInfo(e.getTrace());
1018 break;
1019 case rpc::Trace::EventInfo::Which::HIBERNATABLE_WEB_SOCKET:
1020 eventInfo = tracing::HibernatableWebSocketEventInfo(e.getHibernatableWebSocket());
1021 break;
1022 case rpc::Trace::EventInfo::Which::CUSTOM:
1023 eventInfo = tracing::CustomEventInfo(e.getCustom());
1024 break;
1025 case rpc::Trace::EventInfo::Which::NONE:
1026 eventInfo = kj::none;
1027 break;
1028 }
1029 }
1030 
1031 if (reader.hasResponse()) {
1032 fetchResponseInfo = tracing::FetchResponseInfo(reader.getResponse());
1033 }
1034}
1035 
1036namespace tracing {
1037 
1038Attribute::Attribute(kj::ConstString name, Value&& value)
1039 : name(kj::mv(name)),
1040 value(kj::arr(kj::mv(value))) {}
1041 
1042Attribute::Attribute(kj::ConstString name, Values&& value)
1043 : name(kj::mv(name)),
1044 value(kj::mv(value)) {}
1045 
1046namespace {
1047kj::Array<Attribute::Value> readValues(const rpc::Trace::Attribute::Reader& reader) {
1048 // There should always be a value and it always have at least one entry in the list.
1049 KJ_ASSERT(reader.hasValue());
1050 auto value = reader.getValue();
1051 return KJ_MAP(v, value) { return deserializeTagValue(v); };
1052}
1053 
1054kj::Maybe<FetchResponseInfo> readReturnInfo(const rpc::Trace::Return::Reader& reader) {
1055 auto info = reader.getInfo();
1056 switch (info.which()) {
1057 case rpc::Trace::Return::Info::EMPTY:
1058 return kj::none;
1059 case rpc::Trace::Return::Info::FETCH: {
1060 return kj::Maybe(FetchResponseInfo(info.getFetch()));
1061 }
1062 }
1063 KJ_UNREACHABLE;
1064}
1065} // namespace
1066 
1067Attribute::Attribute(rpc::Trace::Attribute::Reader reader)
1068 : name(kj::str(reader.getName())),
1069 value(readValues(reader)) {}
1070 
1071void Attribute::copyTo(rpc::Trace::Attribute::Builder builder) const {
1072 builder.setName(name.asPtr());
1073 auto vec = builder.initValue(value.size());
1074 for (size_t n = 0; n < value.size(); n++) {
1075 serializeTagValue(vec[n], value[n]);
1076 }
1077}
1078 
1079Attribute Attribute::clone() const {
1080 return Attribute(name.clone(), KJ_MAP(v, value) { return spanTagClone(v); });
1081}
1082 
1083kj::String Attribute::toString() const {
1084 return kj::str("Attribute: ", name, ", ", value);
1085}
1086 
1087Return::Return(kj::Maybe<FetchResponseInfo> info): info(kj::mv(info)) {}
1088Return::Return(rpc::Trace::Return::Reader reader): info(readReturnInfo(reader)) {}
1089 
1090void Return::copyTo(rpc::Trace::Return::Builder builder) const {
1091 KJ_IF_SOME(fetchInfo, info) {
1092 auto infoBuilder = builder.initInfo();
1093 fetchInfo.copyTo(infoBuilder.initFetch());
1094 }
1095}
1096 
1097Return Return::clone() const {
1098 KJ_IF_SOME(fetchInfo, info) {
1099 return Return(kj::Maybe(fetchInfo.clone()));
1100 }
1101 return Return();
1102}
1103 
1104SpanOpen::SpanOpen(SpanId spanId, kj::ConstString operationName, kj::Maybe<Info> info)
1105 : operationName(kj::mv(operationName)),
1106 info(kj::mv(info)),
1107 spanId(spanId) {}
1108 
1109namespace {
1110kj::Maybe<SpanOpen::Info> readSpanOpenInfo(rpc::Trace::SpanOpen::Reader& reader) {
1111 auto info = reader.getInfo();
1112 switch (info.which()) {
1113 case rpc::Trace::SpanOpen::Info::EMPTY:
1114 return kj::none;
1115 case rpc::Trace::SpanOpen::Info::FETCH: {
1116 return kj::Maybe(FetchEventInfo(info.getFetch()));
1117 }
1118 case rpc::Trace::SpanOpen::Info::JS_RPC: {
1119 return kj::Maybe(JsRpcEventInfo(info.getJsRpc()));
1120 }
1121 case rpc::Trace::SpanOpen::Info::CUSTOM: {
1122 auto custom = info.getCustom();
1123 return kj::Maybe(KJ_MAP(a, custom) { return Attribute(a); });
1124 }
1125 }
1126 KJ_UNREACHABLE;
1127}
1128} // namespace
1129 
1130SpanOpen::SpanOpen(rpc::Trace::SpanOpen::Reader reader)
1131 : operationName(kj::str(reader.getOperationName())),
1132 info(readSpanOpenInfo(reader)),
1133 spanId(reader.getSpanId()) {}
1134 
1135void SpanOpen::copyTo(rpc::Trace::SpanOpen::Builder builder) const {
1136 builder.setOperationName(operationName.asPtr());
1137 builder.setSpanId(spanId);
1138 KJ_IF_SOME(i, info) {
1139 auto infoBuilder = builder.initInfo();
1140 KJ_SWITCH_ONEOF(i) {
1141 KJ_CASE_ONEOF(fetch, FetchEventInfo) {
1142 fetch.copyTo(infoBuilder.initFetch());
1143 }
1144 KJ_CASE_ONEOF(jsrpc, JsRpcEventInfo) {
1145 jsrpc.copyTo(infoBuilder.initJsRpc());
1146 }
1147 KJ_CASE_ONEOF(custom, CustomInfo) {
1148 auto customBuilder = infoBuilder.initCustom(custom.size());
1149 for (size_t n = 0; n < custom.size(); n++) {
1150 custom[n].copyTo(customBuilder[n]);
1151 }
1152 }
1153 }
1154 }
1155}
1156 
1157SpanOpen SpanOpen::clone() const {
1158 constexpr auto cloneInfo = [](const kj::Maybe<Info>& info) -> kj::Maybe<SpanOpen::Info> {
1159 return info.map([](const Info& info) -> SpanOpen::Info {
1160 KJ_SWITCH_ONEOF(info) {
1161 KJ_CASE_ONEOF(fetch, FetchEventInfo) {
1162 return fetch.clone();
1163 }
1164 KJ_CASE_ONEOF(jsrpc, JsRpcEventInfo) {
1165 return jsrpc.clone();
1166 }
1167 KJ_CASE_ONEOF(custom, CustomInfo) {
1168 return KJ_MAP(attr, custom) { return attr.clone(); };
1169 }
1170 }
1171 KJ_UNREACHABLE;
1172 });
1173 };
1174 return SpanOpen(spanId, operationName.clone(), cloneInfo(info));
1175}
1176 
1177kj::String KJ_STRINGIFY(const SpanOpen::Info& info) {
1178 KJ_SWITCH_ONEOF(info) {
1179 KJ_CASE_ONEOF(fetch, FetchEventInfo) {
1180 return fetch.toString();
1181 }
1182 KJ_CASE_ONEOF(jsrpc, JsRpcEventInfo) {
1183 return jsrpc.toString();
1184 }
1185 KJ_CASE_ONEOF(customInfo, CustomInfo) {
1186 return kj::str(customInfo);
1187 }
1188 }
1189 KJ_UNREACHABLE
1190}
1191 
1192kj::String SpanOpen::toString() const {
1193 return kj::str("SpanOpen:", operationName, ", ", info);
1194}
1195 
1196SpanClose::SpanClose(EventOutcome outcome): outcome(outcome) {}
1197 
1198SpanClose::SpanClose(rpc::Trace::SpanClose::Reader reader): outcome(reader.getOutcome()) {}
1199 
1200void SpanClose::copyTo(rpc::Trace::SpanClose::Builder builder) const {
1201 builder.setOutcome(outcome);
1202}
1203 
1204SpanClose SpanClose::clone() const {
1205 return SpanClose(outcome);
1206}
1207 
1208kj::String SpanClose::toString() const {
1209 return kj::str("SpanClose: ", outcome);
1210}
1211 
1212Onset::Info readOnsetInfo(const rpc::Trace::Onset::Info::Reader& info) {
1213 switch (info.which()) {
1214 case rpc::Trace::Onset::Info::FETCH: {
1215 return FetchEventInfo(info.getFetch());
1216 }
1217 case rpc::Trace::Onset::Info::JS_RPC: {
1218 return JsRpcEventInfo(info.getJsRpc());
1219 }
1220 case rpc::Trace::Onset::Info::CONNECT: {
1221 return ConnectEventInfo(info.getConnect());
1222 }
1223 case rpc::Trace::Onset::Info::SCHEDULED: {
1224 return ScheduledEventInfo(info.getScheduled());
1225 }
1226 case rpc::Trace::Onset::Info::ALARM: {
1227 return AlarmEventInfo(info.getAlarm());
1228 }
1229 case rpc::Trace::Onset::Info::QUEUE: {
1230 return QueueEventInfo(info.getQueue());
1231 }
1232 case rpc::Trace::Onset::Info::EMAIL: {
1233 return EmailEventInfo(info.getEmail());
1234 }
1235 case rpc::Trace::Onset::Info::TRACE: {
1236 return TraceEventInfo(info.getTrace());
1237 }
1238 case rpc::Trace::Onset::Info::HIBERNATABLE_WEB_SOCKET: {
1239 return HibernatableWebSocketEventInfo(info.getHibernatableWebSocket());
1240 }
1241 case rpc::Trace::Onset::Info::CUSTOM: {
1242 return CustomEventInfo();
1243 }
1244 }
1245 KJ_UNREACHABLE;
1246}
1247 
1248void writeOnsetInfo(const Onset::Info& info, rpc::Trace::Onset::Info::Builder& infoBuilder) {
1249 KJ_SWITCH_ONEOF(info) {
1250 KJ_CASE_ONEOF(fetch, FetchEventInfo) {
1251 fetch.copyTo(infoBuilder.initFetch());
1252 }
1253 KJ_CASE_ONEOF(connect, ConnectEventInfo) {
1254 connect.copyTo(infoBuilder.initConnect());
1255 }
1256 KJ_CASE_ONEOF(jsrpc, JsRpcEventInfo) {
1257 jsrpc.copyTo(infoBuilder.initJsRpc());
1258 }
1259 KJ_CASE_ONEOF(scheduled, ScheduledEventInfo) {
1260 scheduled.copyTo(infoBuilder.initScheduled());
1261 }
1262 KJ_CASE_ONEOF(alarm, AlarmEventInfo) {
1263 alarm.copyTo(infoBuilder.initAlarm());
1264 }
1265 KJ_CASE_ONEOF(queue, QueueEventInfo) {
1266 queue.copyTo(infoBuilder.initQueue());
1267 }
1268 KJ_CASE_ONEOF(email, EmailEventInfo) {
1269 email.copyTo(infoBuilder.initEmail());
1270 }
1271 KJ_CASE_ONEOF(trace, TraceEventInfo) {
1272 trace.copyTo(infoBuilder.initTrace());
1273 }
1274 KJ_CASE_ONEOF(hws, HibernatableWebSocketEventInfo) {
1275 hws.copyTo(infoBuilder.initHibernatableWebSocket());
1276 }
1277 KJ_CASE_ONEOF(custom, CustomEventInfo) {
1278 infoBuilder.initCustom();
1279 }
1280 }
1281}
1282 
1283namespace {
1284kj::Maybe<kj::String> getScriptNameFromReader(const rpc::Trace::Onset::Reader& reader) {
1285 if (reader.hasScriptName()) {
1286 return kj::str(reader.getScriptName());
1287 }
1288 return kj::none;
1289}
1290 
1291kj::Maybe<kj::Own<ScriptVersion::Reader>> getScriptVersionFromReader(
1292 const rpc::Trace::Onset::Reader& reader) {
1293 if (reader.hasScriptVersion()) {
1294 return capnp::clone(reader.getScriptVersion());
1295 }
1296 return kj::none;
1297}
1298 
1299kj::Maybe<kj::String> getDispatchNamespaceFromReader(const rpc::Trace::Onset::Reader& reader) {
1300 if (reader.hasDispatchNamespace()) {
1301 return kj::str(reader.getDispatchNamespace());
1302 }
1303 return kj::none;
1304}
1305 
1306kj::Maybe<kj::String> getScriptIdFromReader(const rpc::Trace::Onset::Reader& reader) {
1307 if (reader.hasScriptId()) {
1308 return kj::str(reader.getScriptId());
1309 }
1310 return kj::none;
1311}
1312 
1313kj::Maybe<kj::Array<kj::String>> getScriptTagsFromReader(const rpc::Trace::Onset::Reader& reader) {
1314 if (reader.hasScriptTags()) {
1315 auto tags = reader.getScriptTags();
1316 kj::Vector<kj::String> scriptTags(tags.size());
1317 for (const auto& tag: tags) {
1318 scriptTags.add(kj::str(tag));
1319 }
1320 return kj::Maybe(scriptTags.releaseAsArray());
1321 }
1322 return kj::none;
1323}
1324 
1325kj::Maybe<kj::String> getEntrypointFromReader(const rpc::Trace::Onset::Reader& reader) {
1326 if (reader.hasEntryPoint()) {
1327 return kj::str(reader.getEntryPoint());
1328 }
1329 return kj::none;
1330}
1331 
1332kj::Maybe<tracing::TracePreview> getPreviewFromReader(const rpc::Trace::Onset::Reader& reader) {
1333 if (reader.hasPreview()) {
1334 return tracing::TracePreview(reader.getPreview());
1335 }
1336 return kj::none;
1337}
1338 
1339Onset::WorkerInfo getWorkerInfoFromReader(const rpc::Trace::Onset::Reader& reader) {
1340 return Onset::WorkerInfo{
1341 .executionModel = reader.getExecutionModel(),
1342 .scriptName = getScriptNameFromReader(reader),
1343 .scriptVersion = getScriptVersionFromReader(reader),
1344 .preview = getPreviewFromReader(reader),
1345 .dispatchNamespace = getDispatchNamespaceFromReader(reader),
1346 .scriptId = getScriptIdFromReader(reader),
1347 .scriptTags = getScriptTagsFromReader(reader),
1348 .entrypoint = getEntrypointFromReader(reader),
1349 };
1350}
1351} // namespace
1352 
1353Onset::Onset(
1354 SpanId spanId, Onset::Info&& info, Onset::WorkerInfo&& workerInfo, CustomInfo attributes)
1355 : spanId(spanId),
1356 info(kj::mv(info)),
1357 workerInfo(kj::mv(workerInfo)),
1358 attributes(kj::mv(attributes)) {}
1359 
1360Onset::Onset(rpc::Trace::Onset::Reader reader)
1361 : spanId(reader.getSpanId()),
1362 info(readOnsetInfo(reader.getInfo())),
1363 workerInfo(getWorkerInfoFromReader(reader)),
1364 attributes(KJ_MAP(attr, reader.getAttributes()) { return Attribute(attr); }) {}
1365 
1366void Onset::copyTo(rpc::Trace::Onset::Builder builder) const {
1367 builder.setExecutionModel(workerInfo.executionModel);
1368 builder.setSpanId(spanId);
1369 KJ_IF_SOME(name, workerInfo.scriptName) {
1370 builder.setScriptName(name);
1371 }
1372 KJ_IF_SOME(version, workerInfo.scriptVersion) {
1373 builder.setScriptVersion(*version);
1374 }
1375 KJ_IF_SOME(name, workerInfo.dispatchNamespace) {
1376 builder.setDispatchNamespace(name);
1377 }
1378 KJ_IF_SOME(scriptId, workerInfo.scriptId) {
1379 builder.setScriptId(scriptId);
1380 }
1381 KJ_IF_SOME(tags, workerInfo.scriptTags) {
1382 auto list = builder.initScriptTags(tags.size());
1383 for (size_t i = 0; i < tags.size(); i++) {
1384 list.set(i, tags[i]);
1385 }
1386 }
1387 KJ_IF_SOME(e, workerInfo.entrypoint) {
1388 builder.setEntryPoint(e);
1389 }
1390 KJ_IF_SOME(p, workerInfo.preview) {
1391 p.copyTo(builder.initPreview());
1392 }
1393 auto infoBuilder = builder.initInfo();
1394 writeOnsetInfo(info, infoBuilder);
1395 
1396 auto attributeBuilder = builder.initAttributes(attributes.size());
1397 for (size_t n = 0; n < attributes.size(); n++) {
1398 attributes[n].copyTo(attributeBuilder[n]);
1399 }
1400}
1401 
1402Onset::WorkerInfo Onset::WorkerInfo::clone() const {
1403 return WorkerInfo{
1404 .executionModel = executionModel,
1405 .scriptName = mapCopyString(scriptName),
1406 .scriptVersion = scriptVersion.map([](auto& version) { return capnp::clone(*version); }),
1407 .preview = preview.map([](auto& preview) { return preview.clone(); }),
1408 .dispatchNamespace = mapCopyString(dispatchNamespace),
1409 .scriptId = mapCopyString(scriptId),
1410 .scriptTags =
1411 scriptTags.map([](auto& tags) { return KJ_MAP(tag, tags) { return kj::str(tag); }; }),
1412 .entrypoint = mapCopyString(entrypoint),
1413 };
1414}
1415 
1416EventInfo cloneEventInfo(const EventInfo& info) {
1417 KJ_SWITCH_ONEOF(info) {
1418 KJ_CASE_ONEOF(fetch, FetchEventInfo) {
1419 return fetch.clone();
1420 }
1421 KJ_CASE_ONEOF(connect, ConnectEventInfo) {
1422 return connect.clone();
1423 }
1424 KJ_CASE_ONEOF(jsrpc, JsRpcEventInfo) {
1425 return jsrpc.clone();
1426 }
1427 KJ_CASE_ONEOF(scheduled, ScheduledEventInfo) {
1428 return scheduled.clone();
1429 }
1430 KJ_CASE_ONEOF(alarm, AlarmEventInfo) {
1431 return alarm.clone();
1432 }
1433 KJ_CASE_ONEOF(queue, QueueEventInfo) {
1434 return queue.clone();
1435 }
1436 KJ_CASE_ONEOF(email, EmailEventInfo) {
1437 return email.clone();
1438 }
1439 KJ_CASE_ONEOF(trace, TraceEventInfo) {
1440 return trace.clone();
1441 }
1442 KJ_CASE_ONEOF(hws, HibernatableWebSocketEventInfo) {
1443 return hws.clone();
1444 }
1445 KJ_CASE_ONEOF(custom, CustomEventInfo) {
1446 return CustomEventInfo();
1447 }
1448 }
1449 KJ_UNREACHABLE;
1450}
1451 
1452Onset Onset::clone() const {
1453 return Onset(spanId, cloneEventInfo(info), workerInfo.clone(),
1454 KJ_MAP(attr, attributes) { return attr.clone(); });
1455}
1456 
1457Outcome::Outcome(EventOutcome outcome, kj::Duration cpuTime, kj::Duration wallTime)
1458 : outcome(outcome),
1459 cpuTime(cpuTime),
1460 wallTime(wallTime) {}
1461 
1462Outcome::Outcome(rpc::Trace::Outcome::Reader reader)
1463 : outcome(reader.getOutcome()),
1464 cpuTime(reader.getCpuTime() * kj::MILLISECONDS),
1465 wallTime(reader.getWallTime() * kj::MILLISECONDS) {}
1466 
1467void Outcome::copyTo(rpc::Trace::Outcome::Builder builder) const {
1468 builder.setOutcome(outcome);
1469 builder.setCpuTime(cpuTime / kj::MILLISECONDS);
1470 builder.setWallTime(wallTime / kj::MILLISECONDS);
1471}
1472 
1473Outcome Outcome::clone() const {
1474 return Outcome(outcome, cpuTime, wallTime);
1475}
1476 
1477TailEvent::TailEvent(
1478 SpanContext context, TraceId invocationId, kj::Date timestamp, kj::uint sequence, Event&& event)
1479 : spanContext(kj::mv(context)),
1480 invocationId(invocationId),
1481 timestamp(timestamp),
1482 sequence(sequence),
1483 event(kj::mv(event)) {}
1484 
1485TailEvent::TailEvent(TraceId traceId,
1486 TraceId invocationId,
1487 kj::Maybe<SpanId> spanId,
1488 kj::Date timestamp,
1489 kj::uint sequence,
1490 Event&& event,
1491 kj::Maybe<TraceFlags> traceFlags)
1492 : spanContext(kj::mv(traceId), kj::mv(spanId), kj::mv(traceFlags)),
1493 invocationId(kj::mv(invocationId)),
1494 timestamp(timestamp),
1495 sequence(sequence),
1496 event(kj::mv(event)) {}
1497 
1498namespace {
1499TailEvent::Event readEventFromTailEvent(const rpc::Trace::TailEvent::Reader& reader) {
1500 const auto event = reader.getEvent();
1501 switch (event.which()) {
1502 case rpc::Trace::TailEvent::Event::ONSET: {
1503 return Onset(event.getOnset());
1504 }
1505 case rpc::Trace::TailEvent::Event::OUTCOME: {
1506 return Outcome(event.getOutcome());
1507 }
1508 case rpc::Trace::TailEvent::Event::SPAN_OPEN: {
1509 return SpanOpen(event.getSpanOpen());
1510 }
1511 case rpc::Trace::TailEvent::Event::SPAN_CLOSE: {
1512 return SpanClose(event.getSpanClose());
1513 }
1514 case rpc::Trace::TailEvent::Event::ATTRIBUTE: {
1515 auto listReader = event.getAttribute();
1516 kj::Vector<Attribute> attrs(listReader.size());
1517 for (const auto& reader: listReader) {
1518 attrs.add(Attribute(reader));
1519 }
1520 return attrs.releaseAsArray();
1521 }
1522 case rpc::Trace::TailEvent::Event::RETURN: {
1523 return Return(event.getReturn());
1524 }
1525 case rpc::Trace::TailEvent::Event::DIAGNOSTIC_CHANNEL_EVENT: {
1526 return DiagnosticChannelEvent(event.getDiagnosticChannelEvent());
1527 }
1528 case rpc::Trace::TailEvent::Event::EXCEPTION: {
1529 return Exception(event.getException());
1530 }
1531 case rpc::Trace::TailEvent::Event::LOG: {
1532 return Log(event.getLog());
1533 }
1534 case rpc::Trace::TailEvent::Event::STREAM_DIAGNOSTICS: {
1535 return StreamDiagnosticsEvent(event.getStreamDiagnostics());
1536 }
1537 }
1538 KJ_UNREACHABLE;
1539}
1540} // namespace
1541 
1542TailEvent::TailEvent(rpc::Trace::TailEvent::Reader reader)
1543 : spanContext(SpanContext::fromCapnp(reader.getSpanContext())),
1544 invocationId(TraceId::fromCapnp(reader.getInvocationId())),
1545 timestamp(kj::UNIX_EPOCH + reader.getTimestampNs() * kj::NANOSECONDS),
1546 sequence(reader.getSequence()),
1547 event(readEventFromTailEvent(reader)) {}
1548 
1549void TailEvent::copyTo(rpc::Trace::TailEvent::Builder builder) const {
1550 spanContext.toCapnp(builder.initSpanContext());
1551 invocationId.toCapnp(builder.initInvocationId());
1552 builder.setTimestampNs((timestamp - kj::UNIX_EPOCH) / kj::NANOSECONDS);
1553 builder.setSequence(sequence);
1554 auto eventBuilder = builder.initEvent();
1555 KJ_SWITCH_ONEOF(event) {
1556 KJ_CASE_ONEOF(onset, Onset) {
1557 onset.copyTo(eventBuilder.initOnset());
1558 }
1559 KJ_CASE_ONEOF(outcome, Outcome) {
1560 outcome.copyTo(eventBuilder.initOutcome());
1561 }
1562 KJ_CASE_ONEOF(open, SpanOpen) {
1563 open.copyTo(eventBuilder.initSpanOpen());
1564 }
1565 KJ_CASE_ONEOF(close, SpanClose) {
1566 close.copyTo(eventBuilder.initSpanClose());
1567 }
1568 KJ_CASE_ONEOF(diag, DiagnosticChannelEvent) {
1569 diag.copyTo(eventBuilder.initDiagnosticChannelEvent());
1570 }
1571 KJ_CASE_ONEOF(ex, Exception) {
1572 ex.copyTo(eventBuilder.initException());
1573 }
1574 KJ_CASE_ONEOF(log, Log) {
1575 log.copyTo(eventBuilder.initLog());
1576 }
1577 KJ_CASE_ONEOF(streamDiag, StreamDiagnosticsEvent) {
1578 streamDiag.copyTo(eventBuilder.initStreamDiagnostics());
1579 }
1580 KJ_CASE_ONEOF(ret, Return) {
1581 ret.copyTo(eventBuilder.initReturn());
1582 }
1583 KJ_CASE_ONEOF(attrs, CustomInfo) {
1584 // Mark is a collection of attributes.
1585 auto attrBuilder = eventBuilder.initAttribute(attrs.size());
1586 for (size_t n = 0; n < attrs.size(); n++) {
1587 attrs[n].copyTo(attrBuilder[n]);
1588 }
1589 }
1590 }
1591}
1592 
1593TailEvent TailEvent::clone() const {
1594 constexpr auto cloneEvent = [](const Event& event) -> Event {
1595 KJ_SWITCH_ONEOF(event) {
1596 KJ_CASE_ONEOF(onset, Onset) {
1597 return onset.clone();
1598 }
1599 KJ_CASE_ONEOF(outcome, Outcome) {
1600 return outcome.clone();
1601 }
1602 KJ_CASE_ONEOF(open, SpanOpen) {
1603 return open.clone();
1604 }
1605 KJ_CASE_ONEOF(close, SpanClose) {
1606 return close.clone();
1607 }
1608 KJ_CASE_ONEOF(diag, DiagnosticChannelEvent) {
1609 return diag.clone();
1610 }
1611 KJ_CASE_ONEOF(ex, Exception) {
1612 return ex.clone();
1613 }
1614 KJ_CASE_ONEOF(log, Log) {
1615 return log.clone();
1616 }
1617 KJ_CASE_ONEOF(streamDiag, StreamDiagnosticsEvent) {
1618 return streamDiag.clone();
1619 }
1620 KJ_CASE_ONEOF(ret, Return) {
1621 return ret.clone();
1622 }
1623 KJ_CASE_ONEOF(attrs, CustomInfo) {
1624 return KJ_MAP(attr, attrs) { return attr.clone(); };
1625 }
1626 }
1627 KJ_UNREACHABLE;
1628 };
1629 return TailEvent(spanContext.getTraceId(), invocationId, spanContext.getSpanId(), timestamp,
1630 sequence, cloneEvent(event), spanContext.getTraceFlags());
1631}
1632 
1633SpanOpenData::SpanOpenData(rpc::SpanOpenData::Reader reader)
1634 : spanId(reader.getSpanId()),
1635 parentSpanId(reader.getParentSpanId()),
1636 operationName(kj::str(reader.getOperationName())),
1637 startTime(kj::UNIX_EPOCH + reader.getStartTimeNs() * kj::NANOSECONDS) {}
1638 
1639void SpanOpenData::copyTo(rpc::SpanOpenData::Builder builder) const {
1640 builder.setOperationName(operationName.asPtr());
1641 builder.setStartTimeNs((startTime - kj::UNIX_EPOCH) / kj::NANOSECONDS);
1642 builder.setSpanId(spanId);
1643 builder.setParentSpanId(parentSpanId);
1644}
1645 
1646SpanEndData::SpanEndData(rpc::SpanEndData::Reader reader)
1647 : spanId(reader.getSpanId()),
1648 endTime(kj::UNIX_EPOCH + reader.getEndTimeNs() * kj::NANOSECONDS) {
1649 auto tagsParam = reader.getTags();
1650 tags.reserve(tagsParam.size());
1651 for (auto tagParam: tagsParam) {
1652 tags.insert(kj::ConstString(kj::heapString(tagParam.getKey())),
1653 deserializeTagValue(tagParam.getValue()));
1654 }
1655}
1656 
1657void SpanEndData::copyTo(rpc::SpanEndData::Builder builder) const {
1658 builder.setEndTimeNs((endTime - kj::UNIX_EPOCH) / kj::NANOSECONDS);
1659 builder.setSpanId(spanId);
1660 
1661 auto tagsParam = builder.initTags(tags.size());
1662 auto i = 0;
1663 for (auto& tag: tags) {
1664 auto tagParam = tagsParam[i++];
1665 tagParam.setKey(tag.key.asPtr());
1666 serializeTagValue(tagParam.initValue(), tag.value);
1667 }
1668}
1669} // namespace tracing
1670 
1671// ======================================================================================
1672 
1673SpanBuilder::SpanBuilder(kj::Maybe<kj::Own<SpanObserver>> observer,
1674 kj::ConstString operationName,
1675 kj::Maybe<kj::Date> startTime) {
1676 KJ_IF_SOME(obs, observer) {
1677 // TODO(o11y): Once we report the user tracing spanOpen event as soon as a span is created, we
1678 // should be able to fold this virtual call and just get the timestamp directly.
1679 kj::Date time = startTime.orDefault([&]() { return obs->getTime(); });
1680 // Report spanOpen event for user tracing spans
1681 obs->onOpen(operationName.clone(), time);
1682 span.emplace(kj::mv(operationName), time);
1683 this->observer = kj::mv(obs);
1684 }
1685}
1686 
1687SpanBuilder& SpanBuilder::operator=(SpanBuilder&& other) {
1688 end();
1689 observer = kj::mv(other.observer);
1690 span = kj::mv(other.span);
1691 return *this;
1692}
1693 
1694SpanBuilder::~SpanBuilder() noexcept(false) {
1695 end();
1696}
1697 
1698void SpanBuilder::end() {
1699 KJ_IF_SOME(o, observer) {
1700 KJ_IF_SOME(s, span) {
1701 // TODO(performance): Fold this timer call if we are using I/O time, where we will look up
1702 // I/O time later.
1703 s.endTime = kj::systemPreciseCalendarClock().now();
1704 o->onClose(s.endTime, kj::mv(s.tags), kj::mv(s.logs));
1705 span = kj::none;
1706 }
1707 }
1708}
1709 
1710void SpanBuilder::setOperationName(kj::ConstString operationName) {
1711 KJ_IF_SOME(s, span) {
1712 KJ_IF_SOME(o, observer) {
1713 o->onUpdateName(operationName.clone());
1714 }
1715 s.operationName = kj::mv(operationName);
1716 }
1717}
1718 
1719void SpanBuilder::setTag(kj::ConstString key, TagInitValue value) {
1720 KJ_IF_SOME(s, span) {
1721 // We allow passing a LiteralStringConst or StringPtr so that we don't have to allocate memory
1722 // if we're not being observed.
1723 TagValue v = [](TagInitValue value) -> Span::TagValue {
1724 KJ_SWITCH_ONEOF(value) {
1725 KJ_CASE_ONEOF(str, kj::StringPtr) {
1726 return kj::ConstString(kj::str(str));
1727 }
1728 KJ_CASE_ONEOF(str, kj::LiteralStringConst) {
1729 return kj::ConstString(str);
1730 }
1731 KJ_CASE_ONEOF(str, kj::ConstString) {
1732 return kj::mv(str);
1733 }
1734 KJ_CASE_ONEOF(str, kj::String) {
1735 return kj::ConstString(kj::mv(str));
1736 }
1737 KJ_CASE_ONEOF(val, int64_t) {
1738 return val;
1739 }
1740 KJ_CASE_ONEOF(val, double) {
1741 return val;
1742 }
1743 KJ_CASE_ONEOF(val, bool) {
1744 return val;
1745 }
1746 }
1747 KJ_UNREACHABLE;
1748 }(kj::mv(value));
1749 
1750 auto keyPtr = key.asPtr();
1751 s.tags.upsert(kj::mv(key), kj::mv(v), [keyPtr](TagValue& existingValue, TagValue&& newValue) {
1752 // This is a programming error, but not a serious one. We could alternatively just emit
1753 // duplicate tags and leave the Jaeger UI in charge of warning about them.
1754 [[maybe_unused]] static auto logged = [keyPtr]() {
1755 if (isPredictableModeForTest()) {
1756 // Logging in ERROR level to have this fail loudly during testing.
1757 KJ_LOG(ERROR, "overwriting previous tag", keyPtr);
1758 } else {
1759 KJ_LOG(WARNING, "overwriting previous tag", keyPtr);
1760 }
1761 return true;
1762 }();
1763 existingValue = kj::mv(newValue);
1764 });
1765 }
1766}
1767 
1768void SpanBuilder::addLog(kj::Date timestamp, kj::ConstString key, TagValue value) {
1769 KJ_IF_SOME(s, span) {
1770 if (s.logs.size() >= Span::MAX_LOGS) {
1771 ++s.droppedLogs;
1772 } else {
1773 s.logs.add(Span::Log{.timestamp = timestamp,
1774 .tag = {
1775 .key = kj::mv(key),
1776 .value = kj::mv(value),
1777 }});
1778 }
1779 }
1780}
1781 
1782void TraceContext::setTag(kj::ConstString key, SpanBuilder::TagInitValue value) {
1783 if (!isObserved()) {
1784 return;
1785 }
1786 // Fast path (without string allocations) if only some spans are observed.
1787 if (!span.isObserved()) {
1788 userSpan.setTag(kj::mv(key), kj::mv(value));
1789 return;
1790 }
1791 if (!userSpan.isObserved()) {
1792 span.setTag(kj::mv(key), kj::mv(value));
1793 return;
1794 }
1795 
1796 // We need to duplicate the key and value since both are move-only types.
1797 // Clone the value based on its type.
1798 KJ_SWITCH_ONEOF(value) {
1799 KJ_CASE_ONEOF(s, kj::StringPtr) {
1800 span.setTag(key.clone(), s);
1801 userSpan.setTag(kj::mv(key), s);
1802 }
1803 KJ_CASE_ONEOF(s, kj::String) {
1804 span.setTag(key.clone(), kj::str(s));
1805 userSpan.setTag(kj::mv(key), kj::mv(s));
1806 }
1807 KJ_CASE_ONEOF(s, kj::LiteralStringConst) {
1808 span.setTag(key.clone(), s);
1809 userSpan.setTag(kj::mv(key), s);
1810 }
1811 KJ_CASE_ONEOF(s, kj::ConstString) {
1812 span.setTag(key.clone(), s.clone());
1813 userSpan.setTag(kj::mv(key), kj::mv(s));
1814 }
1815 KJ_CASE_ONEOF(b, bool) {
1816 span.setTag(key.clone(), b);
1817 userSpan.setTag(kj::mv(key), b);
1818 }
1819 KJ_CASE_ONEOF(d, double) {
1820 span.setTag(key.clone(), d);
1821 userSpan.setTag(kj::mv(key), d);
1822 }
1823 KJ_CASE_ONEOF(i, int64_t) {
1824 span.setTag(key.clone(), i);
1825 userSpan.setTag(kj::mv(key), i);
1826 }
1827 }
1828}
1829 
1830Span::TagValue spanTagClone(const Span::TagValue& tag) {
1831 KJ_SWITCH_ONEOF(tag) {
1832 KJ_CASE_ONEOF(str, kj::ConstString) {
1833 return str.clone();
1834 }
1835 KJ_CASE_ONEOF(val, int64_t) {
1836 return val;
1837 }
1838 KJ_CASE_ONEOF(val, double) {
1839 return val;
1840 }
1841 KJ_CASE_ONEOF(val, bool) {
1842 return val;
1843 }
1844 }
1845 KJ_UNREACHABLE;
1846}
1847 
1848using RpcValue = rpc::TagValue;
1849void serializeTagValue(RpcValue::Builder builder, const Span::TagValue& value) {
1850 KJ_SWITCH_ONEOF(value) {
1851 KJ_CASE_ONEOF(b, bool) {
1852 builder.setBool(b);
1853 }
1854 KJ_CASE_ONEOF(i, int64_t) {
1855 builder.setInt64(i);
1856 }
1857 KJ_CASE_ONEOF(d, double) {
1858 builder.setFloat64(d);
1859 }
1860 KJ_CASE_ONEOF(s, kj::ConstString) {
1861 builder.setString(s.asPtr());
1862 }
1863 }
1864}
1865 
1866Span::TagValue deserializeTagValue(RpcValue::Reader value) {
1867 switch (value.which()) {
1868 case RpcValue::BOOL:
1869 return value.getBool();
1870 case RpcValue::FLOAT64:
1871 return value.getFloat64();
1872 case RpcValue::INT64:
1873 return value.getInt64();
1874 case RpcValue::STRING:
1875 return kj::ConstString(kj::heapString(value.getString()));
1876 default:
1877 KJ_UNREACHABLE;
1878 }
1879}
1880 
1881ScopedDurationTagger::ScopedDurationTagger(
1882 SpanBuilder& span, kj::ConstString key, const kj::MonotonicClock& timer)
1883 : span(span),
1884 key(kj::mv(key)),
1885 timer(timer),
1886 startTime(timer.now()) {}
1887 
1888ScopedDurationTagger::~ScopedDurationTagger() noexcept(false) {
1889 auto duration = timer.now() - startTime;
1890 if (isPredictableModeForTest()) {
1891 duration = 0 * kj::NANOSECONDS;
1892 }
1893 span.setTag(kj::mv(key), duration / kj::NANOSECONDS);
1894}
1895 
1896} // namespace workerd