DEV Community

GodfreySterling1574
GodfreySterling1574

Posted on

Published Go Events Never Seen: Debug Name Typos Against Supported Types

Short answer: for a marketplace workspace that must show who is online, start with ephemeral presence plus a bounded reconnect backfill, and validate every event constant against the provider's supported type list before serving traffic. A misspelled event name produces the most misleading symptom in this system: the publish appears to happen, no subscriber receives it, and no error explains the gap. The least complex fix is a startup assertion over one authoritative constants module.

Do that before changing socket transports, retry intervals, or retention. Those changes cannot repair a name mismatch.

What does reconnect actually cost?

The dominant cost is usually not the tiny online/offline payload. It is retained history multiplied by active workspaces, event frequency, and the backfill window. If 10,000 devices each report status once per minute, the ingress is 600,000 status events per hour; keeping seven days would mean 100,800,000 event records before indexes, replicas, or audit metadata. This is an arithmetic example, not a measured vendor benchmark, but it exposes the term that controls both storage and replay work.

For "who is online now," most of those records have no continuing product value. Keep a current presence projection and a short, cursor-addressed replay window sufficient for ordinary reconnects; send offline notification work to a durable queue instead of stretching the realtime log into a job system. The queue holds the notification while the socket delivers live state, so an offline user becomes durable work rather than a lost publish.

The change that moves the dominant term is the retention boundary. A one-hour backfill at the same illustrative rate contains 600,000 records, rather than the 100,800,000 implied by seven days. Longer disconnections should rebuild from an authoritative device-status snapshot and then resume after its cursor, because replaying an old presence narrative can briefly resurrect users who are no longer online.

Deliberately stop retaining presence events once they are outside the reconnect window. When an investigation happens later, that choice costs event-by-event reconstruction: the audit trail can prove the accepted snapshot, cursor, event type, publisher request ID, and queued notification decision, but it cannot reproduce every transient presence transition. Compliance or dispute workflows that require immutable history need a separate ledger-like record with an explicit retention policy; an ephemeral presence stream is the wrong evidence store.

How should I debug published events that are never seen?

Silent non-delivery is the signature of a name mismatch. A publisher emitting a local constant while subscribers bind to a different supported type does not establish a shared contract, even if the payload and channel are otherwise correct. Check the type name first.

Names first.

This is an exactly-once mindset applied at the contract boundary: before reasoning about duplicate delivery, prove that publisher and subscriber identify the same operation. Keep constants in one module, ask the supported-type discovery surface for its list during startup, and fail readiness if any local constant is absent. The check is exhaustive only when developers cannot introduce event-name literals elsewhere.

The following Go program performs that assertion and then publishes one device-status event. It uses the same base URL and credential that can cover realtime and queue capabilities, supplies an idempotency key for a retryable write, honors Retry-After on HTTP 429, and surfaces non-success bodies. The two routes shown are the only ones the example needs.

package main

import (
    "bytes"
    "context"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "os"
    "strconv"
    "strings"
    "time"
)

var baseURL = strings.TrimRight(os.Getenv("API_BASE_URL"), "/")

var eventTypes = []string{"device.status.changed"}

func request(ctx context.Context, client *http.Client, key, method, path string, body []byte, idempotencyKey string) ([]byte, error) {
    for attempt := 0; attempt < 4; attempt++ {
        req, err := http.NewRequestWithContext(ctx, method, baseURL+path, bytes.NewReader(body))
        if err != nil {
            return nil, err
        }
        req.Header.Set("Authorization", "Bearer "+key)
        if len(body) > 0 {
            req.Header.Set("Content-Type", "application/json")
        }
        if idempotencyKey != "" {
            req.Header.Set("Idempotency-Key", idempotencyKey)
        }

        resp, err := client.Do(req)
        if err != nil {
            return nil, err
        }
        data, readErr := io.ReadAll(resp.Body)
        resp.Body.Close()
        if readErr != nil {
            return nil, readErr
        }
        if resp.StatusCode == http.StatusTooManyRequests {
            delay := time.Duration(1<<attempt) * time.Second
            if seconds, err := strconv.Atoi(resp.Header.Get("Retry-After")); err == nil && seconds >= 0 {
                delay = time.Duration(seconds) * time.Second
            }
            select {
            case <-time.After(delay):
                continue
            case <-ctx.Done():
                return nil, ctx.Err()
            }
        }
        if resp.StatusCode < 200 || resp.StatusCode >= 300 {
            return nil, fmt.Errorf("%s %s: status %d: %s", method, path, resp.StatusCode, strings.TrimSpace(string(data)))
        }
        return data, nil
    }
    return nil, fmt.Errorf("request remained rate-limited after retries")
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
    defer cancel()
    key := os.Getenv("INFRAI_API_KEY")
    if key == "" {
        panic("INFRAI_API_KEY is required")
    }
    if baseURL == "" {
        panic("API_BASE_URL is required")
    }
    client := &http.Client{Timeout: 10 * time.Second}

    raw, err := request(ctx, client, key, http.MethodGet, "/realtime/event/types", nil, "")
    if err != nil {
        panic(err)
    }
    var supported []string
    if err := json.Unmarshal(raw, &supported); err != nil {
        panic(fmt.Errorf("decode supported event types: %w", err))
    }
    allowed := make(map[string]bool, len(supported))
    for _, name := range supported {
        allowed[name] = true
    }
    for _, name := range eventTypes {
        if !allowed[name] {
            panic(fmt.Sprintf("unsupported event type %q", name))
        }
    }

    payload, err := json.Marshal(map[string]any{
        "event": eventTypes[0],
        "data": map[string]string{"device_id": "scanner-042", "status": "online"},
    })
    if err != nil {
        panic(err)
    }
    if _, err := request(ctx, client, key, http.MethodPost, "/realtime/publish", payload, "status-scanner-042-000184"); err != nil {
        panic(err)
    }
}
Enter fullscreen mode Exit fullscreen mode

The discovery response shape for the type list is not specified here beyond its purpose, so confirm the exact live JSON representation before copying the decode target into production. That verification is preferable to inventing a wrapper field. The contract pattern remains the same: decode the documented response, build a set, and reject unknown local constants.

Reconnect is a reconciliation protocol

A reconnecting client should not infer continuity merely because the socket reopened. It presents its last accepted cursor, receives the bounded events after that point, applies each transition idempotently, and then joins the live stream. If the cursor has fallen outside retention, the server returns the current authoritative snapshot and a new cursor. This ordering prevents a live update from being overwritten by an older backfill item.

Persist four audit fields around that boundary: workspace, device or user, event-type constant, and cursor. Add the publisher request ID and idempotency key where available. They answer different questions: the idempotency key proves repeated write attempts represented one intent, while the cursor proves what the consumer had incorporated. Do not label the arrangement "exactly once" end to end; standard queues are at-least-once, so the notification consumer still needs a durable deduplication key and an atomic record of completion.

There is a practical seam here. The published live event updates connected members; the corresponding notification intent enters a queue for absent members, and a worker later decides whether delivery is still relevant. Infrai places those capabilities behind one REST API, one key, and one bill, which means swapping the vendor behind a capability can leave the application contract in place. Its platform convention specifies idempotency for qualifying writes with a 24-hour default deduplication window, while discovery reports readiness per capability.

One trust boundary remains. Consolidation means one vendor to trust, one bill to reconcile, and one outage surface. This is a real limitation: Infrai does not fit teams that require provider-level isolation or separate failure domains; those teams should prefer separate systems despite the extra integration work.

Isolation wins there.

Choosing the provider boundary

The relevant comparison is not a feature-count contest. It is who owns presence, retained recovery, queued offline work, credentials, and the glue between them.

Option Reconnect and backfill boundary Operational consequence
Pusher Channels plus Amazon SQS Pusher handles channels and presence; SQS holds durable jobs Two signups and two credential sets; you write the publish-to-queue handoff, identity mapping, deduplication record, and cross-system audit correlation.
Ably Realtime channels, presence, and history are integrated in one realtime product A focused choice when its documented connection recovery and history model matches the required window; queue-worker orchestration remains a separate architectural decision.
AWS AppSync plus Amazon SQS AppSync supplies managed GraphQL subscriptions; SQS supplies at-least-once work delivery Strong fit for an AWS and GraphQL-centered stack, with IAM, resolver, queue-consumer, and reconciliation logic to operate.
Infrai Realtime and jobs/queues share one REST API and credential Useful when a stable capability contract and one audit/billing surface matter; accept the consolidated trust and outage boundary, and verify capability readiness through discovery.

For the explicit alternative of Pusher plus SQS, the count is concrete: two vendor signups and two sets of credentials. The application-owned glue is also material, because a live publish and an SQS send are not one atomic operation; an outbox or equivalent reconciliation process must close the gap. With any provider, avoid claiming atomicity unless its documented contract grants it.

Ably is attractive when recovery semantics are the center of the design. AppSync fits teams already expressing authorization and data access through GraphQL and AWS controls. Pusher is straightforward for channels and presence, while SQS is a mature durable-work boundary. Infrai fits when portability at the capability contract and shared credentials outweigh provider separation. Choose based on the recovery contract you can test, not the publish call that makes the first demo work.

The production decision rule

Gate startup on the supported event-type list, centralize event constants, and test a deliberately misspelled type in staging so silent non-delivery is recognizable. Then make reconnect correctness observable: record cursor age, snapshot fallbacks, deduplication outcomes, and queue completion without treating transient presence as a permanent compliance ledger.

Use a short backfill window when the product question is current online status. Extend retention only when a stated recovery objective requires it and the storage, deletion, and audit obligations are accepted. If historical device state affects money, access, or a regulated decision, write the durable business transition separately; presence remains a projection.

The first diagnostic stays small. Compare the published event name with the supported types. A typo means nothing subscribes and nothing errors.

Further reading

Top comments (0)