KubeMQ
LearnEvents StoreScenarios

Event-Driven State Machine

Implement an order processing state machine driven by persistent events.

This scenario implements an order processing state machine where each state transition is recorded as a persistent event. The full order lifecycle — from creation through delivery — is captured as an immutable event stream that can be replayed to rebuild state at any point.

Architecture

Each service publishes a transition to a per-order Events Store channel. The append-only stream is the source of truth; any consumer can replay it from the first sequence to rebuild the current state.

Transitions fan into one per-order channel; the persistent log is replayed to recover the latest state.

State Machine

The legal transition graph — the rebuilder applies stored events in sequence to walk these edges.

Valid Transitions

FromToEvent
Createdorder.created
CreatedPaidorder.paid
CreatedCancelledorder.cancelled
PaidPickingorder.picking
PickingShippedorder.shipped
ShippedDeliveredorder.delivered
ShippedReturnRequestedorder.return_requested
ReturnRequestedReturnedorder.returned

Implementation

Publish State Transition Events

Each state change is published as a persistent event to a per-order channel.

order_state_machine.go
package main

import (
    "context"
    "fmt"
    "log"

    "github.com/kubemq-io/kubemq-go/v2"
)

type StateTransition struct {
    OrderID   string `json:"orderId"`
    FromState string `json:"fromState"`
    ToState   string `json:"toState"`
    Event     string `json:"event"`
    Actor     string `json:"actor"`
}

func transitionOrder(ctx context.Context, client *kubemq.Client, t StateTransition) error {
    channel := fmt.Sprintf("order-state.%s", t.OrderID)
    body := fmt.Sprintf(`{"orderId":"%s","fromState":"%s","toState":"%s","event":"%s","actor":"%s"}`,
        t.OrderID, t.FromState, t.ToState, t.Event, t.Actor)

    result, err := client.SendEventStore(ctx, kubemq.NewEvent().
        SetChannel(channel).
        SetMetadata(t.Event).
        SetBody([]byte(body)),
    )
    if err != nil {
        return fmt.Errorf("transition failed: %w", err)
    }
    log.Printf("Transition: %s -> %s (seq: %s)", t.FromState, t.ToState, result.EventID)
    return nil
}

func main() {
    ctx := context.Background()
    client, _ := kubemq.NewClient(ctx, kubemq.WithAddress("localhost", 50000))
    defer client.Close()

    transitions := []StateTransition{
        {"ORD-2001", "", "created", "order.created", "customer:C-100"},
        {"ORD-2001", "created", "paid", "order.paid", "payment-service"},
        {"ORD-2001", "paid", "picking", "order.picking", "warehouse-worker:W-5"},
        {"ORD-2001", "picking", "shipped", "order.shipped", "shipping-service"},
        {"ORD-2001", "shipped", "delivered", "order.delivered", "carrier:fedex"},
    }

    for _, t := range transitions {
        transitionOrder(ctx, client, t)
    }
}
order_state_machine.py
import json
from kubemq import PubSubClient, EventStoreMessage

def transition_order(client, order_id, from_state, to_state, event, actor):
    body = json.dumps({
        "orderId": order_id, "fromState": from_state,
        "toState": to_state, "event": event, "actor": actor,
    })
    result = client.publish_event_store(
        EventStoreMessage(
            channel=f"order-state.{order_id}",
            metadata=event,
            body=body.encode("utf-8"),
        )
    )
    print(f"Transition: {from_state} -> {to_state} (ID: {result.id})")

with PubSubClient(address="localhost:50000") as client:
    transitions = [
        ("ORD-2001", "", "created", "order.created", "customer:C-100"),
        ("ORD-2001", "created", "paid", "order.paid", "payment-service"),
        ("ORD-2001", "paid", "picking", "order.picking", "warehouse:W-5"),
        ("ORD-2001", "picking", "shipped", "order.shipped", "shipping-service"),
        ("ORD-2001", "shipped", "delivered", "order.delivered", "carrier:fedex"),
    ]
    for oid, fs, ts, evt, actor in transitions:
        transition_order(client, oid, fs, ts, evt, actor)
order_state_machine.ts
import { KubeMQClient, createEventStoreMessage } from 'kubemq-js';

const client = await KubeMQClient.create({ address: 'localhost:50000' });

interface StateTransition {
  orderId: string; fromState: string; toState: string; event: string; actor: string;
}

async function transitionOrder(t: StateTransition) {
  const result = await client.sendEventStore(
    createEventStoreMessage({
      channel: `order-state.${t.orderId}`,
      metadata: t.event,
      body: JSON.stringify(t),
    })
  );
  console.log(`Transition: ${t.fromState} -> ${t.toState} (ID: ${result.id})`);
}

const transitions: StateTransition[] = [
  { orderId: 'ORD-2001', fromState: '', toState: 'created', event: 'order.created', actor: 'customer:C-100' },
  { orderId: 'ORD-2001', fromState: 'created', toState: 'paid', event: 'order.paid', actor: 'payment-service' },
  { orderId: 'ORD-2001', fromState: 'paid', toState: 'picking', event: 'order.picking', actor: 'warehouse:W-5' },
  { orderId: 'ORD-2001', fromState: 'picking', toState: 'shipped', event: 'order.shipped', actor: 'shipping-service' },
  { orderId: 'ORD-2001', fromState: 'shipped', toState: 'delivered', event: 'order.delivered', actor: 'carrier:fedex' },
];

for (const t of transitions) await transitionOrder(t);
OrderStateMachine.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("state-machine")
    .build();

String[][] transitions = {
    {"ORD-2001", "", "created", "order.created", "customer:C-100"},
    {"ORD-2001", "created", "paid", "order.paid", "payment-service"},
    {"ORD-2001", "paid", "picking", "order.picking", "warehouse:W-5"},
    {"ORD-2001", "picking", "shipped", "order.shipped", "shipping-service"},
    {"ORD-2001", "shipped", "delivered", "order.delivered", "carrier:fedex"},
};

for (String[] t : transitions) {
    String body = String.format(
        "{\"orderId\":\"%s\",\"fromState\":\"%s\",\"toState\":\"%s\",\"event\":\"%s\"}",
        t[0], t[1], t[2], t[3]);
    client.sendEventsStoreMessage(EventStoreMessage.builder()
        .channel("order-state." + t[0])
        .metadata(t[3])
        .body(body.getBytes())
        .build());
    System.out.printf("Transition: %s -> %s%n", t[1], t[2]);
}
client.close();
OrderStateMachine.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var transitions = new[] {
    ("", "created", "order.created"),
    ("created", "paid", "order.paid"),
    ("paid", "picking", "order.picking"),
    ("picking", "shipped", "order.shipped"),
    ("shipped", "delivered", "order.delivered"),
};

foreach (var (from, to, evt) in transitions)
{
    var body = $"{{\"orderId\":\"ORD-2001\",\"fromState\":\"{from}\",\"toState\":\"{to}\"}}";
    await client.SendEventStoreAsync(new EventStoreMessage
    {
        Channel = "order-state.ORD-2001",
        Metadata = evt,
        Body = Encoding.UTF8.GetBytes(body),
    });
    Console.WriteLine($"Transition: {from} -> {to}");
}
OrderStateMachine.kt
val client = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "state-machine"
}

data class Transition(val from: String, val to: String, val event: String)

val transitions = listOf(
    Transition("", "created", "order.created"),
    Transition("created", "paid", "order.paid"),
    Transition("paid", "picking", "order.picking"),
    Transition("picking", "shipped", "order.shipped"),
    Transition("shipped", "delivered", "order.delivered"),
)

client.use {
    for (t in transitions) {
        client.sendEventStore(eventStoreMessage {
            channel = "order-state.ORD-2001"
            metadata = t.event
            body = """{"fromState":"${t.from}","toState":"${t.to}"}""".toByteArray()
        })
        println("Transition: ${t.from} -> ${t.to}")
    }
}
order_state_machine.cc
kubemq::ClientOptions options;
options.set_address("localhost", 50000);
auto client = kubemq::Client::Create(options).value();

struct Transition { std::string from, to, event; };
std::vector<Transition> transitions = {
    {"", "created", "order.created"},
    {"created", "paid", "order.paid"},
    {"paid", "picking", "order.picking"},
    {"picking", "shipped", "order.shipped"},
    {"shipped", "delivered", "order.delivered"},
};

for (const auto& t : transitions) {
    kubemq::EventStoreMessage msg;
    msg.set_channel("order-state.ORD-2001");
    msg.set_metadata(t.event);
    msg.set_body("{\"fromState\":\"" + t.from + "\",\"toState\":\"" + t.to + "\"}");
    client->SendEventStore(msg);
    std::cout << "Transition: " << t.from << " -> " << t.to << std::endl;
}
order_state_machine.rs
use kubemq::prelude::*;
use kubemq::EventStoreBuilder;

#[tokio::main]
async fn main() -> kubemq::Result<()> {
    let client = KubemqClient::builder()
        .host("localhost")
        .port(50000)
        .build()
        .await?;

    let order_id = "ORD-2001";
    let channel = format!("order-state.{order_id}");

    let transitions = [
        ("", "created", "order.created"),
        ("created", "paid", "order.paid"),
        ("paid", "picking", "order.picking"),
        ("picking", "shipped", "order.shipped"),
        ("shipped", "delivered", "order.delivered"),
    ];

    for (from, to, event) in transitions {
        let body = format!(
            r#"{{"orderId":"{order_id}","fromState":"{from}","toState":"{to}"}}"#
        );
        let msg = EventStoreBuilder::new()
            .channel(&channel)
            .metadata(event)
            .body(body.into_bytes())
            .build();

        let result = client.send_event_store(msg).await?;
        println!("Transition: {from} -> {to} (id: {}, sent: {})", result.id, result.sent);
    }

    client.close().await?;
    Ok(())
}
order_state_machine.rb
require 'kubemq'

client = KubeMQ::PubSubClient.new(address: 'localhost:50000', client_id: 'state-machine')
order_id = 'ORD-2001'

transitions = [
  ['', 'created', 'order.created'],
  ['created', 'paid', 'order.paid'],
  ['paid', 'picking', 'order.picking'],
  ['picking', 'shipped', 'order.shipped'],
  ['shipped', 'delivered', 'order.delivered']
]

transitions.each do |from, to, event|
  body = { orderId: order_id, fromState: from, toState: to }.to_json
  msg = KubeMQ::PubSub::EventStoreMessage.new(
    channel: "order-state.#{order_id}",
    metadata: event,
    body: body
  )
  result = client.send_event_store(msg)
  puts "Transition: #{from} -> #{to} (id: #{result.id}, sent: #{result.sent})"
end

client.close
order_state_machine.exs
{:ok, client} =
  KubeMQ.Client.start_link(address: "localhost:50000", client_id: "state-machine")

order_id = "ORD-2001"

transitions = [
  {"", "created", "order.created"},
  {"created", "paid", "order.paid"},
  {"paid", "picking", "order.picking"},
  {"picking", "shipped", "order.shipped"},
  {"shipped", "delivered", "order.delivered"}
]

for {from, to, event} <- transitions do
  body = Jason.encode!(%{orderId: order_id, fromState: from, toState: to})

  event_msg =
    KubeMQ.EventStore.new(
      channel: "order-state.#{order_id}",
      metadata: event,
      body: body
    )

  {:ok, result} = KubeMQ.Client.send_event_store(client, event_msg)
  IO.puts("Transition: #{from} -> #{to} (id: #{result.id}, sent: #{result.sent})")
end

KubeMQ.Client.close(client)

Rebuild Current State

Subscribe with StartFromFirst to replay the transition history and determine the current state.

seq=1 order.created  -> state=created
seq=2 order.paid     -> state=paid
seq=3 order.picking  -> state=picking
seq=4 order.shipped  -> state=shipped
seq=5 order.delivered -> state=delivered

Current state: delivered

Production Considerations

Was this page helpful?

On this page