SIGN IN SIGN UP

perf(stream): stop per-session stream groups and bound the viewer backlog (#1420)

* perf(stream): stop writing per-session mem-live groups, bound the viewer group

The engine's file-based stream store rewrites a whole (stream, group)
scope on every dirty flush tick, the same whole-scope-rewrite cost as
state. observe.ts and compress.ts wrote every raw and compressed
observation into a per-session mem-live group (group id = sessionId)
in addition to the shared viewer group, but grepping src/, plugin/,
integrations/, packages/ and the viewer for stream::get, stream::list,
stream::delete or any WS subscribe targeting a session-id group turned
up zero reads anywhere: the viewer only ever joins stream/mem-live/viewer.
That write was pure waste, doubling the persisted cost of every
observation for nothing, so this drops it and removes the now-dead
STREAM.group helper.

The remaining shared "viewer" group is written on every observation via
one persisted stream::set call (the synthetic/non-auto-compress path in
observe.ts; the stream::send calls elsewhere only publish to live
subscribers and never touch storage, confirmed against the engine's
stream::send handler) and was never pruned, so a long-running daemon
accumulates one entry per observation forever and every viewer
join/reconnect replays the whole backlog. src/state/viewer-stream.ts
tracks written item ids in memory and, every 50 writes, deletes the
oldest items beyond AGENTMEMORY_VIEWER_STREAM_MAX (default 500) via
stream::delete, batched rather than per-write. Delete failures are
logged and swallowed (Promise.allSettled) so pruning never blocks
observation ingestion.

Did not implement a one-time boot cleanup of pre-existing per-session
groups: the engine has stream::list_groups to enumerate groups cheaply,
but stream::list/get_group returns only stored values, never item ids,
and stream::delete requires an item id per call with no group-level or
bulk delete anywhere in the stream or state workers. Reconstructing ids
from nested value fields across every historical session group would be
fragile and, at the reported scale (600+ sessions, up to 500 items
each), not cheap - the opposite of what this fix is for. Per the task's
own conditional, reporting instead of building it.

* perf(stream): trim the stored viewer backlog after restarts

The prune tracker only knew about items written by the running
process, so every item stored before a restart stayed in the viewer
group for good, and that is the backlog existing stores already have.
At boot the tracker is now seeded from one listing of the viewer group,
oldest first, and anything beyond the cap is deleted in batches of 100
in the background.

Pruning no longer holds up the observation that triggered it.

* chore: credit co-authors

The persisted per-session stream groups were measured and reported in
discussion #1358.

Co-authored-by: Vladimir <7998636+MarvinFS@users.noreply.github.com>

* fix(viewer-stream): retry failed prunes and reject malformed stream cap

Keep a failed stream::delete id tracked instead of dropping it, so the
next prune retries it rather than letting the backlog grow past the
limit forever. Recheck for overflow once an in-progress prune finishes
if writes arrived while it was busy, instead of losing that check when
the write counter resets. Seed and requeue large id lists in bounded
chunks instead of spreading them into function arguments, which could
throw once a stored backlog got large enough. Reject a non-empty
AGENTMEMORY_VIEWER_STREAM_MAX value with trailing junk (e.g.
"500000junk") and fall back to the documented default instead of
silently taking the numeric prefix.

* perf(viewer-stream): delete overflow newest-first, seed in listing order

The engine's stream store removes items with IndexMap::shift_remove, whose
cost is proportional to how many entries sit after the removed one. Pruning
took the oldest tracked ids first and deleted them front to back, so each
delete shifted almost the whole remaining store; a large boot backlog made
every stream::set queue behind a prune that ran quadratic in the store size.
Deleting the same oldest-ids selection but issuing the deletes newest-of-the-
overflow-first means each removal only shifts the retained tail, bounded by
the configured cap instead of the store size.

The boot seed sorted the listed backlog by the hook's client-supplied
timestamp before choosing what counts as oldest. An unparseable or skewed
timestamp could make the newest item look oldest and get pruned first. The
engine already returns the group in insertion order, so the seed now trusts
that order directly and never parses timestamps.

AGENTMEMORY_VIEWER_STREAM_MAX accepted any integer, including negative
values that emptied the tracker (and deleted every stored item) on the next
prune, and values under 200 that silently shrank the backlog below what the
viewer's join sync displays. The getter now rejects negative values (falls
back to the default) and floors any other value at 200.

The boot stream::list call had no timeout, so an oversized backlog — or a
list response that exceeds the SDK socket's payload limit and closes the
connection without rejecting the pending call — could hang the seed
indefinitely during startup. It now carries a 5s timeoutMs; a failure of any
kind is already treated as "skip seeding this boot".

trackViewerStreamItem and the seed now share one id set, so a write that
races the boot seed for the same id is tracked once instead of twice.

The PR description's claim that the old per-session groups "would need ids
the stored values do not carry reliably" is wrong: both writers set
item_id = obsId with observation.id === obsId, so the ids are reliable.
Leaving those groups alone is a scope decision, not an ids problem.

---------

Co-authored-by: Vladimir <7998636+MarvinFS@users.noreply.github.com>
R
Rohit Ghumare committed
92d4500ff3eda8de3e699b6af60f39a136a6ef98
Parent: c314c7b
Committed by GitHub <noreply@github.com> on 9/28/2026, 2:53:56 PM