hakocluster/lib.rs
1//! hakocluster: dispatcher over N hakodb instances (Fase 2).
2//!
3//! Full design: hakodb/hakocluster#1.
4//! - `Cluster::open[_with_config]` — one `Hako` per data dir + full-mesh
5//! `socket_sync` peering (one connection per pair, `i` dials `j > i`).
6//! - Reads (`get`/`query`) fan out round-robin across healthy instances.
7//! - Writes (`put`/`put_owned`/`delete`) route to the designated writer,
8//! `instances[0]` by convention.
9//! - Fase 2 adds: stagger policy (same interval + spaced opens, or
10//! per-instance intervals) and a lag guard (eject replicas trailing the
11//! writer by more than `max_replica_lag_versions`, re-admit at half).
12//! - Replicated writes join each instance's normal write/flush queue (that
13//! is `SocketSync`'s own behavior); the cluster adds no flush paths.
14//!
15//! Non-unix: `socket_sync` compiles out (same rule as hakodb), so peering
16//! is unavailable — `open` with N > 1 fails closed; N = 1 works as a
17//! degenerate single-node cluster.
18
19pub mod ffi;
20
21use std::collections::{HashMap, HashSet};
22use std::path::{Path, PathBuf};
23use std::sync::{
24 Arc,
25 atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering},
26};
27
28use hakodb::config::HakoConfig;
29use hakodb::document::hako_doc::HakoDoc;
30use hakodb::engine::Hako;
31use hakodb::query::query::Query;
32
33#[cfg(unix)]
34use hakodb::socket_sync::SocketSync;
35
36/// Flush-cadence stagger policy (issue #1: flush cadence is the ONLY
37/// knob — never skip the queue, no direct-flush paths).
38#[derive(Debug, Clone)]
39pub enum StaggerPolicy {
40 /// Option A (default): same interval everywhere; instance opens spaced
41 /// by `offset_ms` so group-commit phases don't coincide. The engine's
42 /// `last_sync` starts at WAL open, so spaced opens = phased flushes.
43 /// Zero coordination protocol.
44 StaggeredStart { offset_ms: u64 },
45 /// Option B: per-instance intervals (e.g. co-prime-ish 5/7/11). Spreads
46 /// load without start-order dependence; durability lag per instance is
47 /// slightly uneven (bounded by the engine's 30s clamp). Length must
48 /// equal the instance count.
49 PerInstance(Vec<u64>),
50 /// Option C: Manual durability everywhere; the deployer calls
51 /// `tick_flush()` at its own cadence and the cluster flushes one
52 /// instance per call in rotation, so fsync storms never coincide.
53 /// Requires `durability_mode = Manual` (refused otherwise — an
54 /// Interval engine would flush behind the rotation's back). Maximum
55 /// control, caller-driven: no background thread, works everywhere.
56 ManualRotation,
57}
58
59impl Default for StaggerPolicy {
60 fn default() -> Self {
61 Self::StaggeredStart { offset_ms: 1 }
62 }
63}
64
65/// Fase 2 cluster configuration.
66#[derive(Debug, Clone)]
67pub struct ClusterConfig {
68 /// Durability for every instance (tests use Interval: the socket tailer
69 /// reads the WAL file, so only flushed bytes replicate live).
70 pub durability_mode: hakodb::config::DurabilityMode,
71 /// Group-commit window used when `stagger` doesn't override it
72 /// (StaggeredStart and the default).
73 pub group_commit_interval_ms: u64,
74 /// Flush-cadence stagger across instances (default: 1ms-spaced starts).
75 pub stagger: StaggerPolicy,
76 /// Lag guard: eject a replica from fan-out when it trails the writer
77 /// by more than this many versions (version-map delta). Versions are
78 /// write-micros, so the delta doubles as staleness: a healthy replica
79 /// trails by ~ the socket tail cadence (500ms = 500_000). Default
80 /// 5_000_000 (~10x the tail) tolerates jitter without flapping.
81 /// Re-admit at half the threshold (hysteresis). `None` disables.
82 /// The writer (index 0) never ejects.
83 pub max_replica_lag_versions: Option<u64>,
84 /// Directory holding one `instance-{i}.sock` per member.
85 pub sock_dir: PathBuf,
86}
87
88impl Default for ClusterConfig {
89 fn default() -> Self {
90 Self {
91 durability_mode: hakodb::config::DurabilityMode::Interval,
92 group_commit_interval_ms: 5,
93 stagger: StaggerPolicy::StaggeredStart { offset_ms: 1 },
94 max_replica_lag_versions: Some(5_000_000),
95 sock_dir: PathBuf::from("socks"),
96 }
97 }
98}
99
100impl std::fmt::Debug for Cluster {
101 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
102 f.debug_struct("Cluster")
103 .field("instances", &self.instances.len())
104 .field("reads", &self.read_counts())
105 .finish()
106 }
107}
108
109/// One clustered engine: N instances, one logical dataset.
110pub struct Cluster {
111 instances: Vec<Instance>,
112 /// Round-robin cursor for read fan-out.
113 rr: AtomicUsize,
114 /// Designated writer index (manual failover moves it).
115 writer_index: AtomicUsize,
116 /// Promotion epoch (1-based; 0 = initial writer, never promoted).
117 epoch: AtomicU64,
118 /// Append-only promotion audit (promotions are rare; Vec is fine).
119 promotions: std::sync::Mutex<Vec<Promotion>>,
120 /// Option C rotation cursor for [`Cluster::tick_flush`].
121 flush_rr: AtomicUsize,
122 /// Lag-guard threshold (`None` = disabled). Stored so
123 /// `refresh_health` stays a pure read of instance state.
124 max_lag: Option<u64>,
125 /// Runtime hosting the socket tasks (unix only): owned when `open`
126 /// runs outside any runtime (plain sync callers, tests), shared when
127 /// called inside one (servers like hakobackend run `#[tokio::main]` —
128 /// nesting runtimes panics, and blocking is illegal there, so dials
129 /// are spawned instead of awaited; peering converges asynchronously).
130 #[cfg(unix)]
131 _rt: Rt,
132}
133
134#[cfg(unix)]
135enum Rt {
136 Owned(tokio::runtime::Runtime),
137 Shared(tokio::runtime::Handle),
138}
139
140struct Instance {
141 db: Arc<Hako>,
142 /// Reads served (fan-out accounting; phase 2 least-busy input).
143 reads: AtomicU64,
144 /// False = ejected by the lag guard, skipped by fan-out.
145 healthy: AtomicBool,
146 #[cfg(unix)]
147 _sync: std::sync::Arc<SocketSync>,
148}
149
150/// One promotion record: which epoch moved the writer where, and when.
151/// Operators use epochs to order promotions (a higher epoch supersedes);
152/// the actual fence is the read-only flags [`Cluster::promote`] sets.
153#[derive(Debug, Clone, Copy, PartialEq, Eq)]
154pub struct Promotion {
155 /// 1-based promotion counter (0 = no promotion yet, initial writer).
156 pub epoch: u64,
157 /// New designated writer index.
158 pub writer: usize,
159 /// Wall millis at promotion (operator audit only, never a lease).
160 pub at_ms: u64,
161}
162
163/// Per-instance health from [`Cluster::refresh_health`].
164#[derive(Debug, Clone, PartialEq, Eq)]
165pub struct ReplicaHealth {
166 /// Instance index (0 = designated writer, always healthy).
167 pub index: usize,
168 /// Max version-map delta vs the writer (0 = converged). Versions are
169 /// write-micros, so this doubles as staleness in micros.
170 pub lag_versions: u64,
171 /// In fan-out rotation or ejected.
172 pub healthy: bool,
173}
174
175impl Cluster {
176 /// Open with default config; socket dir is `socks/` next to `paths[0]`.
177 pub fn open(paths: &[&str]) -> Result<Self, String> {
178 let sock_dir = Path::new(paths.first().ok_or("need ≥1 path")?)
179 .parent()
180 .map(Path::to_path_buf)
181 .unwrap_or_else(|| PathBuf::from("."))
182 .join("socks");
183 Self::open_with_config(
184 paths,
185 ClusterConfig {
186 sock_dir,
187 ..ClusterConfig::default()
188 },
189 )
190 }
191
192 /// Open one `Hako` per data dir and mesh the instances with
193 /// `socket_sync` (unix only — see module docs).
194 pub fn open_with_config(paths: &[&str], cfg: ClusterConfig) -> Result<Self, String> {
195 if paths.is_empty() {
196 return Err("need ≥1 path".into());
197 }
198 #[cfg(not(unix))]
199 if paths.len() > 1 {
200 return Err(
201 "socket peering is unix-only: N>1 clusters need a unix host".into(),
202 );
203 }
204
205 // Resolve the per-instance group-commit window. StaggeredStart
206 // shares one value; PerInstance overrides per index (length
207 // validated up front — fail before opening anything).
208 // ManualRotation forces Manual durability (validated, not
209 // silently overridden — an Interval engine would flush behind
210 // the rotation's back, defeating the point).
211 let intervals: Vec<u64> = match &cfg.stagger {
212 StaggerPolicy::StaggeredStart { .. } => {
213 vec![cfg.group_commit_interval_ms; paths.len()]
214 }
215 StaggerPolicy::PerInstance(v) => {
216 if v.len() != paths.len() {
217 return Err(format!(
218 "PerInstance needs one interval per path (got {} for {})",
219 v.len(),
220 paths.len()
221 ));
222 }
223 v.clone()
224 }
225 StaggerPolicy::ManualRotation => {
226 if cfg.durability_mode != hakodb::config::DurabilityMode::Manual {
227 return Err(
228 "ManualRotation needs durability_mode = Manual".into(),
229 );
230 }
231 vec![cfg.group_commit_interval_ms; paths.len()]
232 }
233 };
234 let stagger_gap =
235 match cfg.stagger {
236 StaggerPolicy::StaggeredStart { offset_ms } => offset_ms,
237 StaggerPolicy::PerInstance(_) | StaggerPolicy::ManualRotation => 0,
238 };
239
240 let mut dbs = Vec::with_capacity(paths.len());
241 for (n, p) in paths.iter().enumerate() {
242 if n > 0 && stagger_gap > 0 {
243 // Phase the group-commit windows: each WAL's last_sync
244 // starts at its own open, so spaced opens = phased flushes.
245 std::thread::sleep(std::time::Duration::from_millis(stagger_gap));
246 }
247 let mut hc = HakoConfig::default();
248 hc.durability_mode = cfg.durability_mode;
249 hc.group_commit_interval_ms = intervals[n];
250 dbs.push(Arc::new(
251 Hako::open(p, hc).map_err(|e| format!("open {p}: {e}"))?,
252 ));
253 }
254
255 #[cfg(unix)]
256 {
257 let rt = match tokio::runtime::Handle::try_current() {
258 Ok(h) => Rt::Shared(h),
259 Err(_) => Rt::Owned(
260 tokio::runtime::Runtime::new()
261 .map_err(|e| format!("tokio runtime: {e}"))?,
262 ),
263 };
264 std::fs::create_dir_all(&cfg.sock_dir)
265 .map_err(|e| format!("sock_dir: {e}"))?;
266 let socks: Vec<PathBuf> = (0..dbs.len())
267 .map(|i| cfg.sock_dir.join(format!("instance-{i}.sock")))
268 .collect();
269 // Serve all, then dial the mesh (one connection per pair:
270 // i dials j > i; traffic is bidirectional per connection).
271 // ponytail: serve() needs a reactor context (UnixListener::
272 // from_std) — see the mesh block below for how each runtime
273 // shape provides it.
274 let mut instances = Vec::with_capacity(dbs.len());
275 for db in &dbs {
276 instances.push(Instance {
277 db: db.clone(),
278 reads: AtomicU64::new(0),
279 healthy: AtomicBool::new(true),
280 _sync: std::sync::Arc::new(SocketSync::new(db.clone(), vec![])),
281 });
282 }
283 // Fail-closed both sides (synchronous, deterministic from
284 // boot): only the designated writer accepts local writes.
285 // Replicated ingest bypasses the flag by design, so replicas
286 // keep converging while read-only.
287 for (n, inst) in instances.iter().enumerate() {
288 inst.db.set_read_only(n != 0);
289 }
290 // Mesh setup: serve needs a reactor context
291 // (UnixListener::from_std) in both cases.
292 let syncs: Vec<(std::sync::Arc<SocketSync>, PathBuf)> = instances
293 .iter()
294 .zip(socks.iter())
295 .map(|(inst, sock)| (inst._sync.clone(), sock.clone()))
296 .collect();
297 let mesh = async move {
298 for (sync, sock) in &syncs {
299 sync.serve(sock.to_string_lossy().as_ref()).map_err(|e| {
300 format!("serve {}: {e}", sock.display())
301 })?;
302 }
303 for i in 0..syncs.len() {
304 for j in (i + 1)..syncs.len() {
305 let path = syncs[j].1.to_string_lossy().into_owned();
306 // ponytail: short retry for boot-order races (the
307 // other side serves a moment later). Peer RESTART
308 // healing is out of phase-1 scope (no re-meshing);
309 // the backend driver layer retries for that.
310 let mut last = String::new();
311 let mut ok = false;
312 for _ in 0..5 {
313 match syncs[i].0.dial(&path).await {
314 Ok(_) => {
315 ok = true;
316 break;
317 }
318 Err(e) => {
319 last = e.to_string();
320 tokio::time::sleep(
321 std::time::Duration::from_secs(1),
322 )
323 .await;
324 }
325 }
326 }
327 if !ok {
328 return Err(format!(
329 "dial {}: {last} (mesh incomplete)",
330 syncs[j].1.display()
331 ));
332 }
333 }
334 }
335 Ok::<(), String>(())
336 };
337 match &rt {
338 Rt::Owned(r) => r.block_on(mesh)?,
339 Rt::Shared(h) => {
340 h.spawn(async move {
341 if let Err(e) = mesh.await {
342 eprintln!("[hakocluster] mesh failed: {e}");
343 }
344 });
345 }
346 }
347 return Ok(Self {
348 instances,
349 rr: AtomicUsize::new(0),
350 writer_index: AtomicUsize::new(0),
351 epoch: AtomicU64::new(0),
352 promotions: std::sync::Mutex::new(Vec::new()),
353 flush_rr: AtomicUsize::new(0),
354 max_lag: cfg.max_replica_lag_versions,
355 _rt: rt,
356 });
357 }
358
359 #[cfg(not(unix))]
360 Ok(Self {
361 instances: dbs
362 .into_iter()
363 .enumerate()
364 .map(|(n, db)| {
365 db.set_read_only(n != 0);
366 Instance {
367 db,
368 reads: AtomicU64::new(0),
369 healthy: AtomicBool::new(true),
370 }
371 })
372 .collect(),
373 rr: AtomicUsize::new(0),
374 writer_index: AtomicUsize::new(0),
375 epoch: AtomicU64::new(0),
376 promotions: std::sync::Mutex::new(Vec::new()),
377 flush_rr: AtomicUsize::new(0),
378 max_lag: cfg.max_replica_lag_versions,
379 })
380 }
381
382 /// All instances, index 0 conventionally the designated writer.
383 pub fn instances(&self) -> Vec<Arc<Hako>> {
384 self.instances.iter().map(|i| i.db.clone()).collect()
385 }
386
387 /// The designated writer. All cluster writes route here.
388 pub fn writer(&self) -> &Arc<Hako> {
389 &self.instances[self.writer_index.load(Ordering::Relaxed)].db
390 }
391
392 /// Designated writer index (0 by convention; [`Self::promote`] moves it).
393 pub fn writer_index(&self) -> usize {
394 self.writer_index.load(Ordering::Relaxed)
395 }
396
397 /// Manual failover: move the designated writer to `index`. Returns
398 /// the promotion epoch (1-based; re-promoting the current writer is
399 /// a no-op returning the current epoch without logging). Every other
400 /// instance is set read-only (in-process, so always reachable — no
401 /// partial-failure story here). NO fencing and NO auto-detect: the
402 /// operator must fence the old writer first; two live writers diverge
403 /// under LWW. Lease granting (auto-failover) needs balancer HA first
404 /// and lives there, not here (see hakobalancer).
405 pub fn promote(&self, index: usize) -> Result<u64, String> {
406 if index >= self.instances.len() {
407 return Err(format!(
408 "promote: index {index} out of range (n={})",
409 self.instances.len()
410 ));
411 }
412 if index == self.writer_index.load(Ordering::Relaxed) {
413 return Ok(self.epoch.load(Ordering::Relaxed));
414 }
415 for (n, inst) in self.instances.iter().enumerate() {
416 inst.db.set_read_only(n != index);
417 }
418 self.writer_index.store(index, Ordering::Relaxed);
419 let epoch = self.epoch.fetch_add(1, Ordering::Relaxed) + 1;
420 let at_ms = std::time::SystemTime::now()
421 .duration_since(std::time::UNIX_EPOCH)
422 .map(|d| d.as_millis() as u64)
423 .unwrap_or(0);
424 self.promotions
425 .lock()
426 .unwrap()
427 .push(Promotion { epoch, writer: index, at_ms });
428 Ok(epoch)
429 }
430
431 /// Current promotion epoch (0 = initial writer, never promoted).
432 pub fn epoch(&self) -> u64 {
433 self.epoch.load(Ordering::Relaxed)
434 }
435
436 /// Promotion audit log, oldest first.
437 pub fn promotion_log(&self) -> Vec<Promotion> {
438 self.promotions.lock().unwrap().clone()
439 }
440
441 /// Option C rotation tick: flush ONE instance (round-robin) so its
442 /// buffered Manual writes reach the WAL file (and the socket tailer).
443 /// The deployer calls this at its own flush cadence; fsync storms
444 /// never coincide because only one instance flushes per tick.
445 pub fn tick_flush(&self) -> Result<(), String> {
446 let n = self.instances.len();
447 let i = self.flush_rr.fetch_add(1, Ordering::Relaxed) % n;
448 self.instances[i]
449 .db
450 .flush()
451 .map_err(|e| format!("tick_flush instance {i}: {e}"))
452 }
453
454 /// Member count.
455 pub fn instance_count(&self) -> usize {
456 self.instances.len()
457 }
458
459 /// Next read replica index, round-robin over healthy instances
460 /// (ejected ones skipped; the writer always qualifies). Driver
461 /// entry-point for fan-out without reimplementing picking.
462 pub fn read_index(&self) -> usize {
463 let n = self.instances.len();
464 // ponytail: wrapping_add, not checked math — a counter that runs
465 // for centuries is the only overflow story, and modulo is safe.
466 let start = self.rr.fetch_add(1, Ordering::Relaxed);
467 for k in 0..n {
468 let i = (start + k) % n;
469 if self.instances[i].healthy.load(Ordering::Relaxed) {
470 return i;
471 }
472 }
473 // Unreachable (writer never ejects) — fail closed to the writer.
474 self.writer_index.load(Ordering::Relaxed)
475 }
476
477 /// Record a served read for fan-out accounting (drivers that pick
478 /// via read_index() and serve through their own handles call this;
479 /// Cluster::get/query do it internally). Out-of-range is a no-op.
480 pub fn note_read(&self, index: usize) {
481 if let Some(inst) = self.instances.get(index) {
482 inst.reads.fetch_add(1, Ordering::Relaxed);
483 }
484 }
485
486 /// Reads served per instance (fan-out accounting).
487 pub fn read_counts(&self) -> Vec<u64> {
488 self.instances
489 .iter()
490 .map(|i| i.reads.load(Ordering::Relaxed))
491 .collect()
492 }
493
494 /// Total live socket peerings (0 off-unix).
495 pub fn peer_count(&self) -> usize {
496 #[cfg(unix)]
497 return self.instances.iter().map(|i| i._sync.peer_count()).sum();
498 #[cfg(not(unix))]
499 return 0;
500 }
501
502 /// Recompute replica lag vs the writer and apply the guard: eject past
503 /// `max_replica_lag_versions`, re-admit at half (hysteresis). The
504 /// designated writer never ejects. `None` disables ejection (lag still
505 /// reported). Explicit call — no background thread in phase 2; drive
506 /// it from the deployer's own tick.
507 pub fn refresh_health(&self) -> Vec<ReplicaHealth> {
508 let w = self.writer_index.load(Ordering::Relaxed);
509 let base = self.instances[w].db.get_version_map();
510 let mut out = Vec::with_capacity(self.instances.len());
511 for (n, inst) in self.instances.iter().enumerate() {
512 let lag = if n == w {
513 0
514 } else {
515 let m = inst.db.get_version_map();
516 base.iter()
517 .map(|(col, wv)| wv.saturating_sub(*m.get(col).unwrap_or(&0)) as u64)
518 .max()
519 .unwrap_or(0)
520 };
521 let healthy = match (n == w, self.max_lag) {
522 (true, _) => true,
523 (false, None) => true,
524 (false, Some(max)) => {
525 let cur = inst.healthy.load(Ordering::Relaxed);
526 // ponytail: hysteresis in one expression — eject past
527 // max, re-admit at/below half, otherwise hold state.
528 if lag > max {
529 false
530 } else if lag <= max / 2 {
531 true
532 } else {
533 cur
534 }
535 }
536 };
537 inst.healthy.store(healthy, Ordering::Relaxed);
538 out.push(ReplicaHealth {
539 index: n,
540 lag_versions: lag,
541 healthy,
542 });
543 }
544 out
545 }
546
547 /// Instances currently in fan-out rotation.
548 pub fn healthy_count(&self) -> usize {
549 self.instances
550 .iter()
551 .filter(|i| i.healthy.load(Ordering::Relaxed))
552 .count()
553 }
554
555 /// Next read replica, round-robin over healthy instances (ejected ones
556 /// are skipped; the writer always qualifies).
557 fn pick(&self) -> &Instance {
558 &self.instances[self.read_index()]
559 }
560
561 // --- Write path: designated writer only (issue #1, phase 1). ---
562
563 /// Route a write to the designated writer.
564 pub fn put(&self, col: &str, id: &str, doc: &HakoDoc) -> Result<String, String> {
565 self.writer()
566 .put(col, id, doc)
567 .map_err(|e| e.to_string())
568 }
569
570 /// Owned-doc variant (skips the deep clone when the caller owns the doc).
571 pub fn put_owned(&self, col: &str, id: &str, doc: HakoDoc) -> Result<String, String> {
572 self.writer()
573 .put_owned(col, id, doc)
574 .map_err(|e| e.to_string())
575 }
576
577 /// Route a delete to the designated writer.
578 pub fn delete(&self, col: &str, id: &str) -> Result<String, String> {
579 self.writer().delete(col, id).map_err(|e| e.to_string())
580 }
581
582 // --- Read path: round-robin fan-out (issue #1, phase 1). ---
583
584 /// Fan-out point read. No read-your-write across instances in phase 1:
585 /// replicas trail the writer by <= the socket tail interval.
586 pub fn get(
587 &self,
588 collection: &str,
589 doc_id: &str,
590 ) -> Result<Option<HakoDoc>, String> {
591 let inst = self.pick();
592 inst.reads.fetch_add(1, Ordering::Relaxed);
593 inst.db.get(collection, doc_id).map_err(|e| e.to_string())
594 }
595
596 /// Fan-out query (same lag contract as [`Self::get`]).
597 pub fn query(&self, query: Query) -> Result<Vec<(String, HakoDoc)>, String> {
598 let inst = self.pick();
599 inst.reads.fetch_add(1, Ordering::Relaxed);
600 inst.db.query(query).map_err(|e| e.to_string())
601 }
602}
603
604#[cfg(unix)]
605impl Drop for Cluster {
606 fn drop(&mut self) {
607 for i in &self.instances {
608 i._sync.stop();
609 }
610 // Runtime drop aborts the accept/tail tasks.
611 }
612}
613
614// --- Multidatabase registry (hakocluster#6): one process, N named
615// databases. Fixed at open (reload re-opens; no runtime membership
616// mutation). Each database is a full `Cluster` with its own mesh —
617// NEVER meshed across databases (converging divergent data is
618// corruption shaped as operation, see #2). Unknown names are ALWAYS
619// 404/None, even opt-in auto-create does not exist here: storage may
620// be fresh-empty on open (as today), but the NAME must be declared.
621
622/// Database name gate: same discipline as collection segments
623/// (`[A-Za-z0-9_-]`, 1–128, no `__` prefix). The name becomes a sock
624/// subdir, so traversal shapes are refused here, not downstream.
625pub fn valid_db_name(s: &str) -> bool {
626 !s.is_empty()
627 && s.len() <= 128
628 && s.bytes().all(|c| c.is_ascii_alphanumeric() || c == b'_' || c == b'-')
629 && !s.starts_with("__")
630}
631
632/// One named database declaration for [`Databases::open`].
633#[derive(Debug, Clone)]
634pub struct DbSpec {
635 pub name: String,
636 pub paths: Vec<String>,
637 pub config: ClusterConfig,
638}
639
640/// N named databases in one process. Selection is by exact name;
641/// there is no default database (a silent default routes writes
642/// somewhere the operator did not choose).
643pub struct Databases {
644 dbs: HashMap<String, Arc<Cluster>>,
645 /// Declaration order (deterministic listing).
646 order: Vec<String>,
647 /// The single shared runtime, present only when this registry
648 /// minted it (FFI sync callers). Server-embedded use shares the
649 /// server runtime instead — N databases must never mean N runtimes.
650 #[cfg(unix)]
651 _rt: Option<tokio::runtime::Runtime>,
652}
653
654impl Databases {
655 /// Open every declared database. Each gets `sock_root/{name}` as
656 /// its mesh dir (doubles as the mesh boundary: no cross-database
657 /// peering is representable). Fail-closed: bad/duplicate/empty
658 /// declarations and any per-db open failure refuse the whole
659 /// registry (already-opened siblings drop cleanly via `Drop`).
660 pub fn open(specs: Vec<DbSpec>, sock_root: PathBuf) -> Result<Self, String> {
661 if specs.is_empty() {
662 return Err("need ≥1 database".into());
663 }
664 let mut seen = HashSet::new();
665 for s in &specs {
666 if !valid_db_name(&s.name) {
667 return Err(format!("bad database name `{}`", s.name));
668 }
669 if !seen.insert(s.name.clone()) {
670 return Err(format!("duplicate database `{}`", s.name));
671 }
672 if s.paths.is_empty() {
673 return Err(format!("database `{}` needs ≥1 path", s.name));
674 }
675 }
676 // ponytail: open inside the shared context when we mint the
677 // runtime, so every Cluster takes Shared (one runtime total).
678 // No Cluster API change: try_current succeeds inside block_on.
679 #[cfg(unix)]
680 {
681 if tokio::runtime::Handle::try_current().is_ok() {
682 let (dbs, order) = Self::open_all(specs, &sock_root)?;
683 Ok(Self { dbs, order, _rt: None })
684 } else {
685 let rt = tokio::runtime::Runtime::new()
686 .map_err(|e| format!("tokio runtime: {e}"))?;
687 let out = rt.block_on(async { Self::open_all(specs, &sock_root) });
688 let (dbs, order) = out?;
689 Ok(Self { dbs, order, _rt: Some(rt) })
690 }
691 }
692 #[cfg(not(unix))]
693 {
694 let (dbs, order) = Self::open_all(specs, &sock_root)?;
695 Ok(Self { dbs, order })
696 }
697 }
698
699 fn open_all(
700 specs: Vec<DbSpec>,
701 sock_root: &Path,
702 ) -> Result<(HashMap<String, Arc<Cluster>>, Vec<String>), String> {
703 let mut dbs = HashMap::with_capacity(specs.len());
704 let mut order = Vec::with_capacity(specs.len());
705 for s in specs {
706 // Mesh boundary, materialized on all platforms (the mesh
707 // itself is unix-only, but the per-db dir existing
708 // everywhere keeps the invariant observable + fails fast
709 // on an unwritable root).
710 let sub = sock_root.join(&s.name);
711 std::fs::create_dir_all(&sub)
712 .map_err(|e| format!("database `{}` sock dir: {e}", s.name))?;
713 let mut cfg = s.config;
714 cfg.sock_dir = sub;
715 let refs: Vec<&str> = s.paths.iter().map(|p| p.as_str()).collect();
716 let c = Cluster::open_with_config(&refs, cfg)
717 .map_err(|e| format!("database `{}`: {e}", s.name))?;
718 order.push(s.name.clone());
719 dbs.insert(s.name, Arc::new(c));
720 }
721 Ok((dbs, order))
722 }
723
724 /// Exact-name lookup. `None` = unknown (caller 404s). No default.
725 pub fn get(&self, name: &str) -> Option<Arc<Cluster>> {
726 self.dbs.get(name).cloned()
727 }
728
729 /// Declared names in declaration order.
730 pub fn names(&self) -> Vec<String> {
731 self.order.clone()
732 }
733}