File
Blob: src/workerd/io/observer.h
| 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 | |
| 18 | namespace workerd { |
| 19 | |
| 20 | class IoContext; |
| 21 | class WorkerInterface; |
| 22 | class LimitEnforcer; |
| 23 | class TimerChannel; |
| 24 | |
| 25 | class 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. |
| 41 | class 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. |
| 54 | class 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 | |
| 161 | class JsgIsolateObserver: public kj::AtomicRefcounted, public jsg::IsolateObserver {}; |
| 162 | |
| 163 | class 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 | |
| 272 | class 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 | |
| 290 | class 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. |
| 328 | template <typename Observer> |
| 329 | class 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. |
| 346 | class 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 |