Skip to main content

valence_backend_redis/
fleet.rs

1//! Multi-node Redis fleet routing for bench campaigns.
2
3use std::sync::Arc;
4
5use valence_core::backend::DatabaseBackend;
6use valence_core::Result;
7
8use crate::backend::RedisBackend;
9use crate::config::{FleetRedisBackendBuilder, RedisConfig};
10
11/// Routes table operations across standalone Redis nodes by table hash.
12#[derive(Debug)]
13pub struct FleetRedisBackend {
14    backends: Vec<RedisBackend>,
15}
16
17impl FleetRedisBackend {
18    /// Start a builder for explicit fleet wiring.
19    pub fn builder() -> FleetRedisBackendBuilder {
20        FleetRedisBackendBuilder::new()
21    }
22
23    /// Connect using env defaults via builder (shorthand).
24    pub async fn from_env() -> Result<Self> {
25        Self::builder().from_env_defaults().build().await
26    }
27
28    /// Connect to every URL with shared key prefix.
29    pub async fn connect_with_urls(urls: Vec<String>, key_prefix: String) -> Result<Self> {
30        let mut backends = Vec::with_capacity(urls.len());
31        for url in urls {
32            backends.push(
33                RedisBackend::connect_with_config(RedisConfig::new(url, key_prefix.clone()))
34                    .await?,
35            );
36        }
37        Ok(Self { backends })
38    }
39
40    fn backend_for_table(&self, table: &str) -> &RedisBackend {
41        let idx = table_slot_index(table, self.backends.len());
42        &self.backends[idx]
43    }
44}
45
46fn table_slot_index(table: &str, n: usize) -> usize {
47    if n == 0 {
48        return 0;
49    }
50    table
51        .bytes()
52        .fold(0usize, |acc, b| acc.wrapping_add(usize::from(b)))
53        % n
54}
55
56#[async_trait::async_trait]
57impl DatabaseBackend for FleetRedisBackend {
58    fn engine_id(&self) -> &'static str {
59        crate::ENGINE_ID
60    }
61
62    fn capabilities(&self) -> valence_core::BackendCapabilities {
63        self.backends[0].capabilities()
64    }
65
66    async fn execute_compiled_query(
67        &self,
68        compiled: &valence_core::CompiledQuery,
69    ) -> Result<Vec<serde_json::Value>> {
70        self.backends[0].execute_compiled_query(compiled).await
71    }
72
73    async fn ensure_schemaless_table(&self, table: &str) -> Result<()> {
74        for backend in &self.backends {
75            backend.ensure_schemaless_table(table).await?;
76        }
77        Ok(())
78    }
79
80    async fn get_record(&self, table: &str, id: &str) -> Result<Option<serde_json::Value>> {
81        self.backend_for_table(table).get_record(table, id).await
82    }
83
84    async fn create_record(
85        &self,
86        table: &str,
87        content: serde_json::Value,
88    ) -> Result<serde_json::Value> {
89        self.backend_for_table(table)
90            .create_record(table, content)
91            .await
92    }
93
94    async fn update_record(
95        &self,
96        table: &str,
97        id: &str,
98        content: serde_json::Value,
99    ) -> Result<serde_json::Value> {
100        self.backend_for_table(table)
101            .update_record(table, id, content)
102            .await
103    }
104
105    async fn merge_record(
106        &self,
107        table: &str,
108        id: &str,
109        patch: serde_json::Value,
110    ) -> Result<serde_json::Value> {
111        self.backend_for_table(table)
112            .merge_record(table, id, patch)
113            .await
114    }
115
116    async fn upsert_record(
117        &self,
118        table: &str,
119        id: &str,
120        content: serde_json::Value,
121    ) -> Result<serde_json::Value> {
122        self.backend_for_table(table)
123            .upsert_record(table, id, content)
124            .await
125    }
126
127    async fn delete_record(&self, table: &str, id: &str) -> Result<()> {
128        self.backend_for_table(table).delete_record(table, id).await
129    }
130
131    async fn relate_edge(
132        &self,
133        from: &valence_core::RecordId,
134        edge_table: &str,
135        to: &valence_core::RecordId,
136    ) -> Result<()> {
137        self.backend_for_table(from.table())
138            .relate_edge(from, edge_table, to)
139            .await
140    }
141
142    async fn unrelate_edge(
143        &self,
144        from: &valence_core::RecordId,
145        edge_table: &str,
146        to: &valence_core::RecordId,
147    ) -> Result<()> {
148        self.backend_for_table(from.table())
149            .unrelate_edge(from, edge_table, to)
150            .await
151    }
152
153    async fn get_edge_targets(
154        &self,
155        from: &valence_core::RecordId,
156        edge_table: &str,
157    ) -> Result<Vec<valence_core::RecordId>> {
158        self.backend_for_table(from.table())
159            .get_edge_targets(from, edge_table)
160            .await
161    }
162
163    async fn define_unique_index(&self, table: &str, field: &str) -> Result<()> {
164        for backend in &self.backends {
165            backend.define_unique_index(table, field).await?;
166        }
167        Ok(())
168    }
169}
170
171/// Install a fleet backend as `Arc<dyn DatabaseBackend>`.
172pub async fn connect_fleet_arc(
173    builder: FleetRedisBackendBuilder,
174) -> Result<Arc<dyn DatabaseBackend>> {
175    Ok(Arc::new(builder.build().await?))
176}