# Events (/connectors/amqp/concepts/events)



Events are **fire-and-forget** pub/sub over the AMQP 1.0 connector. Attach a sender or receiver to a node whose address begins with `events/<channel>` and the connector binds the link to the KubeMQ **Events** pattern. Every active subscriber receives a **copy** of every message; there is no persistence and no replay.

## Overview [#overview]

The `events/` prefix selects the Events pattern (longest-prefix match, evaluated before `events-store/`). A producer attaches a **sender** to `events/<ch>`; each `TRANSFER` becomes a KubeMQ `SendEvents` with `Store=false`. A consumer attaches a **receiver** to the same node and grants credit; deliveries are **always pre-settled** (`settled=true`) — at-most-once, with no `DISPOSITION` round-trip.

Use Events for real-time notifications, telemetry, and broadcast where a missed message is acceptable. When you need durability and replay, use [Events Store](/connectors/amqp/concepts/events-store) instead.

| Operation      | AMQP action                                                                       | KubeMQ mapping                              |
| -------------- | --------------------------------------------------------------------------------- | ------------------------------------------- |
| Produce        | Attach a **sender** to `events/<ch>`, `TRANSFER` (pre-settle for fire-and-forget) | `SendEvents` (`Store=false`)                |
| Consume        | Attach a **receiver** to `events/<ch>`, grant credit                              | Pre-settled fan-out delivery (at-most-once) |
| Consumer group | Set link property `x-opt-kubemq-group` on the receiver                            | Load-balanced subset of the stream          |

## How it works [#how-it-works]

A published event fans out to every connected receiver on the channel. Receivers that share an `x-opt-kubemq-group` link property split the stream load-balanced; a receiver with no group is a plain fan-out subscriber.

<Mermaid
  chart="`
graph LR
PUB[&#x22;Producer<br/>(sender → events/orders)&#x22;]
CONN[&#x22;AMQP 1.0 connector<br/>:5672&#x22;]
BROKER[&#x22;Message Broker&#x22;]
S1[&#x22;Subscriber A<br/>(no group)&#x22;]
S2[&#x22;Subscriber B<br/>(no group)&#x22;]
G[&#x22;Group 'workers'<br/>(one member receives)&#x22;]

PUB -- &#x22;TRANSFER events/orders&#x22; --> CONN
CONN -- &#x22;SendEvents (Store=false)&#x22; --> BROKER
BROKER -. &#x22;copy&#x22; .-> S1
BROKER -. &#x22;copy&#x22; .-> S2
BROKER -. &#x22;load-balanced&#x22; .-> G

class PUB,S1,S2,G client
class CONN connector
class BROKER broker
`"
/>

*Each event is copied to every plain subscriber; members sharing an `x-opt-kubemq-group` split the stream between them.*

<Callout type="warn">
  **Events at 0 credit are silently dropped.** A message that arrives at a receiver whose link credit is 0 is discarded with no error and no `DISPOSITION` — that is what at-most-once means here. Grant a standing credit and replenish it eagerly, and **subscribe before you publish** (there is no replay to catch up from). The connector counts every drop in `kubemq_amqp10_events_dropped_no_credit_total`.
</Callout>

## Publish and subscribe [#publish-and-subscribe]

Each example below subscribes **first** (a receiver with a large standing credit), waits \~750 ms for the connector's subscription pump to go live, then publishes pre-settled events and drains them. Every client reads the broker endpoint from `KUBEMQ_AMQP_URL` (default `amqp://localhost:5672`).

<Tabs groupId="language" items="['Go','Python','Java','C#','JavaScript','Rust']">
  <Tab value="Go">
    ```go
    package main

    import (
    	"context"
    	"fmt"
    	"log"
    	"os"
    	"time"

    	amqp "github.com/Azure/go-amqp"
    )

    const channel = "amqp10.examples.pubsub"
    const total = 20
    const standingCredit = 100 // never let credit reach 0 — a 0-credit event is dropped

    func amqpURL() string {
    	if v := os.Getenv("KUBEMQ_AMQP_URL"); v != "" {
    		return v
    	}
    	return "amqp://localhost:5672"
    }

    func main() {
    	ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
    	defer cancel()

    	addr := "events/" + channel // events/ prefix → KubeMQ Events pattern

    	conn, err := amqp.Dial(ctx, amqpURL(), nil)
    	if err != nil {
    		log.Fatalf("dial: %v", err)
    	}
    	defer func() { _ = conn.Close() }()
    	session, err := conn.NewSession(ctx, nil)
    	if err != nil {
    		log.Fatalf("new session: %v", err)
    	}

    	// 1. SUBSCRIBE FIRST with standing credit. Events have no replay — a publish
    	//    that beats the subscription is lost forever.
    	receiver, err := session.NewReceiver(ctx, addr, &amqp.ReceiverOptions{Credit: standingCredit})
    	if err != nil {
    		log.Fatalf("new receiver: %v", err)
    	}
    	// The attach reply confirms the link, not that the subscription pump is live.
    	time.Sleep(750 * time.Millisecond)

    	// 2. PUBLISH pre-settled (fire-and-forget) — no DISPOSITION to await.
    	sender, err := session.NewSender(ctx, addr, &amqp.SenderOptions{
    		SettlementMode: amqp.SenderSettleModeSettled.Ptr(),
    	})
    	if err != nil {
    		log.Fatalf("new sender: %v", err)
    	}
    	for i := 0; i < total; i++ {
    		if err := sender.Send(ctx, amqp.NewMessage([]byte(fmt.Sprintf("event-%03d", i))), nil); err != nil {
    			log.Fatalf("publish: %v", err)
    		}
    	}
    	_ = sender.Close(ctx)

    	// 3. RECEIVE. Standing credit drains every event; accept is a no-op on
    	//    pre-settled fan-out but harmless.
    	seen := make(map[string]struct{}, total)
    	for len(seen) < total {
    		msg, err := receiver.Receive(ctx, nil)
    		if err != nil {
    			log.Fatalf("receive: %v", err)
    		}
    		_ = receiver.AcceptMessage(ctx, msg)
    		seen[string(msg.GetData())] = struct{}{}
    	}
    	fmt.Printf("received all %d events\n", len(seen))
    	_ = receiver.Close(ctx)
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python
    import os
    import time

    from proton import Message
    from proton.reactor import AtMostOnce
    from proton.utils import BlockingConnection

    CHANNEL = "amqp10.examples.pubsub"
    TOTAL = 20
    STANDING_CREDIT = 100  # never let credit reach 0 — a 0-credit event is dropped


    def amqp_url() -> str:
        return os.environ.get("KUBEMQ_AMQP_URL", "amqp://localhost:5672")


    def accept_if_unsettled(receiver) -> None:
        # Events fan-out deliveries are pre-settled, so accept() on a settled delivery
        # raises IndexError. This makes accept a true no-op on pre-settled pub/sub.
        if receiver.fetcher.unsettled:
            receiver.accept()


    def main() -> None:
        addr = "events/" + CHANNEL  # events/ prefix → KubeMQ Events pattern
        conn = BlockingConnection(amqp_url())
        try:
            # 1. SUBSCRIBE FIRST with standing credit (events have no replay).
            receiver = conn.create_receiver(addr, credit=STANDING_CREDIT)
            time.sleep(0.75)  # let the subscription pump go live before publishing

            # 2. PUBLISH pre-settled (AtMostOnce) — fire-and-forget, no DISPOSITION.
            sender = conn.create_sender(addr, options=AtMostOnce())
            for i in range(TOTAL):
                sender.send(Message(body=f"event-{i:03d}"))
            sender.close()

            # 3. RECEIVE. Standing credit drains every event.
            seen: set[str] = set()
            while len(seen) < TOTAL:
                msg = receiver.receive(timeout=30.0)
                accept_if_unsettled(receiver)
                seen.add(str(msg.body))
            print(f"received all {len(seen)} events")
            receiver.close()
        finally:
            conn.close()


    if __name__ == "__main__":
        main()
    ```
  </Tab>

  <Tab value="Java">
    ```java
    import java.util.HashSet;
    import java.util.Set;

    import javax.jms.Connection;
    import javax.jms.DeliveryMode;
    import javax.jms.Message;
    import javax.jms.MessageConsumer;
    import javax.jms.MessageProducer;
    import javax.jms.Session;
    import javax.jms.Topic;

    import org.apache.qpid.jms.JmsConnectionFactory;

    public final class Main {
        private static final String CHANNEL = "amqp10.examples.pubsub";
        private static final int TOTAL = 20;

        public static void main(String[] args) throws Exception {
            String url = System.getenv().getOrDefault("KUBEMQ_AMQP_URL", "amqp://localhost:5672");
            String address = "events/" + CHANNEL; // events/ prefix → KubeMQ Events pattern

            JmsConnectionFactory factory = new JmsConnectionFactory(url);
            try (Connection connection = factory.createConnection()) {
                connection.start();
                try (Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE)) {
                    Topic topic = session.createTopic(address);

                    // 1. SUBSCRIBE FIRST (Qpid JMS grants a standing prefetch credit).
                    try (MessageConsumer consumer = session.createConsumer(topic)) {
                        Thread.sleep(750); // let the subscription pump go live

                        // 2. PUBLISH NON_PERSISTENT (fire-and-forget, pre-settled).
                        try (MessageProducer producer = session.createProducer(topic)) {
                            producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
                            for (int i = 0; i < TOTAL; i++) {
                                producer.send(session.createTextMessage(String.format("event-%03d", i)));
                            }
                        }

                        // 3. RECEIVE. The connector re-emits the body as a Data section;
                        //    getBody(String.class) decodes either type as UTF-8.
                        Set<String> seen = new HashSet<>();
                        while (seen.size() < TOTAL) {
                            Message msg = consumer.receive(30_000);
                            if (msg == null) throw new IllegalStateException("timed out");
                            seen.add(msg.getBody(String.class));
                        }
                        System.out.printf("received all %d events%n", seen.size());
                    }
                }
            }
        }
    }
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    using System.Text;
    using Amqp;
    using Amqp.Framing;

    const string channel = "amqp10.examples.pubsub";
    const int total = 20;
    const int standingCredit = 100; // never let credit reach 0 — a 0-credit event is dropped

    static string AmqpUrl() =>
        Environment.GetEnvironmentVariable("KUBEMQ_AMQP_URL") is { Length: > 0 } v
            ? v
            : "amqp://localhost:5672";

    var addr = "events/" + channel; // events/ prefix → KubeMQ Events pattern

    var connection = await Connection.Factory.CreateAsync(new Address(AmqpUrl()));
    try
    {
        var session = new Session(connection);

        // 1. SUBSCRIBE FIRST with standing credit (autoRestore replenishes on settle).
        var receiver = new ReceiverLink(session, "pubsub-receiver", addr);
        receiver.SetCredit(standingCredit, autoRestore: true);
        await Task.Delay(750); // let the subscription pump go live before publishing

        // 2. PUBLISH pre-settled — SndSettleMode.Settled marks every TRANSFER settled.
        var senderAttach = new Attach
        {
            Source = new Source(),
            Target = new Target { Address = addr },
            SndSettleMode = SenderSettleMode.Settled,
        };
        var sender = new SenderLink(session, "pubsub-sender", senderAttach, null);
        for (var i = 0; i < total; i++)
        {
            var message = new Message { BodySection = new Data { Binary = Encoding.UTF8.GetBytes($"event-{i:D3}") } };
            sender.Send(message, TimeSpan.FromSeconds(15));
        }

        // 3. RECEIVE. Standing credit drains every event; Accept is a no-op here.
        var seen = new HashSet<string>();
        while (seen.Count < total)
        {
            var message = receiver.Receive(TimeSpan.FromSeconds(30))
                ?? throw new InvalidOperationException("receive timed out");
            receiver.Accept(message);
            seen.Add(BodyString(message));
        }
        Console.WriteLine($"received all {seen.Count} events");

        await sender.CloseAsync();
        await receiver.CloseAsync();
        await session.CloseAsync();
    }
    finally
    {
        await connection.CloseAsync();
    }

    static string BodyString(Message message) => message.BodySection switch
    {
        Data d => Encoding.UTF8.GetString(d.Binary),
        AmqpValue { Value: byte[] bytes } => Encoding.UTF8.GetString(bytes),
        AmqpValue { Value: string str } => str,
        AmqpValue v => v.Value?.ToString() ?? string.Empty,
        _ => string.Empty,
    };
    ```
  </Tab>

  <Tab value="JavaScript">
    ```typescript
    import {
      Connection,
      ReceiverEvents,
      type EventContext,
      type Receiver,
    } from "rhea-promise";

    const channel = "amqp10.examples.pubsub";
    const total = 20;
    const standingCredit = 100; // never let credit reach 0 — a 0-credit event is dropped

    function sleep(ms: number): Promise<void> {
      return new Promise((resolve) => setTimeout(resolve, ms));
    }

    function bodyToString(body: unknown): string {
      return Buffer.isBuffer(body) ? body.toString("utf8") : String(body);
    }

    async function main(): Promise<void> {
      const url = new URL(process.env["KUBEMQ_AMQP_URL"] ?? "amqp://localhost:5672");
      const address = `events/${channel}`; // events/ prefix → KubeMQ Events pattern

      const connection = new Connection({
        host: url.hostname,
        port: url.port ? Number(url.port) : 5672,
        container_id: `kubemq-amqp10-js-pubsub-${process.pid}`,
        reconnect: false,
      });
      await connection.open();

      try {
        // 1. SUBSCRIBE FIRST. Register the handler before granting credit so no
        //    early delivery is missed (events have no replay).
        const receiver = await connection.createReceiver({
          source: { address },
          credit_window: 0,
          autoaccept: false,
          autosettle: false,
        });

        const seen = new Set<string>();
        const received = drainEvents(receiver, (ctx) => {
          ctx.delivery?.accept(); // no-op for pre-settled fan-out, but harmless
          seen.add(bodyToString(ctx.message?.body));
          return seen.size >= total;
        }, standingCredit, 30_000);

        await sleep(750); // let the subscription pump go live before publishing

        // 2. PUBLISH pre-settled (snd_settle_mode: 1) — fire-and-forget.
        const sender = await connection.createSender({
          target: { address },
          snd_settle_mode: 1,
          autosettle: true,
        });
        for (let i = 0; i < total; i++) {
          sender.send({ body: `event-${String(i).padStart(3, "0")}` });
        }
        await sender.close();

        await received;
        console.log(`received all ${seen.size} events`);
        await receiver.close();
      } finally {
        await connection.close();
      }
    }

    // Grants standing credit and tops it back up as messages arrive so the
    // subscriber is never starved (a 0-credit event is silently dropped).
    function drainEvents(
      receiver: Receiver,
      onMessage: (ctx: EventContext) => boolean,
      credit: number,
      timeoutMs: number,
    ): Promise<void> {
      return new Promise<void>((resolve, reject) => {
        const timer = setTimeout(() => {
          receiver.removeListener(ReceiverEvents.message, handler);
          reject(new Error("timed out waiting for events"));
        }, timeoutMs);
        const handler = (ctx: EventContext): void => {
          if (onMessage(ctx)) {
            clearTimeout(timer);
            receiver.removeListener(ReceiverEvents.message, handler);
            resolve();
            return;
          }
          receiver.addCredit(1); // replenish so standing credit never drains to 0
        };
        receiver.on(ReceiverEvents.message, handler);
        receiver.addCredit(credit);
      });
    }

    main().catch((err) => {
      console.error(err);
      process.exit(1);
    });
    ```
  </Tab>

  <Tab value="Rust">
    ```rust
    use std::collections::HashSet;
    use std::time::Duration;

    use fe2o3_amqp::link::delivery::Delivery;
    use fe2o3_amqp::link::receiver::CreditMode;
    use fe2o3_amqp::{Connection, Receiver, Sender, Session};
    use fe2o3_amqp_types::definitions::SenderSettleMode;
    use fe2o3_amqp_types::messaging::{Body, Message};
    use fe2o3_amqp_types::primitives::Value;

    const CHANNEL: &str = "amqp10.examples.pubsub";
    const TOTAL: usize = 20;
    const STANDING_CREDIT: u32 = 100; // never let credit reach 0 — a 0-credit event is dropped

    fn amqp_url() -> String {
        std::env::var("KUBEMQ_AMQP_URL").unwrap_or_else(|_| "amqp://localhost:5672".to_string())
    }

    fn body_string(msg: &Message<Body<Value>>) -> String {
        let bytes = match &msg.body {
            Body::Data(batch) => batch.iter().flat_map(|d| d.0.to_vec()).collect(),
            Body::Value(v) => match &v.0 {
                Value::Binary(b) => b.to_vec(),
                Value::String(s) => s.clone().into_bytes(),
                other => format!("{other:?}").into_bytes(),
            },
            _ => Vec::new(),
        };
        String::from_utf8_lossy(&bytes).into_owned()
    }

    #[tokio::main]
    async fn main() -> Result<(), Box<dyn std::error::Error>> {
        let addr = format!("events/{CHANNEL}"); // events/ prefix → KubeMQ Events pattern

        let mut connection = Connection::open("amqp10-examples-pubsub", amqp_url().as_str()).await?;
        let mut session = Session::begin(&mut connection).await?;

        // 1. SUBSCRIBE FIRST with standing credit (CreditMode::Auto auto-replenishes).
        let mut receiver = Receiver::builder()
            .name("basic-pubsub-receiver")
            .source(addr.as_str())
            .credit_mode(CreditMode::Auto(STANDING_CREDIT))
            .attach(&mut session)
            .await?;
        tokio::time::sleep(Duration::from_millis(750)).await; // let the pump go live

        // 2. PUBLISH pre-settled (SenderSettleMode::Settled) — fire-and-forget.
        let mut sender = Sender::builder()
            .name("basic-pubsub-sender")
            .target(addr.as_str())
            .sender_settle_mode(SenderSettleMode::Settled)
            .attach(&mut session)
            .await?;
        for i in 0..TOTAL {
            sender.send(format!("event-{i:03}")).await?;
        }
        sender.close().await?;

        // 3. RECEIVE. Standing credit drains every event.
        let mut seen: HashSet<String> = HashSet::with_capacity(TOTAL);
        while seen.len() < TOTAL {
            let delivery: Delivery<Body<Value>> = receiver.recv().await?;
            let _ = receiver.accept(&delivery).await; // no-op on pre-settled, harmless
            seen.insert(body_string(delivery.message()));
        }
        println!("received all {} events", seen.len());

        receiver.close().await?;
        session.end().await?;
        connection.close().await?;
        Ok(())
    }
    ```
  </Tab>
</Tabs>

## Consumer groups [#consumer-groups]

A receiver with **no*&#x2A; group is a plain fan-out subscriber — it receives every event (the KubeMQ default). Set the link property &#x2A;*`x-opt-kubemq-group`** on a receiver's `ATTACH` to join a consumer group: within one group the stream is **load-balanced** across members (each message goes to exactly one member), while different groups each get the full stream independently.

```text
events/orders
   ├── group "g1": receiver-A ┐  (split — no duplicate within g1)
   │               receiver-B ┘
   └── group "g2": receiver-C    (full stream)
```

In Go, set it on the receiver's link properties at attach:

```go
receiver, err := session.NewReceiver(ctx, "events/orders", &amqp.ReceiverOptions{
	Credit:     100,
	Properties: map[string]any{"x-opt-kubemq-group": "g1"},
})
```

The `x-opt-kubemq-group` property is honored on `events`, `events-store`, and the RPC consume patterns. The Go, Python, C#, JavaScript, and Rust clients can all set it on the receiver's `ATTACH`.

<Callout type="info">
  **Qpid JMS (Java) cannot join a consumer group today.** The connector advertises no `SHARED-SUBS` capability, so `createSharedConsumer` / `createSharedDurableConsumer&#x60; throws &#x2A;"Remote peer does not support shared subscriptions"*, and Qpid JMS exposes no API to set the `x-opt-kubemq-group` link property directly. Java is fan-out only on Events; the other five languages support groups fully.
</Callout>

## Related [#related]

<Cards>
  <Card title="Events Store" href="/connectors/amqp/concepts/events-store" description="Durable, replayable pub/sub — resume after a disconnect and choose a start position." />

  <Card title="Queues" href="/connectors/amqp/concepts/queues" description="Competing-consumer work queues with at-least-once delivery and settlement." />

  <Card title="Address mapping" href="/connectors/amqp/reference/address-mapping" description="The full address grammar, longest-prefix matching, and the events/ prefix." />
</Cards>
