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-actionsfor 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.
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
}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}")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}`);
}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());
}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}");
}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}")
}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;
}
}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(())
}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}"
enddef 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
endSubscribe for Compliance Replay
The compliance dashboard replays the full audit history using StartFromFirst.
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)
}),
)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(),
)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),
});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());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}");
}client.subscribeToEventsStore {
channel = "audit.user-actions"
group = "compliance-reader"
startPosition = StartPosition.StartFromFirst
}.collect { msg ->
println("[Compliance] seq=${msg.sequence} action=${msg.metadata}")
}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; });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?;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{: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
Related
- Configure Retention for compliance retention windows
- Event Sourcing for state reconstruction from audit logs
- Events Store Reference for message structure details
Was this page helpful?