diff --git a/src/api/compress.rs b/src/api/compress.rs index 220cd92..ad93f71 100644 --- a/src/api/compress.rs +++ b/src/api/compress.rs @@ -299,12 +299,15 @@ async fn compress_json( } else { (saved_bytes as f64) * 100.0 / (original_size as f64) }; - let skip_charge = req.compression_rate == Some(100) - && req.target_size_bytes.is_none() - && format_in == format_out - && req.max_width.is_none() - && req.max_height.is_none(); - let charge_units = anonymous_reserved || (!skip_charge && compressed_size < original_size); + let charge_units = anonymous_reserved + || quota::output_consumes_unit( + req.compression_rate, + format_in == format_out, + req.max_width.is_some() || req.max_height.is_some(), + req.target_size_bytes.is_some(), + original_size, + compressed_size, + ); let task_id = Uuid::new_v4(); let file_id = Uuid::new_v4(); @@ -583,12 +586,14 @@ async fn compress_direct( } else { (saved_bytes as f64) * 100.0 / (original_size as f64) }; - let skip_charge = req.compression_rate == Some(100) - && req.target_size_bytes.is_none() - && format_in == format_out - && req.max_width.is_none() - && req.max_height.is_none(); - let charge_units = !skip_charge && compressed_size < original_size; + let charge_units = quota::output_consumes_unit( + req.compression_rate, + format_in == format_out, + req.max_width.is_some() || req.max_height.is_some(), + req.target_size_bytes.is_some(), + original_size, + compressed_size, + ); let task_id = Uuid::new_v4(); let file_id = Uuid::new_v4(); diff --git a/src/services/quota.rs b/src/services/quota.rs index cd39879..06f35f2 100644 --- a/src/services/quota.rs +++ b/src/services/quota.rs @@ -451,7 +451,8 @@ pub async fn settle_anonymous_task_reservation( AND f.status = 'completed' AND f.compressed_size < f.original_size AND NOT ( - tasks.compression_rate = 100 + -- NULL means the caller did not request the explicit 100% passthrough. + COALESCE(tasks.compression_rate = 100, false) AND f.original_format = f.output_format AND tasks.max_width IS NULL AND tasks.max_height IS NULL @@ -514,6 +515,19 @@ pub async fn settle_anonymous_task_reservation( Ok(Some(refundable)) } +pub(crate) fn output_consumes_unit( + compression_rate: Option, + same_format: bool, + has_resize: bool, + has_target_size: bool, + original_size: u64, + output_size: u64, +) -> bool { + let is_unmetered_passthrough = + compression_rate == Some(100) && same_format && !has_resize && !has_target_size; + !is_unmetered_passthrough && output_size < original_size +} + fn refundable_reserved_units(reserved: i32, total_files: i32, consumed_units: i32) -> u32 { let reserved = reserved.max(0); let total_files = total_files.max(0); @@ -556,6 +570,125 @@ fn utc8_date() -> NaiveDate { #[cfg(test)] mod tests { use super::*; + use crate::config::Config; + use crate::services::mail::Mailer; + use sqlx::postgres::PgPoolOptions; + use std::sync::Arc; + use tokio::sync::Semaphore; + + struct AnonymousSettlementFixture<'a> { + session_id: &'a str, + ip: IpAddr, + date: NaiveDate, + reserved_units: u32, + compression_rate: Option, + total_files: usize, + completed_files: usize, + } + + async fn build_test_state( + pool: sqlx::PgPool, + database_url: String, + redis_url: String, + ) -> AppState { + let mut config = Config::from_env().expect("load quota test config"); + config.database_url = database_url; + config.redis_url = redis_url; + config.mail_enabled = false; + config.mail_log_links_when_disabled = false; + config.anon_daily_units = 10; + let redis = redis::Client::open(config.redis_url.clone()) + .expect("create quota test Redis client") + .get_connection_manager() + .await + .expect("connect quota test Redis"); + AppState { + mailer: Arc::new(Mailer::new(&config).expect("create disabled quota test mailer")), + image_processing_semaphore: Arc::new(Semaphore::new(2)), + runtime_policy_cache: crate::services::settings::RuntimePolicyCache::new(), + storage_cache: crate::services::storage::StorageCache::new(), + config, + db: pool, + redis, + } + } + + async fn insert_anonymous_settlement_fixture( + pool: &sqlx::PgPool, + fixture: AnonymousSettlementFixture<'_>, + ) -> Uuid { + assert!(fixture.completed_files <= fixture.total_files); + let task_id = Uuid::new_v4(); + sqlx::query( + r#" + INSERT INTO tasks ( + id, session_id, client_ip, status, + compression_rate, total_files, completed_files, failed_files, + anonymous_units_reserved, anonymous_quota_date + ) VALUES ( + $1, $2, $3::inet, 'completed', + $4, $5, $6, $7, + $8, $9 + ) + "#, + ) + .bind(task_id) + .bind(fixture.session_id) + .bind(fixture.ip.to_string()) + .bind(fixture.compression_rate) + .bind(fixture.total_files as i32) + .bind(fixture.completed_files as i32) + .bind((fixture.total_files - fixture.completed_files) as i32) + .bind(fixture.reserved_units as i32) + .bind(fixture.date) + .execute(pool) + .await + .expect("insert anonymous settlement task"); + + for index in 0..fixture.total_files { + let completed = index < fixture.completed_files; + sqlx::query( + r#" + INSERT INTO task_files ( + id, task_id, original_name, original_format, output_format, + original_size, compressed_size, status + ) VALUES ( + $1, $2, $3, 'jpeg', 'jpeg', + 100, $4, $5::file_status + ) + "#, + ) + .bind(Uuid::new_v4()) + .bind(task_id) + .bind(format!("fixture-{index}.jpg")) + .bind(completed.then_some(50_i64)) + .bind(if completed { "completed" } else { "failed" }) + .execute(pool) + .await + .expect("insert anonymous settlement file"); + } + task_id + } + + async fn anonymous_quota_counts( + state: &AppState, + session_id: &str, + ip: IpAddr, + date: NaiveDate, + ) -> (i64, i64) { + let mut redis = state.redis.clone(); + let session_count: Option = redis::cmd("GET") + .arg(anonymous_session_key(session_id, date)) + .query_async(&mut redis) + .await + .expect("read anonymous session quota"); + let ip_count: Option = redis::cmd("GET") + .arg(anonymous_ip_key(ip, date)) + .query_async(&mut redis) + .await + .expect("read anonymous IP quota"); + (session_count.unwrap_or(0), ip_count.unwrap_or(0)) + } #[test] fn balance_keeps_redeemed_units_separate_from_plan_usage() { @@ -606,6 +739,29 @@ mod tests { assert_eq!(refundable_reserved_units(10, -1, -2), 0); } + #[test] + fn output_metering_matches_passthrough_contract() { + assert!(output_consumes_unit(None, true, false, false, 100, 50)); + assert!(!output_consumes_unit( + Some(100), + true, + false, + false, + 100, + 50 + )); + assert!(output_consumes_unit( + Some(100), + false, + false, + false, + 100, + 50 + )); + assert!(output_consumes_unit(Some(100), true, true, false, 100, 50)); + assert!(!output_consumes_unit(None, true, false, false, 100, 100)); + } + #[test] fn anonymous_quota_keys_use_the_reserved_date() { let date = NaiveDate::from_ymd_opt(2026, 7, 25).unwrap(); @@ -633,4 +789,194 @@ mod tests { "203.0.113.7" ); } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + #[ignore = "requires isolated IMAGEFORGE_TEST_DATABASE_URL and IMAGEFORGE_TEST_REDIS_URL"] + async fn anonymous_batch_settlement_charges_successes_and_refunds_only_unused_units() { + let database_url = std::env::var("IMAGEFORGE_TEST_DATABASE_URL") + .expect("IMAGEFORGE_TEST_DATABASE_URL must be set"); + assert!( + database_url.to_ascii_lowercase().contains("test"), + "refusing to run destructive integration test outside a test database" + ); + let redis_url = std::env::var("IMAGEFORGE_TEST_REDIS_URL") + .expect("IMAGEFORGE_TEST_REDIS_URL must be set"); + let pool = PgPoolOptions::new() + .max_connections(16) + .connect(&database_url) + .await + .expect("connect quota test database"); + sqlx::migrate!().run(&pool).await.expect("run migrations"); + let state = build_test_state(pool.clone(), database_url, redis_url).await; + let marker = Uuid::new_v4().simple().to_string(); + let mut cleanup = Vec::new(); + + let null_session = format!("quota-null-{marker}"); + let null_ip: IpAddr = "198.51.100.11".parse().expect("parse fixture IP"); + let null_date = reserve_anonymous_units(&state, &null_session, null_ip, 3) + .await + .expect("reserve NULL-rate batch quota"); + let null_task = insert_anonymous_settlement_fixture( + &pool, + AnonymousSettlementFixture { + session_id: &null_session, + ip: null_ip, + date: null_date, + reserved_units: 3, + compression_rate: None, + total_files: 3, + completed_files: 3, + }, + ) + .await; + cleanup.push((null_task, null_session.clone(), null_ip, null_date)); + assert_eq!( + settle_anonymous_task_reservation(&state, null_task) + .await + .expect("settle NULL-rate batch"), + Some(0) + ); + assert_eq!( + anonymous_quota_counts(&state, &null_session, null_ip, null_date).await, + (3, 3) + ); + + let passthrough_session = format!("quota-passthrough-{marker}"); + let passthrough_ip: IpAddr = "198.51.100.12".parse().expect("parse fixture IP"); + let passthrough_date = + reserve_anonymous_units(&state, &passthrough_session, passthrough_ip, 3) + .await + .expect("reserve passthrough batch quota"); + let passthrough_task = insert_anonymous_settlement_fixture( + &pool, + AnonymousSettlementFixture { + session_id: &passthrough_session, + ip: passthrough_ip, + date: passthrough_date, + reserved_units: 3, + compression_rate: Some(100), + total_files: 3, + completed_files: 3, + }, + ) + .await; + cleanup.push(( + passthrough_task, + passthrough_session.clone(), + passthrough_ip, + passthrough_date, + )); + assert_eq!( + settle_anonymous_task_reservation(&state, passthrough_task) + .await + .expect("settle passthrough batch"), + Some(3) + ); + assert_eq!( + anonymous_quota_counts( + &state, + &passthrough_session, + passthrough_ip, + passthrough_date, + ) + .await, + (0, 0) + ); + + let partial_session = format!("quota-partial-{marker}"); + let partial_ip: IpAddr = "198.51.100.13".parse().expect("parse fixture IP"); + let partial_date = reserve_anonymous_units(&state, &partial_session, partial_ip, 3) + .await + .expect("reserve partial batch quota"); + let partial_task = insert_anonymous_settlement_fixture( + &pool, + AnonymousSettlementFixture { + session_id: &partial_session, + ip: partial_ip, + date: partial_date, + reserved_units: 3, + compression_rate: None, + total_files: 3, + completed_files: 2, + }, + ) + .await; + cleanup.push(( + partial_task, + partial_session.clone(), + partial_ip, + partial_date, + )); + assert_eq!( + settle_anonymous_task_reservation(&state, partial_task) + .await + .expect("settle partial batch"), + Some(1) + ); + assert_eq!( + anonymous_quota_counts(&state, &partial_session, partial_ip, partial_date).await, + (2, 2) + ); + + let limit_session = format!("quota-limit-{marker}"); + let limit_ip: IpAddr = "198.51.100.14".parse().expect("parse fixture IP"); + let mut limit_date = None; + for batch in 0..2 { + let date = reserve_anonymous_units(&state, &limit_session, limit_ip, 5) + .await + .expect("reserve consecutive anonymous batch"); + limit_date = Some(date); + let task_id = insert_anonymous_settlement_fixture( + &pool, + AnonymousSettlementFixture { + session_id: &limit_session, + ip: limit_ip, + date, + reserved_units: 5, + compression_rate: None, + total_files: 5, + completed_files: 5, + }, + ) + .await; + cleanup.push((task_id, limit_session.clone(), limit_ip, date)); + assert_eq!( + settle_anonymous_task_reservation(&state, task_id) + .await + .expect("settle consecutive anonymous batch"), + Some(0), + "batch {batch} unexpectedly refunded consumed units" + ); + } + let limit_error = reserve_anonymous_units(&state, &limit_session, limit_ip, 1) + .await + .expect_err("daily anonymous quota was bypassed"); + assert_eq!(limit_error.code, ErrorCode::QuotaExceeded); + assert_eq!( + anonymous_quota_counts( + &state, + &limit_session, + limit_ip, + limit_date.expect("limit quota date"), + ) + .await, + (10, 10) + ); + + let mut redis = state.redis.clone(); + for (task_id, session_id, ip, date) in cleanup { + sqlx::query("DELETE FROM tasks WHERE id = $1") + .bind(task_id) + .execute(&pool) + .await + .expect("delete quota settlement fixture"); + let _: i64 = redis::cmd("DEL") + .arg(anonymous_session_key(&session_id, date)) + .arg(anonymous_ip_key(ip, date)) + .arg(format!("anon_quota_refund:{task_id}")) + .query_async(&mut redis) + .await + .expect("delete quota settlement Redis keys"); + } + } } diff --git a/src/worker/mod.rs b/src/worker/mod.rs index ab0364d..d3c09e1 100644 --- a/src/worker/mod.rs +++ b/src/worker/mod.rs @@ -1086,11 +1086,14 @@ async fn process_task_file( } else { (original_size.saturating_sub(compressed_size) as f64) * 100.0 / (original_size as f64) }; - let skip_charge = compression_rate == Some(100) - && format_in == format_out - && max_width.is_none() - && max_height.is_none(); - let charge_units = !skip_charge && compressed_size < original_size; + let charge_units = quota::output_consumes_unit( + compression_rate, + format_in == format_out, + max_width.is_some() || max_height.is_some(), + false, + original_size, + compressed_size, + ); let object_key = storage::result_attempt_key( ctx.retention_hours as i64,