diff --git a/api/.env.example b/api/.env.example index aac5cfa..67128e9 100644 --- a/api/.env.example +++ b/api/.env.example @@ -18,3 +18,7 @@ SYSTEM_SLICE_URI=at://did:plc:bcgltzqazw5tb6k2g3ttenbj/network.slices.slice/3lym # Logging level RUST_LOG=debug + +# Redis configuration (optional - if not set, falls back to in-memory cache) +REDIS_URL=redis://localhost:6379 +REDIS_TTL_SECONDS=3600 diff --git a/api/.sqlx/query-8fe482ef09e83d77ac7679d98e5c36c10c6b00c30cc289616a6da9c86ac11006.json b/api/.sqlx/query-8fe482ef09e83d77ac7679d98e5c36c10c6b00c30cc289616a6da9c86ac11006.json new file mode 100644 index 0000000..fe90206 --- /dev/null +++ b/api/.sqlx/query-8fe482ef09e83d77ac7679d98e5c36c10c6b00c30cc289616a6da9c86ac11006.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT time_us\n FROM jetstream_cursor\n WHERE id = $1\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "time_us", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "8fe482ef09e83d77ac7679d98e5c36c10c6b00c30cc289616a6da9c86ac11006" +} diff --git a/api/.sqlx/query-ae961e575c6f9d0761401115344ae5993753bd5064fc66344b402aa4ee8cc177.json b/api/.sqlx/query-ae961e575c6f9d0761401115344ae5993753bd5064fc66344b402aa4ee8cc177.json new file mode 100644 index 0000000..0113222 --- /dev/null +++ b/api/.sqlx/query-ae961e575c6f9d0761401115344ae5993753bd5064fc66344b402aa4ee8cc177.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO jetstream_cursor (id, time_us, updated_at)\n VALUES ($1, $2, NOW())\n ON CONFLICT (id)\n DO UPDATE SET time_us = $2, updated_at = NOW()\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Int8" + ] + }, + "nullable": [] + }, + "hash": "ae961e575c6f9d0761401115344ae5993753bd5064fc66344b402aa4ee8cc177" +} diff --git a/api/.sqlx/query-e167d0d1f047f25f6116097ff025352727645c99faf8fc451678d2a130df7e78.json b/api/.sqlx/query-e167d0d1f047f25f6116097ff025352727645c99faf8fc451678d2a130df7e78.json new file mode 100644 index 0000000..525a463 --- /dev/null +++ b/api/.sqlx/query-e167d0d1f047f25f6116097ff025352727645c99faf8fc451678d2a130df7e78.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO jetstream_cursor (id, time_us, updated_at)\n VALUES ($1, $2, NOW())\n ON CONFLICT (id)\n DO UPDATE SET time_us = $2, updated_at = NOW()\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Int8" + ] + }, + "nullable": [] + }, + "hash": "e167d0d1f047f25f6116097ff025352727645c99faf8fc451678d2a130df7e78" +} diff --git a/api/CLAUDE.md b/api/CLAUDE.md index 79b2f8c..dbaf2ed 100644 --- a/api/CLAUDE.md +++ b/api/CLAUDE.md @@ -85,6 +85,8 @@ cargo build --release stable pagination - **OAuth DPoP authentication** integrated with AIP server for ATProto authentication +- **Multi-tier caching** with Redis (if configured) or in-memory fallback for + performance optimization ### Module Organization @@ -95,6 +97,8 @@ cargo build --release - `src/jetstream.rs` - Real-time event processing from ATProto firehose - `src/sync.rs` - Bulk synchronization operations with ATProto relay - `src/auth.rs` - OAuth verification and DPoP authentication setup +- `src/cache.rs` - Generic caching interface and in-memory cache implementation +- `src/redis_cache.rs` - Redis cache implementation for distributed caching - `src/errors.rs` - Error type definitions (reference for new errors) ## Error Handling @@ -131,3 +135,42 @@ helper functionality. - After updating, run `cargo check` to fix errors and warnings - Don't use dead code, if it's not used remove it + +## Caching Architecture + +The application uses a flexible caching system that supports both Redis and in-memory caching with automatic fallback. + +### Cache Configuration + +Configure caching via environment variables: + +```bash +# Redis configuration (optional) +REDIS_URL=redis://localhost:6379 +REDIS_TTL_SECONDS=3600 +``` + +If `REDIS_URL` is not set, the application automatically falls back to in-memory caching. + +### Cache Types and TTLs + +- **Actor Cache** (Jetstream): No TTL (permanent cache for slice actors) +- **Lexicon Cache**: 2 hours (7200s) - lexicons change infrequently +- **Domain Cache**: 4 hours (14400s) - slice domain mappings rarely change +- **Collections Cache**: 2 hours (7200s) - slice collections change infrequently +- **Auth Cache**: 5 minutes (300s) - OAuth tokens and AT Protocol sessions +- **DID Resolution Cache**: 24 hours (86400s) - DID documents change rarely + +### Cache Implementation + +- `src/cache.rs` - Defines the `Cache` trait and `SliceCache` wrapper with domain-specific methods +- `src/redis_cache.rs` - Redis implementation of the `Cache` trait +- Both implementations provide the same interface through `SliceCache` +- Cache keys use prefixed formats (e.g., `actor:{did}:{slice_uri}`, `oauth_userinfo:{token}`) + +### Cache Usage Patterns + +- **Jetstream Consumer**: Creates 4 separate cache instances for actors, lexicons, domains, and collections +- **Auth System**: Uses dedicated auth cache for OAuth and AT Protocol session caching +- **Actor Resolution**: Caches DID resolution results to avoid repeated lookups +- **Automatic Fallback**: Redis failures automatically fall back to in-memory caching without errors diff --git a/api/Cargo.lock b/api/Cargo.lock index 1a5b8da..0dc8998 100644 --- a/api/Cargo.lock +++ b/api/Cargo.lock @@ -4,9 +4,9 @@ version = 4 [[package]] name = "addr2line" -version = "0.24.2" +version = "0.25.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dfbe277e56a376000877090da837660b4427aad530e3028d44e0bffe4f89a1c1" +checksum = "1b5d307320b3181d6d7954e663bd7c774a838b8220fe0593c86d9fb09f498b4b" dependencies = [ "gimli", ] @@ -32,12 +32,6 @@ version = "0.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" -[[package]] -name = "android-tzdata" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e999941b234f3131b00bc13c22d06e8c5ff726d1b6318ac7eb276997bbb4fef0" - [[package]] name = "android_system_properties" version = "0.1.5" @@ -49,9 +43,9 @@ dependencies = [ [[package]] name = "anyhow" -version = "1.0.99" +version = "1.0.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b0674a1ddeecb70197781e945de4b3b8ffb61fa939a5597bcf48503737663100" +checksum = "a23eb6b1614318a8071c9b2521f36b424b2c83db5eb3a0fead4a6c0809af6e61" [[package]] name = "anymap2" @@ -59,6 +53,12 @@ version = "0.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d301b3b94cb4b2f23d7917810addbbaff90738e0ca2be692bd027e70d7e0330c" +[[package]] +name = "arc-swap" +version = "1.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69f7f8c3906b62b754cd5326047894316021dcfe5a194c8ea52bdd94934a3457" + [[package]] name = "async-trait" version = "0.1.89" @@ -87,9 +87,9 @@ checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" [[package]] name = "atproto-client" -version = "0.13.0" +version = "0.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c34ed7ebeec01cd7775c1c7841838c142d123d403983ed7179b31850435b5c7c" +checksum = "9f388d83aa9552d7c7b80cc131558fb86be5a46a104742fb4ccf38e22a3c8620" dependencies = [ "anyhow", "atproto-identity", @@ -101,7 +101,7 @@ dependencies = [ "reqwest-middleware", "serde", "serde_json", - "thiserror 2.0.14", + "thiserror 2.0.16", "tokio", "tracing", "urlencoding", @@ -109,9 +109,9 @@ dependencies = [ [[package]] name = "atproto-identity" -version = "0.13.0" +version = "0.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b956c07726fce812630be63c5cb31b1961cbb70f0a05614278523102d78c3a48" +checksum = "65e405e13a96ce91d1e832f56b90ae3f1fcbe29a5d9731e417074bfe1df71a5f" dependencies = [ "anyhow", "async-trait", @@ -128,7 +128,7 @@ dependencies = [ "serde", "serde_ipld_dagcbor", "serde_json", - "thiserror 2.0.14", + "thiserror 2.0.16", "tokio", "tracing", "urlencoding", @@ -136,9 +136,9 @@ dependencies = [ [[package]] name = "atproto-jetstream" -version = "0.13.0" +version = "0.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7b1897fb2f7c6d02d46f7b8d25d653c141cee4a68a10efd135d46201a95034db" +checksum = "b36d0d4fec207d04563bdb151ca793dbfbb9c4708b5c9320d65da4b564b186d2" dependencies = [ "anyhow", "async-trait", @@ -147,7 +147,7 @@ dependencies = [ "http", "serde", "serde_json", - "thiserror 2.0.14", + "thiserror 2.0.16", "tokio", "tokio-util", "tokio-websockets", @@ -159,9 +159,9 @@ dependencies = [ [[package]] name = "atproto-oauth" -version = "0.13.0" +version = "0.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3ea205901c33d074a1b498591d0511bcd788b6772ec0ca6e09a92c4327ddbdff" +checksum = "3f7a82388b59f83c2c141af434afa96d41841e4e7c1b361afc83193f2ee75d63" dependencies = [ "anyhow", "async-trait", @@ -183,7 +183,7 @@ dependencies = [ "serde_ipld_dagcbor", "serde_json", "sha2", - "thiserror 2.0.14", + "thiserror 2.0.16", "tokio", "tracing", "ulid", @@ -191,9 +191,9 @@ dependencies = [ [[package]] name = "atproto-record" -version = "0.13.0" +version = "0.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0550f74423ca745132dc07ba1cb01f2f08243a7bf7497f5c3f2185a774c92ca2" +checksum = "accec4d63f3f653e947b051b8c9f56f6d939d7a646e44b845eeabc03c713dad5" dependencies = [ "anyhow", "atproto-identity", @@ -202,7 +202,7 @@ dependencies = [ "serde", "serde_ipld_dagcbor", "serde_json", - "thiserror 2.0.14", + "thiserror 2.0.16", ] [[package]] @@ -305,11 +305,20 @@ dependencies = [ "syn 2.0.106", ] +[[package]] +name = "backon" +version = "1.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "592277618714fbcecda9a02ba7a8781f319d26532a88553bbacc77ba5d2b3a8d" +dependencies = [ + "fastrand", +] + [[package]] name = "backtrace" -version = "0.3.75" +version = "0.3.76" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6806a6321ec58106fea15becdad98371e28d92ccbc7c8f1b3b6dd724fe8f1002" +checksum = "bb531853791a215d7c62a30daf0dde835f381ab5de4589cfe7c649d2cbe92bd6" dependencies = [ "addr2line", "cfg-if", @@ -317,7 +326,7 @@ dependencies = [ "miniz_oxide", "object", "rustc-demangle", - "windows-targets 0.52.6", + "windows-link 0.2.0", ] [[package]] @@ -352,9 +361,9 @@ checksum = "55248b47b0caf0546f7988906588779981c43bb1bc9d0c44087278f80cdb44ba" [[package]] name = "bitflags" -version = "2.9.1" +version = "2.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1b8e56985ec62d17e9c1001dc89c88ecd7dc08e47eba5ec7c29c7b5eeecde967" +checksum = "2261d10cca569e4643e526d8dc2e62e433cc8aba21ab764233731f8d369bf394" dependencies = [ "serde", ] @@ -397,10 +406,11 @@ dependencies = [ [[package]] name = "cc" -version = "1.2.33" +version = "1.2.39" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3ee0f8803222ba5a7e2777dd72ca451868909b1ac410621b676adf07280e9b5f" +checksum = "e1354349954c6fc9cb0deab020f27f783cf0b604e8bb754dc4658ecf0d29c35f" dependencies = [ + "find-msvc-tools", "jobserver", "libc", "shlex", @@ -408,9 +418,9 @@ dependencies = [ [[package]] name = "cfg-if" -version = "1.0.1" +version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9555578bc9e57714c812a1f84e4fc5b4d21fcb063490c624de019f7464c91268" +checksum = "2fd1289c04a9ea8cb22300a459a72a385d7c73d3259e2ed7dcb2af674838cfa9" [[package]] name = "cfg_aliases" @@ -420,17 +430,16 @@ checksum = "613afe47fcd5fac7ccf1db93babcb082c5994d996f20b8b159f2ad1658eb5724" [[package]] name = "chrono" -version = "0.4.41" +version = "0.4.42" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c469d952047f47f91b68d1cba3f10d63c11d73e4636f24f08daf0278abf01c4d" +checksum = "145052bdd345b87320e369255277e3fb5152762ad123a901ef5c262dd38fe8d2" dependencies = [ - "android-tzdata", "iana-time-zone", "js-sys", "num-traits", "serde", "wasm-bindgen", - "windows-link", + "windows-link 0.2.0", ] [[package]] @@ -447,6 +456,20 @@ dependencies = [ "unsigned-varint", ] +[[package]] +name = "combine" +version = "4.6.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba5a308b75df32fe02788e748662718f03fde005016435c444eea572398219fd" +dependencies = [ + "bytes", + "futures-core", + "memchr", + "pin-project-lite", + "tokio", + "tokio-util", +] + [[package]] name = "concurrent-queue" version = "2.5.0" @@ -725,12 +748,12 @@ checksum = "877a4ace8713b0bcf2a4e7eec82529c029f1d0619886d18145fea96c3ffe5c0f" [[package]] name = "errno" -version = "0.3.13" +version = "0.3.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "778e2ac28f6c47af28e4907f13ffd1e1ddbd400980a9abd7c8df189bf578a5ad" +checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.1", ] [[package]] @@ -771,6 +794,12 @@ dependencies = [ "subtle", ] +[[package]] +name = "find-msvc-tools" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ced73b1dacfc750a6db6c0a0c3a3853c8b41997e2e2c563dc90804ae6867959" + [[package]] name = "flume" version = "0.11.1" @@ -811,9 +840,9 @@ checksum = "00b0228411908ca8685dba7fc2cdd70ec9990a6e753e89b6ac91a84c40fbaf4b" [[package]] name = "form_urlencoded" -version = "1.2.1" +version = "1.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e13624c2627564efccf4934284bdd98cbaa14e79b0b5a141218e507b3a823456" +checksum = "cb4cb245038516f5f85277875cdaa4f7d2c9a0fa0468de06ed190163b1581fcf" dependencies = [ "percent-encoding", ] @@ -918,20 +947,6 @@ dependencies = [ "slab", ] -[[package]] -name = "generator" -version = "0.8.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d18470a76cb7f8ff746cf1f7470914f900252ec36bbc40b569d74b1258446827" -dependencies = [ - "cc", - "cfg-if", - "libc", - "log", - "rustversion", - "windows", -] - [[package]] name = "generic-array" version = "0.14.7" @@ -966,15 +981,15 @@ dependencies = [ "js-sys", "libc", "r-efi", - "wasi 0.14.2+wasi-0.2.4", + "wasi 0.14.7+wasi-0.2.4", "wasm-bindgen", ] [[package]] name = "gimli" -version = "0.31.1" +version = "0.32.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "07e28edb80900c19c28f1072f2e8aeca7fa06b23cd4169cefe1af5aa3260783f" +checksum = "e629b9b98ef3dd8afe6ca2bd0f89306cec16d43d907889945bc5d6687f2f13c7" [[package]] name = "group" @@ -1017,13 +1032,19 @@ dependencies = [ "foldhash", ] +[[package]] +name = "hashbrown" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5419bdc4f6a9207fbeba6d11b604d481addf78ecd10c11ad51e76c2f6482748d" + [[package]] name = "hashlink" version = "0.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7382cf6263419f2d8df38c55d7da83da5c18aef87fc7a7fc1fb1e344edfe14c1" dependencies = [ - "hashbrown", + "hashbrown 0.15.5", ] [[package]] @@ -1056,7 +1077,7 @@ dependencies = [ "once_cell", "rand 0.9.2", "ring", - "thiserror 2.0.14", + "thiserror 2.0.16", "tinyvec", "tokio", "tracing", @@ -1079,7 +1100,7 @@ dependencies = [ "rand 0.9.2", "resolv-conf", "smallvec", - "thiserror 2.0.14", + "thiserror 2.0.16", "tokio", "tracing", ] @@ -1159,13 +1180,14 @@ checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" [[package]] name = "hyper" -version = "1.6.0" +version = "1.7.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cc2b571658e38e0c01b1fdca3bbbe93c00d3d71693ff2770043f8c29bc7d6f80" +checksum = "eb3aa54a13a0dfe7fbe3a59e0c76093041720fdc77b110cc0fc260fafb4dc51e" dependencies = [ + "atomic-waker", "bytes", "futures-channel", - "futures-util", + "futures-core", "h2", "http", "http-body", @@ -1173,6 +1195,7 @@ dependencies = [ "httpdate", "itoa", "pin-project-lite", + "pin-utils", "smallvec", "tokio", "want", @@ -1213,9 +1236,9 @@ dependencies = [ [[package]] name = "hyper-util" -version = "0.1.16" +version = "0.1.17" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8d9b05277c7e8da2c93a568989bb6207bef0112e8d17df7a6eda4a3cf143bc5e" +checksum = "3c6995591a8f1380fcb4ba966a252a4b29188d51d2b89e3a252f5305be65aea8" dependencies = [ "base64 0.22.1", "bytes", @@ -1239,9 +1262,9 @@ dependencies = [ [[package]] name = "iana-time-zone" -version = "0.1.63" +version = "0.1.64" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b0c919e5debc312ad217002b8048a17b7d83f80703865bbfcfebb0458b0b27d8" +checksum = "33e57f83510bb73707521ebaffa789ec8caf86f9657cad665b092b581d40e9fb" dependencies = [ "android_system_properties", "core-foundation-sys", @@ -1349,9 +1372,9 @@ dependencies = [ [[package]] name = "idna" -version = "1.0.3" +version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "686f825264d630750a544639377bae737628043f20d38bbc029e8f29ea968a7e" +checksum = "3b0875f23caa03898994f6ddc501886a45c7d3d62d04d2d90788d47be1b1e4de" dependencies = [ "idna_adapter", "smallvec", @@ -1370,19 +1393,19 @@ dependencies = [ [[package]] name = "indexmap" -version = "2.10.0" +version = "2.11.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fe4cd85333e22411419a0bcae1297d25e58c9443848b11dc6a86fefe8c78a661" +checksum = "4b0f83760fb341a774ed326568e19f5a863af4a952def8c39f9ab92fd95b88e5" dependencies = [ "equivalent", - "hashbrown", + "hashbrown 0.16.0", ] [[package]] name = "io-uring" -version = "0.7.9" +version = "0.7.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d93587f37623a1a17d94ef2bc9ada592f5465fe7732084ab7beefabe5c77c0c4" +checksum = "046fa2d4d00aea763528b4950358d0ead425372445dc8ff86312b3c69ff7727b" dependencies = [ "bitflags", "cfg-if", @@ -1446,9 +1469,9 @@ dependencies = [ [[package]] name = "js-sys" -version = "0.3.77" +version = "0.3.81" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1cfaf33c695fc6e08064efbc1f72ec937429614f25eef83af942d0e227c3a28f" +checksum = "ec48937a97411dcb524a265206ccd4c90bb711fca92b2792c407f268825b9305" dependencies = [ "once_cell", "wasm-bindgen", @@ -1479,9 +1502,9 @@ dependencies = [ [[package]] name = "libc" -version = "0.2.175" +version = "0.2.176" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a82ae493e598baaea5209805c49bbf2ea7de956d50d7da0da1164f9c6d28543" +checksum = "58f929b4d672ea937a23a1ab494143d968337a5f47e56d0815df1e0890ddf174" [[package]] name = "libm" @@ -1491,9 +1514,9 @@ checksum = "f9fbbcab51052fe104eb5e5d351cf728d30a5be1fe14d9be8a3b097481fb97de" [[package]] name = "libredox" -version = "0.1.9" +version = "0.1.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "391290121bad3d37fbddad76d8f5d1c1c314cfc646d143d7e07a3086ddff0ce3" +checksum = "416f7e718bdb06000964960ffa43b4335ad4012ae8b99060261aa4a8088d5ccb" dependencies = [ "bitflags", "libc", @@ -1512,9 +1535,9 @@ dependencies = [ [[package]] name = "linux-raw-sys" -version = "0.9.4" +version = "0.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cd945864f07fe9f5371a27ad7b52a172b4b499999f1d97574c9fa68373937e12" +checksum = "df1d3c3b53da64cf5760482273a98e575c651a67eec7f77df96b5b642de8f039" [[package]] name = "litemap" @@ -1534,22 +1557,9 @@ dependencies = [ [[package]] name = "log" -version = "0.4.27" +version = "0.4.28" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "13dc2df351e3202783a1fe0d44375f7295ffb4049267b0f3018346dc122a1d94" - -[[package]] -name = "loom" -version = "0.7.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "419e0dc8046cb947daa77eb95ae174acfbddb7673b4151f56d1eed8e93fbfaca" -dependencies = [ - "cfg-if", - "generator", - "scoped-tls", - "tracing", - "tracing-subscriber", -] +checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432" [[package]] name = "lru" @@ -1557,7 +1567,7 @@ version = "0.12.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "234cf4f4a04dc1f57e24b96cc0cd600cf2af460d4161ac5ecdd0af8e1f3b2a38" dependencies = [ - "hashbrown", + "hashbrown 0.15.5", ] [[package]] @@ -1568,11 +1578,11 @@ checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" [[package]] name = "matchers" -version = "0.1.0" +version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8263075bb86c5a1b1427b5ae862e8889656f126e9f77c484496e8b47cf5c5558" +checksum = "d1525a2a28c7f4fa0fc98bb91ae755d1e2d1505079e05539e35bc876b5d65ae9" dependencies = [ - "regex-automata 0.1.10", + "regex-automata", ] [[package]] @@ -1593,9 +1603,9 @@ dependencies = [ [[package]] name = "memchr" -version = "2.7.5" +version = "2.7.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "32a282da65faaf38286cf3be983213fcf1d2e2a58700e808f83f4ea9a4804bc0" +checksum = "f52b00d39961fc5b2736ea853c9cc86238e165017a493d1d5c8eac6bdc4cc273" [[package]] name = "mime" @@ -1635,20 +1645,19 @@ dependencies = [ [[package]] name = "moka" -version = "0.12.10" +version = "0.12.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a9321642ca94a4282428e6ea4af8cc2ca4eac48ac7a6a4ea8f33f76d0ce70926" +checksum = "8261cd88c312e0004c1d51baad2980c66528dfdb2bee62003e643a4d8f86b077" dependencies = [ "crossbeam-channel", "crossbeam-epoch", "crossbeam-utils", - "loom", + "equivalent", "parking_lot", "portable-atomic", "rustc_version", "smallvec", "tagptr", - "thiserror 1.0.69", "uuid", ] @@ -1710,12 +1719,21 @@ dependencies = [ [[package]] name = "nu-ansi-term" -version = "0.46.0" +version = "0.50.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d4a28e057d01f97e61255210fcff094d74ed0466038633e95017f5beb68e4399" +dependencies = [ + "windows-sys 0.52.0", +] + +[[package]] +name = "num-bigint" +version = "0.4.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "77a8165726e8236064dbb45459242600304b42a5ea24ee2948e18e023bf7ba84" +checksum = "a5e44f723f1133c9deac646763579fdb3ac745e418f2a7af9cd0c431da1f20b9" dependencies = [ - "overload", - "winapi", + "num-integer", + "num-traits", ] [[package]] @@ -1767,9 +1785,9 @@ dependencies = [ [[package]] name = "object" -version = "0.36.7" +version = "0.37.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "62948e14d923ea95ea2c7c86c71013138b66525b86bdc08d2dcc262bdb497b87" +checksum = "ff76201f031d8863c38aa7f905eca4f53abbfa15f609db4277d44cd8938f33fe" dependencies = [ "memchr", ] @@ -1828,12 +1846,6 @@ dependencies = [ "vcpkg", ] -[[package]] -name = "overload" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b15813163c1d831bf4a13c3610c05c0d03b39feb07f7e09fa234dac9b15aaf39" - [[package]] name = "p256" version = "0.13.2" @@ -1900,9 +1912,9 @@ dependencies = [ [[package]] name = "percent-encoding" -version = "2.3.1" +version = "2.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e3148f5046208a5d56bcfc03053e3ca6334e51da8dfb19b6cdc8b306fae3283e" +checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] name = "pin-project-lite" @@ -1951,9 +1963,9 @@ checksum = "f84267b20a16ea918e43c6a88433c2d54fa145c92a811b5b047ccbe153674483" [[package]] name = "potential_utf" -version = "0.1.2" +version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e5a7c30837279ca13e7c867e9e40053bc68740f988cb07f7ca6df43cc734b585" +checksum = "84df19adbe5b5a0782edcab45899906947ab039ccf4573713735ee7de1e6b08a" dependencies = [ "zerovec", ] @@ -1979,18 +1991,18 @@ dependencies = [ [[package]] name = "proc-macro2" -version = "1.0.97" +version = "1.0.101" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d61789d7719defeb74ea5fe81f2fdfdbd28a803847077cecce2ff14e1472f6f1" +checksum = "89ae43fd86e4158d6db51ad8e2b80f313af9cc74f5c0e03ccb87de09998732de" dependencies = [ "unicode-ident", ] [[package]] name = "quinn" -version = "0.11.8" +version = "0.11.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "626214629cda6781b6dc1d316ba307189c85ba657213ce642d9c77670f8202c8" +checksum = "b9e20a958963c291dc322d98411f541009df2ced7b5a4f2bd52337638cfccf20" dependencies = [ "bytes", "cfg_aliases", @@ -1999,8 +2011,8 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls", - "socket2 0.5.10", - "thiserror 2.0.14", + "socket2 0.6.0", + "thiserror 2.0.16", "tokio", "tracing", "web-time", @@ -2008,9 +2020,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.12" +version = "0.11.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "49df843a9161c85bb8aae55f101bc0bac8bcafd637a620d9122fd7e0b2f7422e" +checksum = "f1906b49b0c3bc04b5fe5d86a77925ae6524a19b816ae38ce1e426255f1d8a31" dependencies = [ "bytes", "getrandom 0.3.3", @@ -2021,7 +2033,7 @@ dependencies = [ "rustls", "rustls-pki-types", "slab", - "thiserror 2.0.14", + "thiserror 2.0.16", "tinyvec", "tracing", "web-time", @@ -2029,16 +2041,16 @@ dependencies = [ [[package]] name = "quinn-udp" -version = "0.5.13" +version = "0.5.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fcebb1209ee276352ef14ff8732e24cc2b02bbac986cd74a4c81bcb2f9881970" +checksum = "addec6a0dcad8a8d96a771f815f0eaf55f9d1805756410b39f5fa81332574cbd" dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.5.10", + "socket2 0.6.0", "tracing", - "windows-sys 0.59.0", + "windows-sys 0.60.2", ] [[package]] @@ -2115,6 +2127,31 @@ dependencies = [ "getrandom 0.3.3", ] +[[package]] +name = "redis" +version = "0.32.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15965fbccb975c38a08a68beca6bdb57da9081cd0859417c5975a160d968c3cb" +dependencies = [ + "arc-swap", + "backon", + "bytes", + "cfg-if", + "combine", + "futures-channel", + "futures-util", + "itoa", + "num-bigint", + "percent-encoding", + "pin-project-lite", + "ryu", + "sha1_smol", + "socket2 0.6.0", + "tokio", + "tokio-util", + "url", +] + [[package]] name = "redox_syscall" version = "0.5.17" @@ -2126,47 +2163,32 @@ dependencies = [ [[package]] name = "regex" -version = "1.11.2" +version = "1.11.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "23d7fd106d8c02486a8d64e778353d1cffe08ce79ac2e82f540c86d0facf6912" +checksum = "8b5288124840bee7b386bc413c487869b360b2b4ec421ea56425128692f2a82c" dependencies = [ "aho-corasick", "memchr", - "regex-automata 0.4.9", - "regex-syntax 0.8.5", -] - -[[package]] -name = "regex-automata" -version = "0.1.10" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6c230d73fb8d8c1b9c0b3135c5142a8acee3a0558fb8db5cf1cb65f8d7862132" -dependencies = [ - "regex-syntax 0.6.29", + "regex-automata", + "regex-syntax", ] [[package]] name = "regex-automata" -version = "0.4.9" +version = "0.4.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "809e8dc61f6de73b46c85f4c96486310fe304c434cfa43669d7b40f711150908" +checksum = "833eb9ce86d40ef33cb1306d8accf7bc8ec2bfea4355cbdebb3df68b40925cad" dependencies = [ "aho-corasick", "memchr", - "regex-syntax 0.8.5", + "regex-syntax", ] [[package]] name = "regex-syntax" -version = "0.6.29" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f162c6dd7b008981e4d40210aca20b4bd0f9b60ca9271061b07f78537722f2e1" - -[[package]] -name = "regex-syntax" -version = "0.8.5" +version = "0.8.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2b15c43186be67a4fd63bee50d0303afffcef381492ebe2c5d87f324e1b8815c" +checksum = "caf4aa5b0f434c91fe5c7f1ecb6a5ece2130b02ad2a590589dda5146df959001" [[package]] name = "reqwest" @@ -2245,9 +2267,9 @@ dependencies = [ [[package]] name = "resolv-conf" -version = "0.7.4" +version = "0.7.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "95325155c684b1c89f7765e30bc1c42e4a6da51ca513615660cb8a62ef9a88e3" +checksum = "6b3789b30bd25ba102de4beabd95d21ac45b69b1be7d14522bab988c526d6799" [[package]] name = "rfc6979" @@ -2316,22 +2338,22 @@ dependencies = [ [[package]] name = "rustix" -version = "1.0.8" +version = "1.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "11181fbabf243db407ef8df94a6ce0b2f9a733bd8be4ad02b4eda9602296cac8" +checksum = "cd15f8a2c5551a84d56efdc1cd049089e409ac19a3072d5037a17fd70719ff3e" dependencies = [ "bitflags", "errno", "libc", "linux-raw-sys", - "windows-sys 0.60.2", + "windows-sys 0.61.1", ] [[package]] name = "rustls" -version = "0.23.31" +version = "0.23.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c0ebcbd2f03de0fc1122ad9bb24b127a5a6cd51d72604a3f3c50ac459762b6cc" +checksum = "cd3c25631629d034ce7cd9940adc9d45762d46de2b0f57193c4443b92c6d4d40" dependencies = [ "once_cell", "ring", @@ -2350,7 +2372,7 @@ dependencies = [ "openssl-probe", "rustls-pki-types", "schannel", - "security-framework 3.3.0", + "security-framework 3.5.0", ] [[package]] @@ -2365,9 +2387,9 @@ dependencies = [ [[package]] name = "rustls-webpki" -version = "0.103.4" +version = "0.103.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0a17884ae0c1b773f1ccd2bd4a8c72f16da897310a98b0e84bf349ad5ead92fc" +checksum = "8572f3c2cb9934231157b45499fc41e1f58c589fdfb81a844ba873265e80f8eb" dependencies = [ "ring", "rustls-pki-types", @@ -2388,19 +2410,13 @@ checksum = "28d3b2b1366ec20994f1fd18c3c594f05c5dd4bc44d8bb0c1c632c8d6829481f" [[package]] name = "schannel" -version = "0.1.27" +version = "0.1.28" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f29ebaa345f945cec9fbbc532eb307f0fdad8161f281b6369539c8d84876b3d" +checksum = "891d81b926048e76efe18581bf793546b4c0eaf8448d72be8de2bbee5fd166e1" dependencies = [ - "windows-sys 0.59.0", + "windows-sys 0.61.1", ] -[[package]] -name = "scoped-tls" -version = "1.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e1cf6437eb19a8f4a6cc0f7dca544973b0b78843adbfeb3683d1a94a0024a294" - [[package]] name = "scopeguard" version = "1.2.0" @@ -2437,9 +2453,9 @@ dependencies = [ [[package]] name = "security-framework" -version = "3.3.0" +version = "3.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "80fb1d92c5028aa318b4b8bd7302a5bfcf48be96a37fc6fc790f806b0004ee0c" +checksum = "cc198e42d9b7510827939c9a15f5062a0c913f3371d765977e586d2fe6c16f4a" dependencies = [ "bitflags", "core-foundation 0.10.1", @@ -2450,9 +2466,9 @@ dependencies = [ [[package]] name = "security-framework-sys" -version = "2.14.0" +version = "2.15.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "49db231d56a190491cb4aeda9527f1ad45345af50b0851622a7adb8c03b01c32" +checksum = "cc1f0cbffaac4852523ce30d8bd3c5cdc873501d96ff467ca09b6767bb8cd5c0" dependencies = [ "core-foundation-sys", "libc", @@ -2460,33 +2476,44 @@ dependencies = [ [[package]] name = "semver" -version = "1.0.26" +version = "1.0.27" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "56e6fa9c48d24d85fb3de5ad847117517440f6beceb7798af16b4a87d616b8d0" +checksum = "d767eb0aabc880b29956c35734170f26ed551a859dbd361d140cdbeca61ab1e2" [[package]] name = "serde" -version = "1.0.219" +version = "1.0.227" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5f0e2c6ed6606019b4e29e69dbaba95b11854410e5347d525002456dbbb786b6" +checksum = "80ece43fc6fbed4eb5392ab50c07334d3e577cbf40997ee896fe7af40bba4245" dependencies = [ + "serde_core", "serde_derive", ] [[package]] name = "serde_bytes" -version = "0.11.17" +version = "0.11.19" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8437fd221bde2d4ca316d61b90e337e9e702b3820b87d63caa9ba6c02bd06d96" +checksum = "a5d440709e79d88e51ac01c4b72fc6cb7314017bb7da9eeff678aa94c10e3ea8" dependencies = [ "serde", + "serde_core", +] + +[[package]] +name = "serde_core" +version = "1.0.227" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a576275b607a2c86ea29e410193df32bc680303c82f31e275bbfcafe8b33be5" +dependencies = [ + "serde_derive", ] [[package]] name = "serde_derive" -version = "1.0.219" +version = "1.0.227" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5b0276cf7f2c73365f7157c8123c21cd9a50fbbd844757af28ca1f5925fc2a00" +checksum = "51e694923b8824cf0e9b382adf0f60d4e05f348f357b38833a3fa5ed7c2ede04" dependencies = [ "proc-macro2", "quote", @@ -2495,22 +2522,22 @@ dependencies = [ [[package]] name = "serde_html_form" -version = "0.2.7" +version = "0.2.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9d2de91cf02bbc07cde38891769ccd5d4f073d22a40683aa4bc7a95781aaa2c4" +checksum = "b2f2d7ff8a2140333718bb329f5c40fc5f0865b84c426183ce14c97d2ab8154f" dependencies = [ "form_urlencoded", "indexmap", "itoa", "ryu", - "serde", + "serde_core", ] [[package]] name = "serde_ipld_dagcbor" -version = "0.6.3" +version = "0.6.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "99600723cf53fb000a66175555098db7e75217c415bdd9a16a65d52a19dcc4fc" +checksum = "46182f4f08349a02b45c998ba3215d3f9de826246ba02bb9dddfe9a2a2100778" dependencies = [ "cbor4ii", "ipld-core", @@ -2520,24 +2547,26 @@ dependencies = [ [[package]] name = "serde_json" -version = "1.0.142" +version = "1.0.145" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "030fedb782600dcbd6f02d479bf0d817ac3bb40d644745b769d6a96bc3afc5a7" +checksum = "402a6f66d8c709116cf22f558eab210f5a50187f702eb4d7e5ef38d9a7f1c79c" dependencies = [ "itoa", "memchr", "ryu", "serde", + "serde_core", ] [[package]] name = "serde_path_to_error" -version = "0.1.17" +version = "0.1.20" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "59fab13f937fa393d08645bf3a84bdfe86e296747b506ada67bb15f10f218b2a" +checksum = "10a9ff822e371bb5403e391ecd83e182e0e77ba7f6fe0160b795797109d1b457" dependencies = [ "itoa", "serde", + "serde_core", ] [[package]] @@ -2573,6 +2602,12 @@ dependencies = [ "digest", ] +[[package]] +name = "sha1_smol" +version = "1.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbfa15b3dddfee50a0fff136974b3e1bde555604ba463834a7eb7deb6417705d" + [[package]] name = "sha2" version = "0.10.9" @@ -2646,6 +2681,7 @@ dependencies = [ "chrono", "dotenvy", "futures-util", + "redis", "regex", "reqwest", "reqwest-chain", @@ -2674,7 +2710,7 @@ dependencies = [ "regex", "serde", "serde_json", - "thiserror 2.0.14", + "thiserror 2.0.16", "unicode-segmentation", ] @@ -2756,7 +2792,7 @@ dependencies = [ "futures-intrusive", "futures-io", "futures-util", - "hashbrown", + "hashbrown 0.15.5", "hashlink", "indexmap", "log", @@ -2769,7 +2805,7 @@ dependencies = [ "serde_json", "sha2", "smallvec", - "thiserror 2.0.14", + "thiserror 2.0.16", "tokio", "tokio-stream", "tracing", @@ -2854,7 +2890,7 @@ dependencies = [ "smallvec", "sqlx-core", "stringprep", - "thiserror 2.0.14", + "thiserror 2.0.16", "tracing", "uuid", "whoami", @@ -2893,7 +2929,7 @@ dependencies = [ "smallvec", "sqlx-core", "stringprep", - "thiserror 2.0.14", + "thiserror 2.0.16", "tracing", "uuid", "whoami", @@ -2919,7 +2955,7 @@ dependencies = [ "serde", "serde_urlencoded", "sqlx-core", - "thiserror 2.0.14", + "thiserror 2.0.16", "tracing", "url", "uuid", @@ -3048,15 +3084,15 @@ checksum = "7b2093cf4c8eb1e67749a6762251bc9cd836b6fc171623bd0a9d324d37af2417" [[package]] name = "tempfile" -version = "3.20.0" +version = "3.23.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e8a64e3985349f2441a1a9ef0b853f869006c3855f2cda6862a94d26ebb9d6a1" +checksum = "2d31c77bdf42a745371d260a26ca7163f1e0924b64afa0b688e61b5a9fa02f16" dependencies = [ "fastrand", "getrandom 0.3.3", "once_cell", "rustix", - "windows-sys 0.59.0", + "windows-sys 0.61.1", ] [[package]] @@ -3070,11 +3106,11 @@ dependencies = [ [[package]] name = "thiserror" -version = "2.0.14" +version = "2.0.16" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0b0949c3a6c842cbde3f1686d6eea5a010516deb7085f79db747562d4102f41e" +checksum = "3467d614147380f2e4e374161426ff399c91084acd2363eaf549172b3d5e60c0" dependencies = [ - "thiserror-impl 2.0.14", + "thiserror-impl 2.0.16", ] [[package]] @@ -3090,9 +3126,9 @@ dependencies = [ [[package]] name = "thiserror-impl" -version = "2.0.14" +version = "2.0.16" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cc5b44b4ab9c2fdd0e0512e6bece8388e214c0749f5862b114cc5b7a25daf227" +checksum = "6c5e1be1c48b9172ee610da68fd9cd2770e7a4056cb3fc98710ee6906f0c7960" dependencies = [ "proc-macro2", "quote", @@ -3120,9 +3156,9 @@ dependencies = [ [[package]] name = "tinyvec" -version = "1.9.0" +version = "1.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "09b3661f17e86524eccd4371ab0429194e0d7c008abb45f7a7495b1719463c71" +checksum = "bfa5fdc3bce6191a1dbc8c02d5c8bffcf557bafa17c124c5264a458f1b0613fa" dependencies = [ "tinyvec_macros", ] @@ -3176,9 +3212,9 @@ dependencies = [ [[package]] name = "tokio-rustls" -version = "0.26.2" +version = "0.26.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8e727b36a1a0e8b74c376ac2211e40c2c8af09fb4013c60d910495810f008e9b" +checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" dependencies = [ "rustls", "tokio", @@ -3335,14 +3371,14 @@ dependencies = [ [[package]] name = "tracing-subscriber" -version = "0.3.19" +version = "0.3.20" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e8189decb5ac0fa7bc8b96b7cb9b2701d60d48805aca84a238004d665fcc4008" +checksum = "2054a14f5307d601f88daf0553e1cbf472acc4f2c51afab632431cdcd72124d5" dependencies = [ "matchers", "nu-ansi-term", "once_cell", - "regex", + "regex-automata", "sharded-slab", "smallvec", "thread_local", @@ -3405,9 +3441,9 @@ checksum = "5c1cb5db39152898a79168971543b1cb5020dff7fe43c8dc468b0885f5e29df5" [[package]] name = "unicode-ident" -version = "1.0.18" +version = "1.0.19" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5a5f39404a5da50712a4c1eecf25e90dd62b613502b7e925fd4e4d19b5c96512" +checksum = "f63a545481291138910575129486daeaf8ac54aee4387fe7906919f7830c7d9d" [[package]] name = "unicode-normalization" @@ -3444,13 +3480,14 @@ checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" [[package]] name = "url" -version = "2.5.4" +version = "2.5.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "32f8b686cadd1473f4bd0117a5d28d36b1ade384ea9b5069a1c40aefed7fda60" +checksum = "08bc136a29a3d1758e07a9cca267be308aeebf5cfd5a10f3f67ab2097683ef5b" dependencies = [ "form_urlencoded", "idna", "percent-encoding", + "serde", ] [[package]] @@ -3473,9 +3510,9 @@ checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be" [[package]] name = "uuid" -version = "1.18.0" +version = "1.18.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f33196643e165781c20a5ead5582283a7dacbb87855d867fbc2df3f81eddc1be" +checksum = "2f87b8aa10b915a06587d0dec516c282ff295b475d94abf425d62b57710070a2" dependencies = [ "getrandom 0.3.3", "js-sys", @@ -3518,11 +3555,20 @@ checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" [[package]] name = "wasi" -version = "0.14.2+wasi-0.2.4" +version = "0.14.7+wasi-0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "883478de20367e224c0090af9cf5f9fa85bed63a95c1abf3afc5c083ebc06e8c" +dependencies = [ + "wasip2", +] + +[[package]] +name = "wasip2" +version = "1.0.1+wasi-0.2.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9683f9a5a998d873c0d21fcbe3c083009670149a8fab228644b8bd36b2c48cb3" +checksum = "0562428422c63773dad2c345a1882263bbf4d65cf3f42e90921f787ef5ad58e7" dependencies = [ - "wit-bindgen-rt", + "wit-bindgen", ] [[package]] @@ -3533,21 +3579,22 @@ checksum = "b8dad83b4f25e74f184f64c43b150b91efe7647395b42289f38e50566d82855b" [[package]] name = "wasm-bindgen" -version = "0.2.100" +version = "0.2.104" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1edc8929d7499fc4e8f0be2262a241556cfc54a0bea223790e71446f2aab1ef5" +checksum = "c1da10c01ae9f1ae40cbfac0bac3b1e724b320abfcf52229f80b547c0d250e2d" dependencies = [ "cfg-if", "once_cell", "rustversion", "wasm-bindgen-macro", + "wasm-bindgen-shared", ] [[package]] name = "wasm-bindgen-backend" -version = "0.2.100" +version = "0.2.104" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2f0a0651a5c2bc21487bde11ee802ccaf4c51935d0d3d42a6101f98161700bc6" +checksum = "671c9a5a66f49d8a47345ab942e2cb93c7d1d0339065d4f8139c486121b43b19" dependencies = [ "bumpalo", "log", @@ -3559,9 +3606,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-futures" -version = "0.4.50" +version = "0.4.54" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "555d470ec0bc3bb57890405e5d4322cc9ea83cebb085523ced7be4144dac1e61" +checksum = "7e038d41e478cc73bae0ff9b36c60cff1c98b8f38f8d7e8061e79ee63608ac5c" dependencies = [ "cfg-if", "js-sys", @@ -3572,9 +3619,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro" -version = "0.2.100" +version = "0.2.104" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7fe63fc6d09ed3792bd0897b314f53de8e16568c2b3f7982f468c0bf9bd0b407" +checksum = "7ca60477e4c59f5f2986c50191cd972e3a50d8a95603bc9434501cf156a9a119" dependencies = [ "quote", "wasm-bindgen-macro-support", @@ -3582,9 +3629,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro-support" -version = "0.2.100" +version = "0.2.104" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8ae87ea40c9f689fc23f209965b6fb8a99ad69aeeb0231408be24920604395de" +checksum = "9f07d2f20d4da7b26400c9f4a0511e6e0345b040694e8a75bd41d578fa4421d7" dependencies = [ "proc-macro2", "quote", @@ -3595,9 +3642,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-shared" -version = "0.2.100" +version = "0.2.104" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1a05d73b933a847d6cccdda8f838a22ff101ad9bf93e33684f39c1f5f0eece3d" +checksum = "bad67dc8b2a1a6e5448428adec4c3e84c43e561d8c9ee8a9e5aabeb193ec41d1" dependencies = [ "unicode-ident", ] @@ -3617,9 +3664,9 @@ dependencies = [ [[package]] name = "web-sys" -version = "0.3.77" +version = "0.3.81" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "33b6dd2ef9186f1f2072e409e99cd22a975331a6b3591b12c764e0e55c60d5d2" +checksum = "9367c417a924a74cae129e6a2ae3b47fabb1f8995595ab474029da749a8be120" dependencies = [ "js-sys", "wasm-bindgen", @@ -3669,79 +3716,24 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dd7cf3379ca1aac9eea11fba24fd7e315d621f8dfe35c8d7d2be8b793726e07d" -[[package]] -name = "winapi" -version = "0.3.9" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5c839a674fcd7a98952e593242ea400abe93992746761e38641405d28b00f419" -dependencies = [ - "winapi-i686-pc-windows-gnu", - "winapi-x86_64-pc-windows-gnu", -] - -[[package]] -name = "winapi-i686-pc-windows-gnu" -version = "0.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" - -[[package]] -name = "winapi-x86_64-pc-windows-gnu" -version = "0.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f" - -[[package]] -name = "windows" -version = "0.61.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9babd3a767a4c1aef6900409f85f5d53ce2544ccdfaa86dad48c91782c6d6893" -dependencies = [ - "windows-collections", - "windows-core", - "windows-future", - "windows-link", - "windows-numerics", -] - -[[package]] -name = "windows-collections" -version = "0.2.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3beeceb5e5cfd9eb1d76b381630e82c4241ccd0d27f1a39ed41b2760b255c5e8" -dependencies = [ - "windows-core", -] - [[package]] name = "windows-core" -version = "0.61.2" +version = "0.62.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c0fdd3ddb90610c7638aa2b3a3ab2904fb9e5cdbecc643ddb3647212781c4ae3" +checksum = "6844ee5416b285084d3d3fffd743b925a6c9385455f64f6d4fa3031c4c2749a9" dependencies = [ "windows-implement", "windows-interface", - "windows-link", - "windows-result", - "windows-strings", -] - -[[package]] -name = "windows-future" -version = "0.2.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fc6a41e98427b19fe4b73c550f060b59fa592d7d686537eebf9385621bfbad8e" -dependencies = [ - "windows-core", - "windows-link", - "windows-threading", + "windows-link 0.2.0", + "windows-result 0.4.0", + "windows-strings 0.5.0", ] [[package]] name = "windows-implement" -version = "0.60.0" +version = "0.60.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a47fddd13af08290e67f4acabf4b459f647552718f683a7b415d290ac744a836" +checksum = "edb307e42a74fb6de9bf3a02d9712678b22399c87e6fa869d6dfcd8c1b7754e0" dependencies = [ "proc-macro2", "quote", @@ -3750,9 +3742,9 @@ dependencies = [ [[package]] name = "windows-interface" -version = "0.59.1" +version = "0.59.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bd9211b69f8dcdfa817bfd14bf1c97c9188afa36f4750130fcdf3f400eca9fa8" +checksum = "c0abd1ddbc6964ac14db11c7213d6532ef34bd9aa042c2e5935f59d7908b46a5" dependencies = [ "proc-macro2", "quote", @@ -3766,14 +3758,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5e6ad25900d524eaabdbbb96d20b4311e1e7ae1699af4fb28c17ae66c80d798a" [[package]] -name = "windows-numerics" +name = "windows-link" version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9150af68066c4c5c07ddc0ce30421554771e528bde427614c61038bc2c92c2b1" -dependencies = [ - "windows-core", - "windows-link", -] +checksum = "45e46c0661abb7180e7b9c281db115305d49ca1709ab8242adf09666d2173c65" [[package]] name = "windows-registry" @@ -3781,9 +3769,9 @@ version = "0.5.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5b8a9ed28765efc97bbc954883f4e6796c33a06546ebafacbabee9696967499e" dependencies = [ - "windows-link", - "windows-result", - "windows-strings", + "windows-link 0.1.3", + "windows-result 0.3.4", + "windows-strings 0.4.2", ] [[package]] @@ -3792,7 +3780,16 @@ version = "0.3.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "56f42bd332cc6c8eac5af113fc0c1fd6a8fd2aa08a0119358686e5160d0586c6" dependencies = [ - "windows-link", + "windows-link 0.1.3", +] + +[[package]] +name = "windows-result" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7084dcc306f89883455a206237404d3eaf961e5bd7e0f312f7c91f57eb44167f" +dependencies = [ + "windows-link 0.2.0", ] [[package]] @@ -3801,7 +3798,16 @@ version = "0.4.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "56e6c93f3a0c3b36176cb1327a4958a0353d5d166c2a35cb268ace15e91d3b57" dependencies = [ - "windows-link", + "windows-link 0.1.3", +] + +[[package]] +name = "windows-strings" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7218c655a553b0bed4426cf54b20d7ba363ef543b52d515b3e48d7fd55318dda" +dependencies = [ + "windows-link 0.2.0", ] [[package]] @@ -3837,7 +3843,16 @@ version = "0.60.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2f500e4d28234f72040990ec9d39e3a6b950f9f22d3dba18416c35882612bcb" dependencies = [ - "windows-targets 0.53.3", + "windows-targets 0.53.4", +] + +[[package]] +name = "windows-sys" +version = "0.61.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6f109e41dd4a3c848907eb83d5a42ea98b3769495597450cf6d153507b166f0f" +dependencies = [ + "windows-link 0.2.0", ] [[package]] @@ -3873,11 +3888,11 @@ dependencies = [ [[package]] name = "windows-targets" -version = "0.53.3" +version = "0.53.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d5fe6031c4041849d7c496a8ded650796e7b6ecc19df1a431c1a363342e5dc91" +checksum = "2d42b7b7f66d2a06854650af09cfdf8713e427a439c97ad65a6375318033ac4b" dependencies = [ - "windows-link", + "windows-link 0.2.0", "windows_aarch64_gnullvm 0.53.0", "windows_aarch64_msvc 0.53.0", "windows_i686_gnu 0.53.0", @@ -3888,15 +3903,6 @@ dependencies = [ "windows_x86_64_msvc 0.53.0", ] -[[package]] -name = "windows-threading" -version = "0.1.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b66463ad2e0ea3bbf808b7f1d371311c80e115c0b71d60efc142cafbcfb057a6" -dependencies = [ - "windows-link", -] - [[package]] name = "windows_aarch64_gnullvm" version = "0.48.5" @@ -4046,13 +4052,10 @@ dependencies = [ ] [[package]] -name = "wit-bindgen-rt" -version = "0.39.0" +name = "wit-bindgen" +version = "0.46.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6f42320e61fe2cfd34354ecb597f86f413484a798ba44a8ca1165c58d42da6c1" -dependencies = [ - "bitflags", -] +checksum = "f17a85883d4e6d00e8a97c586de764dabcc06133f7f1d55dce5cdc070ad7fe59" [[package]] name = "writeable" @@ -4086,18 +4089,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.26" +version = "0.8.27" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1039dd0d3c310cf05de012d8a39ff557cb0d23087fd44cad61df08fc31907a2f" +checksum = "0894878a5fa3edfd6da3f88c4805f4c8558e2b996227a3d864f47fe11e38282c" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.26" +version = "0.8.27" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9ecf5b4cc5364572d7f4c329661bcc82724222973f2cab6f050a4e5c22f75181" +checksum = "88d2b8d9c68ad2b9e4340d7832716a4d21a22a1154777ad56ea55c51a9cf3831" dependencies = [ "proc-macro2", "quote", @@ -4184,9 +4187,9 @@ dependencies = [ [[package]] name = "zstd-sys" -version = "2.0.15+zstd.1.5.7" +version = "2.0.16+zstd.1.5.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "eb81183ddd97d0c74cedf1d50d85c8d08c1b8b68ee863bdee9e706eedba1a237" +checksum = "91e19ebc2adc8f83e43039e79776e3fda8ca919132d68a1fed6a5faca2683748" dependencies = [ "cc", "pkg-config", diff --git a/api/Cargo.toml b/api/Cargo.toml index 3539be2..be46857 100644 --- a/api/Cargo.toml +++ b/api/Cargo.toml @@ -49,10 +49,10 @@ futures-util = "0.3" async-trait = "0.1" # AT Protocol client -atproto-client = "0.13.0" -atproto-identity = "0.13.0" -atproto-oauth = "0.13.0" -atproto-jetstream = "0.13.0" +atproto-client = "0.12.0" +atproto-identity = "0.12.0" +atproto-oauth = "0.12.0" +atproto-jetstream = "0.12.0" # Middleware for HTTP requests with retry logic @@ -62,3 +62,6 @@ reqwest-chain = "1.0.0" # Job queue sqlxmq = "0.6" regex = "1.11.2" + +# Redis for caching +redis = { version = "0.32", features = ["tokio-comp", "connection-manager"] } diff --git a/api/flake.nix b/api/flake.nix index f9dab1f..8bfc7af 100644 --- a/api/flake.nix +++ b/api/flake.nix @@ -73,7 +73,7 @@ # Fix for linker issues in Nix CC = "${pkgs.stdenv.cc}/bin/cc"; - + cargoExtraArgs = "--bin slices"; }; diff --git a/api/src/actor_resolver.rs b/api/src/actor_resolver.rs index 1f2d2a5..7b98a13 100644 --- a/api/src/actor_resolver.rs +++ b/api/src/actor_resolver.rs @@ -5,6 +5,10 @@ use atproto_identity::{ web::query as web_query, }; use thiserror::Error; +use std::sync::Arc; +use tokio::sync::Mutex; +use serde::{Serialize, Deserialize}; +use crate::cache::SliceCache; #[derive(Error, Debug)] pub enum ActorResolverError { @@ -18,7 +22,7 @@ pub enum ActorResolverError { InvalidSubject, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, Serialize, Deserialize)] pub struct ActorData { pub did: String, pub handle: Option, @@ -26,6 +30,66 @@ pub struct ActorData { } pub async fn resolve_actor_data(client: &Client, did: &str) -> Result { + resolve_actor_data_cached(client, did, None).await +} + +pub async fn resolve_actor_data_cached( + client: &Client, + did: &str, + cache: Option>> +) -> Result { + // Try cache first if provided + if let Some(cache) = &cache { + let cached_result = { + let mut cache_lock = cache.lock().await; + cache_lock.get_cached_did_resolution(did).await + }; + + if let Ok(Some(actor_data_value)) = cached_result + && let Ok(actor_data) = serde_json::from_value::(actor_data_value) { + return Ok(actor_data); + } + } + + // Cache miss - resolve from PLC/web + let actor_data = resolve_actor_data_impl(client, did).await?; + + // Cache the result if cache is provided + if let Some(cache) = &cache + && let Ok(actor_data_value) = serde_json::to_value(&actor_data) { + let mut cache_lock = cache.lock().await; + let _ = cache_lock.cache_did_resolution(did, &actor_data_value).await; + } + + Ok(actor_data) +} + +pub async fn resolve_actor_data_with_retry( + client: &Client, + did: &str, + cache: Option>>, + invalidate_cache_on_retry: bool +) -> Result { + match resolve_actor_data_cached(client, did, cache.clone()).await { + Ok(actor_data) => Ok(actor_data), + Err(e) => { + // If we should invalidate cache on retry and we have a cache + if invalidate_cache_on_retry { + if let Some(cache) = &cache { + let mut cache_lock = cache.lock().await; + let _ = cache_lock.invalidate_did_resolution(did).await; + } + + // Retry once with fresh resolution + resolve_actor_data_cached(client, did, cache).await + } else { + Err(e) + } + } + } +} + +async fn resolve_actor_data_impl(client: &Client, did: &str) -> Result { let (pds_url, handle) = match parse_input(did) { Ok(InputType::Plc(did_str)) => { match plc_query(client, "plc.directory", &did_str).await { diff --git a/api/src/api/xrpc_dynamic.rs b/api/src/api/xrpc_dynamic.rs index 9dfd627..c3500b7 100644 --- a/api/src/api/xrpc_dynamic.rs +++ b/api/src/api/xrpc_dynamic.rs @@ -11,7 +11,7 @@ use chrono::Utc; use serde::Deserialize; use crate::AppState; -use crate::auth::{extract_bearer_token, get_atproto_auth_for_user, verify_oauth_token}; +use crate::auth::{extract_bearer_token, get_atproto_auth_for_user_cached, verify_oauth_token_cached}; use crate::models::{ IndexedRecord, Record, SliceRecordsOutput, SliceRecordsParams, SortField, WhereCondition, }; @@ -526,12 +526,12 @@ async fn dynamic_collection_create_impl( ) -> Result, (StatusCode, Json)> { // Extract and verify OAuth token let token = extract_bearer_token(&headers).map_err(status_to_error_response)?; - let user_info = verify_oauth_token(&token, &state.config.auth_base_url) + let user_info = verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())) .await .map_err(status_to_error_response)?; - // Get AT Protocol DPoP auth and PDS URL - let (dpop_auth, pds_url) = get_atproto_auth_for_user(&token, &state.config.auth_base_url) + // Get AT Protocol DPoP auth and PDS URL (with caching) + let (dpop_auth, pds_url) = get_atproto_auth_for_user_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())) .await .map_err(status_to_error_response)?; @@ -644,12 +644,12 @@ async fn dynamic_collection_update_impl( ) -> Result, (StatusCode, Json)> { // Extract and verify OAuth token let token = extract_bearer_token(&headers).map_err(status_to_error_response)?; - let user_info = verify_oauth_token(&token, &state.config.auth_base_url) + let user_info = verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())) .await .map_err(status_to_error_response)?; - // Get AT Protocol DPoP auth and PDS URL - let (dpop_auth, pds_url) = get_atproto_auth_for_user(&token, &state.config.auth_base_url) + // Get AT Protocol DPoP auth and PDS URL (with caching) + let (dpop_auth, pds_url) = get_atproto_auth_for_user_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())) .await .map_err(status_to_error_response)?; @@ -762,12 +762,12 @@ async fn dynamic_collection_delete_impl( ) -> Result, (StatusCode, Json)> { // Extract and verify OAuth token let token = extract_bearer_token(&headers).map_err(status_to_error_response)?; - let user_info = verify_oauth_token(&token, &state.config.auth_base_url) + let user_info = verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())) .await .map_err(status_to_error_response)?; - // Get AT Protocol DPoP auth and PDS URL - let (dpop_auth, pds_url) = get_atproto_auth_for_user(&token, &state.config.auth_base_url) + // Get AT Protocol DPoP auth and PDS URL (with caching) + let (dpop_auth, pds_url) = get_atproto_auth_for_user_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())) .await .map_err(status_to_error_response)?; diff --git a/api/src/auth.rs b/api/src/auth.rs index 62a9d83..213aff8 100644 --- a/api/src/auth.rs +++ b/api/src/auth.rs @@ -3,6 +3,9 @@ use serde::{Deserialize, Serialize}; use atproto_client::client::DPoPAuth; use atproto_identity::key::KeyData; use atproto_oauth::jwk::WrappedJsonWebKey; +use std::sync::Arc; +use tokio::sync::Mutex; +use crate::cache::SliceCache; #[derive(Serialize, Deserialize, Debug)] pub struct UserInfoResponse { @@ -10,6 +13,13 @@ pub struct UserInfoResponse { pub did: Option, } +#[derive(Serialize, Deserialize, Debug, Clone)] +struct CachedSession { + pds_url: String, + atproto_access_token: String, + dpop_jwk: serde_json::Value, +} + // Extract bearer token from Authorization header pub fn extract_bearer_token(headers: &HeaderMap) -> Result { let auth_header = headers @@ -26,7 +36,29 @@ pub fn extract_bearer_token(headers: &HeaderMap) -> Result { } // Verify OAuth token with auth server -pub async fn verify_oauth_token(token: &str, auth_base_url: &str) -> Result { + +// Verify OAuth token with auth server with optional caching +pub async fn verify_oauth_token_cached( + token: &str, + auth_base_url: &str, + cache: Option>>, +) -> Result { + + // Try cache first if provided + if let Some(cache) = &cache { + let cached_result = { + let mut cache_lock = cache.lock().await; + cache_lock.get_cached_oauth_userinfo(token).await + }; + + if let Ok(Some(user_info_value)) = cached_result { + let user_info: UserInfoResponse = serde_json::from_value(user_info_value) + .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; + return Ok(user_info); + } + } + + // Cache miss - verify with auth server let client = reqwest::Client::new(); let userinfo_url = format!("{}/oauth/userinfo", auth_base_url); @@ -37,7 +69,6 @@ pub async fn verify_oauth_token(token: &str, auth_base_url: &str) -> Result Result>>, ) -> Result<(DPoPAuth, String), StatusCode> { - // First get session info from auth server + + // Try cache first if provided + if let Some(cache) = &cache { + let cached_result = { + let mut cache_lock = cache.lock().await; + cache_lock.get_cached_atproto_session(token).await + }; + + if let Ok(Some(session_value)) = cached_result { + let cached_session: CachedSession = serde_json::from_value(session_value) + .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; + + // Convert cached data back to DPoP auth + let dpop_jwk: WrappedJsonWebKey = serde_json::from_value(cached_session.dpop_jwk) + .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; + + let dpop_private_key_data = KeyData::try_from(dpop_jwk) + .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; + + let dpop_auth = DPoPAuth { + dpop_private_key_data, + oauth_access_token: cached_session.atproto_access_token, + }; + + return Ok((dpop_auth, cached_session.pds_url)); + } + } + + // Cache miss - fetch from auth server let client = reqwest::Client::new(); let session_url = format!("{}/api/atprotocol/session", auth_base_url); @@ -66,7 +136,6 @@ pub async fn get_atproto_auth_for_user( .await .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; - if !session_response.status().is_success() { return Err(StatusCode::UNAUTHORIZED); } @@ -76,38 +145,43 @@ pub async fn get_atproto_auth_for_user( .await .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; - // Extract PDS URL from session let pds_url = session_data["pds_endpoint"] .as_str() - .ok_or({ - StatusCode::INTERNAL_SERVER_ERROR - })? + .ok_or(StatusCode::INTERNAL_SERVER_ERROR)? .to_string(); - // Extract AT Protocol access token from session data let atproto_access_token = session_data["access_token"] .as_str() - .ok_or({ - StatusCode::INTERNAL_SERVER_ERROR - })? + .ok_or(StatusCode::INTERNAL_SERVER_ERROR)? .to_string(); - // Extract DPoP private key from session data - convert JWK to KeyData - let dpop_jwk: WrappedJsonWebKey = serde_json::from_value(session_data["dpop_jwk"].clone()) + let dpop_jwk_value = session_data["dpop_jwk"].clone(); + let dpop_jwk: WrappedJsonWebKey = serde_json::from_value(dpop_jwk_value.clone()) .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; - let dpop_private_key_data = KeyData::try_from(dpop_jwk) .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; - let dpop_auth = DPoPAuth { dpop_private_key_data, - oauth_access_token: atproto_access_token, + oauth_access_token: atproto_access_token.clone(), }; + // Cache the session data if cache is provided (5 minute TTL) + if let Some(cache) = &cache { + let cached_session = CachedSession { + pds_url: pds_url.clone(), + atproto_access_token, + dpop_jwk: dpop_jwk_value, + }; + let session_value = serde_json::to_value(&cached_session) + .map_err(|_e| StatusCode::INTERNAL_SERVER_ERROR)?; + let mut cache_lock = cache.lock().await; + let _ = cache_lock.cache_atproto_session(token, &session_value, 300).await; + } + Ok((dpop_auth, pds_url)) } \ No newline at end of file diff --git a/api/src/cache.rs b/api/src/cache.rs new file mode 100644 index 0000000..9054a0e --- /dev/null +++ b/api/src/cache.rs @@ -0,0 +1,455 @@ +use async_trait::async_trait; +use anyhow::Result; +use serde::{Serialize, Deserialize}; +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; +use tokio::sync::RwLock; +use tracing::{debug, info, warn}; +use std::time::{Duration, Instant}; + +/// Generic cache trait for different backend implementations +#[async_trait] +pub trait Cache: Send + Sync { + /// Get a value from cache + async fn get(&mut self, key: &str) -> Result> + where + T: for<'de> Deserialize<'de> + Send; + + /// Set a value in cache with optional TTL + async fn set(&mut self, key: &str, value: &T, ttl_seconds: Option) -> Result<()> + where + T: Serialize + Send + Sync; + + + /// Delete a key from cache + async fn delete(&mut self, key: &str) -> Result<()>; + + /// Set multiple key-value pairs + async fn set_multiple(&mut self, items: Vec<(&str, &T, Option)>) -> Result<()> + where + T: Serialize + Send + Sync; + + /// Test cache connection/health + async fn ping(&mut self) -> Result; + + /// Get cache info/statistics + async fn get_info(&mut self) -> Result; +} + +/// Cache entry type: (serialized_value, expiry) +type CacheEntry = (String, Option); + +/// In-memory cache implementation with TTL support +pub struct InMemoryCache { + data: Arc>>, + default_ttl_seconds: u64, +} + +impl InMemoryCache { + pub fn new(default_ttl_seconds: Option) -> Self { + Self { + data: Arc::new(RwLock::new(HashMap::new())), + default_ttl_seconds: default_ttl_seconds.unwrap_or(3600), + } + } + +} + +#[async_trait] +impl Cache for InMemoryCache { + async fn get(&mut self, key: &str) -> Result> + where + T: for<'de> Deserialize<'de> + Send, + { + let data = self.data.read().await; + + if let Some((serialized, expiry)) = data.get(key) { + // Check if expired + if let Some(exp) = expiry + && *exp <= Instant::now() { + debug!(cache_key = %key, "Cache entry expired"); + return Ok(None); + } + + match serde_json::from_str::(serialized) { + Ok(value) => { + // Cache hit - no logging needed + Ok(Some(value)) + } + Err(e) => { + warn!( + error = ?e, + cache_key = %key, + "Failed to deserialize cached value" + ); + Ok(None) + } + } + } else { + // Cache miss - no logging needed + Ok(None) + } + } + + async fn set(&mut self, key: &str, value: &T, ttl_seconds: Option) -> Result<()> + where + T: Serialize + Send + Sync, + { + let ttl = ttl_seconds.unwrap_or(self.default_ttl_seconds); + + match serde_json::to_string(value) { + Ok(serialized) => { + let expiry = if ttl > 0 { + Some(Instant::now() + Duration::from_secs(ttl)) + } else { + None // No expiry + }; + + let mut data = self.data.write().await; + data.insert(key.to_string(), (serialized, expiry)); + + debug!( + cache_key = %key, + ttl_seconds = ttl, + "Cached value in memory" + ); + Ok(()) + } + Err(e) => { + warn!( + error = ?e, + cache_key = %key, + "Failed to serialize value for caching" + ); + Ok(()) + } + } + } + + + async fn delete(&mut self, key: &str) -> Result<()> { + let mut data = self.data.write().await; + data.remove(key); + debug!(cache_key = %key, "Deleted key from in-memory cache"); + Ok(()) + } + + async fn set_multiple(&mut self, items: Vec<(&str, &T, Option)>) -> Result<()> + where + T: Serialize + Send + Sync, + { + if items.is_empty() { + return Ok(()); + } + + let mut data = self.data.write().await; + let mut success_count = 0; + + for (key, value, ttl) in &items { + match serde_json::to_string(value) { + Ok(serialized) => { + let ttl_to_use = ttl.unwrap_or(self.default_ttl_seconds); + let expiry = if ttl_to_use > 0 { + Some(Instant::now() + Duration::from_secs(ttl_to_use)) + } else { + None + }; + + data.insert(key.to_string(), (serialized, expiry)); + success_count += 1; + } + Err(e) => { + warn!( + error = ?e, + cache_key = %key, + "Failed to serialize value for bulk caching" + ); + } + } + } + + debug!( + items_count = success_count, + total_items = items.len(), + "Successfully bulk cached items in memory" + ); + Ok(()) + } + + async fn ping(&mut self) -> Result { + // Always healthy for in-memory cache + Ok(true) + } + + async fn get_info(&mut self) -> Result { + let data = self.data.read().await; + let now = Instant::now(); + + let mut total_entries = 0; + let mut expired_entries = 0; + + for (_, expiry) in data.values() { + total_entries += 1; + if let Some(exp) = expiry + && *exp <= now { + expired_entries += 1; + } + } + + Ok(format!( + "InMemoryCache: {} total entries, {} expired, {} active", + total_entries, + expired_entries, + total_entries - expired_entries + )) + } +} + +/// Cache backend enum to avoid dyn trait issues +pub enum CacheBackendImpl { + InMemory(InMemoryCache), + Redis(crate::redis_cache::RedisCache), +} + +impl CacheBackendImpl { + pub async fn get(&mut self, key: &str) -> Result> + where + T: for<'de> Deserialize<'de> + Send, + { + match self { + CacheBackendImpl::InMemory(cache) => cache.get(key).await, + CacheBackendImpl::Redis(cache) => cache.get(key).await, + } + } + + pub async fn set(&mut self, key: &str, value: &T, ttl_seconds: Option) -> Result<()> + where + T: Serialize + Send + Sync, + { + match self { + CacheBackendImpl::InMemory(cache) => cache.set(key, value, ttl_seconds).await, + CacheBackendImpl::Redis(cache) => cache.set(key, value, ttl_seconds).await, + } + } + + pub async fn delete(&mut self, key: &str) -> Result<()> { + match self { + CacheBackendImpl::InMemory(cache) => cache.delete(key).await, + CacheBackendImpl::Redis(cache) => cache.delete(key).await, + } + } + + pub async fn set_multiple(&mut self, items: Vec<(&str, &T, Option)>) -> Result<()> + where + T: Serialize + Send + Sync, + { + match self { + CacheBackendImpl::InMemory(cache) => cache.set_multiple(items).await, + CacheBackendImpl::Redis(cache) => cache.set_multiple(items).await, + } + } + + pub async fn ping(&mut self) -> Result { + match self { + CacheBackendImpl::InMemory(cache) => cache.ping().await, + CacheBackendImpl::Redis(cache) => cache.ping().await, + } + } + + pub async fn get_info(&mut self) -> Result { + match self { + CacheBackendImpl::InMemory(cache) => cache.get_info().await, + CacheBackendImpl::Redis(cache) => cache.get_info().await, + } + } +} + +/// Cache-specific helper methods for slice operations +pub struct SliceCache { + cache: CacheBackendImpl, +} + +impl SliceCache { + pub fn new(cache: CacheBackendImpl) -> Self { + Self { cache } + } + + /// Actor cache methods + pub async fn is_actor(&mut self, did: &str, slice_uri: &str) -> Result> { + let key = format!("actor:{}:{}", did, slice_uri); + self.cache.get::(&key).await + } + + pub async fn cache_actor_exists(&mut self, did: &str, slice_uri: &str) -> Result<()> { + let key = format!("actor:{}:{}", did, slice_uri); + self.cache.set(&key, &true, None).await + } + + pub async fn remove_actor(&mut self, did: &str, slice_uri: &str) -> Result<()> { + let key = format!("actor:{}:{}", did, slice_uri); + self.cache.delete(&key).await + } + + pub async fn preload_actors(&mut self, actors: Vec<(String, String)>) -> Result<()> { + if actors.is_empty() { + return Ok(()); + } + + let items: Vec<(String, bool, Option)> = actors + .into_iter() + .map(|(did, slice_uri)| { + (format!("actor:{}:{}", did, slice_uri), true, None) + }) + .collect(); + + let items_ref: Vec<(&str, &bool, Option)> = items + .iter() + .map(|(key, value, ttl)| (key.as_str(), value, *ttl)) + .collect(); + + self.cache.set_multiple(items_ref).await + } + + /// Lexicon cache methods + pub async fn cache_lexicons(&mut self, slice_uri: &str, lexicons: &Vec) -> Result<()> { + let key = format!("lexicons:{}", slice_uri); + let lexicons_ttl = 7200; // 2 hours for lexicons + self.cache.set(&key, lexicons, Some(lexicons_ttl)).await + } + + pub async fn get_lexicons(&mut self, slice_uri: &str) -> Result>> { + let key = format!("lexicons:{}", slice_uri); + self.cache.get::>(&key).await + } + + /// Domain cache methods + pub async fn cache_slice_domain(&mut self, slice_uri: &str, domain: &str) -> Result<()> { + let key = format!("domain:{}", slice_uri); + let domain_ttl = 14400; // 4 hours for domains + self.cache.set(&key, &domain.to_string(), Some(domain_ttl)).await + } + + pub async fn get_slice_domain(&mut self, slice_uri: &str) -> Result> { + let key = format!("domain:{}", slice_uri); + self.cache.get::(&key).await + } + + /// Collections cache methods + pub async fn cache_slice_collections(&mut self, slice_uri: &str, collections: &HashSet) -> Result<()> { + let key = format!("collections:{}", slice_uri); + let collections_ttl = 7200; // 2 hours for collections + self.cache.set(&key, collections, Some(collections_ttl)).await + } + + pub async fn get_slice_collections(&mut self, slice_uri: &str) -> Result>> { + let key = format!("collections:{}", slice_uri); + self.cache.get::>(&key).await + } + + /// Utility methods + pub async fn ping(&mut self) -> Result { + self.cache.ping().await + } + + pub async fn get_info(&mut self) -> Result { + self.cache.get_info().await + } + + /// Auth cache methods (5 minute TTL) + pub async fn get_cached_oauth_userinfo(&mut self, token: &str) -> Result> { + let key = format!("oauth_userinfo:{}", token); + self.cache.get(&key).await + } + + pub async fn cache_oauth_userinfo(&mut self, token: &str, userinfo: &serde_json::Value, ttl_seconds: u64) -> Result<()> { + let key = format!("oauth_userinfo:{}", token); + self.cache.set(&key, userinfo, Some(ttl_seconds)).await + } + + pub async fn get_cached_atproto_session(&mut self, token: &str) -> Result> { + let key = format!("atproto_session:{}", token); + self.cache.get(&key).await + } + + pub async fn cache_atproto_session(&mut self, token: &str, session: &serde_json::Value, ttl_seconds: u64) -> Result<()> { + let key = format!("atproto_session:{}", token); + self.cache.set(&key, session, Some(ttl_seconds)).await + } + + /// DID resolution cache methods (24 hour TTL - DIDs change infrequently) + pub async fn get_cached_did_resolution(&mut self, did: &str) -> Result> { + let key = format!("did_resolution:{}", did); + self.cache.get(&key).await + } + + pub async fn cache_did_resolution(&mut self, did: &str, actor_data: &serde_json::Value) -> Result<()> { + let key = format!("did_resolution:{}", did); + let ttl_seconds = 24 * 60 * 60; // 24 hours + self.cache.set(&key, actor_data, Some(ttl_seconds)).await + } + + pub async fn invalidate_did_resolution(&mut self, did: &str) -> Result<()> { + let key = format!("did_resolution:{}", did); + self.cache.delete(&key).await + } + + /// Generic get/set for custom caching needs + pub async fn get(&mut self, key: &str) -> Result> + where + T: for<'de> Deserialize<'de> + Send, + { + self.cache.get(key).await + } + + pub async fn set(&mut self, key: &str, value: &T, ttl_seconds: Option) -> Result<()> + where + T: Serialize + Send + Sync, + { + self.cache.set(key, value, ttl_seconds).await + } +} + +/// Cache backend configuration +#[derive(Debug, Clone)] +pub enum CacheBackend { + InMemory { ttl_seconds: Option }, + Redis { url: String, ttl_seconds: Option }, +} + +/// Cache factory for creating cache instances +pub struct CacheFactory; + +impl CacheFactory { + /// Create a cache instance based on configuration + pub async fn create_cache(backend: CacheBackend) -> Result { + match backend { + CacheBackend::InMemory { ttl_seconds } => { + let ttl_display = ttl_seconds.map(|t| format!("{}s", t)).unwrap_or_else(|| "default".to_string()); + info!("Creating in-memory cache with TTL: {}", ttl_display); + Ok(CacheBackendImpl::InMemory(InMemoryCache::new(ttl_seconds))) + } + CacheBackend::Redis { url, ttl_seconds } => { + info!("Attempting to create Redis cache at: {}", url); + match crate::redis_cache::RedisCache::new(&url, ttl_seconds).await { + Ok(redis_cache) => { + info!("✓ Created Redis cache successfully"); + Ok(CacheBackendImpl::Redis(redis_cache)) + } + Err(e) => { + warn!( + error = ?e, + "Failed to create Redis cache, falling back to in-memory" + ); + Ok(CacheBackendImpl::InMemory(InMemoryCache::new(ttl_seconds))) + } + } + } + } + } + + /// Create a SliceCache with the specified backend + pub async fn create_slice_cache(backend: CacheBackend) -> Result { + let cache = Self::create_cache(backend).await?; + Ok(SliceCache::new(cache)) + } +} \ No newline at end of file diff --git a/api/src/errors.rs b/api/src/errors.rs index ad756ac..d78a7f6 100644 --- a/api/src/errors.rs +++ b/api/src/errors.rs @@ -64,6 +64,9 @@ pub enum AppError { #[error("error-slices-app-8 Forbidden: {0}")] Forbidden(String), + + #[error("error-slices-app-9 Cache error: {0}")] + Cache(#[from] anyhow::Error), } impl From for AppError { @@ -99,6 +102,7 @@ impl IntoResponse for AppError { AppError::DatabaseConnection(e) => (StatusCode::INTERNAL_SERVER_ERROR, "InternalServerError", e.to_string()), AppError::Migration(e) => (StatusCode::INTERNAL_SERVER_ERROR, "InternalServerError", e.to_string()), AppError::ServerBind(e) => (StatusCode::INTERNAL_SERVER_ERROR, "InternalServerError", e.to_string()), + AppError::Cache(e) => (StatusCode::INTERNAL_SERVER_ERROR, "InternalServerError", e.to_string()), }; let body = Json(serde_json::json!({ diff --git a/api/src/jetstream.rs b/api/src/jetstream.rs index 3f32afc..00d59f0 100644 --- a/api/src/jetstream.rs +++ b/api/src/jetstream.rs @@ -2,10 +2,10 @@ use atproto_jetstream::{Consumer, ConsumerTaskConfig, EventHandler, JetstreamEve use async_trait::async_trait; use anyhow::Result; use chrono::Utc; -use std::collections::{HashMap, HashSet}; +use std::collections::HashSet; use std::sync::Arc; -use tokio::sync::RwLock; -use tracing::{error, info}; +use tokio::sync::{Mutex, RwLock}; +use tracing::{error, info, warn}; use reqwest::Client; use crate::actor_resolver::resolve_actor_data; @@ -14,29 +14,32 @@ use crate::jetstream_cursor::PostgresCursorHandler; use crate::models::{Record, Actor}; use crate::errors::SliceError; use crate::logging::{Logger, LogLevel}; +use crate::cache::{SliceCache, CacheFactory, CacheBackend}; pub struct JetstreamConsumer { consumer: Consumer, database: Database, http_client: Client, - slice_collections: Arc>>>, - slice_domains: Arc>>, - actor_cache: Arc>>, - slice_lexicons: Arc>>>, + actor_cache: Arc>, + lexicon_cache: Arc>, + domain_cache: Arc>, + collections_cache: Arc>, pub event_count: Arc, cursor_handler: Option>, + slices_list: Arc>>, } // Event handler that implements the EventHandler trait struct SliceEventHandler { database: Database, http_client: Client, - slice_collections: Arc>>>, - slice_domains: Arc>>, event_count: Arc, - actor_cache: Arc>>, - slice_lexicons: Arc>>>, + actor_cache: Arc>, + lexicon_cache: Arc>, + domain_cache: Arc>, + collections_cache: Arc>, cursor_handler: Option>, + slices_list: Arc>>, } #[async_trait] @@ -101,16 +104,143 @@ impl EventHandler for SliceEventHandler { } impl SliceEventHandler { + /// Check if DID is an actor for the given slice + async fn is_actor_cached(&self, did: &str, slice_uri: &str) -> Result, anyhow::Error> { + match self.actor_cache.lock().await.is_actor(did, slice_uri).await { + Ok(result) => Ok(result), + Err(e) => { + warn!( + error = ?e, + did = did, + slice_uri = slice_uri, + "Actor cache error" + ); + Ok(None) + } + } + } + + /// Cache that an actor exists + async fn cache_actor_exists(&self, did: &str, slice_uri: &str) { + if let Err(e) = self.actor_cache.lock().await.cache_actor_exists(did, slice_uri).await { + warn!( + error = ?e, + did = did, + slice_uri = slice_uri, + "Failed to cache actor exists" + ); + } + } + + /// Remove actor from cache + async fn remove_actor_from_cache(&self, did: &str, slice_uri: &str) { + if let Err(e) = self.actor_cache.lock().await.remove_actor(did, slice_uri).await { + warn!( + error = ?e, + did = did, + slice_uri = slice_uri, + "Failed to remove actor from cache" + ); + } + } + + /// Get slice collections from cache with database fallback + async fn get_slice_collections(&self, slice_uri: &str) -> Result>, anyhow::Error> { + // Try cache first + let cache_result = { + let mut cache = self.collections_cache.lock().await; + cache.get_slice_collections(slice_uri).await + }; + + match cache_result { + Ok(Some(collections)) => Ok(Some(collections)), + Ok(None) => { + // Cache miss - load from database + match self.database.get_slice_collections_list(slice_uri).await { + Ok(collections) => { + let collections_set: HashSet = collections.into_iter().collect(); + // Cache the result + let _ = self.collections_cache.lock().await.cache_slice_collections(slice_uri, &collections_set).await; + Ok(Some(collections_set)) + } + Err(e) => Err(e.into()) + } + } + Err(e) => Err(e) + } + } + + /// Get slice domain from cache with database fallback + async fn get_slice_domain(&self, slice_uri: &str) -> Result, anyhow::Error> { + // Try cache first + let cache_result = { + let mut cache = self.domain_cache.lock().await; + cache.get_slice_domain(slice_uri).await + }; + + match cache_result { + Ok(Some(domain)) => Ok(Some(domain)), + Ok(None) => { + // Cache miss - load from database + match self.database.get_slice_domain(slice_uri).await { + Ok(Some(domain)) => { + // Cache the result + let _ = self.domain_cache.lock().await.cache_slice_domain(slice_uri, &domain).await; + Ok(Some(domain)) + } + Ok(None) => Ok(None), + Err(e) => Err(e.into()) + } + } + Err(e) => Err(e) + } + } + + /// Get slice lexicons from cache with database fallback + async fn get_slice_lexicons(&self, slice_uri: &str) -> Result>, anyhow::Error> { + // Try cache first + let cache_result = { + let mut cache = self.lexicon_cache.lock().await; + cache.get_lexicons(slice_uri).await + }; + + match cache_result { + Ok(Some(lexicons)) => Ok(Some(lexicons)), + Ok(None) => { + // Cache miss - load from database + match self.database.get_lexicons_by_slice(slice_uri).await { + Ok(lexicons) if !lexicons.is_empty() => { + // Cache the result + let _ = self.lexicon_cache.lock().await.cache_lexicons(slice_uri, &lexicons).await; + Ok(Some(lexicons)) + } + Ok(_) => Ok(None), // Empty lexicons + Err(e) => Err(e.into()) + } + } + Err(e) => Err(e) + } + } async fn handle_commit_event( &self, did: &str, commit: atproto_jetstream::JetstreamEventCommit, ) -> Result<()> { - let slice_collections = self.slice_collections.read().await; - let slice_domains = self.slice_domains.read().await; - let slice_lexicons = self.slice_lexicons.read().await; - - for (slice_uri, collections) in slice_collections.iter() { + // Get all slices from cached list + let slices = self.slices_list.read().await.clone(); + + // Process each slice + for slice_uri in slices { + // Get collections for this slice (with caching) + let collections = match self.get_slice_collections(&slice_uri).await { + Ok(Some(collections)) => collections, + Ok(None) => continue, // No collections for this slice + Err(e) => { + error!("Failed to get collections for slice {}: {}", slice_uri, e); + continue; + } + }; + if collections.contains(&commit.collection) { // Special handling for network.slices.lexicon records // These should only be indexed to the slice specified in their JSON data @@ -125,72 +255,53 @@ impl SliceEventHandler { continue; } } - // Get the domain for this slice - let domain = match slice_domains.get(slice_uri) { - Some(d) => d, - None => continue, // No domain, skip + // Get the domain for this slice (with caching) + let domain = match self.get_slice_domain(&slice_uri).await { + Ok(Some(domain)) => domain, + Ok(None) => continue, // No domain, skip + Err(e) => { + error!("Failed to get domain for slice {}: {}", slice_uri, e); + continue; + } }; - + // Check if this is a primary collection (starts with slice domain) - let is_primary_collection = commit.collection.starts_with(domain); - + let is_primary_collection = commit.collection.starts_with(&domain); + // For external collections, check actor status BEFORE expensive validation if !is_primary_collection { - let cache_key = (did.to_string(), slice_uri.clone()); - let is_actor = { - let cache = self.actor_cache.read().await; - cache.get(&cache_key).copied() - }; - - let is_actor: Result = match is_actor { - Some(cached_result) => Ok(cached_result), - None => { + let is_actor = match self.is_actor_cached(did, &slice_uri).await { + Ok(Some(cached_result)) => cached_result, + Ok(None) => { // Cache miss means this DID is not an actor we've synced // For external collections, we only care about actors we've already added - // Don't cache negative results to avoid memory bloat - Ok(false) - } - }; - - match is_actor { - Ok(false) => { - // Not an actor - skip validation entirely for external collections - continue; - } - Ok(true) => { - // Actor found - continue to validation + false } Err(e) => { error!("Error checking actor status: {}", e); continue; } + }; + + if !is_actor { + // Not an actor - skip validation entirely for external collections + continue; } } - + // Get lexicons for validation (after actor check for external collections) - let lexicons = match slice_lexicons.get(slice_uri) { - Some(lexicons) => lexicons.clone(), - None => { - // Fallback: Try to load fresh lexicons from database for this slice - info!("No cached lexicons for slice {} - attempting database fallback", slice_uri); - match self.database.get_lexicons_by_slice(slice_uri).await { - Ok(fresh_lexicons) if !fresh_lexicons.is_empty() => { - info!("✓ Loaded fresh lexicons for slice {} from database", slice_uri); - // Cache the fresh lexicons for future use - { - let mut lexicons_cache = self.slice_lexicons.write().await; - lexicons_cache.insert(slice_uri.clone(), fresh_lexicons.clone()); - } - fresh_lexicons - } - _ => { - info!("No lexicons found for slice {} - skipping validation", slice_uri); - continue; - } - } + let lexicons = match self.get_slice_lexicons(&slice_uri).await { + Ok(Some(lexicons)) => lexicons, + Ok(None) => { + info!("No lexicons found for slice {} - skipping validation", slice_uri); + continue; + } + Err(e) => { + error!("Failed to get lexicons for slice {}: {}", slice_uri, e); + continue; } }; - + // Validate the record against the slice's lexicons let validation_result = match slices_lexicon::validate_record(lexicons.clone(), &commit.collection, commit.record.clone()) { Ok(_) => { @@ -198,62 +309,28 @@ impl SliceEventHandler { true } Err(e) => { - info!("Validation failed with cached validator for collection {} in slice {}: {} - trying database fallback", - commit.collection, slice_uri, e); - - // Try database fallback in case lexicons were updated - match self.database.get_lexicons_by_slice(slice_uri).await { - Ok(fresh_lexicons) if !fresh_lexicons.is_empty() => { - match slices_lexicon::validate_record(fresh_lexicons.clone(), &commit.collection, commit.record.clone()) { - Ok(_) => { - info!("✓ Record validated with fresh lexicons for collection {} in slice {}", - commit.collection, slice_uri); - // Update cache with fresh lexicons - { - let mut lexicons_cache = self.slice_lexicons.write().await; - lexicons_cache.insert(slice_uri.clone(), fresh_lexicons); - } - true - } - Err(fresh_e) => { - let message = format!("Validation failed for collection {} in slice {}", commit.collection, slice_uri); - error!("✗ {}: {}", message, fresh_e); - Logger::global().log_jetstream_with_slice(LogLevel::Warn, &message, Some(serde_json::json!({ - "collection": commit.collection, - "slice_uri": slice_uri, - "did": did - })), Some(slice_uri)); - false - } - } - } - Ok(_) => { - // Empty lexicons - skip logging as this is expected for many slices - false - } - Err(_) => { - // Database error - skip logging for missing lexicons - false - } - } + let message = format!("Validation failed for collection {} in slice {}", commit.collection, slice_uri); + error!("✗ {}: {}", message, e); + Logger::global().log_jetstream_with_slice(LogLevel::Warn, &message, Some(serde_json::json!({ + "collection": commit.collection, + "slice_uri": slice_uri, + "did": did + })), Some(&slice_uri)); + false } }; - + if !validation_result { continue; // Skip this slice if validation fails } - + if is_primary_collection { // Primary collection - ensure actor exists and index ALL records info!("✓ Primary collection {} for slice {} (domain: {}) - indexing record", commit.collection, slice_uri, domain); // Ensure actor exists for primary collections - let cache_key = (did.to_string(), slice_uri.clone()); - let is_cached = { - let cache = self.actor_cache.read().await; - cache.contains_key(&cache_key) - }; + let is_cached = matches!(self.is_actor_cached(did, &slice_uri).await, Ok(Some(_))); if !is_cached { // Actor not in cache - create it @@ -274,8 +351,7 @@ impl SliceEventHandler { error!("Failed to create actor {}: {}", did, e); } else { // Add to cache after successful database insert - let mut cache = self.actor_cache.write().await; - cache.insert(cache_key, true); + self.cache_actor_exists(did, &slice_uri).await; info!("✓ Created actor {} for slice {}", did, slice_uri); } } @@ -284,9 +360,9 @@ impl SliceEventHandler { } } } - + let uri = format!("at://{}/{}/{}", did, commit.collection, commit.rkey); - + let record = Record { uri: uri.clone(), cid: commit.cid.clone(), @@ -296,12 +372,12 @@ impl SliceEventHandler { indexed_at: Utc::now(), slice_uri: Some(slice_uri.clone()), }; - + match self.database.upsert_record(&record).await { Ok(is_insert) => { - let message = if is_insert { + let message = if is_insert { format!("Record inserted in {}", commit.collection) - } else { + } else { format!("Record updated in {}", commit.collection) }; let operation = if is_insert { "insert" } else { "update" }; @@ -311,7 +387,7 @@ impl SliceEventHandler { "slice_uri": slice_uri, "did": did, "record_type": "primary" - })), Some(slice_uri)); + })), Some(&slice_uri)); } Err(e) => { let message = "Failed to insert/update record"; @@ -322,21 +398,21 @@ impl SliceEventHandler { "did": did, "error": e.to_string(), "record_type": "primary" - })), Some(slice_uri)); + })), Some(&slice_uri)); return Err(anyhow::anyhow!("Database error: {}", e)); } } - - info!("✓ Successfully indexed {} record from primary collection: {}", + + info!("✓ Successfully indexed {} record from primary collection: {}", commit.operation, uri); break; } else { // External collection - we already checked actor status, so just index - info!("✓ External collection {} - DID {} is actor in slice {} - indexing", + info!("✓ External collection {} - DID {} is actor in slice {} - indexing", commit.collection, did, slice_uri); - + let uri = format!("at://{}/{}/{}", did, commit.collection, commit.rkey); - + let record = Record { uri: uri.clone(), cid: commit.cid.clone(), @@ -346,12 +422,12 @@ impl SliceEventHandler { indexed_at: Utc::now(), slice_uri: Some(slice_uri.clone()), }; - + match self.database.upsert_record(&record).await { Ok(is_insert) => { - let message = if is_insert { + let message = if is_insert { format!("Record inserted in {}", commit.collection) - } else { + } else { format!("Record updated in {}", commit.collection) }; let operation = if is_insert { "insert" } else { "update" }; @@ -361,7 +437,7 @@ impl SliceEventHandler { "slice_uri": slice_uri, "did": did, "record_type": "external" - })), Some(slice_uri)); + })), Some(&slice_uri)); } Err(e) => { let message = "Failed to insert/update record"; @@ -372,18 +448,18 @@ impl SliceEventHandler { "did": did, "error": e.to_string(), "record_type": "external" - })), Some(slice_uri)); + })), Some(&slice_uri)); return Err(anyhow::anyhow!("Database error: {}", e)); } } - - info!("✓ Successfully indexed {} record from external collection: {}", + + info!("✓ Successfully indexed {} record from external collection: {}", commit.operation, uri); break; } } } - + Ok(()) } @@ -394,34 +470,49 @@ impl SliceEventHandler { ) -> Result<()> { let uri = format!("at://{}/{}/{}", did, commit.collection, commit.rkey); - // Get slices that track this collection - let slice_collections = self.slice_collections.read().await; - let slice_domains = self.slice_domains.read().await; - let actor_cache = self.actor_cache.read().await; + // Get all slices from cached list + let slices = self.slices_list.read().await.clone(); let mut relevant_slices: Vec = Vec::new(); - for (slice_uri, collections) in slice_collections.iter() { + for slice_uri in slices { + // Get collections for this slice (with caching) + let collections = match self.get_slice_collections(&slice_uri).await { + Ok(Some(collections)) => collections, + Ok(None) => continue, // No collections for this slice + Err(e) => { + error!("Failed to get collections for slice {}: {}", slice_uri, e); + continue; + } + }; + if !collections.contains(&commit.collection) { continue; } - // Get the domain for this slice - let domain = match slice_domains.get(slice_uri) { - Some(d) => d, - None => continue, + // Get the domain for this slice (with caching) + let domain = match self.get_slice_domain(&slice_uri).await { + Ok(Some(domain)) => domain, + Ok(None) => continue, // No domain, skip + Err(e) => { + error!("Failed to get domain for slice {}: {}", slice_uri, e); + continue; + } }; // Check if this is a primary collection (starts with slice domain) - let is_primary_collection = commit.collection.starts_with(domain); + let is_primary_collection = commit.collection.starts_with(&domain); if is_primary_collection { // Primary collection - always process deletes relevant_slices.push(slice_uri.clone()); } else { // External collection - only process if DID is an actor in this slice - let cache_key = (did.to_string(), slice_uri.clone()); - if actor_cache.get(&cache_key).copied().unwrap_or(false) { + let is_actor = match self.is_actor_cached(did, &slice_uri).await { + Ok(Some(cached_result)) => cached_result, + _ => false, + }; + if is_actor { relevant_slices.push(slice_uri.clone()); } } @@ -466,9 +557,7 @@ impl SliceEventHandler { if deleted > 0 { info!("✓ Cleaned up actor {} from slice {} (no records remaining)", did, slice_uri); // Remove from cache - let cache_key = (did.to_string(), slice_uri.clone()); - let mut cache = self.actor_cache.write().await; - cache.remove(&cache_key); + self.remove_actor_from_cache(did, slice_uri).await; } } Err(e) => { @@ -512,18 +601,20 @@ impl SliceEventHandler { } impl JetstreamConsumer { - /// Create a new Jetstream consumer with optional cursor support + /// Create a new Jetstream consumer with optional cursor support and Redis cache /// /// # Arguments /// * `database` - Database connection for slice configurations and record storage /// * `jetstream_hostname` - Optional custom jetstream hostname /// * `cursor_handler` - Optional cursor handler for resumable event processing /// * `initial_cursor` - Optional starting cursor position (time_us) to resume from + /// * `redis_url` - Optional Redis URL for caching (falls back to in-memory if not provided) pub async fn new( database: Database, jetstream_hostname: Option, cursor_handler: Option>, initial_cursor: Option, + redis_url: Option, ) -> Result { let config = ConsumerTaskConfig { user_agent: "slice-server/1.0".to_string(), @@ -541,148 +632,116 @@ impl JetstreamConsumer { let consumer = Consumer::new(config); let http_client = Client::new(); + // Determine cache backend based on Redis URL + let cache_backend = if let Some(redis_url) = redis_url { + CacheBackend::Redis { url: redis_url, ttl_seconds: None } + } else { + CacheBackend::InMemory { ttl_seconds: None } + }; + + // Create cache instances + let actor_cache = Arc::new(Mutex::new( + CacheFactory::create_slice_cache(cache_backend.clone()).await + .map_err(|e| SliceError::JetstreamError { + message: format!("Failed to create actor cache: {}", e) + })? + )); + + let lexicon_cache = Arc::new(Mutex::new( + CacheFactory::create_slice_cache(cache_backend.clone()).await + .map_err(|e| SliceError::JetstreamError { + message: format!("Failed to create lexicon cache: {}", e) + })? + )); + + let domain_cache = Arc::new(Mutex::new( + CacheFactory::create_slice_cache(cache_backend.clone()).await + .map_err(|e| SliceError::JetstreamError { + message: format!("Failed to create domain cache: {}", e) + })? + )); + + let collections_cache = Arc::new(Mutex::new( + CacheFactory::create_slice_cache(cache_backend).await + .map_err(|e| SliceError::JetstreamError { + message: format!("Failed to create collections cache: {}", e) + })? + )); + Ok(Self { consumer, database, http_client, - slice_collections: Arc::new(RwLock::new(HashMap::new())), - slice_domains: Arc::new(RwLock::new(HashMap::new())), - actor_cache: Arc::new(RwLock::new(HashMap::new())), - slice_lexicons: Arc::new(RwLock::new(HashMap::new())), + actor_cache, + lexicon_cache, + domain_cache, + collections_cache, event_count: Arc::new(std::sync::atomic::AtomicU64::new(0)), cursor_handler, + slices_list: Arc::new(RwLock::new(Vec::new())), }) } - /// Load slice configurations to know which collections to index + /// Load slice configurations pub async fn load_slice_configurations(&self) -> Result<(), SliceError> { - info!("Loading slice configurations for Jetstream indexing"); - - // Get all slices that have lexicon definitions - let slices = self.database.get_all_slices().await?; - info!("Found {} total slices in database", slices.len()); - - let mut collections_map = HashMap::new(); - let mut domains_map = HashMap::new(); - let mut total_collections = 0; - - for slice_uri in &slices { - info!("Checking slice: {}", slice_uri); - - // Get the domain for this slice - if let Ok(Some(domain)) = self.database.get_slice_domain(slice_uri).await { - info!("Slice {} has domain: {}", slice_uri, domain); - domains_map.insert(slice_uri.clone(), domain.clone()); - - // Get collections defined in this slice's lexicons - let collections = self.database.get_slice_collections_list(slice_uri).await?; - - if !collections.is_empty() { - // Categorize collections as primary or external - let mut primary = Vec::new(); - let mut external = Vec::new(); - - for collection in &collections { - if collection.starts_with(&domain) { - primary.push(collection.clone()); - } else { - external.push(collection.clone()); - } - } - - info!("Slice {} has {} primary collections: {:?}", slice_uri, primary.len(), primary); - info!("Slice {} has {} external collections: {:?}", slice_uri, external.len(), external); - - total_collections += collections.len(); - collections_map.insert(slice_uri.clone(), collections.into_iter().collect()); - } else { - info!("Slice {} has no collections defined (no lexicons or empty lexicons)", slice_uri); - } - } else { - info!("Slice {} has no domain defined - skipping", slice_uri); - } - } - - let mut slice_collections = self.slice_collections.write().await; - *slice_collections = collections_map; - - let mut slice_domains = self.slice_domains.write().await; - *slice_domains = domains_map; - - // Load lexicons for each slice - let mut lexicons_map = HashMap::new(); - for slice_uri in slice_collections.keys() { - info!("Loading lexicons for slice: {}", slice_uri); - - // Get all lexicons for this slice - match self.database.get_lexicons_by_slice(slice_uri).await { - Ok(lexicons) if !lexicons.is_empty() => { - lexicons_map.insert(slice_uri.clone(), lexicons); - info!("✓ Loaded lexicons for slice {}", slice_uri); - } - Ok(_) => { - info!("No lexicons found for slice {}", slice_uri); - } - Err(e) => { - error!("Failed to load lexicons for slice {}: {}", slice_uri, e); - } - } - } + info!("Jetstream consumer now uses on-demand loading with caching"); - let mut slice_lexicons = self.slice_lexicons.write().await; - *slice_lexicons = lexicons_map; + // Get all slices and update cached list + let slices = self.database.get_all_slices().await?; + *self.slices_list.write().await = slices.clone(); + info!("Found {} total slices in database - data will be loaded on-demand", slices.len()); - info!("Jetstream consumer will monitor {} total collections across {} slices with {} lexicon sets loaded", - total_collections, slice_collections.len(), slice_lexicons.len()); - Ok(()) } /// Preload actor cache to avoid database hits during event processing async fn preload_actor_cache(&self) -> Result<(), SliceError> { info!("Preloading actor cache..."); - + let actors = self.database.get_all_actors().await?; info!("Found {} actors to cache", actors.len()); - - let mut cache = self.actor_cache.write().await; - cache.clear(); // Clear existing cache - for (did, slice_uri) in actors { - cache.insert((did, slice_uri), true); + + match self.actor_cache.lock().await.preload_actors(actors).await { + Ok(_) => { + info!("✓ Actor cache preloaded successfully"); + Ok(()) + } + Err(e) => { + warn!(error = ?e, "Failed to preload actors to cache"); + Ok(()) // Don't fail startup if preload fails + } } - - info!("Actor cache preloaded with {} entries", cache.len()); - Ok(()) } - + /// Start consuming events from Jetstream pub async fn start_consuming(&self, cancellation_token: CancellationToken) -> Result<(), SliceError> { info!("Starting Jetstream consumer"); - + // Load initial slice configurations self.load_slice_configurations().await?; - + // Preload actor cache self.preload_actor_cache().await?; - + // Create and register the event handler let handler = Arc::new(SliceEventHandler { database: self.database.clone(), http_client: self.http_client.clone(), - slice_collections: self.slice_collections.clone(), - slice_domains: self.slice_domains.clone(), event_count: self.event_count.clone(), actor_cache: self.actor_cache.clone(), - slice_lexicons: self.slice_lexicons.clone(), + lexicon_cache: self.lexicon_cache.clone(), + domain_cache: self.domain_cache.clone(), + collections_cache: self.collections_cache.clone(), cursor_handler: self.cursor_handler.clone(), + slices_list: self.slices_list.clone(), }); - + self.consumer.register_handler(handler).await .map_err(|e| SliceError::JetstreamError { message: format!("Failed to register event handler: {}", e), })?; - + // Start periodic status reporting let event_count_for_status = self.event_count.clone(); tokio::spawn(async move { @@ -693,7 +752,7 @@ impl JetstreamConsumer { info!("Jetstream consumer status: {} total events processed", count); } }); - + // Start the consumer info!("Starting Jetstream background consumer..."); let result = self.consumer.run_background(cancellation_token).await @@ -718,18 +777,19 @@ impl JetstreamConsumer { pub fn start_configuration_reloader(consumer: Arc) { tokio::spawn(async move { let mut interval = tokio::time::interval(tokio::time::Duration::from_secs(300)); // Reload every 5 minutes - + interval.tick().await; // Skip first immediate tick + loop { interval.tick().await; - + if let Err(e) = consumer.load_slice_configurations().await { error!("Failed to reload slice configurations: {}", e); } - + if let Err(e) = consumer.preload_actor_cache().await { error!("Failed to reload actor cache: {}", e); } } }); } -} \ No newline at end of file +} diff --git a/api/src/jobs.rs b/api/src/jobs.rs index 9aec028..a670992 100644 --- a/api/src/jobs.rs +++ b/api/src/jobs.rs @@ -5,8 +5,11 @@ use uuid::Uuid; use crate::sync::SyncService; use crate::models::BulkSyncParams; use crate::logging::LogLevel; +use crate::cache; use serde_json::json; use tracing::{info, error}; +use std::sync::Arc; +use tokio::sync::Mutex; /// Payload for sync jobs #[derive(Debug, Clone, Serialize, Deserialize)] @@ -65,16 +68,25 @@ async fn sync_job(mut current_job: CurrentJob) -> Result<(), Box, + pub auth_cache: Arc>, } #[tokio::main] @@ -126,6 +130,7 @@ async fn main() -> Result<(), AppError> { let jetstream_connected_clone = jetstream_connected.clone(); tokio::spawn(async move { let jetstream_hostname = env::var("JETSTREAM_HOSTNAME").ok(); + let redis_url = env::var("REDIS_URL").ok(); let cursor_write_interval = env::var("JETSTREAM_CURSOR_WRITE_INTERVAL_SECS") .unwrap_or_else(|_| "5".to_string()) .parse::() @@ -179,12 +184,13 @@ async fn main() -> Result<(), AppError> { cursor_write_interval, )); - // Create consumer with cursor support + // Create consumer with cursor support and Redis cache let consumer_result = JetstreamConsumer::new( database_for_jetstream.clone(), jetstream_hostname.clone(), Some(cursor_handler.clone()), initial_cursor, + redis_url.clone(), ).await; let consumer_arc = match consumer_result { @@ -231,11 +237,23 @@ async fn main() -> Result<(), AppError> { } }); + // Create auth cache for token/session caching (5 minute TTL) + let redis_url = env::var("REDIS_URL").ok(); + let auth_cache_backend = if let Some(redis_url) = redis_url { + cache::CacheBackend::Redis { url: redis_url, ttl_seconds: Some(300) } + } else { + cache::CacheBackend::InMemory { ttl_seconds: Some(300) } + }; + let auth_cache = Arc::new(Mutex::new( + cache::CacheFactory::create_slice_cache(auth_cache_backend).await? + )); + let state = AppState { database: database.clone(), database_pool: pool, config, jetstream_connected, + auth_cache, }; // Build application with routes diff --git a/api/src/redis_cache.rs b/api/src/redis_cache.rs new file mode 100644 index 0000000..d684793 --- /dev/null +++ b/api/src/redis_cache.rs @@ -0,0 +1,259 @@ +use redis::{Client, AsyncCommands}; +use redis::aio::ConnectionManager; +use anyhow::Result; +use tracing::{debug, error, warn}; +use serde::{Serialize, Deserialize}; +use async_trait::async_trait; +use crate::cache::Cache; + +/// Generic Redis cache for scalable caching across multiple instances +pub struct RedisCache { + conn: ConnectionManager, + default_ttl_seconds: u64, +} + +impl RedisCache { + /// Create a new Redis cache + /// + /// # Arguments + /// * `redis_url` - Redis connection URL (e.g., "redis://localhost:6379") + /// * `default_ttl_seconds` - Default time-to-live for cache entries (default: 3600 = 1 hour) + pub async fn new(redis_url: &str, default_ttl_seconds: Option) -> Result { + let client = Client::open(redis_url)?; + let conn = ConnectionManager::new(client).await?; + + Ok(Self { + conn, + default_ttl_seconds: default_ttl_seconds.unwrap_or(3600), + }) + } + + /// Get a value from cache + /// + /// Returns: + /// - Some(T) if key exists and can be deserialized + /// - None if key doesn't exist or deserialization fails + pub async fn get_value(&mut self, key: &str) -> Result> + where + T: for<'de> Deserialize<'de>, + { + match self.conn.get::<_, Option>(key).await { + Ok(Some(value)) => { + match serde_json::from_str::(&value) { + Ok(parsed) => { + // Cache hit - no logging needed + Ok(Some(parsed)) + } + Err(e) => { + error!( + error = ?e, + cache_key = %key, + "Failed to deserialize cached value" + ); + // Remove corrupted entry + let _ = self.conn.del::<_, ()>(key).await; + Ok(None) + } + } + } + Ok(None) => { + // Cache miss - no logging needed + Ok(None) + } + Err(e) => { + error!( + error = ?e, + cache_key = %key, + "Redis error during get" + ); + // Return cache miss on Redis error + Ok(None) + } + } + } + + /// Set a value in cache with optional TTL + pub async fn set_value(&mut self, key: &str, value: &T, ttl_seconds: Option) -> Result<()> + where + T: Serialize, + { + let ttl = ttl_seconds.unwrap_or(self.default_ttl_seconds); + + match serde_json::to_string(value) { + Ok(serialized) => { + match self.conn.set_ex::<_, _, ()>(key, serialized, ttl).await { + Ok(_) => { + debug!( + cache_key = %key, + ttl_seconds = ttl, + "Cached value in Redis" + ); + Ok(()) + } + Err(e) => { + error!( + error = ?e, + cache_key = %key, + "Failed to cache value in Redis" + ); + // Don't fail the operation if Redis is down + Ok(()) + } + } + } + Err(e) => { + error!( + error = ?e, + cache_key = %key, + "Failed to serialize value for caching" + ); + Ok(()) + } + } + } + + /// Check if a key exists in cache + pub async fn key_exists(&mut self, key: &str) -> Result { + match self.conn.exists(key).await { + Ok(exists) => { + debug!(cache_key = %key, exists = exists, "Redis exists check"); + Ok(exists) + } + Err(e) => { + error!( + error = ?e, + cache_key = %key, + "Redis error during exists check" + ); + Ok(false) + } + } + } + + /// Delete a key from cache + pub async fn delete_key(&mut self, key: &str) -> Result<()> { + match self.conn.del::<_, ()>(key).await { + Ok(_) => { + debug!(cache_key = %key, "Deleted key from Redis cache"); + Ok(()) + } + Err(e) => { + error!( + error = ?e, + cache_key = %key, + "Failed to delete key from Redis cache" + ); + Ok(()) + } + } + } + + /// Set multiple key-value pairs using pipeline for efficiency + pub async fn set_multiple_values(&mut self, items: Vec<(&str, &T, Option)>) -> Result<()> + where + T: Serialize, + { + if items.is_empty() { + return Ok(()); + } + + let mut pipe = redis::pipe(); + let mut serialization_errors = 0; + + for (key, value, ttl) in &items { + match serde_json::to_string(value) { + Ok(serialized) => { + let ttl_to_use = ttl.unwrap_or(self.default_ttl_seconds); + pipe.set_ex(key, serialized, ttl_to_use); + } + Err(e) => { + error!( + error = ?e, + cache_key = %key, + "Failed to serialize value for bulk caching" + ); + serialization_errors += 1; + } + } + } + + match pipe.query_async::<()>(&mut self.conn).await { + Ok(_) => { + debug!( + items_count = items.len() - serialization_errors, + serialization_errors = serialization_errors, + "Successfully bulk cached items in Redis" + ); + Ok(()) + } + Err(e) => { + error!( + error = ?e, + items_count = items.len(), + "Failed to bulk cache items in Redis" + ); + Ok(()) + } + } + } + + /// Test Redis connection + pub async fn ping(&mut self) -> Result { + match self.conn.ping::().await { + Ok(response) => Ok(response == "PONG"), + Err(e) => { + error!(error = ?e, "Redis ping failed"); + Ok(false) + } + } + } + + /// Get cache statistics (for monitoring) + pub async fn get_info(&mut self) -> Result { + match redis::cmd("INFO").arg("memory").query_async::(&mut self.conn).await { + Ok(info) => Ok(info), + Err(e) => { + warn!(error = ?e, "Failed to get Redis info"); + Ok("Redis info unavailable".to_string()) + } + } + } +} + +#[async_trait] +impl Cache for RedisCache { + async fn get(&mut self, key: &str) -> Result> + where + T: for<'de> Deserialize<'de> + Send, + { + self.get_value(key).await + } + + async fn set(&mut self, key: &str, value: &T, ttl_seconds: Option) -> Result<()> + where + T: Serialize + Send + Sync, + { + self.set_value(key, value, ttl_seconds).await + } + + + async fn delete(&mut self, key: &str) -> Result<()> { + self.delete_key(key).await + } + + async fn set_multiple(&mut self, items: Vec<(&str, &T, Option)>) -> Result<()> + where + T: Serialize + Send + Sync, + { + self.set_multiple_values(items).await + } + + async fn ping(&mut self) -> Result { + RedisCache::ping(self).await + } + + async fn get_info(&mut self) -> Result { + RedisCache::get_info(self).await + } +} + diff --git a/api/src/sync.rs b/api/src/sync.rs index c1b7f3a..308f1d1 100644 --- a/api/src/sync.rs +++ b/api/src/sync.rs @@ -8,13 +8,16 @@ use atproto_identity::{ resolve::{resolve_subject, HickoryDnsResolver}, }; -use crate::actor_resolver::resolve_actor_data; +use crate::actor_resolver::{resolve_actor_data_cached, resolve_actor_data_with_retry}; +use crate::cache::SliceCache; use crate::database::Database; use crate::errors::SyncError; use crate::models::{Actor, Record}; use crate::logging::LogLevel; use crate::logging::Logger; use serde_json::json; +use std::sync::Arc; +use tokio::sync::Mutex; use uuid::Uuid; @@ -66,32 +69,34 @@ pub struct SyncService { client: Client, database: Database, relay_endpoint: String, - atp_cache: std::sync::Arc>>, + cache: Option>>, logger: Option, job_id: Option, user_did: Option, } impl SyncService { - pub fn new(database: Database, relay_endpoint: String) -> Self { + + pub fn with_cache(database: Database, relay_endpoint: String, cache: Arc>) -> Self { Self { client: Client::new(), database, relay_endpoint, - atp_cache: std::sync::Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())), + cache: Some(cache), logger: None, job_id: None, user_did: None, } } - /// Create a new SyncService with logging enabled for a specific job - pub fn with_logging(database: Database, relay_endpoint: String, logger: Logger, job_id: Uuid, user_did: String) -> Self { + + /// Create a new SyncService with logging and cache enabled for a specific job + pub fn with_logging_and_cache(database: Database, relay_endpoint: String, logger: Logger, job_id: Uuid, user_did: String, cache: Arc>) -> Self { Self { client: Client::new(), database, relay_endpoint, - atp_cache: std::sync::Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())), + cache: Some(cache), logger: Some(logger), job_id: Some(job_id), user_did: Some(user_did), @@ -126,8 +131,7 @@ impl SyncService { let all_collections = [&primary_collections[..], &external_collections[..]].concat(); - // Clear cache at start of each backfill operation - self.clear_atp_cache(); + // DID resolution cache is now handled by SliceCache let all_repos = if let Some(provided_repos) = repos { info!("📋 Using {} provided repositories", provided_repos.len()); @@ -474,7 +478,26 @@ impl SyncService { async fn fetch_records_for_repo_collection_with_atp_map(&self, repo: &str, collection: &str, atp_map: &std::collections::HashMap, slice_uri: &str) -> Result, SyncError> { let atp_data = atp_map.get(repo).ok_or_else(|| SyncError::Generic(format!("No ATP data found for repo: {}", repo)))?; - self.fetch_records_for_repo_collection(repo, collection, &atp_data.pds, slice_uri).await + + match self.fetch_records_for_repo_collection(repo, collection, &atp_data.pds, slice_uri).await { + Ok(records) => Ok(records), + Err(SyncError::ListRecords { status }) if (400..600).contains(&status) => { + // 4xx/5xx error from PDS - try invalidating cache and retrying once + debug!("PDS error {} for repo {}, attempting cache invalidation and retry", status, repo); + + match resolve_actor_data_with_retry(&self.client, repo, self.cache.clone(), true).await { + Ok(fresh_actor_data) => { + debug!("Successfully re-resolved actor data for {}, retrying with PDS: {}", repo, fresh_actor_data.pds); + self.fetch_records_for_repo_collection(repo, collection, &fresh_actor_data.pds, slice_uri).await + } + Err(e) => { + debug!("Failed to re-resolve actor data for {}: {:?}", repo, e); + Err(SyncError::ListRecords { status }) // Return original error + } + } + } + Err(e) => Err(e), // Other errors (network, etc.) - don't retry + } } async fn fetch_records_for_repo_collection(&self, repo: &str, collection: &str, pds_url: &str, slice_uri: &str) -> Result, SyncError> { @@ -605,21 +628,13 @@ impl SyncService { async fn resolve_atp_data(&self, did: &str) -> Result { debug!("Resolving ATP data for DID: {}", did); - { - let cache = self.atp_cache.lock().unwrap(); - if let Some(cached_data) = cache.get(did) { - debug!("Using cached ATP data for DID: {}", did); - return Ok(cached_data.clone()); - } - } - let dns_resolver = HickoryDnsResolver::create_resolver(&[]); match resolve_subject(&self.client, &dns_resolver, did).await { Ok(resolved_did) => { debug!("Successfully resolved subject: {}", resolved_did); - let actor_data = resolve_actor_data(&self.client, &resolved_did).await + let actor_data = resolve_actor_data_cached(&self.client, &resolved_did, self.cache.clone()).await .map_err(|e| SyncError::Generic(e.to_string()))?; let atp_data = AtpData { @@ -628,12 +643,6 @@ impl SyncService { handle: actor_data.handle, }; - // Cache the result - { - let mut cache = self.atp_cache.lock().unwrap(); - cache.insert(did.to_string(), atp_data.clone()); - } - Ok(atp_data) } Err(e) => { @@ -664,10 +673,6 @@ impl SyncService { Ok(()) } - pub fn clear_atp_cache(&self) { - let mut cache = self.atp_cache.lock().unwrap(); - cache.clear(); - } /// Get external collections for a slice (collections that don't start with the slice's domain) async fn get_external_collections_for_slice(&self, slice_uri: &str) -> Result, SyncError> { diff --git a/api/src/xrpc/com/atproto/repo/upload_blob.rs b/api/src/xrpc/com/atproto/repo/upload_blob.rs index 11c3d16..394ccfc 100644 --- a/api/src/xrpc/com/atproto/repo/upload_blob.rs +++ b/api/src/xrpc/com/atproto/repo/upload_blob.rs @@ -9,10 +9,10 @@ pub async fn handler( let headers = request.headers().clone(); let token = auth::extract_bearer_token(&headers)?; - let _user_info = auth::verify_oauth_token(&token, &state.config.auth_base_url).await?; + let _user_info = auth::verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let (dpop_auth, pds_url) = - auth::get_atproto_auth_for_user(&token, &state.config.auth_base_url).await?; + auth::get_atproto_auth_for_user_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let mime_type = headers .get("content-type") diff --git a/api/src/xrpc/network/slices/slice/create_oauth_client.rs b/api/src/xrpc/network/slices/slice/create_oauth_client.rs index 38e7d57..ffc3b44 100644 --- a/api/src/xrpc/network/slices/slice/create_oauth_client.rs +++ b/api/src/xrpc/network/slices/slice/create_oauth_client.rs @@ -102,7 +102,7 @@ pub async fn handler( } let token = auth::extract_bearer_token(&headers)?; - let user_info = auth::verify_oauth_token(&token, &state.config.auth_base_url).await?; + let user_info = auth::verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let user_did = user_info.sub; diff --git a/api/src/xrpc/network/slices/slice/delete_oauth_client.rs b/api/src/xrpc/network/slices/slice/delete_oauth_client.rs index 3fe6456..8e01a1a 100644 --- a/api/src/xrpc/network/slices/slice/delete_oauth_client.rs +++ b/api/src/xrpc/network/slices/slice/delete_oauth_client.rs @@ -21,7 +21,7 @@ pub async fn handler( Json(params): Json, ) -> Result, AppError> { let token = auth::extract_bearer_token(&headers)?; - auth::verify_oauth_token(&token, &state.config.auth_base_url).await?; + auth::verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let oauth_client = state .database diff --git a/api/src/xrpc/network/slices/slice/get_oauth_clients.rs b/api/src/xrpc/network/slices/slice/get_oauth_clients.rs index 2834a98..7b93166 100644 --- a/api/src/xrpc/network/slices/slice/get_oauth_clients.rs +++ b/api/src/xrpc/network/slices/slice/get_oauth_clients.rs @@ -55,7 +55,7 @@ pub async fn handler( Query(params): Query, ) -> Result, AppError> { let token = auth::extract_bearer_token(&headers)?; - auth::verify_oauth_token(&token, &state.config.auth_base_url).await?; + auth::verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let clients = state .database diff --git a/api/src/xrpc/network/slices/slice/start_sync.rs b/api/src/xrpc/network/slices/slice/start_sync.rs index 146b565..0b29c7d 100644 --- a/api/src/xrpc/network/slices/slice/start_sync.rs +++ b/api/src/xrpc/network/slices/slice/start_sync.rs @@ -24,7 +24,7 @@ pub async fn handler( Json(params): Json, ) -> Result, AppError> { let token = auth::extract_bearer_token(&headers)?; - let user_info = auth::verify_oauth_token(&token, &state.config.auth_base_url).await?; + let user_info = auth::verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let user_did = user_info.sub; let slice_uri = params.slice; diff --git a/api/src/xrpc/network/slices/slice/sync_user_collections.rs b/api/src/xrpc/network/slices/slice/sync_user_collections.rs index ab11df4..f54dc1e 100644 --- a/api/src/xrpc/network/slices/slice/sync_user_collections.rs +++ b/api/src/xrpc/network/slices/slice/sync_user_collections.rs @@ -20,7 +20,7 @@ pub async fn handler( Json(params): Json, ) -> Result, AppError> { let token = auth::extract_bearer_token(&headers)?; - let user_info = auth::verify_oauth_token(&token, &state.config.auth_base_url).await?; + let user_info = auth::verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let user_did = user_info.did.unwrap_or(user_info.sub); @@ -38,7 +38,7 @@ pub async fn handler( ); let sync_service = - crate::sync::SyncService::new(state.database.clone(), state.config.relay_endpoint.clone()); + crate::sync::SyncService::with_cache(state.database.clone(), state.config.relay_endpoint.clone(), state.auth_cache.clone()); let result = sync_service .sync_user_collections(&user_did, ¶ms.slice, params.timeout_seconds) diff --git a/api/src/xrpc/network/slices/slice/update_oauth_client.rs b/api/src/xrpc/network/slices/slice/update_oauth_client.rs index 7d59a05..f20c818 100644 --- a/api/src/xrpc/network/slices/slice/update_oauth_client.rs +++ b/api/src/xrpc/network/slices/slice/update_oauth_client.rs @@ -73,7 +73,7 @@ pub async fn handler( Json(params): Json, ) -> Result, AppError> { let token = auth::extract_bearer_token(&headers)?; - auth::verify_oauth_token(&token, &state.config.auth_base_url).await?; + auth::verify_oauth_token_cached(&token, &state.config.auth_base_url, Some(state.auth_cache.clone())).await?; let oauth_client = state .database diff --git a/docker-compose.yml b/docker-compose.yml index c637b21..e37985f 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -28,6 +28,18 @@ services: networks: - default + redis: + image: redis:7-alpine + ports: + - "6379:6379" + volumes: + - redis_data:/data + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 5s + timeout: 3s + retries: 5 + aip: image: ghcr.io/bigmoves/aip/aip-sqlite:main-e445b82 environment: @@ -54,4 +66,5 @@ services: volumes: postgres_data: + redis_data: aip_data: