// Corium control-plane protocol (see docs/design/protocol.md).
//
// Values never travel as protobuf messages: every `bytes` payload below is
// the Corium composite wire encoding produced by `corium-protocol::codec`.
syntax = "proto3";
package corium.v1;
// ---------------------------------------------------------------------------
// Shared messages
// ---------------------------------------------------------------------------
// Names a database view for a request: current, as-of, since, or history.
message DbViewSpec {
string db = 1;
oneof view {
uint64 as_of = 2;
uint64 since = 3;
bool history = 4;
}
}
message TransactRequest {
string db = 1;
uint32 protocol_version = 2;
// Composite-encoded EDN vector of transaction forms (map and list forms).
bytes tx_data = 3;
}
message TransactResponse {
// Basis of the database before this transaction.
uint64 basis_before = 1;
// Basis including this transaction (its `t`).
uint64 basis_t = 2;
int64 tx_instant = 3;
// Composite-encoded EDN map of tempid string -> allocated entity id (long).
bytes tempids = 4;
// Composite-encoded datom list for the transaction.
bytes tx_data = 5;
}
message SubscribeRequest {
string db = 1;
uint32 protocol_version = 2;
// The server backfills reports with `t > from_basis_t`, then streams live.
uint64 from_basis_t = 3;
}
// First item of every subscription: everything a peer needs to fold reports.
message Handshake {
uint64 basis_t = 1;
uint64 index_basis_t = 2;
// Composite-encoded schema attributes and ident registry.
bytes schema = 3;
// Interval between server heartbeats on this stream, in milliseconds.
// Peers treat silence for a few multiples of this as a dead transactor
// and fail over; 0 (older servers) disables the timeout.
uint64 heartbeat_interval_ms = 4;
}
message TxReport {
uint64 t = 1;
int64 tx_instant = 2;
// Composite-encoded datom list.
bytes datoms = 3;
}
message IndexBasis {
uint64 index_basis_t = 1;
}
message Heartbeat {
uint64 basis_t = 1;
}
// Tx-reports, index-basis announcements, and heartbeats multiplexed on the
// subscription stream.
message SubscribeItem {
oneof item {
Handshake handshake = 1;
TxReport report = 2;
IndexBasis index_basis = 3;
Heartbeat heartbeat = 4;
}
}
message SyncRequest {
string db = 1;
// Basis to wait for; 0 waits for the current basis.
uint64 t = 2;
}
message SyncResponse {
uint64 basis_t = 1;
}
message StatusRequest {
string db = 1;
}
message StatusResponse {
uint64 basis_t = 1;
uint64 index_basis_t = 2;
string lease_owner = 3;
uint64 lease_version = 4;
int64 lease_expires_unix_ms = 5;
uint64 datom_count = 6;
uint64 entity_count = 7;
uint64 attribute_count = 8;
uint64 transaction_count = 9;
uint64 transaction_failure_count = 10;
uint64 transaction_queue_depth = 11;
uint64 index_lag = 12;
uint64 indexing_runs = 13;
uint64 gc_runs = 14;
uint64 gc_swept_blobs = 15;
// Client endpoint advertised by the lease owner (HA discovery).
string lease_owner_endpoint = 16;
}
// ---------------------------------------------------------------------------
// Transactor service (peers -> transactor)
// ---------------------------------------------------------------------------
service Transactor {
rpc Transact(TransactRequest) returns (TransactResponse);
rpc Subscribe(SubscribeRequest) returns (stream SubscribeItem);
rpc Sync(SyncRequest) returns (SyncResponse);
rpc Status(StatusRequest) returns (StatusResponse);
}
// ---------------------------------------------------------------------------
// Catalog service (admin)
// ---------------------------------------------------------------------------
message CreateDatabaseRequest {
string db = 1;
// Composite-encoded EDN vector of Datomic-style attribute maps.
bytes schema = 2;
}
message CreateDatabaseResponse {
// False when the database already existed.
bool created = 1;
}
message DeleteDatabaseRequest {
string db = 1;
}
message DeleteDatabaseResponse {
bool deleted = 1;
}
message ForkDatabaseRequest {
// Source database.
string db = 1;
// Name for the new fork.
string target = 2;
// Transaction basis to fork at; 0 forks at the source's current basis.
uint64 as_of_t = 3;
}
message ForkDatabaseResponse {
// False when the target already existed (nothing was forked).
bool created = 1;
// The fork's basis: the source transaction it duplicates through.
uint64 basis_t = 2;
}
message ListDatabasesRequest {}
message ListDatabasesResponse {
repeated string dbs = 1;
}
message GcDeletedDatabasesRequest {
// Minimum age of unreachable blobs before deletion. Omitted uses the
// server default; an explicit 0 requests an immediate sweep.
optional uint64 retention_millis = 1;
}
message GcDeletedDatabasesResponse {
uint64 swept_blobs = 1;
}
message RequestIndexRequest {
string db = 1;
}
message RequestIndexResponse {
// Index basis after the publication; when the published indexes already
// covered every committed transaction, the current index basis.
uint64 index_basis_t = 1;
}
// Per-database runtime overrides of the index-publication pacing configured
// by the transactor's flags. Only set fields change; a request with no
// fields set reads the current policy.
message SetIndexPolicyRequest {
string db = 1;
optional uint64 interval_ms = 2;
optional uint32 backoff = 3;
optional uint64 tail_threshold = 4;
optional uint64 tail_deadline_ms = 5;
}
// The policy in effect after applying the request.
message SetIndexPolicyResponse {
uint64 interval_ms = 1;
uint32 backoff = 2;
uint64 tail_threshold = 3;
uint64 tail_deadline_ms = 4;
}
service Catalog {
rpc CreateDatabase(CreateDatabaseRequest) returns (CreateDatabaseResponse);
rpc DeleteDatabase(DeleteDatabaseRequest) returns (DeleteDatabaseResponse);
// Creates a new database duplicating an existing one at a transaction
// basis (point-in-time fork).
rpc ForkDatabase(ForkDatabaseRequest) returns (ForkDatabaseResponse);
rpc ListDatabases(ListDatabasesRequest) returns (ListDatabasesResponse);
rpc GcDeletedDatabases(GcDeletedDatabasesRequest) returns (GcDeletedDatabasesResponse);
// Publishes the database's covering indexes now, bypassing pacing.
rpc RequestIndex(RequestIndexRequest) returns (RequestIndexResponse);
// Adjusts (or reads) the database's index-publication pacing at runtime.
rpc SetIndexPolicy(SetIndexPolicyRequest) returns (SetIndexPolicyResponse);
}
// ---------------------------------------------------------------------------
// Peer server service (thin clients -> hosted peer)
// ---------------------------------------------------------------------------
message QueryRequest {
// Database views bound positionally to the query's `$`, `$2`, ... inputs.
repeated DbViewSpec dbs = 1;
// Composite-encoded EDN query form.
bytes query = 2;
// Composite-encoded EDN vector of the non-database arguments, in `:in` order.
bytes args = 3;
// Maximum datoms the execution may touch; 0 uses the server limit.
uint64 fuel = 4;
}
// Result shape, sent on the first chunk so clients can reassemble.
enum ResultShape {
RESULT_SHAPE_UNSPECIFIED = 0;
RESULT_SHAPE_RELATION = 1;
RESULT_SHAPE_COLLECTION = 2;
RESULT_SHAPE_TUPLE = 3;
RESULT_SHAPE_SCALAR = 4;
}
message QueryResultChunk {
ResultShape shape = 1;
// Composite-encoded EDN. Relations/collections stream vectors of rows to
// concatenate; tuple/scalar results arrive in a single chunk.
bytes rows = 2;
bool last = 3;
}
message PullRequest {
DbViewSpec db = 1;
// Composite-encoded EDN pull pattern.
bytes pattern = 2;
// Composite-encoded EDN entity position: long, `#eid`, ident, or lookup ref.
bytes eid = 3;
}
message PullResponse {
bytes result = 1;
}
message DatomsRequest {
DbViewSpec db = 1;
// One of "eavt", "aevt", "avet", "vaet".
string index = 2;
// Composite-encoded EDN vector of leading components in index order.
bytes components = 3;
// Maximum datoms to return; 0 uses the server limit.
uint64 limit = 4;
}
message DatomChunk {
// Composite-encoded datom list.
bytes datoms = 1;
bool last = 2;
}
message TxRangeRequest {
string db = 1;
uint64 start = 2;
// Exclusive end; 0 means open-ended.
uint64 end = 3;
}
message TxChunk {
repeated TxReport txes = 1;
bool last = 2;
}
message DbStatsRequest {
DbViewSpec db = 1;
}
message DbStatsResponse {
uint64 basis_t = 1;
uint64 datom_count = 2;
uint64 entity_count = 3;
uint64 attribute_count = 4;
}
service PeerServer {
rpc Query(QueryRequest) returns (stream QueryResultChunk);
rpc Pull(PullRequest) returns (PullResponse);
rpc Transact(TransactRequest) returns (TransactResponse);
rpc Datoms(DatomsRequest) returns (stream DatomChunk);
rpc TxRange(TxRangeRequest) returns (stream TxChunk);
rpc DbStats(DbStatsRequest) returns (DbStatsResponse);
rpc Subscribe(SubscribeRequest) returns (stream SubscribeItem);
}