Skip to content

Support scaleout session deletion - #1774

Open
malatewang wants to merge 7 commits into
mainfrom
scaleout_deletion
Open

malatewang wants to merge 7 commits into
mainfrom
scaleout_deletion

Conversation

@malatewang

@malatewang malatewang commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

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.]

  • Bug fix (non-breaking change which fixes an issue)
  • New feature (non-breaking change which adds functionality)
  • Breaking change (fix or feature that would cause existing functionality to not work as expected)
  • Refactor (does not change functionality, e.g., code style improvements, linting)
  • Documentation update
  • Project Maintenance (updates to build scripts, CI, etc., that do not affect the main project)
  • Security (improves security without changing functionality)

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.]

  • Unit Test
  • Integration Test
  • End-to-end Test
  • Test Script (please provide)
  • Manual verification (list step-by-step instructions)

Test Results: [Attach logs, screenshots, or relevant output]

Checklist

[Please delete options that are not relevant.]

  • I have signed the commit(s) within this pull request
  • My code follows the style guidelines of this project (See STYLE_GUIDE.md)
  • I have performed a self-review of my own code
  • I have commented my code
  • I have made corresponding changes to the documentation
  • My changes generate no new warnings
  • I have added unit tests that prove my fix is effective or that my feature works
  • New and existing unit tests pass locally with my changes
  • Any dependent changes have been merged and published in downstream modules
  • I have checked my code and corrected any misspellings

Maintainer Checklist

  • Confirmed all checks passed
  • Contributor has signed the commit(s)
  • Reviewed the code
  • Run, Tested, and Verified the change(s) work as expected

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".]

@malatewang
malatewang force-pushed the scaleout_deletion branch 2 times, most recently from a8dc213 to f9da45a Compare October 6, 2026 20:53

@marvinyu-memverge marvinyu-memverge left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed at f9da45a. The lease + recheck shape looks right to me; one ask.

  1. stop() can hang if it lands while the worker is in its idle scan. The scan runs inside the worker task, so if stop() queues its None sentinel while _queue_pending_session_deletions is awaiting get_sessions_by_status, the scan's keys go into the queue behind the sentinel. The worker takes None and returns, those items never get task_done(), and self._deletion_queue.join() never returns, so resources.close() isn't reached either. It needs at least one Deleted row 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. Checking if not self._started: return between the query and the put_nowait loop 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.

@edwinyyyu edwinyyyu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 Deleted sessions and probes one lease per key.
  • The scan runs only after DELETION_IDLE_SECONDS of 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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@malatewang

Copy link
Copy Markdown
Contributor Author

@marvinyu-memverge The shutdown race you reported is addressed in a095d3f: _queue_pending_session_deletions checks _started after its last await, and test_stop_finishes_when_idle_scan_finds_pending_deletion covers the scan/stop interleaving. The PR description also now says Part of #1749, so that issue remains open for the durable job, backoff, and terminal failure work.

@edwinyyyu

Copy link
Copy Markdown
Contributor

A follow-up question on the deletion queue. This is non-blocking and optional.

With the database scan in place, a session's Deleted status is the real to-do list. The in-memory queue's remaining job is to start a deletion as soon as delete_session is called instead of at the next scan. Is there another reason to keep it?

If not, an asyncio.Event could do that job instead. delete_session marks the session and sets the event. The worker deletes everything pending whenever the event is set, or every DELETION_IDLE_SECONDS, and checks a stop flag between sessions. This:

  • removes the None sentinel, the from_worker flag, and the _started recheck;
  • closes the remaining way to hang stop(), where delete_session puts a session on the queue behind the sentinel while stop() waits in join();
  • makes stop() wait for at most the deletion in progress. Sessions it skips stay Deleted, so the next scan picks them up.

The catch is that every deletion would then come from the scan, which only has keys, so key_to_session would become a required argument of start(). The server already passes it in server/api_v2/mcp.py. Three test call sites call start() without it: test_memmachine_integration.py, test_memmachine_delete_session.py, and the fixture in main/conftest.py. Does anything else call start() without a converter?


🤖 Written by Claude Code (Claude Opus 5.5) on behalf of @edwinyyyu.

@xiongzubiao

Copy link
Copy Markdown
Contributor

Should it be in feat/horizontal-scaling branch?

This branch has not been deployed

No deployments
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.

Project deletion never completes when semantic memory is disabled: the deletion worker requests the semantic service unconditionally and gives up

4 participants