use std::collections::HashMap;
use std::path::Path;
use std::sync::Arc;
use std::thread::JoinHandle;
use horon::{
Horon, HoronConfig, HoronResult, DurabilityMode,
};
#[derive(Clone)]
pub struct GeoStore {
inner: Arc<Horon>,
}
impl GeoStore {
pub fn open<P: AsRef<Path>>(path: P) -> HoronResult<Self> {
Ok(Self {
inner: Arc::new(Horon::open(path)?),
})
}
pub fn open_with_config<P: AsRef<Path>>(
path: P,
config: HoronConfig,
) -> HoronResult<Self> {
Ok(Self {
inner: Arc::new(Horon::open_with_config(path, config)?),
})
}
pub fn open_embedded<P: AsRef<Path>>(path: P) -> HoronResult<Self> {
Self::open_with_config(path, HoronConfig {
auto_compact_threshold: 1_000,
durability: DurabilityMode::Fsync,
..Default::default()
})
}
pub fn open_server<P: AsRef<Path>>(path: P) -> HoronResult<Self> {
Self::open_with_config(path, HoronConfig {
auto_compact_threshold: 50_000,
wal_batch_size: 64,
wal_flush_interval_ms: 100,
durability: DurabilityMode::Batched,
..Default::default()
})
}
pub fn open_analytics<P: AsRef<Path>>(path: P) -> HoronResult<Self> {
Self::open_with_config(path, HoronConfig {
auto_compact_threshold: 100_000,
wal_batch_size: 256,
wal_flush_interval_ms: 500,
durability: DurabilityMode::Relaxed,
..Default::default()
})
}
pub fn put(&self, key: &str, data: &[u8]) -> HoronResult<()> {
self.inner.put(key, data)
}
pub fn get(&self, key: &str) -> HoronResult<Vec<u8>> {
self.inner.get(key)
}
pub fn remove(&self, key: &str) -> HoronResult<()> {
self.inner.remove(key)
}
pub fn exists(&self, key: &str) -> bool {
self.inner.exists(key)
}
pub fn set_meta(&self, key: &str, name: &str, value: &str) -> HoronResult<()> {
self.inner.set_meta(key, name, value)
}
pub fn get_meta(&self, key: &str) -> HoronResult<HashMap<String, String>> {
self.inner.get_meta(key)
}
pub fn children(&self, path: &str) -> HoronResult<Vec<String>> {
self.inner.children(path)
}
pub fn list(&self, prefix: &str) -> HoronResult<Vec<String>> {
self.inner.list(prefix)
}
pub fn nearest(&self, coords: &[g_math::fixed_point::FixedPoint])
-> HoronResult<(String, g_math::fixed_point::FixedPoint)> {
self.inner.nearest(coords)
}
pub fn neighbors(&self, key: &str, k: usize) -> HoronResult<Vec<String>> {
self.inner.neighbors(key, k)
}
pub fn compact(&self) -> HoronResult<bool> {
self.inner.compact()
}
pub fn compact_async(&self) -> JoinHandle<HoronResult<bool>> {
self.inner.compact_async()
}
pub fn flush(&self) -> HoronResult<()> {
self.inner.flush()
}
pub fn len(&self) -> usize {
self.inner.len()
}
pub fn is_empty(&self) -> bool {
self.inner.is_empty()
}
pub fn wal_len(&self) -> u32 {
self.inner.wal_len()
}
}
fn main() {
println!("=== Service Registry ===\n");
{
let db = GeoStore::open_server("/tmp/geostore_services.htt").unwrap();
db.put("/org/platform/auth", b"auth-service-v2.3").unwrap();
db.put("/org/platform/gateway", b"api-gateway-v1.8").unwrap();
db.put("/org/ml/inference", b"bert-serving-v4.1").unwrap();
db.put("/org/ml/training", b"trainer-gpu-v2.0").unwrap();
db.put("/org/data/ingest", b"kafka-bridge-v3.2").unwrap();
db.put("/org/data/warehouse", b"clickhouse-proxy-v1.1").unwrap();
db.set_meta("/org/ml/inference", "gpu", "a100").unwrap();
db.set_meta("/org/ml/inference", "replicas", "8").unwrap();
db.set_meta("/org/data/ingest", "throughput", "500k-msg-sec").unwrap();
let ml_services = db.children("/org/ml").unwrap();
println!(" ML services: {:?}", ml_services);
let all = db.list("/org").unwrap();
println!(" All registered: {} services", all.len());
let (closest, dist) = db.nearest(&fp(&[0.0, 0.0, 0.0, 0.0])).unwrap();
println!(" Nearest to origin: {} (dist={:.4})", closest, dist);
let neighbors = db.neighbors("/org/platform/auth", 3).unwrap();
println!(" Auth neighbors: {:?}", neighbors);
db.flush().unwrap();
}
println!("\n=== Config Store ===\n");
{
let db = GeoStore::open_embedded("/tmp/geostore_config.htt").unwrap();
db.put("/config/global/timeout_ms", b"5000").unwrap();
db.put("/config/global/retry_count", b"3").unwrap();
db.put("/config/global/log_level", b"info").unwrap();
db.put("/config/prod/timeout_ms", b"10000").unwrap();
db.put("/config/prod/log_level", b"warn").unwrap();
db.put("/config/staging/log_level", b"debug").unwrap();
db.put("/config/prod/auth/timeout_ms", b"30000").unwrap();
db.put("/config/prod/auth/max_sessions", b"10000").unwrap();
fn resolve_config(db: &GeoStore, env: &str, service: &str, key: &str) -> String {
let paths = [
format!("/config/{}/{}/{}", env, service, key),
format!("/config/{}/{}", env, key),
format!("/config/global/{}", key),
];
for path in &paths {
if let Ok(data) = db.get(path) {
return String::from_utf8_lossy(&data).to_string();
}
}
"<unset>".to_string()
}
let timeout = resolve_config(&db, "prod", "auth", "timeout_ms");
let log_lvl = resolve_config(&db, "prod", "auth", "log_level");
let retries = resolve_config(&db, "prod", "auth", "retry_count");
println!(" prod/auth/timeout_ms = {} (service override)", timeout);
println!(" prod/auth/log_level = {} (env override)", log_lvl);
println!(" prod/auth/retry_count = {} (global default)", retries);
let staging_keys = db.list("/config/staging").unwrap();
println!(" Staging overrides: {:?}", staging_keys);
db.flush().unwrap();
}
println!("\n=== Concurrent Content Store ===\n");
{
let db = GeoStore::open_server("/tmp/geostore_content.htt").unwrap();
let mut handles = vec![];
for team in ["frontend", "backend", "infra", "security"] {
let db = db.clone(); let team = team.to_string();
handles.push(std::thread::spawn(move || {
for i in 0..10 {
let key = format!("/docs/{}/page_{}", team, i);
let content = format!("{} documentation page {}", team, i);
db.put(&key, content.as_bytes()).unwrap();
db.set_meta(&key, "team", &team).unwrap();
db.set_meta(&key, "version", &i.to_string()).unwrap();
}
}));
}
for _ in 0..4 {
let db = db.clone();
handles.push(std::thread::spawn(move || {
for i in 0..10 {
let _ = db.get(&format!("/docs/frontend/page_{}", i));
let _ = db.exists(&format!("/docs/backend/page_{}", i));
}
}));
}
for h in handles {
h.join().unwrap();
}
println!(" Total entries: {}", db.len());
println!(" WAL entries: {}", db.wal_len());
let compact_handle = db.compact_async();
let teams = db.children("/docs").unwrap();
println!(" Team namespaces: {:?}", teams);
for team in &teams {
let pages = db.children(team).unwrap();
println!(" {} has {} pages", team, pages.len());
}
compact_handle.join().unwrap().unwrap();
println!(" After compaction: WAL entries = {}", db.wal_len());
db.flush().unwrap();
}
let _ = std::fs::remove_file("/tmp/geostore_services.htt");
let _ = std::fs::remove_file("/tmp/geostore_config.htt");
let _ = std::fs::remove_file("/tmp/geostore_content.htt");
println!("\nDone.");
}
fn fp(vals: &[f64]) -> Vec<g_math::fixed_point::FixedPoint> {
vals.iter().map(|&v| g_math::fixed_point::FixedPoint::from_f64(v)).collect()
}