# CESQL Routing (/connectors/cloudevents/how-to/cesql-routing)



CESQL (CloudEvents SQL) lets the KubeMQ routing table make fan-out decisions from a CloudEvent's **attributes** — its `type`, `source`, `subject`, and extensions — instead of matching on the channel name. Publishers send to one channel; the server evaluates CESQL rules against the `ce_*` tags and routes each event to the matching destinations.

## Overview [#overview]

The CloudEvents connector maps every CE attribute onto a KubeMQ message tag with a `ce_` prefix (`type` → `ce_type`, `source` → `ce_source`, and so on). The server-side **routing table** can carry CESQL rules alongside regular regex rules: each rule with `keyType: "cesql"` is a boolean expression evaluated against those attributes. When the expression returns `true`, the rule's `routes` are applied — potentially fanning the event across events, events store, and queues at once.

This is **content-based routing**: the decision comes from the message's attributes, not its body or channel name. It is configured entirely on the server (no client code change), and any message carrying `ce_*` tags — including gRPC or REST messages that set `ce_specversion` — is eligible, not just CE-connector traffic.

| Concept           | Detail                                                                                  |
| ----------------- | --------------------------------------------------------------------------------------- |
| Rule type         | `keyType: "cesql"` entry in the routing table                                           |
| Evaluated against | CloudEvent attributes reconstructed from `ce_*` tags                                    |
| Eligible messages | Any message with a `ce_specversion` tag                                                 |
| Destinations      | Standard routing syntax (`events:`, `events_store:`, `queues:`) with `{ce_*}` templates |
| Failure mode      | Fail-open — a failing rule is skipped, delivery continues                               |

<Callout type="info">
  CESQL routing is part of the **server-wide routing engine**, not a CE-only feature. This guide covers the CESQL-specific behavior.
</Callout>

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

A publisher sends a CloudEvent to a source channel; the connector tags it with `ce_*` attributes and hands it to the routing engine, which evaluates each CESQL rule and fans the event out to every matching destination channel.

<Mermaid
  chart="`
graph LR
PUB[&#x22;Publisher&#x22;]
CE[&#x22;CloudEvents connector<br/>:9090&#x22;]
RT[&#x22;Routing table<br/>(CESQL eval)&#x22;]
C1{{&#x22;events:order-archive&#x22;}}
C2{{&#x22;events:alert-stream&#x22;}}
C3{{&#x22;events:all-events&#x22;}}

PUB -- &#x22;POST /ce/send/event&#x22; --> CE
CE -- &#x22;ce_* tags&#x22; --> RT
RT -- &#x22;type = 'order'&#x22; --> C1
RT -- &#x22;type = 'alert'&#x22; --> C2
RT -- &#x22;type LIKE '...%'&#x22; --> C3

class CE connector
class RT broker
class C1,C2,C3 events
class PUB client
`"
/>

*The routing table evaluates each CESQL rule against the event's attributes and fans it out to every matching channel.*

## Configuration [#configuration]

CESQL rules live in the routing table. Set `keyType` to `"cesql"` and put a CESQL expression in `key`; the `routes` field uses the standard routing syntax and supports `{ce_*}` templates.

```json
[
  {
    "key": "type = 'com.example.order.created'",
    "keyType": "cesql",
    "routes": "events_store:order-archive;queues:order-processing"
  },
  {
    "key": "source = 'audit-service'",
    "keyType": "cesql",
    "routes": "events_store:{ce_source}-log"
  },
  {
    "key": "type LIKE 'com.example.%' AND source = 'critical-service'",
    "keyType": "cesql",
    "routes": "events:alerts"
  }
]
```

Enable routing and supply the table inline, from a file, or from a URL:

```toml
[Routing]
  Enable = true
  # Inline JSON (note the escaped single quotes inside the CESQL expression):
  Data = '[{"key":"type = '\''com.example.order.created'\''","keyType":"cesql","routes":"queues:order-processing"}]'
  # Or load from a file:
  # FilePath = "/etc/kubemq/routes.json"
  # Or from a config service:
  # URL = "http://config-service/routes"
  AutoReload = 0
```

<Callout type="warn">
  In CESQL expressions, reference attributes by their **plain CloudEvents name** (`type`, `source`, `subject`, `id`) with **no** `ce_` prefix. The `ce_` prefix only appears in `{ce_*}` template placeholders on the `routes` side.
</Callout>

## Supported operators [#supported-operators]

KubeMQ uses the CloudEvents SDK CESQL parser. The following operators and functions are available:

| Category         | Operators / Functions                                                      | Example                              |
| ---------------- | -------------------------------------------------------------------------- | ------------------------------------ |
| Comparison       | `=`, `!=`, `<>`, `<`, `<=`, `>`, `>=`                                      | `type = 'order.created'`             |
| Logical          | `AND`, `OR`, `NOT`                                                         | `type = 'order' AND source = '/app'` |
| String match     | `LIKE` (with `%` wildcard)                                                 | `type LIKE 'order.%'`                |
| Existence        | `EXISTS`                                                                   | `EXISTS priority`                    |
| Set              | `IN`                                                                       | `type IN ('a', 'b', 'c')`            |
| String functions | `CONCAT`, `LENGTH`, `LOWER`, `UPPER`, `TRIM`, `LEFT`, `RIGHT`, `SUBSTRING` | `LENGTH(type) > 10`                  |
| Type checks      | `IS_BOOL`, `IS_INT`                                                        | `IS_INT(source)`                     |
| Math             | `+`, `-`, `*`, `/`, `%`                                                    | integer attribute comparisons        |

## Template substitution [#template-substitution]

Route destinations can embed CE attributes with `{ce_*}` placeholders, letting one rule fan out to dynamically named channels:

| Placeholder    | Substituted with                            |
| -------------- | ------------------------------------------- |
| `{ce_type}`    | The CloudEvent `type` attribute             |
| `{ce_source}`  | The CloudEvent `source` attribute           |
| `{ce_subject}` | The CloudEvent `subject` attribute          |
| `{ce_id}`      | The CloudEvent `id` attribute               |
| `{ce_*}`       | Any CE attribute or extension from the tags |

For example, this rule routes each event to a channel named after its type:

```json
{
  "key": "type LIKE 'com.example.%'",
  "keyType": "cesql",
  "routes": "queues:{ce_type}"
}
```

A message with `type: "com.example.order.created"` is routed to `queues:com.example.order.created`. Semicolons in attribute values are stripped during substitution to prevent route injection.

## Evaluation scope [#evaluation-scope]

CESQL rules are evaluated against **any message that carries `ce_*` tags**, regardless of which connector produced it:

* Messages from the CE connector automatically have `ce_*` tags and are evaluated.
* Messages sent over gRPC or REST that include `ce_specversion` and other `ce_*` tags are also evaluated.
* Messages **without** `ce_*` tags skip CESQL rules entirely — they are only matched against regex rules.

CESQL and regex rules coexist in the same table. Each message is evaluated against all rules in order: CESQL rules match on attributes, regex rules match on the channel name.

## Error behavior [#error-behavior]

CESQL evaluation is **fail-open**:

* If an expression fails to evaluate (for example, a type mismatch), that rule is skipped, a warning is logged, and the remaining rules still run.
* Invalid expressions are caught when the routing table is loaded; the offending entry is skipped with a warning.
* The original message delivery is never blocked by a routing failure.

## Usage [#usage]

CESQL routing is configured on the server, so the client side is just a normal CloudEvent publish. The examples below publish events whose `type` matches the rules above, then subscribe to the destination channels to confirm the routing took effect.

<Callout type="info">
  These examples require KubeMQ to be running **with the CESQL routing rules configured** (the three-rule table shown in the comments below). They publish to a source channel and read from the routed destination channels.
</Callout>

<Tabs groupId="language" items="['curl','C#','Go','Java','JavaScript','Python','Ruby','Rust']">
  <Tab value="curl">
    ```bash
    # Publish a CloudEvent — the server's CESQL rules route it by `type`.
    curl -X POST http://localhost:9090/ce/send/event \
      -H "Content-Type: application/cloudevents+json" \
      -d '{
        "specversion": "1.0",
        "type": "com.kubemq.examples.routing.order",
        "source": "kubemq-ce-curl-cesql",
        "subject": "routing-source",
        "datacontenttype": "application/json",
        "data": {"description": "Event of type order"}
      }'

    # Read the routed destination channel (CESQL routed `order` -> events:order-archive).
    curl -N "http://localhost:9090/ce/subscribe/events?client_id=curl-order-sub&channel=order-archive"
    ```
  </Tab>

  <Tab value="C#">
    ```csharp
    // Example: routing/CesqlRouting — publish events for server-side CESQL routing.
    // Requires KubeMQ configured with CESQL rules (see Go example for config).
    // Run: dotnet run
    using CloudNative.CloudEvents;
    using CloudNative.CloudEvents.SystemTextJson;
    using System.Net.Http.Headers;
    using System.Text;
    using System.Text.Json;

    static string ServerUrl() => Environment.GetEnvironmentVariable("KUBEMQ_CE_URL") ?? "http://localhost:9090";

    var base_ = ServerUrl();
    var formatter = new JsonEventFormatter();
    Console.WriteLine("CESQL Routing Example — C#");
    Console.WriteLine("Requires KubeMQ with CESQL routing configured.\n");

    var results = new System.Collections.Concurrent.ConcurrentQueue<string>();

    // Subscribe to routed destination channels.
    async Task StartSubscriber(string ch, string clientId, int maxEvents)
    {
        using var http = new HttpClient { Timeout = Timeout.InfiniteTimeSpan };
        var url = $"{base_}/ce/subscribe/events?client_id={clientId}&channel={Uri.EscapeDataString(ch)}";
        using var req = new HttpRequestMessage(HttpMethod.Get, url);
        req.Headers.Accept.Add(new MediaTypeWithQualityHeaderValue("text/event-stream"));
        using var resp = await http.SendAsync(req, HttpCompletionOption.ResponseHeadersRead);
        using var reader = new StreamReader(await resp.Content.ReadAsStreamAsync(), Encoding.UTF8);
        string? evType = null, data = null, line; int count = 0;
        while ((line = await reader.ReadLineAsync()) != null && count < maxEvents)
        {
            if (line == "") {
                if (evType == "cloudevent" && data != null) {
                    var ce = JsonSerializer.Deserialize<JsonElement>(data);
                    results.Enqueue($"  [{ch}] type={ce.GetProperty("type")}");
                    count++;
                }
                evType = null; data = null;
            }
            else if (line.StartsWith("event:")) evType = line[6..].Trim();
            else if (line.StartsWith("data:")) data = line[5..].Trim();
        }
    }

    var subTasks = new List<Task> {
        Task.Run(() => StartSubscriber("order-archive", "csharp-order-sub", 1)),
        Task.Run(() => StartSubscriber("alert-stream",  "csharp-alert-sub", 1)),
        Task.Run(() => StartSubscriber("all-events",    "csharp-all-sub",   3)),
    };
    await Task.Delay(500);

    using var httpClient = new HttpClient();
    var eventTypes = new[] {
        "com.kubemq.examples.routing.order",
        "com.kubemq.examples.routing.alert",
        "com.kubemq.examples.routing.info",
    };
    foreach (var evType in eventTypes)
    {
        var ev = new CloudEvent {
            Id = Guid.NewGuid().ToString(),
            Type = evType,
            Source = new Uri("urn:kubemq-ce-csharp-cesql"),
            Subject = "routing-source",
            DataContentType = "application/json",
            Data = new { description = $"Event of type {evType}" },
        };
        var bytes = formatter.EncodeStructuredModeMessage(ev, out var ct);
        using var c = new ByteArrayContent(bytes.ToArray());
        c.Headers.ContentType = MediaTypeHeaderValue.Parse(ct.ToString());
        var r = await httpClient.PostAsync($"{base_}/ce/send/event", c);
        Console.WriteLine($"Published type={evType} (status={r.StatusCode})");
    }

    await Task.Delay(2000);
    while (results.TryDequeue(out var msg)) Console.WriteLine(msg);
    Console.WriteLine("\nCESQL routing demonstration complete.");
    ```
  </Tab>

  <Tab value="Go">
    ```go
    // Example: routing/cesql-routing
    //
    // Demonstrates KubeMQ CESQL routing with CloudEvents.
    // CESQL routing is SERVER-SIDE configuration — this example shows how to
    // publish events whose attributes match CESQL expressions, then verifies
    // routing worked by subscribing to the target channels.
    //
    // Configure KubeMQ with the following routing rules before running:
    //
    // [Routing]
    //   Enable = true
    //   Data = '[
    //     {"key":"type = '\''com.kubemq.examples.routing.order'\''","keyType":"cesql","routes":"events:order-archive"},
    //     {"key":"type = '\''com.kubemq.examples.routing.alert'\''","keyType":"cesql","routes":"events:alert-stream"},
    //     {"key":"type LIKE '\''com.kubemq.examples.routing.%'\''","keyType":"cesql","routes":"events:all-events"}
    //   ]'
    //
    // Run: go run ./routing/cesql-routing/main.go
    package main

    import (
    	"bufio"
    	"encoding/json"
    	"fmt"
    	"net/http"
    	"os"
    	"strings"
    	"time"

    	cloudevents "github.com/cloudevents/sdk-go/v2"
    )

    func serverURL() string {
    	if u := os.Getenv("KUBEMQ_CE_URL"); u != "" {
    		return u
    	}
    	return "http://localhost:9090"
    }

    func sendEvent(base, evType, channel string) {
    	event := cloudevents.NewEvent()
    	event.SetType(evType)
    	event.SetSource("kubemq-ce-go-cesql")
    	event.SetSubject(channel)
    	_ = event.SetData(cloudevents.ApplicationJSON, map[string]string{
    		"description": fmt.Sprintf("Event of type %s", evType),
    	})

    	body, _ := json.Marshal(event)
    	req, _ := http.NewRequest("POST", base+"/ce/send/event", strings.NewReader(string(body)))
    	req.Header.Set("Content-Type", "application/cloudevents+json")
    	resp, err := http.DefaultClient.Do(req)
    	if err != nil {
    		fmt.Printf("send error: %v\n", err)
    		return
    	}
    	defer resp.Body.Close()
    	fmt.Printf("Published type=%s to channel=%s (status=%d)\n", evType, channel, resp.StatusCode)
    }

    func subscribeAndPrint(base, channel, clientID string, maxMsgs int) {
    	sseURL := fmt.Sprintf("%s/ce/subscribe/events?client_id=%s&channel=%s",
    		base, clientID, channel)
    	req, _ := http.NewRequest("GET", sseURL, nil)
    	req.Header.Set("Accept", "text/event-stream")
    	req.Header.Set("Cache-Control", "no-cache")

    	client := &http.Client{Timeout: 0}
    	resp, err := client.Do(req)
    	if err != nil {
    		fmt.Printf("[%s] subscribe error: %v\n", channel, err)
    		return
    	}
    	defer resp.Body.Close()

    	count := 0
    	scanner := bufio.NewScanner(resp.Body)
    	var evType, data string
    	for scanner.Scan() {
    		line := scanner.Text()
    		if line == "" {
    			if evType == "cloudevent" && data != "" {
    				var ce map[string]interface{}
    				_ = json.Unmarshal([]byte(data), &ce)
    				fmt.Printf("  [%s] type=%v\n", channel, ce["type"])
    				count++
    				if count >= maxMsgs {
    					return
    				}
    			}
    			evType, data = "", ""
    			continue
    		}
    		if strings.HasPrefix(line, ":") {
    			continue
    		}
    		if strings.HasPrefix(line, "event:") {
    			evType = strings.TrimSpace(strings.TrimPrefix(line, "event:"))
    		} else if strings.HasPrefix(line, "data:") {
    			data = strings.TrimSpace(strings.TrimPrefix(line, "data:"))
    		}
    	}
    }

    func main() {
    	base := serverURL()

    	// Subscribe to routed destination channels.
    	go subscribeAndPrint(base, "order-archive", "go-cesql-order-sub", 1)
    	go subscribeAndPrint(base, "alert-stream", "go-cesql-alert-sub", 1)
    	go subscribeAndPrint(base, "all-events", "go-cesql-all-sub", 3)

    	time.Sleep(500 * time.Millisecond)

    	// Publish events — CESQL rules on the server route them to target channels.
    	sendEvent(base, "com.kubemq.examples.routing.order", "routing-source")
    	sendEvent(base, "com.kubemq.examples.routing.alert", "routing-source")
    	sendEvent(base, "com.kubemq.examples.routing.info", "routing-source")

    	// Allow routing to complete.
    	time.Sleep(2 * time.Second)
    	fmt.Println("\nCESQL routing demonstration complete.")
    }
    ```
  </Tab>

  <Tab value="Java">
    ```java
    // Example: routing/cesql-routing
    //
    // Demonstrates server-side CESQL routing. Events with different type values
    // are published to the source channel. KubeMQ routes them based on CESQL
    // expressions to different destination channels.
    //
    // Requires KubeMQ configured with CESQL routing rules, e.g.:
    // [Routing]
    //   Enable = true
    //   Data = '[{"key":"type = ''com.kubemq.examples.routing.order''","keyType":"cesql","routes":"events:order-archive"},...]'
    //
    // Run: mvn compile exec:java
    package io.kubemq.examples.routing.cesqlrouting;

    import com.fasterxml.jackson.databind.ObjectMapper;
    import io.cloudevents.CloudEvent;
    import io.cloudevents.core.builder.CloudEventBuilder;
    import io.cloudevents.core.format.EventFormat;
    import io.cloudevents.core.provider.EventFormatProvider;
    import io.cloudevents.jackson.JsonFormat;

    import java.io.BufferedReader;
    import java.io.InputStreamReader;
    import java.net.HttpURLConnection;
    import java.net.URI;
    import java.net.URL;
    import java.net.http.HttpClient;
    import java.net.http.HttpRequest;
    import java.net.http.HttpResponse;
    import java.nio.charset.StandardCharsets;
    import java.time.OffsetDateTime;
    import java.util.List;
    import java.util.Map;
    import java.util.UUID;
    import java.util.concurrent.ArrayBlockingQueue;
    import java.util.concurrent.BlockingQueue;
    import java.util.concurrent.TimeUnit;

    public class Main {
        static String serverUrl() {
            String u = System.getenv("KUBEMQ_CE_URL");
            return (u != null && !u.isEmpty()) ? u : "http://localhost:9090";
        }
        static final ObjectMapper MAPPER = new ObjectMapper();

        static void startSubscriber(String base, String channel, String clientId,
                                     int maxEvents, BlockingQueue<String> results) {
            String sseUrl = base + "/ce/subscribe/events?client_id=" + clientId + "&channel=" + channel;
            Thread.ofVirtual().start(() -> {
                try {
                    HttpURLConnection conn = (HttpURLConnection) new URL(sseUrl).openConnection();
                    conn.setRequestProperty("Accept", "text/event-stream");
                    conn.setReadTimeout(15_000);
                    try (BufferedReader reader = new BufferedReader(
                            new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))) {
                        String line; String evType = null, data = null; int count = 0;
                        while ((line = reader.readLine()) != null && count < maxEvents) {
                            if (line.isEmpty()) {
                                if ("cloudevent".equals(evType) && data != null) {
                                    Map<?, ?> ce = MAPPER.readValue(data, Map.class);
                                    results.offer("  [" + channel + "] type=" + ce.get("type"));
                                    count++;
                                }
                                evType = null; data = null;
                            } else if (line.startsWith("event:")) evType = line.substring(6).trim();
                            else if (line.startsWith("data:")) data = line.substring(5).trim();
                        }
                    }
                } catch (Exception e) { /* read ended */ }
            });
        }

        public static void main(String[] args) throws Exception {
            String base = serverUrl();
            System.out.println("CESQL Routing Example — Java");

            EventFormatProvider.getInstance().registerFormat(new JsonFormat());
            EventFormat format = EventFormatProvider.getInstance().resolveFormat(JsonFormat.CONTENT_TYPE);
            HttpClient httpClient = HttpClient.newHttpClient();

            // Subscribe to routed destination channels.
            BlockingQueue<String> results = new ArrayBlockingQueue<>(20);
            startSubscriber(base, "order-archive", "java-order-sub", 1, results);
            startSubscriber(base, "alert-stream",  "java-alert-sub", 1, results);
            startSubscriber(base, "all-events",    "java-all-sub",   3, results);
            Thread.sleep(500);

            // Publish events with different type values; CESQL rules route them.
            List<String> eventTypes = List.of(
                    "com.kubemq.examples.routing.order",
                    "com.kubemq.examples.routing.alert",
                    "com.kubemq.examples.routing.info"
            );
            for (String evType : eventTypes) {
                CloudEvent event = CloudEventBuilder.v1()
                        .withId(UUID.randomUUID().toString())
                        .withType(evType)
                        .withSource(URI.create("kubemq-ce-java-cesql"))
                        .withSubject("routing-source")
                        .withDataContentType("application/json")
                        .withTime(OffsetDateTime.now())
                        .withData("application/json",
                                MAPPER.writeValueAsBytes(Map.of("description", "Event of type " + evType)))
                        .build();
                HttpResponse<String> resp = httpClient.send(
                        HttpRequest.newBuilder().uri(URI.create(base + "/ce/send/event"))
                                .POST(HttpRequest.BodyPublishers.ofByteArray(format.serialize(event)))
                                .header("Content-Type", "application/cloudevents+json").build(),
                        HttpResponse.BodyHandlers.ofString());
                System.out.println("Published type=" + evType + " (status=" + resp.statusCode() + ")");
            }

            // Collect routed events (up to 5 with timeout).
            Thread.sleep(2000);
            for (int i = 0; i < 5; i++) {
                String msg = results.poll(100, TimeUnit.MILLISECONDS);
                if (msg != null) System.out.println(msg);
            }
            System.out.println("\nCESQL routing demonstration complete.");
        }
    }
    ```
  </Tab>

  <Tab value="JavaScript">
    ```typescript
    /**
     * Example: routing/cesql-routing — publish events for CESQL server-side routing.
     * Requires KubeMQ configured with CESQL rules (see Go example for config).
     * Run: npx tsx routing/cesql-routing/index.ts
     */
    import { CloudEvent, HTTP } from 'cloudevents';
    import EventSource from 'eventsource';

    function serverUrl(): string {
      return process.env.KUBEMQ_CE_URL ?? 'http://localhost:9090';
    }

    function subscribeAndCollect(
      base: string, channel: string, clientId: string, max: number, results: string[],
    ): void {
      const url = `${base}/ce/subscribe/events?client_id=${clientId}&channel=${encodeURIComponent(channel)}`;
      const es = new EventSource(url);
      let count = 0;
      es.addEventListener('cloudevent', (evt: MessageEvent) => {
        const ce = JSON.parse(evt.data) as Record<string, unknown>;
        results.push(`  [${channel}] type=${ce.type}`);
        count++;
        if (count >= max) es.close();
      });
      es.addEventListener('error', (err) => {
        console.error('SSE error:', err);
        es.close();
      });
    }

    async function main(): Promise<void> {
      const base = serverUrl();
      console.log('CESQL Routing Example — JavaScript/TypeScript');

      const results: string[] = [];
      subscribeAndCollect(base, 'order-archive', 'js-order-sub', 1, results);
      subscribeAndCollect(base, 'alert-stream', 'js-alert-sub', 1, results);
      subscribeAndCollect(base, 'all-events', 'js-all-sub', 3, results);

      await new Promise((r) => setTimeout(r, 500));

      for (const evType of [
        'com.kubemq.examples.routing.order',
        'com.kubemq.examples.routing.alert',
        'com.kubemq.examples.routing.info',
      ]) {
        const event = new CloudEvent({
          type: evType,
          source: 'kubemq-ce-js-cesql',
          subject: 'routing-source',
          datacontenttype: 'application/json',
          data: { description: `Event of type ${evType}` },
        });
        const msg = HTTP.structured(event);
        const resp = await fetch(`${base}/ce/send/event`, {
          method: 'POST',
          headers: msg.headers as Record<string, string>,
          body: msg.body as string,
        });
        console.log(`Published type=${evType} (status=${resp.status})`);
      }

      await new Promise((r) => setTimeout(r, 2000));
      for (const r of results) console.log(r);
      console.log('\nCESQL routing demonstration complete.');
    }

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

  <Tab value="Python">
    ```python
    # Example: routing/cesql_routing — publish events for server-side CESQL routing.
    #
    # CESQL routing is SERVER-SIDE configuration. Required KubeMQ config (TOML):
    # [Routing]
    #   Enable = true
    #   Data = '[
    #     {"key":"type = '\''com.kubemq.examples.routing.order'\''","keyType":"cesql","routes":"events:order-archive"},
    #     {"key":"type = '\''com.kubemq.examples.routing.alert'\''","keyType":"cesql","routes":"events:alert-stream"},
    #     {"key":"type LIKE '\''com.kubemq.examples.routing.%'\''","keyType":"cesql","routes":"events:all-events"}
    #   ]'
    from __future__ import annotations

    import json
    import os
    import threading
    import time

    import requests
    from cloudevents.v1.conversion import to_structured
    from cloudevents.v1.http import CloudEvent


    def server_url() -> str:
        return os.environ.get("KUBEMQ_CE_URL", "http://localhost:9090")


    def subscribe_and_print(base: str, channel: str, client_id: str,
                            results: list[str], max_msgs: int) -> None:
        sse_url = (f"{base}/ce/subscribe/events"
                   f"?client_id={client_id}&channel={channel}")
        with requests.get(sse_url, stream=True, timeout=None,
                          headers={"Accept": "text/event-stream"}) as resp:
            ev_type = data = ""
            count = 0
            for line in resp.iter_lines(decode_unicode=True):
                if line == "":
                    if ev_type == "cloudevent" and data:
                        ce = json.loads(data)
                        results.append(f"  [{channel}] type={ce.get('type')}")
                        count += 1
                        if count >= max_msgs:
                            return
                    ev_type = data = ""
                    continue
                if line.startswith(":"):
                    continue
                if line.startswith("event:"):
                    ev_type = line[6:].strip()
                elif line.startswith("data:"):
                    data = line[5:].strip()


    def main() -> None:
        base = server_url()
        print("CESQL Routing Example — Python\n")

        results: list[str] = []
        # Subscribe to routed channels.
        for ch, cid, n in [("order-archive", "py-order-sub", 1),
                            ("alert-stream", "py-alert-sub", 1),
                            ("all-events", "py-all-sub", 3)]:
            t = threading.Thread(
                target=subscribe_and_print,
                args=(base, ch, cid, results, n),
                daemon=True,
            )
            t.start()

        time.sleep(0.5)

        for ev_type in [
            "com.kubemq.examples.routing.order",
            "com.kubemq.examples.routing.alert",
            "com.kubemq.examples.routing.info",
        ]:
            event = CloudEvent(
                attributes={
                    "type": ev_type,
                    "source": "kubemq-ce-python-cesql",
                    "subject": "routing-source",
                    "datacontenttype": "application/json",
                },
                data={"description": f"Event of type {ev_type}"},
            )
            headers, body = to_structured(event)
            resp = requests.post(f"{base}/ce/send/event", data=body,
                                  headers=dict(headers), timeout=10)
            print(f"Published type={ev_type} (status={resp.status_code})")

        time.sleep(2)
        for r in results:
            print(r)
        print("\nCESQL routing demonstration complete.")


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

  <Tab value="Ruby">
    ```ruby
    # Example: routing/cesql_routing — publish events for server-side CESQL routing.
    # Requires KubeMQ configured with CESQL rules (see Go example for config).
    require "net/http"; require "uri"; require "json"; require "timeout"; require "securerandom"; require "cloud_events"
    def server_url = ENV.fetch("KUBEMQ_CE_URL", "http://localhost:9090")

    base = server_url; sdk = CloudEvents::HttpBinding.default
    puts "CESQL Routing Example — Ruby\n"

    results = Queue.new

    [["order-archive","ruby-order-sub",1],["alert-stream","ruby-alert-sub",1],["all-events","ruby-all-sub",3]].each do |ch, cid, max|
      Thread.new do
        uri = URI("#{base}/ce/subscribe/events?client_id=#{cid}&channel=#{URI.encode_www_form_component(ch)}")
        count = 0
        Net::HTTP.start(uri.host, uri.port) do |http|
          req = Net::HTTP::Get.new(uri); req["Accept"] = "text/event-stream"
          http.request(req) do |resp|
            ev_type = nil; data = nil
            resp.read_body do |chunk|
              chunk.each_line do |line|
                line.chomp!
                if line.empty?
                  if ev_type == "cloudevent" && data
                    ce = JSON.parse(data)
                    results.push("  [#{ch}] type=#{ce['type']}")
                    count += 1
                    Thread.exit if count >= max
                  end
                  ev_type = nil; data = nil
                elsif line.start_with?("event:") then ev_type = line.sub("event:","").strip
                elsif line.start_with?("data:") then data = line.sub("data:","").strip
                end
              end
            end
          end
        end
      end
    end

    sleep 0.5

    %w[com.kubemq.examples.routing.order com.kubemq.examples.routing.alert com.kubemq.examples.routing.info].each do |ev_type|
      ev = CloudEvents::Event::V1.new(
        id: SecureRandom.uuid, type: ev_type,
        source: URI("urn:kubemq-ce-ruby-cesql"), subject: "routing-source",
        spec_version: "1.0",
        data_content_type: CloudEvents::ContentType.new("application/json"),
        data: JSON.generate({ description: "Event of type #{ev_type}" }))
      enc_h, enc_b = sdk.encode_event(ev, structured_format: "json")
      uri = URI("#{base}/ce/send/event")
      Net::HTTP.start(uri.host, uri.port) do |http|
        req = Net::HTTP::Post.new(uri); enc_h.each{|k,v|req[k]=v}; req.body=enc_b
        res = http.request(req)
        puts "Published type=#{ev_type} (status=#{res.code})"
      end
    end

    sleep 2
    5.times { r = begin; Timeout.timeout(0.1) { results.pop }; rescue Timeout::Error; nil; end; puts r if r }
    puts "\nCESQL routing demonstration complete."
    ```
  </Tab>

  <Tab value="Rust">
    ```rust
    //! Example: routing/cesql-routing
    //!
    //! Publishes events with different type values for server-side CESQL routing.
    //! Requires KubeMQ configured with CESQL routing rules.
    //!
    //! Run: cargo run -p cesql-routing
    use bytes::Bytes;
    use cloudevents::{EventBuilder, EventBuilderV10};
    use futures_util::StreamExt;
    use reqwest::Client;
    use serde_json::{json, Value};
    use std::env;
    use tokio::sync::mpsc;
    use uuid::Uuid;

    fn server_url() -> String {
        env::var("KUBEMQ_CE_URL").unwrap_or_else(|_| "http://localhost:9090".to_string())
    }

    async fn start_subscriber(client: Client, base: String, channel: String, client_id: String, max: usize, tx: mpsc::Sender<String>) {
        let url = format!("{}/ce/subscribe/events?client_id={}&channel={}", base, client_id, channel);
        let stream = client.get(&url)
            .header("Accept", "text/event-stream")
            .send().await.expect("SSE connect").bytes_stream();
        let mut stream = Box::pin(stream);
        let mut ev_type = String::new(); let mut data = String::new();
        let mut buffer = String::new(); let mut count = 0usize;
        while let Some(chunk) = stream.next().await {
            let chunk: Bytes = chunk.unwrap_or_default();
            buffer.push_str(&String::from_utf8_lossy(&chunk));
            while let Some(pos) = buffer.find('\n') {
                let line = buffer[..pos].trim_end_matches('\r').to_string();
                buffer = buffer[pos + 1..].to_string();
                if line.is_empty() {
                    if ev_type == "cloudevent" && !data.is_empty() {
                        let ce: Value = serde_json::from_str(&data).unwrap_or(Value::Null);
                        let _ = tx.send(format!("  [{}] type={}", channel, ce["type"])).await;
                        count += 1; if count >= max { return; }
                    }
                    ev_type.clear(); data.clear();
                } else if line.starts_with(':') {
                } else if let Some(v) = line.strip_prefix("event:") { ev_type = v.trim().to_string(); }
                else if let Some(v) = line.strip_prefix("data:") { data = v.trim().to_string(); }
            }
        }
    }

    #[tokio::main]
    async fn main() -> Result<(), Box<dyn std::error::Error>> {
        let base = server_url();
        let client = Client::new();
        println!("CESQL Routing Example — Rust");

        let (tx, mut rx) = mpsc::channel::<String>(20);

        // Subscribe to routed destination channels.
        tokio::spawn(start_subscriber(client.clone(), base.clone(), "order-archive".to_string(), "rust-order-sub".to_string(), 1, tx.clone()));
        tokio::spawn(start_subscriber(client.clone(), base.clone(), "alert-stream".to_string(),  "rust-alert-sub".to_string(), 1, tx.clone()));
        tokio::spawn(start_subscriber(client.clone(), base.clone(), "all-events".to_string(),    "rust-all-sub".to_string(),   3, tx.clone()));
        tokio::time::sleep(tokio::time::Duration::from_millis(500)).await;

        // Publish events with different type values.
        let event_types = [
            "com.kubemq.examples.routing.order",
            "com.kubemq.examples.routing.alert",
            "com.kubemq.examples.routing.info",
        ];
        for ev_type in &event_types {
            let event = EventBuilderV10::new()
                .id(Uuid::new_v4().to_string())
                .ty(*ev_type)
                .source("urn:kubemq-ce-rust-cesql")
                .subject("routing-source")
                .data("application/json", json!({"description": format!("Event of type {}", ev_type)}))
                .build()?;
            let body = serde_json::to_string(&event)?;
            let resp = client.post(format!("{}/ce/send/event", base))
                .header("Content-Type", "application/cloudevents+json")
                .body(body).send().await?;
            println!("Published type={} (status={})", ev_type, resp.status());
        }

        tokio::time::sleep(tokio::time::Duration::from_secs(2)).await;
        while let Ok(msg) = rx.try_recv() { println!("{}", msg); }
        println!("\nCESQL routing demonstration complete.");
        Ok(())
    }
    ```
  </Tab>
</Tabs>

## Related [#related]

<Cards>
  <Card title="Events" href="/connectors/cloudevents/how-to/events" description="Publish and subscribe to CloudEvents — the pattern CESQL rules route between." />

  <Card title="Channel resolution" href="/connectors/cloudevents/how-to/channel-resolution" description="How subject, channel, and ClientID map onto KubeMQ channels." />

  <Card title="CE to KubeMQ mapping" href="/connectors/cloudevents/reference/ce-to-kubemq-mapping" description="The full CloudEvents attribute to ce_* tag mapping table." />
</Cards>
