Skip to main content

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}