Databases and engines¶
A Database is a prefix in an object store, or a local directory, holding versioned manifests that
map index names to segment files. A collection is an index of the same name, so the client and these
building blocks work on the same databases.
Opening a database¶
open_database(url) opens a database without a client, and creates a local one if it does not exist. It
takes the same URLs as connect, plus options, cache_dir and the
build options.
The database API is async in Rust and runs on Tokio. From synchronous code, completr::block_on runs one
of its futures on completr's internal runtime; do not call it from async code.
Transactions¶
A transaction changes one or more indexes and commits them atomically, by writing the next manifest create-only.
Transaction method |
Does |
|---|---|
append(index, documents, deletes=(), vectors=None) |
Adds a segment of upserts and deletes to an index, creating the index if needed. |
overwrite(index, documents, vectors=None) |
Replaces the whole index with one base segment. |
drop_index(index) |
Removes an index. |
set_max_score(index, max_score) |
Pins the score that normalises to 1.0. |
set_metadata(key, value) |
Sets or, with None, removes a string in the manifest. Collections keep their settings here. |
strict(True) |
Fails with ConflictError instead of rebasing when an index it changes was changed since read_version. |
max_retries(n) |
Retries of the create-only manifest write before giving up with ConflictError (32). |
commit() |
Builds and uploads the segments, then commits; returns the new manifest as a dict. |
A commit that loses a race rebases onto the newer version and retries. In strict mode it fails instead:
from completr import ConflictError
txn = db.begin()
txn.strict()
other = db.begin()
other.append("products", [{"id": "fryer", "text": "Air Fryer", "popularity": 0.8}])
other.commit()
txn.overwrite("products", [{"id": "cards", "text": "Playing Cards"}])
try:
txn.commit()
except ConflictError as error:
print(error)
db.versions(), db.latest_version(), db.index_names(version) and db.manifest(version) read the
history. db.begin(version) starts a transaction from an older version.
Snapshots¶
db.open_index(name, version=None) returns one version of one index as an Index, for scripts and tests:
Engines that follow a database¶
A serving process asks the database for an engine. The engine loads every index now, and each call to
sync() loads the newest version, downloading only segments it does not already hold. It returns the new
version, or None when nothing changed.
In Rust, the Replica is separate from the Engine it publishes to; in Python, db.engine() combines
them, and Replica(database, engine) is also available. Nothing syncs a low-level engine for you: call
sync() periodically, for example from a thread:
import threading
import time
def follow(interval=5.0):
while True:
time.sleep(interval)
engine.sync()
threading.Thread(target=follow, daemon=True).start()
A replica lists only manifests newer than its version, downloads only new segments, rebuilds only changed
indexes, and publishes them atomically. With cache_dir, it removes cached files no current index uses
after each sync.
Compaction¶
db.compact(index) merges small delta segments tier by tier: fanout (4) same-level segments merge into
one of the next level. An index is rebuilt into a single base when more than max_hidden_fraction (0.25)
of its documents are superseded or deleted, or when there are more than max_segments (16) segments and
no tiered merge is available. Transactions you commit directly are not compacted; run db.compact after
them:
Cleanup¶
db.cleanup(keep_versions=10, older_than_seconds=3600) deletes old manifests, keeping the newest
keep_versions, and segment files that no retained version references. It only touches objects older
than older_than_seconds, measured on the store's clock. See
Serverless deployment.
Leases¶
A lease is an expiring lock stored in the database, for jobs that only one process should run at a time.
acquire_lease(name, owner, ttl_seconds) returns a Lease, or None while another owner holds it.
Its generation grows with every new holder, so it serves as a fencing token.
lease = db.acquire_lease("nightly-rebuild", "worker-1", ttl_seconds=60.0)
print(lease.generation, db.acquire_lease("nightly-rebuild", "worker-2", ttl_seconds=60.0))
lease.renew(60.0) # False once the lease was lost
lease.release()
The ingestor uses a lease named ingestor.
Memory¶
While a replica switches versions, the old and new segments of an index are both mapped. When a database
holds many indexes, pass group_separator to db.engine(). The replica then switches indexes group by
group, grouping by the part after the last separator, and releases replaced segments before it loads the
next group. Serving memory holds at most one group twice, never the whole database.
txn = db.begin()
txn.append("shared/en", [{"id": "ps5", "text": "PlayStation 5 Console", "popularity": 0.9}])
txn.append("acme/en", [{"id": "ps5", "text": "PlayStation 5 Console Bundle", "popularity": 0.9}])
txn.commit()
engine = db.engine(group_separator="/") # "acme/en", "shared/en": all "en" indexes switch together
print([(s.text, s.layer) for s in engine.complete("play", ["shared/en", "acme/en"])])
asyncio¶
open_database_async returns an AsyncDatabase whose storage operations are awaitable. They run in
worker threads through asyncio.to_thread.
import asyncio
from completr.lowlevel import open_database_async
async def main():
db = await open_database_async("./data")
txn = await db.begin()
txn.append("products", [{"id": "kb-1", "text": "Wireless Keyboard", "popularity": 0.9}])
await db.commit(txn)
engine = await db.engine()
await engine.sync() # call periodically
print([s.text for s in engine.complete("wirel", ["products"])])
asyncio.run(main())
AsyncDatabase method |
Wraps |
|---|---|
await versions(), await latest_version() |
Database.versions(), Database.latest_version() |
await index_names(version), await manifest(version) |
Database.index_names, Database.manifest |
await begin(version), await commit(txn) |
Database.begin, Transaction.commit |
await open_index(name, version, **options) |
Database.open_index |
await engine(**options) |
Database.engine, returning an AsyncEngine |
await submit(changes), await pending_change_sets() |
Database.submit, Database.pending_change_sets |
await compact(index, **policy), await cleanup(**policy) |
Database.compact, Database.cleanup |
await run_ingestor(ingestor) |
One round of Ingestor.run_once() |
An AsyncEngine's sync() is awaitable; complete and every other engine method are the synchronous
Engine methods. db.database returns the underlying synchronous Database.
Transaction.append and Transaction.overwrite build a segment on the calling thread, which blocks the
event loop for large batches: wrap them in asyncio.to_thread, or run bulk loads from a separate worker.
ChangeSet.upsert only collects documents; await db.submit(changes) builds and writes them in a worker
thread.