Skip to content

Tracking: move multitenancy from sharded (one process per collection, per-tenant tables and shard keys) to non-sharded (shared structures served by any process) #1574

Description

@edwinyyyu

Where to start (status 2026-10-02)

Storage layer. Segment store: done (#1661 on main). Qdrant and Milvus: #1733 -> #1734 -> #1735 -> #1736 in review (SQL collection registry, an incarnation per collection life, fenced handles, claimed purge); they close #1524, #1525 and #1563 on merge. SQLite stores: #1536 stays with the SQLite store fixes (#1460 to #1675), and they are HOST or PROCESS scope by declaration (#1757).

Above the stores, in priority order. Each is a tracking issue with a one-glance table (what is wrong, best fix, order) and a PR split, and each child issue shows its tracker at the top.

  1. Tracking: make the session lifecycle safe across replicas with a registry row and durable jobs #1755 Session lifecycle: the sessions row as registry plus durable jobs. Tracks Session data manager loses concurrent config updates and leaks IntegrityError on concurrent create #1543, Stale segment-store handles need cache eviction and an API status mapping #1571, Project deletion never completes when semantic memory is disabled: the deletion worker requests the semantic service unconditionally and gives up #1575, A search or add on a nonexistent project creates it (session row, partition and collection), so deletion is not final and the search 404 path is unreachable #1576, Deleting a project that is in use returns 204 and then never completes: the deletion worker gives up on SessionInUseError, leaves a partial deletion, and only a restart retries #1577, semantic_memory.enabled is honored only at startup: search, add and list still build and query the semantic stack whenever a request names the semantic type #1600, Session lifecycle is single-process: locks, in-use checks and row-vs-storage creation are not arbitrated across replicas #1655, [Bug]: Deleting a session whose memory is not open creates its storage first, only to delete it #1728, Session deletion runs in a per-process queue: a crash loses the job until a restart, and every replica re-enqueues every Deleted session at boot #1749.
  2. Tracking: background work across replicas with one claimable-job primitive (ingestion, deletion, purge) #1756 Background work: one claimable-job primitive for provision, delete, ingestion and purge. Tracks Semantic background ingestion runs in every replica with no claim, so N replicas ingest the same messages N times #1745, Session deletion runs in a per-process queue: a crash loses the job until a restart, and every replica re-enqueues every Deleted session at boot #1749, [Feat]: Make the semantic ingestion poll interval configurable (currently hardcoded at 2s) #1699, Episode add is not atomic across the store, episodic memory, and semantic history; partial failures are unrecoverable without a full scan #1738, Add SQL leased read-write locks #1722.
  3. Tracking: declare a concurrency scope per component so single-process components stay usable without blocking multi-replica deployments #1757 Concurrency scope per component, so declarative memory and short-term memory stay available for single-replica deployments without a routing layer. Tracks Short-term memory is process-local: recent episodes are never reloaded, replicas diverge, and the persisted summary is last-writer-wins #1748, Declarative memory and the Neo4j semantic backend keep per-process caches that are not refreshed when another replica writes #1753, Semantic background ingestion runs in every replica with no claim, so N replicas ingest the same messages N times #1745.
  4. Tracking: replica-safe defaults in one change (count cache, config cache, MCP, config API, chart) #1758 Replica-safe defaults in one change. Tracks Episode count cache is per process and on by default, so episode counts go stale across replicas #1746, Semantic config cache is per process and on by default: stale categories across replicas, and a stale replica overwrites a set's embedder and language model #1747, [Bug]: Semantic session manager and semantic service each build their own config cache, so set type deletes and registrations do not reach ingestion #1730, MCP default identity derives from the host MAC address, so a client that sends no ids gets a different project on every machine #1750, MCP over HTTP is stateful: a session id issued by one replica is unknown to the others #1751, The config API mutates one process and rewrites the whole configuration file with secrets in plaintext #1752, Deployment artifacts assume one replica: Helm hard-codes replicas: 1 with no affinity or worker count, compose cannot scale, and all replicas would share one log file #1754.

Design reference: #1579 (deferred as a whole; its vocabulary and target shapes are what the tracking issues cite).

Labels. horizontal scaling: wrong or unsafe with more than one server process on the same backends. concurrency: wrong under concurrent requests or processes regardless of replica count (races, lost updates, unarbitrated read-then-write). recovery: a failure or crash leaves durable state nothing repairs (partial writes, lost jobs, no retry or terminal state). An issue carries every label that applies.

Why this issue exists

The issues linked below were filed one at a time, most of them as bugs. Read individually they look like an accumulation of unrelated defects in the storage layer. Read together they are one change of model: the codebase was built for sharded horizontal scaling (each logical collection or partition owned by exactly one process, and separated from other tenants by a physical structure of its own), and the requirement is now non-sharded scaling (any process serves any tenant, and tenants are rows or values inside structures they share). This issue states that change once, records why the previous model was there, classifies each linked issue against it, and tracks what remains.

The model the codebase was built on, and why

Two decisions, both explicit and documented when they were made.

1. Ownership by process sharding. #1201 (merged 2026-04-02) set the VectorStore contract that is still in common/vector_store/vector_store.py today:

A given logical collection identified by a (namespace, name) pair must be managed by at most one process at a time. The consumer is responsible for sharding names across processes.

Its description states the guarantees plainly: "safe concurrent usage on a single process" and "safe concurrent sharded usage (one node per logical collection) on multiple processes". With ownership fixed per process, lifecycle operations (create, open-or-create, delete) could be read-check-write sequences guarded by per-process asyncio locks, and backends without transactions or unique constraints (Qdrant, Milvus) could keep their catalog of logical collections inside the backend itself. The multi-process guarantee was correct under its contract and was never broken: strict sharding always worked. The restriction itself is what became the limitation (#1524, #1525).

2. Per-tenant physical structures. #1205 (merged 2026-05-08) scoped the event memory "by partition key: should allow sharding and horizontal scaling", and on PostgreSQL gave every segment-store partition its own pair of child tables under LIST-partitioned parents, created and dropped by per-tenant DDL. #1292 (merged 2026-05-02) gave every SQLite vector collection its own tables, for a stated reason: "Partition keys are avoided in favor of per-collection tables, since sqlite-vec ANN indexes may not support them." Qdrant's is_distributed flag gives every logical collection its own shard key. The payoff was real: deleting a tenant was one DROP rather than a row-scanning DELETE, one tenant's rows never interleaved with another's in a heap or index, and backend index limitations were sidestepped. Note that #1201 already rejected one native collection per tenant ("most vector databases recommend against one collection per tenant, and empirical tests show that many collections induces significant overhead"); the per-tenant structure was the shard key or child table beneath a shared native collection, not the collection itself.

Both decisions fit the deployments they served: a single server process, embedded SQLite as a first-class backend, an opt-in pre-GA event backend, and tenant counts at which per-tenant DDL was cheap. Nothing below should be read as those decisions having been mistakes at the time.

What changed

Two requirements the sharded model does not meet:

Why the sharded model fails them

Measured, in the linked issues and in the design doc:

Cost of the model Evidence
Per-tenant DDL on shared PostgreSQL parents #1544's second half, and the design doc in #1545: create and delete DDL deadlocked with writers and with each other, on upstream and on every intermediate fix, in two cycle shapes; the class has nowhere to live once there is no lifecycle DDL. (#1544's first half, a CASCADE that removed the shared foreign key, is a plain bug; see the classification below.)
Per-tenant partitions and prepared statements #1546: after five executions PostgreSQL switches to a generic plan that locks every child partition per read: 9 locks per read with a custom plan, 319 with the generic plan at 316 partitions, HTTP 500 out of shared memory at 128-way concurrency. Each per-tenant table also costs about 5 catalog relations and disk files, which does not survive the required tenant counts.
Per-tenant shard keys on Qdrant #1564: 450 ms per tenant admission, 45 s for 100 tenants against 0.3 s with payload partitioning, 505 segments against 5, and cluster mode required.
Per-process ownership #1524, #1525: the moment two processes manage the same name, lock-guarded read-check-write lifecycle has no arbiter. #1566 (folded into #1564): the in-process lock table was a leftover of that serialisation.

The comparison that settled the segment store (raw asyncpg, 40 tenants x 2000 row-pairs, pgvector:pg16, 2 CPUs; full table in the design doc):

layout ingest seed read tenant create catalog cost
PARTITION OF children 8.7-9.2k pairs/s 0.43-0.46 ms ~12 ms DDL ~5 relations/tenant
shared tables 15k pairs/s 0.21 ms 0.006 ms none

Per-tenant tables won only small-tenant removal latency (9 ms drop against 58 ms row delete), which does not justify their catalog cost at the required tenant counts. This is the industry's pool model: O(1) logical deletion via a registry plus asynchronous batched reclamation; O(1) physical drops exist only in per-tenant-resource (silo) architectures, which the tenant-count requirement excludes.

The target model

The layer above the stores that ties these together (a tenant record with a step table and a reconciler, UUID-only store keys minted per tenant lifetime, a single tenant handle, and prompt purge as best-effort) is proposed in design/tenant_lifecycle.md under #1579, with the resource-level contracts the stores expose to it. The bullets below are the store-level targets that proposal builds on.

Status by component

Each row names the issues that carry the work; fixes and their state are tracked on those issues, not here.

Component Sharded model Target Issues
Segment store, PostgreSQL LIST-partitioned parents, per-tenant child tables and DDL shared tables keyed by incarnation; registry and purge queue; writers pin the registry row #1544, #1546, #1549
Segment store, SQLite shared tables keyed by partition key (never sharded), no fencing same architecture on every dialect #1549
Vector store lifecycle, all backends per-backend ad-hoc catalogs, per-process locks; contract forbids multi-process management collection metadata arbitrated by a cross-process authority, so any process can create, open and delete; the contract states what each store supports #1524, #1525
Vector store create path and interface native container created inside create_collection (a dual write across registry and backend, safe only because names are content-addressed); per-collection VectorStoreCollectionConfig promising per-tenant schemas containers provisioned before serving; create_collection is a registry insert; collection shape configured per deployment #1572, #1573
Qdrant is_distributed shard key per logical collection; native name re-derived from sha256(config) on every open; in-process lock table payload partitioning; native name pinned at creation; lifecycle arbitrated by the cross-process authority; reclamation of dead partition keys #1564, #1565, #1563, #1562 (#1566 folded into #1564)
SQLite vector stores per-collection tables; stale handles reach the successor incarnation-scoped native resources (per-collection tables stay; the store's stated reason is that sqlite-vec's index may not support a partition key) #1536
Milvus per-record partition_key field filtered on queries (already value-based); catalog in a dummy-vector __registry collection catalog moved to the cross-process authority #1524, #1535
Session data manager read-then-write on session rows, safe only with one process per session database-arbitrated create and update #1543
Schema provisioning every process runs create_all at boot; no migration path for a layout change one provisioning owner; migrations carry layout changes such as the segment store's #1570
Above the stores process-local handle caches; no domain-error-to-status mapping evict and re-resolve on a stale handle; explicit API status #1571

Found along the way, not sharding-related, but blocking correctness on SQLite: #1568 (per-connection state registered per store on caller-supplied engines) and #1542 (engines built without WAL, busy timeout, or explicit write transactions). Contract issues that matter more once native collections are shared across tenants: #1534 (unbounded property-key cardinality, undeclared-property mandate) and #1535 (indexed_properties_schema is part of collection identity); #1573 narrows that to whether per-collection configuration belongs in the interface at all.

How to read the linked issues

Done when

  • The "at most one process" clause is gone from vector_store.py, replaced by a scope every store declares and honours.
  • No backend creates a per-tenant physical structure by default; tenant creation is a row write everywhere, and the per-tenant tier is an explicit opt-in.
  • Every lifecycle operation is safe from any process on the same backend: create races resolve by constraint, delete is idempotent and fences in-flight writers, purge is concurrent-safe and bounded.
  • A stale handle cannot reach a successor tenant on any backend (raising where the backend can pin a registry row, landing in a dead incarnation where it cannot), and every backend has reclamation for dead tenants.
  • Native containers are provisioned before serving and creating a logical collection is a single registry insert; collection shape is configured per deployment, one schema per container.
  • Existing deployments have a stated path for each layout change: a migration through Schema provisioning runs from every process's boot path: create_all races on cold boot and never evolves a table, Alembic runs destructive migrations on first use #1570's infrastructure, or an explicit recreate-schema statement of the kind the segment store overhaul makes for the pre-GA event backend.

Out of scope

Investigated and written by Claude (Claude Code), filed from the account of the user who commissioned the investigation.

Activity

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

Metadata

Metadata

Assignees

Labels

horizontal scalingWrong or unsafe when more than one server process serves the same backends (replicas or workers)keep-openPrevents the auto-close task from closing this issue.refactorCode refactoring that doesn't add new features.

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions