From c241f1ce9361d7c58599db3f964319e899ccf850 Mon Sep 17 00:00:00 2001 From: Joep Meindertsma Date: Thu, 10 Sep 2026 10:51:58 +0200 Subject: [PATCH 1/3] Separate native node lifetime from the HTTP adapter --- CHANGELOG.md | 1 + TESTING_COVERAGE.md | 13 +++ desktop/README.md | 14 +++ desktop/src/lib.rs | 31 ++++--- planning/README.md | 2 +- planning/atomic-lib-runtime.md | 33 ++++++- server/src/appstate.rs | 6 ++ server/src/serve.rs | 142 +++++++++++++++++++++++++++---- server/tests/it/http_optional.rs | 104 ++++++++++++++++++++++ server/tests/it/main.rs | 1 + 10 files changed, 316 insertions(+), 31 deletions(-) create mode 100644 server/tests/it/http_optional.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index 19d96cea17..0efdff551f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,7 @@ See [STATUS.md](server/STATUS.md) to learn more about which features will remain ## UNRELEASED +- Separate embedded node startup from the optional HTTP adapter. Tauri keeps native services alive if HTTP stops; its frontend still requires HTTP/WS. This starts the HTTP-optional runtime migration ([#1196](https://github.com/ontola/atomic-server/issues/1196), [#749](https://github.com/ontola/atomic-server/issues/749)). - The outbox drains over a live Iroh link too (`sync::peer::LivePeerCommitTransport`): a device with no hub in reach delivers its queued writes to a paired peer as signed `COMMIT` frames, which the peer validates and applies like a hub diff --git a/TESTING_COVERAGE.md b/TESTING_COVERAGE.md index 6bee8ca6b0..4eb7c74911 100644 --- a/TESTING_COVERAGE.md +++ b/TESTING_COVERAGE.md @@ -308,6 +308,19 @@ installed metadata and a dependency sentinel remain untouched. Corepack must have this pnpm version cached or be able to download it. A stub Cargo command verifies Clippy dispatch, staged input, and failure propagation; this fixture does not compile the Rust workspace. +## HTTP-optional native node lifecycle + +`server/tests/it/http_optional.rs` reserves the configured HTTP port, starts +`serve::run_node`, creates and queries a document through the native store/runtime, +and explicitly attempts HTTP binding. After that bind fails, native reads, +edits and change events still work. A second test keeps the hosted embedder +ready hook before HTTP binding. `serve::lifecycle_tests` verifies that +ending the owned flush worker releases redb and preserves its final write. +These are Rust glue tests, not desktop UI acceptance. Tauri still starts the +HTTP/WS adapter for its frontend. Missing: native frontend reads/commits/events, +attachment bytes, restore and peer sync through Tauri with no HTTP listener; +process-global Iroh teardown/restart is also not covered by this extraction. + ## Browser WebRTC transport (issue #1396) `browser/lib/src/webrtc-transport.test.ts` covers frame fragmentation/order, diff --git a/desktop/README.md b/desktop/README.md index 4a448a55a3..d3d3bbcb51 100644 --- a/desktop/README.md +++ b/desktop/README.md @@ -12,6 +12,20 @@ cargo tauri dev cargo tauri build ``` +## Node and HTTP lifecycle + +The embedded node is initialized by `atomic_server_lib::serve::run_node`. +Tauri binds its native commands to the shared `AtomicNode`, then explicitly +starts `serve_http` for the current webview. A stopped HTTP adapter does not +stop the native node's flush worker and peer tasks while the app is open. + +This is the first step of [HTTP-optional local runtime](../planning/atomic-lib-runtime.md#phase-7-tauri--android-without-loopback). +The current frontend **still requires HTTP/WebSocket**. There is no no-HTTP +user setting yet: reads, commits, subscriptions and attachment access must +move to native commands/events before removing the listener or the Actix +build dependency. Hosted servers keep the existing `serve` / `serve_with_hook` +entry points. + ## Running in development `cargo tauri dev` starts the front-end for you — `beforeDevCommand` in `tauri.conf.json` runs `pnpm -C browser/data-browser dev:tauri`, and the app points at `localhost:6747` (`devUrl`). diff --git a/desktop/src/lib.rs b/desktop/src/lib.rs index 1e38267dc3..13e8cde137 100644 --- a/desktop/src/lib.rs +++ b/desktop/src/lib.rs @@ -89,7 +89,7 @@ const PAIR_LINK_RETRY_WINDOW: std::time::Duration = std::time::Duration::from_se /// A handle on the embedded node, captured once it has booted. #[derive(Default)] struct EmbeddedNode { - store: std::sync::OnceLock, + runtime: std::sync::OnceLock, config_file: std::sync::OnceLock, /// Why the node never came up, if it didn't. The server runs on its own /// thread, so a boot failure there used to be an unwind into nothing: the @@ -102,8 +102,8 @@ struct EmbeddedNode { impl EmbeddedNode { /// The store, or an explanation of why there isn't one. fn require_store(&self) -> Result { - if let Some(store) = self.store.get() { - return Ok(store.clone()); + if let Some(node) = self.runtime.get() { + return Ok(node.db().clone()); } Err(match self.startup_error.get() { @@ -310,11 +310,7 @@ mod vault_ipc { } pub fn store_of(node: &std::sync::Arc) -> Result { - node - .store - .get() - .ok_or_else(|| "The local node has not finished starting up.".to_string()) - .cloned() + node.require_store() } /// Run vault work on a thread of its own. @@ -709,12 +705,21 @@ pub fn run() { .expect("TLS verifier initialization did not complete"); let rt = actix_rt::Runtime::new().unwrap(); - // The hook hands us the store once it's up, so `adopt_agent` can point - // the node's identity at the signed-in user. - let result = rt.block_on(atomic_server_lib::serve::serve_with_hook( + // Native commands receive the shared runtime before the HTTP adapter + // starts, so `adopt_agent` can set the node's signed-in identity. + let result = rt.block_on(atomic_server_lib::serve::run_node( config_clone, - |appstate| { - let _ = node_for_server.store.set(appstate.store.clone()); + |appstate| async { + let _ = node_for_server.runtime.set(appstate.node()); + // The current webview still needs HTTP/WS. It is an adapter, + // explicitly started after native access is ready; replacing the + // frontend's transport will let this call become optional. + if let Err(error) = atomic_server_lib::serve::serve_http(appstate).await { + eprintln!("[node] the HTTP adapter stopped: {error}"); + } + // The window and native commands outlive the HTTP adapter. Keep + // peer tasks and durable flushing alive even if its port is busy. + std::future::pending().await }, )); diff --git a/planning/README.md b/planning/README.md index 7233cc573c..d7895cbede 100644 --- a/planning/README.md +++ b/planning/README.md @@ -88,7 +88,7 @@ browser flow; standalone recovery remains self-managed. | [`authorization-sync.md`](./authorization-sync.md) | **Draft.** P1 done, P2 partial (`classify_auth_impact` exists). P3 open: no `AuthorizationProof`, signer still auto-inserted into `write`, `genesis_signer()` has no non-test caller. P4 open: `SYNC_PUSH` imports raw Loro on a drive-level verdict. | | [`unified-data-layer.md`](./unified-data-layer.md) | **Partial.** Browser/JS: one ingress, one outbox, one subscription model. Atomic writes, outbox, ingress entry points and immutable read/save subscriptions shipped. Remaining: consumer migration (38 `addResource` call sites, 42 `CommitBuilder` refs) and boundary extraction. | | [`loro-source-of-truth.md`](./loro-source-of-truth.md) | **Partial.** Sparse `datatypes` map + Phase 2a–2c shipped (`Tree::Resources` is a derived cache). Remaining: drop the untagged heuristic and the 3-way `build_state_doc` fallback, snapshot backfill (Phase 4), Phase 1.6 `Value` reshape (~966 sites), Flutter undo. | -| [`atomic-lib-runtime.md`](./atomic-lib-runtime.md) | **Partial.** `AtomicNode` is the binding runtime; the WASM `ClientDb` is its only adapter. Flutter and desktop still hold a raw `Db`; the node has no blob API; `NodeConfig` was cut. Library-owned durable flush is the first slice. Local KV FTS landed in [`local-search.md`](./local-search.md). | +| [`atomic-lib-runtime.md`](./atomic-lib-runtime.md) | **Partial.** `AtomicNode` is the binding runtime; the WASM `ClientDb` and Tauri's native handle bind it. Native storage/identity/durability now live in atomic_lib with a standalone no-Actix CI gate; node startup is separate from the HTTP adapter; the Tauri frontend still needs HTTP/WS. The node has no blob API; `NodeConfig` was cut. Open: native frontend transport and the other bindings. Local KV FTS landed in [`local-search.md`](./local-search.md). | | [`genesis-self-verifying.md`](./genesis-self-verifying.md) | **Partial.** Server and browser mint and verify inline genesis certs. Remaining: DataRoute verify UI, `genesis` propval immutability, cert-signed `drive` in `check_rights`. | | [`drive-reconciliation.md`](./drive-reconciliation.md) | **Partial.** Core in `lib/src/sync/rbsr.rs` + TS mirror; **on the WS wire** as the stateless text frames `RBSR_FP`/`RBSR_ITEMS` (full-VV fallback). Not on Iroh; fingerprints still O(range); canonical cross-impl hash unspecified, so the hash-first probe rarely matches. | | [`auditability-loro-history.md`](./auditability-loro-history.md) | **Mostly built.** `Tree::Envelopes`, `attribute_history`, `/history-attribution`, the History Verified badge, and since 2026-09-15 envelope replication over `SYNC_PUSH` and vault packs. Next: secondary indexes, session certificates, header-only envelopes. | diff --git a/planning/atomic-lib-runtime.md b/planning/atomic-lib-runtime.md index 30a76537e7..f291a37b92 100644 --- a/planning/atomic-lib-runtime.md +++ b/planning/atomic-lib-runtime.md @@ -1,6 +1,6 @@ # Atomic Lib Runtime: HTTP-Optional Local Node -> **Status:** Partial. `AtomicNode` in `lib/src/runtime/` is the binding runtime (`from_db`, `db`, the agent accessors, `query`, `apply_commit` under `IngestPolicy::{Hub, Peer, LocalCache}`); the WASM `ClientDb` is its only adapter, and the unused `open` / `get` / `mutate` / `subscribe` / `sync_with_peer` surface was cut back 2026-09-04. Open: binding the remaining adapters (#1277 / #1241). +> **Status:** Partial. `AtomicNode` in `lib/src/runtime/` is the binding runtime (`from_db`, `db`, the agent accessors, `query`, `apply_commit` under `IngestPolicy::{Hub, Peer, LocalCache}`); WASM `ClientDb` and Tauri's embedded native handle bind it, and the unused `open` / `get` / `mutate` / `subscribe` / `sync_with_peer` surface was cut back 2026-09-04. Open: binding the remaining adapters (#1277 / #1241). ## Status @@ -732,6 +732,37 @@ Tests: ### Phase 7: Tauri / Android Without Loopback +Active implementation (2026-09-10, `codex/tauri-http-optional`): + +- [x] Separate the existing embedding lifecycle from HTTP binding; keep the + managed-node `serve_with_hook` contract intact. +- [x] Bind Tauri's existing native operations to `AtomicNode`. +- [x] Prove native CRUD/query survives an unavailable HTTP port, and that + opting into HTTP still reports binding failures. +- [ ] Replace frontend local HTTP/WS operations with native commands/events, + including reads, commits, live updates, search and attachment bytes. +- [ ] Move remaining server-owned bootstrap/plugins into shared runtime + services, then make Actix an optional Tauri build dependency. +- [ ] Verify fresh install, sign-in/restore, drive switching, offline restart, + attachment access and peer sync in the desktop UI with no listener bound. + +Implemented boundary: `serve::run_node(config, adapter)` initializes the node +and keeps its services alive while the adapter future runs; +`serve::serve_http(appstate)` binds the optional HTTP/WS listener. +`serve_with_hook` remains source-compatible for hosted/managed nodes. Tauri +installs `appstate.node()` before starting HTTP and keeps the native lifecycle +alive if HTTP stops. Its Vault bridge now shares the same startup-error path +as other native commands. The periodic flush worker is owned and joined, with +a final flush when its lifecycle ends. + +The first extraction still uses `AppState` to bootstrap existing plugins and +Actix actors. It must not be described as removing Actix or as making the +current Tauri frontend work without HTTP. Do not expose a no-HTTP user setting +until that frontend acceptance gate passes. + +Related: #1196 (local-first SDK/API), #749 (origin independence), #1277 and +#1241 (native bindings), and the accepted runtime-boundary decision. + - Add Tauri commands for node get/query/mutate/blob operations. - Add an event stream from node events to the webview. - Point Tauri frontend at native node commands for local data. diff --git a/server/src/appstate.rs b/server/src/appstate.rs index c8078ef9fd..ba50202aef 100644 --- a/server/src/appstate.rs +++ b/server/src/appstate.rs @@ -46,6 +46,12 @@ pub struct AppState { } impl AppState { + /// Shared runtime boundary for native adapters. Clones share this store's + /// identity, event channels and durable data with the HTTP/WS adapter. + pub fn node(&self) -> atomic_lib::runtime::AtomicNode { + atomic_lib::runtime::AtomicNode::from_db(self.store.clone()) + } + /// Creates the AppState (the server's context available in Handlers). /// Initializes or opens a store on disk. /// Creates a new agent, if necessary. diff --git a/server/src/serve.rs b/server/src/serve.rs index 0b9f7270bf..ebc8d9f05d 100644 --- a/server/src/serve.rs +++ b/server/src/serve.rs @@ -201,6 +201,29 @@ pub async fn serve_with_hook( ) -> AtomicServerResult<()> where F: FnOnce(&crate::appstate::AppState), +{ + run_node(config, |appstate| async move { + on_ready(&appstate); + serve_http(appstate).await + }) + .await +} + +/// Run the embedded node without binding HTTP. The adapter owns the future's +/// lifetime: keep it pending while handling native commands. +/// Storage, plugins, indexing, flushes and peer transport are initialized once +/// and kept alive until `run` finishes. Call [`serve_http`] inside `run` only +/// when an HTTP/WS listener is wanted. +/// +/// Transitional boundary: bootstrap still uses `AppState` and Actix actors. +/// Native operations bind to `appstate.node()` (`AtomicNode`); moving the +/// remaining server-owned services into atomic_lib is a separate migration. +/// Must run inside an Actix system, like `serve_with_hook`. Peer transport still +/// has process-global state; this is not a restartable/multi-node host API. +pub async fn run_node(config: crate::config::Config, run: F) -> AtomicServerResult<()> +where + F: FnOnce(crate::appstate::AppState) -> Fut, + Fut: std::future::Future>, { println!( "Atomic-server {} \nUse --help for instructions. Visit https://docs.atomicdata.dev and https://github.com/atomicdata-dev/atomic-server for more info.", @@ -320,9 +343,12 @@ where }); } - // The durable-flush tick that makes Durability::None commits survive a - // crash is owned by `atomic_lib` (`Db::init_redb_file` spawns it), so the - // desktop and Flutter bindings get it without remembering to. + // Native adapters need durability even when no HTTP listener runs. + // Own the worker so returning/cancelling this lifecycle stops and joins it. + // (`Db::init_redb_file` also spawns its own best-effort tick so no binding + // that opens storage directly can forget durability; this worker adds the + // clean stop/join semantics native lifecycles need.) + let _flush_worker = DurableFlush::start(appstate.store.clone())?; // Start Iroh peer-to-peer transport let _iroh_router = { @@ -374,13 +400,25 @@ where #[cfg(feature = "wasm-plugins")] crate::plugins::sync_worker::spawn(appstate.clone()); - // Embedder hook: the store, indexes and transports are up, but the HTTP - // server hasn't started accepting connections yet. A managed-node wrapper - // (atomic-saas/managed-node) uses this to flip the `managed` flag, install - // its sync admission policy, and spawn its control-plane tasks. The open - // server passes a no-op (see `serve`), so it never phones home. - on_ready(&appstate); + // Embedder hook runs inside `run` (see `serve_with_hook`): the store, + // indexes and transports are up, but no HTTP listener is required to run + // it any more. A managed-node wrapper (atomic-saas/managed-node) uses it + // to flip the `managed` flag, install its sync admission policy, and + // spawn its control-plane tasks. The open server passes a no-op, so it + // never phones home. + let result = run(appstate).await; + tracing::info!("Node adapter stopped"); + if let Some(guard) = tracing_chrome_flush_guard { + guard.flush(); + } + result +} +/// Optional HTTP/WS adapter for an already initialized node. A bind failure +/// does not invalidate other clones of the native node handle. The caller +/// keeps `run_node` alive for as long as any adapter needs its services. +pub async fn serve_http(appstate: crate::appstate::AppState) -> AtomicServerResult<()> { + let config = appstate.config.clone(); let server = HttpServer::new(move || { let cors = Cors::permissive().expose_headers([SERVER_VERSION_HEADER]); @@ -503,17 +541,53 @@ where .await?; } - tracing::info!("Cleaning up"); - // Cleanup, runs when server is stopped - // Note that more cleanup code is in Appstate::exit - if let Some(guard) = tracing_chrome_flush_guard { - guard.flush() - } - tracing::info!("Server stopped"); Ok(()) } +/// Fsync runs off the async executor. Dropping the lifecycle wakes the worker +/// immediately, performs a final flush and joins it; the old detached loop +/// kept the database open forever after a failed bind or an embedder exit. +struct DurableFlush { + stop: Option>, + thread: Option>, +} + +impl DurableFlush { + fn start(store: atomic_lib::Db) -> std::io::Result { + let (stop, rx) = std::sync::mpsc::channel(); + let thread = std::thread::Builder::new() + .name("durable-flush".into()) + .spawn(move || loop { + let stopping = !matches!( + rx.recv_timeout(std::time::Duration::from_millis(100)), + Err(std::sync::mpsc::RecvTimeoutError::Timeout) + ); + if let Err(error) = store.flush() { + tracing::warn!("durable flush failed: {error}"); + } + if stopping { + break; + } + })?; + Ok(Self { + stop: Some(stop), + thread: Some(thread), + }) + } +} + +impl Drop for DurableFlush { + fn drop(&mut self) { + self.stop.take(); + if let Some(thread) = self.thread.take() { + if thread.join().is_err() { + tracing::error!("durable-flush thread panicked"); + } + } + } +} + /// Amount of seconds before server shuts down connections after SIGTERM signal const TIMEOUT: u64 = 15; @@ -553,3 +627,39 @@ const BANNER: &str = r#" / /_/ / /_/ /_/ / / / / / / / /__/_____(__ ) __/ / | |/ / __/ / \__,_/\__/\____/_/ /_/ /_/_/\___/ /____/\___/_/ |___/\___/_/ "#; + +#[cfg(test)] +mod lifecycle_tests { + use super::DurableFlush; + use atomic_lib::{urls, Resource, Storelike, Value}; + + #[actix_web::test] + async fn dropping_flush_worker_releases_database_and_persists_final_write() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("node.redb"); + let blobs = dir.path().join("blobs"); + let db = atomic_lib::Db::init_redb_file(&path, None, &blobs) + .await + .unwrap(); + let worker = DurableFlush::start(db.clone()).unwrap(); + let mut resource = Resource::new("did:ad:flush-test".into()); + resource + .set_unsafe(urls::NAME.into(), Value::String("Durable".into())) + .unwrap(); + db.add_resource_opts(&resource, false, false, true) + .await + .unwrap(); + drop(db); + // The worker owns the final Db reference. A detached loop would keep + // redb locked and reopening below would fail with DatabaseAlreadyOpen. + drop(worker); + let reopened = atomic_lib::Db::init_redb_file(&path, None, &blobs) + .await + .unwrap(); + let resource = reopened + .get_resource(&resource.get_subject()) + .await + .unwrap(); + assert_eq!(resource.get(urls::NAME).unwrap().to_string(), "Durable"); + } +} diff --git a/server/tests/it/http_optional.rs b/server/tests/it/http_optional.rs new file mode 100644 index 0000000000..ce7bcfd2c9 --- /dev/null +++ b/server/tests/it/http_optional.rs @@ -0,0 +1,104 @@ +//! The embedding lifecycle must not bind HTTP before handing control to a native adapter. +use atomic_lib::{agents::ForAgent, storelike::Query, urls, Resource, Storelike, Subject, Value}; +use atomic_server_lib::{config, serve}; + +#[actix_web::test] +async fn native_node_works_when_http_port_is_occupied() { + // Hold the configured port for the entire test: startup cannot secretly + // bind it, and opting into HTTP must fail without destroying local data. + let occupied = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let unique = format!("http_optional_{}", atomic_lib::utils::random_string(10)); + let mut config = config::build_temp_config(&unique).unwrap(); + config.opts.ip = std::net::Ipv4Addr::LOCALHOST.into(); + config.opts.port = occupied.local_addr().unwrap().port().into(); + + serve::run_node(config, |appstate| async move { + let node = appstate.node(); + let mut events = node.db().subscribe_events(); + let (_agent, drive) = node.db().setup("Native user").await?; + let drive = Subject::from(drive); + let mut document = Resource::new("did:ad:placeholder".into()); + document.set_unsafe(urls::NAME.into(), Value::String("No HTTP needed".into()))?; + document.set_unsafe(urls::PARENT.into(), Value::AtomicUrl(drive.clone()))?; + let created = document.save_as_genesis(node.db()).await?; + let subject = created.commit.subject; + + let found = node + .query(&Query { + property: Some(urls::PARENT.into()), + value: Some(Value::AtomicUrl(drive)), + for_agent: ForAgent::Sudo, + ..Query::new() + }) + .await?; + assert!(found.subjects.contains(&subject)); + + let error = serve::serve_http(appstate) + .await + .expect_err("the occupied HTTP port must fail"); + assert!(error.to_string().contains("Cannot bind"), "{error}"); + let stored = node.db().get_resource(&subject).await?; + assert_eq!(stored.get(urls::NAME)?.to_string(), "No HTTP needed"); + let mut edited = stored; + edited + .set( + urls::NAME.into(), + Value::String("Still writable".into()), + node.db(), + ) + .await?; + edited.save(node.db()).await?; + assert_eq!( + node.db() + .get_resource(&subject) + .await? + .get(urls::NAME)? + .to_string(), + "Still writable" + ); + let mut observed = false; + while let Ok(event) = events.try_recv() { + if let atomic_lib::db::DbEvent::Changed { + subject: changed, .. + } = event + { + observed |= changed.pure_id() == subject.pure_id(); + } + } + assert!( + observed, + "native subscribers must see local mutations without WebSocket" + ); + node.db().flush()?; + Ok(()) + }) + .await + .expect("native startup and CRUD must be independent of HTTP binding"); +} + +#[actix_web::test] +async fn hosted_embedder_hook_still_runs_before_http_binding() { + let occupied = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let unique = format!( + "http_optional_hook_{}", + atomic_lib::utils::random_string(10) + ); + let mut config = config::build_temp_config(&unique).unwrap(); + config.opts.ip = std::net::Ipv4Addr::LOCALHOST.into(); + config.opts.port = occupied.local_addr().unwrap().port().into(); + let called = std::cell::Cell::new(false); + let error = serve::serve_with_hook(config, |appstate| { + assert!(appstate.node().agent().is_some()); + appstate + .managed + .store(true, std::sync::atomic::Ordering::Relaxed); + called.set(true); + }) + .await + .expect_err("binding must fail after the embedder configures the node"); + assert!( + called.get(), + "managed embedders must get the ready hook before bind" + ); + assert!(error.to_string().contains("Cannot bind"), "{error}"); +} diff --git a/server/tests/it/main.rs b/server/tests/it/main.rs index 9b0723e088..3fd2878ee0 100644 --- a/server/tests/it/main.rs +++ b/server/tests/it/main.rs @@ -11,6 +11,7 @@ mod drive_presence; mod drive_presence_shared; mod file_search_repro; mod history_attribution; +mod http_optional; mod iroh_pairing; mod loro_ephemeral_sync; mod multi_client_sync; From 6eece74b9f57b4d9d18f0df632a6ab414badf81c Mon Sep 17 00:00:00 2001 From: Joep Meindertsma Date: Thu, 10 Sep 2026 20:14:50 +0200 Subject: [PATCH 2/3] Move native storage identity and durability into atomic_lib --- .dagger/src/index.ts | 3 + CHANGELOG.md | 1 + Cargo.lock | 2 +- TESTING_COVERAGE.md | 12 ++- desktop/README.md | 4 +- desktop/src/lib.rs | 2 +- desktop/src/system_tray.rs | 6 +- lib/Cargo.toml | 1 + lib/src/runtime/durable_flush.rs | 78 ++++++++++++++ lib/src/runtime/identity.rs | 77 ++++++++++++++ lib/src/runtime/mod.rs | 16 +-- lib/src/runtime/node.rs | 30 +++++- lib/tests/check-native-runtime.sh | 12 +++ lib/tests/native_runtime.rs | 162 ++++++++++++++++++++++++++++++ planning/atomic-lib-runtime.md | 26 ++++- server/src/appstate.rs | 114 ++------------------- server/src/serve.rs | 81 +-------------- 17 files changed, 423 insertions(+), 204 deletions(-) create mode 100644 lib/src/runtime/durable_flush.rs create mode 100644 lib/src/runtime/identity.rs create mode 100644 lib/tests/check-native-runtime.sh create mode 100644 lib/tests/native_runtime.rs diff --git a/.dagger/src/index.ts b/.dagger/src/index.ts index 2a2d40e1c7..78b91225ab 100644 --- a/.dagger/src/index.ts +++ b/.dagger/src/index.ts @@ -1505,6 +1505,9 @@ export class AtomicServer { `--test-threads ${this.hostKnobs.nextestTestThreads} ` + `--retries ${this.hostKnobs.nextestRetries}`, ]) + // Compile the native runtime separately: workspace feature unification + // must not hide a dependency on the hosted server/Actix adapter. + .withExec(['sh', 'lib/tests/check-native-runtime.sh']) .stdout() ); } diff --git a/CHANGELOG.md b/CHANGELOG.md index 0efdff551f..30c65c641f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,7 @@ See [STATUS.md](server/STATUS.md) to learn more about which features will remain ## UNRELEASED +- Move native storage opening, persisted identity loading and durable flushing into `atomic_lib::runtime`; reject damaged identity configs without replacing the key. Add a standalone no-Actix runtime CI check. - Separate embedded node startup from the optional HTTP adapter. Tauri keeps native services alive if HTTP stops; its frontend still requires HTTP/WS. This starts the HTTP-optional runtime migration ([#1196](https://github.com/ontola/atomic-server/issues/1196), [#749](https://github.com/ontola/atomic-server/issues/749)). - The outbox drains over a live Iroh link too (`sync::peer::LivePeerCommitTransport`): a device with no hub in reach delivers its queued writes to a paired peer as diff --git a/Cargo.lock b/Cargo.lock index 557994848c..0ee0350fc0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1283,7 +1283,6 @@ dependencies = [ "sha1 0.11.0", "simple-server-timing-header", "static-files 0.3.1", - "tempfile", "text-splitter", "tokio", "tokio-tungstenite 0.29.0", @@ -1401,6 +1400,7 @@ dependencies = [ "serde_json", "sha2 0.10.9", "sled", + "tempfile", "tokio", "tokio-tungstenite 0.29.0", "toml 1.1.2+spec-1.1.0", diff --git a/TESTING_COVERAGE.md b/TESTING_COVERAGE.md index 4eb7c74911..92f9de7412 100644 --- a/TESTING_COVERAGE.md +++ b/TESTING_COVERAGE.md @@ -314,13 +314,23 @@ verifies Clippy dispatch, staged input, and failure propagation; this fixture do `serve::run_node`, creates and queries a document through the native store/runtime, and explicitly attempts HTTP binding. After that bind fails, native reads, edits and change events still work. A second test keeps the hosted embedder -ready hook before HTTP binding. `serve::lifecycle_tests` verifies that +ready hook before HTTP binding. `runtime::durable_flush::tests` in atomic_lib +verifies that ending the owned flush worker releases redb and preserves its final write. These are Rust glue tests, not desktop UI acceptance. Tauri still starts the HTTP/WS adapter for its frontend. Missing: native frontend reads/commits/events, attachment bytes, restore and peer sync through Tauri with no HTTP listener; process-global Iroh teardown/restart is also not covered by this extraction. +`lib/tests/native_runtime.rs` runs on plain Tokio with only `db-redb,config`: +originless storage + identity + signed creation survive close/reopen, a retained +config recreates an identity in a replacement database, legacy secrets resolve +to the same key-derived DID, and malformed existing config is not overwritten. +`lib/tests/check-native-runtime.sh` rejects atomic-server/Actix in the normal +core dependency graph and compiles/runs those tests outside workspace feature +unification. `rustTest` runs this isolation gate after the workspace suite. +The gate is not a claim that the Tauri dependency graph is already server-free. + ## Browser WebRTC transport (issue #1396) `browser/lib/src/webrtc-transport.test.ts` covers frame fragmentation/order, diff --git a/desktop/README.md b/desktop/README.md index d3d3bbcb51..63b89dbed5 100644 --- a/desktop/README.md +++ b/desktop/README.md @@ -15,7 +15,9 @@ cargo tauri build ## Node and HTTP lifecycle The embedded node is initialized by `atomic_server_lib::serve::run_node`. -Tauri binds its native commands to the shared `AtomicNode`, then explicitly +Storage opening, identity loading and periodic durable flushing are provided +by `atomic_lib::runtime`, with no Actix dependency in that core. Tauri binds +its native commands to the shared `AtomicNode`, then explicitly starts `serve_http` for the current webview. A stopped HTTP adapter does not stop the native node's flush worker and peer tasks while the app is open. diff --git a/desktop/src/lib.rs b/desktop/src/lib.rs index 13e8cde137..10f63e82c8 100644 --- a/desktop/src/lib.rs +++ b/desktop/src/lib.rs @@ -736,7 +736,7 @@ pub fn run() { { let menu = crate::menu::build(app.handle())?; app.handle().set_menu(menu)?; - system_tray::setup(app, &config)?; + system_tray::setup(app, &config.get_origin(), &config.config_dir)?; } Ok(()) diff --git a/desktop/src/system_tray.rs b/desktop/src/system_tray.rs index b68d26b727..fd589b1051 100644 --- a/desktop/src/system_tray.rs +++ b/desktop/src/system_tray.rs @@ -12,7 +12,7 @@ use tauri_plugin_opener::OpenerExt; /// a second one. You get two identical icons in the tray, only this one having /// a menu (the config can't attach one), and both vanish together when the /// process dies — which reads as a rendering glitch rather than two real icons. -pub fn setup(app: &mut App, config: &atomic_server_lib::config::Config) -> tauri::Result<()> { +pub fn setup(app: &mut App, origin: &str, config_dir: &std::path::Path) -> tauri::Result<()> { let open = MenuItem::with_id(app, "open", "Open", true, None::<&str>)?; let browser = MenuItem::with_id(app, "browser", "Open in browser", true, None::<&str>)?; let config_item = MenuItem::with_id(app, "config", "Config folder", true, None::<&str>)?; @@ -22,8 +22,8 @@ pub fn setup(app: &mut App, config: &atomic_server_lib::config::Config) -> tauri let menu = Menu::with_items(app, &[&open, &browser, &config_item, &docs, &sep, &quit])?; - let origin = config.get_origin(); - let config_dir = config.config_dir.to_str().unwrap().to_string(); + let origin = origin.to_owned(); + let config_dir = config_dir.to_string_lossy().into_owned(); TrayIconBuilder::new() .icon(app.default_window_icon().unwrap().clone()) diff --git a/lib/Cargo.toml b/lib/Cargo.toml index 06985b2cc0..f9b8b10d00 100644 --- a/lib/Cargo.toml +++ b/lib/Cargo.toml @@ -73,6 +73,7 @@ wasm-bindgen-futures = { version = "0.4.72", optional = true } web-sys = { version = "0.3.99", optional = true, features = ["DomException", "FileSystemDirectoryHandle", "FileSystemFileHandle", "FileSystemGetFileOptions", "FileSystemSyncAccessHandle", "FileSystemReadWriteOptions", "StorageManager", "WorkerGlobalScope", "WorkerNavigator"] } [dev-dependencies] +tempfile = "3" criterion = { version = "0.8.2", features = ["async_tokio"] } iai = "0.1" lazy_static = "1" diff --git a/lib/src/runtime/durable_flush.rs b/lib/src/runtime/durable_flush.rs new file mode 100644 index 0000000000..f805cdda69 --- /dev/null +++ b/lib/src/runtime/durable_flush.rs @@ -0,0 +1,78 @@ +//! Owned native durability worker, independent of HTTP and async executors. + +/// Fsync runs off the async executor. Dropping the lifecycle wakes the worker +/// immediately, performs a final flush and joins it; the old detached loop +/// kept the database open forever after a failed bind or an embedder exit. +#[must_use = "Keep the flush guard alive for the lifetime of the node"] +pub struct DurableFlush { + stop: Option>, + thread: Option>, +} + +impl DurableFlush { + pub(crate) fn start(store: crate::Db) -> std::io::Result { + let (stop, rx) = std::sync::mpsc::channel(); + let thread = std::thread::Builder::new() + .name("durable-flush".into()) + .spawn(move || loop { + let stopping = !matches!( + rx.recv_timeout(std::time::Duration::from_millis(100)), + Err(std::sync::mpsc::RecvTimeoutError::Timeout) + ); + if let Err(error) = store.flush() { + tracing::warn!("durable flush failed: {error}"); + } + if stopping { + break; + } + })?; + Ok(Self { + stop: Some(stop), + thread: Some(thread), + }) + } +} + +impl Drop for DurableFlush { + fn drop(&mut self) { + self.stop.take(); + if let Some(thread) = self.thread.take() { + if thread.join().is_err() { + tracing::error!("durable-flush thread panicked"); + } + } + } +} + +#[cfg(all(test, feature = "db-redb"))] +mod tests { + use super::DurableFlush; + use crate::{urls, Resource, Storelike, Value}; + + #[tokio::test] + async fn dropping_flush_worker_releases_database_and_persists_final_write() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("node.redb"); + let blobs = dir.path().join("blobs"); + let db = crate::Db::init_redb_file(&path, None, &blobs) + .await + .unwrap(); + let worker = DurableFlush::start(db.clone()).unwrap(); + let mut resource = Resource::new("did:ad:flush-test".into()); + resource + .set_unsafe(urls::NAME.into(), Value::String("Durable".into())) + .unwrap(); + db.add_resource_opts(&resource, false, false, true) + .await + .unwrap(); + drop(db); + // The worker owns the final Db reference. A detached loop would keep + // redb locked and reopening below would fail with DatabaseAlreadyOpen. + drop(worker); + let reopened = crate::Db::init_redb_file(&path, None, &blobs) + .await + .unwrap(); + let resource = reopened.get_resource(resource.get_subject()).await.unwrap(); + assert_eq!(resource.get(urls::NAME).unwrap().to_string(), "Durable"); + } +} diff --git a/lib/src/runtime/identity.rs b/lib/src/runtime/identity.rs new file mode 100644 index 0000000000..0607d987c4 --- /dev/null +++ b/lib/src/runtime/identity.rs @@ -0,0 +1,77 @@ +//! Shared native identity bootstrap for adapters; no server configuration or HTTP. +use super::AtomicNode; +use crate::{agents::Agent, config::SharedConfig, errors::AtomicResult, Storelike}; + +impl AtomicNode { + /// Load this node's persisted identity, including legacy subject migration. + /// Only a missing config creates a new identity. A damaged or unreadable + /// existing config fails without overwriting it. Does not create a drive. + pub async fn load_or_create_agent( + &self, + config_path: &std::path::Path, + name: &str, + ) -> AtomicResult<()> { + let store = self.db(); + tracing::info!("Setting default agent"); + + let agent = if config_path.try_exists()? { + // An existing but damaged/unreadable file is not a fresh install. + // Never replace its identity merely because parsing/loading failed. + let agent_config = crate::config::read_config(Some(config_path))?; + // Agent::from_secret owns legacy URL-to-DID migration for every + // adapter; do not duplicate that conversion in runtime bootstrap. + let agent = Agent::from_secret(&agent_config.shared.agent_secret)?; + + match store.get_resource(&agent.subject.clone()).await { + Ok(_) => agent, + Err(e) => { + if agent.subject.is_local() { + // If there is an agent in the config, but not in the store, + // That probably means that the DB has been erased and only the config file exists. + // This means that the Agent from the Config file should be recreated, using its private key. + tracing::info!("Agent not retrievable, but config was found. Recreating Agent in new store."); + + let mut recreated_agent = Agent::new_from_private_key( + Some(name), + &agent.private_key.ok_or("No private key found")?, + )?; + recreated_agent.initial_drive = agent.initial_drive; + store.add_resource(&recreated_agent.to_resource()?).await?; + + recreated_agent + } else { + return Err(format!( + "An agent is present in {:?}, but this agent cannot be retrieved. Either make sure the agent is retrievable, or remove it from your config. {}", + config_path, e, + ).into()); + } + } + } + } else { + let agent = store.create_agent(Some(name)).await?; + let cfg = crate::config::Config { + shared: SharedConfig { + agent_secret: agent.build_secret()?, + initial_drive: agent.initial_drive.clone().map(|s| s.to_string()), + }, + client: None, + }; + + cfg.save(config_path)?; + + // Never log the agent secret: on Android it would land in logcat + // (bug reports, `adb logcat` history), and on servers in log + // aggregators. The secret lives only in the config file. + tracing::warn!( + "No existing config found, created a new Config at {:?}. To sign in from another device, use the pairing/sign-in flow in the browser app, or copy the agent secret from that file.", + config_path + ); + + agent + }; + + tracing::info!("Default Agent is set: {}", &agent.subject); + store.set_default_agent(agent); + Ok(()) + } +} diff --git a/lib/src/runtime/mod.rs b/lib/src/runtime/mod.rs index c72c3fe7d5..05bad3f0ee 100644 --- a/lib/src/runtime/mod.rs +++ b/lib/src/runtime/mod.rs @@ -2,13 +2,17 @@ //! (HTTP, WebSocket, Iroh, WASM, FFI, Flutter) bind to instead of wrapping //! [`crate::Db`] themselves. //! -//! Slice 1 (`planning/atomic-lib-runtime.md`, `planning/completed/runtime-boundary-decision.md`) -//! is a thin wrapper: every method delegates to code that already existed, so -//! there is no behaviour change — only a named place for it. The surface is -//! deliberately no wider than what binds to it today (the WASM `ClientDb`): -//! `from_db`, `db`, the agent accessors, `query` and `apply_commit` under an -//! [`IngestPolicy`]. +//! WASM binds existing storage through `from_db`. Native adapters can open +//! local storage, load a persisted identity and own a durable-flush worker +//! here without a server crate, HTTP origin or Actix executor. mod node; pub use node::{AtomicNode, IngestPolicy}; + +#[cfg(not(target_arch = "wasm32"))] +mod durable_flush; +#[cfg(not(target_arch = "wasm32"))] +pub use durable_flush::DurableFlush; +#[cfg(all(feature = "config", not(target_arch = "wasm32")))] +mod identity; diff --git a/lib/src/runtime/node.rs b/lib/src/runtime/node.rs index d0c6775d94..d3fbc94957 100644 --- a/lib/src/runtime/node.rs +++ b/lib/src/runtime/node.rs @@ -10,8 +10,8 @@ //! with a storage config, `get`, `mutate`, `subscribe`, `sync_with_peer`. //! Three days later nothing but the WASM binding had bound to it, and the //! WASM binding used four methods. The surface was cut down to those on -//! 2026-09-04; the seam stays, and grows again when a second adapter binds -//! to it (`planning/atomic-lib-runtime.md`). +//! 2026-09-04. Native startup and durability were added in September 2026 +//! when the server/Tauri adapter began consuming them. use crate::{ agents::Agent, @@ -63,8 +63,32 @@ pub struct AtomicNode { } impl AtomicNode { + /// Open durable native storage and bootstrap the bundled models/search + /// index. An origin is only needed for legacy URL resources or a hosted + /// adapter; native DID-only nodes pass `None`. Does not start listeners, + /// create an identity, or contact peers. + #[cfg(all(feature = "db-redb", not(target_arch = "wasm32")))] + pub async fn open_local( + data_path: &std::path::Path, + blobs_path: &std::path::Path, + origin: Option, + ) -> AtomicResult { + Ok(Self::from_db( + Db::init_redb_file(data_path, origin, blobs_path).await?, + )) + } + + /// Start the shared 100ms fsync worker. Keep the guard alive while this + /// node is running; dropping it joins the worker after a final flush. + /// Start one guard per store, not one per adapter or node clone. + #[cfg(not(target_arch = "wasm32"))] + pub fn start_durable_flush(&self) -> AtomicResult { + Ok(super::DurableFlush::start(self.db.clone())?) + } + /// Wrap an already-opened store. Adapters construct `Db` themselves - /// (the WASM `ClientDb`, the server's `AppState`) and bind here. + /// (such as the WASM `ClientDb`) and bind here; native hosts can use + /// `open_local` instead. pub fn from_db(db: Db) -> Self { Self { db } } diff --git a/lib/tests/check-native-runtime.sh b/lib/tests/check-native-runtime.sh new file mode 100644 index 0000000000..2cb4685f94 --- /dev/null +++ b/lib/tests/check-native-runtime.sh @@ -0,0 +1,12 @@ +#!/bin/sh +# Run from the workspace root. Isolate core features from server feature unification. +set -eu +runtime_deps=$(mktemp) +trap 'rm -f "$runtime_deps"' EXIT HUP INT TERM +cargo tree --locked -p atomic_lib --no-default-features --features db-redb,config \ + --edges normal --prefix none --format '{p}' > "$runtime_deps" +if grep -E '^(atomic-server|actix(-[^ ]+)?) ' "$runtime_deps"; then + echo 'Native runtime must not depend on atomic-server or Actix.' >&2 + exit 1 +fi +cargo test --locked -p atomic_lib --no-default-features --features db-redb,config --test native_runtime diff --git a/lib/tests/native_runtime.rs b/lib/tests/native_runtime.rs new file mode 100644 index 0000000000..b7fac09a5e --- /dev/null +++ b/lib/tests/native_runtime.rs @@ -0,0 +1,162 @@ +//! Core startup must work without a server crate, HTTP origin, or Actix runtime. +#![cfg(all(feature = "db-redb", feature = "config", not(target_arch = "wasm32")))] + +use atomic_lib::{runtime::AtomicNode, urls, Resource, Storelike, Value}; + +#[tokio::test] +async fn originless_node_preserves_identity_and_data_across_restart() { + let dir = tempfile::tempdir().unwrap(); + let data = dir.path().join("data"); + let blobs = dir.path().join("blobs"); + let identity = dir.path().join("config.toml"); + let node = AtomicNode::open_local(&data, &blobs, None).await.unwrap(); + node.load_or_create_agent(&identity, "Native user") + .await + .unwrap(); + let agent = node.agent().unwrap().subject; + assert_eq!(node.db().get_base_domain(), None); + let flush = node.start_durable_flush().unwrap(); + let drive = node.db().create_drive("Local drive").await.unwrap(); + let mut draft = Resource::new("did:ad:placeholder".into()); + draft + .set_unsafe(urls::NAME.into(), Value::String("Without a server".into())) + .unwrap(); + draft + .set_unsafe(urls::PARENT.into(), Value::AtomicUrl(drive.into())) + .unwrap(); + let saved = draft.save_as_genesis(node.db()).await.unwrap(); + let subject = saved.commit.subject; + drop(node); + drop(flush); + + let reopened = AtomicNode::open_local(&data, &blobs, None).await.unwrap(); + reopened + .load_or_create_agent(&identity, "Another label") + .await + .unwrap(); + assert_eq!(reopened.agent().unwrap().subject, agent); + assert_eq!( + reopened + .db() + .get_resource(&subject) + .await + .unwrap() + .get(urls::NAME) + .unwrap() + .to_string(), + "Without a server" + ); +} + +#[tokio::test] +async fn invalid_identity_config_is_not_replaced_with_a_new_identity() { + let dir = tempfile::tempdir().unwrap(); + let identity = dir.path().join("config.toml"); + let invalid = "a damaged existing identity file"; + std::fs::write(&identity, invalid).unwrap(); + let node = AtomicNode::open_local(&dir.path().join("data"), &dir.path().join("blobs"), None) + .await + .unwrap(); + assert!(node + .load_or_create_agent(&identity, "Native user") + .await + .is_err()); + assert!(node.agent().is_none()); + assert_eq!(std::fs::read_to_string(identity).unwrap(), invalid); +} + +#[tokio::test] +async fn existing_identity_can_bootstrap_a_replacement_database() { + let dir = tempfile::tempdir().unwrap(); + let identity = dir.path().join("config.toml"); + let first = AtomicNode::open_local(&dir.path().join("first"), &dir.path().join("blobs"), None) + .await + .unwrap(); + first + .load_or_create_agent(&identity, "Original") + .await + .unwrap(); + let mut original = first.agent().unwrap(); + original.initial_drive = Some( + first + .db() + .create_drive("Remembered drive") + .await + .unwrap() + .into(), + ); + let mut persisted = atomic_lib::config::read_config(Some(&identity)).unwrap(); + persisted.shared.agent_secret = original.build_secret().unwrap(); + persisted.shared.initial_drive = original.initial_drive.as_ref().map(ToString::to_string); + persisted.save(&identity).unwrap(); + let saved_config = std::fs::read(&identity).unwrap(); + let replacement = AtomicNode::open_local( + &dir.path().join("replacement"), + &dir.path().join("blobs"), + None, + ) + .await + .unwrap(); + replacement + .load_or_create_agent(&identity, "Restored") + .await + .unwrap(); + let restored = replacement.agent().unwrap(); + assert_eq!(restored.subject, original.subject); + assert_eq!(restored.public_key, original.public_key); + assert_eq!(restored.initial_drive, original.initial_drive); + assert!(replacement + .db() + .get_resource(&restored.subject) + .await + .is_ok()); + assert_eq!(std::fs::read(identity).unwrap(), saved_config); +} + +#[tokio::test] +async fn legacy_identity_migrates_to_did_without_changing_key_or_client_config() { + use atomic_lib::{ + agents::Agent, + config::{ClientConfig, Config, SharedConfig}, + }; + let dir = tempfile::tempdir().unwrap(); + let identity = dir.path().join("config.toml"); + let original = Agent::new(Some("Legacy")).unwrap(); + // Write the actual historical wire form, rather than serializing an + // Agent that may already normalize the subject during construction. + let legacy_secret = atomic_lib::agents::encode_base64( + &serde_json::to_vec(&serde_json::json!({ + "privateKey": original.private_key.as_ref().unwrap(), + "subject": format!("https://atomicdata.dev/agents/{}", original.public_key), + })) + .unwrap(), + ); + Config { + shared: SharedConfig { + agent_secret: legacy_secret, + initial_drive: None, + }, + client: Some(ClientConfig { + server_url: "https://example.test".into(), + }), + } + .save(&identity) + .unwrap(); + let node = AtomicNode::open_local(&dir.path().join("data"), &dir.path().join("blobs"), None) + .await + .unwrap(); + node.load_or_create_agent(&identity, "Migrated") + .await + .unwrap(); + let migrated = node.agent().unwrap(); + assert_eq!(migrated.subject, original.subject); + assert_eq!(migrated.public_key, original.public_key); + let config = atomic_lib::config::read_config(Some(&identity)).unwrap(); + assert_eq!(config.client.unwrap().server_url, "https://example.test"); + assert_eq!( + Agent::from_secret(&config.shared.agent_secret) + .unwrap() + .subject, + original.subject + ); +} diff --git a/planning/atomic-lib-runtime.md b/planning/atomic-lib-runtime.md index f291a37b92..e33900cad4 100644 --- a/planning/atomic-lib-runtime.md +++ b/planning/atomic-lib-runtime.md @@ -741,6 +741,9 @@ Active implementation (2026-09-10, `codex/tauri-http-optional`): opting into HTTP still reports binding failures. - [ ] Replace frontend local HTTP/WS operations with native commands/events, including reads, commits, live updates, search and attachment bytes. +- [x] Move native storage opening, identity loading and durable flushing into + `atomic_lib::runtime`; exercise them without HTTP or Actix under standalone + `db-redb,config` features and enforce the core dependency boundary in CI. - [ ] Move remaining server-owned bootstrap/plugins into shared runtime services, then make Actix an optional Tauri build dependency. - [ ] Verify fresh install, sign-in/restore, drive switching, offline restart, @@ -755,7 +758,28 @@ alive if HTTP stops. Its Vault bridge now shares the same startup-error path as other native commands. The periodic flush worker is owned and joined, with a final flush when its lifecycle ends. -The first extraction still uses `AppState` to bootstrap existing plugins and +End state: the default Tauri dependency graph must contain neither +`atomic-server` nor Actix. HTTP hosting can be a separate opt-in adapter; +keeping it mandatory behind an environment flag does not meet this gate. + +Remaining direct coupling: + +- `desktop/src/lib.rs` still uses server CLI/config and `run_node`/`serve_http`. +- `server::AppState` still registers feature plugins and Actix commit/presence + actors. Native startup must own equivalent shared services, not duplicate + their logic in Tauri. +- The frontend still reads, commits, subscribes and loads attachment bytes + through the local HTTP/WS adapter. +- Iroh is already in atomic_lib, but its global endpoint/router lifetime must + be owned independently before native shutdown/restart is promised. + +`AtomicNode::open_local`, `load_or_create_agent` and `start_durable_flush` now +provide the server-free storage/identity/durability path. The hosted adapter +uses those same operations. Identity loading fails on damaged existing config +rather than silently replacing the key. The tray accepts ordinary values and +no longer names the server configuration type. + +The current embedding still uses `AppState` to bootstrap existing plugins and Actix actors. It must not be described as removing Actix or as making the current Tauri frontend work without HTTP. Do not expose a no-HTTP user setting until that frontend acceptance gate passes. diff --git a/server/src/appstate.rs b/server/src/appstate.rs index ba50202aef..182c59a4e3 100644 --- a/server/src/appstate.rs +++ b/server/src/appstate.rs @@ -5,7 +5,7 @@ use crate::{ commit_monitor::CommitMonitor, config::Config, errors::AtomicServerResult, handlers::web_sockets::IndexStatusBroadcast, plugins, }; -use atomic_lib::{agents::Agent, commit::CommitResponse, config::SharedConfig, Storelike}; +use atomic_lib::commit::CommitResponse; #[cfg(feature = "wasm-plugins")] use crate::plugins::wasm; @@ -66,12 +66,13 @@ impl AppState { tracing::warn!("Development mode is enabled. This will use staging environments for services like LetsEncrypt."); } - let mut store = atomic_lib::Db::init_redb_file( + let node = atomic_lib::runtime::AtomicNode::open_local( &config.store_path, - Some(config.get_origin()), &config.uploads_path, + Some(config.get_origin()), ) .await?; + let mut store = node.db().clone(); crate::blob_storage::configure(&mut store).await?; // Before anything reads or writes a secret, so nothing is stored in @@ -163,7 +164,9 @@ impl AppState { } } - set_default_agent(&config, &store).await?; + atomic_lib::runtime::AtomicNode::from_db(store.clone()) + .load_or_create_agent(&config.config_file_path, "server") + .await?; let should_init = !&config.store_path.exists() || config.initialize; // If the store is empty, populate the core models (classes, properties, etc.). @@ -288,106 +291,3 @@ impl Drop for AppState { } } } - -/// Create a new agent if it does not yet exist. -async fn set_default_agent(config: &Config, store: &impl Storelike) -> AtomicServerResult<()> { - tracing::info!("Setting default agent"); - - let agent = match atomic_lib::config::read_config(Some(&config.config_file_path)) { - Ok(agent_config) => { - let mut agent = Agent::from_secret(&agent_config.shared.agent_secret)?; - - // Migrate old-format agent subjects (e.g. "https://atomicdata.dev/agents/...") - // to the new "did:ad:" format. Old configs stored the agent subject as an - // HTTP URL on atomicdata.dev, but the agent's keys are local. During invite token - // verification the old URL would resolve to an external resource with a different - // public key, causing "Invalid signature" errors. - let needs_migration = agent - .subject - .as_str() - .starts_with("https://atomicdata.dev/agents/") - || agent - .subject - .as_str() - .starts_with("http://atomicdata.dev/agents/"); - if needs_migration { - let private_key = agent - .private_key - .clone() - .ok_or("No private key found on agent to migrate")?; - let migrated = Agent::new_from_private_key(Some("server"), &private_key)?; - tracing::info!( - "Migrating agent subject from old format '{}' to new format '{}'", - agent.subject, - migrated.subject - ); - agent = migrated; - - // Update the config file so the migration only happens once - let cfg = atomic_lib::config::Config { - shared: SharedConfig { - agent_secret: agent.build_secret()?, - initial_drive: agent.initial_drive.clone().map(|s| s.to_string()), - }, - client: agent_config.client, - }; - cfg.save(&config.config_file_path)?; - tracing::info!( - "Config file updated with migrated agent at {:?}", - config.config_file_path - ); - } - - match store.get_resource(&agent.subject.clone()).await { - Ok(_) => agent, - Err(e) => { - if agent.subject.is_local() { - // If there is an agent in the config, but not in the store, - // That probably means that the DB has been erased and only the config file exists. - // This means that the Agent from the Config file should be recreated, using its private key. - tracing::info!("Agent not retrievable, but config was found. Recreating Agent in new store."); - - let recreated_agent = Agent::new_from_private_key( - "server".into(), - &agent.private_key.ok_or("No private key found")?, - )?; - store.add_resource(&recreated_agent.to_resource()?).await?; - - recreated_agent - } else { - return Err(format!( - "An agent is present in {:?}, but this agent cannot be retrieved. Either make sure the agent is retrievable, or remove it from your config. {}", - config.config_file_path, e, - ).into()); - } - } - } - } - Err(_no_config) => { - let agent = store.create_agent(Some("server")).await?; - let cfg = atomic_lib::config::Config { - shared: SharedConfig { - agent_secret: agent.build_secret()?, - initial_drive: agent.initial_drive.clone().map(|s| s.to_string()), - }, - client: None, - }; - - cfg.save(&config.config_file_path)?; - - // Never log the agent secret: on Android it would land in logcat - // (bug reports, `adb logcat` history), and on servers in log - // aggregators. The secret lives only in the config file. - tracing::warn!( - "No existing config found, created a new Config at {:?}. To sign in from another device, use the pairing/sign-in flow in the browser app, or copy the agent secret from that file.", - &config.config_file_path - ); - - agent - } - }; - - tracing::info!("Default Agent is set: {}", &agent.subject); - store.set_default_agent(agent); - Ok(()) -} diff --git a/server/src/serve.rs b/server/src/serve.rs index ebc8d9f05d..427d690f2d 100644 --- a/server/src/serve.rs +++ b/server/src/serve.rs @@ -348,7 +348,7 @@ where // (`Db::init_redb_file` also spawns its own best-effort tick so no binding // that opens storage directly can forget durability; this worker adds the // clean stop/join semantics native lifecycles need.) - let _flush_worker = DurableFlush::start(appstate.store.clone())?; + let _flush_worker = appstate.node().start_durable_flush()?; // Start Iroh peer-to-peer transport let _iroh_router = { @@ -545,49 +545,6 @@ pub async fn serve_http(appstate: crate::appstate::AppState) -> AtomicServerResu Ok(()) } -/// Fsync runs off the async executor. Dropping the lifecycle wakes the worker -/// immediately, performs a final flush and joins it; the old detached loop -/// kept the database open forever after a failed bind or an embedder exit. -struct DurableFlush { - stop: Option>, - thread: Option>, -} - -impl DurableFlush { - fn start(store: atomic_lib::Db) -> std::io::Result { - let (stop, rx) = std::sync::mpsc::channel(); - let thread = std::thread::Builder::new() - .name("durable-flush".into()) - .spawn(move || loop { - let stopping = !matches!( - rx.recv_timeout(std::time::Duration::from_millis(100)), - Err(std::sync::mpsc::RecvTimeoutError::Timeout) - ); - if let Err(error) = store.flush() { - tracing::warn!("durable flush failed: {error}"); - } - if stopping { - break; - } - })?; - Ok(Self { - stop: Some(stop), - thread: Some(thread), - }) - } -} - -impl Drop for DurableFlush { - fn drop(&mut self) { - self.stop.take(); - if let Some(thread) = self.thread.take() { - if thread.join().is_err() { - tracing::error!("durable-flush thread panicked"); - } - } - } -} - /// Amount of seconds before server shuts down connections after SIGTERM signal const TIMEOUT: u64 = 15; @@ -627,39 +584,3 @@ const BANNER: &str = r#" / /_/ / /_/ /_/ / / / / / / / /__/_____(__ ) __/ / | |/ / __/ / \__,_/\__/\____/_/ /_/ /_/_/\___/ /____/\___/_/ |___/\___/_/ "#; - -#[cfg(test)] -mod lifecycle_tests { - use super::DurableFlush; - use atomic_lib::{urls, Resource, Storelike, Value}; - - #[actix_web::test] - async fn dropping_flush_worker_releases_database_and_persists_final_write() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("node.redb"); - let blobs = dir.path().join("blobs"); - let db = atomic_lib::Db::init_redb_file(&path, None, &blobs) - .await - .unwrap(); - let worker = DurableFlush::start(db.clone()).unwrap(); - let mut resource = Resource::new("did:ad:flush-test".into()); - resource - .set_unsafe(urls::NAME.into(), Value::String("Durable".into())) - .unwrap(); - db.add_resource_opts(&resource, false, false, true) - .await - .unwrap(); - drop(db); - // The worker owns the final Db reference. A detached loop would keep - // redb locked and reopening below would fail with DatabaseAlreadyOpen. - drop(worker); - let reopened = atomic_lib::Db::init_redb_file(&path, None, &blobs) - .await - .unwrap(); - let resource = reopened - .get_resource(&resource.get_subject()) - .await - .unwrap(); - assert_eq!(resource.get(urls::NAME).unwrap().to_string(), "Durable"); - } -} From 2ff25390bdce4f3d740cc9883d625d9e98f1cddf Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 18 Sep 2026 05:32:49 +0000 Subject: [PATCH 3/3] test: raise timeout for the collaborative-Markdown-schema conversion test First test in the file to load @tiptap/markdown and the full collaborative editor schema graph; the cold module transform can exceed vitest's default 5s under CI load. Matches the existing precedent in vitest.integration.config.ts for the same class of timing issue. --- .../views/File/convertFileToDocument.test.ts | 60 +++++++++++-------- 1 file changed, 34 insertions(+), 26 deletions(-) diff --git a/browser/data-browser/src/views/File/convertFileToDocument.test.ts b/browser/data-browser/src/views/File/convertFileToDocument.test.ts index 7a1b7a4b7c..53504798d6 100644 --- a/browser/data-browser/src/views/File/convertFileToDocument.test.ts +++ b/browser/data-browser/src/views/File/convertFileToDocument.test.ts @@ -111,32 +111,40 @@ describe('plainTextToTiptapJson', () => { }); describe('fileContentsToTiptapJson', () => { - it('uses the collaborative Markdown schema so Markdown formatting becomes document nodes', async () => { - const json = await fileContentsToTiptapJson( - '# Heading\n\nThis is **bold**.', - 'markdown', - new Store(), - ); - - expect(json).toMatchObject({ - type: 'doc', - content: [ - { - type: 'heading', - attrs: { level: 1 }, - content: [{ text: 'Heading' }], - }, - { - type: 'paragraph', - content: [ - { text: 'This is ' }, - { text: 'bold', marks: [{ type: 'bold' }] }, - { text: '.' }, - ], - }, - ], - }); - }); + it( + 'uses the collaborative Markdown schema so Markdown formatting becomes document nodes', + async () => { + const json = await fileContentsToTiptapJson( + '# Heading\n\nThis is **bold**.', + 'markdown', + new Store(), + ); + + expect(json).toMatchObject({ + type: 'doc', + content: [ + { + type: 'heading', + attrs: { level: 1 }, + content: [{ text: 'Heading' }], + }, + { + type: 'paragraph', + content: [ + { text: 'This is ' }, + { text: 'bold', marks: [{ type: 'bold' }] }, + { text: '.' }, + ], + }, + ], + }); + }, + // First test in this file to load @tiptap/markdown and the full + // collaborative editor schema graph; that cold module transform can + // exceed the default 5s under CI load (see vitest.integration.config.ts + // for the same reasoning applied elsewhere). + 15000, + ); it('keeps blank Markdown and an unpaired marker as document text', async () => { await expect(