Skip to content

Many writers

Every write through a client is a commit, and commits from many processes contend on the next manifest. For many processes writing small changes at a high rate, the building blocks have an inbox instead: any process submits a ChangeSet without committing it, and an Ingestor commits pending change sets in order. Only the process holding the ingestor lease acts, so many processes can run one safely.

Run ingestors outside serving processes

Building segments and compacting indexes use far more memory and CPU than serving. Build and compaction peaks belong to the ingestor, so run ingestors in separate processes (a worker, a cron job or completr ingest), not inside the processes that answer completion requests.

Role Runs Where
Serving engine = db.engine(), then engine.sync() every few seconds; engine.complete(...) per request, or a client from connect Every process that answers completions
Submitting db.submit(changes) Any process that changes data: an API handler, a queue consumer, a batch job
Ingesting Ingestor(db, owner).run_once() in a loop One or more worker processes, outside serving processes
Maintenance db.cleanup(...) on a schedule A cron job or the ingestor's worker

Submitting changes

A ChangeSet holds upserts and deletes for one or more indexes, applied together. db.submit(changes) writes it to the database's _inbox/ and returns its id. It does not wait for a commit. Change sets from one process apply in submission order, and later change sets win per document id.

from completr.lowlevel import ChangeSet, Ingestor, open_database

db = open_database("./data")
changes = ChangeSet()
changes.upsert("products", [{"id": "kb-1", "text": "Wireless Keyboard", "popularity": 0.9, "contexts": ["peripherals"]}])
changes.delete("products", ["p42"])
change_set_id = db.submit(changes)
print(db.pending_change_sets())
use completr::{key_id, ChangeSet, Document};

let mut changes = ChangeSet::new();
changes
    .upsert("products", [Document::keyed("kb-1", "Wireless Keyboard", 0.9)])
    .delete("products", [key_id("p42")]);
database.submit(changes).await?;
1

ChangeSet.upsert takes vectors= like Collection.add. A change set can also override a shared index for one tenant, as layers do:

changes = ChangeSet()
changes.upsert("acme", [{"id": "psvr", "text": "PlayStation VR2 Headset"}])   # rename for acme only
changes.delete("acme", ["cards"])                                              # hide for acme only
db.submit(changes)

Running ingestors

An Ingestor commits pending change sets while it holds the database's ingestor lease. Every round takes or renews the lease, folds up to max_change_sets change sets into one commit, deletes them from the inbox, and compacts the indexes it touched (disable with compact=False). You may run several ingestors for availability: only the lease holder acts, and another takes over within lease_ttl_seconds (30 by default) when it stops.

ingestor = Ingestor(db, "ingestor-1")
print(ingestor.run_once())
print(ingestor.run_once())
ingestor.release()
use completr::{IngestStep, Ingestor};

let mut ingestor = Ingestor::new(database.clone(), "ingestor-1");
if let IngestStep::Committed { version, .. } = ingestor.run_once().await? {
    println!("committed version {version}");
}
ingestor.release().await?;
{'step': 'committed', 'version': 1, 'change_sets': 2, 'documents': 2}
{'step': 'idle'}
step Meaning
standby Another ingestor holds the lease.
idle The inbox was empty.
committed Change sets were committed; the result also has version, change_sets and documents.

In a real deployment, the ingestor runs in its own process and calls run_once() in a loop:

import time

ingestor = Ingestor(db, "worker-1", lease_ttl_seconds=30.0)
try:
    while True:
        step = ingestor.run_once()   # {'step': 'standby' | 'idle' | 'committed', ...}
        if step["step"] != "committed":
            time.sleep(1.0)
finally:
    ingestor.release()

Commits carry the lease generation as a fencing token, so an ingestor that lost its lease cannot commit, and a change set left over after a crash is never applied twice. Call run_once() well within lease_ttl_seconds: it renews the lease once a third of the TTL has passed, and loses it after the TTL. The command-line tool runs the same loop: completr s3://my-bucket/completions ingest --interval 1.

Serving processes see the commits on their next sync, whether they use an engine from db.engine() or a client from completr.connect:

import completr

print([(s.id, s.text, s.kind) for s in completr.connect("./data")["products"].complete("wirel", contexts=["peripherals"])])
[('kb-1', 'Wireless Keyboard', 'prefix')]

Under asyncio

AsyncDatabase.run_ingestor runs one round in a worker thread:

import asyncio

from completr.lowlevel import Ingestor, open_database_async

async def ingest():
    db = await open_database_async("s3://my-bucket/completions")
    ingestor = Ingestor(db.database, "worker-1")
    try:
        while True:
            step = await db.run_ingestor(ingestor)
            if step["step"] != "committed":
                await asyncio.sleep(1.0)
    finally:
        ingestor.release()

A runnable Rust example of this flow is live_updates.rs.