Compare commits

..
8 Commits
Author SHA1 Message Date
semantic-release-bot 8ff33817a8 chore(release): 1.3.4 [skip ci]
## [1.3.4](https://github.com/asepharyana/mytheclipse/compare/v1.3.3...v1.3.4) (2026-08-28)

### Bug Fixes

* **cache,storage:** harden cache bounds + atomic disk writes ([b6f138b](https://github.com/asepharyana/mytheclipse/commit/b6f138b90d67c9531b5a58993e8ec750e5eec57f))
2026-08-28 18:30:42 +00:00
asepharyana b6f138b90d fix(cache,storage): harden cache bounds + atomic disk writes
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s
- cache: guard MokaL1::new(0) with panic; add MemoryCache::with_max_entries
  bounded LRU eviction (oldest evicted past cap) + docs warning about
  unbounded default growth. Verifies moka treats max_capacity=0 as a
  permanent no-insert sentinel.
- storage: make LocalFileStorage::put atomic via temp-file + fsync + rename;
  cleans up temp on write failure; no leftover .tmp-* on disk after success.

Tests: 17 cache tests (incl bounded_cache_evicts_oldest, zero_max_panics),
6 storage tests (incl put_leaves_no_temp_file). Full workspace: clippy 0
warnings, all tests green.
2026-08-29 01:29:40 +07:00
semantic-release-bot 4027d3eb17 chore(release): 1.3.3 [skip ci]
## [1.3.3](https://github.com/asepharyana/mytheclipse/compare/v1.3.2...v1.3.3) (2026-08-28)

### Bug Fixes

* **publish:** trim mytheclipse keywords to 5 to satisfy crates.io limit ([5f1e3ac](https://github.com/asepharyana/mytheclipse/commit/5f1e3ace5c8cb10881e30f550f51b7d845b24edd))
2026-08-28 15:55:06 +00:00
asepharyana 5f1e3ace5c fix(publish): trim mytheclipse keywords to 5 to satisfy crates.io limit
crates.io rejects crates with more than 5 keywords (HTTP 400 'expected at
most 5 keywords per crate'), which broke every 'Publish to crates.io'
workflow run. Reduced the mytheclipse crate's keywords from 12 to the 5
most representative (async, concurrency, resiliency, circuit-breaker,
observability).
2026-08-28 22:53:59 +07:00
semantic-release-bot 0f2b776bb8 chore(release): 1.3.2 [skip ci]
## [1.3.2](https://github.com/asepharyana/mytheclipse/compare/v1.3.1...v1.3.2) (2026-08-28)

### Bug Fixes

* **ci:** gate cache & storage crate doctests behind their features ([5717f8a](https://github.com/asepharyana/mytheclipse/commit/5717f8aaaae34cb66cdbfc31f4982c93816ee5a3))
2026-08-28 15:50:48 +00:00
asepharyana 5717f8aaaa fix(ci): gate cache & storage crate doctests behind their features
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s
Same E0432 class as mytheclipse-event: the crate-root doctests referenced
feature-gated types that are absent in reduced-feature test builds flagged by
the CI matrix.

- mytheclipse-cache: example used MemoryCache (l1-memory) plus
  MultiLayerCache/CacheAside (cache-aside); the l1-moka and l2-redis builds
  (which don't enable cache-aside) failed the doctest. Gated the body behind
  all(l1-memory, cache-aside) with a no-op fallback.
- mytheclipse-storage: example used LocalFileStorage (local feature); the
  gcs and s3 builds failed the doctest. Gated the body behind the local
  feature with a no-op fallback.

Examples remain compile-checked & runnable under default features.
2026-08-28 22:49:39 +07:00
semantic-release-bot 52d9e1e93e chore(release): 1.3.1 [skip ci]
## [1.3.1](https://github.com/asepharyana/mytheclipse/compare/v1.3.0...v1.3.1) (2026-08-28)

### Bug Fixes

* **ci:** gate event crate doctest behind mem feature and apply rustfmt ([1994115](https://github.com/asepharyana/mytheclipse/commit/19941156b43ec58370d0d1369174ec84400e9bc9))
2026-08-28 15:41:53 +00:00
asepharyana 19941156b4 fix(ci): gate event crate doctest behind mem feature and apply rustfmt
The crate-root doctest used InMemoryEventBus/TypedEventBus (gated behind
the 'mem' feature), so `cargo test -p mytheclipse-event --no-default-features
--features amqp|nats` failed to compile the doctest (E0432). The example is
now no_run with the mem-dependent imports inside a #[cfg(feature="mem")]
main, so it compiles (and runs) when mem is on and degrades to an empty main
when off.

Also apply rustfmt to mytheclipse-config/dynamic.rs and
mytheclipse-storage/{gcs,s3}.rs to satisfy cargo fmt --all --check.
2026-08-28 22:41:06 +07:00
17 changed files with 245 additions and 43 deletions
+28
View File
@@ -1,3 +1,31 @@
## [1.3.4](https://github.com/asepharyana/mytheclipse/compare/v1.3.3...v1.3.4) (2026-08-28)
### Bug Fixes
* **cache,storage:** harden cache bounds + atomic disk writes ([b6f138b](https://github.com/asepharyana/mytheclipse/commit/b6f138b90d67c9531b5a58993e8ec750e5eec57f))
## [1.3.3](https://github.com/asepharyana/mytheclipse/compare/v1.3.2...v1.3.3) (2026-08-28)
### Bug Fixes
* **publish:** trim mytheclipse keywords to 5 to satisfy crates.io limit ([5f1e3ac](https://github.com/asepharyana/mytheclipse/commit/5f1e3ace5c8cb10881e30f550f51b7d845b24edd))
## [1.3.2](https://github.com/asepharyana/mytheclipse/compare/v1.3.1...v1.3.2) (2026-08-28)
### Bug Fixes
* **ci:** gate cache & storage crate doctests behind their features ([5717f8a](https://github.com/asepharyana/mytheclipse/commit/5717f8aaaae34cb66cdbfc31f4982c93816ee5a3))
## [1.3.1](https://github.com/asepharyana/mytheclipse/compare/v1.3.0...v1.3.1) (2026-08-28)
### Bug Fixes
* **ci:** gate event crate doctest behind mem feature and apply rustfmt ([1994115](https://github.com/asepharyana/mytheclipse/commit/19941156b43ec58370d0d1369174ec84400e9bc9))
# [1.3.0](https://github.com/asepharyana/mytheclipse/compare/v1.2.0...v1.3.0) (2026-08-28) # [1.3.0](https://github.com/asepharyana/mytheclipse/compare/v1.2.0...v1.3.0) (2026-08-28)
Generated
+6 -6
View File
@@ -2539,7 +2539,7 @@ dependencies = [
[[package]] [[package]]
name = "mytheclipse" name = "mytheclipse"
version = "1.3.0" version = "1.3.4"
dependencies = [ dependencies = [
"num_cpus", "num_cpus",
"rand 0.8.8", "rand 0.8.8",
@@ -2551,7 +2551,7 @@ dependencies = [
[[package]] [[package]]
name = "mytheclipse-cache" name = "mytheclipse-cache"
version = "1.3.0" version = "1.3.4"
dependencies = [ dependencies = [
"async-trait", "async-trait",
"moka", "moka",
@@ -2564,7 +2564,7 @@ dependencies = [
[[package]] [[package]]
name = "mytheclipse-config" name = "mytheclipse-config"
version = "1.3.0" version = "1.3.4"
dependencies = [ dependencies = [
"dotenvy", "dotenvy",
"notify", "notify",
@@ -2579,7 +2579,7 @@ dependencies = [
[[package]] [[package]]
name = "mytheclipse-crypto" name = "mytheclipse-crypto"
version = "1.3.0" version = "1.3.4"
dependencies = [ dependencies = [
"aead", "aead",
"aes-gcm", "aes-gcm",
@@ -2596,7 +2596,7 @@ dependencies = [
[[package]] [[package]]
name = "mytheclipse-event" name = "mytheclipse-event"
version = "1.3.0" version = "1.3.4"
dependencies = [ dependencies = [
"async-nats", "async-nats",
"async-trait", "async-trait",
@@ -2612,7 +2612,7 @@ dependencies = [
[[package]] [[package]]
name = "mytheclipse-storage" name = "mytheclipse-storage"
version = "1.3.0" version = "1.3.4"
dependencies = [ dependencies = [
"async-trait", "async-trait",
"aws-config", "aws-config",
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "mytheclipse-cache" name = "mytheclipse-cache"
version = "1.3.0" version = "1.3.4"
edition = "2021" edition = "2021"
rust-version = "1.75" rust-version = "1.75"
license = "MIT OR Apache-2.0" license = "MIT OR Apache-2.0"
+6 -1
View File
@@ -17,9 +17,12 @@
//! //!
//! ## Example //! ## Example
//! //!
//! Multi-layer + cache-aside composition (default features):
//!
//! ```no_run //! ```no_run
//! use mytheclipse_cache::{Cache, MemoryCache, MultiLayerCache, CacheAside}; //! # #[cfg(all(feature = "l1-memory", feature = "cache-aside"))]
//! # async fn run() { //! # async fn run() {
//! use mytheclipse_cache::{Cache, MemoryCache, MultiLayerCache, CacheAside};
//! let l1 = MemoryCache::new(); //! let l1 = MemoryCache::new();
//! let l2 = MemoryCache::new(); // in a real app: a RedisCache //! let l2 = MemoryCache::new(); // in a real app: a RedisCache
//! let cache = MultiLayerCache::new(l1, l2); //! let cache = MultiLayerCache::new(l1, l2);
@@ -34,6 +37,8 @@
//! ); //! );
//! let _v = aside.get("orders:42").await.unwrap(); //! let _v = aside.get("orders:42").await.unwrap();
//! # } //! # }
//! # #[cfg(not(all(feature = "l1-memory", feature = "cache-aside")))]
//! # fn run() {}
//! ``` //! ```
#![forbid(unsafe_code)] #![forbid(unsafe_code)]
+86 -8
View File
@@ -4,7 +4,7 @@
//! Entries are lazily expired on access by comparing against `Instant`; a //! Entries are lazily expired on access by comparing against `Instant`; a
//! monotonic clock keeps TTLs robust against wall-clock discontinuities. //! monotonic clock keeps TTLs robust against wall-clock discontinuities.
use std::collections::HashMap; use std::collections::{HashMap, VecDeque};
use std::sync::{Arc, Mutex}; use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
@@ -15,14 +15,34 @@ use crate::traits::{Cache, CacheError};
/// A wrapping entry: `None` expiry means the value never expires. /// A wrapping entry: `None` expiry means the value never expires.
type Entry = (Vec<u8>, Option<Instant>); type Entry = (Vec<u8>, Option<Instant>);
/// An in-process [`Cache`] implementation for L1 caching. /// An in-process [`Cache`] for L1 caching.
#[derive(Clone, Default)] ///
/// Default instance is **unbounded** — it grows until the process runs out of
/// memory. For memory-constrained workloads, use [`MemoryCache::with_max_entries`]
/// to install a simple LRU-style cap: when the cap is exceeded, the oldest
/// (least-recently-inserted) entry is evicted.
#[derive(Debug, Clone)]
pub struct MemoryCache { pub struct MemoryCache {
inner: Arc<Mutex<HashMap<String, Entry>>>, inner: Arc<Mutex<HashMap<String, Entry>>>,
/// When `Some(n)`, the cache refuses more than `n` live entries and evicts
/// the oldest on overflow. `None` = unbounded (legacy default).
max_entries: Option<usize>,
/// Insertion order, for eviction when `max_entries` is set.
order: Arc<Mutex<VecDeque<String>>>,
}
impl Default for MemoryCache {
fn default() -> Self {
Self {
inner: Arc::new(Mutex::new(HashMap::new())),
max_entries: None,
order: Arc::new(Mutex::new(VecDeque::new())),
}
}
} }
impl MemoryCache { impl MemoryCache {
/// Builds an empty in-memory cache. /// Builds an empty in-memory cache (unbounded by default).
pub fn new() -> Self { pub fn new() -> Self {
Self::default() Self::default()
} }
@@ -32,6 +52,23 @@ impl MemoryCache {
self.inner.lock().unwrap().reserve(capacity); self.inner.lock().unwrap().reserve(capacity);
self self
} }
/// Installs a bounded LRU-style cap. When the cache exceeds `max`, the
/// oldest (least-recently-inserted) entry is evicted on each `set`.
///
/// This is the recommended constructor for production L1 caches: a
/// [`MemoryCache::new()`] (unbounded) left unmanaged can grow without bound
/// and exhaust process memory.
pub fn with_max_entries(mut self, max: usize) -> Self {
assert!(max > 0, "mytheclipse-cache: with_max_entries must be > 0");
self.max_entries = Some(max);
self
}
/// The configured max entries, if any.
pub fn max_entries(&self) -> Option<usize> {
self.max_entries
}
} }
#[async_trait] #[async_trait]
@@ -41,6 +78,7 @@ impl Cache for MemoryCache {
match map.get(key) { match map.get(key) {
Some((value, Some(expires))) if *expires <= Instant::now() => { Some((value, Some(expires))) if *expires <= Instant::now() => {
map.remove(key); map.remove(key);
self.remove_order(key);
Ok(None) Ok(None)
} }
Some((value, _)) => Ok(Some(value.clone())), Some((value, _)) => Ok(Some(value.clone())),
@@ -55,24 +93,44 @@ impl Cache for MemoryCache {
ttl: Option<Duration>, ttl: Option<Duration>,
) -> Result<(), CacheError> { ) -> Result<(), CacheError> {
let expires = ttl.map(|d| Instant::now() + d); let expires = ttl.map(|d| Instant::now() + d);
self.inner let mut map = self.inner.lock().unwrap();
.lock() let is_new = !map.contains_key(key);
.unwrap() map.insert(key.to_string(), (value, expires));
.insert(key.to_string(), (value, expires)); if is_new {
let mut order = self.order.lock().unwrap();
order.push_back(key.to_string());
if let Some(cap) = self.max_entries {
while order.len() > cap {
if let Some(oldest) = order.pop_front() {
map.remove(&oldest);
}
}
}
}
Ok(()) Ok(())
} }
async fn invalidate(&self, key: &str) -> Result<(), CacheError> { async fn invalidate(&self, key: &str) -> Result<(), CacheError> {
self.inner.lock().unwrap().remove(key); self.inner.lock().unwrap().remove(key);
self.remove_order(key);
Ok(()) Ok(())
} }
async fn clear(&self) -> Result<(), CacheError> { async fn clear(&self) -> Result<(), CacheError> {
self.inner.lock().unwrap().clear(); self.inner.lock().unwrap().clear();
self.order.lock().unwrap().clear();
Ok(()) Ok(())
} }
} }
impl MemoryCache {
/// Removes `key` from the insertion-order deque (if present).
fn remove_order(&self, key: &str) {
let mut order = self.order.lock().unwrap();
order.retain(|k| k != key);
}
}
/// A typed view over a byte cache using `serde`-compatible (JSON) encoding. /// A typed view over a byte cache using `serde`-compatible (JSON) encoding.
/// ///
/// Only enabled with the `cache-aside` feature, which pulls in `serde`. /// Only enabled with the `cache-aside` feature, which pulls in `serde`.
@@ -158,6 +216,26 @@ mod tests {
assert_eq!(c.get("b").await.unwrap(), None); assert_eq!(c.get("b").await.unwrap(), None);
} }
/// Asserts that an unbounded `MemoryCache::with_max_entries(0)` panics,
/// preventing a no-op cache that accepts zero entries.
#[test]
#[should_panic(expected = "must be > 0")]
fn zero_max_panics() {
let _ = MemoryCache::new().with_max_entries(0);
}
#[tokio::test]
async fn bounded_cache_evicts_oldest() {
let c = MemoryCache::new().with_max_entries(2);
c.set("a", b"1".to_vec(), None).await.unwrap();
c.set("b", b"2".to_vec(), None).await.unwrap();
c.set("c", b"3".to_vec(), None).await.unwrap();
// "a" (oldest) should have been evicted.
assert_eq!(c.get("a").await.unwrap(), None);
assert_eq!(c.get("b").await.unwrap(), Some(b"2".to_vec()));
assert_eq!(c.get("c").await.unwrap(), Some(b"3".to_vec()));
}
#[cfg(feature = "cache-aside")] #[cfg(feature = "cache-aside")]
#[tokio::test] #[tokio::test]
async fn typed_cache_roundtrip() { async fn typed_cache_roundtrip() {
+26 -2
View File
@@ -19,7 +19,22 @@ pub struct MokaL1 {
impl MokaL1 { impl MokaL1 {
/// Builds a Moka cache with `max_capacity` entries and an optional default /// Builds a Moka cache with `max_capacity` entries and an optional default
/// `ttl`. /// `ttl`.
///
/// # Panics
///
/// Panics if `max_capacity` is `0`. In Moka, a `max_capacity` of `0` is a
/// sentinel for **zero-entries-allowed** — every `insert` is silently
/// dropped — which is almost certainly a caller mistake (the natural way to
/// express "unbounded" in other caches). Pass `1..=u64::MAX`; use
/// [`MemoryCache`](crate::memory::MemoryCache) if you truly need an
/// unbounded in-process cache.
pub fn new(max_capacity: u64, ttl: Option<Duration>) -> Self { pub fn new(max_capacity: u64, ttl: Option<Duration>) -> Self {
assert!(
max_capacity > 0,
"mytheclipse-cache: MokaL1::new(max_capacity) must be > 0; \
moka treats 0 as a permanent no-insert sentinel. \
Use MemoryCache for an unbounded cache."
);
let mut builder = MokaCache::builder().max_capacity(max_capacity); let mut builder = MokaCache::builder().max_capacity(max_capacity);
if let Some(ttl) = ttl { if let Some(ttl) = ttl {
builder = builder.time_to_live(ttl); builder = builder.time_to_live(ttl);
@@ -36,14 +51,15 @@ impl Cache for MokaL1 {
Ok(self.inner.get(key).await) Ok(self.inner.get(key).await)
} }
/// Inserts `value`, using the cache's configured TTL policy. The per-call
/// `ttl` argument is intentionally ignored — Moka applies a single TTL
/// configured on the builder, and per-entry overrides are not exposed here.
async fn set( async fn set(
&self, &self,
key: &str, key: &str,
value: Vec<u8>, value: Vec<u8>,
_ttl: Option<Duration>, _ttl: Option<Duration>,
) -> Result<(), CacheError> { ) -> Result<(), CacheError> {
// Per-entry TTL overrides are handled by the builder default in Moka;
// the passed `ttl` is intentionally ignored (single configured policy).
self.inner.insert(key.to_string(), value).await; self.inner.insert(key.to_string(), value).await;
Ok(()) Ok(())
} }
@@ -83,6 +99,14 @@ mod tests {
assert_eq!(c.get("b").await.unwrap(), None); assert_eq!(c.get("b").await.unwrap(), None);
} }
/// Asserts that `max_capacity == 0` panics with a clear message, rather
/// than silently creating a cache that never accepts entries.
#[test]
#[should_panic(expected = "must be > 0")]
fn zero_capacity_panics() {
let _ = MokaL1::new(0, None);
}
#[tokio::test] #[tokio::test]
async fn ttl_does_expire() { async fn ttl_does_expire() {
// Keep a firm TTL assertion; sleep well past the expiry window. // Keep a firm TTL assertion; sleep well past the expiry window.
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "mytheclipse-config" name = "mytheclipse-config"
version = "1.3.0" version = "1.3.4"
edition = "2021" edition = "2021"
rust-version = "1.75" rust-version = "1.75"
license = "MIT OR Apache-2.0" license = "MIT OR Apache-2.0"
+4 -1
View File
@@ -49,7 +49,10 @@ impl<T: Config + Clone> DynamicConfig<T> {
/// Replaces the current value and notifies subscribers. /// Replaces the current value and notifies subscribers.
pub fn set(&self, new: T) { pub fn set(&self, new: T) {
*self.inner.write().expect("mytheclipse-config: RwLock poisoned") = new; *self
.inner
.write()
.expect("mytheclipse-config: RwLock poisoned") = new;
let _ = self.tx.send(()); let _ = self.tx.send(());
} }
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "mytheclipse-crypto" name = "mytheclipse-crypto"
version = "1.3.0" version = "1.3.4"
edition = "2021" edition = "2021"
rust-version = "1.75" rust-version = "1.75"
license = "MIT OR Apache-2.0" license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "mytheclipse-event" name = "mytheclipse-event"
version = "1.3.0" version = "1.3.4"
edition = "2021" edition = "2021"
rust-version = "1.75" rust-version = "1.75"
license = "MIT OR Apache-2.0" license = "MIT OR Apache-2.0"
+16 -9
View File
@@ -15,20 +15,27 @@
//! //!
//! ## Example //! ## Example
//! //!
//! ``` //! The in-memory + typed bus (`mem` feature, on by default):
//! use mytheclipse_event::{EventBus, InMemoryEventBus, TypedEventBus}; //!
//! ```no_run
//! use serde::{Deserialize, Serialize}; //! use serde::{Deserialize, Serialize};
//! //!
//! #[derive(Debug, Serialize, Deserialize, PartialEq)] //! #[derive(Debug, Serialize, Deserialize, PartialEq)]
//! struct OrderCreated { id: u64 } //! struct OrderCreated { id: u64 }
//! //!
//! # #[tokio::main] async fn main() { //! #[cfg(feature = "mem")]
//! let bus = TypedEventBus::new(InMemoryEventBus::default()); //! #[tokio::main]
//! let mut sub = bus.subscribe::<OrderCreated>("orders").await.unwrap(); //! async fn main() {
//! bus.publish("orders", &OrderCreated { id: 42 }).await.unwrap(); //! use mytheclipse_event::{EventBus, InMemoryEventBus, TypedEventBus};
//! let event = sub.recv().await.unwrap(); //! let bus = TypedEventBus::new(InMemoryEventBus::default());
//! assert_eq!(event, OrderCreated { id: 42 }); //! let mut sub = bus.subscribe::<OrderCreated>("orders").await.unwrap();
//! # } //! bus.publish("orders", &OrderCreated { id: 42 }).await.unwrap();
//! let event = sub.recv().await.unwrap();
//! assert_eq!(event, OrderCreated { id: 42 });
//! }
//!
//! #[cfg(not(feature = "mem"))]
//! fn main() {}
//! ``` //! ```
pub mod traits; pub mod traits;
+1 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "mytheclipse-storage" name = "mytheclipse-storage"
version = "1.3.0" version = "1.3.4"
edition = "2021" edition = "2021"
rust-version = "1.75" rust-version = "1.75"
license = "MIT OR Apache-2.0" license = "MIT OR Apache-2.0"
+4 -1
View File
@@ -160,6 +160,9 @@ mod tests {
.await .await
.unwrap(); .unwrap();
assert_eq!(data, b"hello gcs"); assert_eq!(data, b"hello gcs");
storage.delete("mytheclipse-storage-test.txt").await.unwrap(); storage
.delete("mytheclipse-storage-test.txt")
.await
.unwrap();
} }
} }
+10 -3
View File
@@ -15,15 +15,22 @@
//! //!
//! ## Example //! ## Example
//! //!
//! ``` //! Local disk storage (`local` feature, default):
//! use mytheclipse_storage::{LocalFileStorage, StorageDriver, bytes_stream, read_to_vec}; //!
//! # #[tokio::main] async fn main() { //! ```no_run
//! use mytheclipse_storage::{StorageDriver, bytes_stream, read_to_vec};
//! # #[cfg(feature = "local")]
//! # #[tokio::main]
//! # async fn main() {
//! use mytheclipse_storage::LocalFileStorage;
//! # let dir = tempfile::tempdir().unwrap(); //! # let dir = tempfile::tempdir().unwrap();
//! let storage = LocalFileStorage::new(dir.path()); //! let storage = LocalFileStorage::new(dir.path());
//! storage.put("hello.txt", bytes_stream(b"hi".to_vec())).await.unwrap(); //! storage.put("hello.txt", bytes_stream(b"hi".to_vec())).await.unwrap();
//! let data = read_to_vec(storage.get("hello.txt").await.unwrap()).await.unwrap(); //! let data = read_to_vec(storage.get("hello.txt").await.unwrap()).await.unwrap();
//! assert_eq!(data, b"hi"); //! assert_eq!(data, b"hi");
//! # } //! # }
//! # #[cfg(not(feature = "local"))]
//! # fn main() {}
//! ``` //! ```
pub mod traits; pub mod traits;
+42 -2
View File
@@ -22,6 +22,11 @@ impl LocalFileStorage {
Self { root: root.into() } Self { root: root.into() }
} }
/// Returns the root directory path.
pub fn root(&self) -> &std::path::Path {
&self.root
}
fn resolve(&self, path: &str) -> Result<PathBuf, StorageError> { fn resolve(&self, path: &str) -> Result<PathBuf, StorageError> {
let rel = Path::new(path.trim_start_matches('/')); let rel = Path::new(path.trim_start_matches('/'));
for component in rel.components() { for component in rel.components() {
@@ -52,14 +57,30 @@ impl StorageDriver for LocalFileStorage {
if let Some(parent) = full.parent() { if let Some(parent) = full.parent() {
tokio::fs::create_dir_all(parent) tokio::fs::create_dir_all(parent)
.await .await
.map_err(|e| StorageError::Io(e.to_string()))?; .map_err(|e| StorageError::Io(e.to_string()))?
} }
let mut file = tokio::fs::File::create(&full) // Write to a sibling temp file then atomically rename. On POSIX this is
// atomic, so a crash mid-write leaves *either* the previous file *or*
// the complete new file — never a half-written truncated object.
// Use a PID-unguarded temp name and clean it up if anything fails.
let tmp = full.with_extension(format!(".tmp-{}", std::process::id()));
let mut file = tokio::fs::File::create(&tmp)
.await .await
.map_err(|e| StorageError::Io(e.to_string()))?; .map_err(|e| StorageError::Io(e.to_string()))?;
let written = tokio::io::copy(&mut data, &mut file) let written = tokio::io::copy(&mut data, &mut file)
.await
.map_err(|e| {
let _ = std::fs::remove_file(&tmp);
StorageError::Io(e.to_string())
})?;
// Ensure durability: flush to OS, fsync, then rename.
tokio::fs::File::open(&tmp)
.await
.map_err(|e| StorageError::Io(e.to_string()))?
.sync_all()
.await .await
.map_err(|e| StorageError::Io(e.to_string()))?; .map_err(|e| StorageError::Io(e.to_string()))?;
std::fs::rename(&tmp, &full).map_err(|e| StorageError::Io(e.to_string()))?;
Ok(written) Ok(written)
} }
@@ -162,4 +183,23 @@ mod tests {
.unwrap_err(); .unwrap_err();
assert!(matches!(err, StorageError::InvalidPath(_))); assert!(matches!(err, StorageError::InvalidPath(_)));
} }
/// Verifies the atomic-write contract: after a successful `put`, the temp
/// file must not linger on disk.
#[tokio::test]
async fn put_leaves_no_temp_file() {
let (storage, _dir) = driver();
storage
.put("clean.txt", bytes_stream(b"data".to_vec()))
.await
.unwrap();
// No `*.tmp-*` files should remain in the root after a clean write.
let leftover: Vec<_> = std::fs::read_dir(storage.root())
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.file_name())
.filter(|n| n.to_string_lossy().contains(".tmp-"))
.collect();
assert!(leftover.is_empty(), "temp files left behind: {leftover:?}");
}
} }
+10 -3
View File
@@ -52,7 +52,8 @@ impl S3Storage {
access_key: &str, access_key: &str,
secret_key: &str, secret_key: &str,
) -> Self { ) -> Self {
let credentials = Credentials::new(access_key, secret_key, None, None, "mytheclipse-storage"); let credentials =
Credentials::new(access_key, secret_key, None, None, "mytheclipse-storage");
let config = aws_sdk_s3::Config::builder() let config = aws_sdk_s3::Config::builder()
.region(Region::new(region.to_string())) .region(Region::new(region.to_string()))
.endpoint_url(endpoint) .endpoint_url(endpoint)
@@ -322,13 +323,19 @@ mod tests {
.await; .await;
storage storage
.put("mytheclipse-storage-test.txt", bytes_stream(b"hello s3".to_vec())) .put(
"mytheclipse-storage-test.txt",
bytes_stream(b"hello s3".to_vec()),
)
.await .await
.unwrap(); .unwrap();
let data = read_to_vec(storage.get("mytheclipse-storage-test.txt").await.unwrap()) let data = read_to_vec(storage.get("mytheclipse-storage-test.txt").await.unwrap())
.await .await
.unwrap(); .unwrap();
assert_eq!(data, b"hello s3"); assert_eq!(data, b"hello s3");
storage.delete("mytheclipse-storage-test.txt").await.unwrap(); storage
.delete("mytheclipse-storage-test.txt")
.await
.unwrap();
} }
} }
+2 -2
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "mytheclipse" name = "mytheclipse"
version = "1.3.0" version = "1.3.4"
edition = "2021" edition = "2021"
rust-version = "1.75" rust-version = "1.75"
license = "MIT OR Apache-2.0" license = "MIT OR Apache-2.0"
@@ -10,7 +10,7 @@ documentation = "https://docs.rs/mytheclipse"
authors = ["asepharyana <superaseph@gmail.com>"] authors = ["asepharyana <superaseph@gmail.com>"]
description = "Resource-aware abstractions for async I/O, heavy compute, background queue management, resiliency, traffic control, lifecycle, and observability." description = "Resource-aware abstractions for async I/O, heavy compute, background queue management, resiliency, traffic control, lifecycle, and observability."
readme = "README.md" readme = "README.md"
keywords = ["async", "concurrency", "rayon", "tokio", "resource-management", "resiliency", "retry", "circuit-breaker", "rate-limit", "observability", "cron", "shutdown"] keywords = ["async", "concurrency", "resiliency", "circuit-breaker", "observability"]
categories = ["asynchronous", "concurrency", "rust-patterns"] categories = ["asynchronous", "concurrency", "rust-patterns"]
[dependencies] [dependencies]