Decentralised Art Server
High-performance C++ backend that exposes HTML interface and a secure REST API for managing Performative Transactions entities
Loading...
Searching...
No Matches
events_runtime.hpp
Go to the documentation of this file.
1#pragma once
2
3#include <atomic>
4#include <cstdint>
5#include <filesystem>
6#include <memory>
7#include <optional>
8#include <string>
9#include <vector>
10
11#include "native.h"
12#include <asio.hpp>
13
14#include "sqlite/wal_store.hpp"
15
17#include "event_projector.hpp"
18#include "sqlite_hot_store.hpp"
19
20namespace dcn::events
21{
22 constexpr std::size_t DEFAULT_PROJECT_BATCH_SIZE = 256;
23
24 std::int64_t reorgLookbackStart(std::int64_t next_from_block, std::size_t reorg_window_blocks);
25
26
28 {
29 std::filesystem::path hot_db_path;
30
31 int chain_id = 1;
32 bool ingestion_enabled = false;
33 std::vector<std::shared_ptr<IEmittedLogSource>> sources;
34 std::optional<std::int64_t> start_block = std::nullopt;
35 unsigned int poll_interval_ms = 5000;
36 unsigned int confirmations = 12;
37 unsigned int block_batch_size = 500;
38
39 std::size_t reorg_window_blocks = 2048;
40
41 unsigned int projector_interval_ms = 200;
42 unsigned int prune_interval_ms = 5000;
43 unsigned int wal_checkpoint_interval_ms = 15 * 1000;
44 };
45
47 {
48 public:
49 EventRuntime(asio::io_context & io_context, EventRuntimeConfig config);
51
52 EventRuntime(const EventRuntime &) = delete;
53 EventRuntime & operator=(const EventRuntime &) = delete;
54
57
58 void start();
59 void requestStop();
60 asio::awaitable<void> stop();
61 bool running() const;
62 bool ingestionEnabled() const;
63
64 void addProjector(std::unique_ptr<IEventProjector> projector);
66 const asio::strand<asio::io_context::executor_type> & writeStrand() const { return _write_strand; }
67
68 asio::awaitable<storage::sqlite::WalCheckpointStats> checkpointWal(storage::sqlite::WalCheckpointMode mode) const override;
69
70 private:
71 asio::awaitable<void> _sleepFor(const std::uint64_t ms) const;
72
73 asio::awaitable<std::optional<std::int64_t>> _storeLoadNextFromBlock(int chain_id) const;
74 asio::awaitable<std::optional<std::uint64_t>> _storeLoadNextLocalSeq(int chain_id) const;
75 asio::awaitable<bool> _storeSaveNextLocalSeq(int chain_id, std::uint64_t next_seq, std::int64_t now_ms) const;
76 asio::awaitable<std::vector<std::int64_t>> _storeLoadReorgWindowBlocks(
77 int chain_id,
78 std::int64_t from_block,
79 std::int64_t to_block) const;
80 asio::awaitable<bool> _storeIngestBatch(
81 int chain_id,
82 std::vector<RawChainLog> raw_events,
83 std::vector<DecodedEvent> decoded_events,
84 std::vector<ChainBlockInfo> block_infos,
85 std::int64_t next_from_block,
86 std::int64_t now_ms,
87 std::optional<std::uint64_t> next_local_seq = std::nullopt) const;
88 asio::awaitable<bool> _storeApplyFinality(
89 int chain_id,
90 FinalityHeights heights,
91 std::int64_t now_ms,
92 std::size_t reorg_window_blocks) const;
93
94 asio::awaitable<void> _runSourceIngestionLoop(std::shared_ptr<IEmittedLogSource> source);
95 asio::awaitable<void> _runProjectorLoop();
96 asio::awaitable<void> _runPruneLoop();
97 asio::awaitable<void> _runMaintenanceLoop();
98 asio::awaitable<void> _waitForLoops();
99
100 private:
101 asio::io_context & _io_context;
102 EventRuntimeConfig _config;
103 asio::strand<asio::io_context::executor_type> _write_strand;
104
105 std::shared_ptr<SQLiteHotStore> _store;
106 std::unique_ptr<IEventDecoder> _decoder;
107
108 std::vector<std::unique_ptr<IEventProjector>> _projectors;
109
110 std::atomic<bool> _stop_requested{false};
111 std::atomic<bool> _running{false};
112
113 std::atomic<std::size_t> _active_loop_count{0};
114 };
115}
const asio::strand< asio::io_context::executor_type > & writeStrand() const
Definition events_runtime.hpp:66
EventRuntime(asio::io_context &io_context, EventRuntimeConfig config)
Definition events_runtime.cpp:29
EventRuntime(EventRuntime &&)=delete
asio::awaitable< storage::sqlite::WalCheckpointStats > checkpointWal(storage::sqlite::WalCheckpointMode mode) const override
Definition events_runtime.cpp:241
void requestStop()
Definition events_runtime.cpp:110
void start()
Definition events_runtime.cpp:51
bool running() const
Definition events_runtime.cpp:139
EventRuntime & operator=(const EventRuntime &)=delete
asio::awaitable< void > stop()
Definition events_runtime.cpp:120
void addProjector(std::unique_ptr< IEventProjector > projector)
Definition events_runtime.cpp:149
EventRuntime(const EventRuntime &)=delete
EventRuntime & operator=(EventRuntime &&)=delete
SQLiteHotStore & projectionStore()
Definition events_runtime.cpp:156
bool ingestionEnabled() const
Definition events_runtime.cpp:144
~EventRuntime()
Definition events_runtime.cpp:40
Definition sqlite_hot_store.hpp:19
Definition wal_store.hpp:10
Definition config.hpp:8
Definition decoded_event.hpp:11
constexpr std::size_t DEFAULT_PROJECT_BATCH_SIZE
Definition events_runtime.hpp:22
std::int64_t reorgLookbackStart(std::int64_t next_from_block, std::size_t reorg_window_blocks)
Definition events_runtime.cpp:23
WalCheckpointMode
Definition wal.hpp:9
Definition events_runtime.hpp:28
unsigned int wal_checkpoint_interval_ms
Definition events_runtime.hpp:43
unsigned int confirmations
Definition events_runtime.hpp:36
bool ingestion_enabled
Definition events_runtime.hpp:32
std::filesystem::path hot_db_path
Definition events_runtime.hpp:29
unsigned int block_batch_size
Definition events_runtime.hpp:37
std::size_t reorg_window_blocks
Definition events_runtime.hpp:39
unsigned int projector_interval_ms
Definition events_runtime.hpp:41
unsigned int prune_interval_ms
Definition events_runtime.hpp:42
std::optional< std::int64_t > start_block
Definition events_runtime.hpp:34
unsigned int poll_interval_ms
Definition events_runtime.hpp:35
int chain_id
Definition events_runtime.hpp:31
std::vector< std::shared_ptr< IEmittedLogSource > > sources
Definition events_runtime.hpp:33
Definition events_ingest.hpp:12