From 2f293123641ece7c692dbc754d0bf2847475fe92 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 12:36:23 +0200 Subject: [PATCH] feat(storage): Add CQL high-volume backend Support Cassandra and Scylla through revision-guarded lightweight transactions, native TTL, and optional authentication and TLS. Add the Cassandra devservice, CI setup, schema, documentation, and backend integration tests. --- .github/workflows/ci.yml | 25 + Cargo.lock | 216 ++++- Cargo.toml | 1 + README.md | 18 + devservices/config.yml | 30 +- devservices/run-cassandra.sh | 27 + objectstore-server/config/cql.example.yaml | 28 + objectstore-server/src/config.rs | 62 +- objectstore-service/Cargo.toml | 2 + objectstore-service/docs/architecture.md | 5 + objectstore-service/src/backend/cql.rs | 853 ++++++++++++++++++ .../src/backend/cql/schema.cql | 13 + objectstore-service/src/backend/cql/tests.rs | 696 ++++++++++++++ objectstore-service/src/backend/mod.rs | 12 +- objectstore-service/src/backend/tiered.rs | 3 +- objectstore-test/src/server.rs | 5 +- 16 files changed, 1974 insertions(+), 22 deletions(-) create mode 100644 devservices/run-cassandra.sh create mode 100644 objectstore-server/config/cql.example.yaml create mode 100644 objectstore-service/src/backend/cql.rs create mode 100644 objectstore-service/src/backend/cql/schema.cql create mode 100644 objectstore-service/src/backend/cql/tests.rs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index bb92ee76..f75e959f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -101,6 +101,20 @@ jobs: google/cloud-sdk \ /run-bigtable.sh + - name: Start Cassandra + run: | + docker run -d \ + --name cassandra \ + -p 9042:9042 \ + -e CASSANDRA_DC=datacenter1 \ + -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch \ + -e CASSANDRA_BROADCAST_RPC_ADDRESS=127.0.0.1 \ + -e MAX_HEAP_SIZE=512M \ + -e HEAP_NEWSIZE=100M \ + -v $PWD/devservices/run-cassandra.sh:/run-cassandra.sh:ro \ + -v $PWD/objectstore-service/src/backend/cql/schema.cql:/schema.cql:ro \ + cassandra:5.0.9 /bin/bash /run-cassandra.sh + - name: Build GCS Emulator run: | docker build -t gcs-emulator-local https://github.com/getsentry/google-cloud-storage-testbench.git#main @@ -137,6 +151,17 @@ jobs: - uses: taiki-e/install-action@cargo-llvm-cov + - name: Wait for Cassandra schema + run: | + deadline=$((SECONDS + 180)) + until docker exec cassandra cqlsh --connect-timeout=2 --request-timeout=5 localhost -e 'SELECT revision FROM objectstore.objects LIMIT 1'; do + if (( SECONDS >= deadline )); then + docker logs cassandra + exit 1 + fi + sleep 2 + done + - name: Run Rust doctests run: cargo test --workspace --all-features --doc --locked diff --git a/Cargo.lock b/Cargo.lock index aa2c00a7..e9ee753c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -585,6 +585,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" @@ -862,6 +868,54 @@ dependencies = [ "syn 2.0.118", ] +[[package]] +name = "darling" +version = "0.24.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed17f5901b6630b993ca003def43f2f8ef4014fc13b047b57aad617ff32bc2ec" +dependencies = [ + "darling_core", + "darling_macro", +] + +[[package]] +name = "darling_core" +version = "0.24.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6837e2cf7485aaae18f86181d2f0e9a7ed297a025e220aeabf63fdebd3a2ddff" +dependencies = [ + "ident_case", + "proc-macro2", + "quote", + "strsim", + "syn 3.0.6", +] + +[[package]] +name = "darling_macro" +version = "0.24.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2ac7135c3ef02b2f7833bbeb1be5ba7f966dcde8a87c6b87f65a778d71a02785" +dependencies = [ + "darling_core", + "quote", + "syn 3.0.6", +] + +[[package]] +name = "dashmap" +version = "6.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6361d5c062261c78a176addb82d4c821ae42bed6089de0e12603cd25de2059c" +dependencies = [ + "cfg-if", + "crossbeam-utils", + "hashbrown 0.14.5", + "lock_api", + "once_cell", + "parking_lot_core", +] + [[package]] name = "data-encoding" version = "2.11.0" @@ -1087,7 +1141,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -1507,6 +1561,12 @@ dependencies = [ "tracing", ] +[[package]] +name = "hashbrown" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" + [[package]] name = "hashbrown" version = "0.15.5" @@ -1835,7 +1895,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.5.10", + "socket2 0.6.4", "system-configuration", "tokio", "tower-service", @@ -1949,6 +2009,12 @@ dependencies = [ "zerovec", ] +[[package]] +name = "ident_case" +version = "1.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9e0384b61958566e926dc50660321d12159025e767c18e043daf26b70104c39" + [[package]] name = "idna" version = "1.1.0" @@ -2054,6 +2120,15 @@ dependencies = [ "either", ] +[[package]] +name = "itertools" +version = "0.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b4baf93f58d4425749ca49a51c50ebab072c5df6994d08fed93541c331481dc" +dependencies = [ + "either", +] + [[package]] name = "itoa" version = "1.0.18" @@ -2286,6 +2361,15 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +[[package]] +name = "lz4_flex" +version = "0.14.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ecbdfe44b1bd960b68170b417450a628c43f7cf56bb3c5317e61cb230ee7f226" +dependencies = [ + "twox-hash", +] + [[package]] name = "mappings" version = "0.7.2" @@ -2500,7 +2584,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -2957,6 +3041,8 @@ dependencies = [ "regex", "reqwest 0.13.4", "ring", + "rustls", + "scylla", "sentry", "serde", "serde_json", @@ -3409,7 +3495,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "03da047801ff44bb6a4d407d4860c05fd70bb81714e6b2f3812603d5b145b042" dependencies = [ "heck", - "itertools", + "itertools 0.14.0", "log", "multimap", "petgraph", @@ -3430,7 +3516,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" dependencies = [ "anyhow", - "itertools", + "itertools 0.14.0", "proc-macro2", "quote", "syn 2.0.118", @@ -3516,7 +3602,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls", - "socket2 0.5.10", + "socket2 0.6.4", "thiserror", "tokio", "tracing", @@ -3554,7 +3640,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.5.10", + "socket2 0.6.4", "tracing", "windows-sys 0.60.2", ] @@ -3676,6 +3762,15 @@ dependencies = [ "rand 0.10.2", ] +[[package]] +name = "rand_pcg" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b48ac3f7ffaab7fac4d2376632268aa5f89abdb55f7ebf8f4d11fffccb2320f7" +dependencies = [ + "rand_core 0.9.5", +] + [[package]] name = "rand_xoshiro" version = "0.7.0" @@ -3981,7 +4076,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -4040,7 +4135,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -4097,6 +4192,83 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" +[[package]] +name = "scylla" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29eebcb7e34257f8ce01aaa1469644825bd4b4995d401c362be65d994e3feaa3" +dependencies = [ + "arc-swap", + "async-trait", + "bytes", + "chrono", + "dashmap", + "futures", + "hashbrown 0.17.1", + "itertools 0.15.0", + "rand 0.9.4", + "rand_pcg", + "rustls", + "scylla-cql", + "scylla-cql-core", + "serde", + "serde_json", + "smallvec", + "socket2 0.6.4", + "thiserror", + "tokio", + "tokio-rustls", + "tracing", + "uuid", +] + +[[package]] +name = "scylla-cql" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a398b0e78fb3872c4d5afc74e461d6b80dd69d030e13ddc7442e5f32c482e0da" +dependencies = [ + "byteorder", + "bytes", + "chrono", + "itertools 0.15.0", + "lz4_flex", + "scylla-cql-core", + "snap", + "stable_deref_trait", + "thiserror", + "tokio", + "uuid", + "yoke", +] + +[[package]] +name = "scylla-cql-core" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72ca9de2eb08d04a9c85002ac179c51353c91152cd2cd58fc3d7cbc1eee8ab16" +dependencies = [ + "byteorder", + "bytes", + "chrono", + "itertools 0.15.0", + "scylla-macros", + "thiserror", + "uuid", +] + +[[package]] +name = "scylla-macros" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c82c9c67cb4912cefc8cb2a68ebdcbe0273b66510c2bb06581ae6b68d98498aa" +dependencies = [ + "darling", + "proc-macro2", + "quote", + "syn 3.0.6", +] + [[package]] name = "sec1" version = "0.7.3" @@ -4151,7 +4323,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5b55fb86dfd3a2f5f76ea78310a88f96c4ea21a3031f8d212443d56123fd0521" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -4545,6 +4717,12 @@ version = "1.15.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8ed6a63f02c8539c91a8685a86f4099661ba3da017932f6ebbea6de3f0fa7c90" +[[package]] +name = "snap" +version = "1.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "199905e6153d6405f9728fe44daace35f8f837bbf830bb6e85fbd5828709a886" + [[package]] name = "socket2" version = "0.5.10" @@ -4562,7 +4740,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -4613,6 +4791,12 @@ dependencies = [ "yansi", ] +[[package]] +name = "strsim" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" + [[package]] name = "subtle" version = "2.6.1" @@ -4698,7 +4882,7 @@ dependencies = [ "getrandom 0.4.3", "once_cell", "rustix", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -5148,6 +5332,12 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" +[[package]] +name = "twox-hash" +version = "2.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "86a801b3cea342a06d468c8710662aa29e5e05e4f5c0d62f00bbb7f2ad7941c2" + [[package]] name = "typenum" version = "1.20.1" @@ -5491,7 +5681,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 3282e21c..8196fa15 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -82,6 +82,7 @@ regex = "1.12.4" reqwest = { version = "0.13.4", default-features = false } ring = "0.17.14" rustls = { version = "0.23.40", default-features = false } +scylla = { version = "1.9.0", features = ["rustls-023"] } secrecy = "0.10.3" sentry = "0.48.3" sentry-options = "1.2.4" diff --git a/README.md b/README.md index a19a25a4..e2367990 100644 --- a/README.md +++ b/README.md @@ -193,6 +193,24 @@ column families. - For **Google Cloud Storage** (GCS), a test bucket is already configured in the dev container. - For **SeaweedFS** (S3-compatible), a public test bucket is created on startup. + - For **Cassandra**, the official `cassandra:5.0.9` image initializes an + `objectstore` keyspace and `objects` table, exposing CQL on `localhost:9042`. + Startup can take several minutes; wait for the service to become healthy. + +To use Cassandra as the high-volume tier, run with +[`objectstore-server/config/cql.example.yaml`](objectstore-server/config/cql.example.yaml). +For production, provision the keyspace with your intended replication settings, +then apply [`schema.cql`](objectstore-service/src/backend/cql/schema.cql) in that +keyspace. The backend does not create or migrate production schema. +All Objectstore instances accessing it must use the same local datacenter. +See the [CQL module documentation](https://getsentry.github.io/objectstore/objectstore_service/backend/cql/index.html) +for TTL limits and TLS requirements. + +Direct CQL backend tests run against the Cassandra devservice, without adding a +CQL variant of the tiered-storage suite. To run them against a separately +provisioned Cassandra or Scylla instance, set `CQL_TEST_CONFIG` to the path of a +JSON `CqlConfig` (the fields under `high_volume` in the example, without `type`), +then run `cargo test -p objectstore-service --all-features backend::cql`. The emulators example config is pre-configured with a tiered configuration using both backends. To use it: diff --git a/devservices/config.yml b/devservices/config.yml index a32d1d0e..c25b968c 100644 --- a/devservices/config.yml +++ b/devservices/config.yml @@ -7,6 +7,8 @@ x-sentry-service-config: description: Objectstore server bigtable: description: Google Cloud BigTable emulator + cassandra: + description: Apache Cassandra CQL storage gcs: description: Google Cloud Storage emulator seaweedfs: @@ -14,13 +16,39 @@ x-sentry-service-config: modes: default: [] containerized: [objectstore] - full: [bigtable, gcs, seaweedfs] + full: [bigtable, cassandra, gcs, seaweedfs] x-programs: devserver: command: cargo run -- services: + cassandra: + image: cassandra:5.0.9 + command: /bin/bash /run-cassandra.sh + environment: + CASSANDRA_DC: datacenter1 + CASSANDRA_ENDPOINT_SNITCH: GossipingPropertyFileSnitch + # The driver discovers this address after contacting the forwarded port. + CASSANDRA_BROADCAST_RPC_ADDRESS: 127.0.0.1 + MAX_HEAP_SIZE: 512M + HEAP_NEWSIZE: 100M + healthcheck: + test: ["CMD", "cqlsh", "localhost", "-e", "SELECT revision FROM objectstore.objects LIMIT 1"] + interval: 5s + timeout: 10s + retries: 36 + start_period: 30s + ports: + - 127.0.0.1:9042:9042 + volumes: + - ./run-cassandra.sh:/run-cassandra.sh:ro + - ../objectstore-service/src/backend/cql/schema.cql:/schema.cql:ro + networks: + - devservices + labels: + - orchestrator=devservices + restart: unless-stopped objectstore: image: ghcr.io/getsentry/objectstore:nightly environment: diff --git a/devservices/run-cassandra.sh b/devservices/run-cassandra.sh new file mode 100644 index 00000000..1bc60c44 --- /dev/null +++ b/devservices/run-cassandra.sh @@ -0,0 +1,27 @@ +#!/bin/bash +set -euo pipefail + +docker-entrypoint.sh cassandra -f & +CASSANDRA_PID=$! +trap 'kill "$CASSANDRA_PID" 2>/dev/null || true' EXIT +trap 'exit 143' TERM +trap 'exit 130' INT + +ready=false +STARTUP_DEADLINE=$((SECONDS + 180)) +while (( SECONDS < STARTUP_DEADLINE )); do + kill -0 "$CASSANDRA_PID" 2>/dev/null || exit 1 + if cqlsh --connect-timeout=2 --request-timeout=5 localhost -e 'SELECT release_version FROM system.local' >/dev/null 2>&1; then + ready=true + break + fi + sleep 2 +done +if [ "$ready" != true ]; then + echo 'Cassandra did not become ready within the startup deadline' >&2 + exit 1 +fi + +cqlsh localhost -e "CREATE KEYSPACE IF NOT EXISTS objectstore WITH replication = {'class': 'NetworkTopologyStrategy', 'datacenter1': 1}" +cqlsh localhost -k objectstore -f /schema.cql +wait "$CASSANDRA_PID" diff --git a/objectstore-server/config/cql.example.yaml b/objectstore-server/config/cql.example.yaml new file mode 100644 index 00000000..43937af6 --- /dev/null +++ b/objectstore-server/config/cql.example.yaml @@ -0,0 +1,28 @@ +# Requires devservices full mode; use a port different from the containerized server. +http_addr: 0.0.0.0:18888 + +storage: + type: tiered + high_volume: + type: cql + nodes: [localhost:9042] + keyspace: objectstore + table_name: objects + local_datacenter: datacenter1 + request_timeout: 5s + # Optional authentication and verified TLS: + # username: objectstore + # password: provided-through-environment + # tls_ca_bundle: /etc/objectstore/cql-ca.pem + # Certificates must include the node IP addresses as subject alternative names. + long_term: + type: gcs + endpoint: http://localhost:8087 + bucket: test-bucket + +auth: + enforce: false + +logging: + level: debug + format: auto diff --git a/objectstore-server/src/config.rs b/objectstore-server/src/config.rs index beb87d9d..7a96c42f 100644 --- a/objectstore-server/src/config.rs +++ b/objectstore-server/src/config.rs @@ -461,7 +461,7 @@ pub struct Config { /// # Environment Variables /// /// - `OS__STORAGE__TYPE` — backend type (`filesystem`, `tiered`, `gcs`, `bigtable`, - /// `s3compatible`) + /// `cql`, `s3compatible`) /// - Additional fields depending on the type (see [`StorageConfig`]) /// /// For tiered storage, sub-backend fields are nested under `high_volume` and `long_term`: @@ -1086,7 +1086,9 @@ mod tests { let StorageConfig::Tiered(c) = &dbg!(&config).storage else { panic!("expected tiered storage"); }; - let HighVolumeStorageConfig::BigTable(hv) = &c.high_volume; + let HighVolumeStorageConfig::BigTable(hv) = &c.high_volume else { + panic!("expected bigtable high_volume"); + }; assert_eq!(hv.project_id, "my-project"); assert_eq!(hv.rpc_timeout, Duration::from_secs(2)); let MultipartUploadStorageConfig::Gcs(lt) = &c.long_term else { @@ -1115,7 +1117,9 @@ mod tests { let StorageConfig::Tiered(c) = &dbg!(&config).storage else { panic!("expected tiered storage"); }; - let HighVolumeStorageConfig::BigTable(hv) = &c.high_volume; + let HighVolumeStorageConfig::BigTable(hv) = &c.high_volume else { + panic!("expected bigtable high_volume"); + }; assert_eq!(hv.project_id, "my-project"); assert_eq!(hv.instance_name, "my-instance"); assert_eq!(hv.table_name, "my-table"); @@ -1129,6 +1133,58 @@ mod tests { }); } + #[test] + fn cql_storage_via_yaml() { + figment::Jail::expect_with(|jail| { + jail.create_file( + "cql.yaml", + r#" +storage: + type: cql + nodes: [localhost:9042] + keyspace: objectstore + table_name: objects + local_datacenter: datacenter1 +"#, + )?; + let config = Config::load(Some(Path::new("cql.yaml"))).unwrap(); + let StorageConfig::Cql(cql) = config.storage else { + panic!("expected CQL") + }; + assert_eq!(cql.nodes, ["localhost:9042"]); + assert_eq!(cql.request_timeout, Duration::from_secs(5)); + assert!(cql.username.is_none()); + assert!(cql.tls_ca_bundle.is_none()); + Ok(()) + }); + } + + #[test] + fn tiered_cql_storage_via_env() { + figment::Jail::expect_with(|jail| { + jail.set_env("OS__STORAGE__TYPE", "tiered"); + jail.set_env("OS__STORAGE__HIGH_VOLUME__TYPE", "cql"); + jail.set_env("OS__STORAGE__HIGH_VOLUME__NODES", "[localhost:9042]"); + jail.set_env("OS__STORAGE__HIGH_VOLUME__KEYSPACE", "objectstore"); + jail.set_env("OS__STORAGE__HIGH_VOLUME__TABLE_NAME", "objects"); + jail.set_env("OS__STORAGE__HIGH_VOLUME__LOCAL_DATACENTER", "datacenter1"); + jail.set_env("OS__STORAGE__HIGH_VOLUME__REQUEST_TIMEOUT", "3s"); + jail.set_env("OS__STORAGE__LONG_TERM__TYPE", "filesystem"); + jail.set_env("OS__STORAGE__LONG_TERM__PATH", "/data/lt"); + let config = Config::load(None).unwrap(); + let StorageConfig::Tiered(tiered) = config.storage else { + panic!("expected tiered") + }; + let HighVolumeStorageConfig::Cql(cql) = tiered.high_volume else { + panic!("expected CQL") + }; + assert_eq!(cql.local_datacenter, "datacenter1"); + assert_eq!(cql.nodes, ["localhost:9042"]); + assert_eq!(cql.request_timeout, Duration::from_secs(3)); + Ok(()) + }); + } + #[test] fn storage_cogs_via_env() { figment::Jail::expect_with(|jail| { diff --git a/objectstore-service/Cargo.toml b/objectstore-service/Cargo.toml index d70938d9..c8690fb1 100644 --- a/objectstore-service/Cargo.toml +++ b/objectstore-service/Cargo.toml @@ -30,6 +30,8 @@ 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"] } ring = { workspace = true } +rustls = { workspace = true, features = ["ring", "std", "tls12"] } +scylla = { workspace = true } sentry = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index f195610f..418809bd 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -85,6 +85,11 @@ given object size is below this threshold. See [`backend::StorageConfig`] for available backend implementations. +[`CqlBackend`](backend::cql::CqlBackend) supports Cassandra and Scylla as either +standalone storage or the high-volume tier. It uses lightweight transactions +within one configured datacenter and native TTL for expiration. See [`backend::cql`] +for schema provisioning, consistency, connection security, and expiration limits. + ## Redirect Tombstones For large objects, `TieredStorage` stores a **redirect tombstone** in the diff --git a/objectstore-service/src/backend/cql.rs b/objectstore-service/src/backend/cql.rs new file mode 100644 index 00000000..6ae60945 --- /dev/null +++ b/objectstore-service/src/backend/cql.rs @@ -0,0 +1,853 @@ +//! Cassandra and Scylla storage using portable CQL and lightweight transactions. +//! +//! # Schema and consistency +//! +//! Pre-create a keyspace and apply `cql/schema.cql` in it before starting Objectstore. +//! Each `(namespace, path)` partition holds one inline object, redirect, or upload marker. +//! Objects and markers occupy separate namespaces. Every mutation is a lightweight +//! transaction (LWT) conditioned on a fresh revision UUID. Ordinary CQL writes must not be +//! mixed with these transactions. Conditional UPDATE creates rows without INSERT liveness +//! markers, so no immortal primary-key-only rows survive expiration. +//! +//! Reads use LOCAL_SERIAL; writes use LOCAL_SERIAL consensus and LOCAL_QUORUM commits. +//! All clients for a keyspace must use the same datacenter. Routing never fails over to another +//! datacenter. LWT introduces extra round trips, including on reads. Mutations are neither +//! automatically replayed nor speculatively executed; ambiguous failures remain errors. +//! +//! # Expiration +//! +//! The absolute deadline is checked against the operation's access time, with equality still +//! live. Native per-cell TTL removes expired data; disk reclamation occurs later during +//! compaction. All live columns share one TTL, and renewal rewrites the entire row (including +//! payload). Scylla's non-portable per-row TTL extension is not used. Application and database +//! clocks must be synchronized. Deadlines must precede `2038-01-19T03:14:06Z`, and native TTL +//! must not exceed 630720000 seconds. Unsupported deadlines are rejected, never clamped. +//! +//! # Connections +//! +//! Optional username/password authentication and rustls TLS are supported. A TLS CA bundle +//! enables certificate verification; node certificates must include node IP addresses in their +//! subject alternative names, as required by the driver. Client certificates are not supported. + +use std::fmt; +use std::path::PathBuf; +use std::sync::Arc; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use anyhow::{Context as _, ensure}; +use bytes::Bytes; +use futures_util::TryStreamExt as _; +use objectstore_types::metadata::Metadata; +use objectstore_types::range::ByteRange; +use objectstore_types::time::Timestamp; +use rustls::pki_types::{CertificateDer, pem::PemObject as _}; +use scylla::client::execution_profile::ExecutionProfile; +use scylla::client::session::Session; +use scylla::client::session_builder::SessionBuilder; +use scylla::errors::{DbError, ExecutionError, RequestAttemptError}; +use scylla::policies::host_filter::DcHostFilter; +use scylla::policies::load_balancing::DefaultPolicy; +use scylla::policies::retry::FallthroughRetryPolicy; +use scylla::response::query_result::QueryResult; +use scylla::statement::prepared::PreparedStatement; +use scylla::statement::{Consistency, SerialConsistency}; +use scylla::value::{CqlValue, Row}; +use serde::{Deserialize, Serialize}; +use uuid::Uuid; + +use super::common::{ + self, Backend, DeleteResponse, ExpiryUpdate, GetResponse, HighVolumeBackend, MetadataResponse, + PutResponse, SetExpiryResponse, TieredGet, TieredMetadata, TieredUpdate, TieredWrite, + Tombstone, +}; +use crate::change_stream::{ + ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream, +}; +use crate::error::{Error, ErrorKind, Result, ResultExt as _}; +use crate::id::ObjectId; +use crate::stream::{ChunkedBytes, ClientStream}; + +const OBJECTS: &str = "objects"; +const UPLOADS: &str = "uploads"; +const CAS_ATTEMPTS: usize = 3; +const MAX_TTL: u64 = 630_720_000; +const MAX_EXPIRATION: u64 = 2_147_483_646; + +/// Connection configuration for [`CqlBackend`]. +/// +/// The keyspace and table must already exist. All clients must use the same local datacenter. +/// +/// ```yaml +/// storage: +/// type: cql +/// nodes: [localhost:9042] +/// keyspace: objectstore +/// table_name: objects +/// local_datacenter: datacenter1 +/// ``` +#[derive(Clone, Deserialize, Serialize)] +pub struct CqlConfig { + /// Nonempty list of contact points, optionally including port (default 9042). + pub nodes: Vec, + /// Pre-created keyspace, using a lowercase CQL identifier. + pub keyspace: String, + /// Pre-created table, using a lowercase CQL identifier. + pub table_name: String, + /// Datacenter to which all requests and connections are restricted. + pub local_datacenter: String, + /// Optional username. Must be configured together with `password`. + pub username: Option, + /// Optional password. Redacted from debug output. + pub password: Option, + /// PEM CA bundle enabling verified TLS. `None` selects plaintext transport. + pub tls_ca_bundle: Option, + /// Timeout for each CQL request. Defaults to five seconds. + #[serde(default = "default_request_timeout", with = "humantime_serde")] + pub request_timeout: Duration, + /// Optional per-backend storage cost attribution. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cogs: Option, +} + +impl fmt::Debug for CqlConfig { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("CqlConfig") + .field("nodes", &self.nodes) + .field("keyspace", &self.keyspace) + .field("table_name", &self.table_name) + .field("local_datacenter", &self.local_datacenter) + .field("username", &self.username.as_ref().map(|_| "[redacted]")) + .field("password", &self.password.as_ref().map(|_| "[redacted]")) + .field("tls_ca_bundle", &self.tls_ca_bundle) + .field("request_timeout", &self.request_timeout) + .field("cogs", &self.cogs) + .finish() + } +} + +fn default_request_timeout() -> Duration { + Duration::from_secs(5) +} + +impl CqlConfig { + fn validate(&self) -> anyhow::Result<()> { + ensure!( + !self.nodes.is_empty() && self.nodes.iter().all(|s| !s.trim().is_empty()), + "CQL nodes must not be empty" + ); + for identifier in [&self.keyspace, &self.table_name] { + ensure!( + valid_identifier(identifier), + "invalid CQL keyspace/table identifier" + ); + } + ensure!( + !self.local_datacenter.trim().is_empty(), + "CQL local_datacenter is required" + ); + ensure!( + self.username.is_some() == self.password.is_some(), + "CQL username and password must be supplied together" + ); + ensure!( + !self.request_timeout.is_zero(), + "CQL request_timeout must be positive" + ); + Ok(()) + } +} + +fn valid_identifier(value: &str) -> bool { + // Lowercase names avoid CQL's implicit case folding. Values are always bound separately. + !value.is_empty() + && value.len() <= 48 + && value.as_bytes()[0].is_ascii_lowercase() + && value + .bytes() + .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == b'_') +} + +/// High-volume storage backed by Cassandra or Scylla. +pub struct CqlBackend { + session: Session, + read: PreparedStatement, + read_metadata: PreparedStatement, + write: PreparedStatement, + delete: PreparedStatement, + change_stream: Arc, +} + +impl fmt::Debug for CqlBackend { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("CqlBackend").finish_non_exhaustive() + } +} + +#[derive(Debug)] +enum Record { + // Metadata-only reads leave payload empty; only full reads can be used for renewal. + Object(Metadata, Bytes), + Redirect(Tombstone), + Upload(Timestamp), +} + +impl Record { + fn expires_at(&self) -> Option { + match self { + Self::Object(metadata, _) => metadata.time_expires, + Self::Redirect(t) => t.time_expires, + Self::Upload(deadline) => Some(*deadline), + } + } + + fn live(&self, now: Timestamp) -> bool { + self.expires_at().is_none_or(|deadline| deadline >= now) + } + + fn redirect(&self, now: Timestamp) -> Option<&Tombstone> { + match self { + Self::Redirect(t) if self.live(now) => Some(t), + _ => None, + } + } +} + +struct StoredRecord { + revision: Uuid, + record: Record, +} + +// Both SELECTs have this prefix. The full SELECT appends the payload column. +type Head = (i8, Uuid, Option, Option, Option); +type FullRow = ( + i8, + Uuid, + Option, + Option, + Option, + Option>, +); + +fn corrupt(message: &'static str) -> Error { + Error::new(ErrorKind::CorruptData, message) +} + +/// Converts logical expiry to native TTL, including the logical deadline second. +fn native_ttl(deadline: Option, now_seconds: u64) -> Result { + let Some(deadline) = deadline else { + return Ok(0); + }; + if deadline.as_secs() >= MAX_EXPIRATION { + return Err(Error::new( + ErrorKind::InvalidMetadata, + "CQL expiration must precede 2038-01-19T03:14:06Z", + )); + } + let ttl = (deadline.as_secs() + 1).saturating_sub(now_seconds).max(1); + if ttl > MAX_TTL { + return Err(Error::new( + ErrorKind::InvalidMetadata, + "CQL TTL exceeds 630720000 seconds", + )); + } + Ok(ttl as i32) +} + +fn wall_clock_seconds() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("clock before Unix epoch") + .as_secs() +} + +fn execution_error(error: ExecutionError) -> Error { + use RequestAttemptError as Attempt; + let kind = match &error { + ExecutionError::RequestTimeout(_) => ErrorKind::BackendTimeout, + ExecutionError::EmptyPlan | ExecutionError::ConnectionPoolError(_) => { + ErrorKind::BackendUnavailable + } + ExecutionError::LastAttemptError(Attempt::BrokenConnectionError(_)) => { + ErrorKind::BackendUnavailable + } + ExecutionError::LastAttemptError(Attempt::DbError(db, _)) => match db { + DbError::ReadTimeout { .. } | DbError::WriteTimeout { .. } => ErrorKind::BackendTimeout, + DbError::Unavailable { .. } | DbError::IsBootstrapping => ErrorKind::BackendUnavailable, + DbError::Overloaded | DbError::RateLimitReached { .. } => ErrorKind::BackendRateLimited, + _ => ErrorKind::BackendFailure, + }, + _ => ErrorKind::BackendFailure, + }; + Error::with_source(kind, error) +} + +fn applied_value(value: Option<&Option>) -> Result { + match value { + Some(Some(CqlValue::Boolean(applied))) => Ok(*applied), + _ => Err(corrupt("missing or invalid CQL [applied] result")), + } +} + +fn applied(result: QueryResult) -> Result { + let rows = result.into_rows_result().kind(ErrorKind::CorruptData)?; + let column = rows + .column_specs() + .iter() + .position(|c| c.name() == "[applied]") + .ok_or_else(|| corrupt("missing CQL [applied] column"))?; + let row = rows.single_row::().kind(ErrorKind::CorruptData)?; + // Cassandra returns only [applied] on success; Scylla also returns old columns. + applied_value(row.columns.get(column)) +} + +impl CqlBackend { + /// Connects to the configured datacenter and prepares statements against an existing table. + /// + /// Returns an error for invalid configuration, TLS material, connection failure, or a missing + /// or incompatible schema. Does not create or migrate production tables. + pub async fn new(config: CqlConfig, streams: &ChangeStreamFactory) -> anyhow::Result { + config.validate()?; + let policy = DefaultPolicy::builder() + .prefer_datacenter(config.local_datacenter.clone()) + .permit_dc_failover(false) + .build(); + let profile = ExecutionProfile::builder() + .consistency(Consistency::LocalQuorum) + .serial_consistency(Some(SerialConsistency::LocalSerial)) + .request_timeout(Some(config.request_timeout)) + .load_balancing_policy(policy) + .retry_policy(Arc::new(FallthroughRetryPolicy)) + .speculative_execution_policy(None) + .build() + .into_handle(); + let mut builder = SessionBuilder::new() + .known_nodes(&config.nodes) + .host_filter(Arc::new(DcHostFilter::new(config.local_datacenter.clone()))) + .default_execution_profile_handle(profile); + if let (Some(user), Some(password)) = (&config.username, &config.password) { + builder = builder.user(user, password); + } + if let Some(path) = &config.tls_ca_bundle { + let mut roots = rustls::RootCertStore::empty(); + for cert in CertificateDer::pem_file_iter(path).with_context(|| "read CQL CA bundle")? { + roots.add(cert.with_context(|| "parse CQL CA certificate")?)?; + } + ensure!(!roots.is_empty(), "CQL CA bundle contains no certificates"); + let tls = rustls::ClientConfig::builder_with_provider(Arc::new( + rustls::crypto::ring::default_provider(), + )) + .with_safe_default_protocol_versions()? + .with_root_certificates(roots) + .with_no_client_auth(); + builder = builder.tls_context(Some(Arc::new(tls))); + } + let session = builder + .build() + .await + .with_context(|| "connect to CQL backend")?; + // Identifiers are validated above and quoted to also permit reserved words. + let table = format!("\"{}\".\"{}\"", config.keyspace, config.table_name); + let columns = "kind, revision, expires_at, metadata, redirect"; + let mut read = session + .prepare(format!( + "SELECT {columns}, payload FROM {table} WHERE namespace = ? AND path = ?" + )) + .await?; + let mut read_metadata = session + .prepare(format!( + "SELECT {columns} FROM {table} WHERE namespace = ? AND path = ?" + )) + .await?; + read.set_consistency(Consistency::LocalSerial); + read_metadata.set_consistency(Consistency::LocalSerial); + let write = session.prepare(format!("UPDATE {table} USING TTL ? SET kind = ?, revision = ?, expires_at = ?, metadata = ?, redirect = ?, payload = ? WHERE namespace = ? AND path = ? IF revision = ?")).await?; + let delete = session + .prepare(format!( + "DELETE FROM {table} WHERE namespace = ? AND path = ? IF revision = ?" + )) + .await?; + Ok(Self { + session, + read, + read_metadata, + write, + delete, + change_stream: streams.build(config.cogs.as_ref()), + }) + } + + async fn read_record( + &self, + namespace: &str, + id: &ObjectId, + payload: bool, + ) -> Result> { + let path = id.as_storage_path().to_string(); + let query = if payload { + &self.read + } else { + &self.read_metadata + }; + let result = self + .session + .execute_unpaged(query, (namespace, &path)) + .await + .map_err(execution_error)?; + let rows = result.into_rows_result().kind(ErrorKind::CorruptData)?; + let row = if payload { + rows.maybe_first_row::() + .kind(ErrorKind::CorruptData)? + } else { + rows.maybe_first_row::() + .kind(ErrorKind::CorruptData)? + .map(|(kind, revision, expiry, metadata, redirect)| { + (kind, revision, expiry, metadata, redirect, None) + }) + }; + let Some((kind, revision, expiry, metadata, redirect, bytes)) = row else { + return Ok(None); + }; + let expiry = expiry + .map(|seconds| { + Timestamp::from_unix_secs( + seconds + .try_into() + .map_err(|_| corrupt("negative CQL expiry"))?, + ) + .map_err(|_| corrupt("invalid CQL expiry")) + }) + .transpose()?; + let record = match (namespace, kind) { + (OBJECTS, 0) => { + if redirect.is_some() || (payload && bytes.is_none()) { + return Err(corrupt("invalid inline CQL row")); + } + let mut metadata: Metadata = + serde_json::from_str(&metadata.ok_or_else(|| corrupt("missing CQL metadata"))?) + .kind(ErrorKind::CorruptData)?; + metadata.time_expires = expiry; + Record::Object(metadata, bytes.unwrap_or_default().into()) + } + (OBJECTS, 1) => { + if metadata.is_some() || bytes.is_some() { + return Err(corrupt("invalid CQL redirect row")); + } + let target = redirect + .as_deref() + .and_then(ObjectId::from_storage_path) + .ok_or_else(|| corrupt("invalid CQL redirect target"))?; + Record::Redirect(Tombstone { + target, + time_expires: expiry, + }) + } + (UPLOADS, 2) => { + Record::Upload(expiry.ok_or_else(|| corrupt("missing upload marker expiry"))?) + } + _ => return Err(corrupt("invalid CQL row kind")), + }; + Ok(Some(StoredRecord { revision, record })) + } + + /// Writes all cells under one TTL and revision condition. `None` expects absence. + async fn write_record( + &self, + namespace: &str, + id: &ObjectId, + expected: Option, + record: &Record, + expiry_update: bool, + ) -> Result { + let path = id.as_storage_path().to_string(); + let expiry = record.expires_at(); + let ttl = native_ttl(expiry, wall_clock_seconds())?; + let (kind, metadata, redirect, payload) = match record { + Record::Object(metadata, payload) => { + let mut metadata = metadata.clone(); + metadata.size = Some(payload.len()); + ( + 0_i8, + Some(serde_json::to_string(&metadata).kind(ErrorKind::InvalidMetadata)?), + None, + Some(payload.as_ref()), + ) + } + Record::Redirect(t) => (1, None, Some(t.target.as_storage_path().to_string()), None), + Record::Upload(_) => (2, None, None, None), + }; + let size = (namespace.len() + + path.len() + + 1 + + 16 + + expiry.map_or(0, |_| 8) + + metadata.as_ref().map_or(0, String::len) + + redirect.as_ref().map_or(0, String::len) + + payload.map_or(0, <[u8]>::len)) as u64; + let result = self + .session + .execute_unpaged( + &self.write, + ( + ttl, + kind, + Uuid::new_v4(), + expiry.map(|t| t.as_secs() as i64), + metadata.as_deref(), + redirect.as_deref(), + payload, + namespace, + &path, + expected, + ), + ) + .await + .map_err(execution_error)?; + let applied = applied(result)?; + if applied && namespace == OBJECTS { + if expiry_update { + self.change_stream.update(id, expiry); + } else { + self.change_stream.write(id, size, expiry); + } + } + Ok(applied) + } + + async fn delete_record(&self, namespace: &str, id: &ObjectId, revision: Uuid) -> Result { + let path = id.as_storage_path().to_string(); + let result = self + .session + .execute_unpaged(&self.delete, (namespace, &path, revision)) + .await + .map_err(execution_error)?; + let applied = applied(result)?; + if applied && namespace == OBJECTS { + self.change_stream.delete(id); + } + Ok(applied) + } + + async fn replace(&self, namespace: &str, id: &ObjectId, record: &Record) -> Result<()> { + for _ in 0..CAS_ATTEMPTS { + let current = self.read_record(namespace, id, false).await?; + if self + .write_record(namespace, id, current.map(|r| r.revision), record, false) + .await? + { + return Ok(()); + } + } + Err(contention()) + } +} + +fn contention() -> Error { + Error::new( + ErrorKind::BackendFailure, + "CQL conditional mutation contention exhausted", + ) +} + +#[async_trait::async_trait] +impl Backend for CqlBackend { + fn name(&self) -> &'static str { + "cql" + } + + async fn put_object( + &self, + id: &ObjectId, + metadata: &Metadata, + mut stream: ClientStream, + _access_time: Timestamp, + ) -> Result { + let mut payload = ChunkedBytes::new(0); + while let Some(chunk) = stream.try_next().await? { + payload.push(chunk); + } + self.replace( + OBJECTS, + id, + &Record::Object(metadata.clone(), payload.into_bytes()), + ) + .await + } + + async fn get_object( + &self, + id: &ObjectId, + access_time: Timestamp, + range: Option, + ) -> Result { + match self.get_tiered_object(id, access_time, range).await? { + TieredGet::Object(m, r, p) => Ok(Some((m, r, p))), + TieredGet::NotFound => Ok(None), + TieredGet::Tombstone(_) => Err(ErrorKind::UnexpectedTombstone.into()), + } + } + + async fn get_metadata( + &self, + id: &ObjectId, + access_time: Timestamp, + ) -> Result { + match self.get_tiered_metadata(id, access_time).await? { + TieredMetadata::Object(m) => Ok(Some(m)), + TieredMetadata::NotFound => Ok(None), + TieredMetadata::Tombstone(_) => Err(ErrorKind::UnexpectedTombstone.into()), + } + } + + async fn set_expiry( + &self, + id: &ObjectId, + target: ExpiryUpdate, + access_time: Timestamp, + ) -> Result { + self.compare_and_update(id, None, TieredUpdate::SetExpiry(target), access_time) + .await + } + + async fn delete_object( + &self, + id: &ObjectId, + _access_time: Timestamp, + ) -> Result { + for _ in 0..CAS_ATTEMPTS { + let Some(current) = self.read_record(OBJECTS, id, false).await? else { + return Ok(()); + }; + if self.delete_record(OBJECTS, id, current.revision).await? { + return Ok(()); + } + } + Err(contention()) + } + + async fn join(&self) { + flush_change_stream(&self.change_stream).await; + } +} + +#[async_trait::async_trait] +impl HighVolumeBackend for CqlBackend { + async fn create_upload_marker( + &self, + revision: &ObjectId, + time_expires: Timestamp, + ) -> Result<()> { + self.replace(UPLOADS, revision, &Record::Upload(time_expires)) + .await + } + + async fn has_upload_marker(&self, revision: &ObjectId, access_time: Timestamp) -> Result { + Ok(self + .read_record(UPLOADS, revision, false) + .await? + .is_some_and(|r| r.record.live(access_time))) + } + + async fn delete_upload_marker( + &self, + revision: &ObjectId, + access_time: Timestamp, + ) -> Result { + let Some(row) = self.read_record(UPLOADS, revision, false).await? else { + return Ok(false); + }; + if !row.record.live(access_time) { + return Ok(false); + } + self.delete_record(UPLOADS, revision, row.revision).await + } + + async fn put_non_tombstone( + &self, + id: &ObjectId, + metadata: &Metadata, + payload: Bytes, + access_time: Timestamp, + ) -> Result> { + let record = Record::Object(metadata.clone(), payload); + for _ in 0..CAS_ATTEMPTS { + let row = self.read_record(OBJECTS, id, false).await?; + if let Some(t) = row.as_ref().and_then(|r| r.record.redirect(access_time)) { + return Ok(Some(t.clone())); + } + if self + .write_record(OBJECTS, id, row.map(|r| r.revision), &record, false) + .await? + { + return Ok(None); + } + } + Err(contention()) + } + + async fn get_tiered_object( + &self, + id: &ObjectId, + access_time: Timestamp, + range: Option, + ) -> Result { + let Some(row) = self.read_record(OBJECTS, id, true).await? else { + return Ok(TieredGet::NotFound); + }; + if !row.record.live(access_time) { + return Ok(TieredGet::NotFound); + } + match row.record { + Record::Object(metadata, payload) => { + let total = payload.len() as u64; + let (range, payload) = match range { + None => (None, payload), + Some(range) => { + let range = range + .resolve(total) + .ok_or(ErrorKind::RangeNotSatisfiable { total })?; + let bytes = payload.slice(range.start as usize..=range.end as usize); + (Some(range), bytes) + } + }; + Ok(TieredGet::Object( + metadata, + range, + crate::stream::single(payload), + )) + } + Record::Redirect(t) => Ok(TieredGet::Tombstone(t)), + Record::Upload(_) => Err(corrupt("upload marker in object namespace")), + } + } + + async fn get_tiered_metadata( + &self, + id: &ObjectId, + access_time: Timestamp, + ) -> Result { + let Some(row) = self.read_record(OBJECTS, id, false).await? else { + return Ok(TieredMetadata::NotFound); + }; + if !row.record.live(access_time) { + return Ok(TieredMetadata::NotFound); + } + match row.record { + Record::Object(m, _) => Ok(TieredMetadata::Object(m)), + Record::Redirect(t) => Ok(TieredMetadata::Tombstone(t)), + Record::Upload(_) => Err(corrupt("upload marker in object namespace")), + } + } + + async fn delete_non_tombstone( + &self, + id: &ObjectId, + access_time: Timestamp, + ) -> Result> { + for _ in 0..CAS_ATTEMPTS { + let Some(row) = self.read_record(OBJECTS, id, false).await? else { + return Ok(None); + }; + if let Some(t) = row.record.redirect(access_time) { + return Ok(Some(t.clone())); + } + if self.delete_record(OBJECTS, id, row.revision).await? { + return Ok(None); + } + } + Err(contention()) + } + + async fn compare_and_write( + &self, + id: &ObjectId, + current: Option<&ObjectId>, + write: TieredWrite, + access_time: Timestamp, + ) -> Result { + let next_target = write.target().cloned(); + let record = match write { + TieredWrite::Object(m, p) => Some(Record::Object(m, p)), + TieredWrite::Tombstone(t) => Some(Record::Redirect(t)), + TieredWrite::Delete => None, + }; + for _ in 0..CAS_ATTEMPTS { + let row = self.read_record(OBJECTS, id, false).await?; + let actual = row + .as_ref() + .and_then(|r| r.record.redirect(access_time)) + .map(|t| &t.target); + if actual != current { + return Ok(actual == next_target.as_ref()); + } + let applied = match (&record, row) { + (Some(record), row) => { + self.write_record(OBJECTS, id, row.map(|r| r.revision), record, false) + .await? + } + (None, Some(row)) => self.delete_record(OBJECTS, id, row.revision).await?, + (None, None) => true, + }; + if applied { + return Ok(true); + } + } + Err(contention()) + } + + async fn compare_and_update( + &self, + id: &ObjectId, + current: Option<&ObjectId>, + update: TieredUpdate, + access_time: Timestamp, + ) -> Result { + let Some(mut row) = self.read_record(OBJECTS, id, true).await? else { + return Ok(SetExpiryResponse::NotFound); + }; + if !row.record.live(access_time) { + return Ok(SetExpiryResponse::NotFound); + } + let created = match &row.record { + Record::Object(m, _) if current.is_none() => m.time_created, + Record::Redirect(t) if current == Some(&t.target) => None, + _ => return Ok(SetExpiryResponse::Rejected), + }; + let Some(old_expiry) = row.record.expires_at() else { + return Ok(SetExpiryResponse::Rejected); + }; + let TieredUpdate::SetExpiry(target) = update; + let Some(deadline) = target.resolve(created, access_time)? else { + return Ok(SetExpiryResponse::Rejected); + }; + native_ttl(Some(deadline), wall_clock_seconds())?; + if old_expiry >= deadline { + return Ok(SetExpiryResponse::Satisfied(deadline)); + } + match &mut row.record { + Record::Object(m, _) => { + m.expiration_policy = common::extended_expiration_policy( + m.expiration_policy, + m.time_created, + old_expiry, + deadline, + )?; + m.time_expires = Some(deadline); + } + Record::Redirect(t) => t.time_expires = Some(deadline), + Record::Upload(_) => unreachable!(), + } + Ok( + if self + .write_record(OBJECTS, id, Some(row.revision), &row.record, true) + .await? + { + SetExpiryResponse::Satisfied(deadline) + } else { + SetExpiryResponse::Rejected + }, + ) + } +} + +#[cfg(test)] +mod tests; diff --git a/objectstore-service/src/backend/cql/schema.cql b/objectstore-service/src/backend/cql/schema.cql new file mode 100644 index 00000000..010c2c2f --- /dev/null +++ b/objectstore-service/src/backend/cql/schema.cql @@ -0,0 +1,13 @@ +-- Run in the pre-created keyspace. Change the table name if table_name differs. +-- All mutations must use LWT. Do not write directly with ordinary INSERT/UPDATE. +CREATE TABLE IF NOT EXISTS objects ( + namespace text, + path text, + kind tinyint, + revision uuid, + expires_at bigint, + metadata text, + payload blob, + redirect text, + PRIMARY KEY ((namespace, path)) +) WITH default_time_to_live = 0; diff --git a/objectstore-service/src/backend/cql/tests.rs b/objectstore-service/src/backend/cql/tests.rs new file mode 100644 index 00000000..e0050c25 --- /dev/null +++ b/objectstore-service/src/backend/cql/tests.rs @@ -0,0 +1,696 @@ +//! These tests use the Cassandra devservice. CQL_TEST_CONFIG may name a JSON CqlConfig +//! file to run the same suite against another pre-provisioned Cassandra or Scylla instance. + +use super::*; +use crate::backend::common::ExpiryTarget; +use crate::id::ObjectContext; +use crate::stream; +use objectstore_types::metadata::ExpirationPolicy; +use objectstore_types::scope::{Scope, Scopes}; + +fn test_config() -> CqlConfig { + if let Some(path) = std::env::var_os("CQL_TEST_CONFIG") { + return serde_json::from_slice(&std::fs::read(path).unwrap()).unwrap(); + } + CqlConfig { + nodes: vec!["localhost:9042".into()], + keyspace: "objectstore".into(), + table_name: "objects".into(), + local_datacenter: "datacenter1".into(), + username: None, + password: None, + tls_ca_bundle: None, + request_timeout: default_request_timeout(), + cogs: None, + } +} + +async fn backend() -> anyhow::Result { + CqlBackend::new(test_config(), &ChangeStreamFactory::default()).await +} + +fn id() -> ObjectId { + ObjectId::random(ObjectContext { + usecase: "testing".into(), + scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]), + }) +} + +fn expiring(now: Timestamp, seconds: u64) -> Metadata { + Metadata { + time_created: Some(now), + time_expires: Some(now + Duration::from_secs(seconds)), + expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(seconds)), + ..Metadata::default() + } +} + +async fn put( + backend: &CqlBackend, + id: &ObjectId, + metadata: &Metadata, + bytes: &'static str, +) -> Result<()> { + backend + .put_object(id, metadata, stream::single(bytes), Timestamp::now()) + .await +} + +async fn body(backend: &CqlBackend, id: &ObjectId) -> Result> { + let (_, _, stream) = backend + .get_object(id, Timestamp::now(), None) + .await? + .unwrap(); + stream::read_to_vec(stream).await +} + +#[test] +fn configuration_validation_and_redaction() { + let mut config = test_config(); + config.validate().unwrap(); + config.username = Some("private-user".into()); + assert!(config.validate().is_err()); + config.password = Some("private-password".into()); + config.validate().unwrap(); + let debug = format!("{config:?}"); + assert!(!debug.contains("private-user")); + assert!(!debug.contains("private-password")); + config.table_name = "objects; DROP TABLE objects".into(); + assert!(config.validate().is_err()); + for name in ["", "UPPER", "1table", "with space", "quoted\"name"] { + assert!(!valid_identifier(name)); + } + assert!(valid_identifier("objects_2")); + config.table_name = "objects".into(); + config.nodes.clear(); + assert!(config.validate().is_err()); + config.nodes.push("localhost".into()); + config.local_datacenter.clear(); + assert!(config.validate().is_err()); + config.local_datacenter = "datacenter1".into(); + config.request_timeout = Duration::ZERO; + assert!(config.validate().is_err()); +} + +#[test] +fn ttl_boundaries() { + let now = 1_800_000_000; + let timestamp = |seconds| Some(Timestamp::from_unix_secs(seconds).unwrap()); + assert_eq!(native_ttl(None, now).unwrap(), 0); + assert_eq!(native_ttl(timestamp(now - 1), now).unwrap(), 1); + assert_eq!(native_ttl(timestamp(now), now).unwrap(), 1); + assert_eq!(native_ttl(timestamp(now + 10), now).unwrap(), 11); + assert_eq!( + native_ttl(timestamp(MAX_EXPIRATION - 1), now).unwrap() as u64, + MAX_EXPIRATION - now + ); + assert_eq!( + native_ttl(timestamp(MAX_EXPIRATION), now) + .unwrap_err() + .kind(), + ErrorKind::InvalidMetadata + ); + assert_eq!( + native_ttl(timestamp(MAX_TTL - 1), 0).unwrap() as u64, + MAX_TTL + ); + assert_eq!( + native_ttl(timestamp(MAX_TTL), 0).unwrap_err().kind(), + ErrorKind::InvalidMetadata + ); +} + +#[test] +fn lwt_applied_decoding() { + for columns in [ + vec![Some(CqlValue::Boolean(true))], + vec![Some(CqlValue::Boolean(true)), None, Some(CqlValue::Int(2))], + ] { + assert!(applied_value(columns.first()).unwrap()); + } + assert!(!applied_value(Some(&Some(CqlValue::Boolean(false)))).unwrap()); + for column in [None, Some(None), Some(Some(CqlValue::Int(1)))] { + assert_eq!( + applied_value(column.as_ref()).unwrap_err().kind(), + ErrorKind::CorruptData + ); + } +} + +#[tokio::test] +async fn rejects_invalid_tls_bundle() -> anyhow::Result<()> { + let pem = tempfile::NamedTempFile::new()?; + let mut config = test_config(); + config.tls_ca_bundle = Some(pem.path().into()); + assert!( + CqlBackend::new(config, &ChangeStreamFactory::default()) + .await + .unwrap_err() + .to_string() + .contains("no certificates") + ); + Ok(()) +} + +#[tokio::test] +async fn object_roundtrip_ranges_and_delete() -> anyhow::Result<()> { + let backend = backend().await?; + let id = id(); + let now = Timestamp::now(); + assert!(backend.get_object(&id, now, None).await?.is_none()); + assert!(backend.get_metadata(&id, now).await?.is_none()); + backend.delete_object(&id, now).await?; + let mut metadata = expiring(now, 3600); + metadata.content_type = "text/plain".into(); + metadata.custom.insert("key".into(), "value".into()); + put(&backend, &id, &metadata, "hello world").await?; + let actual = backend.get_metadata(&id, now).await?.unwrap(); + assert_eq!(actual.time_expires, metadata.time_expires); + assert_eq!(actual.time_created, metadata.time_created); + assert_eq!(actual.custom, metadata.custom); + assert_eq!(actual.content_type, metadata.content_type); + assert_eq!(actual.size, Some(11)); + assert_eq!(body(&backend, &id).await?, b"hello world"); + for (range, expected, start) in [ + (ByteRange::Bounded(1, 3), "ell", 1), + (ByteRange::From(6), "world", 6), + (ByteRange::Last(3), "rld", 8), + ] { + let (_, range, stream) = backend.get_object(&id, now, Some(range)).await?.unwrap(); + let range = range.unwrap(); + assert_eq!(range.start, start); + assert_eq!(range.total, 11); + assert_eq!(stream::read_to_vec(stream).await?, expected.as_bytes()); + } + assert_eq!( + backend + .get_object(&id, now, Some(ByteRange::From(99))) + .await + .err() + .unwrap() + .kind(), + ErrorKind::RangeNotSatisfiable { total: 11 } + ); + put(&backend, &id, &Metadata::default(), "").await?; + assert!(body(&backend, &id).await?.is_empty()); + assert_eq!(backend.get_metadata(&id, now).await?.unwrap().size, Some(0)); + backend.delete_object(&id, now).await?; + assert!(backend.get_metadata(&id, now).await?.is_none()); + backend.delete_object(&id, now).await?; + Ok(()) +} + +#[tokio::test] +async fn expiry_outcomes_and_policy() -> anyhow::Result<()> { + let backend = backend().await?; + let id = id(); + let now = Timestamp::now(); + let metadata = expiring(now, 60); + let deadline = metadata.time_expires.unwrap(); + let update = ExpiryTarget::At(now + Duration::from_secs(120)); + assert_eq!( + backend.set_expiry(&id, update.into(), now).await?, + SetExpiryResponse::NotFound + ); + put(&backend, &id, &Metadata::default(), "manual").await?; + assert_eq!( + backend.set_expiry(&id, update.into(), now).await?, + SetExpiryResponse::Rejected + ); + put(&backend, &id, &metadata, "payload").await?; + assert!(backend.get_metadata(&id, deadline).await?.is_some()); + assert!( + backend + .get_metadata(&id, deadline + Duration::from_secs(1)) + .await? + .is_none() + ); + assert_eq!( + backend + .set_expiry(&id, update.into(), deadline + Duration::from_secs(1)) + .await?, + SetExpiryResponse::NotFound + ); + assert_eq!( + backend + .set_expiry(&id, ExpiryTarget::At(deadline).into(), now) + .await?, + SetExpiryResponse::Satisfied(deadline) + ); + assert_eq!( + backend + .set_expiry( + &id, + ExpiryUpdate { + target: ExpiryTarget::At(deadline), + max: Some(Duration::from_secs(1)) + }, + now + ) + .await + .unwrap_err() + .kind(), + ErrorKind::InvalidMetadata + ); + let extended = now + Duration::from_secs(120); + assert_eq!( + backend + .set_expiry( + &id, + ExpiryTarget::FromCreation(Duration::from_secs(120)).into(), + now + ) + .await?, + SetExpiryResponse::Satisfied(extended) + ); + let actual = backend.get_metadata(&id, now).await?.unwrap(); + assert_eq!(actual.time_expires, Some(extended)); + assert_eq!( + actual.expiration_policy, + ExpirationPolicy::TimeToLive(Duration::from_secs(120)) + ); + assert_eq!(body(&backend, &id).await?, b"payload"); + let tti = Metadata { + expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_secs(60)), + ..metadata + }; + put(&backend, &id, &tti, "tti").await?; + backend.set_expiry(&id, update.into(), now).await?; + assert_eq!( + backend + .get_metadata(&id, now) + .await? + .unwrap() + .expiration_policy, + tti.expiration_policy + ); + let far_future = Timestamp::from_unix_secs(MAX_EXPIRATION).unwrap(); + assert_eq!( + backend + .set_expiry(&id, ExpiryTarget::At(far_future).into(), now) + .await + .unwrap_err() + .kind(), + ErrorKind::InvalidMetadata + ); + let no_creation = Metadata { + time_created: None, + ..tti + }; + put(&backend, &id, &no_creation, "no creation").await?; + assert_eq!( + backend + .set_expiry( + &id, + ExpiryTarget::FromCreation(Duration::from_secs(120)).into(), + now + ) + .await?, + SetExpiryResponse::Rejected + ); + backend.delete_object(&id, now).await?; + Ok(()) +} + +/// Poll native removal without applying the backend's logical expiry filter. +async fn wait_for_expiration(backend: &CqlBackend, id: &ObjectId) -> anyhow::Result<()> { + tokio::time::timeout(Duration::from_secs(15), async { + loop { + if backend.read_record(OBJECTS, id, true).await?.is_none() { + return Ok::<_, Error>(()); + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + }) + .await??; + Ok(()) +} + +#[tokio::test] +async fn native_ttl_renewal_and_nonexpiring_transitions() -> anyhow::Result<()> { + let backend = backend().await?; + let now = Timestamp::now(); + let expires = id(); + let renewed = id(); + let manual = id(); + let short = expiring(now, 2); + put(&backend, &expires, &Metadata::default(), "formerly manual").await?; + put(&backend, &expires, &short, "expires").await?; + put(&backend, &renewed, &short, "renewed").await?; + put(&backend, &manual, &short, "temporary").await?; + put(&backend, &manual, &Metadata::default(), "permanent").await?; + backend + .set_expiry( + &renewed, + ExpiryTarget::At(now + Duration::from_secs(60)).into(), + now, + ) + .await?; + wait_for_expiration(&backend, &expires).await?; + assert_eq!(body(&backend, &renewed).await?, b"renewed"); + assert_eq!(body(&backend, &manual).await?, b"permanent"); + // A new incarnation can be created after every cell from the old one expires. + put(&backend, &expires, &Metadata::default(), "recreated").await?; + assert_eq!(body(&backend, &expires).await?, b"recreated"); + for id in [&expires, &renewed, &manual] { + backend.delete_object(id, now).await?; + } + Ok(()) +} + +#[tokio::test] +async fn redirects_and_conditional_transitions() -> anyhow::Result<()> { + let backend = backend().await?; + let key = id(); + let now = Timestamp::now(); + let first = Tombstone { + target: id(), + time_expires: Some(now + Duration::from_secs(60)), + }; + let second = Tombstone { + target: id(), + time_expires: None, + }; + let write = || TieredWrite::Tombstone(first.clone()); + assert!(backend.compare_and_write(&key, None, write(), now).await?); + assert!(backend.compare_and_write(&key, None, write(), now).await?); + assert_eq!( + backend + .put_non_tombstone( + &key, + &Metadata::default(), + Bytes::from_static(b"blocked"), + now + ) + .await?, + Some(first.clone()) + ); + assert_eq!( + backend.delete_non_tombstone(&key, now).await?, + Some(first.clone()) + ); + assert!( + matches!(backend.get_tiered_object(&key, now, None).await?, TieredGet::Tombstone(t) if t == first) + ); + assert!( + matches!(backend.get_tiered_metadata(&key, now).await?, TieredMetadata::Tombstone(t) if t == first) + ); + assert_eq!( + backend.get_metadata(&key, now).await.unwrap_err().kind(), + ErrorKind::UnexpectedTombstone + ); + assert!( + !backend + .compare_and_write(&key, Some(&second.target), TieredWrite::Delete, now) + .await? + ); + let extension = ExpiryTarget::At(now + Duration::from_secs(120)); + assert_eq!( + backend + .compare_and_update(&key, None, TieredUpdate::SetExpiry(extension.into()), now) + .await?, + SetExpiryResponse::Rejected + ); + assert_eq!( + backend + .compare_and_update( + &key, + Some(&second.target), + TieredUpdate::SetExpiry(extension.into()), + now + ) + .await?, + SetExpiryResponse::Rejected + ); + assert_eq!( + backend + .compare_and_update( + &key, + Some(&first.target), + TieredUpdate::SetExpiry( + ExpiryTarget::FromCreation(Duration::from_secs(120)).into() + ), + now + ) + .await?, + SetExpiryResponse::Rejected + ); + assert_eq!( + backend + .compare_and_update( + &key, + Some(&first.target), + TieredUpdate::SetExpiry(extension.into()), + now + ) + .await?, + SetExpiryResponse::Satisfied(now + Duration::from_secs(120)) + ); + assert!( + backend + .compare_and_write( + &key, + Some(&first.target), + TieredWrite::Tombstone(second.clone()), + now + ) + .await? + ); + // Repeating the swap recognizes its destination even though the old target no longer matches. + assert!( + backend + .compare_and_write( + &key, + Some(&first.target), + TieredWrite::Tombstone(second.clone()), + now + ) + .await? + ); + assert!( + backend + .compare_and_write( + &key, + Some(&second.target), + TieredWrite::Object(Metadata::default(), Bytes::from_static(b"inline")), + now + ) + .await? + ); + assert_eq!(body(&backend, &key).await?, b"inline"); + assert!( + backend + .compare_and_write(&key, None, TieredWrite::Delete, now) + .await? + ); + assert!( + backend + .compare_and_write(&key, None, TieredWrite::Delete, now) + .await? + ); + assert!( + backend + .compare_and_write( + &key, + None, + TieredWrite::Object(Metadata::default(), Bytes::from_static(b"new")), + now + ) + .await? + ); + // Ordinary Backend writes and deletes also replace redirects atomically. + assert!(backend.compare_and_write(&key, None, write(), now).await?); + put(&backend, &key, &Metadata::default(), "ordinary").await?; + assert_eq!(body(&backend, &key).await?, b"ordinary"); + backend.compare_and_write(&key, None, write(), now).await?; + backend.delete_object(&key, now).await?; + Ok(()) +} + +#[tokio::test] +async fn expired_redirects_do_not_block_mutations() -> anyhow::Result<()> { + let backend = backend().await?; + let now = Timestamp::now(); + let expired_at = now + Duration::from_secs(60); + let access = expired_at + Duration::from_secs(1); + for operation in 0..3 { + let key = id(); + backend + .compare_and_write( + &key, + None, + TieredWrite::Tombstone(Tombstone { + target: id(), + time_expires: Some(expired_at), + }), + now, + ) + .await?; + assert!(matches!( + backend.get_tiered_object(&key, access, None).await?, + TieredGet::NotFound + )); + match operation { + 0 => assert!( + backend + .put_non_tombstone( + &key, + &Metadata::default(), + Bytes::from_static(b"new"), + access + ) + .await? + .is_none() + ), + 1 => assert!(backend.delete_non_tombstone(&key, access).await?.is_none()), + _ => assert!( + backend + .compare_and_write( + &key, + None, + TieredWrite::Object(Metadata::default(), Bytes::from_static(b"new")), + access + ) + .await? + ), + } + backend.delete_object(&key, now).await?; + } + Ok(()) +} + +#[tokio::test] +async fn stale_revision_cannot_overwrite_delete_or_renew_new_data() -> anyhow::Result<()> { + let backend = backend().await?; + let key = id(); + let now = Timestamp::now(); + put(&backend, &key, &expiring(now, 60), "old").await?; + let mut stale = backend.read_record(OBJECTS, &key, true).await?.unwrap(); + put(&backend, &key, &expiring(now, 60), "new").await?; + if let Record::Object(metadata, _) = &mut stale.record { + metadata.time_expires = Some(now + Duration::from_secs(120)); + } + assert!( + !backend + .write_record(OBJECTS, &key, Some(stale.revision), &stale.record, true) + .await? + ); + assert!(!backend.delete_record(OBJECTS, &key, stale.revision).await?); + assert_eq!(body(&backend, &key).await?, b"new"); + backend.delete_object(&key, now).await?; + assert!( + !backend + .write_record(OBJECTS, &key, Some(stale.revision), &stale.record, true) + .await? + ); + assert!(backend.get_metadata(&key, now).await?.is_none()); + Ok(()) +} + +#[tokio::test] +async fn competing_redirects_have_one_winner() -> anyhow::Result<()> { + let backend = backend().await?; + let key = id(); + let now = Timestamp::now(); + let first = TieredWrite::Tombstone(Tombstone { + target: id(), + time_expires: None, + }); + let second = TieredWrite::Tombstone(Tombstone { + target: id(), + time_expires: None, + }); + let (a, b) = tokio::join!( + backend.compare_and_write(&key, None, first, now), + backend.compare_and_write(&key, None, second, now) + ); + assert_ne!(a?, b?); + backend.delete_object(&key, now).await?; + Ok(()) +} + +#[tokio::test] +async fn upload_marker_isolation_expiry_and_single_consumption() -> anyhow::Result<()> { + let backend = backend().await?; + let key = id(); + let now = Timestamp::now(); + let deadline = now + Duration::from_secs(60); + assert!(!backend.has_upload_marker(&key, now).await?); + assert!(!backend.delete_upload_marker(&key, now).await?); + backend.create_upload_marker(&key, deadline).await?; + assert!(backend.get_metadata(&key, now).await?.is_none()); + assert!(backend.has_upload_marker(&key, deadline).await?); + assert!( + !backend + .has_upload_marker(&key, deadline + Duration::from_secs(1)) + .await? + ); + assert!( + !backend + .delete_upload_marker(&key, deadline + Duration::from_secs(1)) + .await? + ); + put(&backend, &key, &Metadata::default(), "object").await?; + let (a, b) = tokio::join!( + backend.delete_upload_marker(&key, now), + backend.delete_upload_marker(&key, now) + ); + assert_ne!(a?, b?); + assert!(!backend.delete_upload_marker(&key, now).await?); + assert_eq!(body(&backend, &key).await?, b"object"); + backend.delete_object(&key, now).await?; + Ok(()) +} + +#[cfg(feature = "storage-cogs")] +#[tokio::test] +async fn change_stream_only_reports_applied_object_mutations() -> anyhow::Result<()> { + use objectstore_inventory_tracker::OpType; + let (streams, producer) = crate::change_stream::dummy_factory(); + let mut config = test_config(); + config.cogs = Some(CostTrackerStreamConfig { + shared_resource_id: "cql_objectstore".into(), + sample_rate: 1.0, + }); + let backend = CqlBackend::new(config, &streams).await?; + let key = id(); + let now = Timestamp::now(); + put(&backend, &key, &expiring(now, 60), "payload").await?; + let stale = backend.read_record(OBJECTS, &key, true).await?.unwrap(); + backend + .set_expiry( + &key, + ExpiryTarget::At(now + Duration::from_secs(120)).into(), + now, + ) + .await?; + assert!( + !backend + .write_record(OBJECTS, &key, Some(stale.revision), &stale.record, true) + .await? + ); + // An already-satisfied update and upload markers emit no events. + backend + .set_expiry( + &key, + ExpiryTarget::At(now + Duration::from_secs(60)).into(), + now, + ) + .await?; + backend + .create_upload_marker(&key, now + Duration::from_secs(60)) + .await?; + backend.delete_upload_marker(&key, now).await?; + backend.delete_object(&key, now).await?; + backend.join().await; + let records = producer.records(); + assert_eq!(records.len(), 3); + assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records[0].shared_resource_id, "cql_objectstore"); + assert_eq!(records[0].app_feature, "testing"); + assert!(records[0].size.unwrap() > 7); + assert_eq!(records[1].op_type, OpType::Update); + assert_eq!(records[2].op_type, OpType::Delete); + assert_eq!(records[0].record_id, records[2].record_id); + Ok(()) +} diff --git a/objectstore-service/src/backend/mod.rs b/objectstore-service/src/backend/mod.rs index a0b6b81d..564f7b67 100644 --- a/objectstore-service/src/backend/mod.rs +++ b/objectstore-service/src/backend/mod.rs @@ -1,7 +1,7 @@ //! Storage backend implementations. //! //! This module contains the [`Backend`](common::Backend) trait and its -//! implementations. Each backend adapts a specific storage system (BigTable, +//! implementations. Each backend adapts a specific storage system (BigTable, Cassandra/Scylla, //! GCS, local filesystem, S3-compatible) to a uniform interface that //! [`StorageService`](crate::StorageService) consumes. //! @@ -17,6 +17,7 @@ pub mod bigtable; pub mod changelog; pub mod common; pub mod counting; +pub mod cql; mod extensions; pub mod gcs; pub mod in_memory; @@ -51,6 +52,9 @@ pub enum StorageConfig { /// [Google Bigtable]: https://cloud.google.com/bigtable BigTable(bigtable::BigTableConfig), + /// Cassandra or Scylla storage backend (type `"cql"`). + Cql(cql::CqlConfig), + /// Tiered storage backend (type `"tiered"`). /// /// Routes objects across two backends based on size: small objects go to @@ -92,6 +96,7 @@ async fn from_leaf_config( ), StorageConfig::Gcs(c) => Box::new(gcs::GcsBackend::new(c, streams).await?), StorageConfig::BigTable(c) => Box::new(bigtable::BigTableBackend::new(c, streams).await?), + StorageConfig::Cql(c) => Box::new(cql::CqlBackend::new(c, streams).await?), StorageConfig::Tiered(_) => anyhow::bail!("nested tiered storage is not supported"), }) } @@ -99,10 +104,12 @@ async fn from_leaf_config( /// Configuration for the high-volume backend in a [`tiered::TieredStorageConfig`]. /// /// Only backends that implement [`common::HighVolumeBackend`] are valid here. -/// Currently this is limited to BigTable. +/// Supported implementations are BigTable and CQL (Cassandra or Scylla). #[derive(Debug, Clone, Deserialize, Serialize)] #[serde(tag = "type", rename_all = "lowercase")] pub enum HighVolumeStorageConfig { + /// Cassandra or Scylla backend. + Cql(cql::CqlConfig), /// [Google Bigtable] backend. /// /// [Google Bigtable]: https://cloud.google.com/bigtable @@ -115,6 +122,7 @@ async fn hv_from_config( streams: &ChangeStreamFactory, ) -> Result> { Ok(match config { + HighVolumeStorageConfig::Cql(c) => Box::new(cql::CqlBackend::new(c, streams).await?), HighVolumeStorageConfig::BigTable(c) => { Box::new(bigtable::BigTableBackend::new(c, streams).await?) } diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 97bee063..ca639922 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -187,8 +187,7 @@ fn new_long_term_revision(id: &ObjectId) -> ObjectId { pub struct TieredStorageConfig { /// Backend for high-volume, small objects. /// - /// Must be a backend that implements [`HighVolumeBackend`] (currently - /// only BigTable). + /// Must be a backend that implements [`HighVolumeBackend`] (BigTable or CQL). pub high_volume: HighVolumeStorageConfig, /// Backend for large, long-term objects. /// diff --git a/objectstore-test/src/server.rs b/objectstore-test/src/server.rs index 459a6a65..4843be1d 100644 --- a/objectstore-test/src/server.rs +++ b/objectstore-test/src/server.rs @@ -142,6 +142,9 @@ fn replace_fs_paths(config: &mut StorageConfig, tempdirs: &mut Vec) { tempdirs.push(dir); } } - StorageConfig::S3Compatible(_) | StorageConfig::Gcs(_) | StorageConfig::BigTable(_) => {} + StorageConfig::S3Compatible(_) + | StorageConfig::Gcs(_) + | StorageConfig::BigTable(_) + | StorageConfig::Cql(_) => {} } }