hydra-node
Safe HaskellSafe-Inferred
LanguageGHC2021

Hydra.Events.SQLiteBased

Description

A SQLite-backed event source and sink.

This is the recommended persistence backend for new deployments. See withSQLiteEventStore which handles migration from the legacy file-based store automatically.

Architecture

Events are stored in a single events table with an integer primary key (event_id) and a BLOB column (event_data) containing CBOR-encoded event data (via ToCBOR / FromCBOR). The database uses WAL journal mode with synchronous=NORMAL to avoid per-write fsyncs while still syncing at WAL checkpoints.

Schema migrations

The schema version is tracked in PRAGMA user_version and migrated on open (see applyMigrations). Version 1 stored event data as JSON; opening a version 1 database re-encodes every row to CBOR in one transaction and runs VACUUM afterwards to reclaim the freed space. A row that fails to decode aborts the migration (and thereby node startup) with EventDecodingException, rolling back to an intact version 1 database. The legacy file-based store (JSON lines) is migrated by decoding each line as JSON and inserting CBOR.

Async write-behind

To keep persistence off the hot path, writes use an async write-behind strategy. $sel:putEvent:EventSink and $sel:putEvents:EventSink encode events eagerly to strict ByteString and enqueue them into a bounded TBQueue. A background writer thread drains the queue and batch-inserts rows using executeMany inside a single transaction, amortising WAL frame writes across multiple events.

The last-seen event id TVar is updated atomically at enqueue time (not write time), so de-duplication and source-of-truth tracking remain correct even though the physical write is deferred.

All SQLite writes go through the single writer thread, preventing concurrent access races. Operations that need data flushed (rotation, reads) use a flush marker: a TMVar is enqueued and the caller blocks until the writer thread has processed all preceding items and signalled it. $sel:sourceEvents:EventSource auto-flushes before reading, so callers always see all enqueued events.

Tradeoffs

  • Writer thread crash surfacing: The background writer is linked to the calling thread. If it dies (e.g. SQLite I/O error), the exception propagates immediately rather than leaving the node silently stalled. Use withSQLiteEventStore which handles cleanup (flush + cancel) on exit.
  • Data loss on hard crash: Events in the queue that have not yet been flushed to SQLite are lost on SIGKILL, OOM, or power loss. This is acceptable because the L1 chain is the source of truth — the node replays missed events from chain on restart.
  • Rotation ordering: $sel:rotate:EventStore flushes the write queue synchronously, archives the current database to old-state/hydra-logId.db via VACUUM INTO, then performs DELETE + INSERT. This is safe because rotation is only called from the single-threaded event processing loop (processStateChanges), so no concurrent enqueues can occur between the flush and the rotation write. The archive is taken before the DELETE, so a backup failure aborts rotation and leaves the events intact.
  • Separate read connection: $sel:sourceEvents:EventSource streams over a dedicated connection. It can run concurrently on API server threads (client history replay), and VACUUM INTO fails with "SQL statements in progress" if a streaming statement is open on the same connection as the rotation. WAL mode makes readers on a separate connection safe.
Synopsis

Documentation

data SQLiteLog Source #

Instances

Instances details
ToJSON SQLiteLog Source # 
Instance details

Defined in Hydra.Events.SQLiteBased

Methods

toJSON :: SQLiteLog -> Value

toEncoding :: SQLiteLog -> Encoding

toJSONList :: [SQLiteLog] -> Value

toEncodingList :: [SQLiteLog] -> Encoding

omitField :: SQLiteLog -> Bool

Generic SQLiteLog Source # 
Instance details

Defined in Hydra.Events.SQLiteBased

Associated Types

type Rep SQLiteLog :: Type -> Type Source #

Show SQLiteLog Source # 
Instance details

Defined in Hydra.Events.SQLiteBased

Eq SQLiteLog Source # 
Instance details

Defined in Hydra.Events.SQLiteBased

type Rep SQLiteLog Source # 
Instance details

Defined in Hydra.Events.SQLiteBased

type Rep SQLiteLog = D1 ('MetaData "SQLiteLog" "Hydra.Events.SQLiteBased" "hydra-node-2.3.0-1cgalYNmLJQC1YXlYVq7mp" 'False) (C1 ('MetaCons "MigratingFromFileBased" 'PrefixI 'True) (S1 ('MetaSel ('Just "legacyFile") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedStrict) (Rec0 FilePath)) :+: (C1 ('MetaCons "MigrationSkipped" 'PrefixI 'True) (S1 ('MetaSel ('Just "legacyFile") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedStrict) (Rec0 FilePath)) :+: C1 ('MetaCons "MigrationComplete" 'PrefixI 'True) (S1 ('MetaSel ('Just "legacyFile") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedStrict) (Rec0 FilePath))))

type WriteItem e = Either (TMVar IO ()) (Word64, e) Source #

Items in the write-behind queue: either an event to insert or a flush marker that the writer thread signals after processing all preceding items.

Events are queued unencoded and CBOR-encoded on the writer thread: the encoding of e.g. a SnapshotRequested event carrying a large UTxO otherwise sits on the node loop between processing a ReqSn and broadcasting the AckSn. The bounded queue briefly pins event values instead of compact bytes, but the writer drains whole-queue batches so the window is short.

withSQLiteEventStore :: forall e a. (ToCBOR e, FromCBOR e, FromJSON e, HasEventId e) => Tracer IO SQLiteLog -> FilePath -> FilePath -> (EventStore e IO -> IO a) -> IO a Source #

Bracket-style wrapper around mkSQLiteEventStore. Creates the database, schema, and writer thread, runs the callback, then flushes queued writes and cancels the writer thread on exit. The writer thread is linked so that crashes surface immediately in the calling thread.

If a legacy state file exists at legacyStateFile, events are migrated into SQLite automatically before the callback runs.

Flushing of the async write queue and reinitialisation of the last-seen event id are handled internally: $sel:sourceEvents:EventSource auto-flushes before reading, $sel:rotate:EventStore flushes before deleting, migration reinitialises the event id TVar, and this bracket flushes on exit.

mkSQLiteEventStore :: forall e. (ToCBOR e, FromCBOR e, FromJSON e, HasEventId e) => FilePath -> IO (Connection, EventStore e IO, IO (), IO (), IO ()) Source #

Create an EventStore backed by a SQLite database at the given file path. The database and schema are created on first use if they do not exist. Returns (conn, store, flush, reinitLastSeen, cleanup). Internal — prefer withSQLiteEventStore which handles cleanup, migration, and flushing automatically.

writerLoop :: ToCBOR e => Connection -> TBQueue IO (WriteItem e) -> IO () Source #

Background writer that drains the queue and batch-inserts into SQLite. Each iteration blocks for at least one item, then flushes everything available. Events are CBOR-encoded here, off the caller's thread, then batch-inserted in a single transaction, and any flush markers in the batch are signalled. Encode errors surface as writer thread crashes, which are linked to the node.

flushWriteQueue :: TBQueue IO (WriteItem e) -> IO () Source #

Block until all items currently in the write queue have been flushed to SQLite. Sends a flush marker through the queue and waits for the writer thread to signal completion.

migrateFromFileBased :: forall e. (FromJSON e, ToCBOR e, HasEventId e) => Proxy e -> Tracer IO SQLiteLog -> FilePath -> Connection -> IO () -> IO () Source #

Migrate events from a legacy newline-delimited JSON file into SQLite. Writes directly to the database, bypassing the async write queue (migration runs at startup before the node processes inputs). After inserting, calls reinitLastSeen to sync the in-memory event id TVar with the database.

Safe to call when the legacy file does not exist (no-op). Not safe to re-run: duplicate event ids will cause a primary key constraint violation.

On success the legacy file is renamed to path.migrated so that subsequent node restarts skip the migration step automatically.

nextVersion :: Int Source #

Current schema version. Bump this and add a migration step to migrateStep whenever the schema changes.

type ReencodeRow = Word64 -> ByteString -> IO ByteString Source #

Re-encode a single event row given its event id and stored bytes, used by the version 1 (JSON) to version 2 (CBOR) migration. Must throw when the row cannot be decoded.

initSchema :: Connection -> ReencodeRow -> IO () Source #

Initialise connection pragmas, then create or migrate the schema to nextVersion using SQLite's built-in user_version pragma.

configurePragmas :: Connection -> IO () Source #

getSchemaVersion :: Connection -> IO Int Source #

Read the schema version from PRAGMA user_version (0 for a fresh DB).

setSchemaVersion :: Connection -> Int -> IO () Source #

applyMigrations :: Connection -> ReencodeRow -> Int -> IO () Source #

Apply all pending migrations from version v up to nextVersion. Each step runs together with its version bump in one transaction (PRAGMA user_version is transactional), so a crash or decoding failure mid-migration rolls back to a well-defined version.

migrateStep :: Connection -> ReencodeRow -> Int -> IO () Source #

Individual migration steps. Pattern-match on the source version.

reencodeAllEvents :: Connection -> ReencodeRow -> IO () Source #

Re-encode all event rows using the given ReencodeRow function (the version 1 JSON to version 2 CBOR migration). Rows are processed in batches of ascending event id so memory stays bounded for large databases. Runs inside the caller's transaction.

createEventsTable :: Connection -> IO () Source #

selectLastEventId :: Connection -> IO [Only Word64] Source #

getEventsASC :: Connection -> IO Statement Source #

insertEvent :: Connection -> (Word64, ByteString) -> IO () Source #

insertEvents :: Connection -> [(Word64, ByteString)] -> IO () Source #

deleteAllEvents :: Connection -> IO () Source #

backupDatabase :: Connection -> FilePath -> Word64 -> IO () Source #

Archive the current database before rotation removes the events, into an old-state subdirectory next to the database, with the log id inserted before the extension (e.g. old-state/hydra-42.db). Uses VACUUM INTO so the snapshot reflects all committed (WAL) data in a single self-contained file, regardless of WAL checkpoint state. The destination is removed first if present (e.g. a re-rotation at the same log id), since VACUUM INTO requires it not to exist.