File
Blob: src/workerd/api/analytics-engine.c++
| 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 | #include "analytics-engine.h" |
| 6 | |
| 7 | #include <workerd/api/analytics-engine-impl.h> |
| 8 | #include <workerd/api/analytics-engine.capnp.h> |
| 9 | #include <workerd/io/io-context.h> |
| 10 | |
| 11 | namespace workerd::api { |
| 12 | |
| 13 | void AnalyticsEngine::writeDataPoint( |
| 14 | jsg::Lock& js, jsg::Optional<api::AnalyticsEngine::AnalyticsEngineEvent> event) { |
| 15 | auto& context = IoContext::current(); |
| 16 | |
| 17 | context.getLimitEnforcer().newAnalyticsEngineRequest(); |
| 18 | |
| 19 | // Optimization: For non-actors, which never have output locks, avoid the overhead of |
| 20 | // awaitIo() and such by not going back to the event loop at all. |
| 21 | KJ_IF_SOME(promise, context.waitForOutputLocksIfNecessary()) { |
| 22 | context.awaitIo( |
| 23 | js, kj::mv(promise), [this, self = JSG_THIS, event = kj::mv(event)](jsg::Lock& js) mutable { |
| 24 | writeDataPointNoOutputLock(js, kj::mv(event)); |
| 25 | }); |
| 26 | } else { |
| 27 | writeDataPointNoOutputLock(js, kj::mv(event)); |
| 28 | } |
| 29 | } |
| 30 | |
| 31 | void AnalyticsEngine::writeDataPointNoOutputLock( |
| 32 | jsg::Lock& js, jsg::Optional<api::AnalyticsEngine::AnalyticsEngineEvent>&& event) { |
| 33 | auto& context = IoContext::current(); |
| 34 | |
| 35 | context.writeLogfwdr(logfwdrChannel, [&](capnp::AnyPointer::Builder ptr) { |
| 36 | api::AnalyticsEngineEvent::Builder aeEvent = ptr.initAs<api::AnalyticsEngineEvent>(); |
| 37 | |
| 38 | aeEvent.setAccountId(static_cast<int64_t>(ownerId)); |
| 39 | aeEvent.setTimestamp(now()); |
| 40 | aeEvent.setDataset(dataset.asBytes()); |
| 41 | aeEvent.setSchemaVersion(version); |
| 42 | // `index1` should default to the empty string (`""`). |
| 43 | // The optional call to `setIndexes()`below assumes defaults, if any. |
| 44 | aeEvent.setIndex1(""_kj.asBytes()); |
| 45 | |
| 46 | kj::StringPtr errorPrefix = "writeDataPoint(): "_kj; |
| 47 | KJ_IF_SOME(ev, event) { |
| 48 | KJ_IF_SOME(indexes, ev.indexes) { |
| 49 | setIndexes<api::AnalyticsEngineEvent::Builder>(aeEvent, indexes, errorPrefix); |
| 50 | } |
| 51 | KJ_IF_SOME(blobs, ev.blobs) { |
| 52 | setBlobs<api::AnalyticsEngineEvent::Builder>(aeEvent, blobs, errorPrefix); |
| 53 | } |
| 54 | KJ_IF_SOME(doubles, ev.doubles) { |
| 55 | setDoubles<api::AnalyticsEngineEvent::Builder>(aeEvent, doubles, errorPrefix); |
| 56 | } |
| 57 | } |
| 58 | }); |
| 59 | } |
| 60 | } // namespace workerd::api |