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
| From | To | Event |
|---|---|---|
| — | Created | order.created |
| Created | Paid | order.paid |
| Created | Cancelled | order.cancelled |
| Paid | Picking | order.picking |
| Picking | Shipped | order.shipped |
| Shipped | Delivered | order.delivered |
| Shipped | ReturnRequested | order.return_requested |
| ReturnRequested | Returned | order.returned |
Implementation
Publish State Transition Events
Each state change is published as a persistent event to a per-order channel.
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)
}
}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)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);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();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}");
}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}")
}
}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;
}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(())
}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{: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: deliveredProduction Considerations
Related
- Event Sourcing for the foundational pattern
- Replay Events for state reconstruction techniques
- Cross-Service Sync for multi-service state coordination
Was this page helpful?