DEV Community

RaffertyBarrett4726
RaffertyBarrett4726

Posted on

Node.js Waiting Room Queue — Live Position Updates and Backfill

The page says admission_lag_seconds is rising. The marketplace support queue still looks alive, but buyers who reconnect see an old place in line and some agents open video rooms for customers who are no longer next. The on-call question is not "is WebSocket up?" It is "can the database state be reconstructed at every waiting browser?"

TL;DR: publish the complete ordered queue whenever authoritative state changes, then let each browser find its own position. Keep the database authoritative, treat the realtime channel as a disposable view, and backfill that view after reconnect. For a Node.js service that also needs scoped video-room tokens, this is the least complex design that survives a lost connection without turning each customer's position into a separate stream.

The important alert should fire before customers report a frozen number: compare the database queue revision with the last revision acknowledged by the publication path. A gap is actionable. A raw connection count usually is not.

How should waiting room queue position updates go live?

Work backward from the symptom. Suppose the database has queue revision 8842, the last published revision is 8839, and the oldest unpublished change is 47 seconds old. Those are example values, not measured service guarantees, but they identify the failure boundary: three authoritative transitions have not reached the channel view.

The queue row should carry a stable entry ID, enqueue time, state, and monotonic revision. On every enqueue, cancellation, timeout, or admission, commit the database change first. A publisher then emits one ordered snapshot containing the revision and stable entry IDs. Waiting browser buyer_73 scans that snapshot, finds its ID, and renders the zero-based index as a human position. If the ID is absent, the client asks the application for fresh authoritative state instead of guessing. A tempting first design is to store and emit a position on every row, but that position becomes wrong whenever anything ahead of it leaves. The order is the fact; position is a projection.

Stop there.

Publishing one queue snapshot beats publishing one position per person. A removal near the front changes every later position. Per-person fan-out turns one cancellation into many writes and creates partially updated views; a single revisioned snapshot gives every subscriber the same ordering fact. The payload grows with queue length, so a very large queue may need a windowed snapshot or a different partitioning scheme. Do not quietly adopt those optimizations before the ordinary snapshot becomes a demonstrated bottleneck.

The admission transaction is also where the marketplace workflow crosses into video. Once an entry becomes eligible, persist an admission record, create the video room, and issue a token scoped to that room and participant. The token is an authorization artifact, not evidence that the customer still owns a queue position. Reconnect logic checks the admission record and requests a fresh scoped token through the application; it does not replay an old token from channel history.

WebRTC handles the media session. It does not make the waiting queue authoritative, and the signaling or data channel should not be asked to do that job.

Reconnect is a state transfer, not a replay trick

A reconnecting client needs a snapshot plus a revision, not every transient position it missed. On connection, fetch or receive the latest queue view, replace local state, and resume updates only from a newer revision. Reject duplicates and older revisions. This makes duplicate delivery boring, which is exactly what a runbook wants.

There is a narrow race between snapshot acquisition and subscription. Close it with a cursor or by subscribing first, buffering revisions, applying the snapshot, then draining only buffered revisions greater than the snapshot revision. The exact transport API varies, but the invariant does not: after recovery, the rendered list must correspond to one database revision.

Do not acknowledge publication merely because an application handler started. Record the last successfully published revision. Then expose at least these service-owned signals:

  • current database queue revision;
  • last successfully published revision;
  • age of the oldest unpublished revision;
  • reconnect backfill success and failure counts;
  • count of clients reporting a revision older than the current snapshot.

Those measurements describe correctness. Connection totals and message rates remain useful capacity signals, but they cannot prove that buyer_73 sees the right place in line.

Put the metric on the same path as the repair signal

An alert that only reaches the on-call dashboard leaves customers staring at stale state. A useful instrumentation change queries the metric and publishes the resulting document to an operations channel, where an internal console can show the same revision gap that paged the engineer. Infrai is one option here because both operations use a plain REST API, so there is no client SDK version to coordinate with the Node.js application. Its public discovery surface requires no key and returns the full request JSON Schema, allowing deployment to validate the publish template before the bridge starts instead of maintaining a second hand-written contract.

Infrai also puts 295 routes across 20 modules under one API key and one bill. For this workflow, that means realtime publication, metrics, and room-token issuance share one credential rotation policy instead of opening separate messaging, monitoring, and video credential lifecycles.

The runnable Go bridge below deliberately accepts PUBLISH_TEMPLATE_JSON from deployment configuration. That template must be validated against the live public discovery schema and contain the JSON string __METRICS_JSON__ at the location assigned to the event payload. This keeps the example honest: the supplied interface does not establish publish-body field names, so the code does not invent them. The bridge queries the verified metrics route, injects that exact JSON response, then publishes it with an idempotency key. It retries rate limits, honors Retry-After, and surfaces other response bodies.

package main

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

func call(ctx context.Context, client *http.Client, baseURL, key, method, path string, body []byte, idem string) ([]byte, error) {
    for attempt := 0; attempt < 5; 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 idem != "" {
            req.Header.Set("Idempotency-Key", idem)
        }

        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 >= 200 && resp.StatusCode < 300 {
            return data, nil
        }
        if resp.StatusCode != http.StatusTooManyRequests {
            return nil, fmt.Errorf("%s %s: status %d: %s", method, path, resp.StatusCode, data)
        }

        delay := time.Second << attempt
        if seconds, err := strconv.Atoi(resp.Header.Get("Retry-After")); err == nil && seconds >= 0 {
            delay = time.Duration(seconds) * time.Second
        }
        select {
        case <-ctx.Done():
            return nil, ctx.Err()
        case <-time.After(delay):
        }
    }
    return nil, fmt.Errorf("rate limit persisted after retries")
}

func main() {
    key := os.Getenv("INFRAI_API_KEY")
    baseURL := strings.TrimRight(os.Getenv("INFRAI_BASE_URL"), "/")
    template := os.Getenv("PUBLISH_TEMPLATE_JSON")
    if key == "" || baseURL == "" || template == "" {
        panic("INFRAI_API_KEY, INFRAI_BASE_URL, and PUBLISH_TEMPLATE_JSON are required")
    }
    if strings.Count(template, `"__METRICS_JSON__"`) != 1 {
        panic(`template must contain exactly one "__METRICS_JSON__" JSON string`)
    }

    ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
    defer cancel()
    client := &http.Client{Timeout: 15 * time.Second}

    metrics, err := call(ctx, client, baseURL, key, http.MethodGet, "/metrics/query", nil, "")
    if err != nil {
        panic(err)
    }
    if !json.Valid(metrics) {
        panic("metrics response was not valid JSON")
    }
    publishBody := []byte(strings.Replace(template, `"__METRICS_JSON__"`, string(metrics), 1))
    if !json.Valid(publishBody) {
        panic("rendered publish body was not valid JSON")
    }

    digest := sha256.Sum256(publishBody)
    idempotencyKey := "queue-metric-" + hex.EncodeToString(digest[:])
    result, err := call(ctx, client, baseURL, key, http.MethodPost, "/realtime/publish", publishBody, idempotencyKey)
    if err != nil {
        panic(err)
    }
    fmt.Println(string(result))
}
Enter fullscreen mode Exit fullscreen mode

In a Node.js system, run this logic as a small worker or translate the same HTTP contract directly; the operational point is independent of language. The application commits queue state, observability exposes the gap, and realtime distributes that metric and the queue view. One credential crosses the handoff.

There is a cost. Combining the paths means one vendor to trust, one bill, and one outage surface. Keep database queue transitions independent of channel availability, and retain a reconciliation worker that republishes the newest revision after service recovery.

Choosing the managed transport without hiding the trade-offs

The sensible comparison is recovery behavior and operational ownership, not a feature-count contest.

Option What fits this queue What the team still owns
Unified REST provider Plain REST calls cover metrics and realtime behind one key; useful when the Node.js service should avoid another SDK lifecycle Database authority, snapshot schema, revision reconciliation, browser recovery, and dependence on one combined provider
Pusher Channels Mature channel-oriented product with documented presence and client events A separate observability account, two credential sets with Datadog, and glue that queries metrics then republishes the result
Ably Realtime messaging with documented connection-state recovery and history features Queue semantics and database reconciliation; a separate metrics integration if Datadog remains the monitor
PubNub Publish/subscribe plus documented message persistence and presence Application-level ordering meaning, authoritative backfill, and observability handoff
Self-hosted NATS Direct control over subjects, retention choices, and deployment topology Broker operation, browser-facing gateway and authentication, metrics pipeline, upgrades, and on-call capacity

The common alternative here is Datadog plus Pusher. It requires two signups, two sets of credentials, and glue code that queries Datadog, transforms the result, and publishes through Pusher. That separation can be desirable: a realtime-provider incident need not erase the monitoring view, and each vendor can be replaced independently. It also gives the on-call two control planes to inspect during one queue incident.

Ably's recovery model deserves attention when continuity across brief disconnects matters more than minimizing integration surfaces. PubNub is a reasonable candidate when its persistence and presence model already matches the rest of the application. NATS is attractive for a team that already operates it and needs control over retention or topology. None of them removes the core responsibility: the marketplace database decides who is next and who may receive a scoped room token.

Run a failure drill before selecting. Disconnect a browser between revisions, cancel an entry near the head, publish the same revision twice, and make the channel unavailable while admitting a customer. The pass condition is a converged browser view and exactly one durable admission decision, not a pretty throughput chart.

Thresholds can page the team into making things worse

Alert on sustained divergence, but derive the threshold from the service's stated queue update objective and measured normal publication delay. A threshold of one missing revision for one second may page on harmless scheduling jitter. A threshold longer than the customer's patience detects the incident after support tickets arrive.

False positives have a direct cost. An on-call engineer may restart a healthy publisher, trigger a needless full backfill, or delay work on a real video admission problem. I have been paged by missed jobs and duplicate deliveries; both failure classes taught the same operational lesson: an alert needs a durable unit of work and a deduplication boundary, or the response can compound the incident. Start with a ticket-level signal while collecting the normal distribution, then promote it to a page only when the revision gap and age together predict customer-visible staleness. Document the suppression and recovery conditions in the runbook, including the exact revision evidence required before anyone restarts a publisher.

Noise burns judgment.

The final decision rule is short: choose the transport whose reconnect contract you can test, but keep queue truth and admission idempotency in your database. Publish one revisioned queue view on change. Backfill on reconnect. Page on durable divergence, not motion.

Further reading

Top comments (0)