corium-protocol 0.1.21

Corium gRPC network protocol
Documentation
// 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);
}