exoware-server
Serve the Exoware API.
Status
exoware-server is ALPHA software and is not yet recommended for production use. Developers should expect breaking changes and occasional instability.
Overview
exoware-server provides a backend-less ConnectRPC server for the Exoware API.
Implement the storage capability traits for your backend, wrap them in AppState,
and call connect_stack to get a ready-to-serve router with ingest, query,
prune, retention, and stream services. Backends that implement every capability
automatically implement the StoreEngine compatibility facade.
Split deployments can instead mount ingest_service, query_stack,
prune_service, retention_service, or stream_service with the narrower
component state. The stream service accepts an in-process StreamNotifier;
StreamHub is the local default.
Reduce uses native DataFusion aggregation and streams completed groups in bounded
frames. Groups remain in memory by default. QueryState::with_runtime accepts a
DataFusion RuntimeEnv to configure its native memory pool and spill storage.
The pool accounts for worker aggregate state; backend buffers, transient input
batches, transport frames, and client results have separate memory ownership.
Unordered aggregation consumes its input before producing results.
use Bytes;
use PrunePolicyDocument;
use ;
use Future;
// Implement the capabilities your component serves:
// Sequence:
// fn current_sequence(&self) -> u64;
//
// Ingest:
// fn put_batch(&self, kvs: Vec<(Bytes, Bytes)>) -> impl Future<Output = Result<u64, String>> + Send + '_;
//
// Query:
// type RangeScan: RangeScan;
// fn get(&self, key: Bytes) -> impl Future<Output = Result<(Option<Vec<u8>>, QueryExtra), String>> + Send + '_;
// fn range_scan(&self, start: Bytes, end: Bytes, limit: usize, forward: bool) -> impl Future<Output = Result<Self::RangeScan, String>> + Send + '_;
// fn get_many(&self, keys: Vec<Bytes>) -> impl Future<Output = Result<(Vec<(Vec<u8>, Option<Vec<u8>>)>, QueryExtra), String>> + Send + '_;
//
// Prune:
// fn apply_prune_policies(&self, document: PrunePolicyDocument) -> impl Future<Output = Result<(), String>> + Send + '_;
//
// Log:
// fn get_batch(&self, sequence_number: u64) -> impl Future<Output = Result<Option<Vec<(Bytes, Bytes)>>, String>> + Send + '_;
// fn oldest_retained_batch(&self) -> impl Future<Output = Result<Option<u64>, String>> + Send + '_;
//
// Retention:
// fn set_retention(&self, policy: Option<RetentionPolicy>) -> impl Future<Output = Result<Option<u64>, String>> + Send + '_;