Guia — IBrokerProvider
Um broker de mensageria novo, ativado por provider: num spec.brokers.
Exemplo mínimo (publisher-only)
using Pipevine.Abstractions;
using Pipevine.Core.Configuration;
using Pipevine.Messaging.Abstractions;
using Pipevine.Model.Consumers;
using Pipevine.Model.Resources;
// A chave usada no YAML: `provider: acmemq`
[Provider("acmemq", ProviderKind.Broker)]
public sealed class AcmeMessageProvider : IMessageProvider
{
public IMessagePublisher CreatePublisher(BrokerDefinition broker)
{
var config = PipevineNodeJsonConverter.ToJsonNode(broker.Config);
var url = config?["url"]?.GetValue<string>()
?? throw new InvalidOperationException($"Broker '{broker.Id}' requer 'url'.");
return new AcmePublisher(url);
}
// Um broker write-only é válido — lance para quem tentar consumir dele.
public IMessageConsumerTransport CreateConsumerTransport(BrokerDefinition broker, ConsumerDefinition consumer) =>
throw new NotSupportedException($"Broker '{broker.Id}' (acmemq): este provider é publisher-only.");
}
internal sealed class AcmePublisher(string url) : IMessagePublisher, IAsyncDisposable
{
public async ValueTask PublishAsync(MessageEnvelope envelope, string topic, string? partitionKey, CancellationToken ct)
{
// Envelope já vem no formato CloudEvents 1.0 — serialize envelope.Payload e envie.
// Publisher confirms não são opcionais: não retorne sucesso sem o broker ter de fato
// aceitado a mensagem. Retries internos são esperados; só depois de esgotados é que uma
// exceção propaga.
}
public ValueTask DisposeAsync() => ValueTask.CompletedTask;
}
Ativação
spec:
brokers:
- { id: events, provider: acmemq, url: "${env:ACME_MQ_URL}" }
events:
- { name: ItemCreated, broker: events, topic: item.created }
routes:
- id: createItem
pipeline:
- { type: publish, event: ItemCreated, payload: vars.item, mode: direct }
services.AddKeyedSingleton<IMessageProvider, AcmeMessageProvider>("acmemq");
// mais services.AddPipevineMessaging() do core, que registra o step genérico `publish`
mode: outbox (no step publish) exige que o provider participe do outbox transacional
(IOutboxStore) — mode: direct não exige nada além do publisher acima.
Envelope e serialização
O envelope padrão é CloudEvents 1.0 — o que dá interoperabilidade e serve
de base natural para a documentação AsyncAPI gerada. Um broker com formato nativo (Avro no Kafka,
por exemplo) pode oferecer seu próprio IMessageSerializer.
Para onde ir daqui
- Visão geral de extensibilidade
- Mensageria e eventos — como
brokers/events/publish/spec.consumerssão usados do ponto de vista de quem consome, não implementa, um broker.