File
Blob: src/workerd/api/analytics-engine.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 | |
| 7 | #include <workerd/io/io-util.h> |
| 8 | #include <workerd/jsg/jsg.h> |
| 9 | |
| 10 | namespace workerd::api { |
| 11 | |
| 12 | // Analytics Engine is a tool for customers to get telemetry about anything |
| 13 | // using Workers. The data points gathered from the edge are stored into |
| 14 | // ClickHouse and can be queried through the Analytics Engine's SQL API. |
| 15 | // |
| 16 | // The generated data points are encoded through the |
| 17 | // analytics_engine_event.capnp format and sent to logfwdr for them to enter |
| 18 | // the Data Pipeline. Each data point consists of an array of index values, 20 |
| 19 | // numeric fields (doubles) and 20 text fields (blobs), alongside some |
| 20 | // metadata. Aside from ordinality and maximum length, the semantics of the |
| 21 | // `blobs` and `doubles` fields are left up to applications submitting |
| 22 | // messages. |
| 23 | // |
| 24 | // https://blog.cloudflare.com/workers-analytics-engine/ |
| 25 | class AnalyticsEngine: public jsg::Object { |
| 26 | public: |
| 27 | explicit AnalyticsEngine( |
| 28 | uint logfwdrChannel, kj::String dataset, int64_t version, uint32_t ownerId) |
| 29 | : logfwdrChannel(logfwdrChannel), |
| 30 | dataset(kj::mv(dataset)), |
| 31 | version(version), |
| 32 | ownerId(ownerId) {} |
| 33 | struct AnalyticsEngineEvent { |
| 34 | // An array of values for the user-defined indexes, that provide a way for |
| 35 | // users to improve the efficiency of common queries. In addition, by |
| 36 | // default, the sampling key includes all the indexes in the list. This |
| 37 | // gives users some control over the way data is sampled. |
| 38 | jsg::Optional<kj::Array<kj::Maybe<kj::OneOf<kj::Array<byte>, kj::String>>>> indexes; |
| 39 | jsg::Optional<kj::Array<double>> doubles; |
| 40 | jsg::Optional<kj::Array<kj::Maybe<kj::OneOf<kj::Array<byte>, kj::String>>>> blobs; |
| 41 | |
| 42 | // The ordering of the elements within the `doubles` and `blobs` fields matters, insofar as the |
| 43 | // elements within each array will be unrolled, and based on the element's ordinality, the |
| 44 | // `.set{Blob,Data}{1..20}(element)` method is invoked. Because of this, these arrays can each |
| 45 | // have a maximum of 20 elements. |
| 46 | |
| 47 | JSG_STRUCT(indexes, doubles, blobs); |
| 48 | JSG_STRUCT_TS_OVERRIDE(AnalyticsEngineDataPoint); |
| 49 | }; |
| 50 | |
| 51 | // Send an Analytics Engine-compatible event to the configured logfwdr socket. |
| 52 | // Like logfwdr itself, `writeDataPoint` makes no delivery guarantees. |
| 53 | void writeDataPoint( |
| 54 | jsg::Lock& js, jsg::Optional<api::AnalyticsEngine::AnalyticsEngineEvent> event); |
| 55 | |
| 56 | JSG_RESOURCE_TYPE(AnalyticsEngine) { |
| 57 | JSG_METHOD(writeDataPoint); |
| 58 | JSG_TS_ROOT(); |
| 59 | JSG_TS_OVERRIDE(AnalyticsEngineDataset); |
| 60 | } |
| 61 | |
| 62 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 63 | tracker.trackField("dataset", dataset); |
| 64 | } |
| 65 | |
| 66 | private: |
| 67 | double millisToNanos(double m) { |
| 68 | return m * 1000000; |
| 69 | } |
| 70 | |
| 71 | // Called within writeDataPoint after waiting for output locks |
| 72 | void writeDataPointNoOutputLock( |
| 73 | jsg::Lock& js, jsg::Optional<api::AnalyticsEngine::AnalyticsEngineEvent>&& event); |
| 74 | |
| 75 | uint logfwdrChannel; |
| 76 | kj::String dataset; |
| 77 | int64_t version; |
| 78 | uint32_t ownerId; |
| 79 | |
| 80 | uint64_t now() { |
| 81 | return millisToNanos(dateNow()); |
| 82 | } |
| 83 | }; |
| 84 | #define EW_ANALYTICS_ENGINE_ISOLATE_TYPES \ |
| 85 | ::workerd::api::AnalyticsEngine, ::workerd::api::AnalyticsEngine::AnalyticsEngineEvent |
| 86 | |
| 87 | } // namespace workerd::api |