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
This commit is contained in:
Dae-Cheol Noh
2026-09-22 21:16:40 +09:00
committed by GitHub
parent 1f04a12abf
commit 3c6c88b2e7
3 changed files with 194 additions and 9 deletions
+88 -3
View File
@@ -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<E>) -> Result<QueuedPayload, TargetError> {
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<Arc<AMQPConnection>, 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<serde_json::Value> {
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()))
}
+55 -6
View File
@@ -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<EntityTarget<serde_json::Value>
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<u8>, 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::<Value>::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();
}
+51
View File
@@ -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.