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