Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "gardal"
version = "0.0.1-alpha.7"
version = "0.0.1-alpha.8"
edition = "2024"
license = "Apache-2.0 OR MIT"
authors = ["Ahmed Farghal <me@asoli.dev>"]
Expand Down
30 changes: 15 additions & 15 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,29 +25,29 @@ Add this to your `Cargo.toml`:

```toml
[dependencies]
gardal = "0.0.1-alpha.4"
gardal = "0.0.1-alpha.8"

# For async support
gardal = { version = "0.0.1-alpha.4", features = ["async"] }
gardal = { version = "0.0.1-alpha.8", features = ["async"] }

# For high-performance timing
gardal = { version = "0.0.1-alpha.4", features = ["quanta"] }
gardal = { version = "0.0.1-alpha.8", features = ["quanta"] }

# For high-resolution async timers
gardal = { version = "0.0.1-alpha.4", features = ["async", "tokio-hrtime"] }
gardal = { version = "0.0.1-alpha.8", features = ["async", "tokio-hrtime"] }
```

### Basic Usage

```rust
use gardal::{Limit, TokenBucket};
use gardal::{Limit, LocalTokenBucket, StdClock};
use nonzero_ext::nonzero;

// Create a token bucket: 10 tokens per second, burst of 20
let bucket = TokenBucket::new(Limit::per_second_and_burst(
let bucket = LocalTokenBucket::new(Limit::per_second_and_burst(
nonzero!(10u32),
nonzero!(20u32),
));
), StdClock);

// Consume 5 tokens
match bucket.consume(nonzero!(5u32)) {
Expand All @@ -61,15 +61,15 @@ match bucket.consume(nonzero!(5u32)) {
```rust
use futures::{StreamExt, stream};
use gardal::futures::StreamExt as GardalStreamExt;
use gardal::{Limit, TokenBucket, AtomicSharedStorage, QuantaClock};
use gardal::{Limit, SharedTokenBucket, QuantaClock};
use nonzero_ext::nonzero;

#[tokio::main]
async fn main() {
let limit = Limit::per_second(nonzero!(5u32));
let bucket = TokenBucket::<AtomicSharedStorage, _>::from_parts(
let bucket = SharedTokenBucket::new(
limit,
QuantaClock::default()
QuantaClock,
);

let mut stream = stream::iter(1..=100)
Expand Down Expand Up @@ -115,13 +115,13 @@ Choose the appropriate storage strategy for your use case:
- **`LocalStorage`**: Thread-local storage for single-threaded applications

```rust
use gardal::{TokenBucket, AtomicSharedStorage, Limit};
use gardal::{TokenBucket, StdClock, AtomicSharedStorage, Limit};
use nonzero_ext::nonzero;

// Explicitly specify storage type
let bucket = TokenBucket::<AtomicSharedStorage>::from_parts(
let bucket = TokenBucket::<AtomicSharedStorage, _>::new(
Limit::per_second(nonzero!(10u32)),
gardal::StdClock::default()
StdClock
);
```

Expand All @@ -141,13 +141,13 @@ For applications requiring precise timing in async contexts, enable the `tokio-h
```rust
use futures::{StreamExt, stream};
use gardal::futures::StreamExt as GardalStreamExt;
use gardal::{Limit, TokenBucket};
use gardal::{Limit, SharedTokenBucket, TokioClock};
use nonzero_ext::nonzero;

#[tokio::main]
async fn main() {
let limit = Limit::per_second(nonzero!(1000u32)); // High-frequency rate limiting
let bucket = TokenBucket::new(limit);
let bucket = SharedTokenBucket::new(limit, TokioClock);

let mut stream = stream::iter(1..=10000)
.throttle(bucket)
Expand Down
47 changes: 18 additions & 29 deletions benches/throughput.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,41 +9,38 @@ use gardal::{
use nonzero_ext::nonzero;

fn bench_consume(c: &mut Criterion) {
let clock = quanta::Clock::new();
let limit = Limit::per_second(nonzero!(10_000u32));
let _quanta_thread = quanta::Upkeep::new_with_clock(Duration::from_micros(10), clock.clone())
let _quanta_thread = quanta::Upkeep::new(Duration::from_micros(10))
.start()
.unwrap();
let clock = FastClock::new(clock);
let quanta_tb =
TokenBucket::<PaddedAtomicStorage, _>::from_parts(limit, QuantaClock::default());
let std_tb = TokenBucket::<PaddedAtomicStorage, _>::from_parts(limit, StdClock::default());
let fast_tb = TokenBucket::<LocalStorage, _>::from_parts(limit, clock.clone());
let fast_tb_padded = TokenBucket::<PaddedAtomicStorage, _>::from_parts(limit, clock.clone());
let quanta_tb = TokenBucket::<PaddedAtomicStorage, _>::with_datum(limit, QuantaClock);
let std_tb = TokenBucket::<PaddedAtomicStorage, _>::with_datum(limit, StdClock);
let fast_tb = TokenBucket::<LocalStorage, _>::with_datum(limit, FastClock);
let fast_tb_padded = TokenBucket::<PaddedAtomicStorage, _>::with_datum(limit, FastClock);
std::thread::sleep(Duration::from_secs(1));
let mut group = c.benchmark_group("tokenbucket");
group
.throughput(Throughput::Elements(1))
.sample_size(100)
.bench_function("consume-mock-clock-local-storage", |b| {
let clock = ManualClock::default();
let tb = TokenBucket::<LocalStorage, _>::from_parts(limit, &clock);
let tb = TokenBucket::<LocalStorage, _>::with_datum(limit, &clock);
clock.set(10.0);
b.iter(|| {
let _x = std::hint::black_box(tb.try_consume_one());
});
})
.bench_function("consume-mock-clock-atomic-storage", |b| {
let clock = ManualClock::default();
let tb = TokenBucket::<AtomicStorage, _>::from_parts(limit, &clock);
let tb = TokenBucket::<AtomicStorage, _>::with_datum(limit, &clock);
clock.set(10.0);
b.iter(|| {
tb.consume_one();
});
})
.bench_function("consume-mock-clock-padded-atomic-storage", |b| {
let clock = ManualClock::default();
let tb = TokenBucket::<PaddedAtomicStorage, _>::from_parts(limit, &clock);
let tb = TokenBucket::<PaddedAtomicStorage, _>::with_datum(limit, &clock);
clock.set(10.0);
b.iter(|| {
tb.consume_one();
Expand Down Expand Up @@ -75,19 +72,16 @@ fn bench_consume(c: &mut Criterion) {
const THREADS: u32 = 24;

fn multi_threaded(c: &mut Criterion) {
let clock = quanta::Clock::new();
let _quanta_thread = quanta::Upkeep::new_with_clock(Duration::from_micros(100), clock.clone())
let _quanta_thread = quanta::Upkeep::new(Duration::from_micros(100))
.start()
.unwrap();
let clock = FastClock::new(clock);
let limit = Limit::per_second(nonzero!(10_000u32));
let mut group = c.benchmark_group("multi_threaded");
group
.throughput(Throughput::Elements(1))
.bench_function("padded", |b| {
let tb = Arc::new(TokenBucket::<PaddedAtomicStorage, _>::from_parts(
limit,
clock.clone(),
let tb = Arc::new(TokenBucket::<PaddedAtomicStorage, _>::with_datum(
limit, FastClock,
));
b.iter_custom(|iters| {
let mut children = vec![];
Expand All @@ -108,9 +102,8 @@ fn multi_threaded(c: &mut Criterion) {
})
})
.bench_function("atomic", |b| {
let tb = Arc::new(TokenBucket::<AtomicStorage, _>::from_parts(
limit,
clock.clone(),
let tb = Arc::new(TokenBucket::<AtomicStorage, _>::with_datum(
limit, FastClock,
));
b.iter_custom(|iters| {
let mut children = vec![];
Expand All @@ -133,20 +126,17 @@ fn multi_threaded(c: &mut Criterion) {
}

fn multi_threaded2(c: &mut Criterion) {
let clock = quanta::Clock::new();
let _quanta_thread = quanta::Upkeep::new_with_clock(Duration::from_micros(10), clock.clone())
let _quanta_thread = quanta::Upkeep::new(Duration::from_micros(10))
.start()
.unwrap();
let clock = FastClock::new(clock);
let limit = Limit::per_second(nonzero!(50u32));
let mut group = c.benchmark_group("multi_threaded2");
group
.throughput(Throughput::Elements(1))
.bench_function("padded", |b| {
b.iter_custom(|iters| {
let tb = Arc::new(TokenBucket::<PaddedAtomicStorage, _>::from_parts(
limit,
clock.clone(),
let tb = Arc::new(TokenBucket::<PaddedAtomicStorage, _>::with_datum(
limit, FastClock,
));
let mut children = vec![];
let start = std::time::Instant::now();
Expand All @@ -166,9 +156,8 @@ fn multi_threaded2(c: &mut Criterion) {
})
.bench_function("atomic", |b| {
b.iter_custom(|iters| {
let tb = Arc::new(TokenBucket::<AtomicStorage, _>::from_parts(
limit,
clock.clone(),
let tb = Arc::new(TokenBucket::<AtomicStorage, _>::with_datum(
limit, FastClock,
));
let mut children = vec![];
let start = std::time::Instant::now();
Expand Down
12 changes: 4 additions & 8 deletions benches/time_ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,25 +4,21 @@ use criterion::{Criterion, criterion_group, criterion_main};
use gardal::{Clock, FastClock, QuantaClock, StdClock};

fn time_single_threaded(c: &mut Criterion) {
let c_clock = quanta::Clock::new();
// 1KHz
let _quanta_thread = quanta::Upkeep::new_with_clock(Duration::from_micros(10), c_clock.clone())
let _quanta_thread = quanta::Upkeep::new(Duration::from_micros(10))
.start()
.unwrap();
let mut group = c.benchmark_group("gardal");
group
.sample_size(100)
.bench_function("std-time-getting-instant", |b| {
let clock = StdClock::default();
b.iter(|| clock.now());
b.iter(|| std::hint::black_box(StdClock.now()));
})
.bench_function("quanta-time-getting-instant", |b| {
let clock = QuantaClock::default();
b.iter(|| clock.now());
b.iter(|| std::hint::black_box(QuantaClock.now()));
})
.bench_function("quanta-fast-getting-instant", |b| {
let clock = FastClock::new(c_clock.clone());
b.iter(|| clock.now());
b.iter(|| std::hint::black_box(FastClock.now()));
});
group.finish();
}
Expand Down
10 changes: 5 additions & 5 deletions examples/basic.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,13 @@
use std::time::Duration;

use gardal::{Limit, TokenBucket};
use gardal::{AtomicTokenBucket, Limit, StdClock};
use nonzero_ext::nonzero;

fn main() {
let tb = TokenBucket::new(Limit::per_second_and_burst(
nonzero!(10u32),
nonzero!(20u32),
));
let tb = AtomicTokenBucket::new(
Limit::per_second_and_burst(nonzero!(10u32), nonzero!(20u32)),
StdClock,
);
// after two seconds bucket should be full
std::thread::sleep(Duration::from_secs(2));
assert_eq!(5, tb.consume(nonzero!(5u32)).unwrap().as_u64());
Expand Down
6 changes: 2 additions & 4 deletions examples/fast_clock.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,14 +4,12 @@ use gardal::{FastClock, Limit, TokenBucket};
use nonzero_ext::nonzero;

fn main() {
let clock = quanta::Clock::new();
// Updates at 1Khz
let _quanta_thread = quanta::Upkeep::new_with_clock(Duration::from_millis(1), clock.clone())
let _quanta_thread = quanta::Upkeep::new(Duration::from_millis(1))
.start()
.unwrap();
let clock = FastClock::new(clock);
let limit = Limit::per_second_and_burst(nonzero!(10u32), nonzero!(20u32));
let tb = TokenBucket::with_clock(limit, clock);
let tb = TokenBucket::with_clock(limit, FastClock);
// after two seconds bucket should be full
println!("sleeping for 2 seconds...");
std::thread::sleep(Duration::from_secs(2));
Expand Down
6 changes: 3 additions & 3 deletions examples/streams.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,16 @@ use std::time::Duration;

use futures::StreamExt;
use futures::stream;
use gardal::SharedTokenBucket;
use gardal::futures::StreamExt as GardalStreamExt;
use gardal::{Limit, PaddedAtomicSharedStorage, TokenBucket, TokioClock};
use gardal::{Limit, TokioClock};
use nonzero_ext::nonzero;
use tokio::task::JoinSet;

#[tokio::main(flavor = "multi_thread")]
async fn main() {
let limit = Limit::per_second_and_burst(nonzero!(1000000u32), nonzero!(100u32));
let bucket =
TokenBucket::<PaddedAtomicSharedStorage, _>::from_parts(limit, TokioClock::default());
let bucket = SharedTokenBucket::new(limit, TokioClock);

let mut print_start = tokio::time::Instant::now();
let program_start = print_start;
Expand Down
6 changes: 3 additions & 3 deletions examples/weighted_stream.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
use futures::stream;
use gardal::futures::{StreamExt as GardalStreamExt, WeightedStream};
use gardal::{Limit, LocalStorage, TokenBucket, TokioClock};
use gardal::{Limit, LocalTokenBucket, TokioClock};
use nonzero_ext::nonzero;
use std::num::NonZeroU32;
use tokio_stream::StreamExt;
Expand Down Expand Up @@ -38,7 +38,7 @@ async fn main() {

// Create a throttling limit: 5 tokens per second with a burst of 25
let limit = Limit::per_second_and_burst(nonzero!(5u32), nonzero!(25u32));
let bucket = TokenBucket::<LocalStorage, _>::from_parts(limit, TokioClock::default());
let bucket = LocalTokenBucket::new(limit, TokioClock);

// Create a weighted stream where each task consumes tokens based on its size
let weighted_stream = WeightedStream::new(stream, bucket, |task: &Task| {
Expand Down Expand Up @@ -76,7 +76,7 @@ async fn main() {

let stream2 = stream::iter(tasks2);
let limit2 = Limit::per_second_and_burst(nonzero!(3u32), nonzero!(25u32));
let bucket2 = TokenBucket::<LocalStorage, _>::from_parts(limit2, TokioClock::default());
let bucket2 = LocalTokenBucket::new(limit2, TokioClock);

// Use the extension trait to create a weighted stream
let weighted_stream2 = stream2.throttle_weighted(bucket2, |text: &&str| {
Expand Down
Loading
Loading