# Cross-Service State Sync (/learn/events-store/scenarios/cross-service-sync)



This scenario synchronizes state across independent microservices using Events Store as a shared event log. When the Order Service creates an order, the Inventory Service, Shipping Service, and Notification Service each independently consume the event stream and update their local state.

## Architecture [#architecture]

<Mermaid
  chart="graph TD
  OS[&#x22;Order Service&#x22;]
  ES{{&#x22;Events Store channel<br/>orders.lifecycle&#x22;}}
  INV[&#x22;Inventory Service<br/>Reserve/release stock&#x22;]
  SHIP[&#x22;Shipping Service<br/>Create shipment&#x22;]
  NOTIF[&#x22;Notification Service<br/>Send email/SMS&#x22;]
  ANALYTICS[&#x22;Analytics<br/>Receives all events&#x22;]

  OS -- publish --> ES
  ES -- &#x22;group: inventory&#x22; --> INV
  ES -- &#x22;group: shipping&#x22; --> SHIP
  ES -- &#x22;group: notifications&#x22; --> NOTIF
  ES -. &#x22;no group&#x22; .-> ANALYTICS

  class ES store
  class OS,INV,SHIP,NOTIF,ANALYTICS client"
/>

*One durable `orders.lifecycle` log fans out to independent services, each tracking its own position via a consumer group; the ungrouped Analytics subscriber sees every event.*

### Design Decisions [#design-decisions]

* **Single event channel** — `orders.lifecycle` serves as the shared event log
* **Consumer groups per service** — each service has its own group for independent processing and position tracking
* **Durable subscriptions** — services that restart catch up from their last position automatically
* **Analytics fan-out** — an ungrouped subscriber receives every event for reporting

## Implementation [#implementation]

<Steps>
  <Step>
    ### Publish Order Lifecycle Events [#publish-order-lifecycle-events]

    The Order Service publishes state changes as they occur.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go title="order_service.go"
        func publishOrderEvent(ctx context.Context, client *kubemq.Client,
            orderId, eventType string, data map[string]interface{}) error {

            payload := map[string]interface{}{
                "orderId":   orderId,
                "eventType": eventType,
                "data":      data,
                "timestamp": time.Now().UTC().Format(time.RFC3339),
            }
            body, _ := json.Marshal(payload)

            result, err := client.SendEventStore(ctx, kubemq.NewEvent().
                SetChannel("orders.lifecycle").
                SetMetadata(eventType).
                SetBody(body).
                SetTags(map[string]string{"orderId": orderId}),
            )
            if err != nil {
                return err
            }
            log.Printf("Published %s for %s (seq: %s)", eventType, orderId, result.EventID)
            return nil
        }
        ```
      </Tab>

      <Tab value="Python">
        ```python title="order_service.py"
        def publish_order_event(client, order_id, event_type, data):
            payload = {
                "orderId": order_id,
                "eventType": event_type,
                "data": data,
                "timestamp": datetime.utcnow().isoformat() + "Z",
            }
            result = client.publish_event_store(
                EventStoreMessage(
                    channel="orders.lifecycle",
                    metadata=event_type,
                    body=json.dumps(payload).encode("utf-8"),
                    tags={"orderId": order_id},
                )
            )
            print(f"Published {event_type} for {order_id} (ID: {result.id})")
        ```
      </Tab>

      <Tab value="Node.js">
        ```typescript title="order_service.ts"
        async function publishOrderEvent(
          client: KubeMQClient, orderId: string, eventType: string, data: object
        ) {
          const payload = { orderId, eventType, data, timestamp: new Date().toISOString() };
          const result = await client.sendEventStore(
            createEventStoreMessage({
              channel: 'orders.lifecycle',
              metadata: eventType,
              body: JSON.stringify(payload),
              tags: { orderId },
            })
          );
          console.log(`Published ${eventType} for ${orderId} (ID: ${result.id})`);
        }
        ```
      </Tab>

      <Tab value="Java">
        ```java title="OrderService.java"
        public void publishOrderEvent(PubSubClient client, String orderId,
                String eventType, Map<String, Object> data) {
            var payload = Map.of(
                "orderId", orderId,
                "eventType", eventType,
                "data", data,
                "timestamp", Instant.now().toString());
            var body = objectMapper.writeValueAsBytes(payload);

            var result = client.sendEventsStoreMessage(EventStoreMessage.builder()
                .channel("orders.lifecycle")
                .metadata(eventType)
                .body(body)
                .tags(Map.of("orderId", orderId))
                .build());
            System.out.printf("Published %s for %s%n", eventType, orderId);
        }
        ```
      </Tab>

      <Tab value="C#">
        ```csharp title="OrderService.cs"
        public async Task PublishOrderEventAsync(KubeMQClient client,
            string orderId, string eventType, object data)
        {
            var payload = new { orderId, eventType, data, timestamp = DateTime.UtcNow };
            await client.SendEventStoreAsync(new EventStoreMessage
            {
                Channel = "orders.lifecycle",
                Metadata = eventType,
                Body = JsonSerializer.SerializeToUtf8Bytes(payload),
                Tags = { ["orderId"] = orderId },
            });
            Console.WriteLine($"Published {eventType} for {orderId}");
        }
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin title="OrderService.kt"
        suspend fun publishOrderEvent(client: KubeMQClient, orderId: String,
            eventType: String, data: Map<String, Any>) {
            val payload = mapOf("orderId" to orderId, "eventType" to eventType, "data" to data)
            client.sendEventStore(eventStoreMessage {
                channel = "orders.lifecycle"
                metadata = eventType
                body = Json.encodeToString(payload).toByteArray()
                tags = mapOf("orderId" to orderId)
            })
            println("Published $eventType for $orderId")
        }
        ```
      </Tab>

      <Tab value="C++">
        ```cpp title="order_service.cc"
        void publish_order_event(kubemq::Client& client,
            const std::string& order_id, const std::string& event_type) {
            kubemq::EventStoreMessage msg;
            msg.set_channel("orders.lifecycle");
            msg.set_metadata(event_type);
            msg.set_body("{\"orderId\":\"" + order_id + "\",\"eventType\":\"" + event_type + "\"}");
            msg.set_tag("orderId", order_id);
            client.SendEventStore(msg);
            std::cout << "Published " << event_type << " for " << order_id << std::endl;
        }
        ```
      </Tab>

      <Tab value="Rust">
        ```rust title="order_service.rs"
        async fn publish_order_event(
            client: &KubemqClient,
            order_id: &str,
            event_type: &str,
            data: &str,
        ) -> kubemq::Result<()> {
            let body = format!(
                r#"{{"orderId":"{}","eventType":"{}","data":{}}}"#,
                order_id, event_type, data
            );
            let event = EventStoreBuilder::new()
                .channel("orders.lifecycle")
                .metadata(event_type)
                .body(body.into_bytes())
                .add_tag("orderId", order_id)
                .build();

            let result = client.send_event_store(event).await?;
            println!(
                "Published {} for {} (id: {})",
                event_type, order_id, result.id
            );
            Ok(())
        }
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby title="order_service.rb"
        def publish_order_event(client, order_id, event_type, data)
          payload = {
            orderId: order_id,
            eventType: event_type,
            data: data,
            timestamp: Time.now.utc.iso8601
          }
          msg = KubeMQ::PubSub::EventStoreMessage.new(
            channel: 'orders.lifecycle',
            metadata: event_type,
            body: payload.to_json,
            tags: { 'orderId' => order_id }
          )
          result = client.send_event_store(msg)
          puts "Published #{event_type} for #{order_id} (sent: #{result.sent})"
        end
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir title="order_service.exs"
        def publish_order_event(client, order_id, event_type, data) do
          payload =
            Jason.encode!(%{
              orderId: order_id,
              eventType: event_type,
              data: data,
              timestamp: DateTime.utc_now() |> DateTime.to_iso8601()
            })

          event =
            KubeMQ.EventStore.new(
              channel: "orders.lifecycle",
              metadata: event_type,
              body: payload,
              tags: %{"orderId" => order_id}
            )

          {:ok, result} = KubeMQ.Client.send_event_store(client, event)
          IO.puts("Published #{event_type} for #{order_id} (sent: #{result.sent})")
        end
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Inventory Service Consumer [#inventory-service-consumer]

    The Inventory Service subscribes with its own group to reserve or release stock based on order events.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go title="inventory_service.go"
        sub, err := client.SubscribeToEventsStore(ctx,
            "orders.lifecycle",
            "inventory-service",
            kubemq.StartFromFirst(),
            kubemq.WithOnEvent(func(event *kubemq.Event) {
                switch event.Metadata {
                case "order.created":
                    log.Printf("[Inventory] Reserving stock for order %s",
                        event.Tags["orderId"])
                case "order.cancelled":
                    log.Printf("[Inventory] Releasing stock for order %s",
                        event.Tags["orderId"])
                }
            }),
            kubemq.WithOnError(func(err error) {
                log.Println("[Inventory] Error:", err)
            }),
        )
        ```
      </Tab>

      <Tab value="Python">
        ```python title="inventory_service.py"
        def on_order_event(event):
            order_id = event.tags.get("orderId", "unknown")
            if event.metadata == "order.created":
                print(f"[Inventory] Reserving stock for {order_id}")
            elif event.metadata == "order.cancelled":
                print(f"[Inventory] Releasing stock for {order_id}")

        client.subscribe_to_events_store(
            subscription=EventsStoreSubscription(
                channel="orders.lifecycle",
                group="inventory-service",
                start_position=EventStoreStartPosition.StartFromFirst,
                on_receive_event_callback=on_order_event,
                on_error_callback=lambda e: print(f"[Inventory] Error: {e}"),
            ),
            cancel=CancellationToken(),
        )
        ```
      </Tab>

      <Tab value="Node.js">
        ```typescript title="inventory_service.ts"
        client.subscribeToEventsStore({
          channel: 'orders.lifecycle',
          group: 'inventory-service',
          startPosition: EventStoreStartPosition.StartFromFirst,
          onEvent: (msg) => {
            const orderId = msg.tags?.orderId ?? 'unknown';
            switch (msg.metadata) {
              case 'order.created':
                console.log(`[Inventory] Reserving stock for ${orderId}`);
                break;
              case 'order.cancelled':
                console.log(`[Inventory] Releasing stock for ${orderId}`);
                break;
            }
          },
          onError: (err) => console.error('[Inventory] Error:', err.message),
        });
        ```
      </Tab>

      <Tab value="Java">
        ```java title="InventoryService.java"
        client.subscribeToEventsStore(EventsStoreSubscription.builder()
            .channel("orders.lifecycle")
            .group("inventory-service")
            .startPosition(EventStoreStartPosition.StartFromFirst)
            .onReceiveEventCallback(event -> {
                String orderId = event.getTags().getOrDefault("orderId", "unknown");
                switch (event.getMetadata()) {
                    case "order.created" ->
                        System.out.printf("[Inventory] Reserving stock for %s%n", orderId);
                    case "order.cancelled" ->
                        System.out.printf("[Inventory] Releasing stock for %s%n", orderId);
                }
            })
            .onErrorCallback(err ->
                System.err.println("[Inventory] Error: " + err.getMessage()))
            .build());
        ```
      </Tab>

      <Tab value="C#">
        ```csharp title="InventoryService.cs"
        await foreach (var msg in client.SubscribeToEventsStoreAsync(
            new EventsStoreSubscription
            {
                Channel = "orders.lifecycle",
                Group = "inventory-service",
                StartPosition = EventStoreStartPosition.StartFromFirst,
            }))
        {
            var orderId = msg.Tags.GetValueOrDefault("orderId", "unknown");
            switch (msg.Metadata)
            {
                case "order.created":
                    Console.WriteLine($"[Inventory] Reserving stock for {orderId}");
                    break;
                case "order.cancelled":
                    Console.WriteLine($"[Inventory] Releasing stock for {orderId}");
                    break;
            }
        }
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin title="InventoryService.kt"
        client.subscribeToEventsStore {
            channel = "orders.lifecycle"
            group = "inventory-service"
            startPosition = StartPosition.StartFromFirst
        }.collect { msg ->
            val orderId = msg.tags["orderId"] ?: "unknown"
            when (msg.metadata) {
                "order.created" -> println("[Inventory] Reserving stock for $orderId")
                "order.cancelled" -> println("[Inventory] Releasing stock for $orderId")
            }
        }
        ```
      </Tab>

      <Tab value="C++">
        ```cpp title="inventory_service.cc"
        client->SubscribeToEventsStore("orders.lifecycle", "inventory-service",
            kubemq::StartPosition::StartFromFirst,
            [](const kubemq::EventStoreReceived& msg) {
                if (msg.metadata() == "order.created") {
                    std::cout << "[Inventory] Reserving stock" << std::endl;
                } else if (msg.metadata() == "order.cancelled") {
                    std::cout << "[Inventory] Releasing stock" << std::endl;
                }
            },
            [](const std::string& err) {
                std::cerr << "[Inventory] Error: " << err << std::endl;
            });
        ```
      </Tab>

      <Tab value="Rust">
        ```rust title="inventory_service.rs"
        let sub = client
            .subscribe_to_events_store(
                "orders.lifecycle",
                "inventory-service",
                EventsStoreSubscription::StartFromFirst,
                |event| {
                    Box::pin(async move {
                        let order_id = event
                            .tags
                            .get("orderId")
                            .map(String::as_str)
                            .unwrap_or("unknown");
                        match event.metadata.as_str() {
                            "order.created" => {
                                println!("[Inventory] Reserving stock for {}", order_id)
                            }
                            "order.cancelled" => {
                                println!("[Inventory] Releasing stock for {}", order_id)
                            }
                            _ => {}
                        }
                    })
                },
                Some(|err| {
                    Box::pin(async move { eprintln!("[Inventory] Error: {}", err) })
                }),
            )
            .await?;
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby title="inventory_service.rb"
        cancel = KubeMQ::CancellationToken.new

        sub = KubeMQ::PubSub::EventsStoreSubscription.new(
          channel: 'orders.lifecycle',
          group: 'inventory-service',
          start_position: KubeMQ::PubSub::EventStoreStartPosition::START_FROM_FIRST
        )

        client.subscribe_to_events_store(sub, cancellation_token: cancel, on_error: lambda { |e|
          warn "[Inventory] Error: #{e.message}"
        }) do |event|
          order_id = event.tags.fetch('orderId', 'unknown')
          case event.metadata
          when 'order.created'
            puts "[Inventory] Reserving stock for #{order_id}"
          when 'order.cancelled'
            puts "[Inventory] Releasing stock for #{order_id}"
          end
        end
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir title="inventory_service.exs"
        {:ok, sub} =
          KubeMQ.Client.subscribe_to_events_store(client, "orders.lifecycle",
            start_at: :start_from_first,
            group: "inventory-service",
            on_event: fn event ->
              order_id = Map.get(event.tags, "orderId", "unknown")

              case event.metadata do
                "order.created" -> IO.puts("[Inventory] Reserving stock for #{order_id}")
                "order.cancelled" -> IO.puts("[Inventory] Releasing stock for #{order_id}")
                _ -> :ok
              end
            end,
            on_error: fn err -> IO.warn("[Inventory] Error: #{err}") end
          )
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Deploy Multiple Consumers [#deploy-multiple-consumers]

    Each service subscribes independently with its own group. All services can be started, stopped, and scaled independently.

    ```bash
    # Each service runs independently and tracks its own position
    ./inventory-service &    # group: inventory-service
    ./shipping-service &     # group: shipping-service
    ./notification-service & # group: notification-service
    ./analytics-service &    # no group (receives all events)
    ```

    If a service restarts, it automatically catches up from its last processed position.
  </Step>
</Steps>

## Production Considerations [#production-considerations]

<Accordions>
  <Accordion title="Idempotent Processing">
    Since Events Store provides at-least-once delivery, each consumer must handle duplicate events gracefully. Use the event's sequence number as an idempotency key — track the last processed sequence and skip events at or below that number.
  </Accordion>

  <Accordion title="Error Handling">
    If a consumer fails to process an event, it should log the error and continue processing subsequent events. Do not block the stream for a single failed event. Implement a dead-letter mechanism by publishing failed events to a separate `orders.lifecycle.dlq` channel for manual review.
  </Accordion>

  <Accordion title="Scaling Consumers">
    Within each group, add more instances to distribute the load. For example, the Inventory Service can run 3 instances in the `inventory-service` group — each event goes to exactly one instance. Scale based on event throughput and processing latency.
  </Accordion>

  <Accordion title="Event Schema Evolution">
    Include a `version` field in your event payloads. When the schema changes, consumers should handle both old and new formats. This allows gradual migration without coordinated deployments across all services.
  </Accordion>
</Accordions>

## Related [#related]

* [Consumer Groups](/learn/events-store/tutorials/consumer-groups) for scaling within a service
* [Resume After Disconnect](/learn/events-store/how-to/resume-after-disconnect) for durable position tracking
* [Audit Trail](/learn/events-store/scenarios/audit-trail) for compliance logging across services
