Skip to main content

cobble_data_structure/
structured_remote_compaction_server.rs

1use 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
10/// A remote compaction server pre-configured to resolve structured data type
11/// merge operators (e.g. list) from request metadata.
12///
13/// This wraps [`RemoteCompactionServer`] and automatically registers the
14/// structured merge operator resolver on construction.
15pub 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    /// Register a factory for raw single-column transform specifications.
36    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;