Skip to content

feat(queue): bulk-insert unmergeable drained events instead of one heartbeat each - #128

Open
TimeToBuildBob wants to merge 12 commits into
ActivityWatch:masterfrom
TimeToBuildBob:feat/batched-insert-queue-drain
Open

TimeToBuildBob wants to merge 12 commits into
ActivityWatch:masterfrom
TimeToBuildBob:feat/batched-insert-queue-drain

Conversation

@TimeToBuildBob

Copy link
Copy Markdown
Contributor

Second half of #32 (follow-up to #127). Stacked on #127 → #115: GitHub can't base a PR on a fork branch, so this targets #115's branch and the diff includes #127. Review only the top commit, 1e819f5.

Problem

#127 merges consecutive queued heartbeats, but events that can't merge (changing window titles) still cost one heartbeat request each. A backlog of 300 window changes drains in 300 requests.

Change

Each run of merged events for a bucket is sent as:

  1. First event → heartbeat. It can still merge into the server's existing last event, so batch boundaries don't fragment events.
  2. Middle events → POST /buckets/<id>/events, in chunks of 100.
  3. Last event → heartbeat. This leaves the server's cached last heartbeat on the true last event, so live heartbeats continue it.

The two hazards, and how they're handled

Stale heartbeat cache on aw-server (Python). create_events does not invalidate ServerAPI.last_event; aw-server-rust does invalidate. Suppose the last heartbeat has the same data as the stale cached first event and lands within pulsetime of it. It then merges, and replace_last() overwrites the newest inserted event. I reproduced this against aw-server v0.13.2 with the guard disabled. Queueing a,b,c,a at 1s spacing (pulsetime 10) stored [a(3.5s), a, b], so c was lost. This can only happen when the gap between the first event's end and the last event's start is ≤ pulsetime, so such runs stay heartbeat-only. With the guard, the same input stores [a, b, c, a]. (The server bug deserves its own one-line fix; this PR doesn't depend on it.)

Inserts are not idempotent. If an insert reaches the server but the response is lost (timeout, reset), a blind retry duplicates the chunk. So before a retried insert is sent, the client GETs the chunk's time range and resends only events the server doesn't already have. This costs one extra request, and only on retry. The same reconciliation runs for inserts in the first batch after startup, because a run that crashed before acknowledging that batch may already have delivered it.

Evidence

  • tests/test_requestqueue.py: 26 passed. New tests cover:
    • a 250-event unmergeable backlog drains in 5 requests, in order, each event exactly once;
    • a short run stays heartbeat-only;
    • lost-response retry (Timeout / ConnectionError) does not duplicate;
    • a refused insert is resent in full.
  • Mutation check: replacing the reconcile with data = request.data fails both lost-response tests, and removing the gap guard fails the short-run test.
  • End-to-end against real servers: 300 distinct-title heartbeats behind an existing event, then one live heartbeat. On aw-server-rust v0.14.0 and aw-server v0.13.2: 6 requests (3 heartbeats including the live one, 3 inserts) instead of 301; 300 events stored, in order, no duplicates; the first event merged into the pre-existing one (15s) and the live heartbeat extended the last.
  • mypy clean. ruff check flags only pre-existing issues in cli.py/queries.py, which this PR doesn't touch.

Refs #32

@greptile-apps

greptile-apps Bot commented Oct 2, 2026 •

Copy link
Copy Markdown

RetriggerConfidence Score: 3/5

[Medium risk] Changes how queued events are batched and sent.

The PR is not safe to merge while failed reconciliation can permanently lose or duplicate queued events.

Findings

  1. P1 Undelivered inserts are discarded ▶
  2. P1 Restart can duplicate inserts ▶

Summary

The PR coalesces queued heartbeats and bulk-inserts the middle events of longer unmergeable runs. The latest change bounds retries for an unreadable insert-reconciliation lookup, but its two fallback choices can lose or duplicate events when delivery status is unknown.

  • Adds bounded lookup retries and tests for persistent read denial.
  • Preserves queue progress at the cost of unsafe assumptions about whether an insert was stored.

Diagram

%%{init: {'theme': 'neutral'}}%%
flowchart TD
  A[Insert needs reconciliation] --> B{Lookup succeeds?}
  B -- Yes --> C[Send only missing events]
  B -- No, fewer than 3 failures --> D[Retry lookup]
  D --> B
  B -- No, 3 failures --> E{Attempted this run?}
  E -- Yes --> F[Skip insert: possible data loss]
  E -- No --> G[Send full insert: possible duplicates after restart]
Loading

Reviews (4) · Last reviewed commit: "fix(queue): bound the unverifiable-looku..."

Comment thread aw_client/client.py
Comment thread aw_client/client.py Outdated
@TimeToBuildBob

Copy link
Copy Markdown
Contributor Author

@greptileai review

Comment thread aw_client/client.py Outdated
@TimeToBuildBob

Copy link
Copy Markdown
Contributor Author

@greptileai review

Comment thread aw_client/client.py Outdated
Comment thread aw_client/client.py
@TimeToBuildBob

Copy link
Copy Markdown
Contributor Author

@greptileai review

Comment thread aw_client/client.py
Comment thread aw_client/client.py
@TimeToBuildBob

Copy link
Copy Markdown
Contributor Author

Review convergence: the remaining findings are mutually exclusive

Four review rounds have now flagged each branch of one decision. They cannot all be satisfied at once:

Round Finding Remedy demanded
1 Failed read discards pending inserts don't drop
2 Failed read duplicates inserts don't duplicate
3 Unreadable lookup blocks the queue don't block
4 Undelivered inserts are discarded and Restart can duplicate inserts don't drop and don't duplicate

Rounds 1–4 are consistent with each other only if the client can always determine whether a bulk insert was stored. It can't: POST /buckets/<id>/events is not idempotent, and the only way to resolve a lost response is GET on the same bucket. In the state every one of these findings is written about — the read fails permanently (a persistent 403/404) while writes still succeed — the read carries no information. At that point "don't lose" and "don't duplicate" are not two review findings; they are the two horns of a dilemma and no code can take both. This is the two-generals problem, not a defect in the fallback.

What the code does now (c58f543), for that degenerate state only:

  1. Retry the lookup up to MAX_RECONCILE_FAILURES (3) times with backoff, so a transient or slow lookup is never guessed at.
  2. If it still fails, make progress instead of blocking the run's last heartbeat and every later request.
  3. Choose the side that preserves the invariant that can be known:
    • Insert already attempted this run (its response was lost) → assume it landed; do not resend, because resending can only add events the server may already have. Logged at ERROR.
    • Insert not attempted this run (startup reconcile) → send it, because discarding it can only lose events the server may never have received.

Transient and connection failures — the realistic cases — never enter this path: they retry and reconcile normally, so in practice neither the loss nor the duplicate branch is reached.

What I'm asking for

The three properties are only jointly achievable if the ambiguity is removed, and removing it needs a server-side change (idempotency keys / dedup on create_events) or a reliable read. That's out of scope for this PR. If a maintainer wants the client to prefer one horn explicitly, name it and I'll pin it; otherwise I recommend merging at c58f543.

I'm deliberately not re-triggering review again — the rounds are oscillating between mutually exclusive demands, which is diminishing returns. Leaving this for maintainer adjudication.

@TimeToBuildBob

Copy link
Copy Markdown
Contributor Author

Restacked on carry PR #130 (#130). The base has been changed to master — once #130 merges, this PR's diff will show only the bulk-insert changes. All 31 tests pass on the rebased branch (7448ebe). Merge order: #130 → #127 → #128.

@TimeToBuildBob

Copy link
Copy Markdown
Contributor Author

Status at 7448ebe: the restack onto #130 only brought in master's profile changes. The queue logic is identical to c58f543, so Greptile's 3/5 still covers the current code. Its two P1s are the drop-vs-duplicate pair described in the convergence comment above. Both apply only when a reconcile GET fails permanently while writes keep succeeding.

I'm not pushing another fix round or re-triggering review. Taking one side reopens the other finding, and a real fix needs server-side idempotency for create_events, which is outside this PR. If you want the client pinned to "never drop" or "never duplicate", say which; otherwise this is ready for maintainer review after #130 → #127.

TimeToBuildBob and others added 12 commits October 10, 2026 14:09
When aw-server is unreachable the request queue accumulates one heartbeat per
commit interval. On reconnect they were dispatched one HTTP request at a time,
so a long outage drained slowly and hammered the server just as it came back
(issues ActivityWatch#32 and ActivityWatch#7).

RequestQueue now pops queued requests in batches and pre-merges consecutive
heartbeats for the same bucket client-side with aw_transform.heartbeat_merge
before sending. The merged events still go to the heartbeat endpoint, so the
server-side result is unchanged - a run of identical heartbeats that the
server would have merged into one event now arrives as one request instead of
dozens.

Requests are held in an in-memory batch until the whole batch has been
handled; a transient error retains and retries the batch, and re-sending an
already-delivered heartbeat is a no-op on the server. A non-retryable client
error drops only that request and keeps dispatching the rest.

Test fills the queue with 200 mergeable heartbeats and asserts the drain uses
a single request while the merged event still spans the whole range (no data
loss).

Git-Session-Id: b337
…, per-request retry

Greptile review feedback on ActivityWatch#127:

- **Order**: group by bucket, not by full endpoint, and only merge a
  contiguous run that shares one endpoint. Grouping by endpoint reordered a
  bucket's heartbeats when their pulsetimes differed (e.g. 10, 20, 10 was sent
  as 10, 10, 20), which changes the server-side timeline.
- **Malformed endpoint**: `pulsetime=1..2` matched the regex but `float()`
  raised inside coalescing, outside the dispatch error handler, killing the
  worker. Parsing is now fallible and a malformed pulsetime falls back to
  sending the request verbatim.
- **Partial retry**: dispatch one coalesced request per call and track the
  index, so a transient error retries only the failed request instead of
  replaying heartbeats that already reached the server.
- **Shutdown**: because one request is dispatched per call, `disconnect()`
  observes `stop()` between requests instead of blocking on a whole batch.

Tests: 21 passed (new: pulsetime-order preservation, malformed pulsetime,
partial-retry-does-not-replay).

Git-Session-Id: b337
…timing bound

Review feedback: the 10k test only checked the first request's duration, so it
would pass even if later batches lost heartbeats; assert the summed span across
all posted events instead. Remove the wall-clock limit, which made a functional
test sensitive to runner load.

Git-Session-Id: b337
…ge contiguously

Git-Session-Id: 3b002f8a-7967-5ade-8e69-6c86e374a155
Git-Session-Id: 7c62adb8-40f8-58a7-82c2-170ab3f2baaa
…vityWatch#125/ActivityWatch#130

- tests: pass tmp_path to _fresh_queue (ActivityWatch#130 made queues isolated per test)
- RequestQueue.__init__: keep _queue_write_failing (ActivityWatch#124) after the batch state
- wait_for_queue_empty (ActivityWatch#125): track the in-flight batch instead of _current
- drop an import left unused by the rebase

Git-Session-Id: 68f2e145-fcff-459d-8e16-c3e43269f5fa
…artbeat each

After merging, a drained backlog whose events cannot merge (e.g. changing
window titles) still cost one heartbeat request per event. Send the middle
of each run via POST /buckets/<id>/events in chunks of 100 instead.

- First and last event of a run still go to the heartbeat endpoint: the
  first so it merges into the server's existing last event, the last so the
  server's cached last heartbeat lands on the true last event.
- Runs that fit within pulsetime stay heartbeat-only: aw-server (Python)
  does not invalidate its last_event cache on insert, so the final heartbeat
  could merge into the stale cached event and overwrite inserted ones.
- An insert is not idempotent, so a retried insert (and inserts in the first
  batch after startup, which a crashed run may have delivered) first looks up
  the chunk's time range and only resends events the server lacks.

Refs ActivityWatch#32

Git-Session-Id: 3ed7
Address review: a 4xx on a bulk insert now resends the chunk one event per
request, so only the rejected event is dropped. A permanently failing
pre-retry lookup sends the chunk anyway (possible duplicate beats dropping
never-sent events); connection errors and retryable statuses still retry.

Git-Session-Id: 3ed7
…opping it

An unreadable pre-retry lookup cannot tell whether a bulk insert was already
stored. Resending the chunk then duplicates stored events; dropping it loses
events the server never got. Keep the request queued and retry the lookup
later, so neither happens.

Replaces the send-anyway fallback from f6dd56e. Tests: the duplicate case
fails on the previous head (it resent the stored chunk) and passes now; the
never-stored case stays queued and drains once the lookup succeeds.

Git-Session-Id: f6dd56e
… the queue

A permanently unreadable pre-retry lookup deferred the insert forever,
blocking the run's last heartbeat and every later request. Retry the lookup
up to MAX_RECONCILE_FAILURES times, then make progress without it:

- The insert was already attempted and its response lost: assume it landed.
  Resending could duplicate stored events and inserts are not idempotent, so
  drop the retried insert (logged at ERROR) rather than block the queue.
- Nothing was attempted this run (startup reconcile): the events may never
  have reached the server, so send the chunk rather than discard it.

This keeps both invariants the earlier rounds established: never duplicate,
never discard an unattempted insert, and now never block the queue.

Git-Session-Id: 0f3c
…ityWatch#125/ActivityWatch#130

- send an empty heartbeat payload instead of dropping it: the post guard
  should only skip an insert reconciled down to nothing (caught by ActivityWatch#125's
  test_wait_for_queue_empty_timeout)
- tests: pass tmp_path to _fresh_queue (ActivityWatch#130 made queues isolated per test)
- tests: silence mypy method-assign on the refuse-then-restore stubs
  (CI runs make typecheck)

Git-Session-Id: 68f2e145-fcff-459d-8e16-c3e43269f5fa
@TimeToBuildBob
TimeToBuildBob force-pushed the feat/batched-insert-queue-drain branch from 7448ebe to 9fe4e1f Compare October 10, 2026 14:15
@TimeToBuildBob

Copy link
Copy Markdown
Contributor Author

Rebased onto the updated #127 (which is now on master). #125's test_wait_for_queue_empty_timeout caught a real bug in this PR: the if data: post guard was meant to skip an insert reconciled down to [], but it also silently dropped any heartbeat whose payload was {}. It now skips only empty inserts (9fe4e1f). That commit also fixes the 4 mypy method-assign errors CI's make typecheck would have hit. Locally: 118 tests pass and mypy is clean.

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.

1 participant