Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8ff33817a8 | ||
|
|
b6f138b90d | ||
|
|
4027d3eb17 | ||
|
|
5f1e3ace5c | ||
|
|
0f2b776bb8 | ||
|
|
5717f8aaaa | ||
|
|
52d9e1e93e | ||
|
|
19941156b4 | ||
|
|
b2d3f32d83 | ||
|
|
4d851c5a58 | ||
|
|
db4f277336 |
+28
-28
@@ -60,33 +60,33 @@ jobs:
|
|||||||
flags: "-p mytheclipse --no-default-features --features lifecycle"
|
flags: "-p mytheclipse --no-default-features --features lifecycle"
|
||||||
- name: mytheclipse / observability only
|
- name: mytheclipse / observability only
|
||||||
flags: "-p mytheclipse --no-default-features --features observability"
|
flags: "-p mytheclipse --no-default-features --features observability"
|
||||||
# corex-cache
|
# mytheclipse-cache
|
||||||
- name: corex-cache / default
|
- name: mytheclipse-cache / default
|
||||||
flags: "-p corex-cache"
|
flags: "-p mytheclipse-cache"
|
||||||
- name: corex-cache / l1-moka
|
- name: mytheclipse-cache / l1-moka
|
||||||
flags: "-p corex-cache --no-default-features --features l1-moka"
|
flags: "-p mytheclipse-cache --no-default-features --features l1-moka"
|
||||||
- name: corex-cache / l2-redis
|
- name: mytheclipse-cache / l2-redis
|
||||||
flags: "-p corex-cache --no-default-features --features l1-memory,l2-redis"
|
flags: "-p mytheclipse-cache --no-default-features --features l1-memory,l2-redis"
|
||||||
# corex-storage
|
# mytheclipse-storage
|
||||||
- name: corex-storage / default (local)
|
- name: mytheclipse-storage / default (local)
|
||||||
flags: "-p corex-storage"
|
flags: "-p mytheclipse-storage"
|
||||||
- name: corex-storage / s3
|
- name: mytheclipse-storage / s3
|
||||||
flags: "-p corex-storage --no-default-features --features s3"
|
flags: "-p mytheclipse-storage --no-default-features --features s3"
|
||||||
- name: corex-storage / gcs
|
- name: mytheclipse-storage / gcs
|
||||||
flags: "-p corex-storage --no-default-features --features gcs"
|
flags: "-p mytheclipse-storage --no-default-features --features gcs"
|
||||||
# corex-event
|
# mytheclipse-event
|
||||||
- name: corex-event / default (mem)
|
- name: mytheclipse-event / default (mem)
|
||||||
flags: "-p corex-event"
|
flags: "-p mytheclipse-event"
|
||||||
- name: corex-event / amqp
|
- name: mytheclipse-event / amqp
|
||||||
flags: "-p corex-event --no-default-features --features amqp"
|
flags: "-p mytheclipse-event --no-default-features --features amqp"
|
||||||
- name: corex-event / nats
|
- name: mytheclipse-event / nats
|
||||||
flags: "-p corex-event --no-default-features --features nats"
|
flags: "-p mytheclipse-event --no-default-features --features nats"
|
||||||
# corex-config
|
# mytheclipse-config
|
||||||
- name: corex-config / default
|
- name: mytheclipse-config / default
|
||||||
flags: "-p corex-config"
|
flags: "-p mytheclipse-config"
|
||||||
# corex-crypto
|
# mytheclipse-crypto
|
||||||
- name: corex-crypto / default
|
- name: mytheclipse-crypto / default
|
||||||
flags: "-p corex-crypto"
|
flags: "-p mytheclipse-crypto"
|
||||||
steps:
|
steps:
|
||||||
- uses: actions/checkout@v4
|
- uses: actions/checkout@v4
|
||||||
- uses: dtolnay/rust-toolchain@stable
|
- uses: dtolnay/rust-toolchain@stable
|
||||||
@@ -118,7 +118,7 @@ jobs:
|
|||||||
fail-fast: false
|
fail-fast: false
|
||||||
matrix:
|
matrix:
|
||||||
crate:
|
crate:
|
||||||
[mytheclipse, corex-cache, corex-storage, corex-event, corex-config, corex-crypto]
|
[mytheclipse, mytheclipse-cache, mytheclipse-storage, mytheclipse-event, mytheclipse-config, mytheclipse-crypto]
|
||||||
steps:
|
steps:
|
||||||
- uses: actions/checkout@v4
|
- uses: actions/checkout@v4
|
||||||
- uses: dtolnay/rust-toolchain@stable
|
- uses: dtolnay/rust-toolchain@stable
|
||||||
|
|||||||
@@ -30,7 +30,7 @@ jobs:
|
|||||||
# between publishes avoids hitting crates.io's rate limit.
|
# between publishes avoids hitting crates.io's rate limit.
|
||||||
- name: Publish workspace crates
|
- name: Publish workspace crates
|
||||||
run: |
|
run: |
|
||||||
for crate in mytheclipse corex-cache corex-storage corex-event corex-config corex-crypto; do
|
for crate in mytheclipse mytheclipse-cache mytheclipse-storage mytheclipse-event mytheclipse-config mytheclipse-crypto; do
|
||||||
echo "Publishing $crate..."
|
echo "Publishing $crate..."
|
||||||
cargo publish -p "$crate" --allow-dirty --no-verify
|
cargo publish -p "$crate" --allow-dirty --no-verify
|
||||||
sleep 15
|
sleep 15
|
||||||
|
|||||||
@@ -1,3 +1,38 @@
|
|||||||
|
## [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)
|
||||||
|
|
||||||
|
|
||||||
|
### Features
|
||||||
|
|
||||||
|
* add corex-storage crate for unified storage abstraction ([db4f277](https://github.com/asepharyana/mytheclipse/commit/db4f277336d1e32cb6a2ddd86ac37ae9789fa4f8))
|
||||||
|
|
||||||
# [1.2.0](https://github.com/asepharyana/mytheclipse/compare/v1.1.0...v1.2.0) (2026-08-28)
|
# [1.2.0](https://github.com/asepharyana/mytheclipse/compare/v1.1.0...v1.2.0) (2026-08-28)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Generated
+4471
-15
File diff suppressed because it is too large
Load Diff
+10
-45
@@ -1,45 +1,10 @@
|
|||||||
[package]
|
[workspace]
|
||||||
name = "mytheclipse"
|
members = [
|
||||||
version = "1.1.0"
|
"crates/mytheclipse",
|
||||||
edition = "2021"
|
"crates/mytheclipse-cache",
|
||||||
rust-version = "1.75"
|
"crates/mytheclipse-storage",
|
||||||
license = "MIT OR Apache-2.0"
|
"crates/mytheclipse-event",
|
||||||
repository = "https://github.com/asepharyana/mytheclipse"
|
"crates/mytheclipse-config",
|
||||||
homepage = "https://github.com/asepharyana/mytheclipse"
|
"crates/mytheclipse-crypto",
|
||||||
documentation = "https://docs.rs/mytheclipse"
|
]
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
resolver = "2"
|
||||||
description = "Resource-aware abstractions for async I/O, heavy compute, background queue management, resiliency, traffic control, lifecycle, and observability."
|
|
||||||
readme = "README.md"
|
|
||||||
keywords = ["async", "concurrency", "rayon", "tokio", "resource-management", "resiliency", "retry", "circuit-breaker", "rate-limit", "observability", "cron", "shutdown"]
|
|
||||||
categories = ["asynchronous", "concurrency", "rust-patterns"]
|
|
||||||
|
|
||||||
[dependencies]
|
|
||||||
tokio = { version = "1.53", features = ["full"], optional = true }
|
|
||||||
rayon = { version = "1.12", optional = true }
|
|
||||||
rand = { version = "0.8", optional = true }
|
|
||||||
num_cpus = "1.17"
|
|
||||||
tracing = "0.1"
|
|
||||||
|
|
||||||
[dev-dependencies]
|
|
||||||
tokio = { version = "1.53", features = ["full"] }
|
|
||||||
tracing-subscriber = "0.3"
|
|
||||||
|
|
||||||
[features]
|
|
||||||
default = []
|
|
||||||
io = ["dep:tokio"]
|
|
||||||
compute = ["dep:rayon"]
|
|
||||||
bg = ["dep:tokio"]
|
|
||||||
resiliency = ["dep:tokio", "dep:rand"]
|
|
||||||
traffic = ["dep:tokio"]
|
|
||||||
lifecycle = ["dep:tokio"]
|
|
||||||
observability = ["dep:tokio"]
|
|
||||||
full = ["io", "compute", "bg", "resiliency", "traffic", "lifecycle", "observability"]
|
|
||||||
|
|
||||||
[[example]]
|
|
||||||
name = "main"
|
|
||||||
path = "examples/main.rs"
|
|
||||||
required-features = ["full"]
|
|
||||||
|
|
||||||
[package.metadata.docs.rs]
|
|
||||||
all-features = true
|
|
||||||
rustdoc-args = ["--cfg", "docsrs"]
|
|
||||||
@@ -12,11 +12,11 @@ concern.
|
|||||||
| Crate | Description | Docs |
|
| Crate | Description | Docs |
|
||||||
| :--- | :--- | :--- |
|
| :--- | :--- | :--- |
|
||||||
| [`mytheclipse`](crates/mytheclipse) | Resource-aware execution primitives (async I/O, compute, background queues), resiliency (retry, circuit breaker, timeout), traffic control (rate limiter, backpressure, concurrency limiter), lifecycle (graceful shutdown, cron), and observability (metrics, panic tracking). | [README](crates/mytheclipse/README.md) |
|
| [`mytheclipse`](crates/mytheclipse) | Resource-aware execution primitives (async I/O, compute, background queues), resiliency (retry, circuit breaker, timeout), traffic control (rate limiter, backpressure, concurrency limiter), lifecycle (graceful shutdown, cron), and observability (metrics, panic tracking). | [README](crates/mytheclipse/README.md) |
|
||||||
| [`corex-cache`](crates/corex-cache) | Unified multi-layer (L1/L2) cache abstraction: in-memory or Moka L1, Redis/Valkey L2, cache-aside read-through. | [README](crates/corex-cache/README.md) |
|
| [`mytheclipse-cache`](crates/mytheclipse-cache) | Unified multi-layer (L1/L2) cache abstraction: in-memory or Moka L1, Redis/Valkey L2, cache-aside read-through. | [README](crates/mytheclipse-cache/README.md) |
|
||||||
| [`corex-storage`](crates/corex-storage) | Unified storage & file system abstraction: one driver interface over local disk, S3/MinIO, and Google Cloud Storage, stream-based. | [README](crates/corex-storage/README.md) |
|
| [`mytheclipse-storage`](crates/mytheclipse-storage) | Unified storage & file system abstraction: one driver interface over local disk, S3/MinIO, and Google Cloud Storage, stream-based. | [README](crates/mytheclipse-storage/README.md) |
|
||||||
| [`corex-event`](crates/corex-event) | Unified events & message bus abstraction: in-memory pub/sub dispatcher plus RabbitMQ and NATS broker adapters behind one trait. | [README](crates/corex-event/README.md) |
|
| [`mytheclipse-event`](crates/mytheclipse-event) | Unified events & message bus abstraction: in-memory pub/sub dispatcher plus RabbitMQ and NATS broker adapters behind one trait. | [README](crates/mytheclipse-event/README.md) |
|
||||||
| [`corex-config`](crates/corex-config) | Type-safe, dynamic configuration engine: load `.env`/YAML/JSON/TOML into typed structs, with hot-reload. | [README](crates/corex-config/README.md) |
|
| [`mytheclipse-config`](crates/mytheclipse-config) | Type-safe, dynamic configuration engine: load `.env`/YAML/JSON/TOML into typed structs, with hot-reload. | [README](crates/mytheclipse-config/README.md) |
|
||||||
| [`corex-crypto`](crates/corex-crypto) | Safe hashing (Argon2id), encryption (AES-256-GCM), and JWT tokens, with key rotation support. | [README](crates/corex-crypto/README.md) |
|
| [`mytheclipse-crypto`](crates/mytheclipse-crypto) | Safe hashing (Argon2id), encryption (AES-256-GCM), and JWT tokens, with key rotation support. | [README](crates/mytheclipse-crypto/README.md) |
|
||||||
|
|
||||||
Every crate follows the same philosophy: **one small interface, pluggable
|
Every crate follows the same philosophy: **one small interface, pluggable
|
||||||
backends behind feature flags, and a working default that needs no external
|
backends behind feature flags, and a working default that needs no external
|
||||||
@@ -31,11 +31,11 @@ Each crate is published independently; add the ones you need:
|
|||||||
```toml
|
```toml
|
||||||
[dependencies]
|
[dependencies]
|
||||||
mytheclipse = { version = "1", features = ["full"] }
|
mytheclipse = { version = "1", features = ["full"] }
|
||||||
corex-cache = "0.1"
|
mytheclipse-cache = "0.1"
|
||||||
corex-storage = { version = "0.1", features = ["s3"] }
|
mytheclipse-storage = { version = "0.1", features = ["s3"] }
|
||||||
corex-event = { version = "0.1", features = ["nats"] }
|
mytheclipse-event = { version = "0.1", features = ["nats"] }
|
||||||
corex-config = "0.1"
|
mytheclipse-config = "0.1"
|
||||||
corex-crypto = "0.1"
|
mytheclipse-crypto = "0.1"
|
||||||
```
|
```
|
||||||
|
|
||||||
See each crate's own README (linked above) for usage examples and the full
|
See each crate's own README (linked above) for usage examples and the full
|
||||||
@@ -52,7 +52,7 @@ cargo clippy --workspace --all-features -- -D warnings
|
|||||||
cargo fmt --all --check
|
cargo fmt --all --check
|
||||||
```
|
```
|
||||||
|
|
||||||
Or target a single crate with `-p <name>`, e.g. `cargo test -p corex-cache`.
|
Or target a single crate with `-p <name>`, e.g. `cargo test -p mytheclipse-cache`.
|
||||||
|
|
||||||
## License
|
## License
|
||||||
|
|
||||||
|
|||||||
@@ -1,12 +1,12 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "corex-cache"
|
name = "mytheclipse-cache"
|
||||||
version = "1.2.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"
|
||||||
repository = "https://github.com/asepharyana/mytheclipse"
|
repository = "https://github.com/asepharyana/mytheclipse"
|
||||||
homepage = "https://github.com/asepharyana/mytheclipse"
|
homepage = "https://github.com/asepharyana/mytheclipse"
|
||||||
documentation = "https://docs.rs/corex-cache"
|
documentation = "https://docs.rs/mytheclipse-cache"
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
authors = ["asepharyana <superaseph@gmail.com>"]
|
||||||
description = "Unified multi-layer cache abstraction: L1/L2 caching, cache-aside and auto-refresh, with pluggable backends."
|
description = "Unified multi-layer cache abstraction: L1/L2 caching, cache-aside and auto-refresh, with pluggable backends."
|
||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
# corex-cache
|
# mytheclipse-cache
|
||||||
|
|
||||||
A unified multi-layer cache abstraction so your app isn't locked to one cache
|
A unified multi-layer cache abstraction so your app isn't locked to one cache
|
||||||
provider. Combines an in-process **L1** cache with a distributed **L2** cache
|
provider. Combines an in-process **L1** cache with a distributed **L2** cache
|
||||||
@@ -15,7 +15,7 @@ provider. Combines an in-process **L1** cache with a distributed **L2** cache
|
|||||||
## Usage
|
## Usage
|
||||||
|
|
||||||
```rust
|
```rust
|
||||||
use corex_cache::{Cache, MemoryCache, MultiLayerCache, CacheAside};
|
use mytheclipse_cache::{Cache, MemoryCache, MultiLayerCache, CacheAside};
|
||||||
|
|
||||||
let cache = MultiLayerCache::new(
|
let cache = MultiLayerCache::new(
|
||||||
MemoryCache::new(), // L1
|
MemoryCache::new(), // L1
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
//! # corex-cache
|
//! # mytheclipse-cache
|
||||||
//!
|
//!
|
||||||
//! A unified multi-layer cache abstraction that keeps your application from
|
//! A unified multi-layer cache abstraction that keeps your application from
|
||||||
//! being locked to any single cache provider.
|
//! being locked to any single cache provider.
|
||||||
@@ -17,9 +17,12 @@
|
|||||||
//!
|
//!
|
||||||
//! ## Example
|
//! ## Example
|
||||||
//!
|
//!
|
||||||
|
//! Multi-layer + cache-aside composition (default features):
|
||||||
|
//!
|
||||||
//! ```no_run
|
//! ```no_run
|
||||||
//! use corex_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)]
|
||||||
@@ -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() {
|
||||||
@@ -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.
|
||||||
@@ -101,7 +101,7 @@ mod tests {
|
|||||||
|
|
||||||
/// Integration test requiring a live Redis at `REDIS_URL`
|
/// Integration test requiring a live Redis at `REDIS_URL`
|
||||||
/// (e.g. `redis://127.0.0.1:6379`). Run with:
|
/// (e.g. `redis://127.0.0.1:6379`). Run with:
|
||||||
/// `REDIS_URL=redis://127.0.0.1:6379 cargo test -p corex-cache --features l2-redis -- --ignored` .
|
/// `REDIS_URL=redis://127.0.0.1:6379 cargo test -p mytheclipse-cache --features l2-redis -- --ignored` .
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires a live Redis instance (REDIS_URL)"]
|
#[ignore = "requires a live Redis instance (REDIS_URL)"]
|
||||||
async fn set_get_roundtrip_live() {
|
async fn set_get_roundtrip_live() {
|
||||||
@@ -111,7 +111,7 @@ mod tests {
|
|||||||
.get_multiplexed_tokio_connection()
|
.get_multiplexed_tokio_connection()
|
||||||
.await
|
.await
|
||||||
.expect("connect");
|
.expect("connect");
|
||||||
let cache = RedisCache::with_prefix(conn, "corex_cache_test:".to_string());
|
let cache = RedisCache::with_prefix(conn, "mytheclipse_cache_test:".to_string());
|
||||||
|
|
||||||
cache
|
cache
|
||||||
.set("k", b"v".to_vec(), Some(Duration::from_secs(3600)))
|
.set("k", b"v".to_vec(), Some(Duration::from_secs(3600)))
|
||||||
@@ -1,12 +1,12 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "corex-config"
|
name = "mytheclipse-config"
|
||||||
version = "1.2.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"
|
||||||
repository = "https://github.com/asepharyana/mytheclipse"
|
repository = "https://github.com/asepharyana/mytheclipse"
|
||||||
homepage = "https://github.com/asepharyana/mytheclipse"
|
homepage = "https://github.com/asepharyana/mytheclipse"
|
||||||
documentation = "https://docs.rs/corex-config"
|
documentation = "https://docs.rs/mytheclipse-config"
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
authors = ["asepharyana <superaseph@gmail.com>"]
|
||||||
description = "Type-safe, dynamic configuration engine: load .env/YAML/JSON/TOML into typed structs with hot-reload and validation."
|
description = "Type-safe, dynamic configuration engine: load .env/YAML/JSON/TOML into typed structs with hot-reload and validation."
|
||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
# corex-config
|
# mytheclipse-config
|
||||||
|
|
||||||
Type-safe, dynamic configuration: load `.env`, YAML, JSON, or TOML directly
|
Type-safe, dynamic configuration: load `.env`, YAML, JSON, or TOML directly
|
||||||
into a typed Rust struct, merge multiple sources with environment variables
|
into a typed Rust struct, merge multiple sources with environment variables
|
||||||
@@ -14,7 +14,7 @@ taking priority, and optionally hot-reload when the source files change.
|
|||||||
|
|
||||||
```rust
|
```rust
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
use corex_config::ConfigLoader;
|
use mytheclipse_config::ConfigLoader;
|
||||||
|
|
||||||
#[derive(Debug, Deserialize, Clone)]
|
#[derive(Debug, Deserialize, Clone)]
|
||||||
struct AppConfig {
|
struct AppConfig {
|
||||||
@@ -31,13 +31,13 @@ let config: AppConfig = ConfigLoader::new()
|
|||||||
### Hot-reload
|
### Hot-reload
|
||||||
|
|
||||||
```rust
|
```rust
|
||||||
use corex_config::DynamicConfig;
|
use mytheclipse_config::DynamicConfig;
|
||||||
# use serde::Deserialize;
|
# use serde::Deserialize;
|
||||||
# #[derive(Debug, Deserialize, Clone)] struct AppConfig { port: u16 }
|
# #[derive(Debug, Deserialize, Clone)] struct AppConfig { port: u16 }
|
||||||
|
|
||||||
let cfg = DynamicConfig::<AppConfig>::watch_files(
|
let cfg = DynamicConfig::<AppConfig>::watch_files(
|
||||||
vec!["config.yaml".into()],
|
vec!["config.yaml".into()],
|
||||||
|| corex_config::ConfigLoader::new().merge_file("config.yaml".as_ref())?.build(),
|
|| mytheclipse_config::ConfigLoader::new().merge_file("config.yaml".as_ref())?.build(),
|
||||||
)?;
|
)?;
|
||||||
|
|
||||||
let mut changes = cfg.subscribe();
|
let mut changes = cfg.subscribe();
|
||||||
@@ -43,13 +43,16 @@ impl<T: Config + Clone> DynamicConfig<T> {
|
|||||||
pub fn get(&self) -> T {
|
pub fn get(&self) -> T {
|
||||||
self.inner
|
self.inner
|
||||||
.read()
|
.read()
|
||||||
.expect("corex-config: RwLock poisoned")
|
.expect("mytheclipse-config: RwLock poisoned")
|
||||||
.clone()
|
.clone()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// 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("corex-config: RwLock poisoned") = new;
|
*self
|
||||||
|
.inner
|
||||||
|
.write()
|
||||||
|
.expect("mytheclipse-config: RwLock poisoned") = new;
|
||||||
let _ = self.tx.send(());
|
let _ = self.tx.send(());
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -101,7 +104,7 @@ impl<T: Config + Clone> DynamicConfig<T> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
std::thread::Builder::new()
|
std::thread::Builder::new()
|
||||||
.name("corex-config-watch".into())
|
.name("mytheclipse-config-watch".into())
|
||||||
.spawn(move || {
|
.spawn(move || {
|
||||||
// Keep the watcher alive for the life of this thread.
|
// Keep the watcher alive for the life of this thread.
|
||||||
let _watcher = watcher;
|
let _watcher = watcher;
|
||||||
@@ -115,12 +118,12 @@ impl<T: Config + Clone> DynamicConfig<T> {
|
|||||||
}
|
}
|
||||||
match reload() {
|
match reload() {
|
||||||
Ok(new) => {
|
Ok(new) => {
|
||||||
*inner.write().expect("corex-config: RwLock poisoned") = new;
|
*inner.write().expect("mytheclipse-config: RwLock poisoned") = new;
|
||||||
let _ = tx.send(());
|
let _ = tx.send(());
|
||||||
last_applied = Instant::now();
|
last_applied = Instant::now();
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
tracing::error!("corex-config: hot-reload failed: {e}");
|
tracing::error!("mytheclipse-config: hot-reload failed: {e}");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
//! Shared error type for corex-config.
|
//! Shared error type for mytheclipse-config.
|
||||||
|
|
||||||
/// Errors surfaced while loading or reloading configuration.
|
/// Errors surfaced while loading or reloading configuration.
|
||||||
#[non_exhaustive]
|
#[non_exhaustive]
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
//! # corex-config
|
//! # mytheclipse-config
|
||||||
//!
|
//!
|
||||||
//! A type-safe, dynamic configuration engine: load environment variables,
|
//! A type-safe, dynamic configuration engine: load environment variables,
|
||||||
//! `.env` files, YAML, JSON, or TOML directly into a typed Rust struct, with
|
//! `.env` files, YAML, JSON, or TOML directly into a typed Rust struct, with
|
||||||
@@ -14,7 +14,7 @@
|
|||||||
//!
|
//!
|
||||||
//! ```no_run
|
//! ```no_run
|
||||||
//! use serde::Deserialize;
|
//! use serde::Deserialize;
|
||||||
//! use corex_config::ConfigLoader;
|
//! use mytheclipse_config::ConfigLoader;
|
||||||
//!
|
//!
|
||||||
//! #[derive(Debug, Deserialize, Clone)]
|
//! #[derive(Debug, Deserialize, Clone)]
|
||||||
//! struct AppConfig {
|
//! struct AppConfig {
|
||||||
@@ -1,12 +1,12 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "corex-crypto"
|
name = "mytheclipse-crypto"
|
||||||
version = "1.2.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"
|
||||||
repository = "https://github.com/asepharyana/mytheclipse"
|
repository = "https://github.com/asepharyana/mytheclipse"
|
||||||
homepage = "https://github.com/asepharyana/mytheclipse"
|
homepage = "https://github.com/asepharyana/mytheclipse"
|
||||||
documentation = "https://docs.rs/corex-crypto"
|
documentation = "https://docs.rs/mytheclipse-crypto"
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
authors = ["asepharyana <superaseph@gmail.com>"]
|
||||||
description = "Safe hashing, encryption, and token helpers (Argon2id, AES-256-GCM, JWT/Paseto) with key rotation support."
|
description = "Safe hashing, encryption, and token helpers (Argon2id, AES-256-GCM, JWT/Paseto) with key rotation support."
|
||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
# corex-crypto
|
# mytheclipse-crypto
|
||||||
|
|
||||||
Safe, one-line security helpers that are easy to get wrong when hand-rolled:
|
Safe, one-line security helpers that are easy to get wrong when hand-rolled:
|
||||||
|
|
||||||
@@ -20,7 +20,7 @@ independent.
|
|||||||
## Usage
|
## Usage
|
||||||
|
|
||||||
```rust
|
```rust
|
||||||
use corex_crypto::{PasswordHasher, Encryptor, TokenSigner};
|
use mytheclipse_crypto::{PasswordHasher, Encryptor, TokenSigner};
|
||||||
|
|
||||||
let hasher = PasswordHasher::new();
|
let hasher = PasswordHasher::new();
|
||||||
let hash = hasher.hash("letmein").unwrap();
|
let hash = hasher.hash("letmein").unwrap();
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
//! # corex-crypto
|
//! # mytheclipse-crypto
|
||||||
//!
|
//!
|
||||||
//! Low-level security helpers that are easy to get wrong when hand-rolled:
|
//! Low-level security helpers that are easy to get wrong when hand-rolled:
|
||||||
//!
|
//!
|
||||||
@@ -18,7 +18,7 @@
|
|||||||
//! ## Example
|
//! ## Example
|
||||||
//!
|
//!
|
||||||
//! ```no_run
|
//! ```no_run
|
||||||
//! use corex_crypto::{PasswordHasher, Encryptor, TokenSigner};
|
//! use mytheclipse_crypto::{PasswordHasher, Encryptor, TokenSigner};
|
||||||
//!
|
//!
|
||||||
//! // Hash & verify a password.
|
//! // Hash & verify a password.
|
||||||
//! let hasher = PasswordHasher::new();
|
//! let hasher = PasswordHasher::new();
|
||||||
@@ -63,7 +63,7 @@ pub use token::{Claims, TokenError, TokenSigner};
|
|||||||
|
|
||||||
pub use key_ring::KeyRing;
|
pub use key_ring::KeyRing;
|
||||||
|
|
||||||
/// Errors returned across corex-crypto primitives.
|
/// Errors returned across mytheclipse-crypto primitives.
|
||||||
#[non_exhaustive]
|
#[non_exhaustive]
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub enum CryptoError {
|
pub enum CryptoError {
|
||||||
@@ -1,12 +1,12 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "corex-event"
|
name = "mytheclipse-event"
|
||||||
version = "1.2.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"
|
||||||
repository = "https://github.com/asepharyana/mytheclipse"
|
repository = "https://github.com/asepharyana/mytheclipse"
|
||||||
homepage = "https://github.com/asepharyana/mytheclipse"
|
homepage = "https://github.com/asepharyana/mytheclipse"
|
||||||
documentation = "https://docs.rs/corex-event"
|
documentation = "https://docs.rs/mytheclipse-event"
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
authors = ["asepharyana <superaseph@gmail.com>"]
|
||||||
description = "Unified events & message bus abstraction: in-memory pub/sub dispatcher plus pluggable distributed broker backends (RabbitMQ, NATS)."
|
description = "Unified events & message bus abstraction: in-memory pub/sub dispatcher plus pluggable distributed broker backends (RabbitMQ, NATS)."
|
||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
# corex-event
|
# mytheclipse-event
|
||||||
|
|
||||||
A unified events & message bus abstraction so component-to-component (or
|
A unified events & message bus abstraction so component-to-component (or
|
||||||
service-to-service) communication isn't locked to one transport.
|
service-to-service) communication isn't locked to one transport.
|
||||||
@@ -19,7 +19,7 @@ service-to-service) communication isn't locked to one transport.
|
|||||||
## Usage
|
## Usage
|
||||||
|
|
||||||
```rust
|
```rust
|
||||||
use corex_event::{InMemoryEventBus, TypedEventBus};
|
use mytheclipse_event::{InMemoryEventBus, TypedEventBus};
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize, PartialEq)]
|
#[derive(Debug, Serialize, Deserialize, PartialEq)]
|
||||||
@@ -2,7 +2,7 @@
|
|||||||
//! AMQP 0.9.1 client.
|
//! AMQP 0.9.1 client.
|
||||||
//!
|
//!
|
||||||
//! Topics map to routing keys on a single topic [`lapin::ExchangeKind::Topic`]
|
//! Topics map to routing keys on a single topic [`lapin::ExchangeKind::Topic`]
|
||||||
//! exchange (default name `corex.events`); each subscription declares its own
|
//! exchange (default name `mytheclipse.events`); each subscription declares its own
|
||||||
//! exclusive, auto-delete queue bound to that routing key, matching the
|
//! exclusive, auto-delete queue bound to that routing key, matching the
|
||||||
//! common "fanout via topic exchange" pattern.
|
//! common "fanout via topic exchange" pattern.
|
||||||
|
|
||||||
@@ -107,7 +107,7 @@ impl EventBus for AmqpEventBus {
|
|||||||
.channel
|
.channel
|
||||||
.basic_consume(
|
.basic_consume(
|
||||||
queue.name().as_str(),
|
queue.name().as_str(),
|
||||||
"corex-event-consumer",
|
"mytheclipse-event-consumer",
|
||||||
BasicConsumeOptions::default(),
|
BasicConsumeOptions::default(),
|
||||||
FieldTable::default(),
|
FieldTable::default(),
|
||||||
)
|
)
|
||||||
@@ -145,12 +145,12 @@ mod tests {
|
|||||||
|
|
||||||
/// Requires a live RabbitMQ at `AMQP_URL` (e.g.
|
/// Requires a live RabbitMQ at `AMQP_URL` (e.g.
|
||||||
/// `amqp://guest:guest@127.0.0.1:5672/%2f`). Run with:
|
/// `amqp://guest:guest@127.0.0.1:5672/%2f`). Run with:
|
||||||
/// `AMQP_URL=... cargo test -p corex-event --features amqp -- --ignored`.
|
/// `AMQP_URL=... cargo test -p mytheclipse-event --features amqp -- --ignored`.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires a live RabbitMQ instance (AMQP_URL)"]
|
#[ignore = "requires a live RabbitMQ instance (AMQP_URL)"]
|
||||||
async fn publish_subscribe_roundtrip_live() {
|
async fn publish_subscribe_roundtrip_live() {
|
||||||
let url = std::env::var("AMQP_URL").expect("set AMQP_URL");
|
let url = std::env::var("AMQP_URL").expect("set AMQP_URL");
|
||||||
let bus = AmqpEventBus::connect(&url, "corex_event_test")
|
let bus = AmqpEventBus::connect(&url, "mytheclipse_event_test")
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let mut sub = bus.subscribe("orders.created").await.unwrap();
|
let mut sub = bus.subscribe("orders.created").await.unwrap();
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
//! # corex-event
|
//! # mytheclipse-event
|
||||||
//!
|
//!
|
||||||
//! A unified events & message bus abstraction so component-to-component (or
|
//! A unified events & message bus abstraction so component-to-component (or
|
||||||
//! service-to-service) communication doesn't get locked to one transport.
|
//! service-to-service) communication doesn't get locked to one transport.
|
||||||
@@ -15,20 +15,27 @@
|
|||||||
//!
|
//!
|
||||||
//! ## Example
|
//! ## Example
|
||||||
//!
|
//!
|
||||||
//! ```
|
//! The in-memory + typed bus (`mem` feature, on by default):
|
||||||
//! use corex_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;
|
||||||
@@ -77,7 +77,7 @@ mod tests {
|
|||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
/// Requires a live NATS server at `NATS_URL` (e.g. `nats://127.0.0.1:4222`).
|
/// Requires a live NATS server at `NATS_URL` (e.g. `nats://127.0.0.1:4222`).
|
||||||
/// Run with: `NATS_URL=... cargo test -p corex-event --features nats -- --ignored`.
|
/// Run with: `NATS_URL=... cargo test -p mytheclipse-event --features nats -- --ignored`.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires a live NATS instance (NATS_URL)"]
|
#[ignore = "requires a live NATS instance (NATS_URL)"]
|
||||||
async fn publish_subscribe_roundtrip_live() {
|
async fn publish_subscribe_roundtrip_live() {
|
||||||
@@ -1,12 +1,12 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "corex-storage"
|
name = "mytheclipse-storage"
|
||||||
version = "1.2.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"
|
||||||
repository = "https://github.com/asepharyana/mytheclipse"
|
repository = "https://github.com/asepharyana/mytheclipse"
|
||||||
homepage = "https://github.com/asepharyana/mytheclipse"
|
homepage = "https://github.com/asepharyana/mytheclipse"
|
||||||
documentation = "https://docs.rs/corex-storage"
|
documentation = "https://docs.rs/mytheclipse-storage"
|
||||||
authors = ["asepharyana <superaseph@gmail.com>"]
|
authors = ["asepharyana <superaseph@gmail.com>"]
|
||||||
description = "Unified storage & file system abstraction: one driver interface over Local Disk, S3/MinIO, and Google Cloud Storage, with stream-based upload/download."
|
description = "Unified storage & file system abstraction: one driver interface over Local Disk, S3/MinIO, and Google Cloud Storage, with stream-based upload/download."
|
||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
# corex-storage
|
# mytheclipse-storage
|
||||||
|
|
||||||
A unified storage & file system abstraction so file handling isn't locked to
|
A unified storage & file system abstraction so file handling isn't locked to
|
||||||
one physical location.
|
one physical location.
|
||||||
@@ -14,7 +14,7 @@ requires holding it entirely in memory.
|
|||||||
## Usage
|
## Usage
|
||||||
|
|
||||||
```rust
|
```rust
|
||||||
use corex_storage::{LocalFileStorage, StorageDriver, bytes_stream, read_to_vec};
|
use mytheclipse_storage::{LocalFileStorage, StorageDriver, bytes_stream, read_to_vec};
|
||||||
|
|
||||||
let storage = LocalFileStorage::new("/var/data");
|
let storage = LocalFileStorage::new("/var/data");
|
||||||
storage.put("uploads/report.csv", bytes_stream(data)).await?;
|
storage.put("uploads/report.csv", bytes_stream(data)).await?;
|
||||||
@@ -143,7 +143,7 @@ mod tests {
|
|||||||
|
|
||||||
/// Requires a live GCS bucket with Application Default Credentials
|
/// Requires a live GCS bucket with Application Default Credentials
|
||||||
/// configured, and `GCS_BUCKET` set. Run with:
|
/// configured, and `GCS_BUCKET` set. Run with:
|
||||||
/// `GCS_BUCKET=my-bucket cargo test -p corex-storage --features gcs -- --ignored`.
|
/// `GCS_BUCKET=my-bucket cargo test -p mytheclipse-storage --features gcs -- --ignored`.
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[ignore = "requires a live GCS bucket + Application Default Credentials"]
|
#[ignore = "requires a live GCS bucket + Application Default Credentials"]
|
||||||
async fn put_get_roundtrip_live() {
|
async fn put_get_roundtrip_live() {
|
||||||
@@ -151,15 +151,18 @@ mod tests {
|
|||||||
let storage = GcsStorage::connect(bucket).await.unwrap();
|
let storage = GcsStorage::connect(bucket).await.unwrap();
|
||||||
storage
|
storage
|
||||||
.put(
|
.put(
|
||||||
"corex-storage-test.txt",
|
"mytheclipse-storage-test.txt",
|
||||||
bytes_stream(b"hello gcs".to_vec()),
|
bytes_stream(b"hello gcs".to_vec()),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let data = read_to_vec(storage.get("corex-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 gcs");
|
assert_eq!(data, b"hello gcs");
|
||||||
storage.delete("corex-storage-test.txt").await.unwrap();
|
storage
|
||||||
|
.delete("mytheclipse-storage-test.txt")
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
//! # corex-storage
|
//! # mytheclipse-storage
|
||||||
//!
|
//!
|
||||||
//! A unified storage & file system abstraction so file handling doesn't get
|
//! A unified storage & file system abstraction so file handling doesn't get
|
||||||
//! locked to one physical storage location.
|
//! locked to one physical storage location.
|
||||||
@@ -15,15 +15,22 @@
|
|||||||
//!
|
//!
|
||||||
//! ## Example
|
//! ## Example
|
||||||
//!
|
//!
|
||||||
//! ```
|
//! Local disk storage (`local` feature, default):
|
||||||
//! use corex_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;
|
||||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user