File
Blob: src/workerd/io/trace-stream.h
| 1 | #pragma once |
| 2 | |
| 3 | #include <workerd/io/io-context.h> |
| 4 | #include <workerd/io/trace.h> |
| 5 | #include <workerd/io/tracer.h> |
| 6 | #include <workerd/io/worker-interface.h> |
| 7 | #include <workerd/util/checked-queue.h> |
| 8 | |
| 9 | namespace workerd::tracing { |
| 10 | |
| 11 | // A WorkerInterface::CustomEvent implementation used to deliver streaming tail |
| 12 | // events to a tail worker. |
| 13 | class TailStreamCustomEvent final: public WorkerInterface::CustomEvent { |
| 14 | public: |
| 15 | TailStreamCustomEvent(kj::PromiseFulfillerPair<rpc::TailStreamTarget::Client> paf = |
| 16 | kj::newPromiseAndFulfiller<rpc::TailStreamTarget::Client>()) |
| 17 | : capFulfiller(kj::mv(paf.fulfiller)), |
| 18 | clientCap(kj::mv(paf.promise)) {} |
| 19 | |
| 20 | ~TailStreamCustomEvent() noexcept(false) { |
| 21 | if (capFulfiller->isWaiting()) { |
| 22 | capFulfiller->reject( |
| 23 | KJ_EXCEPTION(DISCONNECTED, "TailStreamCustomEvent was destroyed before completion")); |
| 24 | } |
| 25 | } |
| 26 | |
| 27 | kj::Promise<Result> run(kj::Own<IoContext::IncomingRequest> incomingRequest, |
| 28 | kj::Maybe<kj::StringPtr> entrypointName, |
| 29 | kj::Maybe<Worker::VersionInfo> versionInfo, |
| 30 | Frankenvalue props, |
| 31 | kj::TaskSet& waitUntilTasks, |
| 32 | bool isDynamicDispatch) override; |
| 33 | |
| 34 | kj::Promise<Result> sendRpc(capnp::HttpOverCapnpFactory& httpOverCapnpFactory, |
| 35 | capnp::ByteStreamFactory& byteStreamFactory, |
| 36 | rpc::EventDispatcher::Client dispatcher) override; |
| 37 | |
| 38 | kj::Promise<Result> notSupported() override { |
| 39 | JSG_FAIL_REQUIRE(TypeError, "The receiver is not a tail stream"); |
| 40 | } |
| 41 | |
| 42 | uint16_t getType() override { |
| 43 | return TYPE; |
| 44 | } |
| 45 | |
| 46 | tracing::EventInfo getEventInfo() const override; |
| 47 | |
| 48 | void failed(const kj::Exception& e) override { |
| 49 | capFulfiller->reject(e.clone()); |
| 50 | } |
| 51 | |
| 52 | // Specify the event type for TailStreamCustomEvent (defined internally). |
| 53 | static constexpr uint16_t TYPE = 12; |
| 54 | |
| 55 | rpc::TailStreamTarget::Client getCap() { |
| 56 | auto result = kj::mv(KJ_ASSERT_NONNULL(clientCap, "can only call getCap() once")); |
| 57 | clientCap = kj::none; |
| 58 | return result; |
| 59 | } |
| 60 | |
| 61 | private: |
| 62 | kj::Own<kj::PromiseFulfiller<workerd::rpc::TailStreamTarget::Client>> capFulfiller; |
| 63 | kj::Maybe<rpc::TailStreamTarget::Client> clientCap; |
| 64 | }; |
| 65 | |
| 66 | // A utility class that receives tracing events and generates/reports TailEvents. |
| 67 | class TailStreamWriter final { |
| 68 | public: |
| 69 | // The maximum size of the queue, in bytes. |
| 70 | const size_t maxQueueSize = 2 * 1024 * 1024; |
| 71 | // The estimated overhead of TailEvent wrapping per message. This does not need to be very |
| 72 | // accurate, but should be enough to avoid allocating too much memory/hitting capnp RPC message |
| 73 | // size limits when sending many tiny events. |
| 74 | const size_t tailSerializationOverhead = 64; |
| 75 | |
| 76 | // The initial state of our tail worker writer is that it is pending the first onset event. During |
| 77 | // this time we will only have a collection of WorkerInterface instances. When our first event is |
| 78 | // reported (the onset) we will arrange to acquire tailStream capabilities from each then use |
| 79 | // those to report the initial onset. |
| 80 | using Pending = kj::Array<kj::Own<WorkerInterface>>; |
| 81 | TailStreamWriter(Pending pending, kj::TaskSet& waitUntilTasks); |
| 82 | KJ_DISALLOW_COPY_AND_MOVE(TailStreamWriter); |
| 83 | |
| 84 | void report(const InvocationSpanContext& context, |
| 85 | TailEvent::Event&& event, |
| 86 | kj::Date time, |
| 87 | size_t sizeHint); |
| 88 | |
| 89 | private: |
| 90 | // Instances of Active are refcounted. The TailStreamWriter itself holds the initial ref. Whenever |
| 91 | // events are being dispatched, an additional ref will be held by the outstanding pump promise in |
| 92 | // order to keep the client stub alive long enough for the rpc calls to complete. It is possible |
| 93 | // that the TailStreamWriter will be dropped while pump promises are still pending. |
| 94 | struct Active: public kj::Refcounted { |
| 95 | // Reference to keep the worker interface instance alive. |
| 96 | kj::Maybe<rpc::TailStreamTarget::Client> capability; |
| 97 | bool pumping = false; |
| 98 | bool onsetSeen = false; |
| 99 | // Estimated byte size of the queue, used to drop events to avoid excessive memory usage. |
| 100 | size_t queueSize = 0; |
| 101 | // The number of tail events we had to drop. We'll send a warning indicating this at the end of |
| 102 | // the stream. |
| 103 | uint32_t droppedEvents = 0; |
| 104 | workerd::util::Queue<TailEvent> queue; |
| 105 | |
| 106 | Active(rpc::TailStreamTarget::Client capability): capability(kj::mv(capability)) {} |
| 107 | }; |
| 108 | |
| 109 | struct Closed {}; |
| 110 | |
| 111 | kj::OneOf<Pending, kj::Vector<kj::Own<Active>>, Closed> inner; |
| 112 | kj::TaskSet& waitUntilTasks; |
| 113 | |
| 114 | static kj::Promise<void> pump(kj::Own<Active> current); |
| 115 | // Report an event to the tail stream writer. |
| 116 | // sizeHint: The approximate size of the event, in bytes. |
| 117 | bool reportImpl(TailEvent&& event, size_t sizeHint); |
| 118 | |
| 119 | uint32_t sequence = 0; |
| 120 | bool onsetSeen = false; |
| 121 | bool outcomeSeen = false; |
| 122 | }; |
| 123 | |
| 124 | kj::Maybe<kj::Own<tracing::TailStreamWriter>> initializeTailStreamWriter( |
| 125 | kj::Array<kj::Own<WorkerInterface>> streamingTailWorkers, kj::TaskSet& waitUntilTasks); |
| 126 | |
| 127 | } // namespace workerd::tracing |