Repository navigation
Replay buffered observational-memory results at production's step in Mastra memory replay - #1230
Conversation
PostgreSQL and LibSQL stores keep buffered observation chunks in a JSON column, so reading a record back returns each chunk's `createdAt` and `lastObservedAt` as ISO strings. The snapshot validator requires Dates, so every turn that started after a buffered observation was marked ineligible, and older generations keep their chunks, so the thread never recovered. `normalizeStoredMemoryDates` converts known date fields back to Dates before validation. It accepts only an exact ISO rendering, so any other string still fails validation.
Replay rebuilds memory in Mastra's InMemoryStore, which behaves unlike database stores in ways observational memory depends on. Its `createReflectionGeneration` fills a missing `lastObservedAt` with the current time where PostgreSQL and LibSQL copy it through, so after a replayed reflection Mastra treated the turn's messages as observed, skipped a recorded observer call and the replay diverged. It also hands out live record objects instead of copies, drops reflection metadata, keeps omitted observed-message ids, and activates caller-supplied buffered chunks. The baseline now records whether the source domain is an InMemoryStore. Unless it was, `createIsolatedMemoryReplay` gives the isolated store database record semantics for those methods.
Add `@mastra/pg` 1.25.0, the newest release compatible with `@mastra/core` 1.67.0, as a dev dependency and a test that records six turns on a real PostgreSQL store, including a buffered observation and a reflection. It requires every turn to be eligible and replays the reflection turn from the recorded tape. The test runs when `KITARU_TEST_MASTRA_POSTGRES_URL` is set and uses its own schema. The TypeScript CI job points it at its PostgreSQL service.
The baseline gives Mastra tape-instrumented observer and reflector models, and Mastra persists its resolved OM configuration into the record it initializes. The source store therefore received Kitaru's wrapper instead of the configured model id. A JSON store serialized the wrapped provider client, so a `ModelRouterLanguageModel` grew each record's config from about 4.7 KB to about 464 KB and pushed later reflection evidence over the replay item limit. InMemoryStore kept the wrapper object itself, which the snapshot codec rejects, so no turn after the first was replayable. The capture binding now hands `initializeObservationalMemory` the model values from the source configuration, and the wrapper serializes as that value if it reaches any other JSON path. Initial snapshots and every mutation that carries OM records now apply the identity projection that initialization evidence already used: models become identities and Mastra's built-in extractors become slugs. Stores that return live Extractor objects can now record later turns.
The source-thread lease marked a thread and its resource as permanently
unsafe whenever the adapter itself wrote outside an eligible lease: a
buffered reflection finishing after release, a reply sent while the
previous turn was still finalizing, a native fallback after a Kitaru
error, or one call without memory selectors. Every later turn for that
user then stayed unreplayable until restart, and forever under a shared
lease that keeps markers across process loss.
- A write outside an eligible lease now registers through
`acquire(selector, { waitMs: 0 })` for the length of the write. An
overlapping holder loses eligibility, and nothing outlives the write.
`markUnsafeWrite` is left for writes that cannot register.
- Native fallbacks register each write under the caller's or the
factory's default selectors. A call with no selectors marks nothing
unless it actually writes.
- Finalization joins Mastra's buffered observation and reflection
through `waitForBuffering` before the drain, bounded by the new
`finalizationWaitMs` (60 s). A turn that misses it is ineligible, and
the lease is released before the final session update.
- The process-local lease waits up to `waitMs` for a releasing holder.
- The OM tape tracks a call from its start, so a slow provider no
longer leaves an empty slot.
- A source `memory` option lets the application's `settled()` join
each recorded turn's memory work, and capture then waits only for the
recorded thread's buffering.
A native fallback whose thread or resource came from Mastra's reserved `mastra__threadId`/`mastra__resourceId` request-context keys was treated as selector-unknown, so its writes set the global unknown-writer marker and every thread of every user stopped being replayable. `getNativeSelector` now follows Mastra's precedence: the context keys override the merged default and caller memory options. A turn also kept its source lease while its memory evidence uploaded to Kitaru after the stream closed. With a slow Kitaru API, a reply sent shortly afterwards overlapped the lease and both turns became ineligible. The binding now releases once its storage writes settle and the eligibility check passes, and `drain()` waits for evidence after. The finalization deadline now bounds every wait on OM work: the tracked work that the source Memory's `settled()` joins, the join loop that kept polling a hung buffering operation, and replay, which could wait indefinitely on a buffering operation left hung on the same thread and now fails instead. The lease contract and docs now require a shared backend to end a lease whose holder process died, since Kitaru never renews one.
Replay matched recorded observer and reflector results by call order and count only. Mastra's number of OM calls depends on timing: a slow production observer merges buffering rounds that an instant replay makes separately, so an unchanged replay failed as `mastra_om_call_order`. - Match each replay call to an unused recorded call with the same phase, method and input fingerprint, else the next unused one of its phase. - Serve surplus calls without a provider: an empty observation for the observer, the last reflection for the reflector. Report input mismatches, surplus calls and unused results in an `om_call_divergence` span and `mastra_om_divergence` session metadata. - Fail closed only when a phase has no recorded result at all. - Record failed OM attempts, so a baseline whose observer succeeded after a retry stays eligible and replays the successful result. - Fingerprint OM inputs without message times, dates, ids and per-part creation timestamps, so `om_input_mismatch` reports only real input drift. - Close the session when a memory processor tripwire ends the run, since Mastra then calls neither `onFinish` nor `onError`.
Memory-mutation evidence was checked against the generic recorder bounds of 10,000 JSON values and 1 MiB, and request evidence against a fixed 1 MiB. `saveMessages` also recorded every new message twice, once as arguments and once as the storage result. A single tool result of 1,400 rows with 10 fields therefore made the turn ineligible, and a request carrying a normal PDF or a long reasoning history did the same. `createRequestCapture` also never received the wrapper's `recordingLimits`. Mutation and request evidence now use the replay input's budget of 16 MiB, 200,000 items and 64 levels per node. A saved message that storage returns unchanged is recorded as a `savedMessageRef` with its id and SHA-256. Evidence over the budget becomes a degraded marker that names the exceeded bound, and the node reports it through `evidence_truncated` or `request_evidence_truncated` with the exact reason. Truncation no longer makes the turn ineligible, because replay rebuilds memory and requests from the replay input rather than from these nodes. Unsafe values and credentials still do. Application-supplied `recordingLimits` now bound request evidence on this path. The tool-sized defaults are not applied there, because they would truncate nearly every prompt.
A signed download URL that reached thread history made every later turn's snapshot capture throw "contains URL credentials", so the whole thread became permanently unreplayable. Three different credential checks also disagreed: fragments, userinfo on other schemes, `key=`, `client_secret=`, `&`-escaped and percent-encoded nested URLs were stored while the session stayed eligible, and a credential-looking URL only in the model output failed an otherwise good baseline or replay. `redactUrlCredentials()` in `@zenml-io/kitaru/adapter` is now the one classifier. The Mastra codec, the evidence sanitizer and the strict replay check all use it: each URL keeps its text and only its credentials become `REDACTED`, pagination parameters stay as they are, and redaction never makes a recording ineligible. Declared files map to captured content references wherever a whole value equals a declared URL, including the initial thread history, so a processor that resolves attachment URLs from history gets recorded bytes on replay without the token. Substrings inside prompt text and working memory are no longer rewritten, declared URLs match in their `new URL(url).href` form, and `files` may be a function evaluated for each call.
A slow or hung Kitaru server delayed every recorded memory turn: each step's node upload was awaited inside `onStepFinish` with the client's 30 s timeout, session setup could hold the stream for 30 s, a declared file download without a timeout could stop the turn from answering, and the pre-turn capture waited on the source Memory's `settled()`, which joins background work on every thread. Step uploads now queue in order and run in the background. A baseline's finalization, including failures, runs detached and cancels any evidence upload still pending twice `finalizationWaitMs` after it starts, closing the session as `recording_flush_timeout`. Root span upserts that open and close the session are never cancelled. A baseline waits `sessionSetupWaitMs` (2 s) for its session and `fileCaptureWaitMs` (10 s) for declared files, which now download concurrently. When either runs out the turn answers natively; a late session is closed as `recording_setup_timeout` and a file timeout is reported as `file_capture_timeout`. The native fallback takes over the downloads capture started, so a file is fetched once per turn. Capture now joins only the recorded thread's buffering, which Mastra tracks per thread for the whole process.
Dependabot writes one version into every package.json that declares a dependency, and it cannot scope an update to one pnpm workspace member: `ignore` matches names and versions only, `dependency-type` classifies a name across all members, and a per-directory npm block inside the workspace fails as misconfigured. A `@mastra/*` update would therefore move the adapter's test install and every compat package's pin at once, as #1222 did. Ignore every version update of `@mastra/*` and point at the tested-version table, where new releases are added by hand after their suite passes. Listing the update types keeps security updates flowing.
Drop the typecheck config from compat packages without tests of their own, the copied engines field and version-specific descriptions. Build the expected rejection text from the table, and check that the README and docs version tables list exactly its rows.
With a slow production observer, Mastra runs one buffering round while the actor reads several tool results. An instant replay finished that round at once, so Mastra started new rounds production never ran. Each round seals the messages it covers, which splits later steps into more stored messages and raises OM's pending token count. A turn that ended just under `messageTokens` then crossed it only in replay, and the actor's last prompt differed from production's. Baselines now record how many actor steps had started when each OM result arrived and how long the call took. Replay returns a matching buffered result once the actor reaches that step, or after that duration at the latest, so Mastra sees the round still running at the same steps it did in production. The duration cap bounds any wait when a diverged replay blocks on the round early. Older tapes replay as before.
Move the envelope's field-by-field copy of OM tape entries into `toStoredOMEntry()` beside the entry type, so a new tape field cannot be dropped at the envelope again. Reuse `isCount()` and `getMastraVersion()` instead of new copies.
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 2a932df368
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| }; | ||
| // A replay that diverged may wait on this round at a step production | ||
| // never reached, so it waits no longer than production's observer did. | ||
| const timer = setTimeout(hold.release, maxMs); |
There was a problem hiding this comment.
Release held OM results before memory settlement times out
When a baseline observer's recorded durationMs exceeds finalizationWaitMs but most of that call overlapped a long actor turn, the baseline can remain eligible; if a replacement model then finishes in fewer steps than actorStepsAtResult, this timer holds the OM call for the full recorded duration. stateful-agent.ts calls await runtime.finish() at lines 1567–1578 before omTape.finish() can release the hold, so memory settlement reaches om_settle_timeout first and an otherwise valid replay fails. Release these holds when finalization begins, or bound them so they cannot outlive the settlement deadline.
Useful? React with 👍 / 👎.
| (entry.actorStepsAtResult === undefined || | ||
| isCount(entry.actorStepsAtResult)) && | ||
| (entry.durationMs === undefined || isCount(entry.durationMs)) |
There was a problem hiding this comment.
Validate new timing fields before marking sessions eligible
When a client supplies or replaces a v3 replay envelope containing an invalid optional value such as durationMs: "bad" or a negative actorStepsAtResult, the server's _mastra_om_tape_entry_well_formed() validator (src/kitaru/server/domain/session.py, lines 331–360) ignores these fields and can finalize the session as replay-eligible. These new checks then mark the entire tape malformed, causing the adapter to fail only during replay execution. Align server ingestion with these constraints, or treat invalid optional timing fields as absent.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: a65528412a
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| actorStepsAtResult: actorSteps, | ||
| durationMs: Math.round(performance.now() - startedAt), |
There was a problem hiding this comment.
Include failed retry time in buffered holds
When a buffered OM provider attempt runs for a while, fails, and Mastra retries successfully, groupRecordedCalls() deliberately collapses both attempts into one logical call, but this records durationMs only from the successful attempt. Replay skips the failed attempt and uses that shorter duration as the hold fallback; if the replay never reaches actorStepsAtResult (for example, a replacement actor finishes in fewer steps), it can release the result earlier than the production round completed and start buffering rounds that production suppressed. Measure the logical call from its first attempt or otherwise include failed-attempt and retry time.
Useful? React with 👍 / 👎.
# Conflicts: # src/kitaru/mcp/data/docs_index.json
# Conflicts: # src/kitaru/mcp/data/docs_index.json
Stacked on #1229 (
ci/mastra-memory-tested-range). Retarget todeveloponce that lands.What this fixes
A Mastra memory replay (
createMemoryReplayAgent()) could observe at a step where the recorded turn did not, so the actor's last prompt in replay differed from the baseline.The scenario: a production observer is slow (the test fixture delays it 300 ms). Mastra starts one buffered observation round after the first tool result, and while that round is still running the actor reads three more tool results. Mastra starts no new round while one is running, so production ran exactly one.
In replay, Kitaru serves that observer result from the recorded tape instantly. The round is finished by the next step, so Mastra starts a new round after nearly every step. The tape skips those extra rounds (they have no recorded result), but each round still seals the messages it covers. Sealing makes Mastra store the next step's content as a new message instead of appending to the old one, and each extra stored message adds its per-message overhead to OM's pending token count. Replay ran about +11, +21, +19 tokens over the baseline at steps 2 to 4.
Before the last step, Mastra compares pending tokens with
messageTokens(1200 in the fixture). A baseline that ended within about 20 tokens under it did not observe; the replay crossed it, made a blocking observation, and the actor sawraw= observed=1,2,3,4,5instead of the recordedraw=2,3,4,5 observed=1.(An earlier reading blamed the
data-om-buffering-startmarkers. Mastra's token counter scores everydata-*part as 0, so the cost is the message split from sealing, not the marker.)The change
Replay now reproduces when a buffered result arrived, measured in actor steps, not only what it was.
actorStepsAtResult(how many actor steps had started when the result arrived) anddurationMs(how long the call took). The step counter ticks in Kitaru's per-step request processor, just before each actor model call.durationMs. Mastra can wait on a running round mid-turn (activation waits up to 60 s). In a replay that diverged, for example after a model swap, it could wait on a round production never waited on at that step. With the cap, the worst case is waiting as long as production's own observer took.toStoredOMEntry(), next to the entry type. The previous field-by-field copy instateful-agent.tssilently dropped the new fields on the first attempt.Considered and rejected:
Reviewer Notes
The behavior lives in
packages/mastra/src/om-result-tape.ts:waitForActorSteps()/beginActorStep()hold a matched buffered result. Holding is only safe if the gate is always reachable. The argument: if Mastra waited on the round at step k in production, the result arrived while k steps had started, so the recorded gate is k, and a replay that follows the same steps reaches it before Mastra waits. Only a diverged replay can wait on a closed gate, which is what thedurationMscap is for. Please check this reasoning.serve()holds only buffered calls with a matching fingerprint. Blocking calls never wait (the actor is already waiting on them), and skipped buffered calls still end at once.finish()releases every hold before waiting for pending calls.packages/mastra/src/stateful-agent.tsadds one line,omTape.beginActorStep(), in the request processor. It runs once per step in both recording and replay.Known limits, not addressed here:
Tests
test/om-replay-tolerance.test.ts: the slow-observer scenario is now a shared helper, run at the comfortable size (34) and at a calibrated just-under-threshold size per Mastra core minor (1.67 → 29, 1.68–1.71 → 30). No single size sits just under the threshold on every version: at 29, 1.68 and later record a different baseline. The just-under test also asserts that the baseline stayed under the threshold. A new compat version without a table row fails with a calibration message. Both tests now expect nomastra_om_divergenceat all; before this fix, even the size-34 replay loggedinput_mismatches: 1.test/om-result-tape.test.ts: covers the recorded step count, the hold until the step, and the time cap (fake timers).Reproduction
With Node 22:
pnpm --filter @zenml-io/kitaru-mastra exec vitest run test/om-replay-tolerance.test.tsBoth pass. To see the bug, revert
packages/mastra/src/to the base branch and rerun: both slow-observer tests fail on every version, and the just-under test's last actor prompt israw= observed=1,2,3,4,5instead ofraw=2,3,4,5 observed=1.OM_DEBUG=1writesom-debug.login the working directory. Its[OM:status]lines show replay's pending tokens climbing above the baseline's at steps 2 to 4 without the fix and matching with it.Local checks run: full
packages/mastrasuite on 1.67 and the four compat suites (1.68–1.71), typecheck, lint, build, andjust check. The Postgres replay test was skipped locally (it needsKITARU_TEST_MASTRA_POSTGRES_URL); CI runs it.Current integration validation
The branch includes latest
developand the updated #1229 parent. Conflict resolution retains both buffered-call timing validation and develop's recorded-output presence/stream-shape validation. The generated MCP documentation index was refreshed for the merged docs.just checkandpnpm run build && pnpm run test:builtpassed on Node 26.10.0 (3,635 passed, 10 PostgreSQL-dependent tests skipped without a database). Fresh CI verifies the supported Node and PostgreSQL matrices.