Sockudo
Server

Observability

Monitor Sockudo connections, subscriptions, publishes, fanout, history, recovery, webhooks, and push notifications.

Observability is part of the runtime contract. A realtime system should tell operators when it is connected, degraded, delayed, retrying, or dropping work.

Metrics endpoint

curl http://127.0.0.1:9601/metrics

Scrape the endpoint with Prometheus and label metrics by environment, region, node, adapter, and app where possible. Sockudo emits metrics through the metrics-rs recorder and exposes them through the Prometheus exporter by default.

TCP metrics exporter

For live debugging or sidecar consumers, Sockudo can also fan out metric events over TCP using metrics-exporter-tcp. This exporter streams protobuf-encoded metric events to connected clients; it is useful for local inspection and custom collectors, but Prometheus scraping should remain the primary production monitoring path.

[metrics.tcp_exporter]
enabled = true
host = "127.0.0.1"
port = 5000
buffer_size = 1024

The TCP exporter has bounded buffering by default. When buffers fill, event samples can be dropped to avoid blocking the server, so do not use it as the only source for alerts or SLO dashboards.

Core signals

AreaWatch
Connectionsactive sockets, connection attempts, disconnect reasons, heartbeat failures
Subscriptionssubscribe successes, auth failures, presence joins and leaves
Publishaccepted, rejected, idempotency hits, payload-too-large failures
Fanoutadapter publish latency, adapter receive latency, duplicate suppression
Recoveryresume successes, resume failures, replay counts, buffer misses
Historywrites, reads, retention purges, cursor errors
Webhooksqueued, delivered, failed, retried, dead-lettered
Pushaccepted, scheduled, dispatched, provider errors, publish status outcomes

When ably-compat is enabled, two lock-free operator snapshots supplement Prometheus:

  • /operator/stats/ably-runtime reports queued messages/bytes, overflow and continuity loss, total encoded frames, shared data-projection encodes (data_encoded), fanout counts, replay source and recovery backend calls, duplicate suppression, backend/degraded/reset state, expiry, and derived-filter cache pressure. data_encoded excludes ACK, heartbeat, connection, and channel-control frames, so it is the counter used to verify one encode per active data format.
  • /operator/stats/aggregation reports the bounded stats queue capacity, backlog, accepted, flushed, dropped, and failed observations.

These endpoints expose counters only. They never return credentials, device tokens, channel payloads, or raw encrypted values.

AI Transport

AI Transport observability is domain-blind. Sockudo reads only the well-known extras.ai.transport headers and never labels metrics by channel, run ID, message ID, or client ID. Headers are validated once and then passed as a typed borrowed view to metric and webhook classification. Empty run-client-id and step-client-id unknown-owner sentinels remain unchanged on the wire but are treated as absent when deriving an identity.

Metrics exposed at /metrics include:

  • sockudo_ai_runs_started_total
  • sockudo_ai_runs_ended_total{reason}
  • sockudo_ai_cancel_signals_total
  • sockudo_ai_active_streams
  • sockudo_ai_stream_duration_seconds
  • sockudo_ai_stream_bytes_total
  • sockudo_ai_messages_rejected_total{code}
  • sockudo_ai_messages_unparseable_total
  • sockudo_appends_received_total
  • sockudo_appends_delivered_total
  • sockudo_rollup_ratio
  • sockudo_flush_latency
  • sockudo_history_recovery_success_total{source}
  • sockudo_history_recovery_failures_total{code}
  • sockudo_versioned_message_mutations_total{action,result}
  • sockudo_versioned_message_retrieval_total{surface,result}

The four production signals are run outcomes, stream rate, rejects by code, and rollup efficiency. Use docs/public/grafana/ai-transport-observability.json as the starting Grafana dashboard.

Active-stream gauges use exact atomic insertion/removal deltas. Tracker entries that never receive a terminal event expire at the configured AI orphan TTL and decrement the same gauge exactly. AI lifecycle webhooks are awaited into the configured bounded webhook queue, which supplies backpressure and retry handling; Sockudo does not spawn an unbounded Tokio task per event or retain AI events in the optional in-process batching vector.

Dashboard starter panels:

PanelMetric
Active AI streamssockudo_ai_active_streams
Run outcomessockudo_ai_runs_started_total, sockudo_ai_runs_ended_total
Rejects by codesockudo_ai_messages_rejected_total
Rollup efficiencysockudo_appends_received_total vs sockudo_appends_delivered_total, sockudo_rollup_ratio
Rollup latencysockudo_flush_latency
Recovery healthsockudo_history_recovery_success_total, sockudo_history_recovery_failures_total
Durable statesockudo_history_degraded_channels, sockudo_history_reset_required_channels
Push backlogpush queue/status/provider metrics from the push dashboard

For support escalation, capture the channel, wall-clock time window, verified clientId, and first error code before collecting logs. Do not ask customers for provider payloads unless the codec layer explicitly requires them.

Persist completed AI runs

For Ably-style production persistence, persist completed runs from your backend instead of making Sockudo a domain store:

  1. Enable the ai_run_ended webhook for the app.
  2. On reason: "complete", query the history endpoint for the channel around the webhook time.
  3. Select messages with the same extras.ai.transport.run-id.
  4. Store the reduced transcript in your application database.
app.post("/sockudo/webhooks", async (req, res) => {
  for (const event of req.body.events) {
    if (event.name !== "ai_run_ended" || event.reason !== "complete") continue;

    const history = await sockudo.channelHistory(event.channel, {
      limit: 1000,
      direction: "newest_first",
    });

    await storeCompletedRun({
      runId: event.run_id,
      channel: event.channel,
      items: history.items.filter((item) => {
        return item.extras?.ai?.transport?.["run-id"] === event.run_id;
      }),
    });
  }
  res.sendStatus(200);
});

Logs

Use structured logs for events that operators need to investigate:

{
  "level": "warn",
  "target": "sockudo_push",
  "app_id": "app-id",
  "publish_id": "pub_123",
  "provider": "apns",
  "error": "BadDeviceToken"
}

Avoid logging secrets, raw auth signatures, provider tokens, or encrypted payloads.

Grafana dashboards

Recommended panels:

  • active connections by node
  • connection churn
  • publish rate by app
  • fanout latency histogram
  • subscription auth failure rate
  • recovery success ratio
  • replay buffer pressure
  • queue depth for webhooks and push
  • push provider failure rate by provider
  • APNs, FCM, Web Push latency by outcome
  • AI run outcomes by reason
  • AI stream duration and active streams
  • AI reject codes and unparseable headers
  • AI append rollup efficiency

Alerts

Alert on symptoms operators can act on:

  • readiness failures
  • adapter connection loss
  • rising publish failures
  • high auth rejection rate after deploy
  • recovery success ratio dropping
  • history write failures
  • webhook retry backlog
  • push queue backlog
  • push provider credential failures
  • push delivery status callback failures

Push status workflow

Push publishes are asynchronous. Store the publish_id returned by the API when a business workflow needs support visibility.

const response = await sockudo.publishPush({
  recipients: [{ type: "channel", channel: "orders" }],
  payload: { title: "Order updated", body: "Packed" },
  sync: false,
});

console.log(response.publish_id);

Then inspect status:

curl "https://realtime.example.com/apps/app-id/push/publish/pub_123/status"

Status records should be retained long enough for customer support and incident review.

On this page