diff --git a/Cargo.lock b/Cargo.lock index 4d908fd2..c1a617f1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -306,6 +306,15 @@ dependencies = [ "syn", ] +[[package]] +name = "atoi" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f28d99ec8bfea296261ca1af174f24225171fea9664ba9003cbebee704810528" +dependencies = [ + "num-traits", +] + [[package]] name = "atomic" version = "0.6.1" @@ -518,6 +527,9 @@ name = "bitflags" version = "2.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b4388bee8683e3d04af747c73422af53102d2bd24d9eadb6cbc100baef4b43f8" +dependencies = [ + "serde_core", +] [[package]] name = "blake3" @@ -542,6 +554,15 @@ dependencies = [ "generic-array", ] +[[package]] +name = "block-buffer" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d2f6c7dbe95a6ed67ad9f18e57daf93a2f034c524b99fd2b76d18fdfeb6660aa" +dependencies = [ + "hybrid-array", +] + [[package]] name = "block2" version = "0.6.2" @@ -585,6 +606,12 @@ version = "1.25.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c8efb64bd706a16a1bdde310ae86b351e4d21550d98d056f22f8a7f7a2183fec" +[[package]] +name = "byteorder" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fd0f2584146f6f2ef48085050886acf353beff7305ebd1ae69500e27c67f64b" + [[package]] name = "bytes" version = "1.12.0" @@ -670,6 +697,12 @@ dependencies = [ "cc", ] +[[package]] +name = "cmov" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a" + [[package]] name = "combine" version = "4.6.7" @@ -774,6 +807,21 @@ dependencies = [ "libc", ] +[[package]] +name = "crc" +version = "3.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5eb8a2a1cd12ab0d987a5d5e825195d372001a4094a0376319d5a0ad71c1ba0d" +dependencies = [ + "crc-catalog", +] + +[[package]] +name = "crc-catalog" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "217698eaf96b4a3f0bc4f3662aaa55bdf913cd54d7204591faa790070c6d0853" + [[package]] name = "crc32fast" version = "1.5.0" @@ -807,6 +855,15 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "crossbeam-queue" +version = "0.3.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "03e8bd762f7479489c70ed6c768ddca99d7296857de437a68dcb2a94365b3fae" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-utils" version = "0.8.21" @@ -835,6 +892,24 @@ dependencies = [ "typenum", ] +[[package]] +name = "crypto-common" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ce6e4c961d6cd6c9a86db418387425e8bdeaf05b3c8bc1411e6dca4c252f1453" +dependencies = [ + "hybrid-array", +] + +[[package]] +name = "ctutils" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7d5515a3834141de9eafb9717ad39eea8247b5674e6066c404e8c4b365d2a29e" +dependencies = [ + "cmov", +] + [[package]] name = "curve25519-dalek" version = "4.1.3" @@ -844,7 +919,7 @@ dependencies = [ "cfg-if", "cpufeatures 0.2.17", "curve25519-dalek-derive", - "digest", + "digest 0.10.7", "fiat-crypto", "rustc_version", "subtle", @@ -934,12 +1009,23 @@ version = "0.10.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" dependencies = [ - "block-buffer", + "block-buffer 0.10.4", "const-oid", - "crypto-common", + "crypto-common 0.1.6", "subtle", ] +[[package]] +name = "digest" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" +dependencies = [ + "block-buffer 0.12.1", + "crypto-common 0.2.2", + "ctutils", +] + [[package]] name = "dispatch2" version = "0.3.1" @@ -961,6 +1047,12 @@ dependencies = [ "syn", ] +[[package]] +name = "dotenvy" +version = "0.15.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1aaf95b3e5c8f23aa320147307562d361db0ae0d51242340f558153b4eb2439b" + [[package]] name = "dunce" version = "1.0.5" @@ -974,7 +1066,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ee27f32b5c5292967d2d4a9d7f1e0b0aed2c15daded5a60300e4abb9d8020bca" dependencies = [ "der 0.7.10", - "digest", + "digest 0.10.7", "elliptic-curve", "rfc6979", "signature", @@ -1000,7 +1092,7 @@ dependencies = [ "curve25519-dalek", "ed25519", "serde", - "sha2", + "sha2 0.10.9", "subtle", "zeroize", ] @@ -1010,6 +1102,9 @@ name = "either" version = "1.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "91622ff5e7162018101f2fea40d6ebf4a78bbe5a49736a2020649edf9693679e" +dependencies = [ + "serde", +] [[package]] name = "elegant-departure" @@ -1031,11 +1126,11 @@ checksum = "b5e6043086bf7973472e0c7dff2142ea0b680d30e18d9cc40f267efbf222bd47" dependencies = [ "base16ct", "crypto-bigint", - "digest", + "digest 0.10.7", "ff", "generic-array", "group", - "hkdf", + "hkdf 0.12.4", "pem-rfc7468 0.7.0", "pkcs8", "rand_core 0.6.4", @@ -1090,6 +1185,26 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "etcetera" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "de48cc4d1c1d97a20fd819def54b890cadde72ed3ad0c614822a0a433361be96" +dependencies = [ + "cfg-if", + "windows-sys 0.61.2", +] + +[[package]] +name = "event-listener" +version = "5.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a23add41df1562121a9393cb065eab5146a1242410f23a644851e90cfd669d2" +dependencies = [ + "parking", + "pin-project-lite", +] + [[package]] name = "fancy-regex" version = "0.16.2" @@ -1184,6 +1299,17 @@ dependencies = [ "serde", ] +[[package]] +name = "flume" +version = "0.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e139bc46ca777eb5efaf62df0ab8cc5fd400866427e56c68b22e414e53bd3be" +dependencies = [ + "futures-core", + "futures-sink", + "spin", +] + [[package]] name = "fnv" version = "1.0.7" @@ -1284,6 +1410,17 @@ dependencies = [ "futures-util", ] +[[package]] +name = "futures-intrusive" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1d930c203dd0b6ff06e0201a4a2fe9149b43c684fd4420555b26d21b1a02956f" +dependencies = [ + "futures-core", + "lock_api", + "parking_lot", +] + [[package]] name = "futures-io" version = "0.3.32" @@ -1538,6 +1675,15 @@ dependencies = [ "foldhash 0.2.0", ] +[[package]] +name = "hashlink" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "824e001ac4f3012dd16a264bec811403a67ca9deb6c102fc5049b32c4574b35f" +dependencies = [ + "hashbrown 0.16.1", +] + [[package]] name = "headers" version = "0.4.1" @@ -1550,7 +1696,7 @@ dependencies = [ "http 1.4.2", "httpdate", "mime", - "sha1", + "sha1 0.10.6", ] [[package]] @@ -1656,7 +1802,16 @@ version = "0.12.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7b5f8eb2ad728638ea2c7d47a21db23b7b58a72ed6a38256b8a1849f15fbbdf7" dependencies = [ - "hmac", + "hmac 0.12.1", +] + +[[package]] +name = "hkdf" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4aaa26c720c68b866f2c96ef5c1264b3e6f473fe5d4ce61cd44bbe913e553018" +dependencies = [ + "hmac 0.13.0", ] [[package]] @@ -1665,7 +1820,16 @@ version = "0.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6c49c37c09c17a53d937dfbb742eb3a961d65a994e6bcdcf37e7399d0cc8ab5e" dependencies = [ - "digest", + "digest 0.10.7", +] + +[[package]] +name = "hmac" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6303bc9732ae41b04cb554b844a762b4115a61bfaa81e3e83050991eeb56863f" +dependencies = [ + "digest 0.11.3", ] [[package]] @@ -1751,6 +1915,15 @@ dependencies = [ "serde", ] +[[package]] +name = "hybrid-array" +version = "0.4.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "27f864f10dfb56725ce5ce5472bc52252c8f93a4ab86327122cebf62c5f59a17" +dependencies = [ + "typenum", +] + [[package]] name = "hyper" version = "1.10.1" @@ -2184,7 +2357,7 @@ dependencies = [ "base64", "ed25519-dalek", "getrandom 0.2.17", - "hmac", + "hmac 0.12.1", "js-sys", "p256", "p384", @@ -2193,7 +2366,7 @@ dependencies = [ "rsa", "serde", "serde_json", - "sha2", + "sha2 0.10.9", "signature", "simple_asn1", "zeroize", @@ -2226,6 +2399,17 @@ version = "0.2.16" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981" +[[package]] +name = "libsqlite3-sys" +version = "0.37.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b1f111c8c41e7c61a49cd34e44c7619462967221a6443b0ec299e0ac30cfb9b1" +dependencies = [ + "cc", + "pkg-config", + "vcpkg", +] + [[package]] name = "libz-sys" version = "1.1.29" @@ -2314,6 +2498,16 @@ version = "0.8.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "47e1ffaa40ddd1f3ed91f717a33c8c0ee23fff369e3aa8772b9605cc1d22f4c3" +[[package]] +name = "md-5" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69b6441f590336821bb897fb28fc622898ccceb1d6cea3fde5ea86b090c4de98" +dependencies = [ + "cfg-if", + "digest 0.11.3", +] + [[package]] name = "mediatype" version = "0.21.0" @@ -2959,6 +3153,7 @@ dependencies = [ "sentry", "serde", "serde_json", + "sqlx", "tempfile", "thiserror", "tokio", @@ -3112,7 +3307,7 @@ dependencies = [ "ecdsa", "elliptic-curve", "primeorder", - "sha2", + "sha2 0.10.9", ] [[package]] @@ -3124,7 +3319,7 @@ dependencies = [ "ecdsa", "elliptic-curve", "primeorder", - "sha2", + "sha2 0.10.9", ] [[package]] @@ -3137,6 +3332,12 @@ dependencies = [ "seize", ] +[[package]] +name = "parking" +version = "2.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f38d5652c16fde515bb1ecef450ab0f6a219d619a7274976324d5e377f7dceba" + [[package]] name = "parking_lot" version = "0.12.5" @@ -3911,7 +4112,7 @@ version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8dd2a808d456c4a54e300a23e9f5a67e122c3024119acbfd73e3bf664491cb2" dependencies = [ - "hmac", + "hmac 0.12.1", "subtle", ] @@ -3936,7 +4137,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b8573f03f5883dcaebdfcf4725caa1ecb9c15b2ef50c43a07b816e06799bb12d" dependencies = [ "const-oid", - "digest", + "digest 0.10.7", "num-bigint-dig", "num-integer", "num-traits", @@ -4264,7 +4465,7 @@ dependencies = [ "num", "sentry-options-validation", "serde_json", - "sha1", + "sha1 0.10.6", "thiserror", "tracing", ] @@ -4437,7 +4638,18 @@ checksum = "e3bf829a2d51ab4a5ddf1352d8470c140cadc8301b2ae1789db023f01cedd6ba" dependencies = [ "cfg-if", "cpufeatures 0.2.17", - "digest", + "digest 0.10.7", +] + +[[package]] +name = "sha1" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "aacc4cc499359472b4abe1bf11d0b12e688af9a805fa5e3016f9a386dc2d0214" +dependencies = [ + "cfg-if", + "cpufeatures 0.3.0", + "digest 0.11.3", ] [[package]] @@ -4448,7 +4660,18 @@ checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" dependencies = [ "cfg-if", "cpufeatures 0.2.17", - "digest", + "digest 0.10.7", +] + +[[package]] +name = "sha2" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "446ba717509524cb3f22f17ecc096f10f4822d76ab5c0b9822c5f9c284e825f4" +dependencies = [ + "cfg-if", + "cpufeatures 0.3.0", + "digest 0.11.3", ] [[package]] @@ -4482,7 +4705,7 @@ version = "2.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de" dependencies = [ - "digest", + "digest 0.10.7", "rand_core 0.6.4", ] @@ -4543,6 +4766,9 @@ name = "smallvec" version = "1.15.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8ed6a63f02c8539c91a8685a86f4099661ba3da017932f6ebbea6de3f0fa7c90" +dependencies = [ + "serde", +] [[package]] name = "socket2" @@ -4569,6 +4795,9 @@ name = "spin" version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +dependencies = [ + "lock_api", +] [[package]] name = "spki" @@ -4580,6 +4809,183 @@ dependencies = [ "der 0.7.10", ] +[[package]] +name = "sqlx" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "378620ccc25c62c89d8be1c819e76a88d59bdcc3304733330788948e619bfd71" +dependencies = [ + "sqlx-core", + "sqlx-macros", + "sqlx-mysql", + "sqlx-postgres", + "sqlx-sqlite", +] + +[[package]] +name = "sqlx-core" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "05b44e85bf579a8eeb4ceaa77a3a523baf2bf0e9bac7e40f405d537b5d2d5ccb" +dependencies = [ + "base64", + "bytes", + "cfg-if", + "chrono", + "crc", + "crossbeam-queue", + "either", + "event-listener", + "futures-core", + "futures-intrusive", + "futures-io", + "futures-util", + "hashbrown 0.16.1", + "hashlink", + "indexmap", + "log", + "memchr", + "percent-encoding", + "rustls", + "rustls-native-certs", + "serde", + "serde_json", + "sha2 0.10.9", + "smallvec", + "thiserror", + "tokio", + "tokio-stream", + "tracing", + "url", +] + +[[package]] +name = "sqlx-macros" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bd2b84f2bc39a5705ef27ec785a11c934a41bbd4a24941e257927cddc26b60bf" +dependencies = [ + "proc-macro2", + "quote", + "sqlx-core", + "sqlx-macros-core", + "syn", +] + +[[package]] +name = "sqlx-macros-core" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fb8d96de5fdc85a5c4ec813432b523ec637e80ba98f046555f75f7908ddac7c3" +dependencies = [ + "cfg-if", + "dotenvy", + "either", + "heck", + "hex", + "proc-macro2", + "quote", + "serde", + "serde_json", + "sha2 0.10.9", + "sqlx-core", + "sqlx-mysql", + "sqlx-postgres", + "sqlx-sqlite", + "syn", + "thiserror", + "tokio", + "url", +] + +[[package]] +name = "sqlx-mysql" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "90b8020fe17c5f2c245bfa2505d7ef59c5604839527c740266ad2214acebea27" +dependencies = [ + "bitflags", + "byteorder", + "bytes", + "chrono", + "crc", + "digest 0.11.3", + "dotenvy", + "either", + "futures-core", + "futures-util", + "generic-array", + "log", + "percent-encoding", + "serde", + "sha1 0.11.0", + "sha2 0.11.0", + "sqlx-core", + "thiserror", + "tracing", +] + +[[package]] +name = "sqlx-postgres" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "87a2bdd6e83f6b3ea525ca9fee568030508b58355a43d0b2c1674d5f79dcd65e" +dependencies = [ + "atoi", + "base64", + "bitflags", + "byteorder", + "chrono", + "crc", + "dotenvy", + "etcetera", + "futures-channel", + "futures-core", + "futures-util", + "hex", + "hkdf 0.13.0", + "hmac 0.13.0", + "itoa", + "log", + "md-5", + "memchr", + "rand 0.10.2", + "serde", + "serde_json", + "sha2 0.11.0", + "smallvec", + "sqlx-core", + "stringprep", + "thiserror", + "tracing", + "whoami", +] + +[[package]] +name = "sqlx-sqlite" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "488e99c397a62007e4229aec669a179816339afc6d2620ca6fa420dbee2e982c" +dependencies = [ + "atoi", + "chrono", + "flume", + "form_urlencoded", + "futures-channel", + "futures-core", + "futures-executor", + "futures-intrusive", + "futures-util", + "libsqlite3-sys", + "log", + "percent-encoding", + "serde", + "sqlx-core", + "thiserror", + "tracing", + "url", +] + [[package]] name = "stable_deref_trait" version = "1.2.1" @@ -4612,6 +5018,17 @@ dependencies = [ "yansi", ] +[[package]] +name = "stringprep" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b4df3d392d81bd458a8a621b8bffbd2302a12ffe288a9d931670948749463b1" +dependencies = [ + "unicode-bidi", + "unicode-normalization", + "unicode-properties", +] + [[package]] name = "subtle" version = "2.6.1" @@ -5166,6 +5583,12 @@ version = "2.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dbc4bc3a9f746d862c45cb89d705aa10f187bb96c76001afab07a0d35ce60142" +[[package]] +name = "unicode-bidi" +version = "0.3.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c1cb5db39152898a79168971543b1cb5020dff7fe43c8dc468b0885f5e29df5" + [[package]] name = "unicode-general-category" version = "1.1.0" @@ -5178,6 +5601,21 @@ version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +[[package]] +name = "unicode-normalization" +version = "0.1.25" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5fd4f6878c9cb28d874b009da9e8d183b5abc80117c40bbd187a1fde336be6e8" +dependencies = [ + "tinyvec", +] + +[[package]] +name = "unicode-properties" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7df058c713841ad818f1dc5d3fd88063241cc61f49f5fbea4b951e8cf5a8d71d" + [[package]] name = "unicode-segmentation" version = "1.13.3" @@ -5451,6 +5889,12 @@ dependencies = [ "rustls-pki-types", ] +[[package]] +name = "whoami" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "626c4bac6755d76ffc12cb01b2eac751db1996b9e0041de9aa02c8c211ddc82c" + [[package]] name = "widestring" version = "1.2.1" diff --git a/Cargo.toml b/Cargo.toml index c2a7b65e..df905b79 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -43,7 +43,10 @@ async-trait = "0.1.89" axum = "0.8.9" axum-extra = { version = "0.12.6", default-features = false } base64 = "0.22.1" -bigtable_rs = { git = "https://github.com/getsentry/bigtable_rs.git", rev = "56c310a5c8518d52ae7c7db53825beb078e326b5", default-features = false, features = ["gcp-auth-ring", "tonic-tls-ring"] } +bigtable_rs = { git = "https://github.com/getsentry/bigtable_rs.git", rev = "56c310a5c8518d52ae7c7db53825beb078e326b5", default-features = false, features = [ + "gcp-auth-ring", + "tonic-tls-ring", +] } blake3 = "1.8.5" bytes = "1.12.0" bytesize = "2.4.0" @@ -99,12 +102,16 @@ serde = { version = "1.0.228", features = ["derive"] } serde-vars = "0.3.1" serde_json = "1.0.150" serde_yaml = "0.9.34-deprecated" +sqlx = "0.9.0" sketches-ddsketch = "0.3.1" tempfile = "3.27.0" thiserror = "2.0.18" thread_local = "1.1.9" jemalloc_pprof = "0.9.0" -tikv-jemallocator = { version = "0.7.0", features = ["background_threads", "override_allocator_on_supported_platforms"] } +tikv-jemallocator = { version = "0.7.0", features = [ + "background_threads", + "override_allocator_on_supported_platforms", +] } tikv-jemalloc-ctl = { version = "0.7.0", features = ["stats"] } tokio = "1.52.3" tokio-util = "0.7.18" diff --git a/migrations/sqlite/0001_garbage_collector.sql b/migrations/sqlite/0001_garbage_collector.sql new file mode 100644 index 00000000..33e81758 --- /dev/null +++ b/migrations/sqlite/0001_garbage_collector.sql @@ -0,0 +1,5 @@ +CREATE TABLE IF NOT EXISTS garbage_collector ( + object_id TEXT NOT NULL PRIMARY KEY, + expires_at INTEGER, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP +); diff --git a/objectstore-service/Cargo.toml b/objectstore-service/Cargo.toml index d70938d9..0e8f839c 100644 --- a/objectstore-service/Cargo.toml +++ b/objectstore-service/Cargo.toml @@ -21,18 +21,35 @@ futures-util = { workspace = true } gcp_auth = { workspace = true } humantime = { workspace = true } humantime-serde = { workspace = true } -objectstore-inventory-tracker = { workspace = true, features = ["kafka"], optional = true } +objectstore-inventory-tracker = { workspace = true, features = [ + "kafka", +], optional = true } objectstore-log = { workspace = true } objectstore-metrics = { workspace = true } objectstore-types = { workspace = true } papaya = { workspace = true } quick-xml = { workspace = true, features = ["serialize"] } regex = { workspace = true } -reqwest = { workspace = true, features = ["charset", "http2", "hickory-dns", "json", "multipart", "native-tls-no-alpn", "stream", "system-proxy"] } +reqwest = { workspace = true, features = [ + "charset", + "http2", + "hickory-dns", + "json", + "multipart", + "native-tls-no-alpn", + "stream", + "system-proxy", +] } ring = { workspace = true } sentry = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } +sqlx = { workspace = true, features = [ + "runtime-tokio", + "tls-rustls-ring-native-roots", + "chrono", + "sqlite", +] } tempfile = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true } diff --git a/objectstore-service/src/change_stream/garbage_collector_sqlite.rs b/objectstore-service/src/change_stream/garbage_collector_sqlite.rs new file mode 100644 index 00000000..0c76af0d --- /dev/null +++ b/objectstore-service/src/change_stream/garbage_collector_sqlite.rs @@ -0,0 +1,434 @@ +use crate::change_stream::ChangeStream; +use async_trait::async_trait; +use objectstore_types::time::Timestamp; +use serde::{Deserialize, Serialize}; +use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode}; +use sqlx::{Connection, SqlitePool}; +use std::fmt; +use std::str::FromStr; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; + +#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)] +#[expect(dead_code)] +pub struct SqliteGarbageCollectorConfig { + pub path: String, +} + +impl Default for SqliteGarbageCollectorConfig { + fn default() -> Self { + Self { + path: "./objectstore.db".to_owned(), + } + } +} + +#[expect(dead_code)] +pub struct SqliteGarbageCollectorStream { + pool: SqlitePool, + active_tasks: Arc, +} + +#[expect(dead_code)] +impl SqliteGarbageCollectorStream { + pub async fn new(config: &SqliteGarbageCollectorConfig) -> Result { + let opts = SqliteConnectOptions::from_str(&config.path)? + .journal_mode(SqliteJournalMode::Wal) + .create_if_missing(true); + let pool = SqlitePool::connect_with(opts).await?; + + sqlx::migrate!("./../migrations/sqlite").run(&pool).await?; + + Ok(Self { + pool, + active_tasks: Arc::new(AtomicUsize::new(0)), + }) + } +} + +impl fmt::Debug for SqliteGarbageCollectorStream { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("SqliteGarbageCollector").finish() + } +} + +#[async_trait] +impl ChangeStream for SqliteGarbageCollectorStream { + fn write(&self, id: &crate::id::ObjectId, _size: u64, expires_at: Option) { + let pool = self.pool.clone(); + let id = id.clone(); + + // Increment counter + let active_tasks = Arc::clone(&self.active_tasks); + active_tasks.fetch_add(1, Ordering::SeqCst); + + tokio::spawn(async move { + let _ = async { + let mut connection = pool.acquire().await?; + let mut tx_db = connection.begin().await?; + + sqlx::query("INSERT INTO garbage_collector (object_id, expires_at) VALUES (?, ?)") + .bind(id.as_storage_path().to_string()) + .bind(expires_at.map(|t| i64::try_from(t.as_secs()).ok())) + .execute(&mut *tx_db) + .await?; + + tx_db.commit().await?; + + Ok::<(), sqlx::Error>(()) + } + .await; + + // Decrement counter when done + active_tasks.fetch_sub(1, Ordering::SeqCst); + }); + } + + fn update(&self, id: &crate::id::ObjectId, expires_at: Option) { + let pool = self.pool.clone(); + let id = id.clone(); + + // Increment counter + let active_tasks = Arc::clone(&self.active_tasks); + active_tasks.fetch_add(1, Ordering::SeqCst); + + tokio::spawn(async move { + let _ = async { + let mut connection = pool.acquire().await?; + let mut tx_db = connection.begin().await?; + + sqlx::query("UPDATE garbage_collector SET expires_at = ? WHERE object_id = ?") + .bind(expires_at.map(|t| i64::try_from(t.as_secs()).ok())) + .bind(id.as_storage_path().to_string()) + .execute(&mut *tx_db) + .await?; + + tx_db.commit().await?; + + Ok::<(), sqlx::Error>(()) + } + .await; + + active_tasks.fetch_sub(1, Ordering::SeqCst); + }); + } + + fn delete(&self, id: &crate::id::ObjectId) { + let pool = self.pool.clone(); + let id = id.clone(); + + // Increment counter + let active_tasks = Arc::clone(&self.active_tasks); + active_tasks.fetch_add(1, Ordering::SeqCst); + + tokio::spawn(async move { + let _ = async { + let mut connection = pool.acquire().await?; + let mut tx_db = connection.begin().await?; + + sqlx::query("DELETE FROM garbage_collector WHERE object_id = ?") + .bind(id.as_storage_path().to_string()) + .execute(&mut *tx_db) + .await?; + + tx_db.commit().await?; + + Ok::<(), sqlx::Error>(()) + } + .await; + + active_tasks.fetch_sub(1, Ordering::SeqCst); + }); + } + + async fn join(&self, timeout: std::time::Duration) { + let start = std::time::Instant::now(); + + loop { + if self.active_tasks.load(Ordering::SeqCst) == 0 { + tracing::info!("All garbage collector tasks completed"); + return; + } + + if start.elapsed() > timeout { + let remaining = self.active_tasks.load(Ordering::SeqCst); + tracing::warn!( + "Timeout waiting for garbage collector: {} tasks still pending", + remaining + ); + return; + } + + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + } +} + +#[cfg(test)] +mod tests { + use crate::id::ObjectContext; + + use super::*; + use objectstore_types::scope::Scopes; + use objectstore_types::time::Timestamp; + use std::time::Duration; + use tempfile::NamedTempFile; + + fn create_test_config() -> SqliteGarbageCollectorConfig { + let temp_file = NamedTempFile::with_prefix("objectstore-gc-test").unwrap(); + SqliteGarbageCollectorConfig { + path: temp_file.path().to_str().unwrap().to_string(), + } + } + + fn create_test_id(s: &str) -> crate::id::ObjectId { + crate::id::ObjectId::new( + ObjectContext { + usecase: "test".to_string(), + scopes: Scopes::empty(), + }, + s.to_string(), + ) + } + + #[tokio::test] + async fn test_new_creates_stream() { + let config = create_test_config(); + let result = SqliteGarbageCollectorStream::new(&config).await; + assert!(result.is_ok()); + } + + #[tokio::test] + async fn test_write_increments_task_counter() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let id = create_test_id("test-object-1"); + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 0); + + stream.write(&id, 1024, None); + + // Counter should increment immediately + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 1); + + // Wait for task to complete + stream.join(Duration::from_secs(5)).await; + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn test_multiple_writes() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + for i in 0..5 { + let id = create_test_id(&format!("test-object-{}", i)); + stream.write(&id, 1024 * (i + 1) as u64, None); + } + + // All 5 tasks should be active + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 5); + + // Wait for all to complete + stream.join(Duration::from_secs(5)).await; + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn test_write_with_expiration() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let id = create_test_id("test-object-expiring"); + let expires_at = Some(Timestamp::from_unix_secs(1234567890).unwrap()); + + stream.write(&id, 2048, expires_at); + + stream.join(Duration::from_secs(5)).await; + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 0); + + // Verify the record was inserted + let count: (i64,) = + sqlx::query_as("SELECT COUNT(*) FROM garbage_collector WHERE object_id = ?") + .bind(id.as_storage_path().to_string()) + .fetch_one(&stream.pool) + .await + .unwrap(); + + assert_eq!(count.0, 1); + } + + #[tokio::test] + async fn test_update() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let id = create_test_id("test-update-object"); + let initial_expires = Some(Timestamp::from_unix_secs(1000000).unwrap()); + let updated_expires = Some(Timestamp::from_unix_secs(2000000).unwrap()); + + // First write the record + stream.write(&id, 1024, initial_expires); + stream.join(Duration::from_secs(5)).await; + + // Update it + stream.update(&id, updated_expires); + stream.join(Duration::from_secs(5)).await; + + // Verify the update + let result: (i64,) = + sqlx::query_as("SELECT expires_at FROM garbage_collector WHERE object_id = ?") + .bind(id.as_storage_path().to_string()) + .fetch_one(&stream.pool) + .await + .unwrap(); + + assert_eq!(result.0, 2000000); + } + + #[tokio::test] + async fn test_delete() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let id = create_test_id("test-delete-object"); + + // Write the record + stream.write(&id, 1024, None); + stream.join(Duration::from_secs(5)).await; + + // Verify it exists + let count: (i64,) = + sqlx::query_as("SELECT COUNT(*) FROM garbage_collector WHERE object_id = ?") + .bind(id.as_storage_path().to_string()) + .fetch_one(&stream.pool) + .await + .unwrap(); + assert_eq!(count.0, 1); + + // Delete it + stream.delete(&id); + stream.join(Duration::from_secs(5)).await; + + // Verify it's gone + let count: (i64,) = + sqlx::query_as("SELECT COUNT(*) FROM garbage_collector WHERE object_id = ?") + .bind(id.as_storage_path().to_string()) + .fetch_one(&stream.pool) + .await + .unwrap(); + assert_eq!(count.0, 0); + } + + #[tokio::test] + async fn test_mixed_operations() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let id1 = create_test_id("mixed-1"); + let id2 = create_test_id("mixed-2"); + let id3 = create_test_id("mixed-3"); + + // Write 3 records + stream.write(&id1, 1024, None); + stream.write(&id2, 2048, None); + stream.write(&id3, 4096, None); + + // Update one + stream.update(&id2, Some(Timestamp::from_unix_secs(9999).unwrap())); + + // Delete one + stream.delete(&id1); + + stream.join(Duration::from_secs(5)).await; + + // Verify state + let count: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM garbage_collector") + .fetch_one(&stream.pool) + .await + .unwrap(); + assert_eq!(count.0, 2); // Only id2 and id3 remain + } + + #[tokio::test] + async fn test_join_timeout() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let id = create_test_id("timeout-test"); + stream.write(&id, 1024, None); + + // Very short timeout (should still complete since task is fast) + let start = std::time::Instant::now(); + stream.join(Duration::from_millis(1)).await; + let elapsed = start.elapsed(); + + // Should timeout quickly + assert!(elapsed < Duration::from_secs(1)); + } + + #[tokio::test] + async fn test_concurrent_operations() { + let config = create_test_config(); + let stream = std::sync::Arc::new(SqliteGarbageCollectorStream::new(&config).await.unwrap()); + + let mut handles = vec![]; + + // Spawn 10 concurrent tasks doing writes + for i in 0..10 { + let stream_clone = std::sync::Arc::clone(&stream); + let handle = tokio::spawn(async move { + let id = create_test_id(&format!("concurrent-{}", i)); + stream_clone.write(&id, 1024 * (i + 1) as u64, None); + }); + handles.push(handle); + } + + // Wait for all spawns to complete + for handle in handles { + handle.await.unwrap(); + } + + // All writes should still be processing + let active = stream.active_tasks.load(Ordering::SeqCst); + assert!(active > 0); + + // Wait for all to complete + stream.join(Duration::from_secs(5)).await; + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 0); + + // Verify all records were inserted + let count: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM garbage_collector") + .fetch_one(&stream.pool) + .await + .unwrap(); + assert_eq!(count.0, 10); + } + + #[tokio::test] + async fn test_join_waits_for_completion() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let id = create_test_id("wait-test"); + stream.write(&id, 1024, None); + + // join() should wait until all tasks complete + let start = std::time::Instant::now(); + stream.join(Duration::from_secs(10)).await; + let elapsed = start.elapsed(); + + // Should complete almost immediately (task is fast) + assert!(elapsed < Duration::from_secs(1)); + assert_eq!(stream.active_tasks.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn test_debug_impl() { + let config = create_test_config(); + let stream = SqliteGarbageCollectorStream::new(&config).await.unwrap(); + + let debug_str = format!("{:?}", stream); + assert!(debug_str.contains("SqliteGarbageCollector")); + } +} diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 0c9473f5..703ad2bc 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -20,6 +20,7 @@ use crate::id::ObjectId; #[cfg(feature = "storage-cogs")] mod cost_tracker; mod factory; +mod garbage_collector_sqlite; #[cfg(feature = "storage-cogs")] pub use cost_tracker::CostTrackerStream;