Mensageria
O pacote messaging oferece uma interface única para:
- AWS SNS e SQS;
- Google Cloud Pub/Sub;
- RabbitMQ.
Selecionar o provedor
Para AWS, GCP ou Firebase, mantenha:
COLIBRI_MESSAGING=CLOUD_DEFAULT
O provedor é definido por CLOUD. Para RabbitMQ:
COLIBRI_MESSAGING=RABBITMQ
RABBITMQ_URL=amqp://guest:guest@localhost:5672/
Inicialize antes de criar produtores ou consumidores:
colibri.InitializeApp()
messaging.Initialize()
Publicar
type UserCreated struct {
ID string `json:"id" validate:"required"`
Email string `json:"email" validate:"required,email"`
}
producer := messaging.NewProducer("USERS_CREATED")
err := producer.Publish(ctx, "create", UserCreated{
ID: user.ID,
Email: user.Email,
})
A mensagem inclui ID, aplicação de origem, ação, payload, correlation ID e contexto de autenticação quando disponível.
Consumir
type userCreatedConsumer struct{}
func (userCreatedConsumer) QueueName() string {
return "NOTIFY_USER_CREATED"
}
func (userCreatedConsumer) Consume(
ctx context.Context,
message *messaging.ProviderMessage,
) error {
var payload UserCreated
if err := message.DecodeAndValidateMessage(&payload); err != nil {
return err
}
return notify(ctx, payload)
}
messaging.NewConsumer(userCreatedConsumer{})
O SDK executa Ack quando o consumer retorna nil. Quando retorna erro, executa Nack sem requeue; o comportamento final depende da configuração de DLQ do provedor.
Tentativas e atributos
if attempt := message.DeliveryAttempt(); attempt != nil && *attempt > 3 {
return ErrMaximumAttempts
}
origin := message.Attributes()["origin-service"]
Esses metadados só existem em mensagens recebidas e variam conforme o broker.
Testar sem broker
attempt := 2
message := messaging.NewConsumerMessage(
"create",
UserCreated{ID: "123", Email: "user@example.com"},
map[string]string{"origin-service": "users"},
&attempt,
)
err := userCreatedConsumer{}.Consume(context.Background(), message)
Mantenha consumers idempotentes: brokers podem entregar a mesma mensagem mais de uma vez.