Skip to content
File

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

24.5 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/io-context.h>
6#include <workerd/io/trace-stream.h>
7#include <workerd/io/tracer.h>
8#include <workerd/util/sentry.h>
9#include <workerd/util/thread-scopes.h>
10 
11#include <capnp/message.h> // for capnp::clone()
12 
13namespace workerd {
14 
15namespace {
16 
17// Approximately how much external data we allow in a trace before we start ignoring requests. We
18// want this number to be big enough to be useful for tracing, but small enough to make it hard to
19// DoS the C++ heap -- keeping in mind we can record a trace per handler run during a request. For
20// streaming tail worker, this is the maximum size per tail event.
21// TODO(streaming-tail): Add a clear indicator for events being truncated based on MAX_TRACE_BYTES
22// so that developers can understand why this happens.
23static constexpr size_t MAX_TRACE_BYTES = 256 * 1024;
24 
25tracing::Attribute::Value cloneAttributeValue(const tracing::Attribute::Value& value) {
26 KJ_SWITCH_ONEOF(value) {
27 KJ_CASE_ONEOF(boolean, bool) {
28 return tracing::Attribute::Value(boolean);
29 }
30 KJ_CASE_ONEOF(number, double) {
31 return tracing::Attribute::Value(number);
32 }
33 KJ_CASE_ONEOF(integer, int64_t) {
34 return tracing::Attribute::Value(integer);
35 }
36 KJ_CASE_ONEOF(string, kj::ConstString) {
37 return tracing::Attribute::Value(string.clone());
38 }
39 }
40 KJ_UNREACHABLE;
41}
42} // namespace
43 
44kj::Promise<kj::Own<Trace>> WorkerTracer::onComplete() {
45 KJ_REQUIRE(completeFulfiller == kj::none, "onComplete() can only be called once");
46 
47 auto paf = kj::newPromiseAndFulfiller<kj::Own<Trace>>();
48 completeFulfiller = kj::mv(paf.fulfiller);
49 return kj::mv(paf.promise);
50}
51 
52WorkerTracer::WorkerTracer(kj::Maybe<kj::Rc<kj::Refcounted>> parentPipeline,
53 kj::Own<Trace> trace,
54 PipelineLogLevel pipelineLogLevel,
55 kj::Maybe<kj::Array<tracing::Attribute>> tailAttributes,
56 kj::Maybe<kj::Own<tracing::TailStreamWriter>> maybeTailStreamWriter)
57 : pipelineLogLevel(pipelineLogLevel),
58 trace(kj::mv(trace)),
59 parentPipeline(kj::mv(parentPipeline)),
60 maybeTailStreamWriter(kj::mv(maybeTailStreamWriter)) {
61 KJ_IF_SOME(tags, tailAttributes) {
62 if (tags.size() == 0) {
63 tailAttributes = kj::none;
64 } else {
65 for (auto& tag: tags) {
66 KJ_REQUIRE(tag.value.size() == 1, "tail attributes must contain exactly one value");
67 setWorkerAttribute(tag.name.clone(), cloneAttributeValue(tag.value[0]));
68 }
69 }
70 }
71 this->trace->tailAttributes = kj::mv(tailAttributes);
72}
73 
74WorkerTracer::~WorkerTracer() noexcept(false) {
75 // Report the outcome event, which should have been delivered by now.
76 
77 // Do not attempt to report an outcome event if logging is disabled, as with other event types.
78 if (pipelineLogLevel == PipelineLogLevel::NONE) {
79 return;
80 }
81 
82 // Report the outcome event if STWs are present. All worker events need to call setEventInfo at
83 // the start of the invocation to submit the onset event before any other tail events.
84 KJ_IF_SOME(writer, maybeTailStreamWriter) {
85 KJ_IF_SOME(spanContext, topLevelInvocationSpanContext) {
86 if (markedUnused) {
87 LOG_WARNING_PERIODICALLY("WorkerTracer was marked unused but actually was used");
88 }
89 if (isPredictableModeForTest()) {
90 writer->report(spanContext,
91 tracing::Outcome(trace->outcome, 0 * kj::MILLISECONDS, 0 * kj::MILLISECONDS),
92 completeTime, 0);
93 } else {
94 writer->report(spanContext,
95 tracing::Outcome(trace->outcome, trace->cpuTime, trace->wallTime), completeTime, 0);
96 }
97 } else if (!markedUnused) {
98 // If no span context is available, we have a streaming tail worker set up but shut down the
99 // worker tracer without ever sending an Onset event. In that case we either failed to set up
100 // the Onset properly (indicating a bug – all event types are required to report an Onset at
101 // the start – although this is more likely to manifest as a "Tail stream onset was not
102 // reported" error) or we created a WorkerInterface with WorkerTracer without ever invoking it
103 // (which is not incorrect behavior, but likely indicates inefficient code that sets up
104 // WorkerInterfaces and then ends up not using it due to an error/incorrect parameters; such
105 // error checking should be done beforehand to avoid unused allocations). Report such cases.
106 // Note: If markedUnused is true, this tracer was intentionally not used (e.g., duplicate
107 // alarm request deduplication) and the warning should be suppressed.
108 LOG_ERROR_PERIODICALLY(
109 "destructed WorkerTracer with STW without reporting Onset event", kj::getStackTrace());
110 }
111 }
112 
113 // Report the completed trace, if fulfiller is set up.
114 KJ_IF_SOME(f, completeFulfiller) {
115 f.get()->fulfill(kj::mv(trace));
116 }
117};
118 
119constexpr kj::LiteralStringConst logSizeExceeded =
120 "[\"Log size limit exceeded: More than 256KB of data (across console.log statements, exception, request metadata and headers) was logged during a single request. Subsequent data for this request will not be recorded in logs, appear when tailing this Worker's logs, or in Tail Workers.\"]"_kjc;
121 
122void WorkerTracer::addLog(const tracing::InvocationSpanContext& context,
123 kj::Date timestamp,
124 LogLevel logLevel,
125 kj::String message) {
126 if (pipelineLogLevel == PipelineLogLevel::NONE) {
127 return;
128 }
129 
130 // TODO(streaming-tail): Here we add the log to the trace object and the tail stream writer, if
131 // available. If the given worker stage is only tailed by a streaming tail worker, adding the log
132 // to the buffered trace object is not needed; this will be addressed in a future refactor.
133 KJ_IF_SOME(writer, maybeTailStreamWriter) {
134 // If message is too big on its own, truncate it.
135 size_t messageSize = kj::min(message.size(), MAX_TRACE_BYTES);
136 writer->report(context,
137 {tracing::Log(timestamp, logLevel, kj::str(message.first(messageSize)))}, timestamp,
138 messageSize);
139 }
140 
141 if (trace->exceededLogLimit) {
142 return;
143 }
144 
145 size_t messageSize = sizeof(tracing::Log) + message.size();
146 if (trace->bytesUsed + messageSize > MAX_TRACE_BYTES) {
147 // We use a JSON encoded array/string to match other console.log() recordings:
148 trace->logs.add(timestamp, LogLevel::WARN, kj::str(logSizeExceeded));
149 trace->exceededLogLimit = true;
150 trace->truncated = true;
151 } else {
152 trace->bytesUsed += messageSize;
153 trace->logs.add(timestamp, logLevel, kj::mv(message));
154 }
155}
156 
157void WorkerTracer::addSpanOpen(tracing::SpanId spanId,
158 tracing::SpanId parentSpanId,
159 kj::ConstString operationName,
160 kj::Date startTime) {
161 if (pipelineLogLevel == PipelineLogLevel::NONE) {
162 return;
163 }
164 
165 auto& tailStreamWriter = KJ_UNWRAP_OR_RETURN(maybeTailStreamWriter);
166 auto& topLevelContext = KJ_ASSERT_NONNULL(topLevelInvocationSpanContext);
167 // Compose SpanOpen. An all-zero spanId is interpreted as having no spans above this one, thus we
168 // use the Onset spanId instead (taken from topLevelContext). We go to great lengths to rule out
169 // getting an all-zero spanId by chance (see SpanId::fromEntropy()), so this should be safe.
170 if (parentSpanId == tracing::SpanId::nullId) {
171 parentSpanId = topLevelContext.getSpanId();
172 }
173 size_t spanNameSize = operationName.size();
174 auto spanOpenContext = tracing::InvocationSpanContext(topLevelContext.getTraceId(),
175 topLevelContext.getInvocationId(), parentSpanId, topLevelContext.getTraceFlags());
176 tailStreamWriter->report(
177 spanOpenContext, tracing::SpanOpen(spanId, kj::mv(operationName)), startTime, spanNameSize);
178}
179 
180void WorkerTracer::addSpanClose(tracing::SpanEndData&& span, kj::Maybe<kj::Date> maybeStartTime) {
181 if (pipelineLogLevel == PipelineLogLevel::NONE) {
182 return;
183 }
184 
185 // Note: spans are not available in the buffered tail worker, so we don't need an exceededSpanLimit
186 // variable for it and it can't cause truncation.
187 auto& tailStreamWriter = KJ_UNWRAP_OR_RETURN(maybeTailStreamWriter);
188 
189 adjustSpanTime(span, maybeStartTime);
190 
191 size_t spanTagsSize = 0;
192 for (const Span::TagMap::Entry& tag: span.tags) {
193 spanTagsSize += tag.key.size();
194 KJ_SWITCH_ONEOF(tag.value) {
195 KJ_CASE_ONEOF(str, kj::ConstString) {
196 spanTagsSize += str.size();
197 }
198 KJ_CASE_ONEOF(val, bool) {
199 spanTagsSize++;
200 }
201 // int64_t and double
202 KJ_CASE_ONEOF_DEFAULT {
203 spanTagsSize += sizeof(int64_t);
204 }
205 }
206 }
207 
208 // Compose Attributes and SpanClose, which are available at span completion time and transmitted
209 // together.
210 auto& topLevelContext = KJ_ASSERT_NONNULL(topLevelInvocationSpanContext);
211 auto spanComponentContext = tracing::InvocationSpanContext(topLevelContext.getTraceId(),
212 topLevelContext.getInvocationId(), span.spanId, topLevelContext.getTraceFlags());
213 
214 if (span.tags.size() && spanTagsSize <= MAX_TRACE_BYTES) {
215 tracing::CustomInfo attr = KJ_MAP(tag, span.tags) {
216 return tracing::Attribute(kj::mv(tag.key), kj::mv(tag.value));
217 };
218 tailStreamWriter->report(spanComponentContext, kj::mv(attr), span.endTime, spanTagsSize);
219 }
220 tailStreamWriter->report(spanComponentContext, tracing::SpanClose(), span.endTime, 0);
221}
222 
223void WorkerTracer::addException(const tracing::InvocationSpanContext& context,
224 kj::Date timestamp,
225 kj::String name,
226 kj::String message,
227 kj::Maybe<kj::String> stack) {
228 // TODO(someday): For now, we're using logLevel == none as a hint to avoid doing anything
229 // expensive while tracing. We may eventually want separate configuration for exceptions vs.
230 // logs.
231 if (pipelineLogLevel == PipelineLogLevel::NONE) {
232 return;
233 }
234 
235 size_t messageSize = sizeof(tracing::Exception) + name.size() + message.size();
236 KJ_IF_SOME(s, stack) {
237 messageSize += s.size();
238 }
239 KJ_IF_SOME(writer, maybeTailStreamWriter) {
240 auto maybeTruncatedName = name.first(kj::min(name.size(), MAX_TRACE_BYTES));
241 auto maybeTruncatedMessage =
242 message.first(kj::min(message.size(), MAX_TRACE_BYTES - maybeTruncatedName.size()));
243 kj::Maybe<kj::String> maybeTruncatedStack;
244 auto maybeTruncatedStackSize = 0;
245 KJ_IF_SOME(s, stack) {
246 maybeTruncatedStackSize = kj::min(
247 s.size(), MAX_TRACE_BYTES - maybeTruncatedName.size() - maybeTruncatedMessage.size());
248 maybeTruncatedStack = kj::heapString(s.first(maybeTruncatedStackSize));
249 }
250 writer->report(context,
251 {tracing::Exception(timestamp, kj::str(maybeTruncatedName), kj::str(maybeTruncatedMessage),
252 kj::mv(maybeTruncatedStack))},
253 timestamp,
254 maybeTruncatedName.size() + maybeTruncatedMessage.size() + maybeTruncatedStackSize);
255 }
256 
257 if (trace->exceededExceptionLimit) {
258 return;
259 }
260 
261 if (trace->bytesUsed + messageSize > MAX_TRACE_BYTES) {
262 trace->exceededExceptionLimit = true;
263 trace->truncated = true;
264 trace->exceptions.add(timestamp, kj::str("Error"),
265 kj::str("Trace resource limit exceeded; subsequent exceptions not recorded."), kj::none);
266 } else {
267 trace->bytesUsed += messageSize;
268 trace->exceptions.add(timestamp, kj::mv(name), kj::mv(message), kj::mv(stack));
269 }
270}
271 
272void WorkerTracer::addDiagnosticChannelEvent(const tracing::InvocationSpanContext& context,
273 kj::Date timestamp,
274 kj::String channel,
275 kj::Array<kj::byte> message) {
276 if (pipelineLogLevel == PipelineLogLevel::NONE) {
277 return;
278 }
279 
280 size_t messageSize = sizeof(tracing::DiagnosticChannelEvent) + channel.size() + message.size();
281 KJ_IF_SOME(writer, maybeTailStreamWriter) {
282 // Drop oversized diagnostic channel events instead of truncating them – a truncated message may
283 // not be deserialized correctly.
284 if (messageSize <= MAX_TRACE_BYTES) {
285 writer->report(context,
286 {tracing::DiagnosticChannelEvent(
287 timestamp, kj::str(channel), kj::heapArray<kj::byte>(message))},
288 timestamp, messageSize);
289 }
290 }
291 
292 if (trace->exceededDiagnosticChannelEventLimit) {
293 return;
294 }
295 
296 if (trace->bytesUsed + messageSize > MAX_TRACE_BYTES) {
297 trace->exceededDiagnosticChannelEventLimit = true;
298 trace->truncated = true;
299 trace->diagnosticChannelEvents.add(
300 timestamp, kj::str("workerd.LimitExceeded"), kj::Array<kj::byte>());
301 } else {
302 trace->bytesUsed += messageSize;
303 trace->diagnosticChannelEvents.add(timestamp, kj::mv(channel), kj::mv(message));
304 }
305}
306 
307void WorkerTracer::setEventInfo(
308 IoContext::IncomingRequest& incomingRequest, tracing::EventInfo&& info) {
309 // IoContext is available at this time, capture weakRef.
310 KJ_ASSERT(weakIoContext == kj::none, "tracer can only be used for a single event");
311 weakIoContext = incomingRequest.getContext().getWeakRef();
312 setEventInfoInternal(
313 incomingRequest.getInvocationSpanContext(), incomingRequest.now(), kj::mv(info));
314}
315 
316void WorkerTracer::setEventInfoInternal(
317 const tracing::InvocationSpanContext& context, kj::Date timestamp, tracing::EventInfo&& info) {
318 KJ_ASSERT(trace->eventInfo == kj::none, "tracer can only be used for a single event");
319 
320 // TODO(someday): For now, we're using logLevel == none as a hint to avoid doing anything
321 // expensive while tracing. We may eventually want separate configuration for event info vs.
322 // logs.
323 // TODO(perf): Find a way to allow caller to avoid the cost of generation if the info struct
324 // won't be used?
325 if (pipelineLogLevel == PipelineLogLevel::NONE) {
326 return;
327 }
328 
329 trace->eventTimestamp = timestamp;
330 this->topLevelInvocationSpanContext = context.clone();
331 
332 size_t eventSize = 0;
333 KJ_SWITCH_ONEOF(info) {
334 KJ_CASE_ONEOF(fetch, tracing::FetchEventInfo) {
335 eventSize += fetch.url.size();
336 for (const auto& header: fetch.headers) {
337 eventSize += header.name.size() + header.value.size();
338 }
339 eventSize += fetch.cfJson.size();
340 // Limit STW onset to MAX_TRACE_BYTES, beyond that dispatch a truncated event too.
341 if (eventSize > MAX_TRACE_BYTES) {
342 info = tracing::FetchEventInfo(fetch.method, {}, {}, {});
343 }
344 }
345 KJ_CASE_ONEOF_DEFAULT {}
346 }
347 
348 KJ_IF_SOME(writer, maybeTailStreamWriter) {
349 // Provide WorkerInfo to the streaming tail worker if available. This data is provided when the
350 // WorkerTracer is created, but the actual onset event is the best time to send it.
351 auto workerInfo = tracing::Onset::WorkerInfo{
352 .executionModel = trace->executionModel,
353 .scriptName = mapCopyString(trace->scriptName),
354 .scriptVersion =
355 trace->scriptVersion.map([](auto& scriptVersion) -> kj::Own<ScriptVersion::Reader> {
356 return capnp::clone(*scriptVersion);
357 }),
358 .preview = trace->preview.map([](auto& preview) { return preview.clone(); }),
359 .dispatchNamespace = mapCopyString(trace->dispatchNamespace),
360 .scriptId = mapCopyString(trace->scriptId),
361 .scriptTags = KJ_MAP(tag, trace->scriptTags) { return kj::str(tag); },
362 .entrypoint = mapCopyString(trace->entrypoint),
363 };
364 
365 tracing::SpanId parentSpanId = tracing::SpanId::nullId;
366 KJ_IF_SOME(trigger, context.getParent()) {
367 parentSpanId = trigger.getSpanId();
368 }
369 // Onset needs special handling for spanId: The top-level spanId is zero unless a trigger
370 // context is available. The inner spanId is taken from the invocation
371 // span context, that span is being "opened" with the onset event. All other tail events have it
372 // as its parent span ID, except for recursive SpanOpens (which have the parent span instead)
373 // and Attribute/SpanClose events (which have the spanId opened in the corresponding SpanOpen).
374 auto onsetContext = tracing::InvocationSpanContext(
375 context.getTraceId(), context.getInvocationId(), parentSpanId, context.getTraceFlags());
376 
377 // Not applying size accounting for Onset since it is sent separately
378 writer->report(onsetContext,
379 tracing::Onset(context.getSpanId(), cloneEventInfo(info), kj::mv(workerInfo),
380 attributes.releaseAsArray()),
381 timestamp, 0);
382 }
383 
384 // truncation should only be needed for fetch events, since we only set eventSize there.
385 if (trace->bytesUsed + eventSize > MAX_TRACE_BYTES && eventSize > 0) {
386 trace->truncated = true;
387 trace->logs.add(timestamp, LogLevel::WARN,
388 kj::str("[\"Trace resource limit exceeded; could not capture event info.\"]"));
389 trace->eventInfo =
390 tracing::FetchEventInfo(info.get<tracing::FetchEventInfo>().method, {}, {}, {});
391 } else {
392 trace->bytesUsed += eventSize;
393 trace->eventInfo = kj::mv(info);
394 }
395}
396 
397void WorkerTracer::setOutcome(EventOutcome outcome, kj::Duration cpuTime, kj::Duration wallTime) {
398 trace->outcome = outcome;
399 trace->cpuTime = cpuTime;
400 trace->wallTime = wallTime;
401 
402 // Defer reporting the actual outcome event to the WorkerTracer destructor: The outcome is
403 // reported when the metrics request is deallocated, but with ctx.waitUntil() there might be spans
404 // continuing to exist beyond that point. By the time the WorkerTracer is deallocated, the
405 // IoContext and its task set will be done and any additional spans will have wrapped up.
406 // This is somewhat at odds with the concept of "streaming" events, but benign as the WorkerTracer
407 // wraps up right after the metrics request object in the average case and since the outcome has a
408 // fixed size.
409}
410 
411void WorkerTracer::recordTimestamp(kj::Date timestamp) {
412 if (completeTime == kj::UNIX_EPOCH) {
413 completeTime = timestamp;
414 }
415}
416 
417kj::Date BaseTracer::getTime() {
418 auto& weakIoCtx = KJ_ASSERT_NONNULL(weakIoContext);
419 kj::Date timestamp = kj::UNIX_EPOCH;
420 weakIoCtx->runIfAlive([&timestamp](IoContext& context) { timestamp = context.now(); });
421 if (!weakIoCtx->isValid()) {
422 // This can happen if we the IoContext gets destroyed following an exception, but we still need
423 // to report a time for the return event.
424 if (completeTime != kj::UNIX_EPOCH) {
425 timestamp = completeTime;
426 } else {
427 // Otherwise, we can't actually get an end timestamp that makes sense.
428 if (isPredictableModeForTest()) {
429 KJ_FAIL_ASSERT("reported return event without valid IoContext or completeTime");
430 } else {
431 LOG_WARNING_PERIODICALLY("reported return event without valid IoContext or completeTime");
432 }
433 }
434 }
435 return timestamp;
436}
437 
438void BaseTracer::adjustSpanTime(tracing::SpanEndData& span, kj::Maybe<kj::Date> maybeStartTime) {
439 // To report I/O time, we need the IOContext to still be alive.
440 // weakIoContext is only none if we are tracing via RPC (in this case span times have already been
441 // adjusted) or if we failed to transmit an Onset event (in that case we'll get an error based on
442 // missing topLevelInvocationSpanContext right after).
443 if (weakIoContext != kj::none) {
444 auto& weakIoCtx = KJ_ASSERT_NONNULL(weakIoContext);
445 // startTime is generally available when we are not tracing via RPC, so we can assert that it is
446 // present. For the RPC case, the adjustment will already have been done earlier and it's ok
447 // for maybeStartTime to be none as this code won't run based on weakIoContext being none.
448 kj::Date startTime = KJ_ASSERT_NONNULL(maybeStartTime);
449 weakIoCtx->runIfAlive([this, &span, &startTime](IoContext& context) {
450 if (context.hasCurrentIncomingRequest()) {
451 span.endTime = context.now();
452 } else {
453 // We have an IOContext, but there's no current IncomingRequest. Always log a warning here,
454 // this should not be happening. Still report completeTime as a useful timestamp if
455 // available.
456 bool hasCompleteTime = false;
457 if (completeTime != kj::UNIX_EPOCH) {
458 span.endTime = completeTime;
459 hasCompleteTime = true;
460 } else {
461 span.endTime = startTime;
462 }
463 if (isPredictableModeForTest()) {
464 KJ_FAIL_ASSERT("reported span without current request", hasCompleteTime);
465 } else {
466 LOG_WARNING_PERIODICALLY("reported span without current request");
467 }
468 }
469 });
470 if (!weakIoCtx->isValid()) {
471 // This can happen if we start a customEvent from this event and cancel it after this IoContext
472 // gets destroyed. In that case we no longer have an IoContext available and can't get the
473 // current time, but the outcome timestamp will have already been set. Since the outcome
474 // timestamp is "late enough", simply use that.
475 // TODO(o11y): fix this – spans should not be outliving the IoContext.
476 if (completeTime != kj::UNIX_EPOCH) {
477 span.endTime = completeTime;
478 } else {
479 // Otherwise, we can't actually get an end timestamp that makes sense. Report a zero-duration
480 // span and log a warning (or fail assert in test mode).
481 span.endTime = startTime;
482 if (isPredictableModeForTest()) {
483 KJ_FAIL_ASSERT("reported span after IoContext was deallocated");
484 } else {
485 KJ_LOG(WARNING, "reported span after IoContext was deallocated");
486 }
487 }
488 }
489 }
490}
491 
492void WorkerTracer::setReturn(
493 kj::Maybe<kj::Date> timestamp, kj::Maybe<tracing::FetchResponseInfo> fetchResponseInfo) {
494 // Match the behavior of setEventInfo(). Any resolution of the TODO comments in setEventInfo()
495 // that are related to this check will probably also affect this function.
496 if (pipelineLogLevel == PipelineLogLevel::NONE) {
497 return;
498 }
499 
500 KJ_IF_SOME(writer, maybeTailStreamWriter) {
501 auto& spanContext = KJ_UNWRAP_OR_RETURN(topLevelInvocationSpanContext);
502 
503 // Fall back to weak IoContext if no timestamp is available
504 writer->report(spanContext,
505 tracing::Return({fetchResponseInfo.map([](auto& info) { return info.clone(); })}),
506 timestamp.orDefault([&]() { return getTime(); }), 0);
507 }
508 
509 // Add fetch response info for buffered tail worker
510 KJ_IF_SOME(info, fetchResponseInfo) {
511 KJ_REQUIRE(KJ_REQUIRE_NONNULL(trace->eventInfo).is<tracing::FetchEventInfo>());
512 KJ_ASSERT(trace->fetchResponseInfo == kj::none, "setFetchResponseInfo can only be called once");
513 trace->fetchResponseInfo = kj::mv(info);
514 }
515}
516 
517void BaseTracer::setMakeUserRequestSpanFunc(MakeUserRequestSpanFunc func) {
518 KJ_ASSERT(
519 makeUserRequestSpanFunc == kj::none, "setMakeUserRequestSpanFunc can only be called once");
520 makeUserRequestSpanFunc = kj::mv(func);
521}
522 
523void WorkerTracer::setWorkerAttribute(kj::ConstString key, Span::TagValue value) {
524 attributes.add(tracing::Attribute{kj::mv(key), kj::mv(value)});
525}
526 
527SpanParent BaseTracer::makeUserRequestSpan(
528 tracing::TraceId traceId, kj::Maybe<tracing::TraceFlags> traceFlags) {
529 KJ_IF_SOME(func, makeUserRequestSpanFunc) {
530 return func(kj::mv(traceId), traceFlags);
531 } else {
532 return SpanParent(nullptr);
533 }
534}
535 
536void WorkerTracer::setJsRpcInfo(const tracing::InvocationSpanContext& context,
537 kj::Date timestamp,
538 const kj::ConstString& methodName) {
539 if (pipelineLogLevel == PipelineLogLevel::NONE) {
540 return;
541 }
542 
543 // Update the method name in the already-set JsRpcEventInfo for buffered tail worker compatibility
544 KJ_IF_SOME(info, trace->eventInfo) {
545 KJ_SWITCH_ONEOF(info) {
546 KJ_CASE_ONEOF(jsRpcInfo, tracing::JsRpcEventInfo) {
547 jsRpcInfo.methodName = kj::str(methodName);
548 }
549 KJ_CASE_ONEOF_DEFAULT {}
550 }
551 }
552 
553 KJ_IF_SOME(writer, maybeTailStreamWriter) {
554 auto tag = tracing::Attribute("jsrpc.method"_kjc, methodName.clone());
555 writer->report(context, kj::arr(kj::mv(tag)), timestamp, methodName.size());
556 }
557}
558 
559kj::Own<SpanObserver> UserSpanObserver::newChild() {
560 return kj::refcounted<UserSpanObserver>(kj::addRef(*submitter), spanId, traceId, traceFlags);
561}
562 
563kj::Own<SpanObserver> UserSpanObserver::newChildFromUserCode() {
564 return kj::refcounted<UserSpanObserver>(
565 kj::addRef(*submitter), spanId, traceId, traceFlags, /*fromUserCode=*/true);
566}
567 
568kj::Maybe<tracing::SpanContext> UserSpanObserver::toSpanContext() {
569 if (traceId == nullptr) {
570 return kj::none;
571 }
572 return tracing::SpanContext(traceId, spanId, traceFlags);
573}
574 
575void UserSpanObserver::onClose(
576 kj::Date endTime, Span::TagMap&& tags, kj::Vector<Span::Log>&& logs) {
577 // span logs are not supported in user tracing.
578 (void)logs;
579 if (wasAccepted) {
580 submitter->submitSpanClose(spanId, startTime, endTime, kj::mv(tags));
581 }
582}
583 
584void UserSpanObserver::onOpen(kj::ConstString operationName, kj::Date startTime) {
585 this->startTime = startTime;
586 if (fromUserCode) {
587 wasAccepted =
588 submitter->submitUserSpanOpen(spanId, parentSpanId, kj::mv(operationName), startTime);
589 } else {
590 wasAccepted = submitter->submitSpanOpen(spanId, parentSpanId, kj::mv(operationName), startTime);
591 }
592}
593 
594// Provide I/O time to the tracing system for user spans.
595kj::Date UserSpanObserver::getTime() {
596 return IoContext::current().now();
597}
598 
599tracing::SpanId UserSpanObserver::getSpanId() {
600 return spanId;
601}
602 
603} // namespace workerd