cobble_data_structure/
structured_remote_compaction_server.rs1use bytes::Bytes;
2use cobble::{Config, RemoteCompactionServer, Result};
3use std::net::TcpStream;
4use std::sync::Arc;
5
6use crate::structured_db::{
7 structured_merge_operator_resolver, structured_resolvable_operator_ids,
8};
9
10pub struct StructuredRemoteCompactionServer {
16 inner: Arc<RemoteCompactionServer>,
17}
18
19impl StructuredRemoteCompactionServer {
20 pub fn new(config: Config) -> Result<Self> {
21 let server = RemoteCompactionServer::new(config)?;
22 server.set_merge_operator_resolver(
23 structured_merge_operator_resolver(),
24 structured_resolvable_operator_ids(),
25 );
26 Ok(Self {
27 inner: Arc::new(server),
28 })
29 }
30
31 pub fn supported_merge_operator_ids(&self) -> Vec<String> {
32 self.inner.supported_merge_operator_ids()
33 }
34
35 pub fn register_schema_transform<F, T>(
37 &self,
38 transform_type: impl Into<String>,
39 factory: F,
40 ) -> Result<()>
41 where
42 F: Fn(&[u8]) -> Result<T> + Send + Sync + 'static,
43 T: Fn(Option<Bytes>) -> Result<Option<Bytes>> + Send + Sync + 'static,
44 {
45 self.inner
46 .register_schema_transform(transform_type, factory)
47 }
48
49 pub fn serve(&self, address: &str) -> Result<()> {
50 self.inner.serve(address)
51 }
52
53 pub fn handle_connection(&self, stream: TcpStream) -> Result<()> {
54 self.inner.handle_connection(stream)
55 }
56
57 pub fn inner(&self) -> &RemoteCompactionServer {
58 &self.inner
59 }
60
61 pub fn close(&self) {
62 self.inner.close()
63 }
64}
65
66#[cfg(test)]
67#[path = "../tests/unit/structured_remote_compaction_server.rs"]
68mod tests;