Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,10 +27,14 @@ 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.

### 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`.
- CI e desenvolvimento local passam a usar toolchain Rust pinada em `rust-toolchain.toml` (1.94.1) em vez de `stable` flutuante — os testes `trybuild` de `serverust-macros` comparam a saída literal do rustc e quebravam a cada mudança de formatação de diagnóstico.

- `serverust-cli`: passa a usar `version.workspace = true` em `Cargo.toml`, herdando `workspace.package.version` como os demais crates publicáveis (evita drift de versão do binário `serverust`).
- **BREAKING** (`serverust-telemetry`): `IdempotencyStore::try_acquire` devolve `AcquireOutcome::Acquired(LockToken)` em vez de `Acquired`, e `release`/`complete` passam a exigir o token da aquisição (fencing). Token divergente é no-op de sucesso. Implementações externas de `IdempotencyStore` e qualquer `match` sobre `AcquireOutcome::Acquired` precisam ser ajustados.

Expand All @@ -39,6 +43,7 @@ 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).
- `IdempotencyLayer`: após falha do handler ou de `complete()`, libera o lock `InProgress` via `IdempotencyStore::release`, permitindo que redeliveries do SQS reexecutem o handler dentro do TTL (antes o lock bloqueava reprocessamento por até 24h e a mensagem ia para DLQ sem nova tentativa). `release`/`complete` só mutam o registro se o token bater com o dono corrente — um owner cujo TTL expirou não apaga nem completa o lock de outro worker. Falha de `release` é logada com `tracing::warn` (chave + erro), sem mascarar o erro do handler. Falha de `complete` após sucesso do handler propaga erro ao SQS e libera o lock.
- `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

Expand Down
7 changes: 3 additions & 4 deletions serverust-events/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down
46 changes: 37 additions & 9 deletions serverust-events/src/broker/kafka.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -101,6 +101,7 @@ impl KafkaBrokerConfig {
pub struct KafkaBroker {
producer: FutureProducer<MskIamContext>,
subscriptions: Mutex<Vec<Subscription>>,
unhandled_topic_policy: UnhandledTopicPolicy,
}

/// Registro interno de uma inscrição (handler + tópico).
Expand Down Expand Up @@ -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<String> {
self.subscriptions
Expand All @@ -161,15 +169,35 @@ 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<BoxedHandler> = 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<BoxedHandler>, Vec<String>) = {
let guard = self
.subscriptions
.lock()
.map_err(|_| BrokerError::Subscribe("subscriptions mutex poisoned".into()))?;
let handlers: Vec<BoxedHandler> = 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 self
.unhandled_topic_policy
.on_unhandled(&msg.topic, &registered_topics);
}

for handler in handlers {
handler(msg.clone()).await?;
Expand Down
41 changes: 31 additions & 10 deletions serverust-events/src/broker/lambda.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand All @@ -37,6 +37,7 @@ use super::{BoxedHandler, Broker, BrokerError, BrokerMessage};
/// ou ao producer dedicado em uma futura US.
pub struct LambdaBroker {
subscriptions: Mutex<Vec<Subscription>>,
unhandled_topic_policy: UnhandledTopicPolicy,
}

struct Subscription {
Expand All @@ -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<String> {
self.subscriptions
Expand All @@ -68,24 +76,37 @@ 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, 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> {
for records in event.records.values() {
for raw in records {
let topic = raw.topic.clone().unwrap_or_default();

let handlers: Vec<BoxedHandler> = 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<BoxedHandler>, Vec<String>) = {
let guard = self.subscriptions.lock().map_err(|_| {
BrokerError::Subscribe("subscriptions mutex poisoned".into())
})?;
let handlers: Vec<BoxedHandler> = 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() {
self.unhandled_topic_policy
.on_unhandled(&topic, &registered_topics)?;
continue;
}

Expand Down
82 changes: 82 additions & 0 deletions serverust-events/src/broker/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String>
where
I: IntoIterator<Item = &'a str>,
{
let mut out: Vec<String> = 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
Expand Down Expand Up @@ -109,3 +160,34 @@ impl<B: Broker> Broker for Arc<B> {
(**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");
}
}
22 changes: 18 additions & 4 deletions serverust-events/src/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -108,10 +108,24 @@ 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) => {
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"
);
}
}
}
Err(last_err.unwrap_or_else(|| BrokerError::Subscribe("sem tentativas".to_string())))
Expand Down
33 changes: 33 additions & 0 deletions serverust-events/tests/kafka_broker_dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -73,6 +74,38 @@ async fn dispatch_em_topico_sem_subscriber_e_no_op() {
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,
offset: None,
key: None,
payload: Vec::new(),
headers: HashMap::new(),
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!(
text.contains("orders.created"),
"mensagem deve citar os tópicos inscritos; erro foi: {text}"
);
}

#[tokio::test]
async fn subscribed_topics_lista_topicos_inscritos() {
let broker = make_broker();
Expand Down
Loading
Loading