From 4f36f4a36fe090fb8a2b16d8115cd605a15ba3fe Mon Sep 17 00:00:00 2001 From: oliver Date: Tue, 21 Jul 2026 16:20:31 +0800 Subject: [PATCH] fix(scheduler): reuse one execution session per scheduled job Reuse channel_session_v2 per job_id so cron runs do not create a new oclaw session each time; isolate scheduled turn context to the active turn_uuid so shared sessions do not leak prior run history. Co-authored-by: Cursor --- runtime/direct_loop.py | 9 +++++- runtime/scheduler/service.py | 1 - runtime/scheduler/session_resolver.py | 33 ++++++++++++---------- tests/test_scheduled_session_isolation.py | 34 ++++++++++++++++++----- 4 files changed, 53 insertions(+), 24 deletions(-) diff --git a/runtime/direct_loop.py b/runtime/direct_loop.py index 73ec1554..29699082 100644 --- a/runtime/direct_loop.py +++ b/runtime/direct_loop.py @@ -588,6 +588,14 @@ def _build_model_context( active_turn_uuid: str | None = None, ) -> list[dict[str, Any]]: rows = store.get_messages(session_id=session_id, limit=int(max_messages)) + pb_ctx = prompt_build_context if isinstance(prompt_build_context, dict) else {} + if bool(pb_ctx.get("scheduled_proactive")) and str(active_turn_uuid or "").strip(): + turn_key = str(active_turn_uuid).strip() + rows = [ + r + for r in rows + if str(getattr(r, "turn_uuid", "") or "").strip() == turn_key + ] rows = _guard_tool_results_for_llm_context( store=store, session_id=session_id, @@ -621,7 +629,6 @@ def _build_model_context( # Hook integration: wiki-auto-inject can prepend retrieval snippets # before prompt build when query/topic hints indicate supplemental lookup. try: - pb_ctx = prompt_build_context if isinstance(prompt_build_context, dict) else {} user_text_final = str(user_text or "").strip() wiki_query = str(pb_ctx.get("wiki_query") or "").strip() hook_ctx = { diff --git a/runtime/scheduler/service.py b/runtime/scheduler/service.py index 8526305d..aa9ef1a0 100644 --- a/runtime/scheduler/service.py +++ b/runtime/scheduler/service.py @@ -46,7 +46,6 @@ def enqueue_scheduled_job_run( store, job=job, created_by_user_id=str(getattr(job, "created_by_user_id", "") or ""), - run_id=str(run.id), ) except Exception as exc: store.scheduled_job_run_update( diff --git a/runtime/scheduler/session_resolver.py b/runtime/scheduler/session_resolver.py index 302924c9..684f74dd 100644 --- a/runtime/scheduler/session_resolver.py +++ b/runtime/scheduler/session_resolver.py @@ -51,24 +51,28 @@ def _ensure_administrator_owner(store: Any, *, tenant_id: str) -> dict[str, Any] } -def _create_scheduled_execution_session( +def _get_or_create_scheduled_execution_session( store: Any, *, tenant_id: str, - user_id: str, + job_id: str, job_name: str, - run_id: str = "", ) -> str: - rid = str(run_id or "").strip() + """One execution session per scheduled job (not per run); visible for ops cleanup.""" title = f"Scheduled · {job_name}" - if rid: - title = f"{title} · {rid[:8]}" + jid = str(job_id or "").strip() tid = str(tenant_id or "").strip() - uid = str(user_id or "").strip() - if uid and tid: - sess = store.create_session_for_user(title=title, tenant_id=tid, user_id=uid) - else: - sess = store.create_session(title) + getter = getattr(store, "get_or_create_channel_session_v2", None) + if jid and tid and callable(getter): + return getter( + tenant_id=tid, + channel="scheduled_job", + account_id="job", + external_chat_id=jid, + external_user_id=jid, + session_title=title, + ) + sess = store.create_session(title) return str(sess.id) @@ -255,9 +259,9 @@ def resolve_scheduled_session( *, job: Any, created_by_user_id: str = "", - run_id: str = "", ) -> ResolvedSession: job_name = str(getattr(job, "name", "") or "Scheduled task") + job_id = str(getattr(job, "id", "") or "").strip() ( tenant_id, user_id, @@ -272,12 +276,11 @@ def resolve_scheduled_session( job=job, created_by_user_id=created_by_user_id, ) - execution_session_id = _create_scheduled_execution_session( + execution_session_id = _get_or_create_scheduled_execution_session( store, tenant_id=tenant_id, - user_id=user_id, + job_id=job_id, job_name=job_name, - run_id=run_id, ) return ResolvedSession( session_id=execution_session_id, diff --git a/tests/test_scheduled_session_isolation.py b/tests/test_scheduled_session_isolation.py index 47f989f4..435d378b 100644 --- a/tests/test_scheduled_session_isolation.py +++ b/tests/test_scheduled_session_isolation.py @@ -64,8 +64,9 @@ class ScheduledSessionIsolationTests(unittest.TestCase): def tearDown(self) -> None: self._tmp.cleanup() - def test_resolve_scheduled_session_uses_fresh_execution_session(self) -> None: + def test_resolve_scheduled_session_isolates_from_interactive_and_reuses_per_job(self) -> None: job = mock.MagicMock() + job.id = "job-water-reminder" job.tenant_id = self.tenant_id job.source_session_id = self.interactive_session_id job.delivery_json = json.dumps( @@ -81,18 +82,37 @@ class ScheduledSessionIsolationTests(unittest.TestCase): ) job.name = "喝水提醒" - resolved = resolve_scheduled_session( + resolved_first = resolve_scheduled_session( self.store, job=job, created_by_user_id=self.admin_id, - run_id="run-abc12345", ) - self.assertNotEqual(resolved.session_id, self.interactive_session_id) - self.assertEqual(resolved.source_session_id, self.interactive_session_id) - self.assertEqual(resolved.external_chat_id, "120363012345678@g.us") - rows = self.store.get_messages(session_id=resolved.session_id, limit=10) + resolved_second = resolve_scheduled_session( + self.store, + job=job, + created_by_user_id=self.admin_id, + ) + + self.assertNotEqual(resolved_first.session_id, self.interactive_session_id) + self.assertEqual(resolved_first.session_id, resolved_second.session_id) + self.assertEqual(resolved_first.source_session_id, self.interactive_session_id) + self.assertEqual(resolved_first.external_chat_id, "120363012345678@g.us") + rows = self.store.get_messages(session_id=resolved_first.session_id, limit=10) self.assertEqual(len(rows), 0) + other_job = mock.MagicMock() + other_job.id = "job-other" + other_job.tenant_id = self.tenant_id + other_job.source_session_id = self.interactive_session_id + other_job.delivery_json = job.delivery_json + other_job.name = "Other task" + resolved_other = resolve_scheduled_session( + self.store, + job=other_job, + created_by_user_id=self.admin_id, + ) + self.assertNotEqual(resolved_other.session_id, resolved_first.session_id) + class ScheduledMentionCreatorTests(unittest.TestCase): def setUp(self) -> None: