Skip to content

Count tracked usage locally when DataStream is not connected - #134

Open
ryanechternacht wants to merge 1 commit into
mainfrom
ryan/replicator-track-offline
Open

ryanechternacht wants to merge 1 commit into
mainfrom
ryan/replicator-track-offline

Conversation

@ryanechternacht

@ryanechternacht ryanechternacht commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

Problem

In replicator mode the SDK polls the replicator's health endpoint and treats its ready field as "connected". When a customer's Schematic account is closed or Schematic is unreachable, the replicator stays up and keeps its Redis cache, but reports ready: false.

AsyncSchematic.track (and track_with_reservation) bump the cached company metric so the next flag check sees the usage before the server pushes the real figure. That bump in _update_company_metrics was gated on ds.is_connected(). Flag checks keep evaluating from the cache while the replicator is not ready, so with the gate the cached metric froze: tracked usage was never reflected locally and numeric limits stopped tripping for as long as the replicator stayed not ready.

Findings

Why the gate exists. It was ported from the Go SDK in the original DataStream port (#58), with no discussion of it in that PR's review. In Go it was added with the optimistic metric updates in SchematicHQ/schematic-go#73 (body.Company != nil && c.datastreamConnected), before replicator mode existed, and #82 later rewrote it as IsConnected(). Neither PR gives a reason for the check. It reads as a copy of the "is the stream up" guard used elsewhere, not a correctness requirement.

What the update does when not connected. DataStreamClient.update_company_metrics reads the company from the cache, adds the quantity to the metric whose event_subtype matches, and writes it back under the per-company lock. It never touches the network.

  • Missing cache: if the company is not cached it returns without writing, so nothing is invented for an unknown company.
  • Double counting: the bump is a local prediction of a figure the server owns. The server's next write for the company replaces the metric outright, whether that is a full message, a partial (metrics are upserted by key, replacing the value), or the replicator rewriting Redis once it is ready again. A bump would only add onto a server figure that already includes the same event if the server processed the event and pushed the company back before the bump ran. The bump runs directly after the enqueue and the buffer flushes on a timer, so that is not a realistic ordering, and it is the same ordering whether or not the stream is connected. The real double count risk, a retried track_with_reservation whose event the server drops on its idempotency key, is handled separately by settled_locally and is unchanged.
  • The track event is still enqueued and sent to the API exactly as before.

WebSocket mode. The gate is removed there too. A dropped socket does not clear the cache, and check_flag still answers from cached entities while disconnected, so the same frozen-metric problem applied. When the socket comes back, pushed company data replaces the bumped value.

Change

_update_company_metrics no longer checks is_connected(). It runs whenever DataStream is configured and the event names a company.

Test

  • New tests/custom/test_replicator_track.py builds an AsyncSchematic in replicator mode on a fake Redis, drives one health poll that returns {"ready": false}, and seeds a company at 95 of a 100 unit limit. check_flag allows, track(quantity=10) enqueues the event and moves the cached metric to 105, and the next check_flag denies without calling the API. A second test tracks for an uncached company and checks that nothing is written. The first test fails on main.
  • test_track_updates_company_metrics_when_datastream_not_connected covers a DataStream client that reports not connected (the WebSocket case) with a mock.
  • poetry run pytest -n auto .: 758 passed, 3 skipped. poetry run mypy .: clean. ruff check on the changed files: clean.

Related: #133 removes the same gate from check_flags. The two PRs are independent.

🤖 Generated with Claude Code

track bumps the cached company metric so the next flag check sees the
usage before the server pushes the real figure. The bump was gated on
is_connected(), which in replicator mode is the replicator's ready flag.
A replicator that reports ready: false (account closed, Schematic
unreachable) keeps its cache and flag checks keep evaluating from it, so
with the gate the cached metric froze and numeric limits stopped
tripping for as long as the replicator stayed not ready.

Drop the gate in both replicator and WebSocket mode. The bump only
touches a company already in the cache, and the server's next push for
that company replaces the metric outright, so it cannot double count.
The track event is still sent to the API as before.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants