// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include #include #include #include #include #include // for capnp::clone() namespace workerd { namespace { // Approximately how much external data we allow in a trace before we start ignoring requests. We // want this number to be big enough to be useful for tracing, but small enough to make it hard to // DoS the C++ heap -- keeping in mind we can record a trace per handler run during a request. For // streaming tail worker, this is the maximum size per tail event. // TODO(streaming-tail): Add a clear indicator for events being truncated based on MAX_TRACE_BYTES // so that developers can understand why this happens. static constexpr size_t MAX_TRACE_BYTES = 256 * 1024; tracing::Attribute::Value cloneAttributeValue(const tracing::Attribute::Value& value) { KJ_SWITCH_ONEOF(value) { KJ_CASE_ONEOF(boolean, bool) { return tracing::Attribute::Value(boolean); } KJ_CASE_ONEOF(number, double) { return tracing::Attribute::Value(number); } KJ_CASE_ONEOF(integer, int64_t) { return tracing::Attribute::Value(integer); } KJ_CASE_ONEOF(string, kj::ConstString) { return tracing::Attribute::Value(string.clone()); } } KJ_UNREACHABLE; } } // namespace kj::Promise> WorkerTracer::onComplete() { KJ_REQUIRE(completeFulfiller == kj::none, "onComplete() can only be called once"); auto paf = kj::newPromiseAndFulfiller>(); completeFulfiller = kj::mv(paf.fulfiller); return kj::mv(paf.promise); } WorkerTracer::WorkerTracer(kj::Maybe> parentPipeline, kj::Own trace, PipelineLogLevel pipelineLogLevel, kj::Maybe> tailAttributes, kj::Maybe> maybeTailStreamWriter) : pipelineLogLevel(pipelineLogLevel), trace(kj::mv(trace)), parentPipeline(kj::mv(parentPipeline)), maybeTailStreamWriter(kj::mv(maybeTailStreamWriter)) { KJ_IF_SOME(tags, tailAttributes) { if (tags.size() == 0) { tailAttributes = kj::none; } else { for (auto& tag: tags) { KJ_REQUIRE(tag.value.size() == 1, "tail attributes must contain exactly one value"); setWorkerAttribute(tag.name.clone(), cloneAttributeValue(tag.value[0])); } } } this->trace->tailAttributes = kj::mv(tailAttributes); } WorkerTracer::~WorkerTracer() noexcept(false) { // Report the outcome event, which should have been delivered by now. // Do not attempt to report an outcome event if logging is disabled, as with other event types. if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } // Report the outcome event if STWs are present. All worker events need to call setEventInfo at // the start of the invocation to submit the onset event before any other tail events. KJ_IF_SOME(writer, maybeTailStreamWriter) { KJ_IF_SOME(spanContext, topLevelInvocationSpanContext) { if (markedUnused) { LOG_WARNING_PERIODICALLY("WorkerTracer was marked unused but actually was used"); } if (isPredictableModeForTest()) { writer->report(spanContext, tracing::Outcome(trace->outcome, 0 * kj::MILLISECONDS, 0 * kj::MILLISECONDS), completeTime, 0); } else { writer->report(spanContext, tracing::Outcome(trace->outcome, trace->cpuTime, trace->wallTime), completeTime, 0); } } else if (!markedUnused) { // If no span context is available, we have a streaming tail worker set up but shut down the // worker tracer without ever sending an Onset event. In that case we either failed to set up // the Onset properly (indicating a bug – all event types are required to report an Onset at // the start – although this is more likely to manifest as a "Tail stream onset was not // reported" error) or we created a WorkerInterface with WorkerTracer without ever invoking it // (which is not incorrect behavior, but likely indicates inefficient code that sets up // WorkerInterfaces and then ends up not using it due to an error/incorrect parameters; such // error checking should be done beforehand to avoid unused allocations). Report such cases. // Note: If markedUnused is true, this tracer was intentionally not used (e.g., duplicate // alarm request deduplication) and the warning should be suppressed. LOG_ERROR_PERIODICALLY( "destructed WorkerTracer with STW without reporting Onset event", kj::getStackTrace()); } } // Report the completed trace, if fulfiller is set up. KJ_IF_SOME(f, completeFulfiller) { f.get()->fulfill(kj::mv(trace)); } }; constexpr kj::LiteralStringConst logSizeExceeded = "[\"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; void WorkerTracer::addLog(const tracing::InvocationSpanContext& context, kj::Date timestamp, LogLevel logLevel, kj::String message) { if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } // TODO(streaming-tail): Here we add the log to the trace object and the tail stream writer, if // available. If the given worker stage is only tailed by a streaming tail worker, adding the log // to the buffered trace object is not needed; this will be addressed in a future refactor. KJ_IF_SOME(writer, maybeTailStreamWriter) { // If message is too big on its own, truncate it. size_t messageSize = kj::min(message.size(), MAX_TRACE_BYTES); writer->report(context, {tracing::Log(timestamp, logLevel, kj::str(message.first(messageSize)))}, timestamp, messageSize); } if (trace->exceededLogLimit) { return; } size_t messageSize = sizeof(tracing::Log) + message.size(); if (trace->bytesUsed + messageSize > MAX_TRACE_BYTES) { // We use a JSON encoded array/string to match other console.log() recordings: trace->logs.add(timestamp, LogLevel::WARN, kj::str(logSizeExceeded)); trace->exceededLogLimit = true; trace->truncated = true; } else { trace->bytesUsed += messageSize; trace->logs.add(timestamp, logLevel, kj::mv(message)); } } void WorkerTracer::addSpanOpen(tracing::SpanId spanId, tracing::SpanId parentSpanId, kj::ConstString operationName, kj::Date startTime) { if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } auto& tailStreamWriter = KJ_UNWRAP_OR_RETURN(maybeTailStreamWriter); auto& topLevelContext = KJ_ASSERT_NONNULL(topLevelInvocationSpanContext); // Compose SpanOpen. An all-zero spanId is interpreted as having no spans above this one, thus we // use the Onset spanId instead (taken from topLevelContext). We go to great lengths to rule out // getting an all-zero spanId by chance (see SpanId::fromEntropy()), so this should be safe. if (parentSpanId == tracing::SpanId::nullId) { parentSpanId = topLevelContext.getSpanId(); } size_t spanNameSize = operationName.size(); auto spanOpenContext = tracing::InvocationSpanContext(topLevelContext.getTraceId(), topLevelContext.getInvocationId(), parentSpanId, topLevelContext.getTraceFlags()); tailStreamWriter->report( spanOpenContext, tracing::SpanOpen(spanId, kj::mv(operationName)), startTime, spanNameSize); } void WorkerTracer::addSpanClose(tracing::SpanEndData&& span, kj::Maybe maybeStartTime) { if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } // Note: spans are not available in the buffered tail worker, so we don't need an exceededSpanLimit // variable for it and it can't cause truncation. auto& tailStreamWriter = KJ_UNWRAP_OR_RETURN(maybeTailStreamWriter); adjustSpanTime(span, maybeStartTime); size_t spanTagsSize = 0; for (const Span::TagMap::Entry& tag: span.tags) { spanTagsSize += tag.key.size(); KJ_SWITCH_ONEOF(tag.value) { KJ_CASE_ONEOF(str, kj::ConstString) { spanTagsSize += str.size(); } KJ_CASE_ONEOF(val, bool) { spanTagsSize++; } // int64_t and double KJ_CASE_ONEOF_DEFAULT { spanTagsSize += sizeof(int64_t); } } } // Compose Attributes and SpanClose, which are available at span completion time and transmitted // together. auto& topLevelContext = KJ_ASSERT_NONNULL(topLevelInvocationSpanContext); auto spanComponentContext = tracing::InvocationSpanContext(topLevelContext.getTraceId(), topLevelContext.getInvocationId(), span.spanId, topLevelContext.getTraceFlags()); if (span.tags.size() && spanTagsSize <= MAX_TRACE_BYTES) { tracing::CustomInfo attr = KJ_MAP(tag, span.tags) { return tracing::Attribute(kj::mv(tag.key), kj::mv(tag.value)); }; tailStreamWriter->report(spanComponentContext, kj::mv(attr), span.endTime, spanTagsSize); } tailStreamWriter->report(spanComponentContext, tracing::SpanClose(), span.endTime, 0); } void WorkerTracer::addException(const tracing::InvocationSpanContext& context, kj::Date timestamp, kj::String name, kj::String message, kj::Maybe stack) { // TODO(someday): For now, we're using logLevel == none as a hint to avoid doing anything // expensive while tracing. We may eventually want separate configuration for exceptions vs. // logs. if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } size_t messageSize = sizeof(tracing::Exception) + name.size() + message.size(); KJ_IF_SOME(s, stack) { messageSize += s.size(); } KJ_IF_SOME(writer, maybeTailStreamWriter) { auto maybeTruncatedName = name.first(kj::min(name.size(), MAX_TRACE_BYTES)); auto maybeTruncatedMessage = message.first(kj::min(message.size(), MAX_TRACE_BYTES - maybeTruncatedName.size())); kj::Maybe maybeTruncatedStack; auto maybeTruncatedStackSize = 0; KJ_IF_SOME(s, stack) { maybeTruncatedStackSize = kj::min( s.size(), MAX_TRACE_BYTES - maybeTruncatedName.size() - maybeTruncatedMessage.size()); maybeTruncatedStack = kj::heapString(s.first(maybeTruncatedStackSize)); } writer->report(context, {tracing::Exception(timestamp, kj::str(maybeTruncatedName), kj::str(maybeTruncatedMessage), kj::mv(maybeTruncatedStack))}, timestamp, maybeTruncatedName.size() + maybeTruncatedMessage.size() + maybeTruncatedStackSize); } if (trace->exceededExceptionLimit) { return; } if (trace->bytesUsed + messageSize > MAX_TRACE_BYTES) { trace->exceededExceptionLimit = true; trace->truncated = true; trace->exceptions.add(timestamp, kj::str("Error"), kj::str("Trace resource limit exceeded; subsequent exceptions not recorded."), kj::none); } else { trace->bytesUsed += messageSize; trace->exceptions.add(timestamp, kj::mv(name), kj::mv(message), kj::mv(stack)); } } void WorkerTracer::addDiagnosticChannelEvent(const tracing::InvocationSpanContext& context, kj::Date timestamp, kj::String channel, kj::Array message) { if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } size_t messageSize = sizeof(tracing::DiagnosticChannelEvent) + channel.size() + message.size(); KJ_IF_SOME(writer, maybeTailStreamWriter) { // Drop oversized diagnostic channel events instead of truncating them – a truncated message may // not be deserialized correctly. if (messageSize <= MAX_TRACE_BYTES) { writer->report(context, {tracing::DiagnosticChannelEvent( timestamp, kj::str(channel), kj::heapArray(message))}, timestamp, messageSize); } } if (trace->exceededDiagnosticChannelEventLimit) { return; } if (trace->bytesUsed + messageSize > MAX_TRACE_BYTES) { trace->exceededDiagnosticChannelEventLimit = true; trace->truncated = true; trace->diagnosticChannelEvents.add( timestamp, kj::str("workerd.LimitExceeded"), kj::Array()); } else { trace->bytesUsed += messageSize; trace->diagnosticChannelEvents.add(timestamp, kj::mv(channel), kj::mv(message)); } } void WorkerTracer::setEventInfo( IoContext::IncomingRequest& incomingRequest, tracing::EventInfo&& info) { // IoContext is available at this time, capture weakRef. KJ_ASSERT(weakIoContext == kj::none, "tracer can only be used for a single event"); weakIoContext = incomingRequest.getContext().getWeakRef(); setEventInfoInternal( incomingRequest.getInvocationSpanContext(), incomingRequest.now(), kj::mv(info)); } void WorkerTracer::setEventInfoInternal( const tracing::InvocationSpanContext& context, kj::Date timestamp, tracing::EventInfo&& info) { KJ_ASSERT(trace->eventInfo == kj::none, "tracer can only be used for a single event"); // TODO(someday): For now, we're using logLevel == none as a hint to avoid doing anything // expensive while tracing. We may eventually want separate configuration for event info vs. // logs. // TODO(perf): Find a way to allow caller to avoid the cost of generation if the info struct // won't be used? if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } trace->eventTimestamp = timestamp; this->topLevelInvocationSpanContext = context.clone(); size_t eventSize = 0; KJ_SWITCH_ONEOF(info) { KJ_CASE_ONEOF(fetch, tracing::FetchEventInfo) { eventSize += fetch.url.size(); for (const auto& header: fetch.headers) { eventSize += header.name.size() + header.value.size(); } eventSize += fetch.cfJson.size(); // Limit STW onset to MAX_TRACE_BYTES, beyond that dispatch a truncated event too. if (eventSize > MAX_TRACE_BYTES) { info = tracing::FetchEventInfo(fetch.method, {}, {}, {}); } } KJ_CASE_ONEOF_DEFAULT {} } KJ_IF_SOME(writer, maybeTailStreamWriter) { // Provide WorkerInfo to the streaming tail worker if available. This data is provided when the // WorkerTracer is created, but the actual onset event is the best time to send it. auto workerInfo = tracing::Onset::WorkerInfo{ .executionModel = trace->executionModel, .scriptName = mapCopyString(trace->scriptName), .scriptVersion = trace->scriptVersion.map([](auto& scriptVersion) -> kj::Own { return capnp::clone(*scriptVersion); }), .preview = trace->preview.map([](auto& preview) { return preview.clone(); }), .dispatchNamespace = mapCopyString(trace->dispatchNamespace), .scriptId = mapCopyString(trace->scriptId), .scriptTags = KJ_MAP(tag, trace->scriptTags) { return kj::str(tag); }, .entrypoint = mapCopyString(trace->entrypoint), }; tracing::SpanId parentSpanId = tracing::SpanId::nullId; KJ_IF_SOME(trigger, context.getParent()) { parentSpanId = trigger.getSpanId(); } // Onset needs special handling for spanId: The top-level spanId is zero unless a trigger // context is available. The inner spanId is taken from the invocation // span context, that span is being "opened" with the onset event. All other tail events have it // as its parent span ID, except for recursive SpanOpens (which have the parent span instead) // and Attribute/SpanClose events (which have the spanId opened in the corresponding SpanOpen). auto onsetContext = tracing::InvocationSpanContext( context.getTraceId(), context.getInvocationId(), parentSpanId, context.getTraceFlags()); // Not applying size accounting for Onset since it is sent separately writer->report(onsetContext, tracing::Onset(context.getSpanId(), cloneEventInfo(info), kj::mv(workerInfo), attributes.releaseAsArray()), timestamp, 0); } // truncation should only be needed for fetch events, since we only set eventSize there. if (trace->bytesUsed + eventSize > MAX_TRACE_BYTES && eventSize > 0) { trace->truncated = true; trace->logs.add(timestamp, LogLevel::WARN, kj::str("[\"Trace resource limit exceeded; could not capture event info.\"]")); trace->eventInfo = tracing::FetchEventInfo(info.get().method, {}, {}, {}); } else { trace->bytesUsed += eventSize; trace->eventInfo = kj::mv(info); } } void WorkerTracer::setOutcome(EventOutcome outcome, kj::Duration cpuTime, kj::Duration wallTime) { trace->outcome = outcome; trace->cpuTime = cpuTime; trace->wallTime = wallTime; // Defer reporting the actual outcome event to the WorkerTracer destructor: The outcome is // reported when the metrics request is deallocated, but with ctx.waitUntil() there might be spans // continuing to exist beyond that point. By the time the WorkerTracer is deallocated, the // IoContext and its task set will be done and any additional spans will have wrapped up. // This is somewhat at odds with the concept of "streaming" events, but benign as the WorkerTracer // wraps up right after the metrics request object in the average case and since the outcome has a // fixed size. } void WorkerTracer::recordTimestamp(kj::Date timestamp) { if (completeTime == kj::UNIX_EPOCH) { completeTime = timestamp; } } kj::Date BaseTracer::getTime() { auto& weakIoCtx = KJ_ASSERT_NONNULL(weakIoContext); kj::Date timestamp = kj::UNIX_EPOCH; weakIoCtx->runIfAlive([×tamp](IoContext& context) { timestamp = context.now(); }); if (!weakIoCtx->isValid()) { // This can happen if we the IoContext gets destroyed following an exception, but we still need // to report a time for the return event. if (completeTime != kj::UNIX_EPOCH) { timestamp = completeTime; } else { // Otherwise, we can't actually get an end timestamp that makes sense. if (isPredictableModeForTest()) { KJ_FAIL_ASSERT("reported return event without valid IoContext or completeTime"); } else { LOG_WARNING_PERIODICALLY("reported return event without valid IoContext or completeTime"); } } } return timestamp; } void BaseTracer::adjustSpanTime(tracing::SpanEndData& span, kj::Maybe maybeStartTime) { // To report I/O time, we need the IOContext to still be alive. // weakIoContext is only none if we are tracing via RPC (in this case span times have already been // adjusted) or if we failed to transmit an Onset event (in that case we'll get an error based on // missing topLevelInvocationSpanContext right after). if (weakIoContext != kj::none) { auto& weakIoCtx = KJ_ASSERT_NONNULL(weakIoContext); // startTime is generally available when we are not tracing via RPC, so we can assert that it is // present. For the RPC case, the adjustment will already have been done earlier and it's ok // for maybeStartTime to be none as this code won't run based on weakIoContext being none. kj::Date startTime = KJ_ASSERT_NONNULL(maybeStartTime); weakIoCtx->runIfAlive([this, &span, &startTime](IoContext& context) { if (context.hasCurrentIncomingRequest()) { span.endTime = context.now(); } else { // We have an IOContext, but there's no current IncomingRequest. Always log a warning here, // this should not be happening. Still report completeTime as a useful timestamp if // available. bool hasCompleteTime = false; if (completeTime != kj::UNIX_EPOCH) { span.endTime = completeTime; hasCompleteTime = true; } else { span.endTime = startTime; } if (isPredictableModeForTest()) { KJ_FAIL_ASSERT("reported span without current request", hasCompleteTime); } else { LOG_WARNING_PERIODICALLY("reported span without current request"); } } }); if (!weakIoCtx->isValid()) { // This can happen if we start a customEvent from this event and cancel it after this IoContext // gets destroyed. In that case we no longer have an IoContext available and can't get the // current time, but the outcome timestamp will have already been set. Since the outcome // timestamp is "late enough", simply use that. // TODO(o11y): fix this – spans should not be outliving the IoContext. if (completeTime != kj::UNIX_EPOCH) { span.endTime = completeTime; } else { // Otherwise, we can't actually get an end timestamp that makes sense. Report a zero-duration // span and log a warning (or fail assert in test mode). span.endTime = startTime; if (isPredictableModeForTest()) { KJ_FAIL_ASSERT("reported span after IoContext was deallocated"); } else { KJ_LOG(WARNING, "reported span after IoContext was deallocated"); } } } } } void WorkerTracer::setReturn( kj::Maybe timestamp, kj::Maybe fetchResponseInfo) { // Match the behavior of setEventInfo(). Any resolution of the TODO comments in setEventInfo() // that are related to this check will probably also affect this function. if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } KJ_IF_SOME(writer, maybeTailStreamWriter) { auto& spanContext = KJ_UNWRAP_OR_RETURN(topLevelInvocationSpanContext); // Fall back to weak IoContext if no timestamp is available writer->report(spanContext, tracing::Return({fetchResponseInfo.map([](auto& info) { return info.clone(); })}), timestamp.orDefault([&]() { return getTime(); }), 0); } // Add fetch response info for buffered tail worker KJ_IF_SOME(info, fetchResponseInfo) { KJ_REQUIRE(KJ_REQUIRE_NONNULL(trace->eventInfo).is()); KJ_ASSERT(trace->fetchResponseInfo == kj::none, "setFetchResponseInfo can only be called once"); trace->fetchResponseInfo = kj::mv(info); } } void BaseTracer::setMakeUserRequestSpanFunc(MakeUserRequestSpanFunc func) { KJ_ASSERT( makeUserRequestSpanFunc == kj::none, "setMakeUserRequestSpanFunc can only be called once"); makeUserRequestSpanFunc = kj::mv(func); } void WorkerTracer::setWorkerAttribute(kj::ConstString key, Span::TagValue value) { attributes.add(tracing::Attribute{kj::mv(key), kj::mv(value)}); } SpanParent BaseTracer::makeUserRequestSpan( tracing::TraceId traceId, kj::Maybe traceFlags) { KJ_IF_SOME(func, makeUserRequestSpanFunc) { return func(kj::mv(traceId), traceFlags); } else { return SpanParent(nullptr); } } void WorkerTracer::setJsRpcInfo(const tracing::InvocationSpanContext& context, kj::Date timestamp, const kj::ConstString& methodName) { if (pipelineLogLevel == PipelineLogLevel::NONE) { return; } // Update the method name in the already-set JsRpcEventInfo for buffered tail worker compatibility KJ_IF_SOME(info, trace->eventInfo) { KJ_SWITCH_ONEOF(info) { KJ_CASE_ONEOF(jsRpcInfo, tracing::JsRpcEventInfo) { jsRpcInfo.methodName = kj::str(methodName); } KJ_CASE_ONEOF_DEFAULT {} } } KJ_IF_SOME(writer, maybeTailStreamWriter) { auto tag = tracing::Attribute("jsrpc.method"_kjc, methodName.clone()); writer->report(context, kj::arr(kj::mv(tag)), timestamp, methodName.size()); } } kj::Own UserSpanObserver::newChild() { return kj::refcounted(kj::addRef(*submitter), spanId, traceId, traceFlags); } kj::Own UserSpanObserver::newChildFromUserCode() { return kj::refcounted( kj::addRef(*submitter), spanId, traceId, traceFlags, /*fromUserCode=*/true); } kj::Maybe UserSpanObserver::toSpanContext() { if (traceId == nullptr) { return kj::none; } return tracing::SpanContext(traceId, spanId, traceFlags); } void UserSpanObserver::onClose( kj::Date endTime, Span::TagMap&& tags, kj::Vector&& logs) { // span logs are not supported in user tracing. (void)logs; if (wasAccepted) { submitter->submitSpanClose(spanId, startTime, endTime, kj::mv(tags)); } } void UserSpanObserver::onOpen(kj::ConstString operationName, kj::Date startTime) { this->startTime = startTime; if (fromUserCode) { wasAccepted = submitter->submitUserSpanOpen(spanId, parentSpanId, kj::mv(operationName), startTime); } else { wasAccepted = submitter->submitSpanOpen(spanId, parentSpanId, kj::mv(operationName), startTime); } } // Provide I/O time to the tracing system for user spans. kj::Date UserSpanObserver::getTime() { return IoContext::current().now(); } tracing::SpanId UserSpanObserver::getSpanId() { return spanId; } } // namespace workerd