Skip to content

Standalone recordstore module #75

Description

@moshloop

Plan: standalone recordstore module + multi-process active-writer model

Context

recordstore/ (root module) and cmd/query/recordresults/ (cmd/query module) together form the record-stream store: backends, index, follow, and result-type profiles. Two problems:

  1. Packaging. recordresults lives inside the cmd/query binary module, so library consumers must depend on the CLI module to serve result types. recordstore lives in the root module although nothing in root imports it. Both belong in one nested module, github.com/flanksource/commons-db/recordstore.
  2. Single-process writes. The contract is "one writer per stream, serialized in-process" (recordstore/recordstore.go:19-20). Every guard is process-local: StreamLocks, the sqlite.DB RWMutex/Lease, the per-process kind-table cache (sqlite/tables.go:49-79), the sweeper in every process, and Notifier wake-ups. Two query serve processes, or a CLI writing a cache while a server runs, race on records.sqlite, catalog reconciliation and predecessor copies.

Outcome: exactly one active writer (owner) per logical store, elected by a flock. It advertises itself in a state (pid) file, with a control socket, modelled on gavel/procfile (lock.go, state.go, control.go). Every other process, long-lived or short-lived, publishes NDJSON / NDJSON.gz / Parquet batches into the spool folder the state file names. The owner ingests and upserts them into the live db exactly once. Long-lived non-owners promote themselves when the owner exits.

Decisions (user-confirmed)

  • Embedded owner. The first process to take the flock is the owner, in-process. There is no daemon.
  • Upsert = per-kind policy. KindOptions.OnConflict: Skip|Replace, default Skip (today's behaviour). Replace deletes the old keyed row and re-appends it at the next seq in the same transaction.
  • Separate go.mod only. The new module requires root commons-db. Removing recordstore's query import is a follow-up TODO.

Design rules derived from review (do not re-derive)

  • Correctness rests only on files and SQLite: flock liveness, spool directories committed by rename, and a ledger table committed in the same transaction as the data. The control socket is a latency and observability aid. A dead or unreachable socket never loses or duplicates data.
  • One *sqlite.Backend per process for its lifetime. It switches read-only → writable in place (Promote). Never swap backend pointers: Registry.index (registry.go:149), NewIndexer identity (indexer.go:64) and Results.location (streamref.go:122) all depend on identity.
  • Lock scope = logical store file, not versioned. Derive every path from the configured path P (for example <dir>/records.sqlite), not from <dir>/v5/…. There is one owner across build versions, and the state file tells non-owners which versioned file is live and its catalog version. That is how an older CLI build spools into a newer server's db.

Part A — Module move (3 commits)

Pre-flight:

  • HEAD is ac17f48, 14 ahead of origin. Remote tags reach v0.1.37, which contains root recordstore/. Every push to main auto-tags (release.yml).
  • So the new module must require root >= v0.1.38, and the move must ship in one push. Pushing is the user's or gavel's call. Never push from here.

A1 refactor(query): move profile store/read-hook contracts to query/profilestore

This breaks the module cycle: recordresults → cmd/query/profiles and cmd/query/sessions → recordstore.

  • New query/profilestore/{store.go,read.go}:
    • store.go: Store, VirtualStore, UpdateOptions.
    • read.go: ReadRequest, BeforeExecuteFunc, PrepareReads, ErrProfileDataNotFound, ErrProfileRequestInvalid.
    • Move them from cmd/query/profiles/{store.go,overlay_store.go,store_update.go,before_execute.go}. No aliases.
  • Sweep:
    • First gopls rename profiles.Store → a temporary MovedStore. A plain gopatch rename hits atomic.Bool.Store, Options{Store:} and embedded s.Store.
    • Then a gopatch file that binds the import name. sessions/*_test.go import the package as profilepkg.
    • Call sites include internal/app/{app,runtime,schema_handler,server}.go, schedules/runner.go, sessions/{service,trace}.go, snapshots/catalog.go, recordresults/{catalog,provider}.go, and the profiles package itself.
  • Add var _ profilestore.VirtualStore = (*Registry)(nil) in registry.go.
  • Check: go list -deps ./cmd/query/recordresults | rg commons-db/cmd/query/ shows only itself.

A2 test(recordresults): share fixtures via recordresultstest; move HTTP e2e to cmd/query/recordresultse2e

  • New non-test package recordresults/recordresultstest (same pattern as recordstoretest). It exports the shared fixtures:
    • SampleEvent*, NewKV, NewRegistry*, KVRouter, ForTenant, OpenResults, JobEvent, RegisterJobTypes, InvocationResult, …
    • They come from e2e_test.go:36-104, registry_test.go:32-60, open_e2e_test.go:20-67, views_e2e_test.go:20-91 and hierarchy_e2e_test.go:17-22.
  • Stay in the module: registry_test, views_registry_test, streamref_e2e_test, benchmark*_test.
  • Split so the non-HTTP specs stay, as hierarchy_registry_test, search_registry_test and open_test: hierarchy_e2e, search_e2e, open_e2e.
  • Move (git mv) to cmd/query/recordresultse2e/ (package recordresultse2e, plus a new suite): e2e, e2e_filters, e2e_window, follow_e2e, follow_routing_e2e, views_e2e, views_nulls_last_e2e, and the HTTP halves of the three split files.
  • suite_test.go in recordresults blank-imports query/providers. The library itself does not: capture-only binaries must not pull every provider SDK. provider.go:133 already fails with "link query/providers".

A3 refactor: move recordresults into the recordstore nested module

  • git mv cmd/query/recordresults recordstore/recordresults, then gopatch the import path across ./recordstore ./cmd/query, including gitignored cmd/query/hack/docs-tutorial. Fix comments in query/param_validate.go:43 and recordstore/probe/manager.go:87.
  • recordstore/go.mod:
    • module github.com/flanksource/commons-db/recordstore, go 1.26.1.
    • Requires: root v0.1.38, clicky, commons, atlas, uuid, ginkgo/gomega, and for tests miniredis, clicky/valkey and valkey-go. Copy versions from the root go.mod before tidying.
    • The glebarez/sqlite => clarkmcc/gorm-sqlite replace plus its comment block (root go.mod:404-411).
    • replace github.com/flanksource/commons-db => ..
  • cmd/query/go.mod:
    • require …/recordstore v0.1.38
    • replace …/recordstore => ../../recordstore
  • Root go.mod: go mod tidy. It drops miniredis, clicky/valkey and valkey-go.
  • Build wiring:
    • Root Makefile: GO_MODULES := . recordstore cmd/query, plus loop targets tidy, tidy-check (go mod tidy -diff), build, vet, test.
    • New recordstore/Makefile.
    • Taskfile.yaml: recordstore:{build,lint,test}.
    • Local go work use ./recordstore.
  • CI/release:
    • test.yml: go work init . ./recordstore ./cmd/query (lines 44 and 131), recordstore/go.sum in the cache paths (39-41 and 126-128), comments, and a new make tidy-check step.
    • release.yml: a recordstore/vX tag step before the cmd/query steps (28-34).
  • Docs:
    • start/consuming.md: three modules, the version table, the providers note, go work use, and "never mix root ≤ v0.1.37 with this module".
    • index.mdx:50-56, start/tutorial.md:47,124,139, recordresults/{open,serving}.md, cmd/query/README.md:139-151.
    • Fix the stale sqlite/record-file.md (catalog v4 → v5).

Part B — Multi-process writers

B.1 Files per logical store (configured path P, for example <dir>/records.sqlite)

Path Role
P.lock gofrs/flock (already a root dependency, db/embedded_lifecycle.go:115). The only liveness truth. The kernel releases it on crash, so pid reuse is irrelevant.
P.owner.json Advisory state, written with a unique temp file, fsync and rename. Fields: {instance, pid, host, build, phase: starting|ready|draining, startedAt, heartbeatAt, store: <dir>/v5/records.sqlite, catalogVersion, spool: {dir, manifestFormat: 1, formats: [ndjson, ndjson.gz, parquet]}, socket, backlog, failed, lastError}
P.sock Control socket. If the path is over 100 bytes (sun_path limit), it goes in $TMPDIR/recordstore-<sha1(P)[:12]>.sock; the actual path is always recorded in state.
P.spool/{tmp,incoming,failed,trash} Spool, mode 0700. Not versioned by format: an owner that can't read a manifest fails it loudly into failed/ instead of stranding it.
<dir>/v5/records.sqlite The live db (existing versionedFile). Only the owner runs copyPredecessor and catalog creation.
  • Refuse network filesystems at Open: statfs NFS/SMB/CIFS/9p/virtiofs/FUSE/Ceph on Linux; Fstypename nfs/smbfs/afpfs/webdav/*fuse on darwin. WAL and flock are unreliable there.
  • ndjson and kv-derived indexes: see B.3/P6.

B.2 Roles and lifecycles

  • Open (every process): MkdirAll → fs-type check → TryLock(P.lock).
    • Won → owner:
      1. state starting.
      2. sqlite.Open writable: predecessor copy, catalog, and the ledger via CREATE TABLE IF NOT EXISTS.
      3. Serve the socket.
      4. state ready.
      5. Ingest loop: startup drain, poll incoming at 50ms→500ms with backoff, woken early by a socket ingest poke.
      6. Heartbeat every 5s and the sweeper (owner only).
    • Lost → non-owner:
      1. Read state and wait for phase=ready (bounded by StartWait, default 60s; copyPredecessor can be slow).
      2. If state.catalogVersion != this build, writes still spool (manifests are version-agnostic), but reads fail fast with ErrCatalogVersion{Owner, Build}.
      3. Otherwise open the backend read-only.
      4. Start the promoter: TryLockContext(ctx, 100ms), cancellable on Close.
  • Promotion: lock granted → state starting → backend.Promote(ctx) → serve socket → state ready → drain the spool. This process's in-flight Awaits resolve from the ledger.
  • Owner Close:
    1. Stop the socket.
    2. Final drain, bounded by DrainOnClose.
    3. Remove state and socket, only if they carry our instance.
    4. Unlock.
    5. Post-unlock recheck: while incoming is non-empty and TryLock succeeds, drain and unlock. This closes the publish-after-final-drain race.
  • Same process opening P twice: an in-process registry shares one handle, reference-counted, so the process never spools to itself.

Long-lived non-owner (second query serve, probe host):

  • Meta, Scan and profile reads go through the read-only pool on the live file. WAL makes cross-process reads safe.
  • Append and Seal/Expire/Trim/Delete:
    1. Validate locally (ValidateAppend, ResolveKind, RowKeys, stored-row conversion, an optimistic Meta seal/kind check).
    2. Publish a one-entry batch.
    3. Poke the socket.
    4. Await the ledger outcome, polling the read-only pool every 25ms gated by PRAGMA data_version. The owner's socket ingest{wait} reply may short-circuit this.
    5. Return the exact AppendResult or sentinel error (ErrSealed, ErrCapacity, ErrNotFound, ErrSchemaConflict).
    6. If ctx ends first, return *PendingError{BatchID} (errors.Is(err, ErrPending)): the rows are durable and resumable with Store.Await. Callers must not blindly retry unkeyed kinds.
  • Followers wake through the new ChangeSource. The read-only side polls PRAGMA data_version every 150ms; the owner fires on every ingested stream.

Short-lived writer (a CLI writing a cache file, then exiting):

  • It uses the bulk store.Writer(spool.WriterOptions{Format}), calling Append/Seal then Flush(ctx, {Wait}). Flush is role-agnostic:
    • As owner, it calls AppendBatch directly, and Close drains other processes' spool files before releasing.
    • As non-owner, it publishes one multi-entry manifest and pokes the socket. With Wait false the CLI exits immediately; the batch is durable and the owner or the next owner ingests it.
  • Close in a non-owner with DrainOnClose > 0: if TryLock now succeeds (the owner vanished), drain before exit.

B.3 Spool format (recordstore/spool, pure, no sqlite)

  • Publish:
    1. Write tmp/<instance>-<id>/ containing the data files and manifest.json.
    2. fsync the files and the directory.
    3. Rename the directory to incoming/<createdNanos:020>-<instance>-<producerSeq:012>-<id>. This is the single commit point.
    4. fsync incoming.
  • Manifest DTO (explicit json tags, string enums): {format:1, id, producer{instance, seq, pid, host, build}, created, schemas:[{kind, columns:[{name,type}], key, retention, onConflict}], entries:[{op: append|seal|expire|trim|delete, stream, kind, file, format, rows, sha256, seal, generation, ttl, before}]}
    • Carrying schemas lets the owner ingest from builds whose kinds differ.
    • Zero-row entries are valid (probe.Arm appends nil rows).
  • Codecs: ndjson and ndjson.gz, decoded with UseNumber and no line-length limit. Also Normalize(schema, row), mirroring sqlitetable.Value (table.go:302) so spooled and direct appends store identical values (time in string columns, no double-encoded JSON).
  • Caps: the bulk Writer splits at 50k rows / 64MiB. A single Append is always one manifest (atomic). The owner has a hard cap and sends anything larger to failed/.
  • GC (owner): tmp/* older than 24h by mtime, failed/* after 30d, trash/* always.
  • Producer check: before publishing, fail fast unless state.spool.manifestFormat >= 1 and the chosen format is in formats.

B.4 Ingest (owner) — recordstore/sqlite batch capability

  • Transaction: one sqlite tx per manifest, with a SAVEPOINT per entry. A deterministic entry error (sealed, not_found or generation mismatch, conflict, invalid) rolls back only that entry and is recorded in its result. Manifest-level errors → failed/ plus error.json, and the ledger records failed so waiters fail at once.
  • Ledger: record_spool_batches(batch_id PK, producer, producer_seq, outcome, results JSON, ingested_at) is committed in the same tx.
    • After commit, rename the batch directory to trash/. A crash between the commit and that rename means the batch is re-seen, the ledger hit makes it a no-op, and it goes to trash.
    • Sweep ledger rows after 7d, never while their directory still exists.
  • Ordering fence per producer: ingest (instance, n) only when n == last+1 or the producer is new. Otherwise defer; after 60s, log an error and proceed. Across producers, ingest order applies.
  • Transient errors (ctx, SQLITE BUSY/LOCKED/FULL/IOERR/NOMEM): halt with backoff from 100ms to 30s. After 20 attempts or 10min, send to failed/.
  • appended_at = ingest time (keeps trimTx's monotonic assumption). Document that retention starts at ingest.
  • Foreign schemas: tableFor(foreignSchema) reconciles additively (catalog v5 already supports added columns) before the tx. A key or storage-type mismatch gives conflict (ErrSchemaConflict). The owner's Schemas registry is never mutated.

Implementation phases (each ships and verifies independently)

P0 — sqlite refactor, no behaviour change (recordstore/sqlite/{sqlite,trim}.go, root sqlite/db.go)

  • Split appendTx(ctx, tx, table, stream, write, now) out of commitAppendLocked (sqlite.go:217-252). Add tx-taking sealTx, expireTx and trimStreamTx.
    • Reason: Append does three separate database.Write calls on a non-reentrant mutex with MaxOpenConns(1), so nesting deadlocks.
  • trimBelowTx: Total -= RowsAffected instead of the dense-range recompute (trim.go:97).
  • Writer DSN gains _txlock=immediate (db.go:100), so any cross-process writer serialises instead of returning SQLITE_BUSY_SNAPSHOT.
  • Export sqlite.CatalogVersion and sqlite.VersionedPath(configured).

P1 — OnConflict: Replace

  • schemas.go: OnConflict enum; Replace requires Key.
  • recordstore.go:
    • AppendResult.Replaced, ErrUnsupported, ErrSchemaConflict.
    • Rewrite the "contiguous" and "Total = HighSeq-LowSeq+1" docs. New invariant: the row at HighSeq exists; LowSeq is a lower bound.
  • sqlite/replace.go: in appendTx, delete by key (summing replaced), insert at HighSeq+1…, Total += n - replaced.
  • kv and ndjson refuse Replace kinds with ErrUnsupported: dense seqs in kv/rows.go:317, chunks.go:180 and ndjson line-per-seq.
  • Indexer gap tolerance:
    • copyAfter (indexer.go:219) accepts strictly increasing seqs and checks that it reached source.HighSeq.
    • ImportRequest.Seqs.
    • index.go Import deletes indexed rows by key first (index.go:132).
  • recordresults.ResultType.OnConflict.
  • New recordstoretest/conformance_replace.go, gated by a harness flag.
  • Docs: kinds.md, concepts.md.

P2 — batch ingest capability (recordstore/batch.go, sqlite/{batch,ledger,foreign}.go)

  • Types: Batch, BatchEntry, BatchOp, Producer, BatchResult, BatchError{Code} with Unwrap to sentinels, and a BatchAppender{AppendBatch; BatchOutcome} capability interface. It is not added to the base Backend, matching the optional-capability approach of TODO 2b7746b9.
  • AppendBatch:
    1. Pre-resolve every kind table outside the tx.
    2. Lock the distinct streams in sorted order.
    3. Open one tx with a SAVEPOINT per entry.
    4. Write the ledger row.
    5. Fire change notifications.

P3 — recordstore/spool (dir, manifest, publish, load, codec, codec_ndjson, normalize, writer, gc.go)

P4 — read-only mode + change feed

  • Root sqlite.Options.ReadOnly, ErrReadOnly, DB.EnableWrites(), DB.WatchDataVersion(ctx, every, fn).
  • recordstore/sqlite/readonly.go:
    • Read-only Open requires an existing versioned file, inspects the catalog without writing, and runs no sweeper.
    • Prepare skips purge.
    • Table(kind) adopts from the catalog. If the table or a declared column is missing, submit a declare-only batch and re-adopt once. This also fixes Registry.register (registry.go:373), which writes today.
    • Mutations delegate to Options.Submit. Promote(ctx).
    • Refuse ReadOnly && Derived.
  • sqlite/changes.go plus recordstore.ChangeSource; notify.go subscribes and gains wakeAll.

P5 — recordstore/owner (lock, fstype_{linux,darwin,other}, state, control, store, promote, ingest, classify, client, gc, inprocess.go)

  • Open[T Target](ctx, Options[T]{Path, Open func(ctx, readOnly bool, submit Submitter) (T, error), Format, StartWait, DrainOnClose, Poll, Build}).
  • Store.{Backend, Role, OnRole, Await, Writer, Status, Close}.
  • Exclusive(path) returns *LockedError{Path, Owner *State} (errors.Is ErrLocked). ReadState.
  • Control socket (port of gavel control.go):
    • One JSON request/response per connection, protocol:1, socket mode 0600, 10s deadlines.
    • Actions: status and ingest{ids, wait}.
    • A stale socket is removed when a dial fails. Teardown unlinks the socket only if status returns our own instance.

P6 — recordresults adoption (recordstore/recordresults/open.go)

  • sqlite mode: owner.Open[*sqlite.Backend]; store.Backend() stays both source and index. Add Results.Role().
  • ndjson mode: owner.Exclusive(<dir>/ndjson) before opening; a second process gets ErrLocked carrying the owner's pid.
  • Source/kv mode: Exclusive(index.sqlite). A non-owner uses a private <dir>/v5/private/<instance>/index.sqlite, GC'd when its lock is free. Coordinate with TODO 62d30eba (shared index across routes).
  • Router routes that open sqlite files wrap with owner.Open per route file.
  • Docs: replace the "single-writer contract" section (concepts.md:111), and update backends.mdx, sqlite/record-file.md and a new recordstore/multi-process.md.

P7 — probe identity across processes

  • probe.ManagerOptions{LockDir}: Arm holds TryLock(<LockDir>/<sha256(identity)[:32]>.lock) for the run and returns ErrAlreadyManaged if the lock is held.

P8 — Parquet

  • recordstore/spool/parquet (an optional import whose init registers the codec) on github.com/apache/arrow-go/v18 parquet/pqarrow.
  • Type mapping:
    • string / uuid / status → UTF8.
    • number → INT64 if every value is integral, else DOUBLE. Reads accept INT32/64, FLOAT, DOUBLE and DECIMAL, and never pass int64 through float.
    • datetime → TIMESTAMP(NANOS, UTC). Reads accept ms/µs/ns.
    • json / key_value → JSON bytes.
  • Start with a spike test that verifies timestamp decoding and the v18 FileReader/GetRecordReader API names.

Risks and follow-ups

  • Release sequencing:
    • The root sqlite changes (P0, P4) ship in the root module. They are local via replace, but every tag of the three modules must come from one push.
    • The catalog: ledger and Replace land in v5 without a version bump, but only if v5 is still unpushed when P1/P2 land. Otherwise bump to v6.
  • Performance: a synchronous spooled Append costs about 50–100ms. Hot non-owner writers should use the bulk Writer. Measure before adding fsnotify.
  • WAL growth: Scan holds one read tx across callbacks (rows.go:86-127), so a slow cross-process Tail pins the WAL. Follow-up: one read tx per page plus an idle wal_checkpoint(TRUNCATE).
  • Windows (gofrs LockFileEx, AF_UNIX on Win10+) is untested. Label it experimental until CI covers it.
  • Tracking: after approval, create gavel TODOs: (1) Part A module move; (2) P0–P3 Replace + batch ingest + spool; (3) P4–P7 read-only, owner and adoption; (4) P8 Parquet; (5) follow-up: drop recordstore's query dependency. Link them to 2b7746b9, 2aac1d16, 5dbb5241, 1ee7ebb3 and 62d30eba, which all touch catalog, open.go or index layout.

Verification

  • Module move:
    • make tidy-check vet build test, with GOWORK=off, so the replaces must resolve.
    • task recordstore:test.
    • No cycle: cd recordstore && GOWORK=off go list -deps -test ./... | rg commons-db/cmd/query is empty; the root go list -deps -test ./... | rg 'commons-db/(recordstore|cmd/query)' is empty.
    • No ambiguous import: cd cmd/query && GOWORK=off go list -f '{{.Module.Path}}' github.com/flanksource/commons-db/recordstore/....
    • A CI-identical worktree: go work init . ./recordstore ./cmd/query && gavel test --dry-run shows 4 roots.
    • gavel test --ignore ./e2e --bench=^$.
    • task query:build, then run .bin/query --help.
    • rg cmd/query/recordresults is empty.
  • Unit (ginkgo, per phase):
    • Replace: Total/HighSeq, a replaced LowSeq row, replace then trim, RetainRows renewal, Router→sqlite→derived index mirroring with no UNIQUE violation, follow sees the new seq.
    • Batch: idempotent by ID, entry isolation, foreign column added, key/type conflict, zero-row entry, generation fence, and a deadlock guard (a new kind inside a batch completes within 2s).
    • Spool: round-trip for every ColumnType including int64 max; the validation matrix; tmp never listed.
    • Normalize parity: the same rows direct vs spooled give identical raw columns.
    • Error-classifier table; fs-type magic table.
  • Multi-process integration (recordstore/owner): TestMain re-execs the test binary when RECORDSTORE_CHILD is set, running a JSON step script and emitting JSON lines. The parent drives children via gexec. RECORDSTORE_FAILPOINT names one of after-ledger-commit, before-trash, after-final-drain or during-starting; the child SIGKILLs itself there. Matrix:
    1. Owner SIGKILLed mid-append; the spooler promotes and returns exact contiguous windows.
    2. A CLI child appends while the owner runs.
    3. A CLI with no owner drains a pre-existing backlog, then releases.
    4. A crash after ledger commit: no re-ingest.
    5. Extra-column vs changed-type producer.
    6. 3 spoolers and a killed owner: exactly one owner and no lost batches.
    7. Publish at after-final-drain is still ingested (post-unlock recheck).
    8. A sealed stream in a bulk batch: other entries commit.
    9. Generation fence after delete + recreate.
    10. 50 async appends then a seal: order is preserved.
    11. format:2 manifest → immediate failure plus failed/.
    12. A non-owner Tail wakes in under 1s.
    13. Opening during starting waits.
    14. tmp GC by age.
    15. A second ndjson process gets ErrLocked with the pid.
    16. kv: two processes, private index.
    17. Socket unreachable (removed): ingestion still completes via the poll.
  • Run with gavel test ./recordstore/.... -race goes on the owner/sqlite packages.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

No labels
No labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions