From 3c6c88b2e78f4875458223bf665041cf2671d03f Mon Sep 17 00:00:00 2001 From: Dae-Cheol Noh Date: Tue, 22 Sep 2026 21:16:40 +0900 Subject: [PATCH] fix(amqp): emit S3 events directly in notification records (#8066) * fix(amqp): emit S3 events directly in notification records * test(amqp): avoid queued timestamp equality race --- crates/targets/src/target/amqp.rs | 91 +++++++++++++++++++++++- crates/targets/tests/amqp_integration.rs | 61 ++++++++++++++-- docs/operations/amqp-notifications.md | 51 +++++++++++++ 3 files changed, 194 insertions(+), 9 deletions(-) create mode 100644 docs/operations/amqp-notifications.md diff --git a/crates/targets/src/target/amqp.rs b/crates/targets/src/target/amqp.rs index 1af39301a..e1bcd9941 100644 --- a/crates/targets/src/target/amqp.rs +++ b/crates/targets/src/target/amqp.rs @@ -30,8 +30,8 @@ use crate::{ store::{Key, Store}, target::{ ChannelTargetType, EntityTarget, QueuedPayload, QueuedPayloadMeta, TargetDeliveryCounters, TargetDeliverySnapshot, - TargetTlsState, TargetType, build_queued_payload_with_records, build_target_tls_fingerprint, is_connectivity_error, - open_target_queue_store, persist_queued_payload_to_store, + TargetTlsState, TargetType, build_queued_payload, build_queued_payload_with_records, build_target_tls_fingerprint, + is_connectivity_error, open_target_queue_store, persist_queued_payload_to_store, }, }; use async_trait::async_trait; @@ -404,7 +404,10 @@ where } fn build_queued_payload(&self, event: &EntityTarget) -> Result { - build_queued_payload_with_records(event, vec![event.clone()]) + match self.args.target_type { + TargetType::NotifyEvent => build_queued_payload(event), + TargetType::AuditLog => build_queued_payload_with_records(event, vec![event.clone()]), + } } async fn get_or_connect(&self) -> Result, TargetError> { @@ -676,6 +679,7 @@ mod tests { use serde_json::json; use std::path::PathBuf; use std::sync::Arc; + use std::time::{SystemTime, UNIX_EPOCH}; use uuid::Uuid; fn valid_args() -> AMQPArgs { @@ -713,6 +717,87 @@ mod tests { }) } + fn notification_event() -> EntityTarget { + EntityTarget { + object_name: "incoming%2Fclip+%252F.mp4".to_string(), + bucket_name: "example-bucket".to_string(), + event_name: EventName::ObjectCreatedPut, + data: json!({ + "eventVersion": "2.0", + "eventSource": "aws:s3", + "eventName": "s3:ObjectCreated:Put", + "s3": { + "bucket": {"name": "example-bucket"}, + "object": {"key": "incoming%2Fclip+%252F.mp4", "size": 42, + "eTag": "example-etag", "versionId": "example-version"} + } + }), + } + } + + #[test] + fn notification_records_contain_the_event_directly() { + let target = AMQPTarget::new("notification".to_string(), valid_args()).unwrap(); + let event = notification_event(); + let queued = target.build_queued_payload(&event).unwrap(); + let payload: serde_json::Value = serde_json::from_slice(&queued.body).unwrap(); + + assert_eq!(payload["Records"], json!([event.data])); + assert_eq!(payload["Key"], "example-bucket/incoming/clip %2F.mp4"); + assert_eq!(payload["EventName"], "s3:ObjectCreated:Put"); + assert_eq!(payload["Records"][0]["s3"]["object"]["key"], event.object_name); + } + + #[test] + fn audit_records_preserve_the_entity_envelope() { + let mut args = valid_args(); + args.target_type = TargetType::AuditLog; + let target = AMQPTarget::new("audit".to_string(), args).unwrap(); + let event = test_event(); + let queued = target.build_queued_payload(&event).unwrap(); + let payload: serde_json::Value = serde_json::from_slice(&queued.body).unwrap(); + + assert_eq!(payload["Records"], json!([event.as_ref()])); + } + + #[tokio::test] + async fn queued_notification_preserves_the_flat_record_and_metadata() { + let mut args = unreachable_args(); + args.queue_dir = temp_store_dir("record-shape").to_string_lossy().to_string(); + let target = AMQPTarget::new("notification".to_string(), args.clone()).unwrap(); + let event = notification_event(); + let expected = target.build_queued_payload(&event).unwrap(); + let unix_time_ms = || { + u64::try_from( + SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("system clock should be after the Unix epoch") + .as_millis(), + ) + .expect("current Unix timestamp should fit in u64") + }; + let before_save = unix_time_ms(); + target.save(Arc::new(event.clone())).await.unwrap(); + let after_save = unix_time_ms(); + let store = target.store().unwrap(); + let keys = store.list(); + assert_eq!(keys.len(), 1); + let raw = store.get_raw(&keys[0]).unwrap(); + let queued = QueuedPayload::decode(&raw).unwrap(); + assert_eq!(queued.body, expected.body); + assert_eq!(queued.meta.event_name, expected.meta.event_name); + assert_eq!(queued.meta.bucket_name, expected.meta.bucket_name); + assert_eq!(queued.meta.object_name, expected.meta.object_name); + assert_eq!(queued.meta.content_type, expected.meta.content_type); + assert_eq!(queued.meta.payload_len, expected.meta.payload_len); + assert_eq!(queued.meta.dedup_id, expected.meta.dedup_id); + assert_eq!(queued.meta.failure, expected.meta.failure); + assert!((before_save..=after_save).contains(&queued.meta.queued_at_unix_ms)); + let payload: serde_json::Value = serde_json::from_slice(&queued.body).unwrap(); + assert_eq!(payload["Records"], json!([event.data])); + std::fs::remove_dir_all(args.queue_dir).unwrap(); + } + fn temp_store_dir(name: &str) -> PathBuf { std::env::temp_dir().join(format!("rustfs-amqp-target-{name}-{}", Uuid::new_v4())) } diff --git a/crates/targets/tests/amqp_integration.rs b/crates/targets/tests/amqp_integration.rs index 97a9924b5..b4f7d4fa5 100644 --- a/crates/targets/tests/amqp_integration.rs +++ b/crates/targets/tests/amqp_integration.rs @@ -32,9 +32,9 @@ use lapin::{ use rustfs_s3_types::EventName; use rustfs_targets::Target; use rustfs_targets::check_amqp_broker_available; -use rustfs_targets::target::EntityTarget; use rustfs_targets::target::TargetType; use rustfs_targets::target::amqp::{AMQPArgs, AMQPTarget}; +use rustfs_targets::target::{EntityTarget, QueuedPayload, QueuedPayloadMeta}; use serde_json::Value; use std::sync::Arc; use uuid::Uuid; @@ -67,7 +67,12 @@ fn entity_for(bucket: &str, object: &str) -> Arc bucket_name: bucket.to_string(), object_name: object.to_string(), event_name: EventName::ObjectCreatedPut, - data: serde_json::json!({"bucket": bucket, "object": object}), + data: serde_json::json!({ + "eventVersion": "2.0", + "eventSource": "aws:s3", + "eventName": "s3:ObjectCreated:Put", + "s3": {"bucket": {"name": bucket}, "object": {"key": object, "size": 42}} + }), }) } @@ -102,7 +107,7 @@ async fn bind_queue(queue: &str, routing_key: &str) -> lapin::Channel { channel } -async fn read_one(channel: &lapin::Channel, queue: &str) -> (Value, BasicProperties) { +async fn read_one_raw(channel: &lapin::Channel, queue: &str) -> (Vec, BasicProperties) { let msg = tokio::time::timeout(std::time::Duration::from_secs(5), async { loop { if let Some(msg) = channel @@ -120,8 +125,12 @@ async fn read_one(channel: &lapin::Channel, queue: &str) -> (Value, BasicPropert .expect("message should arrive"); let properties = msg.properties.clone(); - let payload = serde_json::from_slice(&msg.data).expect("message payload should be JSON"); - (payload, properties) + (msg.data.clone(), properties) +} + +async fn read_one(channel: &lapin::Channel, queue: &str) -> (Value, BasicProperties) { + let (body, properties) = read_one_raw(channel, queue).await; + (serde_json::from_slice(&body).expect("message payload should be JSON"), properties) } #[tokio::test] @@ -147,7 +156,7 @@ async fn test_direct_publish_delivers_json_payload() { let (payload, properties) = read_one(&channel, &queue).await; assert_eq!(payload["Key"], "bucket1/object-A"); - assert_eq!(payload["Records"][0]["data"]["bucket"], "bucket1"); + assert_eq!(payload["Records"], serde_json::json!([entity_for("bucket1", "object-A").data])); assert_eq!(properties.content_type().as_ref().map(|s| s.as_str()), Some("application/json")); assert_eq!(*properties.delivery_mode(), Some(2)); @@ -209,6 +218,7 @@ async fn test_queue_replay_delivers_and_removes_stored_payload() { let (payload, properties) = read_one(&channel, &queue).await; assert_eq!(payload["Key"], "bucket1/object-B"); + assert_eq!(payload["Records"], serde_json::json!([entity_for("bucket1", "object-B").data])); assert_eq!(properties.content_type().as_ref().map(|s| s.as_str()), Some("application/json")); assert_eq!(*properties.delivery_mode(), Some(2)); assert_eq!(target.delivery_snapshot().queue_length, 0); @@ -219,3 +229,42 @@ async fn test_queue_replay_delivers_and_removes_stored_payload() { .expect("delete queue"); let _ = std::fs::remove_dir_all(args.queue_dir); } + +#[tokio::test] +#[ignore = "requires running RabbitMQ-compatible AMQP broker"] +async fn test_legacy_queued_notification_replays_original_bytes() { + let routing_key = format!("rustfs.legacy.{}", Uuid::new_v4().simple()); + let queue = format!("rustfs-test-{}", Uuid::new_v4().simple()); + let channel = bind_queue(&queue, &routing_key).await; + let queue_dir = std::env::temp_dir().join(format!("rustfs-amqp-legacy-{}", Uuid::new_v4())); + let mut args = test_args(&routing_key); + args.queue_dir = queue_dir.to_string_lossy().to_string(); + let target = AMQPTarget::::new("legacy".to_string(), args).expect("construct AMQP target"); + let event = entity_for("example-bucket", "old-object"); + // Model the pre-upgrade envelope, including whitespace that reserialization would change. + let body = serde_json::to_vec_pretty(&serde_json::json!({ + "EventName": event.event_name, + "Key": "example-bucket/old-object", + "Records": [event.as_ref()] + })) + .unwrap(); + let meta = QueuedPayloadMeta::new( + event.event_name, + event.bucket_name.clone(), + event.object_name.clone(), + "application/json", + body.len(), + ); + let encoded = QueuedPayload::new(meta, body.clone()).encode().unwrap(); + let store = target.store().expect("store configured"); + let key = store.put_raw(&encoded).expect("persist pre-upgrade payload"); + target.send_from_store(key).await.expect("replay legacy payload"); + let (received, _) = read_one_raw(&channel, &queue).await; + assert_eq!(received, body); + assert!(store.list().is_empty()); + channel + .queue_delete(queue.into(), QueueDeleteOptions::default()) + .await + .expect("delete queue"); + std::fs::remove_dir_all(queue_dir).unwrap(); +} diff --git a/docs/operations/amqp-notifications.md b/docs/operations/amqp-notifications.md new file mode 100644 index 000000000..492f51cf1 --- /dev/null +++ b/docs/operations/amqp-notifications.md @@ -0,0 +1,51 @@ +# AMQP notification record layout + +New AMQP bucket notifications place the S3 event directly in `Records`, matching +Kafka and webhook notification targets. For example, a consumer reads the bucket +from `Records[0].s3.bucket.name` and the object key from +`Records[0].s3.object.key`. Previously these fields were nested under +`Records[0].data` in an internal target envelope. + +This follows the record layout described in the +[AWS S3 notification message structure](https://docs.aws.amazon.com/AmazonS3/latest/userguide/notification-content-structure.html). +It changes the record envelope only; the event's existing fields and values are +preserved. The outer `EventName` and `Key` fields are retained. The outer `Key` +decodes the object name once, while `Records[0].s3.object.key` retains the event's +encoded key. AMQP audit log payloads retain their existing envelope. + +## Upgrading consumers + +Consumers can read event fields at the standard S3 record paths. Consumers built +around the previous AMQP envelope must be updated before upgrading producers. During rolling upgrades or queue draining, consumers +may receive both layouts. A temporary compatibility path can unwrap `record.data` +when present and otherwise use `record` directly. + +Existing disk-queued messages are replayed byte for byte. Upgrading does not +rewrite their bodies. Newly queued notifications use the same flat record as +direct publishing. Retain dual-layout handling until old producers are upgraded +and old queued messages have drained. If rolling back a producer, retain that +handling while any flat messages remain queued or in the broker. + +## Verification + +The regression `notification_records_contain_the_event_directly` fails against +the old serializer because `Records[0]` contains the target envelope. It passes +with the notification serializer using the shared event-record builder. Unit +tests also cover encoded keys, preserved metadata in the queue, and the unchanged +audit envelope: + +```bash +cargo test -p rustfs-targets --lib target::amqp::tests +``` + +Broker integration tests exercise direct publishing, disk-queue replay, reconnect, +and byte-for-byte replay of a pre-upgrade envelope. They use synthetic events and +require an isolated RabbitMQ-compatible broker: + +```bash +RUSTFS_TEST_AMQP_URL='amqp://guest:guest@127.0.0.1:5672/%2f' \ + cargo test -p rustfs-targets --test amqp_integration -- --ignored +``` + +These checks cover the target's serialization and delivery boundary. They do not +exercise the full bucket-notification pipeline or every downstream S3 client.