From ee3ddc9d7f3a711e5a8079c4cab1cfab138084e7 Mon Sep 17 00:00:00 2001 From: Anirudha Bose Date: Mon, 6 Jul 2026 22:31:54 +0530 Subject: [PATCH] feat(screener): cache screen results with a 5-minute TTL --- Cargo.lock | 171 +++++++++++++++++++++++++++++++++++++++++++++++- Cargo.toml | 1 + src/screener.rs | 150 ++++++++++++++++++++++++++++++++++++++---- 3 files changed, 307 insertions(+), 15 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index f3a6601..fcf886d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -380,6 +380,17 @@ dependencies = [ "serde_json", ] +[[package]] +name = "async-lock" +version = "3.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "290f7f2596bd5b78a9fec8088ccd89180d7f9f55b94b0576823bbbdc72ee8311" +dependencies = [ + "event-listener", + "event-listener-strategy", + "pin-project-lite", +] + [[package]] name = "async-trait" version = "0.1.89" @@ -1074,6 +1085,7 @@ dependencies = [ "base64", "dotenvy", "http-body-util", + "moka", "reqwest", "serde_json", "thiserror", @@ -1202,6 +1214,15 @@ dependencies = [ "memchr", ] +[[package]] +name = "concurrent-queue" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ca0197aee26d1ae37445ee532fefce43251d24cc7c166799f4d46817f1d3973" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "const-hex" version = "1.19.1" @@ -1319,6 +1340,30 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "crossbeam-channel" +version = "0.5.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d85363c37faeca707aef026efa9f3b34d077bce547e48f770770625c6013679e" +dependencies = [ + "crossbeam-utils", +] + +[[package]] +name = "crossbeam-epoch" +version = "0.9.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4279b0f31df31b4add8aeae57d40f3b5e53aad2b531e89b982bd75ed1898851d" +dependencies = [ + "crossbeam-utils", +] + +[[package]] +name = "crossbeam-utils" +version = "0.8.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61803da095bee82a81bb1a452ecc25d3b2f1416d1897eb86430c6159ef717c17" + [[package]] name = "crunchy" version = "0.2.4" @@ -1629,6 +1674,27 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "event-listener" +version = "5.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e13b66accf52311f30a0db42147dadea9850cb48cd070028831ae5f5d4b856ab" +dependencies = [ + "concurrent-queue", + "parking", + "pin-project-lite", +] + +[[package]] +name = "event-listener-strategy" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8be9f3dfaaffdae2972880079a491a1a8bb7cbed0b8dd7a347f668b4150a3b93" +dependencies = [ + "event-listener", + "pin-project-lite", +] + [[package]] name = "fastrand" version = "2.4.1" @@ -1839,11 +1905,22 @@ dependencies = [ "cfg-if", "js-sys", "libc", - "r-efi", + "r-efi 5.3.0", "wasip2", "wasm-bindgen", ] +[[package]] +name = "getrandom" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099" +dependencies = [ + "cfg-if", + "libc", + "r-efi 6.0.0", +] + [[package]] name = "group" version = "0.13.0" @@ -2444,6 +2521,15 @@ version = "0.8.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "92daf443525c4cce67b150400bc2316076100ce0b3686209eb8cf3c31612e6f0" +[[package]] +name = "lock_api" +version = "0.4.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "224399e74b87b5f3557511d98dff8b14089b3dadafcab6bb93eab67d3aace965" +dependencies = [ + "scopeguard", +] + [[package]] name = "log" version = "0.4.30" @@ -2524,6 +2610,26 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "moka" +version = "0.12.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "957228ad12042ee839f93c8f257b62b4c0ab5eaae1d4fa60de53b27c9d7c5046" +dependencies = [ + "async-lock", + "crossbeam-channel", + "crossbeam-epoch", + "crossbeam-utils", + "equivalent", + "event-listener", + "futures-util", + "parking_lot", + "portable-atomic", + "smallvec", + "tagptr", + "uuid", +] + [[package]] name = "nu-ansi-term" version = "0.50.3" @@ -2624,6 +2730,35 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "parking" +version = "2.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f38d5652c16fde515bb1ecef450ab0f6a219d619a7274976324d5e377f7dceba" + +[[package]] +name = "parking_lot" +version = "0.12.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93857453250e3077bd71ff98b6a65ea6621a19bb0f559a85248955ac12c45a1a" +dependencies = [ + "lock_api", + "parking_lot_core", +] + +[[package]] +name = "parking_lot_core" +version = "0.9.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2621685985a2ebf1c516881c026032ac7deafcda1a2c9b7850dc81e3dfcb64c1" +dependencies = [ + "cfg-if", + "libc", + "redox_syscall", + "smallvec", + "windows-link", +] + [[package]] name = "paste" version = "1.0.15" @@ -2668,6 +2803,12 @@ dependencies = [ "spki", ] +[[package]] +name = "portable-atomic" +version = "1.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c33a9471896f1c69cecef8d20cbe2f7accd12527ce60845ff44c153bb2a21b49" + [[package]] name = "potential_utf" version = "0.1.5" @@ -2859,6 +3000,12 @@ version = "5.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" +[[package]] +name = "r-efi" +version = "6.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" + [[package]] name = "radium" version = "0.7.0" @@ -2944,6 +3091,15 @@ dependencies = [ "rustversion", ] +[[package]] +name = "redox_syscall" +version = "0.5.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d" +dependencies = [ + "bitflags", +] + [[package]] name = "ref-cast" version = "1.0.25" @@ -3346,6 +3502,12 @@ dependencies = [ "serde_json", ] +[[package]] +name = "scopeguard" +version = "1.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" + [[package]] name = "seahash" version = "4.1.0" @@ -3783,6 +3945,12 @@ dependencies = [ "libc", ] +[[package]] +name = "tagptr" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b2093cf4c8eb1e67749a6762251bc9cd836b6fc171623bd0a9d324d37af2417" + [[package]] name = "tap" version = "1.0.1" @@ -4167,6 +4335,7 @@ version = "1.23.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d258b83ceec21034727ecee8c382cfa6c3e133699b0742c64571814fb420c9f7" dependencies = [ + "getrandom 0.4.3", "js-sys", "wasm-bindgen", ] diff --git a/Cargo.toml b/Cargo.toml index cec5a4a..0757f38 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -33,6 +33,7 @@ aws-sdk-s3 = { version = "1", default-features = false, features = ["default-htt # legacy rustls connector, so no vulnerable TLS stack. aws-config = "1" base64 = "0.22" +moka = { version = "0.12.15", features = ["future"] } [dev-dependencies] http-body-util = "=0.1.3" diff --git a/src/screener.rs b/src/screener.rs index 88f2941..9cd1ec5 100644 --- a/src/screener.rs +++ b/src/screener.rs @@ -10,6 +10,13 @@ //! other error (timeout, misconfiguration, S3 outage) blocks the payment rather than //! letting it through unchecked. //! +//! A short-lived cache sits in front of the `HeadObject`, so repeat screens of the same +//! identifier skip the round-trip: +//! +//! - only definite results (`Allowed`/`Blocked`) are cached, for [`SCREEN_CACHE_TTL`] +//! - the TTL bounds how long a newly restricted address can keep paying +//! - errors are never cached, so an unreachable bucket re-screens every request +//! //! The module is chain agnostic. It screens the identifier string exactly as given and //! knows nothing about the address itself. It returns a plain outcome, never an HTTP //! response; the rail turns that outcome into a status code. @@ -36,12 +43,25 @@ use crate::Config; /// Canary key used to probe the bucket at startup. const CANARY_KEY: &str = "bx402.canary"; +/// How long a definite screen result stays cached. This bounds staleness: a newly +/// restricted address keeps paying until its cached `Allowed` entry expires. +const SCREEN_CACHE_TTL: Duration = Duration::from_secs(300); + +/// Cap on the number of cached identifiers. Screening runs before the payment is verified, +/// so untrusted addresses reach the cache. The cap bounds memory: once the cache is full, +/// the oldest entry is evicted. +const SCREEN_CACHE_MAX_ENTRIES: u64 = 100_000; + /// Screens identifiers against the restricted-address S3 bucket. `Clone` is cheap: the -/// inner `aws_sdk_s3::Client` is `Arc`-backed, like `reqwest::Client`. +/// inner `aws_sdk_s3::Client` and the cache are both `Arc`-backed, so every clone shares +/// one client and one cache. #[derive(Clone)] pub struct RestrictedAddressScreener { client: aws_sdk_s3::Client, bucket: String, + /// Recently screened identifiers, each mapped to its definite result for + /// [`SCREEN_CACHE_TTL`]. Only `Allowed`/`Blocked` are cached, never a [`ScreenError`]. + cache: moka::future::Cache, } /// The two definite answers a screen can give. @@ -67,16 +87,42 @@ pub(crate) struct ScreenError(#[source] Box impl RestrictedAddressScreener { pub(crate) fn new(client: aws_sdk_s3::Client, bucket: String) -> Self { - Self { client, bucket } + Self::with_ttl(client, bucket, SCREEN_CACHE_TTL) + } + + /// Build a screener whose cache entries live for `ttl`. Production uses + /// [`SCREEN_CACHE_TTL`] via [`new`](Self::new); tests inject a short `ttl` to exercise + /// expiry. + fn with_ttl(client: aws_sdk_s3::Client, bucket: String, ttl: Duration) -> Self { + let cache = moka::future::Cache::builder() + .max_capacity(SCREEN_CACHE_MAX_ENTRIES) + .time_to_live(ttl) + .build(); + Self { + client, + bucket, + cache, + } } /// Screen one identifier, exactly as given. /// - /// The caller must pass the already-canonical form (casing differs per chain; see - /// the module docs). The identifier is base64url-encoded into the S3 key. + /// The caller must pass the already-canonical form; casing differs per chain (see the + /// module docs). The result comes from the cache when possible, otherwise from a live + /// `HeadObject`: + /// + /// - cache hit: return the stored result, no S3 call + /// - cache miss: encode the identifier into the S3 key, look it up, cache the answer + /// - [`ScreenError`]: never cached, so the next request checks S3 again pub(crate) async fn screen(&self, identifier: &str) -> Result { - self.head_key(URL_SAFE_NO_PAD.encode(identifier.as_bytes())) - .await + if let Some(hit) = self.cache.get(identifier).await { + return Ok(hit); + } + let screening = self + .head_key(URL_SAFE_NO_PAD.encode(identifier.as_bytes())) + .await?; + self.cache.insert(identifier.to_string(), screening).await; + Ok(screening) } /// Look up one exact S3 key with `HeadObject`: @@ -205,13 +251,21 @@ mod tests { test_screener(endpoint, "restricted-address-bucket") } - /// Screen `"0xanything"` against a mock S3 that answers every `HEAD` with `status`. - async fn screen_against(status: u16) -> Result { + /// Mount a HEAD mock that answers with `status` exactly `hits` times, failing the test + /// on drop if it was called a different number of times. + async fn server_expecting(status: u16, hits: u64) -> MockServer { let server = MockServer::start().await; Mock::given(method("HEAD")) .respond_with(ResponseTemplate::new(status)) + .expect(hits) .mount(&server) .await; + server + } + + /// Screen `"0xanything"` against a mock S3 that answers one `HEAD` with `status`. + async fn screen_against(status: u16) -> Result { + let server = server_expecting(status, 1).await; screener_for(server.uri()).screen("0xanything").await } @@ -278,6 +332,78 @@ mod tests { assert!(screener.screen("0xanything").await.is_err()); } + #[tokio::test] + async fn caches_allowed_result() { + let server = server_expecting(404, 1).await; + let screener = screener_for(server.uri()); + assert_eq!( + screener.screen("0xpayer").await.unwrap(), + Screening::Allowed + ); + assert_eq!( + screener.screen("0xpayer").await.unwrap(), + Screening::Allowed, + "the second screen is served from the cache", + ); + // `.expect(1)` is verified when the server drops: only one HEAD reached S3. + } + + #[tokio::test] + async fn caches_blocked_result() { + let server = server_expecting(200, 1).await; + let screener = screener_for(server.uri()); + assert_eq!( + screener.screen("0xpayer").await.unwrap(), + Screening::Blocked + ); + assert_eq!( + screener.screen("0xpayer").await.unwrap(), + Screening::Blocked, + "a blocked result is cached too", + ); + } + + #[tokio::test] + async fn does_not_cache_errors() { + // Two HEADs expected: an error is never cached, so it re-screens every time. + let server = server_expecting(500, 2).await; + let screener = screener_for(server.uri()); + assert!(screener.screen("0xpayer").await.is_err()); + assert!(screener.screen("0xpayer").await.is_err()); + } + + #[tokio::test] + async fn distinct_identifiers_screen_separately() { + // Two different identifiers are two different keys, so each hits S3 once. + let server = server_expecting(404, 2).await; + let screener = screener_for(server.uri()); + assert_eq!(screener.screen("0xone").await.unwrap(), Screening::Allowed); + assert_eq!(screener.screen("0xtwo").await.unwrap(), Screening::Allowed); + } + + #[tokio::test] + async fn expired_entry_rescreens() { + // A short TTL, so the cached entry expires between the two screens and the second + // one goes back to S3. moka times entries with its own clock, so this uses a brief + // real sleep rather than a paused runtime clock. + let server = server_expecting(404, 2).await; + let screener = RestrictedAddressScreener::with_ttl( + test_client(server.uri()), + "restricted-address-bucket".to_string(), + Duration::from_millis(50), + ); + assert_eq!( + screener.screen("0xpayer").await.unwrap(), + Screening::Allowed + ); + tokio::time::sleep(Duration::from_millis(120)).await; + assert_eq!( + screener.screen("0xpayer").await.unwrap(), + Screening::Allowed, + "the entry expired, so the screen re-reads from S3", + ); + } + #[tokio::test] async fn init_disabled_when_bucket_unset() { let config = crate::Config { @@ -317,13 +443,9 @@ mod tests { assert!(matches!(status, Status::Enabled { .. })); } - /// Run `init_with` against a mock S3 that answers every `HEAD` with `status`. + /// Run `init_with` against a mock S3 that answers one `HEAD` with `status`. async fn init_against(status: u16) -> anyhow::Result<(RestrictedAddressScreener, Status)> { - let server = MockServer::start().await; - Mock::given(method("HEAD")) - .respond_with(ResponseTemplate::new(status)) - .mount(&server) - .await; + let server = server_expecting(status, 1).await; init_with( test_client(server.uri()), "restricted-address-bucket".to_string(),