Skip to content
File

Blob: src/workerd/io/trace-stream.h

cpp128 lines
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 
9namespace workerd::tracing {
10 
11// A WorkerInterface::CustomEvent implementation used to deliver streaming tail
12// events to a tail worker.
13class 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.
67class 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 
124kj::Maybe<kj::Own<tracing::TailStreamWriter>> initializeTailStreamWriter(
125 kj::Array<kj::Own<WorkerInterface>> streamingTailWorkers, kj::TaskSet& waitUntilTasks);
126 
127} // namespace workerd::tracing