1use std::collections::HashMap;
10use std::path::{Path, PathBuf};
11use std::sync::Arc;
12use std::sync::atomic::AtomicU64;
13
14use axum::{
15 extract::{Path as AxPath, State, Query as AxQuery},
16 http::{HeaderMap, StatusCode},
17 response::{IntoResponse, Response, sse::{Event, KeepAlive, Sse}},
18 routing::{delete, get, post},
19 Json, Router,
20};
21use dashmap::DashMap;
22use serde::Deserialize;
23use serde_json::{json, Value};
24use tokio::sync::{broadcast, RwLock};
25use tokio_stream::wrappers::BroadcastStream;
26use tokio_stream::StreamExt as _;
27
28use crate::db::Db;
29use crate::nql;
30use crate::store::Node;
31
32const LOG_CHANNEL_CAP: usize = 512;
35const SUB_CHANNEL_CAP: usize = 256;
36
37type SubKey = (String, u64); type SubVal = (String, String, broadcast::Sender<String>); macro_rules! nlog {
47 ($tx:expr, $($arg:tt)*) => {{
48 let line = format!($($arg)*);
49 println!("{}", line);
50 let _ = $tx.send(line);
51 }};
52}
53
54#[derive(Clone)]
57pub struct Manager {
58 inner: Arc<RwLock<ManagerInner>>,
59 pub token: Option<String>,
60 pub log_tx: broadcast::Sender<String>,
62 subs: Arc<DashMap<SubKey, SubVal>>,
64 sub_ctr: Arc<AtomicU64>,
65 #[cfg(feature = "cast")]
69 pub caster: Option<crate::cast::Caster>,
70}
71
72struct ManagerInner {
73 data_dir: PathBuf,
74 dbs: HashMap<String, Arc<Db>>,
75 tmk: Option<[u8; 32]>,
76 memory_mode: bool,
77}
78
79impl Manager {
80 pub fn new(data_dir: &Path, tmk: Option<[u8; 32]>, token: Option<String>, memory_mode: bool) -> Self {
81 let (log_tx, _) = broadcast::channel(LOG_CHANNEL_CAP);
82 Self {
83 inner: Arc::new(RwLock::new(ManagerInner {
84 data_dir: data_dir.to_path_buf(),
85 dbs: HashMap::new(),
86 tmk,
87 memory_mode,
88 })),
89 token,
90 log_tx,
91 #[cfg(feature = "cast")]
92 caster: None,
93 subs: Arc::new(DashMap::new()),
94 sub_ctr: Arc::new(AtomicU64::new(1)),
95 }
96 }
97
98 fn subscribe(&self, db: &str, nql: String) -> (u64, broadcast::Receiver<String>) {
100 use std::sync::atomic::Ordering;
101 let id = self.sub_ctr.fetch_add(1, Ordering::Relaxed);
102 let (tx, rx) = broadcast::channel(SUB_CHANNEL_CAP);
103 self.subs.insert((db.to_string(), id), (nql, String::new(), tx));
104 (id, rx)
105 }
106
107 fn unsubscribe(&self, db: &str, sub_id: u64) {
109 self.subs.remove(&(db.to_string(), sub_id));
110 }
111
112 fn notify_subscribers(&self, db: &str, db_arc: &Arc<crate::db::Db>) {
114 let keys: Vec<SubKey> = self.subs.iter()
115 .filter(|e| e.key().0 == db)
116 .map(|e| e.key().clone())
117 .collect();
118
119 for key in keys {
120 if let Some(mut entry) = self.subs.get_mut(&key) {
121 let (nql, last_hash, tx) = entry.value_mut();
122 let rows = match crate::nql::query(db_arc, nql) {
124 Ok((rows, _)) => rows,
125 Err(_) => continue,
126 };
127 let new_hash = format!("{:?}", rows.iter().map(|r| r.to_string()).collect::<Vec<_>>());
129 if new_hash == *last_hash { continue; }
130 *last_hash = new_hash;
131 let event = json!({
133 "sub_id": key.1,
134 "db": &key.0,
135 "nql": nql.as_str(),
136 "rows": rows,
137 "count": rows.len(),
138 });
139 let _ = tx.send(event.to_string());
140 }
141 }
142 }
143
144 pub async fn open_all(&self) -> anyhow::Result<()> {
146 let (data_dir, tmk, memory_mode) = {
147 let inner = self.inner.read().await;
148 (inner.data_dir.clone(), inner.tmk, inner.memory_mode)
149 };
150 if memory_mode { return Ok(()); }
152 if !data_dir.exists() {
153 std::fs::create_dir_all(&data_dir)?;
154 return Ok(());
155 }
156 let mut names = vec![];
157 for entry in std::fs::read_dir(&data_dir)? {
158 let entry = entry?;
159 if entry.file_type()?.is_dir() {
160 names.push(entry.file_name().to_string_lossy().to_string());
161 }
162 }
163 let log_tx = self.log_tx.clone();
164 let mut inner = self.inner.write().await;
165 for name in names {
166 let db_path = inner.data_dir.join(&name);
167 let dek = tmk.map(|k| crate::store::Dek::from_tmk(&k, name.as_bytes()));
168 match Db::open(&db_path, dek) {
169 Ok(db) => {
170 nlog!(log_tx, " [nedbd] opened database {:?}", name);
171 let db_arc = Arc::new(db);
172 Db::start_cold_scan(Arc::clone(&db_arc));
173 Db::start_manifest_ticker(Arc::clone(&db_arc), 1000);
175 inner.dbs.insert(name, db_arc);
176 }
177 Err(e) => nlog!(log_tx, " [nedbd] ERROR opening {:?}: {}", name, e),
178 }
179 }
180 Ok(())
181 }
182
183 async fn get_db(&self, name: &str) -> Option<Arc<Db>> {
184 self.inner.read().await.dbs.get(name).cloned()
185 }
186
187 async fn create_db(&self, name: &str) -> anyhow::Result<Arc<Db>> {
188 let (data_dir, tmk, memory_mode) = {
189 let inner = self.inner.read().await;
190 (inner.data_dir.clone(), inner.tmk, inner.memory_mode)
191 };
192 let db = if memory_mode {
193 Arc::new(Db::in_memory())
195 } else {
196 let db_path = data_dir.join(name);
197 let dek = tmk.map(|k| crate::store::Dek::from_tmk(&k, name.as_bytes()));
198 let db = Arc::new(Db::open(&db_path, dek)?);
199 Db::start_cold_scan(Arc::clone(&db));
200 Db::start_manifest_ticker(Arc::clone(&db), 1000);
201 db
202 };
203 self.inner.write().await.dbs.insert(name.to_string(), db.clone());
204 Ok(db)
205 }
206
207 async fn drop_db(&self, name: &str) -> bool {
208 let db = self.inner.write().await.dbs.remove(name);
209 if let Some(db) = db {
210 db.flush_manifest_if_dirty();
212 let data_dir = self.inner.read().await.data_dir.clone();
213 let _ = std::fs::remove_dir_all(data_dir.join(name));
214 true
215 } else {
216 false
217 }
218 }
219
220 pub async fn flush_all(&self) {
222 let inner = self.inner.read().await;
223 for db in inner.dbs.values() {
224 db.flush_all(); }
226 }
227
228 async fn names(&self) -> Vec<String> {
229 self.inner.read().await.dbs.keys().cloned().collect()
230 }
231
232 pub fn log(&self, msg: impl Into<String>) {
234 let line = msg.into();
235 println!("{}", line);
236 let _ = self.log_tx.send(line);
237 }
238
239 fn check_auth(&self, headers: &HeaderMap) -> bool {
240 match &self.token {
241 None => true,
242 Some(required) => {
243 if let Some(auth) = headers.get("authorization") {
244 if let Ok(s) = auth.to_str() {
245 return s == format!("Bearer {}", required);
246 }
247 }
248 false
249 }
250 }
251 }
252}
253
254fn err(status: StatusCode, msg: &str) -> Response {
257 (status, Json(json!({"error": msg}))).into_response()
258}
259
260fn ok(body: Value) -> Response {
261 (StatusCode::OK, Json(body)).into_response()
262}
263
264fn db_seq_head(db: &Db) -> (u64, String) {
268 let seq = db.seq.load(std::sync::atomic::Ordering::SeqCst);
269 let head = db.head();
270 (seq, head)
271}
272
273async fn health(State(mgr): State<Manager>) -> Response {
276 let names = mgr.names().await;
277 let inner = mgr.inner.read().await;
278 ok(json!({
279 "ok": true,
280 "service": "nedbd",
281 "version": env!("CARGO_PKG_VERSION"),
282 "engine": "dag",
283 "memory": inner.memory_mode,
284 "databases": names,
285 "encrypted": inner.tmk.is_some(),
286 }))
287}
288
289async fn list_databases(State(mgr): State<Manager>, headers: HeaderMap) -> Response {
290 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
291 let names = mgr.names().await;
292 let summaries: Vec<Value> = {
293 let inner = mgr.inner.read().await;
294 names.iter().map(|n| {
295 if let Some(db) = inner.dbs.get(n) {
296 let (seq, head) = db_seq_head(db);
297 json!({"name": n, "seq": seq, "head": head, "collections": db.id_index.collections()})
298 } else {
299 json!({"name": n})
300 }
301 }).collect()
302 };
303 ok(json!({"databases": summaries}))
304}
305
306#[derive(Deserialize)]
307struct CreateDbBody { name: String }
308
309async fn create_database(
310 State(mgr): State<Manager>,
311 headers: HeaderMap,
312 Json(body): Json<CreateDbBody>,
313) -> Response {
314 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
315 if body.name.is_empty() { return err(StatusCode::BAD_REQUEST, "name is required"); }
316 match mgr.create_db(&body.name).await {
317 Ok(db) => {
318 let (seq, head) = db_seq_head(&db);
319 (StatusCode::CREATED, Json(json!({"database": {"name": body.name, "seq": seq, "head": head}}))).into_response()
320 }
321 Err(e) => err(StatusCode::INTERNAL_SERVER_ERROR, &e.to_string()),
322 }
323}
324
325async fn get_database(
326 State(mgr): State<Manager>,
327 headers: HeaderMap,
328 AxPath(name): AxPath<String>,
329) -> Response {
330 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
331 match mgr.get_db(&name).await {
332 None => err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
333 Some(db) => {
334 let (seq, head) = db_seq_head(&db);
335 ok(json!({"name": name, "seq": seq, "head": head, "collections": db.id_index.collections()}))
336 }
337 }
338}
339
340async fn drop_database(
341 State(mgr): State<Manager>,
342 headers: HeaderMap,
343 AxPath(name): AxPath<String>,
344) -> Response {
345 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
346 let dropped = mgr.drop_db(&name).await;
347 ok(json!({"dropped": dropped}))
348}
349
350#[derive(Deserialize)]
351struct QueryBody { nql: String }
352
353#[cfg_attr(not(feature = "cast"), allow(dead_code))]
361#[derive(Deserialize)]
362struct CastBody {
363 prompt: String,
364 #[serde(default)]
368 execute: bool,
369}
370
371#[cfg(feature = "cast")]
377async fn cast_prompt(
378 State(mgr): State<Manager>,
379 headers: HeaderMap,
380 AxPath(name): AxPath<String>,
381 Json(body): Json<CastBody>,
382) -> Response {
383 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
384
385 let caster = match &mgr.caster {
386 Some(c) => c,
387 None => return err(
388 StatusCode::SERVICE_UNAVAILABLE,
389 "cast is not enabled; start nedbd with --cast (or NEDBD_CAST=1) \
390 and place model.cast in the data directory",
391 ),
392 };
393
394 let db = match mgr.get_db(&name).await {
395 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
396 Some(db) => db,
397 };
398 if body.prompt.trim().is_empty() {
399 return err(StatusCode::BAD_REQUEST, "prompt is required");
400 }
401
402 let collections = db.id_index.collections();
405 let result = caster.cast_checked(&body.prompt, &collections);
406
407 let parse_err = match nql::parse(&result.nql) {
410 Ok(_) => None,
411 Err(e) => Some(e.to_string()),
412 };
413
414 let (seq, head) = db_seq_head(&db);
415 let mut out = json!({
416 "prompt": body.prompt,
417 "nql": result.nql,
418 "valid": parse_err.is_none(),
419 "collection": result.collection,
420 "collection_known": result.collection_known,
421 "collections": collections,
422 "executed": false,
423 "seq": seq,
424 "head": head,
425 });
426
427 if let Some(d) = &result.drift {
434 out["drift"] = json!(d);
435 }
436
437 if let Some(e) = parse_err {
438 out["error"] = json!(format!("NQL error: {}", e));
441 return (StatusCode::UNPROCESSABLE_ENTITY, Json(out)).into_response();
442 }
443
444 if !result.collection_known {
445 out["error"] = json!(format!(
449 "collection {:?} does not exist in {:?}",
450 result.collection.unwrap_or_default(), name
451 ));
452 return (StatusCode::UNPROCESSABLE_ENTITY, Json(out)).into_response();
453 }
454
455 if !body.execute {
456 return ok(out);
457 }
458
459 let nql_text = out["nql"].as_str().unwrap_or("").to_string();
461 match nql::query(&db, &nql_text) {
462 Ok((rows, count)) => {
463 out["executed"] = json!(true);
464 out["rows"] = json!(rows);
465 out["count"] = json!(count);
466 ok(out)
467 }
468 Err(e) => {
469 out["error"] = json!(format!("NQL error: {}", e));
470 (StatusCode::BAD_REQUEST, Json(out)).into_response()
471 }
472 }
473}
474
475#[cfg(not(feature = "cast"))]
478async fn cast_prompt(
479 State(_mgr): State<Manager>,
480 _headers: HeaderMap,
481 AxPath(_name): AxPath<String>,
482 Json(_body): Json<CastBody>,
483) -> Response {
484 err(
485 StatusCode::NOT_IMPLEMENTED,
486 "this nedbd was built without the `cast` feature; \
487 rebuild with --features cast to enable natural-language planning",
488 )
489}
490
491async fn query_database(
492 State(mgr): State<Manager>,
493 headers: HeaderMap,
494 AxPath(name): AxPath<String>,
495 Json(body): Json<QueryBody>,
496) -> Response {
497 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
498 let db = match mgr.get_db(&name).await {
499 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
500 Some(db) => db,
501 };
502 if body.nql.trim().is_empty() {
503 return err(StatusCode::BAD_REQUEST, "nql is required");
504 }
505 match nql::query(&db, &body.nql) {
506 Ok((rows, count)) => {
507 let (seq, head) = db_seq_head(&db);
508 ok(json!({"rows": rows, "count": count, "seq": seq, "head": head}))
509 }
510 Err(e) => err(StatusCode::BAD_REQUEST, &format!("NQL error: {}", e)),
511 }
512}
513
514#[derive(Deserialize)]
515struct PutBody {
516 coll: String,
517 id: String,
518 doc: Value,
519 caused_by: Option<Vec<serde_json::Value>>,
520 valid_from: Option<String>,
521 valid_to: Option<String>,
522 #[allow(dead_code)] evidence: Option<String>,
523 #[allow(dead_code)] confidence: Option<f64>,
524 #[allow(dead_code)] client: Option<String>,
525 #[allow(dead_code)] nonce: Option<u64>,
526 #[allow(dead_code)] idem: Option<String>,
527}
528
529#[derive(Deserialize)]
530struct LinkBody {
531 frm: String,
532 rel: String,
533 to: String,
534}
535
536async fn put_document(
537 State(mgr): State<Manager>,
538 headers: HeaderMap,
539 AxPath(name): AxPath<String>,
540 Json(body): Json<PutBody>,
541) -> Response {
542 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
543 let db = match mgr.get_db(&name).await {
544 None => {
545 match mgr.create_db(&name).await {
547 Ok(db) => db,
548 Err(e) => return err(StatusCode::INTERNAL_SERVER_ERROR, &e.to_string()),
549 }
550 }
551 Some(db) => db,
552 };
553 if !db.startup_ready.load(std::sync::atomic::Ordering::SeqCst) {
556 return err(StatusCode::SERVICE_UNAVAILABLE,
557 "database startup in progress — reads available, writes retry in a moment");
558 }
559 let caused_by: Vec<String> = body.caused_by.unwrap_or_default()
561 .into_iter()
562 .filter_map(|v| match v {
563 serde_json::Value::String(s) => Some(s),
564 serde_json::Value::Number(n) => {
565 n.as_u64().and_then(|seq| db.get_hash_by_seq(seq))
566 }
567 _ => None,
568 })
569 .collect();
570 let coll = body.coll.clone();
573 let id = body.id.clone();
574 let doc = body.doc.clone();
575 let vf = body.valid_from.clone();
576 let vt = body.valid_to.clone();
577 let db2 = Arc::clone(&db);
578 let result = tokio::task::spawn_blocking(move || {
579 db2.put(&coll, &id, doc, caused_by, vf, vt)
580 }).await;
581 match result {
582 Err(join_err) => err(StatusCode::INTERNAL_SERVER_ERROR, &join_err.to_string()),
583 Ok(Err(e)) => err(StatusCode::INTERNAL_SERVER_ERROR, &e.to_string()),
584 Ok(Ok(node)) => {
585 let (seq, head) = db_seq_head(&db);
586 mgr.notify_subscribers(&name, &db);
587 ok(json!({"ok": true, "doc": node_to_response(&node), "seq": seq, "head": head}))
588 }
589 }
590}
591
592fn node_to_response(node: &Node) -> Value {
593 json!({
594 "_id": node.id,
595 "_hash": node.hash,
596 "_seq": node.seq,
597 "_coll": node.coll,
598 "data": node.data,
599 })
600}
601
602async fn link_document(
603 State(mgr): State<Manager>,
604 headers: HeaderMap,
605 AxPath(name): AxPath<String>,
606 Json(body): Json<LinkBody>,
607) -> Response {
608 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
609 let db = match mgr.get_db(&name).await {
610 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
611 Some(db) => db,
612 };
613 if !db.startup_ready.load(std::sync::atomic::Ordering::SeqCst) {
614 return err(StatusCode::SERVICE_UNAVAILABLE, "startup scan in progress");
615 }
616 match db.link(&body.frm, &body.rel, &body.to) {
617 Ok(()) => {
618 let (seq, head) = db_seq_head(&db);
619 ok(json!({"ok": true, "frm": body.frm, "rel": body.rel, "to": body.to, "seq": seq, "head": head}))
620 }
621 Err(e) => err(StatusCode::BAD_REQUEST, &e.to_string()),
622 }
623}
624
625async fn get_document(
649 State(mgr): State<Manager>,
650 headers: HeaderMap,
651 AxPath((name, coll, id)): AxPath<(String, String, String)>,
652 AxQuery(q): AxQuery<GetRowQuery>,
653) -> Response {
654 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
655 let db = match mgr.get_db(&name).await {
656 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
657 Some(db) => db,
658 };
659 let node = match q.as_of {
660 Some(seq) => db.get_as_of(&coll, &id, seq),
661 None => db.get(&coll, &id),
662 };
663 let (seq, head) = db_seq_head(&db);
674 let row = match node {
675 None => Value::Null,
676 Some(n) => crate::nql::node_to_json(&n),
677 };
678 ok(json!({"row": row, "seq": seq, "head": head}))
679}
680
681#[derive(Deserialize, Default)]
682struct GetRowQuery {
683 as_of: Option<u64>,
684}
685
686async fn delete_document(
687 State(mgr): State<Manager>,
688 headers: HeaderMap,
689 AxPath((name, coll, id)): AxPath<(String, String, String)>,
690) -> Response {
691 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
692 let db = match mgr.get_db(&name).await {
693 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
694 Some(db) => db,
695 };
696 let existed = match db.delete(&coll, &id) {
699 Ok(v) => v,
700 Err(e) => return err(StatusCode::INTERNAL_SERVER_ERROR, &e.to_string()),
701 };
702 let (seq, head) = db_seq_head(&db);
703 ok(json!({"ok": existed, "seq": seq, "head": head}))
704}
705
706#[derive(Deserialize)]
707struct BatchOp {
708 op: String,
709 coll: Option<String>,
710 id: Option<String>,
711 doc: Option<Value>,
712 caused_by: Option<Vec<serde_json::Value>>,
713}
714#[derive(Deserialize)]
715struct BatchBody { ops: Vec<BatchOp> }
716
717async fn batch_operations(
718 State(mgr): State<Manager>,
719 headers: HeaderMap,
720 AxPath(name): AxPath<String>,
721 Json(body): Json<BatchBody>,
722) -> Response {
723 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
724 let db = match mgr.get_db(&name).await {
725 None => match mgr.create_db(&name).await {
726 Ok(db) => db,
727 Err(e) => return err(StatusCode::INTERNAL_SERVER_ERROR, &e.to_string()),
728 },
729 Some(db) => db,
730 };
731
732 if !db.startup_ready.load(std::sync::atomic::Ordering::SeqCst) {
733 return err(StatusCode::SERVICE_UNAVAILABLE,
734 "database startup in progress — reads available, writes retry in a moment");
735 }
736
737 let mut put_ops = vec![];
741 let mut del_ops: Vec<(String, String)> = vec![];
742 let mut op_order: Vec<(&str, usize)> = vec![]; for op in &body.ops {
745 let t = op.op.to_lowercase();
746 match t.as_str() {
747 "put" => {
748 let caused_by: Vec<String> = op.caused_by.clone().unwrap_or_default()
750 .into_iter()
751 .filter_map(|v| match v {
752 serde_json::Value::String(s) => Some(s),
753 serde_json::Value::Number(n) => {
754 n.as_u64().and_then(|seq| db.get_hash_by_seq(seq))
755 }
756 _ => None,
757 })
758 .collect();
759 op_order.push(("put", put_ops.len()));
760 put_ops.push((
761 op.coll.clone().unwrap_or_default(),
762 op.id.clone().unwrap_or_default(),
763 op.doc.clone().unwrap_or(json!({})),
764 caused_by,
765 None::<String>,
766 None::<String>,
767 ));
768 }
769 "del" | "delete" => {
770 op_order.push(("del", del_ops.len()));
771 del_ops.push((
772 op.coll.clone().unwrap_or_default(),
773 op.id.clone().unwrap_or_default(),
774 ));
775 }
776 _ => { op_order.push(("unknown", 0)); }
777 }
778 }
779
780 let put_results = if put_ops.is_empty() {
782 vec![]
783 } else {
784 match db.put_batch(put_ops) {
785 Ok(nodes) => nodes.into_iter().map(|n| json!({"op":"put","id":n.id,"seq":n.seq,"hash":n.hash})).collect(),
786 Err(e) => return err(StatusCode::INTERNAL_SERVER_ERROR, &e.to_string()),
787 }
788 };
789
790 let del_results: Vec<serde_json::Value> = del_ops.iter().map(|(coll, id)| {
792 match db.delete(coll, id) {
793 Ok(existed) => json!({"op":"del","id":id,"ok":existed}),
794 Err(e) => json!({"op":"del","id":id,"error":e.to_string()}),
795 }
796 }).collect();
797
798 let mut results = vec![];
800 for (kind, idx) in &op_order {
801 let r = match *kind {
802 "put" => put_results.get(*idx).cloned().unwrap_or(json!({"op":"put","error":"missing"})),
803 "del" => del_results.get(*idx).cloned().unwrap_or(json!({"op":"del","error":"missing"})),
804 _ => json!({"op": kind, "error": "unknown op"}),
805 };
806 results.push(r);
807 }
808 let (seq, head) = db_seq_head(&db);
809 mgr.notify_subscribers(&name, &db);
811 ok(json!({"results": results, "count": results.len(), "seq": seq, "head": head}))
812}
813
814#[derive(Deserialize)]
815struct IndexBody { coll: String, field: String, kind: Option<String> }
816
817async fn create_index(
818 State(mgr): State<Manager>,
819 headers: HeaderMap,
820 AxPath(name): AxPath<String>,
821 Json(body): Json<IndexBody>,
822) -> Response {
823 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
824 let db = match mgr.get_db(&name).await {
825 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
826 Some(db) => db,
827 };
828 let kind = body.kind.as_deref().unwrap_or("eq");
829 match kind {
830 "sorted" | "eq" => {
831 db.create_sorted_index(&body.coll, &body.field);
832 ok(json!({"ok": true, "coll": body.coll, "field": body.field, "kind": kind}))
833 }
834 _ => err(StatusCode::BAD_REQUEST, &format!("unknown index kind: {}", kind)),
835 }
836}
837
838async fn verify_database(
839 State(mgr): State<Manager>,
840 headers: HeaderMap,
841 AxPath(name): AxPath<String>,
842) -> Response {
843 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
844 let db = match mgr.get_db(&name).await {
845 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
846 Some(db) => db,
847 };
848 let (ok_count, tampered) = db.verify();
849 let (seq, head) = db_seq_head(&db);
850 ok(json!({
851 "ok": tampered.is_empty(),
852 "seq": seq,
853 "head": head,
854 "tamper_evident": true,
855 "objects_checked": ok_count,
856 "tampered": tampered,
857 }))
858}
859
860async fn checkpoint(
861 State(mgr): State<Manager>,
862 headers: HeaderMap,
863 AxPath(name): AxPath<String>,
864) -> Response {
865 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
866 let db = match mgr.get_db(&name).await {
867 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
868 Some(db) => db,
869 };
870 let (seq, head) = db_seq_head(&db);
871 ok(json!({"ok": true, "head": head, "seq": seq}))
873}
874
875#[derive(Deserialize)]
876struct LogQuery { limit: Option<usize> }
877
878async fn get_log(
879 State(mgr): State<Manager>,
880 headers: HeaderMap,
881 AxPath(name): AxPath<String>,
882 AxQuery(q): AxQuery<LogQuery>,
883) -> Response {
884 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
885 let db = match mgr.get_db(&name).await {
886 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
887 Some(db) => db,
888 };
889 let limit = q.limit.unwrap_or(50);
890 let mut log_entries: Vec<Value> = db.objects.all_hashes()
892 .filter_map(|h| db.objects.read(&h).ok())
893 .take(limit)
894 .map(|n| json!({
895 "seq": n.seq, "coll": n.coll, "id": n.id,
896 "hash": n.hash, "ts": n.ts, "op": "put"
897 }))
898 .collect();
899 log_entries.sort_by(|a, b|
900 b["seq"].as_u64().cmp(&a["seq"].as_u64())
901 );
902 log_entries.truncate(limit);
903 let (seq, head) = db_seq_head(&db);
904 ok(json!({"log": log_entries, "seq": seq, "head": head}))
905}
906
907async fn tip_database(
913 State(mgr): State<Manager>,
914 headers: HeaderMap,
915 AxPath(name): AxPath<String>,
916) -> Response {
917 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
918 let db = match mgr.get_db(&name).await {
919 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
920 Some(db) => db,
921 };
922 let (seq, head) = db_seq_head(&db);
923 let tip = db.tip().map(|n| serde_json::to_value(&n).unwrap_or(Value::Null));
924 ok(json!({"tip": tip, "seq": seq, "head": head}))
925}
926
927async fn tip_collection_database(
929 State(mgr): State<Manager>,
930 headers: HeaderMap,
931 AxPath((name, coll)): AxPath<(String, String)>,
932) -> Response {
933 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
934 let db = match mgr.get_db(&name).await {
935 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
936 Some(db) => db,
937 };
938 let (seq, head) = db_seq_head(&db);
939 let tip = db.tip_collection(&coll).map(|n| serde_json::to_value(&n).unwrap_or(Value::Null));
940 ok(json!({"coll": coll, "tip": tip, "seq": seq, "head": head}))
941}
942
943#[derive(Deserialize)]
944struct SinceQuery { after_seq: Option<u64>, limit: Option<usize> }
945
946async fn since_database(
947 State(mgr): State<Manager>,
948 headers: HeaderMap,
949 AxPath(name): AxPath<String>,
950 AxQuery(q): AxQuery<SinceQuery>,
951) -> Response {
952 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
953 let db = match mgr.get_db(&name).await {
954 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
955 Some(db) => db,
956 };
957 let after = q.after_seq.unwrap_or(0);
958 let b = db.since(after, q.limit.unwrap_or(0));
959 let nodes: Vec<Value> = b.nodes.iter()
960 .map(|n| serde_json::to_value(n).unwrap_or(Value::Null))
961 .collect();
962 let (seq, head) = db_seq_head(&db);
963 ok(json!({
964 "nodes": nodes, "count": nodes.len(),
965 "from_seq": b.from_seq, "to_seq": b.to_seq, "head_seq": b.head_seq, "has_more": b.has_more,
966 "seq": seq, "head": head
967 }))
968}
969
970async fn status_database(
973 State(mgr): State<Manager>,
974 headers: HeaderMap,
975 AxPath(name): AxPath<String>,
976) -> Response {
977 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
978 let db = match mgr.get_db(&name).await {
979 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
980 Some(db) => db,
981 };
982 let s = db.scan_status();
983 ok(json!({
984 "ok": true,
985 "scan_complete": s.scan_complete,
986 "tip_seq": s.tip_seq,
987 "indexed_seq_min": s.indexed_seq_min,
988 "indexed_seq_max": s.indexed_seq_max,
989 "indexed_count": s.indexed_count
990 }))
991}
992
993#[derive(Deserialize)]
996struct SubscribeBody { nql: String }
997
998async fn subscribe_query(
999 State(mgr): State<Manager>,
1000 headers: HeaderMap,
1001 AxPath(name): AxPath<String>,
1002 Json(body): Json<SubscribeBody>,
1003) -> Response {
1004 if !mgr.check_auth(&headers) {
1005 return err(StatusCode::UNAUTHORIZED, "unauthorized");
1006 }
1007 let db = match mgr.get_db(&name).await {
1008 None => return err(StatusCode::NOT_FOUND, &format!("database not found: {}", name)),
1009 Some(db) => db,
1010 };
1011
1012 let (sub_id, rx) = mgr.subscribe(&name, body.nql.clone());
1013
1014 if let Ok((rows, _)) = crate::nql::query(&db, &body.nql) {
1016 let init = json!({
1017 "sub_id": sub_id,
1018 "db": &name,
1019 "nql": &body.nql,
1020 "rows": rows,
1021 "count": rows.len(),
1022 "event": "initial",
1023 });
1024 if let Some(mut entry) = mgr.subs.get_mut(&(name.clone(), sub_id)) {
1026 let hash = format!("{:?}", rows);
1027 entry.value_mut().1 = hash;
1028 }
1029 if let Some(entry) = mgr.subs.get(&(name.clone(), sub_id)) {
1031 let _ = entry.value().2.send(init.to_string());
1032 }
1033 }
1034
1035 let stream = BroadcastStream::new(rx).filter_map(|msg| {
1036 match msg {
1037 Ok(line) => Some(Ok::<Event, std::convert::Infallible>(Event::default().data(line))),
1038 Err(_) => None,
1039 }
1040 });
1041 Sse::new(stream)
1042 .keep_alive(KeepAlive::default())
1043 .into_response()
1044}
1045
1046async fn unsubscribe_query(
1047 State(mgr): State<Manager>,
1048 headers: HeaderMap,
1049 AxPath((name, sub_id)): AxPath<(String, u64)>,
1050) -> Response {
1051 if !mgr.check_auth(&headers) { return err(StatusCode::UNAUTHORIZED, "unauthorized"); }
1052 mgr.unsubscribe(&name, sub_id);
1053 ok(json!({"ok": true, "sub_id": sub_id}))
1054}
1055
1056async fn log_events(State(mgr): State<Manager>) -> Sse<impl futures_core::Stream<Item = Result<Event, std::convert::Infallible>>> {
1059 let rx = mgr.log_tx.subscribe();
1060 let stream = BroadcastStream::new(rx).filter_map(|msg| {
1061 match msg {
1062 Ok(line) => Some(Ok::<Event, std::convert::Infallible>(Event::default().data(line))),
1063 Err(_) => None, }
1065 });
1066 Sse::new(stream).keep_alive(KeepAlive::default())
1067}
1068
1069pub fn router(mgr: Manager) -> Router {
1072 Router::new()
1073 .route("/health", get(health))
1074 .route("/events", get(log_events))
1075 .route("/v1/databases", get(list_databases).post(create_database))
1076 .route("/v1/databases/:name", get(get_database).delete(drop_database))
1077 .route("/v1/databases/:name/query", post(query_database))
1078 .route("/v1/databases/:name/cast", post(cast_prompt))
1079 .route("/v1/databases/:name/put", post(put_document))
1080 .route("/v1/databases/:name/link", post(link_document))
1081 .route("/v1/databases/:name/rows/:coll/:id",
1084 get(get_document).delete(delete_document))
1085 .route("/v1/databases/:name/batch", post(batch_operations))
1086 .route("/v1/databases/:name/index", post(create_index))
1087 .route("/v1/databases/:name/verify", get(verify_database))
1088 .route("/v1/databases/:name/checkpoint", post(checkpoint))
1089 .route("/v1/databases/:name/log", get(get_log))
1090 .route("/v1/databases/:name/tip", get(tip_database))
1091 .route("/v1/databases/:name/collections/:coll/tip", get(tip_collection_database))
1092 .route("/v1/databases/:name/since", get(since_database))
1093 .route("/v1/databases/:name/status", get(status_database))
1094 .route("/v1/databases/:name/subscribe", post(subscribe_query))
1095 .route("/v1/databases/:name/subscribe/:sub_id", delete(unsubscribe_query))
1096 .with_state(mgr)
1097}
1098
1099pub async fn run(host: &str, port: u16, data_dir: &str, tmk: Option<[u8; 32]>, token: Option<String>, memory_mode: bool) -> anyhow::Result<()> {
1101 #[cfg(feature = "cast")]
1106 let mut mgr = Manager::new(Path::new(data_dir), tmk, token, memory_mode);
1107 #[cfg(not(feature = "cast"))]
1108 let mgr = Manager::new(Path::new(data_dir), tmk, token, memory_mode);
1109
1110 mgr.open_all().await?;
1111
1112 #[cfg(feature = "cast")]
1116 {
1117 let want = std::env::var("NEDBD_CAST").map(|v| v == "1").unwrap_or(false);
1118 if want {
1119 match crate::cast::Caster::load(Path::new(data_dir)) {
1120 Ok(c) => {
1121 println!(" cast enabled — {:.2}M params, vocab {}, {}",
1122 c.n_params() as f64 / 1e6, c.vocab_size(), c.source());
1123 mgr.caster = Some(c);
1124 }
1125 Err(e) => {
1126 eprintln!(" cast DISABLED — {}", e);
1127 }
1128 }
1129 }
1130 }
1131 #[cfg(feature = "cast")]
1134 let mgr = mgr;
1135
1136 let has_token = mgr.token.is_some();
1137 let mgr_for_shutdown = mgr.clone();
1138 let app = router(mgr);
1139 let addr = format!("{}:{}", host, port).parse::<std::net::SocketAddr>()?;
1140 let banner = format!(r#"
1141 ◆
1142 ╱ ╲ N E D B · DAG ENGINE {}
1143 ◆ ◆ ─────────────────────────────────────────────
1144 ╱ ╲ ╱ ╲ content-addressed · tamper-evident · causal
1145 ◆ ◆ ◆ bi-temporal · replay-protected · encrypted
1146 ╱ ╲ ╱ ╲ ╱ ╲
1147 ◆ ◆ ◆ ◆ © INTERCHAINED LLC × Vex (Interchained AI fleet: GLM · Claude · Opus · Fable · GPT-6)
1148 ╱ ╲ ╱ ╲ ╱ ╲ ╱ ╲ interchained.org · hyperagent.com/refer/J2G6TCD7
1149
1150 ─────────────────────────────────────────────────────────────
1151 listen http://{}
1152 data {}
1153 enc {}
1154 token {}
1155 memory {}
1156 ─────────────────────────────────────────────────────────────
1157"#,
1158 env!("CARGO_PKG_VERSION"),
1159 addr,
1160 data_dir,
1161 if tmk.is_some() { "AES-256-GCM" } else { "off" },
1162 if has_token { "on" } else { "off (set NEDBD_TOKEN to require auth)" },
1163 if memory_mode { "yes — all data lost on exit (NEDBD_MEMORY=1)" } else { "no — durable DAG on disk" }
1164 );
1165 print!("{}", banner);
1166
1167 let listener = tokio::net::TcpListener::bind(addr).await?;
1168
1169 let mgr_hourly = mgr_for_shutdown.clone();
1173 tokio::spawn(async move {
1174 loop {
1175 let now_secs = std::time::SystemTime::now()
1177 .duration_since(std::time::UNIX_EPOCH)
1178 .map(|d| d.as_secs()).unwrap_or(0);
1179 let secs_into_hour = now_secs % 3600;
1180 let sleep_secs = 3600 - secs_into_hour;
1181 tokio::time::sleep(tokio::time::Duration::from_secs(sleep_secs)).await;
1182 mgr_hourly.flush_all().await;
1183 println!(" [nedbd] hourly checkpoint — manifests flushed");
1184 }
1185 });
1186
1187 let shutdown = async {
1189 #[cfg(unix)]
1190 {
1191 use tokio::signal::unix::{signal, SignalKind};
1192 let mut sigterm = signal(SignalKind::terminate()).unwrap();
1193 let mut sigint = signal(SignalKind::interrupt()).unwrap();
1194 tokio::select! {
1195 _ = sigterm.recv() => println!(" [nedbd] SIGTERM — flushing and exiting..."),
1196 _ = sigint.recv() => println!(" [nedbd] SIGINT — flushing and exiting..."),
1197 }
1198 }
1199 #[cfg(not(unix))]
1200 {
1201 tokio::signal::ctrl_c().await.ok();
1202 println!(" [nedbd] shutting down — flushing manifests...");
1203 }
1204 };
1205
1206 axum::serve(listener, app)
1207 .tcp_nodelay(true)
1208 .with_graceful_shutdown(shutdown)
1209 .await?;
1210
1211 mgr_for_shutdown.flush_all().await;
1213 println!(" [nedbd] goodbye");
1214 Ok(())
1215}