Repository navigation
Support scaleout session deletion - #1774
malatewang wants to merge 7 commits into
Conversation
a8dc213 to
f9da45a
Compare
marvinyu-memverge
left a comment
There was a problem hiding this comment.
Reviewed at f9da45a. The lease + recheck shape looks right to me; one ask.
stop()can hang if it lands while the worker is in its idle scan. The scan runs inside the worker task, so ifstop()queues itsNonesentinel while_queue_pending_session_deletionsis awaitingget_sessions_by_status, the scan's keys go into the queue behind the sentinel. The worker takesNoneand returns, those items never gettask_done(), andself._deletion_queue.join()never returns, soresources.close()isn't reached either. It needs at least oneDeletedrow at that moment (a peer mid-delete, or a deletion that keeps failing) and the window is one DB round trip per 30 s of idle, so it's rare, but when it hits, shutdown hangs until the process gets killed. I reproduced it with mocks against this head:stop()times out with the worker task done and one unfinished item; with the scan returning no keys it returns normally. Checkingif not self._started: returnbetween the query and theput_nowaitloop would close it (no await in between, so no race), plus a test along the lines of the repro.
Observation, not an ask: Fixes #1749 will close the issue on merge. This covers the boot fan-out and picking up a crashed replica's work, and the 30 s rescan also retries the deletions #1577 describes as stuck until restart. It doesn't add the durable job row, backoff, or terminal failed state from the issue's Expected section, so if those are planned as a follow-up, Part of #1749 would keep it open.
What I checked: the lock is a 30 s lease that #1722's service renews every 10 s and cancels the body with LeaseLostError if renewal can't keep up, which the worker's except Exception logs and the next scan retries. wait_timeout=timedelta(0) collapses the boot re-enqueue on N replicas to one deleter per session. The service shares the session DB engine, and the lease store supports both relational providers the config allows (PostgreSQL, SQLite). init_global_memory always passes key_to_session to start(), so the idle scan is live in the server.
86de8c3 to
5714319
Compare
5714319 to
a095d3f
Compare
edwinyyyu
left a comment
There was a problem hiding this comment.
Reviewed at a095d3f, against the merged lease lock from #1722.
Blocking: none.
The stop() hang from the earlier review is fixed: the _started check sits after the scan's last await and before the enqueue loop. With a real SQLite-backed SQLLeaseLockService shared by two MemMachine instances, only one instance runs the store deletes for a pending session, and the session row is deleted once.
Note, pre-existing (#1655). An episodic memory instance checks the session row only when it is opened. With session_manager.instance_cache_size at its default of 0, that happens on every request, so new requests see the deletion. A request already running on another replica keeps its instance, and so does an idle cached instance when the cache size is above 0, for up to max_life_time. Below the instance, only the segment store rejects writes from a stale handle (SegmentStorePartitionHandleStaleError), and the event backend writes segments before vectors. The declarative backend, the episode store, and semantic memory have no liveness check, and add_episodes writes the episode store in parallel with the episodic status check. A write in flight when the deletion finishes can therefore land under the old session key and appear in a project re-created with the same org and project IDs. The deletion itself is correct for the data present when it runs; this is the cross-replica session lifecycle gap in #1655, noted here because scale-out deletion relies on it.
Non-blocking, optional. Five inline comments. Also:
- When the process dies after the session row is deleted and before the lease is released, the
session-deletion:<key>lease row stays, because that key is never acquired again. The lease service's docstring asks workloads with one-time keys for periodic expiry cleanup. - Every grant goes through the lease counter row, and on each scan every idle replica loads all
Deletedsessions and probes one lease per key. - The scan runs only after
DELETION_IDLE_SECONDSof empty queue, so a replica that receives a steady stream of deletions never picks up work left by another replica. - Replicas that first use the lease service at the same time race on table creation (#1570). The losing replica's deletion fails and is retried on the next scan.
stop()processes everything queued ahead of its sentinel, so a stop right after a scan works through the whole backlog first.
🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu.
| if self._conf.semantic_memory.enabled: | ||
| tasks.append(self._delete_session_semantic_memory(session)) | ||
|
|
||
| await asyncio.gather(*tasks) |
There was a problem hiding this comment.
Non-blocking, optional. When one of these deletes raises, gather re-raises at once and leaves the others running. The async with then exits and releases the lease while they continue outside it. Reproduced with a real SQLite-backed SQLLeaseLockService: the episodic delete raised SessionInUseError, the lease row was already gone while the episode-store delete was still running, and a second instance took the lease and ran the same delete alongside it. The deletes are idempotent, so the usual cost is duplicate work, and the leftover tasks also outlive stop()'s drain. asyncio.TaskGroup cancels and awaits the remaining deletes before the block exits; the worker's except Exception still catches the ExceptionGroup it raises.
🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu.
| async def _delete_queued_session(self, session: SessionData) -> None: | ||
| lock_service = await self._resources.get_sql_lock_service() | ||
| try: | ||
| async with lock_service.lock( |
There was a problem hiding this comment.
Non-blocking, optional (needs a stalled process). The lease service asks consumers to check the fencing token at side-effect boundaries, and the store deletes here address data by session key. A deleter that stalls longer than SESSION_DELETE_LEASE_DURATION, from a process pause or a blocked event loop, can resume after another replica finished the deletion and a project with the same org and project IDs was created again, and then delete the new project's data. Reasoned from the code, not reproduced. A comment stating the limit covers it for now; the full fix deletes by incarnation, as the segment store's partitions do.
🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu.
| ) | ||
| except TimeoutError: | ||
| try: | ||
| await self._queue_pending_session_deletions(from_worker=True) |
There was a problem hiding this comment.
Non-blocking, optional. A deletion that keeps failing is retried on every scan by whichever replica wins the lease, with a traceback each time and no backoff or terminal state. A session that stays open on the deleting replica raises SessionInUseError and logs a traceback every DELETION_IDLE_SECONDS until it closes. Retrying is an improvement over dropping the job until restart; backoff and a failed state fit the rest of #1749.
🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu.
| """Only a pending session is deleted, and one instance owns it at a time.""" | ||
| lease = asyncio.Lock() | ||
|
|
||
| class SharedLockService: |
There was a problem hiding this comment.
Non-blocking, optional. Every test in this PR uses a fake for the lease lock. A test that builds SQLLeaseLockService on a temporary SQLite file and shares it between two MemMachine instances takes about twenty lines and exercises the real acquire, contention, and release path.
🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu.
| assert cleanup_mock.call_count == 3 | ||
| assert cleanup_mock.call_args_list[0].args[0] == all_ids[:batch_size] | ||
| session_manager.delete_session.assert_awaited_once_with(session_key="test-session") | ||
| resources.get_semantic_service.assert_not_awaited() |
There was a problem hiding this comment.
Non-blocking, optional. This test now runs with semantic memory disabled only, and the per-batch _cleanup_semantic_history assertions were removed, so no unit test covers the enabled branch of the new check in _delete_session_episode_store. The integration test above covers the end-to-end outcome.
🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu.
|
@marvinyu-memverge The shutdown race you reported is addressed in a095d3f: |
|
A follow-up question on the deletion queue. This is non-blocking and optional. With the database scan in place, a session's If not, an
The catch is that every deletion would then come from the scan, which only has keys, so 🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu. |
5abc48c to
68477f9
Compare
|
Should it be in |
Purpose of the change
Support session deletion in scaleout
Description
This change supports session deletion in scale out. The background deletion worker acquires the lock before deletion the session so only one MemMachine perform the actual deletion job. The background job also checks the database periodically so any dead jobs will be picked up by a running instance.
Fixes/Closes
Part of #1749
Fixes #1575
Type of change
[Please delete options that are not relevant.]
How Has This Been Tested?
Please describe the tests that you ran to verify your changes. Provide instructions so we can reproduce. Please also list any relevant details for your test configuration.
[Please delete options that are not relevant.]
Test Results: [Attach logs, screenshots, or relevant output]
Checklist
[Please delete options that are not relevant.]
Maintainer Checklist
Screenshots/Gifs
[If applicable, add screenshots or GIFs that show the changes in action. This is especially helpful for API responses. Otherwise, delete this section or type "N/A".]
Further comments
[Add any other relevant information here, such as potential side effects, future considerations, or any specific questions for the reviewer. Otherwise, type "None".]