Dataset backend adapters
Problem
Section titled “Problem”StorageSession used to branch on if backend == …. Adding a sink meant editing session code.
Solution (P1)
Section titled “Solution (P1)”Protocols + registry:
class DatasetWriteBackend(Protocol): name: str def write(self, session, dataset_key, rows, *, shard_id="0"): ... def flush_buffers(self, session): ... def drain_all_buffers(self, session): ...
class DatasetReadBackend(Protocol): name: str def read_manifest_total_rows(...): ... def read_run_rows(...): ... def iter_run_rows(...): ...Builtins (backend_builtins.py): parquet, postgres, qdrant, neo4j stub.
session_write.write calls write_backend_for(spec.backend); session_read.read_run_rows and row iteration call read_backend_for(spec.backend).
- One dataset — one backend (dual-write forbidden).
- Steps do not branch on backend.
- New sink: adapter +
register_*_backend(+db/<engine>/transport) — no if in session. - DRTML whitelist:
drtml/storage/backend.py. - Multi-worker remote Ray: parquet stores with
backend: localare rejected whenray_addressis set andworkers > 1(orDRTOLLER_REQUIRE_OBJECT_STORE=1). Usebackend: s3(MinIO / AWS / GCS interop) so every actor sees the same keys — Ray object-store guard.
Transport vs adapter
Section titled “Transport vs adapter”| Backend | Adapter | Transport |
|---|---|---|
| parquet | builtin → session_write.write_parquet |
local/S3 backends/ |
| postgres | builtin → postgres_write |
db/postgres/ |
| qdrant | builtin → qdrant_write |
db/qdrant/ |
| neo4j | stub NotImplementedError |
future db/neo4j/ |
dev/tests/storage/test_dataset_backend_adapter.py — fake memory_test backend via the registry.
Cookbook
Section titled “Cookbook”- New dataset backend — full vertical recipe
- New database / engine — same path, DB-oriented entry