| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457 |
- #pragma once
- #include "document.hpp"
- #include "memory_store.hpp"
- #include "database_grpc_impl.hpp"
- #include "persistence/persistence_manager.hpp"
- #include "persistence/history_store.hpp"
- #include "encryption/encryption_manager.hpp"
- #include "events/event_manager.hpp"
- #include "files/file_manager.hpp"
- #include "replication/replication_manager.hpp"
- #include "migrations/migration_runner.hpp"
- #include "views/view_manager.hpp"
- #include "config/collection_config_manager.hpp"
- #include "security/policy_manager.hpp"
- #include "relations/relation_manager.hpp"
- // LMDB storage substrate (v2.0+). DocumentStore + LmdbEnv live alongside the
- // v1.x MemoryStore, which is still the write entry point and a bounded read
- // cache; writes are mirrored into LMDB under MemoryStore's per-collection lock
- // and reads are LMDB-first while the mirror is healthy.
- //
- // This coexistence is deliberate and still current. Retiring MemoryStore means
- // migrating the write handlers onto a WriteCoordinator - see the "Pending"
- // section of docs/ROADMAP.md. Older comments call that "Stage 4"; the "Phase C"
- // label they sometimes carry refers to a homegrown page engine that was never
- // built (LMDB replaced it).
- namespace smartbotic::db::storage {
- class LmdbEnv;
- class DocumentStore;
- class ProjectStoreRegistry;
- }
- #include <nlohmann/json.hpp>
- #include <atomic>
- #include <filesystem>
- #include <memory>
- #include <mutex>
- #include <string>
- #include <string_view>
- #include <unordered_map>
- #include <thread>
- #include <vector>
- namespace grpc {
- class Server;
- }
- namespace smartbotic::database {
- /**
- * Main storage service application.
- * Coordinates all components and manages the gRPC server.
- */
- class DatabaseService {
- public:
- // v2.4 — per-listener configuration. Each binding can have its own
- // TLS material and auth policy. The default 127.0.0.1 listener
- // preserves v2.3 behaviour (plaintext, no auth) so existing local
- // consumers keep working without code change.
- struct ListenerConfig {
- std::string bind = "127.0.0.1";
- uint16_t port = 9004;
- struct TlsConfig {
- bool enabled = false;
- std::filesystem::path cert_path;
- std::filesystem::path key_path;
- // When TLS is enabled but cert+key are missing/unreadable, the
- // server auto-generates a 10-year self-signed cert to
- // <dataDir>/tls/auto_self_signed.{pem,key}. Operators who run
- // a real PKI override by dropping their cert at cert_path.
- bool auto_self_signed_if_missing = true;
- } tls;
- struct AuthConfig {
- bool required = false;
- // v2.7.0 — a key now carries a name, which becomes the
- // PRINCIPAL that access policy is written against.
- //
- // Config accepts either form:
- // "keys": [ {"name": "shadowman", "key": "<base64>"} ]
- // "keys": [ "<base64>" ] // v2.4-v2.6 shape
- //
- // A bare string still authenticates and maps to the reserved
- // principal `unnamed`, so existing deployments keep working; they
- // simply cannot be told apart in a policy until they are named.
- struct NamedKey {
- std::string name; // principal; "unnamed" when config gave a bare string
- std::string key; // the bearer token itself
- };
- // Bearer tokens accepted by this listener. Constant-time
- // compared against the value of the `authorization` gRPC
- // metadata header. Rotation = add new key, distribute,
- // remove old key. Empty list with required=true is a config
- // error caught at startup.
- std::vector<NamedKey> keys;
- } auth;
- };
- struct Config {
- std::string nodeId = "storage-1";
- // v2.4 — listener fleet. Empty means "synthesize a single
- // back-compat listener from the legacy bindAddress + rpcPort
- // fields below". Once any listener is configured here, the
- // legacy fields are ignored.
- std::vector<ListenerConfig> listeners;
- // v2.3 legacy single-listener fields. Kept so existing
- // config.json files continue to work — the config loader
- // synthesises a ListenerConfig from these when `listeners` is
- // empty. Operators who want TLS or auth must use the
- // `listeners[]` schema instead of these.
- std::string bindAddress = "0.0.0.0";
- uint16_t rpcPort = 9004;
- // Bootstrap/service registration
- std::string webapiUrl;
- std::string serviceKey;
- uint32_t heartbeatIntervalSec = 30;
- std::filesystem::path dataDirectory;
- std::filesystem::path keyFilePath;
- // Memory eviction settings
- uint64_t maxMemoryMb = 800; // Max memory before eviction
- uint32_t evictionThresholdPercent = 80; // Start evicting at this %
- uint32_t evictionTargetPercent = 60; // Evict down to this %
- uint32_t evictionCheckIntervalMs = 5000; // How often to check
- // NEW v1.7.0 eviction tuning — mirrors MemoryStore::Config fields.
- // Plumbed here so the config loader can set them and setupComponents()
- // can copy them into the MemoryStore. Schema-only in T2; behavior lands
- // in T3/T4.
- uint32_t evictionChunkSize = 1000;
- uint32_t evictionChunkPauseMs = 50;
- uint32_t maxEvictionPassesPerTrigger = 20;
- uint32_t hotWriteFloorMs = 30000;
- uint32_t memorySoftPercent = 70;
- uint32_t memoryHardPercent = 85;
- uint32_t memoryEmergencyPercent = 95;
- // v1.7.0 T10 — eviction burst event threshold (docs per tick).
- uint32_t evictionBurstThreshold = 10000;
- // v2.4.3 — cap on the share of the resident set one pressure episode
- // may evict. 0 disables. See MemoryStore::Config for the rationale.
- uint32_t evictionMaxEpisodePercent = 50;
- // Persistence settings
- uint32_t walSyncIntervalMs = 100;
- uint32_t snapshotIntervalSec = 3600;
- bool compressionEnabled = true;
- // Full persistence manager config (recovery mode, etc.).
- // DatabaseService::setupComponents() overlays the simpler fields above
- // onto this struct before constructing PersistenceManager.
- PersistenceManager::Config persistenceConfig;
- // File storage settings
- uint64_t maxFileSizeMb = 500;
- std::vector<std::string> allowedFileTypes;
- uint32_t fileCleanupIntervalSec = 3600;
- // v2.8.0 — default file retention per type. See
- // FileManager::Config::defaultTtlSecondsByType for the key syntax.
- std::unordered_map<std::string, uint32_t> fileDefaultTtlByType;
- // Encryption settings
- bool encryptionEnabled = true;
- bool autoGenerateKey = true;
- // Replication settings
- bool replicationEnabled = true;
- std::vector<std::string> peerAddresses;
- std::string conflictResolution = "last_writer_wins";
- // Migration settings
- struct MigrationsConfig {
- bool enabled = true;
- std::filesystem::path directory;
- bool autoApply = true;
- bool failOnError = true;
- } migrations;
- // gRPC server settings (v1.6.2 — concurrency cap + configurable sizes)
- struct GrpcConfig {
- uint32_t maxReceiveMessageSizeMb = 100;
- uint32_t maxSendMessageSizeMb = 100;
- uint32_t resourceQuotaMemoryMb = 256;
- uint32_t maxConcurrentSubscribeStreams = 50;
- uint32_t maxConcurrentFileStreams = 10;
- };
- GrpcConfig grpc;
- };
- explicit DatabaseService(Config config);
- ~DatabaseService();
- /**
- * Initialize the service.
- * Loads configuration, initializes components, and prepares for startup.
- */
- bool initialize();
- /**
- * Start the service.
- * Starts the gRPC server and all background threads.
- */
- void start();
- /**
- * Stop the service.
- * Gracefully shuts down all components.
- */
- void stop();
- /**
- * Wait for the service to stop.
- */
- void wait();
- /**
- * Check if the service is running.
- */
- [[nodiscard]] bool isRunning() const { return running_; }
- /**
- * Signal the service to stop (for signal handlers).
- */
- void signalStop();
- /**
- * Get service statistics.
- */
- [[nodiscard]] nlohmann::json getStats() const;
- /**
- * Access the view manager (for RPC handlers, migrations, etc.).
- */
- ViewManager& viewManager() { return *view_manager_; }
- PolicyManager& policyManager() { return *policy_manager_; }
- /** Current read-only state (atomic, lock-free read). */
- bool isReadOnly() const { return read_only_.load(std::memory_order_acquire); }
- /** Human-readable reason the DB is read-only (empty when writable). */
- std::string readOnlyReason() const;
- /** Toggle read-only state. Called by SetReadOnly RPC and startup logic. */
- void setReadOnly(bool value, const std::string& reason);
- /** Recovery outcome from the last recover() call. Used by GetReadOnlyStatus RPC. */
- const RecoveryOutcome& recoveryOutcome() const { return recovery_outcome_; }
- /** Set the force-readwrite flag (from --force-readwrite CLI arg). */
- void setForceReadwrite(bool value) { force_readwrite_ = value; }
- /**
- * Load configuration from a JSON file.
- */
- static Config loadConfig(const std::filesystem::path& configPath);
- /**
- * Parse configuration from an already-loaded JSON object. Used by the
- * conf.d drop-in loader (see config/config_loader.hpp) after it has
- * merged /etc/smartbotic-database/config.json with every *.json in
- * /etc/smartbotic-database/conf.d/.
- */
- static Config parseConfig(const nlohmann::json& json);
- // v2.3 — project-aware read accessors. `docStore(project)` returns
- // the per-project LmdbDocumentStore or nullptr if the project isn't
- // open (rare unless someone dropped it concurrently). The no-arg
- // `docStore()` defaults to the "default" project for callers still
- // on the v2.0-v2.2 API shape (Stage D will retire it). Health +
- // drift atomics gate any LMDB-first read path.
- smartbotic::db::storage::DocumentStore* docStore() noexcept;
- smartbotic::db::storage::DocumentStore* docStore(std::string_view project) noexcept;
- smartbotic::db::storage::ProjectStoreRegistry* projects() noexcept {
- return projects_.get();
- }
- bool mirrorHealthy() const noexcept {
- return mirror_healthy_.load(std::memory_order_acquire);
- }
- uint64_t mirrorDriftCount() const noexcept {
- return mirror_drift_count_.load(std::memory_order_relaxed);
- }
- // v2.11.0 close-out — drift accrued AFTER the boot path finished. What
- // destructive relation policies gate on; every READ gate keeps using
- // mirrorDriftCount() above. See MemoryStore::markMirrorDriftBaseline().
- uint64_t mirrorDriftSinceReady() const noexcept {
- return store_ == nullptr ? 0 : store_->mirrorDriftSinceBaseline();
- }
- uint64_t mirrorDriftAtReady() const noexcept {
- return store_ == nullptr ? 0 : store_->mirrorDriftBaseline();
- }
- // v2.11.0 T12 review (C2) — replication queueing + Subscribe event
- // publication, extracted out of the MemoryStore persistCallback_
- // lambda (setupComponents()) so a caller that deliberately bypasses
- // that callback can still drive both explicitly, in the same shape
- // the callback itself uses. This is what a cascade's `executeCascade`
- // (relations/relation_cascade.cpp) calls once per child mutation and
- // once for the parent delete, AFTER MemoryStore has been updated —
- // mirroring how v2.3.1 fixed replicated-entry apply by explicitly
- // driving `applyDualWriteMirror` rather than relying on a callback
- // that had already been bypassed for the same "don't double up on
- // WAL/mirror" reason. `collection` is the QUALIFIED name (what
- // `ReplicationEntry`/`DatabaseEvent` both expect); `doc` is null for a
- // delete. Does NOT touch WAL or the LMDB mirror — those are the
- // caller's job, done separately, and calling this a second time for
- // the same mutation would double-queue replication and double-publish
- // the event, so callers must call it exactly once per mutation.
- void notifyReplicationAndEvents(const std::string& collection,
- const std::string& id,
- const std::optional<Document>& doc,
- EventType eventType);
- // v2.3 Stage F — project CRUD entry points (delegate to registry).
- // The registry is the source of truth for listing / creating / dropping.
- // CreateProject is idempotent. DropProject refuses "default". The
- // gRPC handlers in database_grpc_impl.cpp call these.
- std::vector<std::string> listProjects() const;
- bool createProject(const std::string& name, std::string& error);
- bool dropProject(const std::string& name, std::string& error);
- private:
- void setupComponents();
- void startGrpcServer();
- void stopGrpcServer();
- void applyReplicatedEntry(const databasepb::ReplicationEntry& entry);
- bool runMigrations();
- // v2.0 Stage 4 — synchronous backfill of MemoryStore into doc_store_.
- // Runs once at boot, after recovery + migrations, before gRPC accepts
- // traffic. Skips when _meta.schema_version=2 marker exists (the env
- // is already populated by migrate_v1_to_v2). Returns false only on
- // catastrophic failure; per-doc errors flip mirror_healthy_=false but
- // do not abort boot (the read-flip downstream is the gate).
- bool backfillIntoDocStore();
- // v2.4.4 — audit every open project env for documents whose declared
- // `collection` disagrees with the sub-db they occupy. Logs ERROR per
- // affected project and leaves a summary; never fails startup.
- void auditSubdbPlacement();
- // v2.9.0 — arm the write path and query planner with the index declarations
- // persisted in _collection_meta. Must run at boot, after config load.
- void applyIndexDeclarations();
- // v2.11.0 close-out — MemoryStore::TtlExpiryRelationHook. Decides what the
- // TTL sweeper must do about an expiring parent's children, and PERFORMS the
- // cascade itself when one is required, by calling the ordinary
- // relations/relation_cascade.cpp machinery (WAL-before-LMDB) rather than a
- // second, sweeper-local implementation - a hand-rolled cascade that wrote
- // only LMDB would have its child deletions resurrected by the next boot's
- // WAL replay, since MemoryStore is rebuilt from snapshot + WAL and NOT from
- // LMDB.
- //
- // Runs on the cleanup thread with NO MemoryStore lock held (that is
- // MemoryStore::expireDocuments()' phase-2 contract, and it is what makes
- // calling back into MemoryStore here safe). Installed in setupComponents();
- // must never throw (the sweeper catches, but a throw per document would
- // stop expiry making progress).
- MemoryStore::TtlExpiryAction ttlExpiryRelationDecision(
- const std::string& qualifiedCollection, const std::string& id, bool firstAttempt);
- // v2.11.0 T6a — arm each project's LmdbDocumentStore with the relation
- // declarations loaded by relation_manager_->loadFromStore(), grouped by
- // CHILD collection. LmdbDocumentStore has no knowledge of RelationManager
- // and maintains the reverse index for nothing until told to via
- // set_relations() (see Task 3's report) — this is what tells it, once at
- // boot. Task 6b's createRelation/dropRelation RPCs must call
- // set_relations() again on live mutation; this only covers what was
- // already persisted at startup.
- // `afterMigrations` only affects LOGGING (see the self-heal block): a
- // relation whose reverse index sub-db is absent is a fault at boot and the
- // expected state for a relation a migration just declared, and the same code
- // handles both.
- void applyRelationDeclarations(bool afterMigrations = false);
- // v2.3 Stage C — atomic rename of <dataDir>/env/ into
- // <dataDir>/projects/default/env/ when the v2.2 layout is detected
- // and the new layout doesn't yet exist. Idempotent. Refuses to start
- // if both exist (likely operator hand-mess); the operator must
- // resolve manually.
- void migrateLegacyEnvToDefaultProject();
- Config config_;
- std::atomic<bool> running_{false};
- std::atomic<bool> stopRequested_{false};
- // Components
- std::unique_ptr<MemoryStore> store_;
- // v2.3 — multi-project storage substrate. Each project is one
- // LmdbEnv at `<dataDir>/projects/<name>/env/`. The registry holds
- // them in a map keyed by project name, lazy-opens new projects on
- // first write, and is the single source of truth for "which projects
- // exist." Old v2.0-v2.2 installs auto-migrate the legacy
- // `<dataDir>/env/` into `<dataDir>/projects/default/env/` at first
- // boot (see migrateLegacyEnvToDefaultProject).
- std::unique_ptr<smartbotic::db::storage::ProjectStoreRegistry> projects_;
- // v2.0 dual-write health. The persist callback mirrors every user-
- // collection write into doc_store_; on failure we log ERROR, bump the
- // drift counter, and flip mirror_healthy_ to false. Downstream Stage 4
- // tasks that flip reads onto doc_store_ MUST check mirror_healthy_
- // before engaging — a degraded mirror means LMDB is missing writes.
- std::atomic<bool> mirror_healthy_{true};
- std::atomic<uint64_t> mirror_drift_count_{0};
- std::unique_ptr<PersistenceManager> persistence_;
- std::unique_ptr<HistoryStore> history_store_; // v1.9.0 — disk-resident history
- std::unique_ptr<EncryptionManager> encryption_;
- std::unique_ptr<EventManager> events_;
- std::unique_ptr<FileManager> files_;
- std::unique_ptr<ReplicationManager> replication_;
- std::unique_ptr<ViewManager> view_manager_;
- std::unique_ptr<CollectionConfigManager> config_manager_;
- // v2.8.0 — per-project access policy. Owned here so its cache is loaded
- // once, after recovery, alongside ViewManager and CollectionConfigManager.
- std::unique_ptr<PolicyManager> policy_manager_;
- // v2.11.0 T6a — relation declarations. Owned here so its cache is loaded
- // once, after recovery, alongside ViewManager/CollectionConfigManager/PolicyManager.
- std::unique_ptr<RelationManager> relation_manager_;
- // gRPC
- std::unique_ptr<DatabaseGrpcImpl> storageImpl_;
- std::unique_ptr<DatabaseReplicationGrpcImpl> replicationImpl_;
- // v2.4 — one gRPC Server per listener (each one ports.json entry
- // builds a separate Server instance with its own TLS / auth
- // policy). Pre-2.4 there was a single grpcServer_ wired to
- // bindAddress + rpcPort; the back-compat shim in parseConfig
- // synthesises exactly one entry for v2.3-shape configs so the
- // observable behavior is unchanged on upgrade.
- std::vector<std::unique_ptr<grpc::Server>> grpcServers_;
- std::thread serverThread_;
- // Migration runner
- std::unique_ptr<MigrationRunner> migrationRunner_;
- // Read-only runtime state (v1.6.1 — auto-readonly on non-trivial recovery)
- std::atomic<bool> read_only_{false};
- mutable std::mutex reason_mutex_;
- std::string read_only_reason_;
- RecoveryOutcome recovery_outcome_;
- bool force_readwrite_ = false;
- };
- } // namespace smartbotic::database
|