File
Blob: src/workerd/io/trace.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/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 | |
| 18 | namespace workerd { |
| 19 | |
| 20 | namespace tracing { |
| 21 | namespace { |
| 22 | kj::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 | |
| 34 | kj::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 | |
| 47 | void 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 | |
| 55 | void 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 |
| 64 | kj::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 |
| 83 | kj::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 |
| 98 | kj::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 |
| 114 | kj::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 |
| 122 | kj::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 | |
| 129 | namespace { |
| 130 | uint64_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 | |
| 149 | TraceId 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 | |
| 164 | kj::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 | |
| 171 | SpanId SpanId::fromEntropy(kj::Maybe<kj::EntropySource&> entropySource) { |
| 172 | return SpanId(getRandom64Bit(entropySource)); |
| 173 | } |
| 174 | |
| 175 | kj::String KJ_STRINGIFY(const SpanId& id) { |
| 176 | return id; |
| 177 | } |
| 178 | |
| 179 | kj::String KJ_STRINGIFY(const TraceId& id) { |
| 180 | return id; |
| 181 | } |
| 182 | |
| 183 | InvocationSpanContext::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 | |
| 199 | InvocationSpanContext 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 | |
| 207 | InvocationSpanContext 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 | |
| 223 | TraceId TraceId::fromCapnp(rpc::TraceId::Reader reader) { |
| 224 | return TraceId(reader.getLow(), reader.getHigh()); |
| 225 | } |
| 226 | |
| 227 | void TraceId::toCapnp(rpc::TraceId::Builder writer) const { |
| 228 | writer.setLow(low); |
| 229 | writer.setHigh(high); |
| 230 | } |
| 231 | |
| 232 | kj::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 | |
| 253 | void 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 | |
| 262 | InvocationSpanContext 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 | |
| 271 | kj::String KJ_STRINGIFY(const InvocationSpanContext& context) { |
| 272 | return kj::str(context.getTraceId(), "-", context.getInvocationId(), "-", context.getSpanId()); |
| 273 | } |
| 274 | |
| 275 | kj::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 | |
| 311 | kj::String KJ_STRINGIFY(const CustomInfo& customInfo) { |
| 312 | return kj::str( |
| 313 | "CustomInfo: ", kj::strArray(KJ_MAP(attr, customInfo) { return kj::str(attr); }, ", ")); |
| 314 | } |
| 315 | |
| 316 | SpanContext 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 | |
| 331 | void 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 | |
| 342 | kj::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 | |
| 379 | kj::String KJ_STRINGIFY(const SpanContext& context) { |
| 380 | return kj::str(context.getTraceId(), "-", context.getSpanId()); |
| 381 | } |
| 382 | |
| 383 | namespace { |
| 384 | |
| 385 | static 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 | |
| 392 | ConnectEventInfo::ConnectEventInfo() {} |
| 393 | |
| 394 | ConnectEventInfo::ConnectEventInfo(rpc::Trace::ConnectEventInfo::Reader reader) {} |
| 395 | |
| 396 | void ConnectEventInfo::copyTo(rpc::Trace::ConnectEventInfo::Builder builder) const {} |
| 397 | |
| 398 | ConnectEventInfo ConnectEventInfo::clone() const { |
| 399 | return ConnectEventInfo(); |
| 400 | } |
| 401 | |
| 402 | FetchEventInfo::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 | |
| 409 | FetchEventInfo::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 | |
| 418 | void 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 | |
| 429 | FetchEventInfo FetchEventInfo::clone() const { |
| 430 | return FetchEventInfo( |
| 431 | method, kj::str(url), kj::str(cfJson), KJ_MAP(h, headers) { return h.clone(); }); |
| 432 | } |
| 433 | |
| 434 | kj::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 | |
| 440 | FetchEventInfo::Header::Header(kj::String name, kj::String value) |
| 441 | : name(kj::mv(name)), |
| 442 | value(kj::mv(value)) {} |
| 443 | |
| 444 | FetchEventInfo::Header::Header(rpc::Trace::FetchEventInfo::Header::Reader reader) |
| 445 | : name(kj::str(reader.getName())), |
| 446 | value(kj::str(reader.getValue())) {} |
| 447 | |
| 448 | void FetchEventInfo::Header::copyTo(rpc::Trace::FetchEventInfo::Header::Builder builder) const { |
| 449 | builder.setName(name); |
| 450 | builder.setValue(value); |
| 451 | } |
| 452 | |
| 453 | FetchEventInfo::Header FetchEventInfo::Header::clone() const { |
| 454 | return Header(kj::str(name), kj::str(value)); |
| 455 | } |
| 456 | |
| 457 | kj::String FetchEventInfo::Header::toString() const { |
| 458 | return kj::str("FetchEventInfo::Header: ", name, ", ", value); |
| 459 | } |
| 460 | |
| 461 | JsRpcEventInfo::JsRpcEventInfo(kj::String methodName): methodName(kj::mv(methodName)) {} |
| 462 | |
| 463 | JsRpcEventInfo::JsRpcEventInfo(rpc::Trace::JsRpcEventInfo::Reader reader) |
| 464 | : methodName(kj::str(reader.getMethodName())) {} |
| 465 | |
| 466 | void JsRpcEventInfo::copyTo(rpc::Trace::JsRpcEventInfo::Builder builder) const { |
| 467 | builder.setMethodName(methodName); |
| 468 | } |
| 469 | |
| 470 | JsRpcEventInfo JsRpcEventInfo::clone() const { |
| 471 | return JsRpcEventInfo(kj::str(methodName)); |
| 472 | } |
| 473 | |
| 474 | kj::String JsRpcEventInfo::toString() const { |
| 475 | return kj::str("JsRpcEventInfo: ", methodName); |
| 476 | } |
| 477 | |
| 478 | ScheduledEventInfo::ScheduledEventInfo(double scheduledTime, kj::String cron) |
| 479 | : scheduledTime(scheduledTime), |
| 480 | cron(kj::mv(cron)) {} |
| 481 | |
| 482 | ScheduledEventInfo::ScheduledEventInfo(rpc::Trace::ScheduledEventInfo::Reader reader) |
| 483 | : scheduledTime(reader.getScheduledTime()), |
| 484 | cron(kj::str(reader.getCron())) {} |
| 485 | |
| 486 | void ScheduledEventInfo::copyTo(rpc::Trace::ScheduledEventInfo::Builder builder) const { |
| 487 | builder.setScheduledTime(scheduledTime); |
| 488 | builder.setCron(cron); |
| 489 | } |
| 490 | |
| 491 | ScheduledEventInfo ScheduledEventInfo::clone() const { |
| 492 | return ScheduledEventInfo(scheduledTime, kj::str(cron)); |
| 493 | } |
| 494 | |
| 495 | AlarmEventInfo::AlarmEventInfo(kj::Date scheduledTime): scheduledTime(scheduledTime) {} |
| 496 | |
| 497 | AlarmEventInfo::AlarmEventInfo(rpc::Trace::AlarmEventInfo::Reader reader) |
| 498 | : scheduledTime(reader.getScheduledTimeMs() * kj::MILLISECONDS + kj::UNIX_EPOCH) {} |
| 499 | |
| 500 | void AlarmEventInfo::copyTo(rpc::Trace::AlarmEventInfo::Builder builder) const { |
| 501 | builder.setScheduledTimeMs((scheduledTime - kj::UNIX_EPOCH) / kj::MILLISECONDS); |
| 502 | } |
| 503 | |
| 504 | AlarmEventInfo AlarmEventInfo::clone() const { |
| 505 | return AlarmEventInfo(scheduledTime); |
| 506 | } |
| 507 | |
| 508 | QueueEventInfo::QueueEventInfo(kj::String queueName, uint32_t batchSize) |
| 509 | : queueName(kj::mv(queueName)), |
| 510 | batchSize(batchSize) {} |
| 511 | |
| 512 | QueueEventInfo::QueueEventInfo(rpc::Trace::QueueEventInfo::Reader reader) |
| 513 | : queueName(kj::heapString(reader.getQueueName())), |
| 514 | batchSize(reader.getBatchSize()) {} |
| 515 | |
| 516 | void QueueEventInfo::copyTo(rpc::Trace::QueueEventInfo::Builder builder) const { |
| 517 | builder.setQueueName(queueName); |
| 518 | builder.setBatchSize(batchSize); |
| 519 | } |
| 520 | |
| 521 | QueueEventInfo QueueEventInfo::clone() const { |
| 522 | return QueueEventInfo(kj::str(queueName), batchSize); |
| 523 | } |
| 524 | |
| 525 | EmailEventInfo::EmailEventInfo(kj::String mailFrom, kj::String rcptTo, uint32_t rawSize) |
| 526 | : mailFrom(kj::mv(mailFrom)), |
| 527 | rcptTo(kj::mv(rcptTo)), |
| 528 | rawSize(rawSize) {} |
| 529 | |
| 530 | EmailEventInfo::EmailEventInfo(rpc::Trace::EmailEventInfo::Reader reader) |
| 531 | : mailFrom(kj::heapString(reader.getMailFrom())), |
| 532 | rcptTo(kj::heapString(reader.getRcptTo())), |
| 533 | rawSize(reader.getRawSize()) {} |
| 534 | |
| 535 | void EmailEventInfo::copyTo(rpc::Trace::EmailEventInfo::Builder builder) const { |
| 536 | builder.setMailFrom(mailFrom); |
| 537 | builder.setRcptTo(rcptTo); |
| 538 | builder.setRawSize(rawSize); |
| 539 | } |
| 540 | |
| 541 | EmailEventInfo EmailEventInfo::clone() const { |
| 542 | return EmailEventInfo(kj::str(mailFrom), kj::str(rcptTo), rawSize); |
| 543 | } |
| 544 | |
| 545 | namespace { |
| 546 | kj::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 | |
| 551 | kj::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 | |
| 557 | TraceEventInfo::TraceEventInfo(kj::ArrayPtr<const kj::Own<Trace>> traces) |
| 558 | : traces(getTraceItemsFromTraces(traces)) {} |
| 559 | |
| 560 | TraceEventInfo::TraceEventInfo(rpc::Trace::TraceEventInfo::Reader reader) |
| 561 | : traces(getTraceItemsFromReader(reader)) {} |
| 562 | |
| 563 | void 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 | |
| 570 | TraceEventInfo TraceEventInfo::clone() const { |
| 571 | return TraceEventInfo(KJ_MAP(item, traces) { return item.clone(); }); |
| 572 | } |
| 573 | |
| 574 | TracePreview::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 | |
| 579 | TracePreview::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 | |
| 584 | void TracePreview::copyTo(rpc::Trace::TracePreviewInfo::Builder builder) const { |
| 585 | builder.setId(id); |
| 586 | builder.setSlug(slug); |
| 587 | builder.setName(name); |
| 588 | } |
| 589 | |
| 590 | TracePreview TracePreview::clone() const { |
| 591 | return TracePreview(kj::str(id), kj::str(slug), kj::str(name)); |
| 592 | } |
| 593 | |
| 594 | TraceEventInfo::TraceItem::TraceItem(kj::Maybe<kj::String> scriptName) |
| 595 | : scriptName(kj::mv(scriptName)) {} |
| 596 | |
| 597 | TraceEventInfo::TraceItem::TraceItem(rpc::Trace::TraceEventInfo::TraceItem::Reader reader) |
| 598 | : scriptName(kj::str(reader.getScriptName())) {} |
| 599 | |
| 600 | void TraceEventInfo::TraceItem::copyTo( |
| 601 | rpc::Trace::TraceEventInfo::TraceItem::Builder builder) const { |
| 602 | KJ_IF_SOME(name, scriptName) { |
| 603 | builder.setScriptName(name); |
| 604 | } |
| 605 | } |
| 606 | |
| 607 | TraceEventInfo::TraceItem TraceEventInfo::TraceItem::clone() const { |
| 608 | return TraceItem(mapCopyString(scriptName)); |
| 609 | } |
| 610 | |
| 611 | DiagnosticChannelEvent::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 | |
| 617 | DiagnosticChannelEvent::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 | |
| 622 | void 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 | |
| 628 | DiagnosticChannelEvent DiagnosticChannelEvent::clone() const { |
| 629 | return DiagnosticChannelEvent(timestamp, kj::str(channel), kj::heapArray<kj::byte>(message)); |
| 630 | } |
| 631 | |
| 632 | StreamDiagnosticsEvent::StreamDiagnosticsEvent(uint32_t droppedEventsCount) |
| 633 | : droppedEventsCount(droppedEventsCount) {} |
| 634 | |
| 635 | StreamDiagnosticsEvent::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 | |
| 649 | void 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 | |
| 656 | StreamDiagnosticsEvent StreamDiagnosticsEvent::clone() const { |
| 657 | return StreamDiagnosticsEvent(droppedEventsCount); |
| 658 | } |
| 659 | |
| 660 | HibernatableWebSocketEventInfo::HibernatableWebSocketEventInfo(Type type): type(type) {} |
| 661 | |
| 662 | HibernatableWebSocketEventInfo::HibernatableWebSocketEventInfo( |
| 663 | rpc::Trace::HibernatableWebSocketEventInfo::Reader reader) |
| 664 | : type(readFrom(reader)) {} |
| 665 | |
| 666 | void 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 | |
| 684 | HibernatableWebSocketEventInfo 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 | |
| 702 | HibernatableWebSocketEventInfo::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 | |
| 722 | FetchResponseInfo::FetchResponseInfo(uint16_t statusCode): statusCode(statusCode) {} |
| 723 | |
| 724 | FetchResponseInfo::FetchResponseInfo(rpc::Trace::FetchResponseInfo::Reader reader) |
| 725 | : statusCode(reader.getStatusCode()) {} |
| 726 | |
| 727 | void FetchResponseInfo::copyTo(rpc::Trace::FetchResponseInfo::Builder builder) const { |
| 728 | builder.setStatusCode(statusCode); |
| 729 | } |
| 730 | |
| 731 | FetchResponseInfo FetchResponseInfo::clone() const { |
| 732 | return FetchResponseInfo(statusCode); |
| 733 | } |
| 734 | |
| 735 | Log::Log(kj::Date timestamp, LogLevel logLevel, kj::String message) |
| 736 | : timestamp(timestamp), |
| 737 | logLevel(logLevel), |
| 738 | message(kj::mv(message)) {} |
| 739 | |
| 740 | void 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 | |
| 746 | Log Log::clone() const { |
| 747 | return Log(timestamp, logLevel, kj::str(message)); |
| 748 | } |
| 749 | |
| 750 | Exception::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 | |
| 757 | Log::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 | |
| 762 | Exception::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 | |
| 771 | void 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 | |
| 780 | Exception Exception::clone() const { |
| 781 | return Exception(timestamp, kj::str(name), kj::str(message), mapCopyString(stack)); |
| 782 | } |
| 783 | } // namespace tracing |
| 784 | |
| 785 | Trace::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) {} |
| 805 | Trace::Trace(rpc::Trace::Reader reader) { |
| 806 | mergeFrom(reader, PipelineLogLevel::FULL); |
| 807 | } |
| 808 | |
| 809 | Trace::~Trace() noexcept(false) {} |
| 810 | |
| 811 | void 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 | |
| 932 | void 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 | |
| 1036 | namespace tracing { |
| 1037 | |
| 1038 | Attribute::Attribute(kj::ConstString name, Value&& value) |
| 1039 | : name(kj::mv(name)), |
| 1040 | value(kj::arr(kj::mv(value))) {} |
| 1041 | |
| 1042 | Attribute::Attribute(kj::ConstString name, Values&& value) |
| 1043 | : name(kj::mv(name)), |
| 1044 | value(kj::mv(value)) {} |
| 1045 | |
| 1046 | namespace { |
| 1047 | kj::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 | |
| 1054 | kj::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 | |
| 1067 | Attribute::Attribute(rpc::Trace::Attribute::Reader reader) |
| 1068 | : name(kj::str(reader.getName())), |
| 1069 | value(readValues(reader)) {} |
| 1070 | |
| 1071 | void 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 | |
| 1079 | Attribute Attribute::clone() const { |
| 1080 | return Attribute(name.clone(), KJ_MAP(v, value) { return spanTagClone(v); }); |
| 1081 | } |
| 1082 | |
| 1083 | kj::String Attribute::toString() const { |
| 1084 | return kj::str("Attribute: ", name, ", ", value); |
| 1085 | } |
| 1086 | |
| 1087 | Return::Return(kj::Maybe<FetchResponseInfo> info): info(kj::mv(info)) {} |
| 1088 | Return::Return(rpc::Trace::Return::Reader reader): info(readReturnInfo(reader)) {} |
| 1089 | |
| 1090 | void 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 | |
| 1097 | Return Return::clone() const { |
| 1098 | KJ_IF_SOME(fetchInfo, info) { |
| 1099 | return Return(kj::Maybe(fetchInfo.clone())); |
| 1100 | } |
| 1101 | return Return(); |
| 1102 | } |
| 1103 | |
| 1104 | SpanOpen::SpanOpen(SpanId spanId, kj::ConstString operationName, kj::Maybe<Info> info) |
| 1105 | : operationName(kj::mv(operationName)), |
| 1106 | info(kj::mv(info)), |
| 1107 | spanId(spanId) {} |
| 1108 | |
| 1109 | namespace { |
| 1110 | kj::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 | |
| 1130 | SpanOpen::SpanOpen(rpc::Trace::SpanOpen::Reader reader) |
| 1131 | : operationName(kj::str(reader.getOperationName())), |
| 1132 | info(readSpanOpenInfo(reader)), |
| 1133 | spanId(reader.getSpanId()) {} |
| 1134 | |
| 1135 | void 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 | |
| 1157 | SpanOpen 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 | |
| 1177 | kj::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 | |
| 1192 | kj::String SpanOpen::toString() const { |
| 1193 | return kj::str("SpanOpen:", operationName, ", ", info); |
| 1194 | } |
| 1195 | |
| 1196 | SpanClose::SpanClose(EventOutcome outcome): outcome(outcome) {} |
| 1197 | |
| 1198 | SpanClose::SpanClose(rpc::Trace::SpanClose::Reader reader): outcome(reader.getOutcome()) {} |
| 1199 | |
| 1200 | void SpanClose::copyTo(rpc::Trace::SpanClose::Builder builder) const { |
| 1201 | builder.setOutcome(outcome); |
| 1202 | } |
| 1203 | |
| 1204 | SpanClose SpanClose::clone() const { |
| 1205 | return SpanClose(outcome); |
| 1206 | } |
| 1207 | |
| 1208 | kj::String SpanClose::toString() const { |
| 1209 | return kj::str("SpanClose: ", outcome); |
| 1210 | } |
| 1211 | |
| 1212 | Onset::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 | |
| 1248 | void 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 | |
| 1283 | namespace { |
| 1284 | kj::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 | |
| 1291 | kj::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 | |
| 1299 | kj::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 | |
| 1306 | kj::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 | |
| 1313 | kj::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 | |
| 1325 | kj::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 | |
| 1332 | kj::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 | |
| 1339 | Onset::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 | |
| 1353 | Onset::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 | |
| 1360 | Onset::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 | |
| 1366 | void 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 | |
| 1402 | Onset::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 | |
| 1416 | EventInfo 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 | |
| 1452 | Onset Onset::clone() const { |
| 1453 | return Onset(spanId, cloneEventInfo(info), workerInfo.clone(), |
| 1454 | KJ_MAP(attr, attributes) { return attr.clone(); }); |
| 1455 | } |
| 1456 | |
| 1457 | Outcome::Outcome(EventOutcome outcome, kj::Duration cpuTime, kj::Duration wallTime) |
| 1458 | : outcome(outcome), |
| 1459 | cpuTime(cpuTime), |
| 1460 | wallTime(wallTime) {} |
| 1461 | |
| 1462 | Outcome::Outcome(rpc::Trace::Outcome::Reader reader) |
| 1463 | : outcome(reader.getOutcome()), |
| 1464 | cpuTime(reader.getCpuTime() * kj::MILLISECONDS), |
| 1465 | wallTime(reader.getWallTime() * kj::MILLISECONDS) {} |
| 1466 | |
| 1467 | void 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 | |
| 1473 | Outcome Outcome::clone() const { |
| 1474 | return Outcome(outcome, cpuTime, wallTime); |
| 1475 | } |
| 1476 | |
| 1477 | TailEvent::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 | |
| 1485 | TailEvent::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 | |
| 1498 | namespace { |
| 1499 | TailEvent::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 | |
| 1542 | TailEvent::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 | |
| 1549 | void 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 | |
| 1593 | TailEvent 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 | |
| 1633 | SpanOpenData::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 | |
| 1639 | void 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 | |
| 1646 | SpanEndData::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 | |
| 1657 | void 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 | |
| 1673 | SpanBuilder::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 | |
| 1687 | SpanBuilder& SpanBuilder::operator=(SpanBuilder&& other) { |
| 1688 | end(); |
| 1689 | observer = kj::mv(other.observer); |
| 1690 | span = kj::mv(other.span); |
| 1691 | return *this; |
| 1692 | } |
| 1693 | |
| 1694 | SpanBuilder::~SpanBuilder() noexcept(false) { |
| 1695 | end(); |
| 1696 | } |
| 1697 | |
| 1698 | void 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 | |
| 1710 | void 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 | |
| 1719 | void 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 | |
| 1768 | void 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 | |
| 1782 | void 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 | |
| 1830 | Span::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 | |
| 1848 | using RpcValue = rpc::TagValue; |
| 1849 | void 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 | |
| 1866 | Span::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 | |
| 1881 | ScopedDurationTagger::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 | |
| 1888 | ScopedDurationTagger::~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 |