KubeMQ
LearnEvents StoreScenarios

Audit Trail System

Build a compliance-ready audit trail using KubeMQ Events Store for immutable event logging.

This scenario builds a compliance-ready audit trail that records every user action as an immutable event. Events Store's persistence, sequencing, and replay capabilities make it ideal for audit logging where data integrity and traceability are mandatory.

Architecture

Every service writes to one immutable audit channel; each consumer replays from a different start position.

Design Decisions

  • Single audit channel — all services publish to audit.user-actions for a unified timeline
  • Immutable events — once stored, events cannot be modified or deleted
  • Sequence numbers — provide a tamper-evident, monotonically increasing ordering
  • Multiple consumers — compliance, security, and SIEM systems each subscribe with different start positions

Implementation

Define the Audit Event Schema

Every audit event follows a consistent structure for compliance tooling:

{
  "eventType": "user.data.export",
  "actor": "user:U-1001",
  "resource": "orders:ORD-5001",
  "action": "export",
  "outcome": "success",
  "ip": "192.168.1.42",
  "userAgent": "Mozilla/5.0...",
  "timestamp": "2026-03-26T14:30:00Z",
  "metadata": {
    "exportFormat": "csv",
    "recordCount": 150
  }
}

Publish Audit Events from Services

Each service publishes audit events as part of its request handling.

audit_logger.go
func logAuditEvent(ctx context.Context, client *kubemq.Client, event AuditEvent) error {
    body, _ := json.Marshal(event)
    result, err := client.SendEventStore(ctx, kubemq.NewEvent().
        SetChannel("audit.user-actions").
        SetMetadata(event.EventType).
        SetBody(body).
        SetTags(map[string]string{
            "actor":    event.Actor,
            "resource": event.Resource,
            "outcome":  event.Outcome,
        }),
    )
    if err != nil {
        return fmt.Errorf("audit log failed: %w", err)
    }
    log.Printf("Audit logged: %s seq=%s", event.EventType, result.EventID)
    return nil
}
audit_logger.py
def log_audit_event(client, event: dict) -> None:
    result = client.publish_event_store(
        EventStoreMessage(
            channel="audit.user-actions",
            metadata=event["eventType"],
            body=json.dumps(event).encode("utf-8"),
            tags={
                "actor": event["actor"],
                "resource": event["resource"],
                "outcome": event["outcome"],
            },
        )
    )
    print(f"Audit logged: {event['eventType']} ID={result.id}")
audit_logger.ts
async function logAuditEvent(client: KubeMQClient, event: AuditEvent) {
  const result = await client.sendEventStore(
    createEventStoreMessage({
      channel: 'audit.user-actions',
      metadata: event.eventType,
      body: JSON.stringify(event),
      tags: {
        actor: event.actor,
        resource: event.resource,
        outcome: event.outcome,
      },
    })
  );
  console.log(`Audit logged: ${event.eventType} ID=${result.id}`);
}
AuditLogger.java
public void logAuditEvent(PubSubClient client, AuditEvent event) {
    var result = client.sendEventsStoreMessage(
        EventStoreMessage.builder()
            .channel("audit.user-actions")
            .metadata(event.getEventType())
            .body(objectMapper.writeValueAsBytes(event))
            .tags(Map.of(
                "actor", event.getActor(),
                "resource", event.getResource(),
                "outcome", event.getOutcome()))
            .build());
    System.out.printf("Audit logged: %s ID=%s%n", event.getEventType(), result.getId());
}
AuditLogger.cs
public async Task LogAuditEventAsync(KubeMQClient client, AuditEvent auditEvent)
{
    var result = await client.SendEventStoreAsync(new EventStoreMessage
    {
        Channel = "audit.user-actions",
        Metadata = auditEvent.EventType,
        Body = JsonSerializer.SerializeToUtf8Bytes(auditEvent),
        Tags = { ["actor"] = auditEvent.Actor, ["resource"] = auditEvent.Resource },
    });
    Console.WriteLine($"Audit logged: {auditEvent.EventType} ID={result.Id}");
}
AuditLogger.kt
suspend fun logAuditEvent(client: KubeMQClient, event: AuditEvent) {
    val result = client.sendEventStore(eventStoreMessage {
        channel = "audit.user-actions"
        metadata = event.eventType
        body = Json.encodeToString(event).toByteArray()
        tags = mapOf("actor" to event.actor, "resource" to event.resource)
    })
    println("Audit logged: ${event.eventType} ID=${result.id}")
}
audit_logger.cc
void log_audit_event(kubemq::Client& client, const AuditEvent& event) {
    kubemq::EventStoreMessage msg;
    msg.set_channel("audit.user-actions");
    msg.set_metadata(event.event_type);
    msg.set_body(event.to_json());
    msg.set_tag("actor", event.actor);
    msg.set_tag("resource", event.resource);
    auto result = client.SendEventStore(msg);
    if (result.ok()) {
        std::cout << "Audit logged: " << event.event_type << std::endl;
    }
}
audit_logger.rs
async fn log_audit_event(client: &KubemqClient, event: &AuditEvent) -> kubemq::Result<()> {
    let store_event = EventStoreBuilder::new()
        .channel("audit.user-actions")
        .metadata(&event.event_type)
        .body(serde_json::to_vec(event).unwrap())
        .add_tag("actor", &event.actor)
        .add_tag("resource", &event.resource)
        .add_tag("outcome", &event.outcome)
        .build();

    let result = client.send_event_store(store_event).await?;
    println!("Audit logged: {} id={}", event.event_type, result.id);
    Ok(())
}
audit_logger.rb
def log_audit_event(client, event)
  msg = KubeMQ::PubSub::EventStoreMessage.new(
    channel: 'audit.user-actions',
    metadata: event[:event_type],
    body: event.to_json,
    tags: {
      'actor' => event[:actor],
      'resource' => event[:resource],
      'outcome' => event[:outcome]
    }
  )
  result = client.send_event_store(msg)
  puts "Audit logged: #{event[:event_type]} sent=#{result.sent}"
end
audit_logger.exs
def log_audit_event(client, event) do
  store_event =
    KubeMQ.EventStore.new(
      channel: "audit.user-actions",
      metadata: event.event_type,
      body: Jason.encode!(event),
      tags: %{
        "actor" => event.actor,
        "resource" => event.resource,
        "outcome" => event.outcome
      }
    )

  case KubeMQ.Client.send_event_store(client, store_event) do
    {:ok, result} -> IO.puts("Audit logged: #{event.event_type} id=#{result.id}")
    {:error, err} -> IO.puts("Audit log failed: #{err.message}")
  end
end

Subscribe for Compliance Replay

The compliance dashboard replays the full audit history using StartFromFirst.

compliance_dashboard.go
sub, err := client.SubscribeToEventsStore(ctx,
    "audit.user-actions",
    "compliance-reader",
    kubemq.StartFromFirst(),
    kubemq.WithOnEvent(func(event *kubemq.Event) {
        fmt.Printf("[Compliance] seq=%d actor=%s action=%s\n",
            event.Sequence, event.Tags["actor"], event.Metadata)
    }),
    kubemq.WithOnError(func(err error) {
        log.Println("Error:", err)
    }),
)
compliance_dashboard.py
client.subscribe_to_events_store(
    subscription=EventsStoreSubscription(
        channel="audit.user-actions",
        group="compliance-reader",
        start_position=EventStoreStartPosition.StartFromFirst,
        on_receive_event_callback=lambda e: print(
            f"[Compliance] seq={e.sequence} actor={e.tags.get('actor')} "
            f"action={e.metadata}"
        ),
        on_error_callback=lambda e: print(f"Error: {e}"),
    ),
    cancel=CancellationToken(),
)
compliance_dashboard.ts
client.subscribeToEventsStore({
  channel: 'audit.user-actions',
  group: 'compliance-reader',
  startPosition: EventStoreStartPosition.StartFromFirst,
  onEvent: (msg) =>
    console.log(`[Compliance] seq=${msg.sequence} actor=${msg.tags?.actor} action=${msg.metadata}`),
  onError: (err) => console.error('Error:', err.message),
});
ComplianceDashboard.java
client.subscribeToEventsStore(EventsStoreSubscription.builder()
    .channel("audit.user-actions")
    .group("compliance-reader")
    .startPosition(EventStoreStartPosition.StartFromFirst)
    .onReceiveEventCallback(event ->
        System.out.printf("[Compliance] seq=%d actor=%s action=%s%n",
            event.getSequence(), event.getTags().get("actor"), event.getMetadata()))
    .onErrorCallback(err -> System.err.println(err.getMessage()))
    .build());
ComplianceDashboard.cs
await foreach (var msg in client.SubscribeToEventsStoreAsync(
    new EventsStoreSubscription
    {
        Channel = "audit.user-actions",
        Group = "compliance-reader",
        StartPosition = EventStoreStartPosition.StartFromFirst,
    }))
{
    Console.WriteLine($"[Compliance] seq={msg.Sequence} action={msg.Metadata}");
}
ComplianceDashboard.kt
client.subscribeToEventsStore {
    channel = "audit.user-actions"
    group = "compliance-reader"
    startPosition = StartPosition.StartFromFirst
}.collect { msg ->
    println("[Compliance] seq=${msg.sequence} action=${msg.metadata}")
}
compliance_dashboard.cc
client->SubscribeToEventsStore("audit.user-actions", "compliance-reader",
    kubemq::StartPosition::StartFromFirst,
    [](const kubemq::EventStoreReceived& msg) {
        std::cout << "[Compliance] seq=" << msg.sequence()
                  << " action=" << msg.metadata() << std::endl;
    },
    [](const std::string& err) { std::cerr << err << std::endl; });
compliance_dashboard.rs
let sub = client
    .subscribe_to_events_store(
        "audit.user-actions",
        "compliance-reader",
        EventsStoreSubscription::StartFromFirst,
        |event| {
            Box::pin(async move {
                println!(
                    "[Compliance] seq={} actor={} action={}",
                    event.sequence,
                    event.tags.get("actor").map(String::as_str).unwrap_or(""),
                    event.metadata,
                );
            })
        },
        None,
    )
    .await?;
compliance_dashboard.rb
cancel = KubeMQ::CancellationToken.new

sub = KubeMQ::PubSub::EventsStoreSubscription.new(
  channel: 'audit.user-actions',
  group: 'compliance-reader',
  start_position: KubeMQ::PubSub::EventStoreStartPosition::START_FROM_FIRST
)

client.subscribe_to_events_store(sub, cancellation_token: cancel, on_error: lambda { |e|
  puts "Error: #{e.message}"
}) do |event|
  puts "[Compliance] seq=#{event.sequence} actor=#{event.tags['actor']} action=#{event.metadata}"
end
compliance_dashboard.exs
{:ok, _sub} =
  KubeMQ.Client.subscribe_to_events_store(client, "audit.user-actions",
    group: "compliance-reader",
    start_at: :start_from_first,
    on_event: fn event ->
      IO.puts(
        "[Compliance] seq=#{event.sequence} " <>
          "actor=#{Map.get(event.tags, "actor")} action=#{event.metadata}"
      )
    end
  )

Production Considerations

Was this page helpful?

On this page