Skip to content

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.

from completr.lowlevel import open_database

db = open_database("./data")   # or s3://bucket/prefix, gs://..., az://..., memory:///name
use completr::Database;

let database = Database::open("./data", Vec::<(String, String)>::new()).await?;

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.

txn = db.begin()
txn.append("products", [{"id": f"p{i}", "text": f"product {i}", "popularity": 0.5} for i in range(10_000)])
manifest = txn.commit()
print(manifest["version"], list(manifest["indexes"]))
use completr::Document;

let mut txn = database.begin().await?;
txn.append_documents(
    "products",
    (0..10_000).map(|i| Document::keyed(format!("p{i}"), format!("product {i}"), 0.5)),
    [],
)?;
txn.commit().await?;
1 ['products']
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)
commit conflict: products changed since version 1

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:

snapshot = db.open_index("products", version=1)
print(len(snapshot), snapshot.get("fryer"))
10000 None

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.

engine = db.engine()
print(engine.version, engine.names())
print([s.text for s in engine.complete("air", ["products"])])
use completr::{Engine, IndexOptions, Replica};

let engine = Engine::new();
let replica = Replica::new(database.clone(), IndexOptions::default());
replica.sync(&engine).await?;
let hits = engine.complete(&["products"], "air", 10);
2 ['products']
['Air Fryer']

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:

print(db.compact("products", until_done=True))   # the new version, or None if nothing was due
None

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()
1 None

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"])])
[('PlayStation 5 Console Bundle', '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())
['Wireless Keyboard']
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.