From 8ba3b6ec5977a9bcd7d0d498ee641176ad5a94b8 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Sun, 31 May 2026 02:06:15 +0000 Subject: [PATCH 1/4] fix(serverust-events): corrigir perda silenciosa Kafka e loop DLQ no SQS Lambda MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - EventRouter::with_dlq retorna Ok após publish bem-sucedido no DLQ (alinha com DlqLayer; evita redelivery infinito no Lambda ESM) - LambdaBroker e KafkaBroker::dispatch erram quando não há handler inscrito, em vez de no-op que commitaria offset/ack implicitamente - Testes de regressão para os três caminhos Co-authored-by: Jaime Basso --- CHANGELOG.md | 2 + serverust-events/src/broker/kafka.rs | 7 +++ serverust-events/src/broker/lambda.rs | 8 ++- serverust-events/src/router.rs | 15 +++-- .../tests/kafka_broker_dispatch.rs | 8 ++- serverust-events/tests/lambda_broker.rs | 9 ++- serverust-events/tests/retry_policy.rs | 6 +- serverust-events/tests/sqs_consumer.rs | 61 +++++++++++++++++++ 8 files changed, 102 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3efd065..b559ca2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `EventRouter` com `RetryPolicy::Exponential`: o atraso `base_delay * 2^n` passa a usar `Duration::saturating_mul` e expoente limitado a 31, evitando panic por overflow de `Duration` em retentativas longas ou `base_delay` grande. - `SqsBroker::handle_sqs_event` (Lambda ESM + `ReportBatchItemFailures`): mensagens sem handler para a fila do ARN ou sem `event_source_arn` válido passam a entrar em `batchItemFailures` quando há `messageId`, em vez de serem tratadas como sucesso implícito (a Lambda removia da fila sem processamento). +- `EventRouter::with_dlq`: após publicação bem-sucedida no tópico DLQ, o wrapper retorna `Ok(())` (mesma semântica de `DlqLayer`), permitindo ack da mensagem original no Lambda SQS em vez de loop infinito de redelivery. +- `LambdaBroker::handle_kafka_event` e `KafkaBroker::dispatch`: tópico sem handler inscrito retorna erro em vez de no-op silencioso, evitando commit de offset / ack implícito sem processamento. ## [0.3.0] - 2026-05-17 diff --git a/serverust-events/src/broker/kafka.rs b/serverust-events/src/broker/kafka.rs index 36f4f5d..c1c1a56 100644 --- a/serverust-events/src/broker/kafka.rs +++ b/serverust-events/src/broker/kafka.rs @@ -171,6 +171,13 @@ impl KafkaBroker { .map(|s| s.handler.clone()) .collect(); + if handlers.is_empty() { + return Err(BrokerError::Subscribe(format!( + "no handler subscribed for kafka topic '{}'", + msg.topic + ))); + } + for handler in handlers { handler(msg.clone()).await?; } diff --git a/serverust-events/src/broker/lambda.rs b/serverust-events/src/broker/lambda.rs index a140f57..40bf67e 100644 --- a/serverust-events/src/broker/lambda.rs +++ b/serverust-events/src/broker/lambda.rs @@ -68,7 +68,9 @@ impl LambdaBroker { /// 1. Identifica o tópico via `record.topic` (campo do próprio registro). /// 2. Se houver handlers inscritos, decodifica `value` (Base64) e /// despacha como [`BrokerMessage`]. - /// 3. Se não houver handlers inscritos para o tópico, ignora o registro. + /// 3. Se não houver handlers inscritos para o tópico, retorna erro para que a + /// invocação Lambda falhe e os offsets não sejam commitados como sucesso + /// (evita perda silenciosa de mensagens). /// /// O primeiro erro encontrado interrompe o despacho e propaga. pub async fn handle_kafka_event(&self, event: &KafkaEvent) -> Result<(), BrokerError> { @@ -86,7 +88,9 @@ impl LambdaBroker { .collect(); if handlers.is_empty() { - continue; + return Err(BrokerError::Subscribe(format!( + "no handler subscribed for kafka topic '{topic}'" + ))); } let value_b64 = raw.value.as_deref().ok_or_else(|| { diff --git a/serverust-events/src/router.rs b/serverust-events/src/router.rs index 656fef0..05a0d7e 100644 --- a/serverust-events/src/router.rs +++ b/serverust-events/src/router.rs @@ -108,10 +108,17 @@ where // Todas as tentativas falharam — publicar no DLQ se configurado. if let Some(dlq_topic) = &dlq { - if let Err(dlq_err) = broker.publish(dlq_topic, &msg.payload).await { - // DLQ write failure: log to stderr so operators can detect it. - // The original handler error is returned regardless. - eprintln!("[serverust-events] DLQ publish to '{dlq_topic}' failed: {dlq_err}"); + match broker.publish(dlq_topic, &msg.payload).await { + Ok(()) => { + // Alinha com `DlqLayer`: DLQ aceito => ack da mensagem original + // (ex.: Lambda SQS sem entrada em batchItemFailures). + return Ok(()); + } + Err(dlq_err) => { + eprintln!( + "[serverust-events] DLQ publish to '{dlq_topic}' failed: {dlq_err}" + ); + } } } Err(last_err.unwrap_or_else(|| BrokerError::Subscribe("sem tentativas".to_string()))) diff --git a/serverust-events/tests/kafka_broker_dispatch.rs b/serverust-events/tests/kafka_broker_dispatch.rs index 7086cbb..2dfec61 100644 --- a/serverust-events/tests/kafka_broker_dispatch.rs +++ b/serverust-events/tests/kafka_broker_dispatch.rs @@ -59,7 +59,7 @@ async fn dispatch_invoca_handlers_inscritos_no_topico() { } #[tokio::test] -async fn dispatch_em_topico_sem_subscriber_e_no_op() { +async fn dispatch_erro_quando_topico_sem_subscriber() { let broker = make_broker(); let msg = BrokerMessage { topic: "topico.sem.subscriber".to_string(), @@ -70,7 +70,11 @@ async fn dispatch_em_topico_sem_subscriber_e_no_op() { headers: HashMap::new(), timestamp: None, }; - broker.dispatch(msg).await.unwrap(); + let err = broker.dispatch(msg).await.unwrap_err(); + assert!( + format!("{err}").contains("no handler subscribed"), + "erro foi: {err}" + ); } #[tokio::test] diff --git a/serverust-events/tests/lambda_broker.rs b/serverust-events/tests/lambda_broker.rs index 635d5d0..7579cd0 100644 --- a/serverust-events/tests/lambda_broker.rs +++ b/serverust-events/tests/lambda_broker.rs @@ -63,11 +63,14 @@ async fn handle_kafka_event_despacha_registros_para_handlers_inscritos() { } #[tokio::test] -async fn handle_kafka_event_ignora_topico_sem_subscriber() { +async fn handle_kafka_event_erro_quando_topico_sem_subscriber() { let broker = Arc::new(LambdaBroker::new()); - // Nenhum subscriber registrado — não deve panicar nem erroar. let event = fixture(); - broker.handle_kafka_event(&event).await.unwrap(); + let err = broker.handle_kafka_event(&event).await.unwrap_err(); + assert!( + format!("{err}").contains("no handler subscribed"), + "erro foi: {err}" + ); } #[tokio::test] diff --git a/serverust-events/tests/retry_policy.rs b/serverust-events/tests/retry_policy.rs index dff5e05..9ceebeb 100644 --- a/serverust-events/tests/retry_policy.rs +++ b/serverust-events/tests/retry_policy.rs @@ -133,7 +133,7 @@ async fn dead_letter_publica_no_dlq_apos_esgotamento_via_policy() { .await .unwrap(); - let _ = broker.publish("orders", &payload).await; + broker.publish("orders", &payload).await.unwrap(); let dlq_msgs = broker.messages("orders.dlq"); assert_eq!(dlq_msgs.len(), 1); @@ -159,7 +159,7 @@ async fn with_dlq_publica_no_dlq_apos_esgotamento_via_router() { .await .unwrap(); - let _ = broker.publish("orders", &payload).await; + broker.publish("orders", &payload).await.unwrap(); let dlq_msgs = broker.messages("orders.dlq"); assert_eq!(dlq_msgs.len(), 1); @@ -218,7 +218,7 @@ async fn exponential_dead_letter_publica_no_dlq() { .await .unwrap(); - let _ = broker.publish("orders", &payload).await; + broker.publish("orders", &payload).await.unwrap(); let dlq_msgs = broker.messages("orders.dlq"); assert_eq!(dlq_msgs.len(), 1); diff --git a/serverust-events/tests/sqs_consumer.rs b/serverust-events/tests/sqs_consumer.rs index b60f4ff..16dc20d 100644 --- a/serverust-events/tests/sqs_consumer.rs +++ b/serverust-events/tests/sqs_consumer.rs @@ -376,3 +376,64 @@ async fn handler_pode_extrair_state_compartilhado() { ] ); } + +/// Broker composto: subscribe no `SqsBroker`, captura publishes de DLQ. +struct DlqRoutingBroker { + sqs: Arc, + dlq_topic: String, + dlq_payloads: Arc>>>, +} + +#[async_trait::async_trait] +impl Broker for DlqRoutingBroker { + async fn subscribe( + &self, + topic: &str, + handler: serverust_events::broker::BoxedHandler, + ) -> Result<(), BrokerError> { + self.sqs.subscribe(topic, handler).await + } + + async fn publish(&self, topic: &str, payload: &[u8]) -> Result<(), BrokerError> { + if topic == self.dlq_topic { + self.dlq_payloads.lock().unwrap().push(payload.to_vec()); + Ok(()) + } else { + Err(BrokerError::Publish(format!( + "DlqRoutingBroker only publishes to {}", + self.dlq_topic + ))) + } + } +} + +#[tokio::test] +async fn event_router_dlq_ack_apos_publicar_dlq_em_lambda_sqs() { + use serverust_events::retry::RetryPolicy; + + let sqs_broker = Arc::new(SqsBroker::new()); + let dlq_payloads = Arc::new(Mutex::new(Vec::new())); + let broker = Arc::new(DlqRoutingBroker { + sqs: sqs_broker.clone(), + dlq_topic: "orders-dlq".to_string(), + dlq_payloads: dlq_payloads.clone(), + }); + + EventRouter::new() + .subscribe::("orders", |_: Order| async move { + Err(BrokerError::Subscribe("falha".to_string())) + }) + .with_retry(RetryPolicy::immediate(1)) + .with_dlq("orders-dlq") + .attach(broker) + .await + .unwrap(); + + let resp = sqs_broker.handle_sqs_event(&fixture()).await; + assert!( + resp.batch_item_failures.is_empty(), + "DLQ aceito => ack na Lambda, got: {:?}", + resp.batch_item_failures + ); + assert_eq!(dlq_payloads.lock().unwrap().len(), 3); +} From 9f2314a271b9e8c81a70fbd2ac4110112b518c97 Mon Sep 17 00:00:00 2001 From: jaimejunr Date: Sun, 13 Sep 2026 00:31:20 +0000 Subject: [PATCH 2/4] feat(serverust-events): add UnhandledTopicPolicy for Kafka/Lambda brokers Default WarnAndIgnore preserves 0.3.x skip (now with tracing::warn!). Error returns BrokerError::Subscribe with received and registered topics. EventRouter logs handler+DLQ errors via tracing::error! when DLQ publish fails, instead of eprintln! of the DLQ error only. Co-Authored-By: Claude Opus 5 --- CHANGELOG.md | 7 +- serverust-events/Cargo.toml | 7 +- serverust-events/src/broker/kafka.rs | 46 ++++++++--- serverust-events/src/broker/lambda.rs | 47 +++++++---- serverust-events/src/broker/mod.rs | 82 +++++++++++++++++++ serverust-events/src/router.rs | 11 ++- .../tests/kafka_broker_dispatch.rs | 32 +++++++- serverust-events/tests/lambda_broker.rs | 26 +++++- serverust-events/tests/retry_policy.rs | 57 +++++++++++++ serverust-events/tests/sqs_consumer.rs | 52 ++++++++++++ 10 files changed, 325 insertions(+), 42 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b559ca2..73c0e51 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,12 +10,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- `UnhandledTopicPolicy` em `LambdaBroker` e `KafkaBroker` (`with_unhandled_topic_policy`): `WarnAndIgnore` (default) preserva o comportamento 0.3.x de pular o record sem handler e passa a emitir `tracing::warn!` com o tópico; `Error` retorna `BrokerError::Subscribe` com o tópico recebido e a lista de tópicos inscritos. + ### Fixed - `EventRouter` com `RetryPolicy::Exponential`: o atraso `base_delay * 2^n` passa a usar `Duration::saturating_mul` e expoente limitado a 31, evitando panic por overflow de `Duration` em retentativas longas ou `base_delay` grande. - `SqsBroker::handle_sqs_event` (Lambda ESM + `ReportBatchItemFailures`): mensagens sem handler para a fila do ARN ou sem `event_source_arn` válido passam a entrar em `batchItemFailures` quando há `messageId`, em vez de serem tratadas como sucesso implícito (a Lambda removia da fila sem processamento). -- `EventRouter::with_dlq`: após publicação bem-sucedida no tópico DLQ, o wrapper retorna `Ok(())` (mesma semântica de `DlqLayer`), permitindo ack da mensagem original no Lambda SQS em vez de loop infinito de redelivery. -- `LambdaBroker::handle_kafka_event` e `KafkaBroker::dispatch`: tópico sem handler inscrito retorna erro em vez de no-op silencioso, evitando commit de offset / ack implícito sem processamento. +- `EventRouter::with_dlq`: após publicação bem-sucedida no tópico DLQ, o wrapper retorna `Ok(())` (mesma semântica de `DlqLayer`), permitindo ack da mensagem original no Lambda SQS em vez de loop infinito de redelivery. Se o publish na DLQ falhar, o erro original do handler é retornado e ambos os erros (handler e DLQ) são registrados com `tracing::error!`. ## [0.3.0] - 2026-05-17 diff --git a/serverust-events/Cargo.toml b/serverust-events/Cargo.toml index ea2a649..a0b3359 100644 --- a/serverust-events/Cargo.toml +++ b/serverust-events/Cargo.toml @@ -32,7 +32,6 @@ sqs = [ "dep:tower", "dep:pin-project-lite", "dep:serverust-telemetry", - "dep:tracing", ] # Habilita o KafkaBroker (rdkafka + librdkafka C). Off por default — apps que # só consomem Kafka via Lambda trigger não pagam o custo de librdkafka. @@ -67,9 +66,9 @@ pin-project-lite = { version = "0.2", optional = true } # Telemetria reusada pelos Tower layers (feature `sqs`). Default features mantém # o build leve (sem otel/dynamodb). serverust-telemetry = { path = "../serverust-telemetry", version = "0.3.0", optional = true } -# Tracing usado no `ObservabilityLayer` (feature `sqs`) para abrir spans com -# trace_id no formato AWS X-Ray. -tracing = { version = "0.1", optional = true } +# Tracing usado em brokers (`UnhandledTopicPolicy`), `EventRouter` (DLQ) e +# no `ObservabilityLayer` (feature `sqs`) para spans X-Ray. +tracing = "0.1" # Opt-in via feature `kafka` (e o alias `kafka-producer`). # `ssl` é necessário para SASL_SSL e habilita OAUTHBEARER nativamente em diff --git a/serverust-events/src/broker/kafka.rs b/serverust-events/src/broker/kafka.rs index c1c1a56..ce1aeb0 100644 --- a/serverust-events/src/broker/kafka.rs +++ b/serverust-events/src/broker/kafka.rs @@ -18,7 +18,7 @@ use rdkafka::config::ClientConfig; use rdkafka::consumer::ConsumerContext; use rdkafka::producer::{FutureProducer, FutureRecord, ProducerContext}; -use super::{BoxedHandler, Broker, BrokerError, BrokerMessage}; +use super::{BoxedHandler, Broker, BrokerError, BrokerMessage, UnhandledTopicPolicy}; /// Contexto rdkafka que fornece o token IAM MSK via OAUTHBEARER. #[derive(Clone)] @@ -101,6 +101,7 @@ impl KafkaBrokerConfig { pub struct KafkaBroker { producer: FutureProducer, subscriptions: Mutex>, + unhandled_topic_policy: UnhandledTopicPolicy, } /// Registro interno de uma inscrição (handler + tópico). @@ -141,9 +142,16 @@ impl KafkaBroker { Ok(Self { producer, subscriptions: Mutex::new(Vec::new()), + unhandled_topic_policy: UnhandledTopicPolicy::default(), }) } + /// Define a política para records cujo tópico não tem handler inscrito. + pub fn with_unhandled_topic_policy(mut self, policy: UnhandledTopicPolicy) -> Self { + self.unhandled_topic_policy = policy; + self + } + /// Lista os tópicos atualmente inscritos (somente leitura, ordem de inscrição). pub fn subscribed_topics(&self) -> Vec { self.subscriptions @@ -161,21 +169,33 @@ impl KafkaBroker { /// [`BrokerMessage`] e chama `dispatch` para entregar aos handlers. /// Testar o loop real exige broker físico; `dispatch` é testado de /// forma isolada. + /// + /// Sem handler inscrito, aplica [`UnhandledTopicPolicy`]: `WarnAndIgnore` + /// (default) pula a mensagem e emite `tracing::warn!`; `Error` retorna + /// [`BrokerError::Subscribe`] com o tópico recebido e a lista de tópicos + /// inscritos. pub async fn dispatch(&self, msg: BrokerMessage) -> Result<(), BrokerError> { - let handlers: Vec = self - .subscriptions - .lock() - .map_err(|_| BrokerError::Subscribe("subscriptions mutex poisoned".into()))? - .iter() - .filter(|s| s.topic == msg.topic) - .map(|s| s.handler.clone()) - .collect(); + let (handlers, registered_topics): (Vec, Vec) = { + let guard = self.subscriptions.lock().map_err(|_| { + BrokerError::Subscribe("subscriptions mutex poisoned".into()) + })?; + let handlers: Vec = guard + .iter() + .filter(|s| s.topic == msg.topic) + .map(|s| s.handler.clone()) + .collect(); + let registered_topics = if handlers.is_empty() { + super::unique_topics(guard.iter().map(|s| s.topic.as_str())) + } else { + Vec::new() + }; + (handlers, registered_topics) + }; if handlers.is_empty() { - return Err(BrokerError::Subscribe(format!( - "no handler subscribed for kafka topic '{}'", - msg.topic - ))); + return self + .unhandled_topic_policy + .on_unhandled(&msg.topic, ®istered_topics); } for handler in handlers { diff --git a/serverust-events/src/broker/lambda.rs b/serverust-events/src/broker/lambda.rs index 40bf67e..bc6e641 100644 --- a/serverust-events/src/broker/lambda.rs +++ b/serverust-events/src/broker/lambda.rs @@ -24,7 +24,7 @@ use async_trait::async_trait; use aws_lambda_events::event::kafka::KafkaEvent; use base64::Engine; -use super::{BoxedHandler, Broker, BrokerError, BrokerMessage}; +use super::{BoxedHandler, Broker, BrokerError, BrokerMessage, UnhandledTopicPolicy}; /// Broker sink-only para o modo Lambda. /// @@ -37,6 +37,7 @@ use super::{BoxedHandler, Broker, BrokerError, BrokerMessage}; /// ou ao producer dedicado em uma futura US. pub struct LambdaBroker { subscriptions: Mutex>, + unhandled_topic_policy: UnhandledTopicPolicy, } struct Subscription { @@ -49,9 +50,16 @@ impl LambdaBroker { pub fn new() -> Self { Self { subscriptions: Mutex::new(Vec::new()), + unhandled_topic_policy: UnhandledTopicPolicy::default(), } } + /// Define a política para records cujo tópico não tem handler inscrito. + pub fn with_unhandled_topic_policy(mut self, policy: UnhandledTopicPolicy) -> Self { + self.unhandled_topic_policy = policy; + self + } + /// Lista os tópicos atualmente inscritos (ordem de inscrição). pub fn subscribed_topics(&self) -> Vec { self.subscriptions @@ -68,9 +76,10 @@ impl LambdaBroker { /// 1. Identifica o tópico via `record.topic` (campo do próprio registro). /// 2. Se houver handlers inscritos, decodifica `value` (Base64) e /// despacha como [`BrokerMessage`]. - /// 3. Se não houver handlers inscritos para o tópico, retorna erro para que a - /// invocação Lambda falhe e os offsets não sejam commitados como sucesso - /// (evita perda silenciosa de mensagens). + /// 3. Se não houver handlers inscritos para o tópico, aplica + /// [`UnhandledTopicPolicy`]: `WarnAndIgnore` (default) pula o registro + /// e emite `tracing::warn!`; `Error` retorna [`BrokerError::Subscribe`] + /// com o tópico recebido e a lista de tópicos inscritos. /// /// O primeiro erro encontrado interrompe o despacho e propaga. pub async fn handle_kafka_event(&self, event: &KafkaEvent) -> Result<(), BrokerError> { @@ -78,19 +87,27 @@ impl LambdaBroker { for raw in records { let topic = raw.topic.clone().unwrap_or_default(); - let handlers: Vec = self - .subscriptions - .lock() - .map_err(|_| BrokerError::Subscribe("subscriptions mutex poisoned".into()))? - .iter() - .filter(|s| s.topic == topic) - .map(|s| s.handler.clone()) - .collect(); + let (handlers, registered_topics): (Vec, Vec) = { + let guard = self.subscriptions.lock().map_err(|_| { + BrokerError::Subscribe("subscriptions mutex poisoned".into()) + })?; + let handlers: Vec = guard + .iter() + .filter(|s| s.topic == topic) + .map(|s| s.handler.clone()) + .collect(); + let registered_topics = if handlers.is_empty() { + super::unique_topics(guard.iter().map(|s| s.topic.as_str())) + } else { + Vec::new() + }; + (handlers, registered_topics) + }; if handlers.is_empty() { - return Err(BrokerError::Subscribe(format!( - "no handler subscribed for kafka topic '{topic}'" - ))); + self.unhandled_topic_policy + .on_unhandled(&topic, ®istered_topics)?; + continue; } let value_b64 = raw.value.as_deref().ok_or_else(|| { diff --git a/serverust-events/src/broker/mod.rs b/serverust-events/src/broker/mod.rs index b9ff0d3..217049f 100644 --- a/serverust-events/src/broker/mod.rs +++ b/serverust-events/src/broker/mod.rs @@ -49,6 +49,57 @@ pub enum BrokerError { Transport(String), } +/// Política para records cujo tópico não tem handler inscrito. +/// +/// O default [`UnhandledTopicPolicy::WarnAndIgnore`] preserva o comportamento +/// 0.3.x: o record é pulado. A variante [`UnhandledTopicPolicy::Error`] falha +/// o despacho com [`BrokerError::Subscribe`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +pub enum UnhandledTopicPolicy { + /// Pula o record e emite `tracing::warn!` com o tópico recebido. + /// + /// Equivale ao skip 0.3.x; o warn é o único acréscimo no default (o skip + /// era totalmente silencioso). + #[default] + WarnAndIgnore, + /// Retorna [`BrokerError::Subscribe`] com o tópico recebido e a lista de + /// tópicos que têm handler registrado. + Error, +} + +impl UnhandledTopicPolicy { + /// Aplica a política a um tópico sem handler. + pub(crate) fn on_unhandled( + self, + topic: &str, + registered_topics: &[String], + ) -> Result<(), BrokerError> { + match self { + Self::WarnAndIgnore => { + tracing::warn!(topic, "no handler subscribed; skipping record"); + Ok(()) + } + Self::Error => Err(BrokerError::Subscribe(format!( + "no handler subscribed for kafka topic '{topic}'; registered topics: [{}]", + registered_topics.join(", ") + ))), + } + } +} + +pub(crate) fn unique_topics<'a, I>(topics: I) -> Vec +where + I: IntoIterator, +{ + let mut out: Vec = Vec::new(); + for topic in topics { + if !out.iter().any(|existing| existing == topic) { + out.push(topic.to_string()); + } + } + out +} + /// Mensagem entregue a um handler inscrito. /// /// Representa um registro de broker já normalizado para um shape comum @@ -109,3 +160,34 @@ impl Broker for Arc { (**self).publish(topic, payload).await } } + +#[cfg(test)] +mod unhandled_topic_policy_tests { + use super::*; + + #[test] + fn default_is_warn_and_ignore() { + assert_eq!( + UnhandledTopicPolicy::default(), + UnhandledTopicPolicy::WarnAndIgnore + ); + } + + #[test] + fn error_policy_includes_received_and_registered_topics() { + let err = UnhandledTopicPolicy::Error + .on_unhandled("incoming", &["orders".to_string(), "billing".to_string()]) + .expect_err("Error policy deve falhar"); + let msg = format!("{err}"); + assert!(msg.contains("incoming"), "{msg}"); + assert!(msg.contains("orders"), "{msg}"); + assert!(msg.contains("billing"), "{msg}"); + } + + #[test] + fn warn_and_ignore_returns_ok() { + UnhandledTopicPolicy::WarnAndIgnore + .on_unhandled("incoming", &[]) + .expect("WarnAndIgnore deve pular"); + } +} diff --git a/serverust-events/src/router.rs b/serverust-events/src/router.rs index 05a0d7e..faa612f 100644 --- a/serverust-events/src/router.rs +++ b/serverust-events/src/router.rs @@ -115,8 +115,15 @@ where return Ok(()); } Err(dlq_err) => { - eprintln!( - "[serverust-events] DLQ publish to '{dlq_topic}' failed: {dlq_err}" + let handler_error = last_err + .as_ref() + .map(ToString::to_string) + .unwrap_or_else(|| "sem tentativas".to_string()); + tracing::error!( + dlq_topic = %dlq_topic, + handler_error = %handler_error, + dlq_error = %dlq_err, + "DLQ publish failed after handler retries exhausted" ); } } diff --git a/serverust-events/tests/kafka_broker_dispatch.rs b/serverust-events/tests/kafka_broker_dispatch.rs index 2dfec61..f90e896 100644 --- a/serverust-events/tests/kafka_broker_dispatch.rs +++ b/serverust-events/tests/kafka_broker_dispatch.rs @@ -12,6 +12,7 @@ use std::sync::Mutex; use serverust_events::broker::Broker; use serverust_events::broker::BrokerMessage; +use serverust_events::broker::UnhandledTopicPolicy; use serverust_events::broker::kafka::{KafkaBroker, KafkaBrokerConfig}; fn make_broker() -> KafkaBroker { @@ -59,8 +60,28 @@ async fn dispatch_invoca_handlers_inscritos_no_topico() { } #[tokio::test] -async fn dispatch_erro_quando_topico_sem_subscriber() { +async fn dispatch_em_topico_sem_subscriber_e_no_op() { let broker = make_broker(); + let msg = BrokerMessage { + topic: "topico.sem.subscriber".to_string(), + partition: None, + offset: None, + key: None, + payload: Vec::new(), + headers: HashMap::new(), + timestamp: None, + }; + broker.dispatch(msg).await.unwrap(); +} + +#[tokio::test] +async fn dispatch_erro_quando_topico_sem_subscriber() { + let broker = make_broker().with_unhandled_topic_policy(UnhandledTopicPolicy::Error); + let h = |_: BrokerMessage| -> serverust_events::broker::HandlerFuture { + Box::pin(async { Ok(()) }) + }; + broker.subscribe("orders.created", Arc::new(h)).await.unwrap(); + let msg = BrokerMessage { topic: "topico.sem.subscriber".to_string(), partition: None, @@ -71,9 +92,14 @@ async fn dispatch_erro_quando_topico_sem_subscriber() { timestamp: None, }; let err = broker.dispatch(msg).await.unwrap_err(); + let text = format!("{err}"); + assert!( + text.contains("topico.sem.subscriber"), + "mensagem deve citar o tópico recebido; erro foi: {text}" + ); assert!( - format!("{err}").contains("no handler subscribed"), - "erro foi: {err}" + text.contains("orders.created"), + "mensagem deve citar os tópicos inscritos; erro foi: {text}" ); } diff --git a/serverust-events/tests/lambda_broker.rs b/serverust-events/tests/lambda_broker.rs index 7579cd0..82fa9d7 100644 --- a/serverust-events/tests/lambda_broker.rs +++ b/serverust-events/tests/lambda_broker.rs @@ -11,6 +11,7 @@ use std::sync::Mutex; use aws_lambda_events::event::kafka::KafkaEvent; use serde::Deserialize; use serverust_events::broker::Broker; +use serverust_events::broker::UnhandledTopicPolicy; use serverust_events::broker::lambda::LambdaBroker; use serverust_events::router::EventRouter; @@ -63,13 +64,32 @@ async fn handle_kafka_event_despacha_registros_para_handlers_inscritos() { } #[tokio::test] -async fn handle_kafka_event_erro_quando_topico_sem_subscriber() { +async fn handle_kafka_event_ignora_topico_sem_subscriber() { let broker = Arc::new(LambdaBroker::new()); + // Nenhum subscriber registrado — default WarnAndIgnore não panica nem erra. + let event = fixture(); + broker.handle_kafka_event(&event).await.unwrap(); +} + +#[tokio::test] +async fn handle_kafka_event_erro_quando_topico_sem_subscriber() { + let broker = Arc::new( + LambdaBroker::new().with_unhandled_topic_policy(UnhandledTopicPolicy::Error), + ); + let router = + EventRouter::new().subscribe::("other.topic", |_| async { Ok(()) }); + router.attach(broker.clone()).await.unwrap(); + let event = fixture(); let err = broker.handle_kafka_event(&event).await.unwrap_err(); + let msg = format!("{err}"); + assert!( + msg.contains("wallet.credits"), + "mensagem deve citar o tópico recebido; erro foi: {msg}" + ); assert!( - format!("{err}").contains("no handler subscribed"), - "erro foi: {err}" + msg.contains("other.topic"), + "mensagem deve citar os tópicos inscritos; erro foi: {msg}" ); } diff --git a/serverust-events/tests/retry_policy.rs b/serverust-events/tests/retry_policy.rs index 9ceebeb..31c3e25 100644 --- a/serverust-events/tests/retry_policy.rs +++ b/serverust-events/tests/retry_policy.rs @@ -223,3 +223,60 @@ async fn exponential_dead_letter_publica_no_dlq() { let dlq_msgs = broker.messages("orders.dlq"); assert_eq!(dlq_msgs.len(), 1); } + +// --------------------------------------------------------------------------- +// DLQ publish falha: o erro original do handler é retornado. +// --------------------------------------------------------------------------- + +struct FailingDlqBroker { + inner: Arc, + fail_topic: String, +} + +#[async_trait::async_trait] +impl Broker for FailingDlqBroker { + async fn subscribe( + &self, + topic: &str, + handler: serverust_events::broker::BoxedHandler, + ) -> Result<(), BrokerError> { + self.inner.subscribe(topic, handler).await + } + + async fn publish(&self, topic: &str, payload: &[u8]) -> Result<(), BrokerError> { + if topic == self.fail_topic { + return Err(BrokerError::Publish(format!( + "simulated DLQ publish failure to '{topic}'" + ))); + } + self.inner.publish(topic, payload).await + } +} + +#[tokio::test] +async fn dlq_publish_falha_propaga_erro_do_handler() { + let inner = Arc::new(InMemoryBroker::new()); + let broker = Arc::new(FailingDlqBroker { + inner: inner.clone(), + fail_topic: "orders.dlq".to_string(), + }); + let payload = serde_json::to_vec(&OrderEvent { id: 8 }).unwrap(); + + EventRouter::new() + .subscribe::("orders", |_: OrderEvent| async move { + Err(BrokerError::Subscribe("falha do handler".to_string())) + }) + .with_retry(RetryPolicy::immediate(1)) + .with_dlq("orders.dlq") + .attach(broker.clone()) + .await + .unwrap(); + + let result = broker.publish("orders", &payload).await; + let err = result.expect_err("DLQ falhou => Err do handler"); + assert!( + format!("{err}").contains("falha do handler"), + "erro foi: {err}" + ); + assert_eq!(inner.messages("orders.dlq").len(), 0); +} diff --git a/serverust-events/tests/sqs_consumer.rs b/serverust-events/tests/sqs_consumer.rs index 16dc20d..e1ff6f0 100644 --- a/serverust-events/tests/sqs_consumer.rs +++ b/serverust-events/tests/sqs_consumer.rs @@ -437,3 +437,55 @@ async fn event_router_dlq_ack_apos_publicar_dlq_em_lambda_sqs() { ); assert_eq!(dlq_payloads.lock().unwrap().len(), 3); } + +struct FailingDlqBroker { + sqs: Arc, + dlq_topic: String, +} + +#[async_trait::async_trait] +impl Broker for FailingDlqBroker { + async fn subscribe( + &self, + topic: &str, + handler: serverust_events::broker::BoxedHandler, + ) -> Result<(), BrokerError> { + self.sqs.subscribe(topic, handler).await + } + + async fn publish(&self, topic: &str, _payload: &[u8]) -> Result<(), BrokerError> { + Err(BrokerError::Publish(format!( + "simulated DLQ publish failure to '{topic}' (configured dlq: {})", + self.dlq_topic + ))) + } +} + +#[tokio::test] +async fn event_router_dlq_publish_falha_retorna_err_em_lambda_sqs() { + use serverust_events::retry::RetryPolicy; + + let sqs_broker = Arc::new(SqsBroker::new()); + let broker = Arc::new(FailingDlqBroker { + sqs: sqs_broker.clone(), + dlq_topic: "orders-dlq".to_string(), + }); + + EventRouter::new() + .subscribe::("orders", |_: Order| async move { + Err(BrokerError::Subscribe("falha do handler".to_string())) + }) + .with_retry(RetryPolicy::immediate(1)) + .with_dlq("orders-dlq") + .attach(broker) + .await + .unwrap(); + + let resp = sqs_broker.handle_sqs_event(&fixture()).await; + assert_eq!( + resp.batch_item_failures.len(), + 3, + "DLQ falhou => Err do handler => batchItemFailures, got: {:?}", + resp.batch_item_failures + ); +} From b830e3299da876393e60e2b7bdcd8492edfd59b5 Mon Sep 17 00:00:00 2001 From: jaimejunr Date: Sat, 12 Sep 2026 21:35:27 -0300 Subject: [PATCH 3/4] style: cargo fmt em serverust-events Co-Authored-By: Claude Opus 5 --- serverust-events/src/broker/kafka.rs | 7 ++++--- serverust-events/tests/kafka_broker_dispatch.rs | 5 ++++- serverust-events/tests/lambda_broker.rs | 5 ++--- 3 files changed, 10 insertions(+), 7 deletions(-) diff --git a/serverust-events/src/broker/kafka.rs b/serverust-events/src/broker/kafka.rs index ce1aeb0..4993226 100644 --- a/serverust-events/src/broker/kafka.rs +++ b/serverust-events/src/broker/kafka.rs @@ -176,9 +176,10 @@ impl KafkaBroker { /// inscritos. pub async fn dispatch(&self, msg: BrokerMessage) -> Result<(), BrokerError> { let (handlers, registered_topics): (Vec, Vec) = { - let guard = self.subscriptions.lock().map_err(|_| { - BrokerError::Subscribe("subscriptions mutex poisoned".into()) - })?; + let guard = self + .subscriptions + .lock() + .map_err(|_| BrokerError::Subscribe("subscriptions mutex poisoned".into()))?; let handlers: Vec = guard .iter() .filter(|s| s.topic == msg.topic) diff --git a/serverust-events/tests/kafka_broker_dispatch.rs b/serverust-events/tests/kafka_broker_dispatch.rs index f90e896..26e24ec 100644 --- a/serverust-events/tests/kafka_broker_dispatch.rs +++ b/serverust-events/tests/kafka_broker_dispatch.rs @@ -80,7 +80,10 @@ async fn dispatch_erro_quando_topico_sem_subscriber() { let h = |_: BrokerMessage| -> serverust_events::broker::HandlerFuture { Box::pin(async { Ok(()) }) }; - broker.subscribe("orders.created", Arc::new(h)).await.unwrap(); + broker + .subscribe("orders.created", Arc::new(h)) + .await + .unwrap(); let msg = BrokerMessage { topic: "topico.sem.subscriber".to_string(), diff --git a/serverust-events/tests/lambda_broker.rs b/serverust-events/tests/lambda_broker.rs index 82fa9d7..fd354f9 100644 --- a/serverust-events/tests/lambda_broker.rs +++ b/serverust-events/tests/lambda_broker.rs @@ -73,9 +73,8 @@ async fn handle_kafka_event_ignora_topico_sem_subscriber() { #[tokio::test] async fn handle_kafka_event_erro_quando_topico_sem_subscriber() { - let broker = Arc::new( - LambdaBroker::new().with_unhandled_topic_policy(UnhandledTopicPolicy::Error), - ); + let broker = + Arc::new(LambdaBroker::new().with_unhandled_topic_policy(UnhandledTopicPolicy::Error)); let router = EventRouter::new().subscribe::("other.topic", |_| async { Ok(()) }); router.attach(broker.clone()).await.unwrap(); From d3e89791400caf87ba8960f89d8ec1f78ef4df9e Mon Sep 17 00:00:00 2001 From: jaimejunr Date: Sat, 12 Sep 2026 21:35:48 -0300 Subject: [PATCH 4/4] =?UTF-8?q?docs(changelog):=20registrar=20tracing=20co?= =?UTF-8?q?mo=20dep=20n=C3=A3o-opcional=20em=20serverust-events?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 5 --- CHANGELOG.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 73c0e51..d2c7bf8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `UnhandledTopicPolicy` em `LambdaBroker` e `KafkaBroker` (`with_unhandled_topic_policy`): `WarnAndIgnore` (default) preserva o comportamento 0.3.x de pular o record sem handler e passa a emitir `tracing::warn!` com o tópico; `Error` retorna `BrokerError::Subscribe` com o tópico recebido e a lista de tópicos inscritos. +### Changed + +- `serverust-events`: `tracing` deixa de ser dependência opcional (antes só sob a feature `sqs`) — `UnhandledTopicPolicy` e o log de falha da DLQ em `EventRouter` rodam em código sem a feature `sqs`. + ### Fixed - `EventRouter` com `RetryPolicy::Exponential`: o atraso `base_delay * 2^n` passa a usar `Duration::saturating_mul` e expoente limitado a 31, evitando panic por overflow de `Duration` em retentativas longas ou `base_delay` grande.