From b3621106c7ae83ad4fe68832cc2c463c27b0e369 Mon Sep 17 00:00:00 2001 From: huvvi Date: Sat, 3 Oct 2026 14:07:18 +0300 Subject: [PATCH 1/2] refactor(aimdb-sync): remove redundant try-set API --- aimdb-core/src/typed_api.rs | 12 ++-- aimdb-sync/CHANGELOG.md | 10 ++++ aimdb-sync/Cargo.toml | 2 +- aimdb-sync/src/error.rs | 7 +-- aimdb-sync/src/producer.rs | 76 ++++-------------------- aimdb-sync/tests/fork_safety_test.rs | 8 +-- aimdb-sync/tests/integration_test.rs | 20 +++---- aimdb-sync/tests/settable_integration.rs | 33 ---------- examples/sync-api-demo/src/main.rs | 2 +- 9 files changed, 41 insertions(+), 129 deletions(-) diff --git a/aimdb-core/src/typed_api.rs b/aimdb-core/src/typed_api.rs index e677ca3a..d36e50d9 100644 --- a/aimdb-core/src/typed_api.rs +++ b/aimdb-core/src/typed_api.rs @@ -160,13 +160,13 @@ where self.write.push(value); } - /// Non-blocking push. Returns the value back via [`TryProduceError::Full`] - /// if a bounded buffer is at capacity, or [`TryProduceError::Closed`] if - /// the record is shutting down. Use when the caller has a meaningful - /// response to backpressure. + /// Extension point for write handles that may reject a value. /// - /// Overwriting buffers (`SpmcRing`, `SingleLatest`, `Mailbox`) always - /// return `Ok(())`. Use [`produce`](Self::produce) for those. + /// The built-in buffers (`SpmcRing`, `SingleLatest`, `Mailbox`) overwrite + /// and currently always return `Ok(())`; use [`produce`](Self::produce) + /// for them. A future bounded, non-overwriting write handle can return the + /// value through [`TryProduceError::Full`], while a close-aware handle can + /// return it through [`TryProduceError::Closed`]. /// /// # Example /// diff --git a/aimdb-sync/CHANGELOG.md b/aimdb-sync/CHANGELOG.md index da307361..8e564353 100644 --- a/aimdb-sync/CHANGELOG.md +++ b/aimdb-sync/CHANGELOG.md @@ -7,6 +7,13 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Changed (breaking) + +- **Issue #212:** Removed `SyncProducer::try_set`, + `SyncProducer::try_set_value`, and `SyncError::SetTimeout`. Every built-in + buffer overwrites rather than refusing a write, so these APIs duplicated + `set()` and `set_value()` without providing distinct behavior. + ## [0.6.0] - 2026-09-18 ### Added @@ -159,6 +166,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Capacity-related API removed: `AimDbBuilderSyncExt::producer_with_capacity`/`consumer_with_capacity` and the `DEFAULT_SYNC_CHANNEL_CAPACITY` constant are gone. - `SyncConsumer`: `get`, `try_get`, `get_with_timeout`, `get_latest`, and `get_latest_with_timeout` now take `&mut self` (was `&self`) - `SyncConsumer` no longer implements `Clone` or `Sync` (still `Send`). + - `BufferLagged` now reaches callers of `get`, `get_with_timeout`, and + `try_get` as `SyncError::Db(DbError::BufferLagged { .. })`; the old + forwarding task swallowed it. - **Migration.** The replacement for a cloned consumer is a second `handle.consumer()` — but note that it is not a like-for-like swap. Cloning in 0.5.0 shared one stream, so each clone got *a share* of the diff --git a/aimdb-sync/Cargo.toml b/aimdb-sync/Cargo.toml index 37d1a609..fb36c38c 100644 --- a/aimdb-sync/Cargo.toml +++ b/aimdb-sync/Cargo.toml @@ -54,7 +54,7 @@ std = ["aimdb-core/std", "dep:tokio", "dep:aimdb-tokio-adapter", "dep:libc"] tracing = ["aimdb-core/tracing"] log = ["aimdb-core/log"] -# SyncProducer::set_value / try_set_value / set_value_at for Settable types +# SyncProducer::set_value / set_value_at for Settable types data-contracts = ["dep:aimdb-data-contracts"] # Fork detection only, and only where `pthread_atfork` exists. Optional and diff --git a/aimdb-sync/src/error.rs b/aimdb-sync/src/error.rs index ee3b2a8f..8310e3ff 100644 --- a/aimdb-sync/src/error.rs +++ b/aimdb-sync/src/error.rs @@ -33,10 +33,6 @@ pub enum SyncError { message: String, }, - /// Timeout while setting a value. - #[error("Timeout while setting value")] - SetTimeout, - /// Timeout while getting a value. #[error("Timeout while getting value")] GetTimeout, @@ -73,7 +69,7 @@ impl SyncError { // runtime thread that panicked. Neither has a `DbError` behind it. Self::DetachFailed { .. } => DbErrorKind::Internal, - Self::SetTimeout | Self::GetTimeout => DbErrorKind::Retry, + Self::GetTimeout => DbErrorKind::Retry, // Terminal for the same reason RuntimeShutdown is: the runtime // thread is gone and will not come back in this process. @@ -118,7 +114,6 @@ mod tests { DbErrorKind::Internal ); assert_eq!(SyncError::GetTimeout.kind(), DbErrorKind::Retry); - assert_eq!(SyncError::SetTimeout.kind(), DbErrorKind::Retry); assert_eq!(SyncError::RuntimeShutdown.kind(), DbErrorKind::Closed); // Terminal for the same reason: the runtime thread is gone and will // not come back in this process, so a caller must not retry. diff --git a/aimdb-sync/src/producer.rs b/aimdb-sync/src/producer.rs index 41983f1f..cc362d3d 100644 --- a/aimdb-sync/src/producer.rs +++ b/aimdb-sync/src/producer.rs @@ -2,7 +2,6 @@ use crate::runtime::Runtime; use crate::{SyncError, SyncResult}; -use aimdb_core::TryProduceError; use alloc::sync::Weak; use core::fmt::Debug; use core::marker::PhantomData; @@ -25,14 +24,8 @@ use core::marker::PhantomData; /// # #[derive(Clone, Debug, Serialize, Deserialize)] /// # struct Temperature { celsius: f32 } /// # fn example(producer: &SyncProducer) -> SyncResult<()> { -/// // Set value (blocks until sent) +/// // Set a value synchronously /// producer.set(Temperature { celsius: 25.0 })?; -/// -/// // Try to set (non-blocking) -/// match producer.try_set(Temperature { celsius: 27.0 }) { -/// Ok(()) => println!("Success"), -/// Err(_) => println!("Buffer full, try later"), -/// } /// # Ok(()) /// # } /// ``` @@ -109,16 +102,19 @@ where self.runtime()?.check() } - /// Set the value, blocking until it can be sent. + /// Set a value synchronously. /// - /// This call will block the current thread until the value can be sent to the runtime thread. - /// It's guaranteed to deliver the value eventually unless the runtime thread has shut down. + /// Checks that the runtime is usable, resolves the record by key, and + /// pushes directly into its buffer. This does not wait for buffer space; + /// all built-in buffers overwrite according to their configured semantics. /// /// # Errors /// - /// Returns `SyncError::RuntimeShutdown` if the runtime thread has been detached. - /// Returns any error from the underlying `produce()` operation (e.g., record not registered, - /// buffer full, etc.). + /// - [`SyncError::RuntimeShutdown`] if the runtime has shut down. + /// - [`SyncError::ForkedChild`] if this producer was inherited across a + /// Unix `fork()` without its runtime thread. + /// - [`SyncError::Db`] if the key is not registered or names a different + /// record type. /// /// # Example /// @@ -135,7 +131,7 @@ where /// .runtime(Arc::new(TokioAdapter)) /// .attach()?; /// let producer = handle.producer::("my_data")?; - /// producer.set(MyData { value: 42 })?; // blocks until value is sent and produced + /// producer.set(MyData { value: 42 })?; /// # Ok(()) /// # } /// ``` @@ -143,49 +139,6 @@ where let rt = self.runtime()?; rt.db()?.produce(&self.key, value).map_err(SyncError::Db) } - - /// Try to set the value without blocking. - /// - /// Pushes the value directly into the record's buffer. Unlike `set()`, this never - /// blocks: it fails immediately if the buffer is full instead of waiting for space. - /// - /// # Errors - /// - /// Returns `SyncError::SetTimeout` for bounded, non-overwriting buffer - /// implementations if the buffer is full. - /// Returns `SyncError::RuntimeShutdown` if the runtime thread has been detached. - /// - /// # Example - /// - /// ```no_run - /// use aimdb_core::AimDbBuilder; - /// use aimdb_sync::{AimDbBuilderSyncExt, SyncResult}; - /// use aimdb_tokio_adapter::TokioAdapter; - /// use std::sync::Arc; - /// - /// # #[derive(Debug, Clone)] - /// # struct MyData { value: i32 } - /// # fn main() -> SyncResult<()> { - /// let handle = AimDbBuilder::new() - /// .runtime(Arc::new(TokioAdapter)) - /// .attach()?; - /// let producer = handle.producer::("my_data")?; - /// match producer.try_set(MyData { value: 42 }) { - /// Ok(()) => println!("Sent immediately"), - /// Err(_) => println!("Buffer full or runtime shutdown"), - /// } - /// # Ok(()) - /// # } - /// ``` - pub fn try_set(&self, value: T) -> SyncResult<()> { - let rt = self.runtime()?; - let db = rt.db()?; - let producer = db.producer(&self.key)?; - producer.try_produce(value).map_err(|e| match e { - TryProduceError::Full(_) => SyncError::SetTimeout, - TryProduceError::Closed(_) => SyncError::RuntimeShutdown, - }) - } } /// Set-by-primitive verbs for `Settable` types (feature `data-contracts`). @@ -199,7 +152,7 @@ impl SyncProducer where T: aimdb_data_contracts::Settable + Send + 'static + Debug + Clone, { - /// Construct via `T::set(value, now)` and send. Blocking, like [`set`](Self::set). + /// Construct via `T::set(value, now)` and send via [`set`](Self::set). /// /// Stamps with the *caller's* `SystemTime` (sample time at the edge), not /// the engine's `ctx.time()` — use [`set_value_at`](Self::set_value_at) for @@ -240,11 +193,6 @@ where self.set(T::set(value, unix_now_ms())) } - /// Non-blocking variant, like [`try_set`](Self::try_set). - pub fn try_set_value(&self, value: T::Value) -> SyncResult<()> { - self.try_set(T::set(value, unix_now_ms())) - } - /// Explicit-timestamp variant (replay, testing). pub fn set_value_at(&self, value: T::Value, timestamp_ms: u64) -> SyncResult<()> { self.set(T::set(value, timestamp_ms)) diff --git a/aimdb-sync/tests/fork_safety_test.rs b/aimdb-sync/tests/fork_safety_test.rs index 0829af84..ca7fb339 100644 --- a/aimdb-sync/tests/fork_safety_test.rs +++ b/aimdb-sync/tests/fork_safety_test.rs @@ -52,12 +52,8 @@ fn a_forked_child_is_refused_rather_than_silently_dropped() { inherited.set(Reading { value: 1 }), Err(SyncError::ForkedChild) ); - let refused_try = matches!( - inherited.try_set(Reading { value: 2 }), - Err(SyncError::ForkedChild) - ); // The same refusal, asked rather than provoked: `check()` has to agree - // with the two publishes above. + // with the publish above. let check_agrees = matches!(inherited.check(), Err(SyncError::ForkedChild)); // Leak rather than free. The child is about to `_exit`, which reclaims // everything anyway, and `free` is the unsafe act here: it takes the @@ -67,7 +63,7 @@ fn a_forked_child_is_refused_rather_than_silently_dropped() { // the destructor, so there is nothing to lose by not running it. std::mem::forget(inherited); - refused && refused_try && check_agrees + refused && check_agrees }); assert_eq!(code, 0, "the child's publishes should have been refused"); diff --git a/aimdb-sync/tests/integration_test.rs b/aimdb-sync/tests/integration_test.rs index c960e8e2..2653cd97 100644 --- a/aimdb-sync/tests/integration_test.rs +++ b/aimdb-sync/tests/integration_test.rs @@ -179,11 +179,9 @@ fn test_non_blocking_operations() { let result = consumer.try_get(); assert!(matches!(result, Err(SyncError::GetTimeout))); - // Try set (should succeed immediately) + // Set should succeed immediately let test_value = test_value(); - producer - .try_set(test_value.clone()) - .expect("Failed to try_set"); + producer.set(test_value.clone()).expect("Failed to set"); // Use blocking get to ensure we receive the value // (try_get is inherently racy in this test scenario) @@ -266,17 +264,15 @@ fn check_reports_what_a_publish_would_find() { /// Test error handling - runtime shutdown, non-blocking operations #[test] fn test_runtime_shutdown_error_non_blocking() { - let (handle, producer, mut consumer) = setup(BufferCfg::SpmcRing { capacity: 10 }); + let handle = attach(BufferCfg::SpmcRing { capacity: 10 }); + let mut consumer = handle + .consumer::("test.data") + .expect("Failed to create consumer"); // Shut down the runtime handle.detach().expect("Failed to detach"); - // Non-blocking operations should now fail with RuntimeShutdown too - let test_value = test_value(); - - let result = producer.try_set(test_value); - assert!(matches!(result, Err(SyncError::RuntimeShutdown))); - + // Non-blocking reads should now fail with RuntimeShutdown too let result = consumer.try_get(); assert!(matches!(result, Err(SyncError::RuntimeShutdown))); } @@ -358,7 +354,7 @@ fn test_single_latest_semantics() { } // Use get_latest() to drain the channel and get the most recent value. - // The BufferLagged errors occuring during it are ignored + // BufferLagged errors occurring during it are ignored let latest = consumer.get_latest().expect("Failed to get latest"); // Should get the last value (5) since get_latest() drains the channel diff --git a/aimdb-sync/tests/settable_integration.rs b/aimdb-sync/tests/settable_integration.rs index 0b496c35..a89ead64 100644 --- a/aimdb-sync/tests/settable_integration.rs +++ b/aimdb-sync/tests/settable_integration.rs @@ -9,7 +9,6 @@ use aimdb_data_contracts::{SchemaType, Settable}; use aimdb_sync::AimDbBuilderSyncExt; use aimdb_tokio_adapter::{TokioAdapter, TokioRecordRegistrarExt}; use std::sync::Arc; -use std::thread; use std::time::Duration; #[derive(Debug, Clone, PartialEq)] @@ -63,38 +62,6 @@ fn set_value_constructs_produces_and_is_consumed() { handle.detach().expect("failed to detach"); } -#[test] -fn try_set_value_is_non_blocking_and_produces() { - let adapter = Arc::new(TokioAdapter); - let mut builder = AimDbBuilder::new().runtime(adapter); - - builder.configure::("temperature", |reg| { - reg.buffer(BufferCfg::SpmcRing { capacity: 10 }) - .tap(|_ctx, _consumer| async move {}); - }); - - let handle = builder.attach().expect("failed to attach"); - let producer = handle - .producer::("temperature") - .expect("failed to create producer"); - let mut consumer = handle - .consumer::("temperature") - .expect("failed to create consumer"); - - producer - .try_set_value(18.0) - .expect("try_set_value should succeed"); - - thread::sleep(Duration::from_millis(100)); - - let received = consumer - .get_with_timeout(Duration::from_secs(2)) - .expect("failed to consume"); - assert_eq!(received.celsius, 18.0); - - handle.detach().expect("failed to detach"); -} - #[test] fn set_value_at_stamps_the_explicit_timestamp() { let adapter = Arc::new(TokioAdapter); diff --git a/examples/sync-api-demo/src/main.rs b/examples/sync-api-demo/src/main.rs index 78f1bb1d..b365b402 100644 --- a/examples/sync-api-demo/src/main.rs +++ b/examples/sync-api-demo/src/main.rs @@ -168,7 +168,7 @@ fn main() -> Result<(), Box> { println!(" • Pure synchronous context (no #[tokio::main])"); println!(" • Multiple independent consumers"); println!(" • Blocking (get), timeout (get_with_timeout), and non-blocking (try_get) reads"); - println!(" • Blocking (set) and non-blocking (try_set) writes"); + println!(" • Synchronous writes with set"); println!(" • Multi-threaded producer-consumer patterns"); Ok(()) From 9b4cf8f8f03e63197c6eb33ec93f15a6b1db48c9 Mon Sep 17 00:00:00 2001 From: huvvi Date: Mon, 5 Oct 2026 10:35:02 +0300 Subject: [PATCH 2/2] docs(aimdb-sync): clarify sync API behavior --- aimdb-core/src/typed_api.rs | 13 +- aimdb-sync/README.md | 172 +++++++++++---------------- aimdb-sync/src/consumer.rs | 41 +++++-- aimdb-sync/src/fork.rs | 7 +- aimdb-sync/src/handle.rs | 23 +++- aimdb-sync/src/lib.rs | 37 ++---- aimdb-sync/src/producer.rs | 6 +- aimdb-sync/tests/integration_test.rs | 8 +- examples/sync-api-demo/src/main.rs | 4 +- 9 files changed, 157 insertions(+), 154 deletions(-) diff --git a/aimdb-core/src/typed_api.rs b/aimdb-core/src/typed_api.rs index d36e50d9..c11d27c8 100644 --- a/aimdb-core/src/typed_api.rs +++ b/aimdb-core/src/typed_api.rs @@ -160,13 +160,14 @@ where self.write.push(value); } - /// Extension point for write handles that may reject a value. + /// Lets a write handle report that it rejected a value. /// - /// The built-in buffers (`SpmcRing`, `SingleLatest`, `Mailbox`) overwrite - /// and currently always return `Ok(())`; use [`produce`](Self::produce) - /// for them. A future bounded, non-overwriting write handle can return the - /// value through [`TryProduceError::Full`], while a close-aware handle can - /// return it through [`TryProduceError::Closed`]. + /// All current buffers (`SpmcRing`, `SingleLatest`, `Mailbox`) overwrite + /// existing values when full and return `Ok(())`. Use + /// [`produce`](Self::produce) for them. A future buffer that refuses a + /// write when full can return the value through [`TryProduceError::Full`]. + /// A buffer that detects closure can return it through + /// [`TryProduceError::Closed`]. /// /// # Example /// diff --git a/aimdb-sync/README.md b/aimdb-sync/README.md index fb8cd16a..af4d7565 100644 --- a/aimdb-sync/README.md +++ b/aimdb-sync/README.md @@ -1,20 +1,23 @@ # aimdb-sync -Synchronous API wrapper for AimDB - blocking operations for async database. +Synchronous API wrapper for AimDB. ## Overview -`aimdb-sync` provides a synchronous interface to AimDB, enabling blocking operations on the async database. Perfect for FFI bindings, legacy codebases, simple scripts, and situations where async is impractical. +`aimdb-sync` provides a synchronous interface to AimDB for code that does not use an async executor. Perfect for FFI bindings, legacy codebases, simple scripts, and situations where async is impractical. **Key Features:** - **Pure Sync Context**: Works in plain `fn main()` - no `#[tokio::main]` required -- **Blocking Operations**: Familiar sync API (set, get, try_get, etc.) -- **Thread-Safe**: All types are `Send + Sync`, shareable across threads +- **Synchronous API**: Produce and consume records without `async` or `await` +- **Thread-Safe**: `SyncProducer` type is `Send + Sync`, shareable across threads. `SyncConsumer` type is `Send` only, can be moved to another thread. - **Type-Safe**: Full compile-time type safety with generics -- **Timeout Support**: All operations support configurable timeouts +- **Timeout Support**: Blocking consumer read operations support configurable timeouts ## Architecture +Producers write directly to AimDB record buffers. Consumers read from the +buffers through `Reader`. A background thread runs AimDB's async tasks. + ``` ┌──────────────────────────────┐ │ Synchronous Context │ @@ -24,11 +27,9 @@ Synchronous API wrapper for AimDB - blocking operations for async database. │ SyncConsumer ▼ ┌──────────────────────────────┐ -│ Channel Bridge │ -│ (tokio::sync::mpsc + │ -│ std::sync::mpsc) │ +│ AimDB Record Buffer │ └──────────────┬───────────────┘ - │ + | shared with ▼ ┌──────────────────────────────┐ │ Async Context │ @@ -43,9 +44,9 @@ Add to your `Cargo.toml`: ```toml [dependencies] -aimdb-sync = "0.5" -aimdb-core = "1.0" -aimdb-tokio-adapter = "0.5" +aimdb-sync = "0.6" +aimdb-core = "2.0" +aimdb-tokio-adapter = "0.7" ``` ### Basic Example @@ -69,15 +70,15 @@ fn main() -> Result<(), Box> { let adapter = Arc::new(TokioAdapter); let mut builder = AimDbBuilder::new().runtime(adapter); - builder.configure::(|reg| { + builder.configure::("temperature", |reg| { reg.buffer(BufferCfg::SingleLatest); }); let handle = builder.attach()?; // Get sync handles - let producer = handle.producer::()?; - let consumer = handle.consumer::()?; + let producer = handle.producer::("temperature")?; + let mut consumer = handle.consumer::("temperature")?; // Send from one thread let prod_handle = std::thread::spawn(move || { @@ -111,44 +112,22 @@ fn main() -> Result<(), Box> { ## Producer Operations -`SyncProducer` provides blocking send operations: +`SyncProducer` writes directly to the configured record buffer: -### Blocking Send +### Write a Value ```rust -let producer = handle.producer::()?; +let producer = handle.producer::("temperature")?; let temp = Temperature { celsius: 23.5, sensor_id: "sensor-001".to_string() }; -// Blocks until send completes +// Writes directly to the configured record buffer producer.set(temp)?; ``` -### Send with Timeout - -```rust -use std::time::Duration; - -// Block for max 1 second -match producer.set_with_timeout(temp, Duration::from_secs(1)) { - Ok(_) => println!("Sent successfully"), - Err(e) => eprintln!("Timeout or error: {}", e), -} -``` - -### Non-Blocking Send - -```rust -// Returns immediately (doesn't wait for produce to complete) -match producer.try_set(temp) { - Ok(_) => println!("Sent immediately"), - Err(e) => eprintln!("Channel full or error: {}", e), -} -``` - ## Consumer Operations `SyncConsumer` provides blocking receive operations: @@ -156,7 +135,7 @@ match producer.try_set(temp) { ### Blocking Receive ```rust -let consumer = handle.consumer::()?; +let mut consumer = handle.consumer::("temperature")?; // Blocks until value is available let temp = consumer.get()?; @@ -185,6 +164,19 @@ match consumer.try_get() { } ``` +### Latest Value + +`get_latest()` waits for one value, then returns the newest available value. +`get_latest_with_timeout()` applies its timeout only while waiting for the +first value. Both methods skip `BufferLagged` while catching up. + +```rust +let latest = consumer.get_latest()?; + +let latest_with_timeout = + consumer.get_latest_with_timeout(Duration::from_secs(1))?; +``` + ## Multi-Consumer Pattern Multiple consumers can receive from the same record: @@ -200,18 +192,18 @@ fn main() -> Result<(), Box> { let adapter = Arc::new(TokioAdapter); let mut builder = AimDbBuilder::new().runtime(adapter); - builder.configure::(|reg| { + builder.configure::("temperature", |reg| { reg.buffer(BufferCfg::SpmcRing { capacity: 16 }); }); let handle = builder.attach()?; - let producer = handle.producer::()?; + let producer = handle.producer::("temperature")?; // Spawn multiple consumer threads let mut handles = vec![]; for id in 0..3 { - let consumer = handle.consumer::()?; + let mut consumer = handle.consumer::("temperature")?; let handle = std::thread::spawn(move || { loop { match consumer.get_with_timeout(Duration::from_secs(1)) { @@ -244,37 +236,32 @@ fn main() -> Result<(), Box> { ## Thread Safety -All sync types are `Send + Sync`: +`SyncProducer` type is `Send + Sync`, shareable across threads. `SyncConsumer` type is `Send` only, can be moved to another thread: ```rust -use std::sync::Arc; +let producer = handle.producer::("temperature")?; +let mut consumer = handle.consumer::("temperature")?; -let producer = Arc::new(handle.producer::()?); -let consumer = Arc::new(handle.consumer::()?); - -// Share across threads +// `SyncProducer` implements `Clone` let prod_clone = producer.clone(); std::thread::spawn(move || { prod_clone.set(Temperature { celsius: 25.0, sensor_id: "s1".to_string() }).ok(); }); -let cons_clone = consumer.clone(); +// `SyncConsumer` implements `Send`, can be moved to another thread std::thread::spawn(move || { - let value = cons_clone.get().ok(); + consumer.get().ok(); }); ``` ## Error Handling ```rust -use aimdb_core::DbError; +use aimdb_sync::SyncError; match producer.set(temp) { Ok(_) => println!("Success"), - Err(DbError::SetTimeout) => { - eprintln!("Operation timed out"); - } - Err(DbError::RuntimeShutdown) => { + Err(SyncError::RuntimeShutdown) => { eprintln!("Runtime thread has stopped"); } Err(e) => { @@ -284,10 +271,15 @@ match producer.set(temp) { ``` Common error types: -- `DbError::SetTimeout` / `DbError::GetTimeout`: Operation exceeded timeout -- `DbError::RuntimeShutdown`: Runtime thread stopped or channel closed -- `DbError::RecordNotFound`: Type not registered in database -- `DbError::AttachFailed`: Failed to start runtime thread +- `SyncError::GetTimeout`: Operation exceeded timeout +- `SyncError::RuntimeShutdown`: Runtime thread stopped +- `SyncError::ForkedChild`: A forked child no longer has the runtime thread +- `SyncError::Db(DbError::RecordKeyNotFound {...})`: Type not registered in database +- `SyncError::Db(DbError::BufferLagged {...})`: An SPMC consumer missed overwritten values +- `SyncError::AttachFailed`: Failed to start runtime thread + +`set()` reports producer key and type errors. `handle.consumer()` reports +equivalent consumer errors during construction. ## Configuration Options @@ -298,39 +290,26 @@ Choose buffer based on use case: ```rust use aimdb_core::buffer::BufferCfg; -// SPMC Ring: Multiple consumers, bounded history -builder.configure::(|reg| { +// SPMC Ring: Multiple consumers, bounded history, overwrites the oldest value when full +builder.configure::("my-data", |reg| { reg.buffer(BufferCfg::SpmcRing { capacity: 100 }); }); -// SingleLatest: Always get newest value -builder.configure::(|reg| { +// SingleLatest: Always get newest value, each write replaces the previous value +builder.configure::("my-data", |reg| { reg.buffer(BufferCfg::SingleLatest); }); -// Mailbox: Single slot, overwrite -builder.configure::(|reg| { +// Mailbox: One pending value. Reading removes it +builder.configure::("my-data", |reg| { reg.buffer(BufferCfg::Mailbox); }); ``` -### Channel Capacity - -Control the sync bridge channel size: - -```rust -// Default capacity (100) -let producer = handle.producer::()?; -let consumer = handle.consumer::()?; - -// Custom capacity for high-frequency data -let producer = handle.producer_with_capacity::(1000)?; -let consumer = handle.consumer_with_capacity::(1000)?; -``` - ## Shutdown -Database automatically shuts down when dropped: +Call `detach()` to stop and join the runtime thread. Dropping the handle only +signals shutdown and does not wait: ```rust fn main() -> Result<(), Box> { @@ -344,7 +323,7 @@ fn main() -> Result<(), Box> { // Or with timeout // handle.detach_timeout(Duration::from_secs(5))?; - // Or just drop (automatic cleanup with warning) + // Dropping only signals shutdown // drop(handle); Ok(()) @@ -372,7 +351,7 @@ impl LegacyAdapter { let adapter = Arc::new(TokioAdapter); let mut builder = AimDbBuilder::new().runtime(adapter); - builder.configure::(|reg| { + builder.configure::("sensor-data", |reg| { reg.buffer(BufferCfg::SpmcRing { capacity: 100 }); }); @@ -381,15 +360,15 @@ impl LegacyAdapter { } pub fn send_sensor_data(&self, data: SensorData) -> Result<(), String> { - let producer = self.handle.producer::() + let producer = self.handle.producer::("sensor-data") .map_err(|e| e.to_string())?; - producer.set_with_timeout(data, Duration::from_secs(1)) + producer.set(data) .map_err(|e| e.to_string()) } pub fn read_sensor_data(&self) -> Result { - let consumer = self.handle.consumer::() + let mut consumer = self.handle.consumer::("sensor-data") .map_err(|e| e.to_string())?; consumer.get_with_timeout(Duration::from_secs(1)) @@ -416,13 +395,13 @@ fn main() -> Result<(), Box> { let adapter = Arc::new(TokioAdapter); let mut builder = AimDbBuilder::new().runtime(adapter); - builder.configure::(|reg| { + builder.configure::("log-message", |reg| { reg.buffer(BufferCfg::SpmcRing { capacity: 100 }); }); let handle = builder.attach()?; - let producer = handle.producer::()?; - let consumer = handle.consumer::()?; + let producer = handle.producer::("log-message")?; + let mut consumer = handle.consumer::("log-message")?; // Simple loop - no async/await loop { @@ -441,15 +420,10 @@ fn main() -> Result<(), Box> { ## Performance Considerations ### Overhead -- Channel crossing adds ~1-10μs latency -- Background runtime uses dedicated threads -- Memory: One tokio::mpsc channel per producer, one std::mpsc channel per consumer +- `attach()` starts a dedicated background thread that owns and drives the Tokio runtime. ### Optimization Tips -1. **Batch Operations**: Group multiple sets/gets when possible -2. **Avoid Blocking**: Use `try_*` methods in latency-sensitive paths -3. **Channel Capacity**: Tune for expected throughput -4. **Thread Count**: Match runtime_threads to workload +- **Avoid Blocking**: Use `SyncConsumer::try_get()` when the caller thread must not wait. ## Testing @@ -460,8 +434,6 @@ cargo test -p aimdb-sync # Run with logging RUST_LOG=debug cargo test -p aimdb-sync -- --nocapture -# Benchmark -cargo bench -p aimdb-sync ``` ## Complete Examples diff --git a/aimdb-sync/src/consumer.rs b/aimdb-sync/src/consumer.rs index d2ee8fa8..426348c2 100644 --- a/aimdb-sync/src/consumer.rs +++ b/aimdb-sync/src/consumer.rs @@ -100,6 +100,10 @@ where /// values were dropped. Not fatal, the next call resumes. /// - `SyncError::Db` for other errors that occurred during the read /// + /// # Panics + /// + /// Panics if called from inside a Tokio runtime. + /// /// # Example /// /// ```no_run @@ -141,6 +145,10 @@ where /// values were dropped. Not fatal, the next call resumes. /// - `SyncError::Db` for other errors that occurred during the read /// + /// # Panics + /// + /// Panics if called from inside a Tokio runtime. + /// /// # Example /// /// ```no_run @@ -166,6 +174,8 @@ where /// ``` pub fn get_with_timeout(&mut self, timeout: Duration) -> SyncResult { let (handle, reader) = self.reader.enter()?; + // Create the timer inside the async block. Creating it earlier would + // panic because no Tokio reactor is active on this thread yet. let fut = async { tokio::time::timeout(timeout, Self::get_impl(reader)).await }; let res = handle.block_on(fut); res.unwrap_or_else(|_| Err(SyncError::GetTimeout)) @@ -228,11 +238,18 @@ where /// /// The most recent available record of type `T`. /// + /// `BufferLagged` is skipped. Other errors are returned before the first + /// value. After a value is read, draining stops on an error and returns the + /// latest value. + /// /// # Errors - /// Note that the error is only reported if no value was retrieved at all. - /// Errors occuring after that are ignored; the latest obtained value is returned instead. + /// /// - `SyncError::RuntimeShutdown` if the runtime thread has stopped - /// - `SyncError::Db` if another error occured upon the very first read. + /// - `SyncError::Db` if another database error occurs before the first value + /// + /// # Panics + /// + /// Panics if called from inside a Tokio runtime. /// /// # Example /// @@ -260,7 +277,7 @@ where // 1) can simply sequence get_catch_up and try_get - // no one else does it simultaneously thanks to &mut self // 2) if draining ends up with an error, we follow the previous impl - // and return the latest succesfully read value + // and return the latest successfully read value // 3) potentially loops forever if producer keeps producing let oldest = self.get_catch_up(None)?; let latest = self.drain_remaining(oldest); @@ -275,12 +292,22 @@ where /// /// # Arguments /// - /// - `timeout`: Maximum time to wait for the first value + /// - `timeout`: Maximum time to wait for the first value. The drain after + /// that first value is not bounded by this timeout. + /// + /// `BufferLagged` is skipped. Other errors are returned before the first + /// value. After a value is read, draining stops on an error and returns the + /// latest value. /// /// # Errors /// /// - `SyncError::GetTimeout` if the timeout expires before any value arrives /// - `SyncError::RuntimeShutdown` if the runtime thread has stopped + /// - `SyncError::Db` if another database error occurs before the first value + /// + /// # Panics + /// + /// Panics if called from inside a Tokio runtime. /// /// # Example /// @@ -316,7 +343,7 @@ where } // Blocks until get() retrieves a value, or until the deadline is missed. - // Skips BufferLagged occuring in the process, and raises all other errors + // Skips BufferLagged occurring in the process, and raises all other errors fn get_catch_up(&mut self, deadline: Option) -> SyncResult { loop { let res = match deadline { @@ -345,7 +372,7 @@ where match self.try_get() { Ok(next) => cur = next, Err(SyncError::Db(DbError::BufferLagged { .. })) => continue, - // errors occured during draining will be ignored + // errors occurring during draining will be ignored Err(_) => return cur, } } diff --git a/aimdb-sync/src/fork.rs b/aimdb-sync/src/fork.rs index 552feced..ce7d7011 100644 --- a/aimdb-sync/src/fork.rs +++ b/aimdb-sync/src/fork.rs @@ -15,10 +15,9 @@ //! //! # Why `pthread_atfork` and not `getpid` //! -//! Because this sits on the publish path. Measured: `try_set` is 121 ns and -//! `std::process::id()` is 321 ns, so reading the pid per call would cost more -//! than twice the work it guards. A relaxed atomic load does not measurably -//! cost anything. +//! The publish path runs for every value. The `pthread_atfork` handler updates +//! a generation counter, so each publish only reads one relaxed atomic value. +//! This avoids asking the operating system for the process identity each time. //! //! # The process-global caveat //! diff --git a/aimdb-sync/src/handle.rs b/aimdb-sync/src/handle.rs index 213cd335..1f2e9c43 100644 --- a/aimdb-sync/src/handle.rs +++ b/aimdb-sync/src/handle.rs @@ -31,6 +31,10 @@ pub trait AimDbBuilderSyncExt { /// - `DbError::RuntimeError` if the database fails to build /// - `SyncError::AttachFailed` if the runtime thread fails to start /// + /// # Panics + /// + /// Panics if called from inside a Tokio runtime. + /// /// # Example /// /// ```no_run @@ -73,6 +77,10 @@ pub trait AimDbSyncExt { /// /// - `SyncError::AttachFailed` if the runtime thread fails to start /// + /// # Panics + /// + /// Panics if called from inside a Tokio runtime. + /// /// # Example /// /// ```no_run @@ -327,6 +335,10 @@ impl AimDbHandle { /// /// - `T`: The record type, must implement `TypedRecord` /// + /// Creating a producer does not check the key or type. Those checks happen + /// when [`set`](crate::SyncProducer::set) is called. The `SyncResult` + /// return type remains for API compatibility. + /// /// # Example /// /// ```no_run @@ -367,11 +379,14 @@ impl AimDbHandle { /// /// - `T`: The record type, must implement `TypedRecord` /// - /// # Errors (wrapped in SyncError::Db) + /// # Errors /// - /// - `DbError::RecordKeyNotFound` if type `T` was not registered - /// - `DbError::TypeMismatch` if the record type does not match `T` - /// - `DbError::MissingConfiguration` if the corresponding buffer was not configured + /// - `SyncError::ForkedChild` if called in a child process with an + /// inherited handle whose runtime thread did not survive the fork. + /// - `SyncError::RuntimeShutdown` if the database has shut down. + /// - `SyncError::Db(DbError::RecordKeyNotFound)` if the key was not registered. + /// - `SyncError::Db(DbError::TypeMismatch)` if the key names another type. + /// - `SyncError::Db(DbError::MissingConfiguration)` if its buffer was not configured. /// /// # Example /// diff --git a/aimdb-sync/src/lib.rs b/aimdb-sync/src/lib.rs index b3e494fc..0fa76f90 100644 --- a/aimdb-sync/src/lib.rs +++ b/aimdb-sync/src/lib.rs @@ -1,7 +1,7 @@ //! # AimDB Sync API //! -//! Synchronous API wrapper for AimDB that enables blocking operations -//! on the async database. Perfect for FFI, legacy codebases, and simple scripts. +//! Synchronous API wrapper for AimDB. Use it with FFI, legacy code, and simple +//! programs that do not run an async executor. //! //! ## Overview //! @@ -12,8 +12,8 @@ //! ## Features //! //! ### Producer Operations -//! - **`set()`**: Blocking send, waits if channel is full -//! - **`try_set()`**: Non-blocking send, returns immediately +//! - **`set()`**: Synchronously validates then pushes +//! directly into the configured buffer //! //! ### Consumer Operations //! - **`get()`**: Blocking receive, waits for value @@ -29,13 +29,11 @@ //! ## Architecture //! //! ```text -//! User Threads (sync) → Runtime Thread (async) -//! ↓ -//! AimDB (async) -//! ↓ -//! Buffers (SPMC, etc.) -//! ↓ -//! Consumer Threads (sync) +//! User Threads (sync) → AimDB → Buffers (SPMC, etc.) +//! ↓ +//! Consumer Threads (sync) +//! +//! Runtime Thread (async) runs AimDB's async tasks. //! ``` //! //! The runtime thread is created automatically when you call `attach()` on the builder. @@ -71,7 +69,7 @@ //! let producer = handle.producer::("sensor.temp")?; //! let mut consumer = handle.consumer::("sensor.temp")?; //! -//! // Producer: blocking operations +//! // Producer: synchronous direct write //! producer.set(Temperature { celsius: 25.0 })?; //! //! // Consumer: blocking operations @@ -130,20 +128,14 @@ //! - **User threads**: Unlimited - any number of threads can call operations concurrently //! - **Runtime thread**: One dedicated thread named "aimdb-sync-runtime" //! -//! ## Performance -//! -//! - **Latency**: Excellent for <50ms target, not suitable for hard low-latency requirements -//! //! ## Error Handling //! -//! All operations return [`SyncResult`] with facade-specific [`SyncError`] -//! variants: +//! All operations return [`SyncResult`]. The main [`SyncError`] variants are: //! //! - `RuntimeShutdown`: The runtime thread stopped //! - `ForkedChild`: Created before a `fork()`, and this is the child — the //! runtime thread it needs did not survive, so a `set()` that would have //! returned `Ok` into a buffer nobody drains is refused instead (std, Unix) -//! - `SetTimeout`: Producer timeout expired //! - `GetTimeout`: Consumer timeout expired or no data (try_get) //! - `AttachFailed`: Failed to start runtime thread, carrying the `DbError` //! that caused it — so `kind()` reports a bad record graph as @@ -154,11 +146,8 @@ //! //! ### Error Propagation //! -//! Producer errors are propagated synchronously back to the caller: -//! - `set()` blocks until the produce operation completes and returns any errors -//! that occur -//! - `try_set()` returns immediately: `Ok(())` if the record's buffer accepted the -//! value, `SyncError::SetTimeout` if it didn't (bounded, non-overwriting buffer, full) +//! `set()` writes directly to the record buffer and returns any runtime, key, +//! or type error. //! #![cfg_attr(feature = "std", doc = "```no_run")] #![cfg_attr(not(feature = "std"), doc = "```ignore")] diff --git a/aimdb-sync/src/producer.rs b/aimdb-sync/src/producer.rs index cc362d3d..2e53b64e 100644 --- a/aimdb-sync/src/producer.rs +++ b/aimdb-sync/src/producer.rs @@ -104,9 +104,9 @@ where /// Set a value synchronously. /// - /// Checks that the runtime is usable, resolves the record by key, and - /// pushes directly into its buffer. This does not wait for buffer space; - /// all built-in buffers overwrite according to their configured semantics. + /// Checks the runtime, finds the record by key, checks its type, and writes + /// directly to its buffer. This method does not wait for buffer space. + /// Every current buffer overwrites according to its configured behavior. /// /// # Errors /// diff --git a/aimdb-sync/tests/integration_test.rs b/aimdb-sync/tests/integration_test.rs index 2653cd97..c37a438e 100644 --- a/aimdb-sync/tests/integration_test.rs +++ b/aimdb-sync/tests/integration_test.rs @@ -170,9 +170,9 @@ fn test_timeout_operations() { handle.detach().expect("Failed to detach"); } -/// Test non-blocking operations +/// Test an immediate consumer read. #[test] -fn test_non_blocking_operations() { +fn test_non_blocking_consumer_operation() { let (handle, producer, mut consumer) = setup(BufferCfg::SpmcRing { capacity: 10 }); // Try get on empty buffer (should fail) @@ -261,7 +261,7 @@ fn check_reports_what_a_publish_would_find() { )); } -/// Test error handling - runtime shutdown, non-blocking operations +/// Test an immediate consumer read after runtime shutdown. #[test] fn test_runtime_shutdown_error_non_blocking() { let handle = attach(BufferCfg::SpmcRing { capacity: 10 }); @@ -272,7 +272,7 @@ fn test_runtime_shutdown_error_non_blocking() { // Shut down the runtime handle.detach().expect("Failed to detach"); - // Non-blocking reads should now fail with RuntimeShutdown too + // Immediate reads should now fail with RuntimeShutdown too let result = consumer.try_get(); assert!(matches!(result, Err(SyncError::RuntimeShutdown))); } diff --git a/examples/sync-api-demo/src/main.rs b/examples/sync-api-demo/src/main.rs index b365b402..0c84bbdf 100644 --- a/examples/sync-api-demo/src/main.rs +++ b/examples/sync-api-demo/src/main.rs @@ -13,7 +13,7 @@ //! 1. Building an AimDB instance with a typed record //! 2. Attaching the database to get a sync handle //! 3. Creating producers and consumers in sync context -//! 4. Setting and getting values using blocking operations +//! 4. Synchronous writes and blocking or immediate reads //! 5. Multi-threaded producer-consumer patterns //! 6. Clean shutdown with detach() @@ -136,7 +136,7 @@ fn main() -> Result<(), Box> { }; println!(" Main: Setting temperature {:.1}°C", temp.celsius); - // Use blocking send + // Write synchronously if let Err(e) = producer.set(temp) { eprintln!(" Error setting value: {}", e); }