Skip to content
File

Blob: src/workerd/io/observer.h

cpp373 lines
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#pragma once
6// Defines abstract interfaces for observing the activity of various components of the system,
7// e.g. to collect logs and metrics.
8 
9#include <workerd/io/features.capnp.h>
10#include <workerd/io/trace.h>
11#include <workerd/jsg/observer.h>
12#include <workerd/util/sqlite.h>
13 
14#include <kj/refcount.h>
15#include <kj/string.h>
16#include <kj/time.h>
17 
18namespace workerd {
19 
20class IoContext;
21class WorkerInterface;
22class LimitEnforcer;
23class TimerChannel;
24 
25class WebSocketObserver: public kj::Refcounted {
26 public:
27 virtual ~WebSocketObserver() noexcept(false) = default;
28 // Called when a worker sends a message on this WebSocket (includes close messages).
29 virtual void sentMessage(size_t bytes) {};
30 // Called when a worker receives a message on this WebSocket (includes close messages).
31 virtual void receivedMessage(size_t bytes) {};
32};
33 
34// Observes a byte stream. Byte streams which use instances of this observer should call enqueue()
35// and dequeue() once for each chunk that passes through the stream. The order of enqueues should
36// match the order of dequeues.
37//
38// Byte observer implementations can then calculate the current number of chunks and the sum of the
39// size of the chunks in the internal queue by incrementing and decrementing each metric in
40// enqueue() and dequeue() respectively.
41class ByteStreamObserver {
42 public:
43 virtual ~ByteStreamObserver() noexcept(false) = default;
44 // Called when a chunk of size `bytes` is enqueued on the stream.
45 virtual void onChunkEnqueued(size_t bytes) {};
46 // Called when a chunk of size `bytes` is dequeued from the stream (e.g. when a writable byte
47 // stream writes the chunk to its corresponding sink).
48 virtual void onChunkDequeued(size_t bytes) {};
49};
50 
51// Observes a specific request to a specific worker. Also observes outgoing subrequests.
52//
53// Observing anything is optional. Default implementations of all methods observe nothing.
54class RequestObserver: public kj::Refcounted {
55 public:
56 // This is called when the request is converted to a WebSocket connection terminating in a worker.
57 // An optional WebSocket observer may be returned to observe events on the worker's end of the
58 // WebSocket connection.
59 //
60 // This means that, when the returned observer observes a message being sent, the message is being
61 // sent from the worker to the client making the request.
62 virtual kj::Maybe<kj::Own<WebSocketObserver>> tryCreateWebSocketObserver() {
63 return kj::none;
64 };
65 
66 // This is called when a writable byte stream is created whilst processing this request. It will
67 // be destroyed when the corresponding byte stream is destroyed.
68 virtual kj::Maybe<kj::Own<ByteStreamObserver>> tryCreateWritableByteStreamObserver() {
69 return kj::none;
70 }
71 
72 // Invoked when the request is actually delivered.
73 //
74 // If, for some reason, this is not invoked before the object is destroyed, this indicate that
75 // the event was canceled for some reason before delivery. No JavaScript was invoked. In this
76 // case, the request should not be billed.
77 virtual void delivered() {};
78 
79 // Call when no more JavaScript will run on behalf of this request. Note that deferred proxying
80 // may still be in progress.
81 virtual void jsDone() {}
82 
83 // Called to indicate this was a prewarm request. Normal request metrics won't be logged, but
84 // the prewarm metric will be incremented.
85 virtual void setIsPrewarm() {}
86 
87 // Describes the source of a failure
88 enum class FailureSource : uint8_t {
89 // Failure occurred during deferred proxying
90 DEFERRED_PROXY,
91 
92 // Failure occurred elsewhere
93 OTHER,
94 };
95 
96 // Report that the request failed with the given exception. This only needs to be called in
97 // cases where the wrapper created with wrapWorkerInterface() wouldn't otherwise see the
98 // exception, e.g. because it has been replaced with an HTTP error response or because it
99 // occurred asynchronously.
100 virtual void reportFailure(const kj::Exception& e, FailureSource source = FailureSource::OTHER) {}
101 
102 static EventOutcome outcomeFromException(
103 const kj::Exception& e, FailureSource source = FailureSource::OTHER);
104 
105 // Called when an internal exception is observed during this request. Used to track which
106 // internal exception types occurred during a request, for metrics purposes. The same exception
107 // type may be reported multiple times during a single request; implementations should deduplicate.
108 virtual void reportInternalException(
109 const kj::Exception& e, jsg::InternalExceptionObserver::Detail detail) {}
110 
111 // Wrap the given WorkerInterface with a version that collects metrics. This method may only be
112 // called once, and only one method call may be made to the returned interface.
113 //
114 // The returned reference remains valid as long as the observer and `worker` both remain live.
115 virtual WorkerInterface& wrapWorkerInterface(WorkerInterface& worker) {
116 return worker;
117 }
118 
119 // Wrap an HttpClient so that its usage is counted in the request's subrequest stats.
120 virtual kj::Own<WorkerInterface> wrapSubrequestClient(kj::Own<WorkerInterface> client) {
121 return kj::mv(client);
122 }
123 
124 // Wrap an HttpClient so that its usage is counted in the request's actor subrequest count.
125 virtual kj::Own<WorkerInterface> wrapActorSubrequestClient(kj::Own<WorkerInterface> client) {
126 return kj::mv(client);
127 }
128 
129 // Used to record when a worker has used a dynamic dispatch binding.
130 virtual void setHasDispatched() {};
131 
132 virtual SpanParent getSpan() {
133 return nullptr;
134 }
135 
136 virtual void setOutcome(EventOutcome outcome) {}
137 
138 virtual kj::Own<void> addedContextTask() {
139 return kj::Own<void>();
140 }
141 virtual kj::Own<void> addedWaitUntilTask() {
142 return kj::Own<void>();
143 }
144 
145 virtual void setFailedOpen(bool value) {}
146 
147 // Called when the language runtime for this worker encounters a fatal error during this
148 // invocation. Currently used for Pyodide fatal errors, but is language-agnostic and can be used
149 // for other language runtimes in the future.
150 virtual void setWorkerFatal() {}
151 
152 // Called when PythonWorkersInternalError is constructed in JS. Used to track internal errors
153 // in Python workers without relying solely on substring checks.
154 virtual void setPythonWorkersInternalError() {}
155 
156 virtual uint64_t clockRead() {
157 return 0;
158 }
159};
160 
161class JsgIsolateObserver: public kj::AtomicRefcounted, public jsg::IsolateObserver {};
162 
163class IsolateObserver: public kj::AtomicRefcounted {
164 public:
165 virtual ~IsolateObserver() noexcept(false) {}
166 
167 // Called when Worker::Isolate is created.
168 virtual void created() {};
169 
170 // Called when the owning Worker::Script is being destroyed. The IsolateObserver may
171 // live a while longer to handle deferred proxy requests.
172 virtual void evicted() {}
173 
174 virtual void teardownStarted() {}
175 virtual void teardownLockAcquired() {}
176 virtual void teardownFinished() {}
177 
178 // Describes why a worker was started.
179 enum class StartType : uint8_t {
180 // Cold start with active request waiting.
181 COLD,
182 
183 // Started due to prewarm hint (e.g. from TLS SNI); a real request is expected soon.
184 PREWARM,
185 
186 // Started due to preload at process startup.
187 PRELOAD
188 };
189 
190 // Created while parsing a script, to record related metrics.
191 class Parse {
192 public:
193 // Marks the ScriptReplica as finished parsing, which starts reporting of isolate metrics.
194 virtual void done() {}
195 };
196 
197 virtual kj::Own<Parse> parse(StartType startType) const {
198 class FinalParse final: public Parse {};
199 return kj::heap<FinalParse>();
200 }
201 
202 class LockTiming {
203 public:
204 // Called by `Isolate::takeAsyncLock()` when it is blocked by a different isolate lock on the
205 // same thread.
206 virtual void waitingForOtherIsolate(kj::StringPtr id) {}
207 
208 // Call if this is an async lock attempt, before constructing LockRecord.
209 virtual void reportAsyncInfo(
210 uint currentLoad, bool threadWaitingSameLock, uint threadWaitingDifferentLockCount) {}
211 // TODO(cleanup): Should be able to get this data at `tryCreateLockTiming()` time. It'd be
212 // easier if IsolateObserver were an AOP class, and thus had access to the real isolate.
213 
214 virtual void start() {}
215 virtual void stop() {}
216 
217 virtual void locked() {}
218 virtual void gcPrologue() {}
219 virtual void gcEpilogue() {}
220 };
221 
222 // Construct a LockTiming if config.reportScriptLockTiming is true, or if the
223 // request (if any) is being traced.
224 virtual kj::Maybe<kj::Own<LockTiming>> tryCreateLockTiming(
225 kj::OneOf<SpanParent, kj::Maybe<RequestObserver&>> parentOrRequest) const {
226 return kj::none;
227 }
228 
229 // Use like so:
230 //
231 // auto lockTiming = MetricsCollector::ScriptReplica::LockTiming::tryCreate(
232 // script, maybeRequest);
233 // MetricsCollector::ScriptReplica::LockRecord record(lockTiming);
234 // isolate.runInLockScope([&](MyIsolate::Lock& lock) {
235 // record.locked();
236 // });
237 //
238 // And `record()` will report the time spent waiting for the lock (including any asynchronous
239 // time you might insert between the construction of `lockTiming` and `LockRecord()`), plus
240 // the time spent holding the lock for the given ScriptReplica.
241 //
242 // This is a thin wrapper around LockTiming which efficiently handles the case where we don't
243 // want to track timing.
244 class LockRecord {
245 public:
246 explicit LockRecord(kj::Maybe<kj::Own<LockTiming>> lockTimingParam)
247 : lockTiming(kj::mv(lockTimingParam)) {
248 KJ_IF_SOME(l, lockTiming) l.get()->start();
249 }
250 ~LockRecord() noexcept(false) {
251 KJ_IF_SOME(l, lockTiming) l.get()->stop();
252 }
253 KJ_DISALLOW_COPY_AND_MOVE(LockRecord);
254 
255 void locked() {
256 KJ_IF_SOME(l, lockTiming) l.get()->locked();
257 }
258 void gcPrologue() {
259 KJ_IF_SOME(l, lockTiming) l.get()->gcPrologue();
260 }
261 void gcEpilogue() {
262 KJ_IF_SOME(l, lockTiming) l.get()->gcEpilogue();
263 }
264 
265 private:
266 // The presence of `lockTiming` determines whether or not we need to record timing data. If
267 // we have no `lockTiming`, then this LockRecord wrapper is just a big nothingburger.
268 kj::Maybe<kj::Own<LockTiming>> lockTiming;
269 };
270};
271 
272class WorkerObserver: public kj::AtomicRefcounted {
273 public:
274 // Created while executing a script's global scope, to record related metrics.
275 class Startup {
276 public:
277 virtual void done() {}
278 };
279 
280 virtual kj::Own<Startup> startup(IsolateObserver::StartType startType) const {
281 class FinalStartup final: public Startup {};
282 return kj::heap<FinalStartup>();
283 }
284 
285 virtual void teardownStarted() {}
286 virtual void teardownLockAcquired() {}
287 virtual void teardownFinished() {}
288};
289 
290class ActorObserver: public kj::Refcounted, public SqliteObserver {
291 public:
292 // Allows the observer to run in the background, periodically making observations. Owner must
293 // call this and store the promise. `limitEnforcer` is used to collect CPU usage metrics, it
294 // must remain valid as long as the loop is running.
295 virtual kj::Promise<void> flushLoop(TimerChannel& timer, LimitEnforcer& limitEnforcer) {
296 return kj::NEVER_DONE;
297 }
298 
299 virtual void startRequest() {}
300 virtual void endRequest() {}
301 
302 virtual void webSocketAccepted() {}
303 virtual void webSocketClosed() {}
304 virtual void receivedWebSocketMessage(size_t bytes) {}
305 virtual void sentWebSocketMessage(size_t bytes) {}
306 
307 virtual void addCachedStorageReadUnits(uint32_t units) {}
308 virtual void addUncachedStorageReadUnits(uint32_t units) {}
309 virtual void addStorageWriteUnits(uint32_t units) {}
310 virtual void addStorageDeletes(uint32_t count) {}
311 
312 virtual void storageReadCompleted(kj::Duration latency) {}
313 virtual void storageWriteCompleted(kj::Duration latency) {}
314 
315 virtual void inputGateLocked() {}
316 virtual void inputGateReleased() {}
317 virtual void inputGateWaiterAdded() {}
318 virtual void inputGateWaiterRemoved() {}
319 virtual void outputGateLocked() {}
320 virtual void outputGateReleased() {}
321 virtual void outputGateWaiterAdded() {}
322 virtual void outputGateWaiterRemoved() {}
323 
324 virtual void shutdown(uint16_t reasonCode, LimitEnforcer& limitEnforcer) {}
325};
326 
327// RAII object to call `teardownFinished()` on an observer for you.
328template <typename Observer>
329class TeardownFinishedGuard {
330 public:
331 TeardownFinishedGuard(Observer& ref): ref(ref) {}
332 ~TeardownFinishedGuard() noexcept(false) {
333 ref.teardownFinished();
334 }
335 KJ_DISALLOW_COPY_AND_MOVE(TeardownFinishedGuard);
336 
337 private:
338 Observer& ref;
339};
340 
341// Provides counters/observers for various features. The intent is to
342// make it possible to collect metrics on which runtime features are
343// used and how often.
344//
345// There is exactly one instance of this class per worker process.
346class FeatureObserver {
347 public:
348 static kj::Own<FeatureObserver> createDefault();
349 static void init(kj::Own<FeatureObserver> instance);
350 static kj::Maybe<FeatureObserver&> get();
351 
352 // A "Feature" is just an opaque identifier defined in the features.capnp
353 // file.
354 using Feature = workerd::Features;
355 
356 // Called to increment the usage counter for a feature.
357 virtual void use(Feature feature) const {}
358 
359 using CollectCallback = kj::Function<void(Feature, const uint64_t)>;
360 // This method is called from the internal metrics collection mechanism to harvest the
361 // current features and counts that have been recorded by the observer.
362 virtual void collect(CollectCallback&& callback) const {}
363 
364 // Records the use of the feature if a FeatureObserver is available.
365 static inline void maybeRecordUse(Feature feature) {
366 KJ_IF_SOME(observer, get()) {
367 observer.use(feature);
368 }
369 }
370};
371 
372} // namespace workerd