File
Blob: src/workerd/io/tracer.c++
| 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 | |
| 13 | namespace workerd { |
| 14 | |
| 15 | namespace { |
| 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. |
| 23 | static constexpr size_t MAX_TRACE_BYTES = 256 * 1024; |
| 24 | |
| 25 | tracing::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 | |
| 44 | kj::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 | |
| 52 | WorkerTracer::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 | |
| 74 | WorkerTracer::~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 | |
| 119 | constexpr 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 | |
| 122 | void 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 | |
| 157 | void 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 | |
| 180 | void 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 | |
| 223 | void 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 | |
| 272 | void 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 | |
| 307 | void 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 | |
| 316 | void 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 | |
| 397 | void 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 | |
| 411 | void WorkerTracer::recordTimestamp(kj::Date timestamp) { |
| 412 | if (completeTime == kj::UNIX_EPOCH) { |
| 413 | completeTime = timestamp; |
| 414 | } |
| 415 | } |
| 416 | |
| 417 | kj::Date BaseTracer::getTime() { |
| 418 | auto& weakIoCtx = KJ_ASSERT_NONNULL(weakIoContext); |
| 419 | kj::Date timestamp = kj::UNIX_EPOCH; |
| 420 | weakIoCtx->runIfAlive([×tamp](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 | |
| 438 | void 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 | |
| 492 | void 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 | |
| 517 | void BaseTracer::setMakeUserRequestSpanFunc(MakeUserRequestSpanFunc func) { |
| 518 | KJ_ASSERT( |
| 519 | makeUserRequestSpanFunc == kj::none, "setMakeUserRequestSpanFunc can only be called once"); |
| 520 | makeUserRequestSpanFunc = kj::mv(func); |
| 521 | } |
| 522 | |
| 523 | void WorkerTracer::setWorkerAttribute(kj::ConstString key, Span::TagValue value) { |
| 524 | attributes.add(tracing::Attribute{kj::mv(key), kj::mv(value)}); |
| 525 | } |
| 526 | |
| 527 | SpanParent 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 | |
| 536 | void 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 | |
| 559 | kj::Own<SpanObserver> UserSpanObserver::newChild() { |
| 560 | return kj::refcounted<UserSpanObserver>(kj::addRef(*submitter), spanId, traceId, traceFlags); |
| 561 | } |
| 562 | |
| 563 | kj::Own<SpanObserver> UserSpanObserver::newChildFromUserCode() { |
| 564 | return kj::refcounted<UserSpanObserver>( |
| 565 | kj::addRef(*submitter), spanId, traceId, traceFlags, /*fromUserCode=*/true); |
| 566 | } |
| 567 | |
| 568 | kj::Maybe<tracing::SpanContext> UserSpanObserver::toSpanContext() { |
| 569 | if (traceId == nullptr) { |
| 570 | return kj::none; |
| 571 | } |
| 572 | return tracing::SpanContext(traceId, spanId, traceFlags); |
| 573 | } |
| 574 | |
| 575 | void 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 | |
| 584 | void 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. |
| 595 | kj::Date UserSpanObserver::getTime() { |
| 596 | return IoContext::current().now(); |
| 597 | } |
| 598 | |
| 599 | tracing::SpanId UserSpanObserver::getSpanId() { |
| 600 | return spanId; |
| 601 | } |
| 602 | |
| 603 | } // namespace workerd |