From 15ab07cba277c73a83061f3c8fa3726d249ea62c Mon Sep 17 00:00:00 2001 From: cgycorey <4724788+cgycorey@users.noreply.github.com> Date: Sun, 2 Aug 2026 17:20:34 +0100 Subject: [PATCH] fix(ui): enforce durable org identity for runtime followups --- opc/engine.py | 91 ++- opc/layer2_organization/company_mode.py | 34 +- opc/layer2_organization/custom_runtime.py | 16 +- opc/plugins/office_ui/execution_identity.py | 45 +- opc/plugins/office_ui/ws_handler.py | 10 + tests/test_company_recruiter.py | 4 + tests/test_company_runtime_identity.py | 149 +++++ tests/test_org_mode_entrypoint.py | 29 + tests/test_runtime_config_enforcement.py | 587 +++++++++++++++++++- 9 files changed, 943 insertions(+), 22 deletions(-) diff --git a/opc/engine.py b/opc/engine.py index 097faa0..1967013 100644 --- a/opc/engine.py +++ b/opc/engine.py @@ -37,6 +37,7 @@ from opc.core.config import ( company_org_path, get_opc_home, get_project_workplace, + validate_organization_id, ) from opc.core.events import EventBus from opc.core.models import ( @@ -3563,6 +3564,13 @@ class OPCEngine: return False return bool(dict(getattr(task, "metadata", {}) or {}).get("shared_role_session", False)) + @staticmethod + def _normalize_durable_org_id(value: Any) -> str: + try: + return validate_organization_id(value) + except ValueError: + return "" + @staticmethod def _runtime_org_id_for_identity( decision: RouterDecision | None, @@ -3583,7 +3591,6 @@ class OPCEngine: getattr(decision, "org_id", None), task_metadata.get("org_id"), task_metadata.get("organization_id"), - getattr(org_config, "organization_id", None), ): normalized = str(candidate or "").strip() if normalized: @@ -3628,6 +3635,12 @@ class OPCEngine: root_session: bool = False, ) -> Task: assert self.store and self.memory + runtime_company_profile = str( + getattr(decision, "company_profile", "") + or (work_item.metadata or {}).get("company_profile", "") + or getattr(getattr(self.config, "org", None), "company_profile", "") + or "" + ).strip().lower() runtime_org_id = self._runtime_org_id_for_identity( decision, getattr(work_item, "metadata", None), @@ -3685,7 +3698,39 @@ class OPCEngine: set_linked_work_item_id(existing, work_item.work_item_id) existing.session_id = session_id existing.metadata = dict(existing.metadata or {}) - if runtime_org_id: + if runtime_company_profile == "custom": + persisted_org_id = self._normalize_durable_org_id(getattr(existing, "org_id", None)) + incoming_org_id = self._normalize_durable_org_id(runtime_org_id) + if persisted_org_id and incoming_org_id and persisted_org_id != incoming_org_id: + from opc.plugins.office_ui.services.models import ServiceError + raise ServiceError( + "org_id_conflict", + "org_id_conflict", + { + "project_id": self.project_id or "default", + "task_id": str(getattr(work_item, "work_item_id", "") or ""), + "persisted_org_id": persisted_org_id, + "incoming_org_id": incoming_org_id, + "reason": "custom_company_run_org_id_conflict", + }, + ) + resolved_org_id = persisted_org_id or incoming_org_id + if not resolved_org_id: + from opc.plugins.office_ui.services.models import ServiceError + raise ServiceError( + "org_id_required", + "org_id_required", + { + "project_id": self.project_id or "default", + "task_id": str(getattr(work_item, "work_item_id", "") or ""), + "reason": "custom_company_run_requires_durable_org_id", + }, + ) + runtime_org_id = resolved_org_id + existing.org_id = resolved_org_id + existing.metadata["org_id"] = resolved_org_id + existing.metadata["organization_id"] = resolved_org_id + elif runtime_org_id: existing.org_id = runtime_org_id existing.metadata["org_id"] = runtime_org_id existing.metadata["organization_id"] = runtime_org_id @@ -3730,6 +3775,17 @@ class OPCEngine: ) await self.store.save_task(existing) return existing + if runtime_company_profile == "custom" and not runtime_org_id: + from opc.plugins.office_ui.services.models import ServiceError + raise ServiceError( + "org_id_required", + "org_id_required", + { + "project_id": self.project_id or "default", + "task_id": str(getattr(work_item, "work_item_id", "") or ""), + "reason": "custom_company_run_requires_durable_org_id", + }, + ) employee_assignment = dict(topology_seat.get("employee_assignment", {}) or {}) if not employee_assignment and self.org_engine and role_id: preferred_employee_id = str(topology_seat.get("employee_id", "") or "").strip() or None @@ -3780,12 +3836,6 @@ class OPCEngine: owner_execution_copy = build_work_item_owner_execution_copy(work_item) owner_execution_copy.setdefault("delegation_role_session_id", role_session_id) owner_execution_copy["work_kind"] = work_item_turn_type - runtime_company_profile = str( - getattr(decision, "company_profile", "") - or (work_item.metadata or {}).get("company_profile", "") - or getattr(getattr(self.config, "org", None), "company_profile", "") - or "" - ).strip().lower() runtime_identity_metadata = ( { "org_id": runtime_org_id or "", @@ -12704,7 +12754,12 @@ class OPCEngine: seen_employee_ids.add(employee_id) history = "" if self.memory: - organization_id = str(getattr(getattr(self.config, "org", None), "organization_id", "") or "").strip() + organization_id = str( + getattr(delivery_task, "org_id", "") + or (delivery_task.metadata or {}).get("org_id") + or (delivery_task.metadata or {}).get("organization_id") + or "" + ).strip() history = self.memory.employee_evolution.build_employee_delta_context( employee_id, project_id=task.project_id, @@ -13246,13 +13301,19 @@ class OPCEngine: await self._mark_company_runtime_checkpoint_status(checkpoint, status="invalid") return "Could not run self-evolution because the runtime task set could not be restored." + from opc.plugins.office_ui.execution_identity import resolve_delivery_task_org_identity + + organization_id, identity_error = resolve_delivery_task_org_identity( + waiting_task, + payload=payload, + active_org_id=getattr(getattr(self.config, "org", None), "organization_id", ""), + default_org_id=DEFAULT_ORGANIZATION_ID, + ) + if identity_error: + await self._mark_company_runtime_checkpoint_status(checkpoint, status="invalid") + return f"Could not run self-evolution because {identity_error}." + plan = deserialize_company_work_item_runtime_plan(payload.get("company_work_item_plan") or payload.get("plan", {})) - organization_id = str( - getattr(waiting_task, "org_id", "") - or payload.get("organization_id") - or getattr(getattr(self.config, "org", None), "organization_id", "") - or DEFAULT_ORGANIZATION_ID - ).strip() or DEFAULT_ORGANIZATION_ID root_role_id = str( getattr(plan, "final_decider_role_id", "") or plan.metadata.get("final_decider_role_id", "") diff --git a/opc/layer2_organization/company_mode.py b/opc/layer2_organization/company_mode.py index 2fce38a..9e8dbf2 100644 --- a/opc/layer2_organization/company_mode.py +++ b/opc/layer2_organization/company_mode.py @@ -22,7 +22,11 @@ from opc.core.active_task_runs import ( ActiveTaskRunAdmissionClosed, ActiveTaskRunRegistry, ) -from opc.core.config import DEFAULT_EXTERNAL_AGENT_STARTUP_TIMEOUT_SECONDS, DEFAULT_ORGANIZATION_ID +from opc.core.config import ( + DEFAULT_EXTERNAL_AGENT_STARTUP_TIMEOUT_SECONDS, + DEFAULT_ORGANIZATION_ID, + validate_organization_id, +) from opc.core.models import ( AdaptiveRoleProfile, AdaptiveSignalSpec, @@ -1326,6 +1330,26 @@ class CompanyRuntimeSpecBuilder(CompanyRuntimeWorkItemHelper): or "corporate" ).strip() or "corporate" org_config = getattr(self.org_engine.config, "org", None) + selected_org_id = "" + if profile == "custom": + try: + selected_org_id = validate_organization_id(getattr(decision, "org_id", None)) + except ValueError: + selected_org_id = "" + if not selected_org_id: + # A custom-organization run must carry a durable org_id on the + # decision. Never derive it from the process-wide active + # config; fail closed before any work items are created. + from opc.plugins.office_ui.services.models import ServiceError + raise ServiceError( + "org_id_required", + "org_id_required", + { + "company_profile": profile, + "reason": "custom_company_run_requires_durable_org_id", + }, + ) + decision.org_id = selected_org_id metadata: dict[str, Any] = { "source": "work_item_runtime", "execution_mode": "company_mode", @@ -1333,7 +1357,11 @@ class CompanyRuntimeSpecBuilder(CompanyRuntimeWorkItemHelper): "runtime_model": "multi_team_org", "work_item_driven": True, "company_profile": profile, - "organization_id": str(getattr(org_config, "organization_id", "") or "").strip(), + "organization_id": ( + selected_org_id + if profile == "custom" + else str(getattr(org_config, "organization_id", "") or "").strip() + ), "organization_name": str(getattr(org_config, "organization_name", "") or "").strip(), "organization_config_file": str(getattr(org_config, "organization_config_file", "") or "").strip(), "original_request": original_message, @@ -1341,7 +1369,7 @@ class CompanyRuntimeSpecBuilder(CompanyRuntimeWorkItemHelper): "domains": list(getattr(decision, "domains", []) or []), "preferred_agent": getattr(decision, "preferred_agent", None), "requested_sub_tasks": list(getattr(decision, "sub_tasks", []) or []), - "org_id": getattr(decision, "org_id", None), + "org_id": selected_org_id if profile == "custom" else getattr(decision, "org_id", None), } return CompanyRuntimeSpec( profile=profile, diff --git a/opc/layer2_organization/custom_runtime.py b/opc/layer2_organization/custom_runtime.py index 24aa5e8..192c8e9 100644 --- a/opc/layer2_organization/custom_runtime.py +++ b/opc/layer2_organization/custom_runtime.py @@ -64,8 +64,22 @@ class CustomRuntimeRunner: ) -> str: from opc.engine import OPCEngine from opc.layer2_organization.phase_hooks import unregister_dispatcher_wake + from opc.plugins.office_ui.services.models import ServiceError - org_config, resolved_org_id = self._build_org_config(org_id) + normalized_org_id = str(org_id or "").strip() + if not normalized_org_id: + # Isolated org mode must carry a durable org_id; resolving the + # active index here would silently route the run to whichever + # organization is currently loaded. + raise ServiceError( + "org_id_required", + "org_id_required", + { + "project_id": project_id or self.parent.project_id or "default", + "reason": "custom_company_run_requires_durable_org_id", + }, + ) + org_config, resolved_org_id = self._build_org_config(normalized_org_id) normalized_project_id = str(project_id or self.parent.project_id or "default").strip() or "default" shared_store = getattr(self.parent, "store", None) runtime = OPCEngine( diff --git a/opc/plugins/office_ui/execution_identity.py b/opc/plugins/office_ui/execution_identity.py index 7201ab7..b902130 100644 --- a/opc/plugins/office_ui/execution_identity.py +++ b/opc/plugins/office_ui/execution_identity.py @@ -14,7 +14,7 @@ for that identity: from __future__ import annotations from dataclasses import dataclass -from typing import Any +from typing import Any, Mapping from opc.core.config import validate_organization_id from opc.layer2_organization.company_runtime_identity import is_company_runtime_task @@ -198,3 +198,46 @@ def execution_identity_from_task( default_preferred_agent=default_preferred_agent, explicit_exec_mode=explicit, ) + + +def resolve_delivery_task_org_identity( + task: Any | None, + *, + payload: Mapping[str, Any] | None = None, + active_org_id: Any = "", + default_org_id: Any = "", +) -> tuple[str, str]: + """Validate the org identity of a delivery self-evolution task. + + Returns ``(organization_id, error)`` with at most one non-empty. Prefers + ``Task.org_id``, then task metadata org fields; conflicting sources are + rejected. Checkpoint-payload org fields are a last-resort legacy fallback + and never override task/metadata identity. The active configuration org + is only consulted for a confirmed corporate task; custom-org deliveries + without a durable org id fail closed. + """ + metadata = task_metadata(task) + candidates: list[str] = [] + for value in ( + getattr(task, "org_id", None), + metadata.get("org_id"), + metadata.get("organization_id"), + ): + normalized = normalize_org_id(value) + if normalized and normalized not in candidates: + candidates.append(normalized) + if len(candidates) > 1: + return "", "the delivery task org identity conflicts across task and metadata sources" + task_org_id = candidates[0] if candidates else "" + if task_org_id: + return task_org_id, "" + payload_org_id = normalize_org_id( + (payload or {}).get("org_id") + or (payload or {}).get("organization_id") + ) + if payload_org_id: + return payload_org_id, "" + identity = execution_identity_from_task(task) + if identity.is_company: + return normalize_org_id(active_org_id) or normalize_org_id(default_org_id), "" + return "", "the custom-organization delivery task has no durable org identity" diff --git a/opc/plugins/office_ui/ws_handler.py b/opc/plugins/office_ui/ws_handler.py index cac83d2..51e9e7a 100644 --- a/opc/plugins/office_ui/ws_handler.py +++ b/opc/plugins/office_ui/ws_handler.py @@ -7854,6 +7854,16 @@ class WSHandler: parent_task = await run_engine.store.get_task(parent_task_id) except Exception: logger.opt(exception=True).debug("failed to load parent task for delivery feedback reply") + if parent_task is None: + raise ServiceError( + "org_id_required", + "org_id_required", + { + "project_id": pid, + "task_id": parent_task_id, + "reason": "delivery_feedback_requires_durable_parent_task", + }, + ) session_exec_mode = self._normalize_session_exec_mode(self._exec_mode) session_company_profile = self._normalize_session_company_profile(self._company_profile) session_org_id = "" diff --git a/tests/test_company_recruiter.py b/tests/test_company_recruiter.py index 6a4bf6a..28b8fb0 100644 --- a/tests/test_company_recruiter.py +++ b/tests/test_company_recruiter.py @@ -2072,6 +2072,7 @@ class CompanyRecruiterFlowTests(unittest.IsolatedAsyncioTestCase): decision = RouterDecision( mode=ExecutionMode.COMPANY_MODE, company_profile="custom", + org_id="test-org", domains=[], ) runtime_spec = engine.company_runtime_spec_builder.build_spec( @@ -2140,6 +2141,7 @@ class CompanyRecruiterFlowTests(unittest.IsolatedAsyncioTestCase): decision = RouterDecision( mode=ExecutionMode.COMPANY_MODE, company_profile="custom", + org_id="test-org", domains=[], ) runtime_spec = engine.company_runtime_spec_builder.build_spec( @@ -2199,6 +2201,7 @@ class CompanyRecruiterFlowTests(unittest.IsolatedAsyncioTestCase): decision = RouterDecision( mode=ExecutionMode.COMPANY_MODE, company_profile="custom", + org_id="test-org", domains=[], ) runtime_spec = engine.company_runtime_spec_builder.build_spec( @@ -2312,6 +2315,7 @@ class CompanyRecruiterFlowTests(unittest.IsolatedAsyncioTestCase): decision = RouterDecision( mode=ExecutionMode.COMPANY_MODE, company_profile="custom", + org_id="test-org", domains=[], ) runtime_spec = engine.company_runtime_spec_builder.build_spec( diff --git a/tests/test_company_runtime_identity.py b/tests/test_company_runtime_identity.py index dd8aa2b..e3d7a4b 100644 --- a/tests/test_company_runtime_identity.py +++ b/tests/test_company_runtime_identity.py @@ -260,6 +260,155 @@ def test_work_item_chat_resume_uses_canonical_ui_anchor_as_engine_origin() -> No asyncio.run(scenario()) +def test_company_suspend_reply_routes_selected_org_to_engine() -> None: + async def scenario() -> None: + _tasks, checkpoint = _runtime_records() + handler = WSHandler.__new__(WSHandler) + handler._exec_mode = "task" + handler._company_profile = "corporate" + handler._shutting_down = False + handler._active_runtime_children = {} + handler._session_to_task = {} + handler._task_bg_context = {} + handler._company_suspend_reply_locks = {"runtime-session": asyncio.Lock()} + handler.chat_store = None + handler._set_company_runtime_control = AsyncMock() + handler._normalize_session_exec_mode = MagicMock(return_value="task") + handler._normalize_session_company_profile = MagicMock(return_value="corporate") + handler._resolve_task_session_config = MagicMock(return_value=("org", "custom")) + handler._resolve_task_org_id = MagicMock(return_value="selected-org") + handler._extract_checkpoint_metadata = AsyncMock(return_value=None) + handler._sync_task_transcript_messages = AsyncMock() + handler.on_kanban_changed = AsyncMock() + handler._flush_progress = AsyncMock() + run_engine = SimpleNamespace( + project_id="project-a", + process_message=AsyncMock(return_value="resumed"), + ) + target = { + "ui_anchor_task_id": "ui-anchor", + "config_task": SimpleNamespace( + metadata={"exec_mode": "org", "company_profile": "custom"}, + org_id="selected-org", + ), + } + + await handler._process_company_suspend_reply( + ui_task_id="final-decider", + runtime_session_id="runtime-session", + content="continue", + attachment_refs=None, + message_metadata={"ui_force_resume": True}, + user_message_id=None, + user_message_created_at=None, + run_engine=run_engine, + run_project_id="project-a", + target=target, + checkpoint=checkpoint, + lock=handler._company_suspend_reply_locks["runtime-session"], + ) + + call = run_engine.process_message.await_args + assert call.kwargs["org_id"] == "selected-org" + + asyncio.run(scenario()) + + +def test_delivery_feedback_reply_fails_closed_without_durable_parent_task() -> None: + async def scenario() -> None: + handler = WSHandler.__new__(WSHandler) + handler._exec_mode = "task" + handler._company_profile = "corporate" + handler._shutting_down = False + handler._active_runtime_children = {} + handler._session_to_task = {} + handler._task_bg_context = {} + handler._company_delivery_feedback_reply_locks = {} + handler.chat_store = None + handler._store_is_ready = MagicMock(return_value=True) + handler._normalize_session_exec_mode = MagicMock(return_value="task") + handler._normalize_session_company_profile = MagicMock(return_value="corporate") + handler._resolve_task_session_config = MagicMock(return_value=("org", "custom")) + handler._resolve_task_org_id = MagicMock(return_value="") + handler._chat_store_is_ready = MagicMock(return_value=False) + handler._flush_progress = AsyncMock() + run_engine = SimpleNamespace( + project_id="project-a", + store=SimpleNamespace(get_task=AsyncMock(return_value=None)), + run_company_delivery_self_evolution_checkpoint=AsyncMock(return_value="ran"), + ) + + await handler._process_company_delivery_feedback_reply( + parent_task_id="delivery-task", + parent_session_id="delivery-session", + reply_channel_id="session:delivery-task", + content="approved", + attachment_refs=None, + message_metadata=None, + user_message_id=None, + user_message_created_at=None, + run_engine=run_engine, + run_project_id="project-a", + checkpoint=SimpleNamespace(checkpoint_id="cp-delivery", payload={}), + waiting_task_id="delivery-task", + lock=asyncio.Lock(), + ) + + run_engine.run_company_delivery_self_evolution_checkpoint.assert_not_awaited() + handler._resolve_task_org_id.assert_not_called() + + asyncio.run(scenario()) + + +def test_delivery_feedback_reply_proceeds_with_durable_parent_task() -> None: + async def scenario() -> None: + handler = WSHandler.__new__(WSHandler) + handler._exec_mode = "task" + handler._company_profile = "corporate" + handler._shutting_down = False + handler._active_runtime_children = {} + handler._session_to_task = {} + handler._task_bg_context = {} + handler._company_delivery_feedback_reply_locks = {} + handler.chat_store = None + handler._store_is_ready = MagicMock(return_value=True) + handler._normalize_session_exec_mode = MagicMock(return_value="task") + handler._normalize_session_company_profile = MagicMock(return_value="corporate") + handler._resolve_task_session_config = MagicMock(return_value=("org", "custom")) + handler._resolve_task_org_id = MagicMock(return_value="selected-org") + handler._chat_store_is_ready = MagicMock(return_value=False) + handler._flush_progress = AsyncMock() + parent_task = SimpleNamespace(id="delivery-task") + run_engine = SimpleNamespace( + project_id="project-a", + store=SimpleNamespace(get_task=AsyncMock(return_value=parent_task)), + run_company_delivery_self_evolution_checkpoint=AsyncMock(return_value="ran"), + ) + + await handler._process_company_delivery_feedback_reply( + parent_task_id="delivery-task", + parent_session_id="delivery-session", + reply_channel_id="session:delivery-task", + content="approved", + attachment_refs=None, + message_metadata=None, + user_message_id=None, + user_message_created_at=None, + run_engine=run_engine, + run_project_id="project-a", + checkpoint=SimpleNamespace(checkpoint_id="cp-delivery", payload={}), + waiting_task_id="delivery-task", + lock=asyncio.Lock(), + ) + + call = run_engine.run_company_delivery_self_evolution_checkpoint.await_args + assert call is not None + assert call.kwargs["action"] == "approve" + handler._resolve_task_org_id.assert_called_once_with(parent_task) + + asyncio.run(scenario()) + + def test_suspend_reply_missing_custom_org_id_fails_closed_before_engine_call() -> None: async def scenario() -> None: _tasks, checkpoint = _runtime_records() diff --git a/tests/test_org_mode_entrypoint.py b/tests/test_org_mode_entrypoint.py index f372e32..509a228 100644 --- a/tests/test_org_mode_entrypoint.py +++ b/tests/test_org_mode_entrypoint.py @@ -13,6 +13,7 @@ from opc.core.org_config import ( ) from opc.engine import OPCEngine from opc.layer2_organization.custom_runtime import CustomRuntimeRunner +from opc.plugins.office_ui.services.models import ServiceError def test_requested_mode_normalization_keeps_core_company_router_main_compatible() -> None: @@ -61,6 +62,34 @@ def test_custom_runtime_runner_loads_org_storage_without_mutating_parent_config( assert (config_dir / "company_orgs" / "org_lab_config.yaml").exists() +def test_process_message_rejects_org_mode_without_org_id() -> None: + engine = OPCEngine.__new__(OPCEngine) + engine.opc_home = None + engine.config = OPCConfig() + runner = CustomRuntimeRunner(engine) + caught: ServiceError | None = None + + async def _run() -> None: + await runner.process_message( + "run org", + project_id="default", + session_id="session-1", + org_id=None, + preferred_agent=None, + domains=None, + origin_task_id=None, + attachment_refs=None, + message_metadata=None, + ) + + try: + asyncio.run(_run()) + except ServiceError as exc: + caught = exc + assert caught is not None + assert caught.code == "org_id_required" + + def test_process_message_routes_org_mode_to_custom_runner(monkeypatch) -> None: engine = OPCEngine.__new__(OPCEngine) engine._initialized = True diff --git a/tests/test_runtime_config_enforcement.py b/tests/test_runtime_config_enforcement.py index 7c4b530..55789ac 100644 --- a/tests/test_runtime_config_enforcement.py +++ b/tests/test_runtime_config_enforcement.py @@ -11,9 +11,15 @@ from unittest.mock import AsyncMock, MagicMock, patch import yaml from pydantic import ValidationError -from opc.core.config import AgentsConfig, ExternalAgentConfig, OPCConfig +from opc.core.config import ( + AgentsConfig, + DEFAULT_ORGANIZATION_ID, + ExternalAgentConfig, + OPCConfig, +) from opc.core.models import ( DelegationWorkItem, + ExecutionCheckpoint, ExecutionMode, RouterDecision, SessionMessageRecord, @@ -24,7 +30,8 @@ from opc.core.models import ( WorkItemExecutionStrategy, ) from opc.engine import OPCEngine -from opc.layer2_organization.company_mode import CompanyWorkItemExecutor +from opc.layer2_organization.company_mode import CompanyRuntimeSpecBuilder, CompanyWorkItemExecutor +from opc.plugins.office_ui.services.models import ServiceError from opc.layer3_agent.adapters.claude_code import ClaudeCodeAdapter from opc.layer3_agent.adapters.codex_adapter import CodexAdapter from opc.layer3_agent.adapters.cursor_adapter import CursorAdapter @@ -391,6 +398,582 @@ class RuntimeConfigEnforcementTests(unittest.IsolatedAsyncioTestCase): self.assertEqual(task.metadata["org_id"], "selected-org") self.assertEqual(task.metadata["organization_id"], "selected-org") + def test_build_spec_rejects_custom_run_without_durable_org_id(self) -> None: + builder = CompanyRuntimeSpecBuilder( + org_engine=SimpleNamespace( + get_company_profile=lambda: "custom", + config=SimpleNamespace( + org=SimpleNamespace( + organization_id="active-org", + organization_name="Active Org", + organization_config_file="org_active-org_config.yaml", + company_profile="custom", + ) + ), + ) + ) + with self.assertRaises(ServiceError) as ctx: + builder.build_spec( + RouterDecision( + mode=ExecutionMode.COMPANY_MODE, + company_profile="custom", + org_id=None, + ), + original_message="Run the company.", + ) + assert ctx.exception.code == "org_id_required" + + def test_build_spec_serializes_selected_org_not_active_config(self) -> None: + builder = CompanyRuntimeSpecBuilder( + org_engine=SimpleNamespace( + get_company_profile=lambda: "custom", + config=SimpleNamespace( + org=SimpleNamespace( + organization_id="active-org", + organization_name="Active Org", + organization_config_file="org_active-org_config.yaml", + company_profile="custom", + ) + ), + ) + ) + decision = RouterDecision( + mode=ExecutionMode.COMPANY_MODE, + company_profile="custom", + org_id="selected-org", + ) + spec = builder.build_spec(decision, original_message="Run the company.") + assert spec.metadata["org_id"] == "selected-org" + assert spec.metadata["organization_id"] == "selected-org" + + def test_runtime_org_id_for_identity_never_derives_from_active_config(self) -> None: + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + decision = RouterDecision( + mode=ExecutionMode.COMPANY_MODE, + company_profile="custom", + org_id=None, + ) + resolved = OPCEngine._runtime_org_id_for_identity( + decision, + {}, + engine.config.org, + ) + assert resolved is None + + async def test_ensure_runtime_work_item_task_rejects_custom_without_durable_org(self) -> None: + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + engine.store = SimpleNamespace( + get_runtime_task_for_work_item=AsyncMock(return_value=None), + save_delegation_work_item=AsyncMock(), + save_task=AsyncMock(), + link_work_item_runtime_task=AsyncMock(return_value=True), + ) + engine.memory = SimpleNamespace(ensure_session=AsyncMock()) + engine.org_engine = SimpleNamespace( + current_org_version=MagicMock(return_value=1), + current_runtime_topology_version=MagicMock(return_value=1), + ) + engine._requests_explicit_project_knowledge = MagicMock(return_value=False) + work_item = DelegationWorkItem( + work_item_id="wi-no-durable-org", + run_id="run-no-durable-org", + cell_id="team::engineering", + team_instance_id="team-instance-1", + role_id="engineer", + seat_id="seat-engineer", + title="Engineering execution", + summary="Implement the requested change.", + kind="execute", + projection_id="engineering-execute", + metadata={"seat_id": "seat-engineer", "team_id": "team::engineering"}, + ) + with self.assertRaises(ServiceError) as ctx: + await engine._ensure_runtime_work_item_task( + work_item=work_item, + parent_session_id="sess-company", + original_message="Build the thing.", + decision=RouterDecision( + mode=ExecutionMode.COMPANY_MODE, + company_profile="custom", + org_id=None, + ), + runtime_topology={ + "final_decider_role_id": "lead", + "seats": [ + { + "seat_id": "seat-engineer", + "team_id": "team::engineering", + "role_id": "engineer", + "employee_assignment": {"employee_id": "eng-1", "name": "Engineer"}, + "metadata": {"role_name": "Engineer"}, + } + ], + }, + delegation_playbook={}, + secretary_context="", + target_output_dir=None, + origin_channel="cli", + origin_chat_id="", + origin_thread_id="", + origin_task_id=None, + attachment_refs=[], + attachment_context="", + force_native_execution=False, + ) + assert ctx.exception.code == "org_id_required" + engine.store.save_task.assert_not_awaited() + + async def test_ensure_runtime_work_item_task_rejects_existing_custom_task_without_durable_org(self) -> None: + existing = Task( + id="existing-no-org", + title="Engineering execution", + project_id="proj1", + session_id="sess-company:wi-existing-no-org", + assigned_to="engineer", + metadata={ + "execution_mode": "company_mode", + "runtime_model": "multi_team_org", + "work_item_runtime": True, + "work_item_projection_id": "engineering-execute", + "work_item_turn_type": "execute", + "company_profile": "custom", + "delegation_seat_id": "seat-engineer", + }, + ) + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + engine.store = SimpleNamespace( + get_runtime_task_for_work_item=AsyncMock(return_value=existing), + save_delegation_work_item=AsyncMock(), + save_task=AsyncMock(), + link_work_item_runtime_task=AsyncMock(return_value=True), + ) + engine.memory = SimpleNamespace(ensure_session=AsyncMock()) + work_item = DelegationWorkItem( + work_item_id="wi-existing-no-org", + run_id="run-existing-no-org", + cell_id="team::engineering", + team_instance_id="team-instance-1", + role_id="engineer", + seat_id="seat-engineer", + title="Engineering execution", + summary="Implement the requested change.", + kind="execute", + projection_id="engineering-execute", + metadata={"seat_id": "seat-engineer", "team_id": "team::engineering"}, + ) + with self.assertRaises(ServiceError) as ctx: + await engine._ensure_runtime_work_item_task( + work_item=work_item, + parent_session_id="sess-company", + original_message="Build the thing.", + decision=RouterDecision( + mode=ExecutionMode.COMPANY_MODE, + company_profile="custom", + org_id=None, + ), + runtime_topology={ + "final_decider_role_id": "lead", + "seats": [ + { + "seat_id": "seat-engineer", + "team_id": "team::engineering", + "role_id": "engineer", + "employee_assignment": {"employee_id": "eng-1", "name": "Engineer"}, + "metadata": {"role_name": "Engineer"}, + } + ], + }, + delegation_playbook={}, + secretary_context="", + target_output_dir=None, + origin_channel="cli", + origin_chat_id="", + origin_thread_id="", + origin_task_id=None, + attachment_refs=[], + attachment_context="", + force_native_execution=False, + ) + assert ctx.exception.code == "org_id_required" + engine.store.save_task.assert_not_awaited() + + async def test_ensure_runtime_work_item_task_rejects_conflicting_org_ids(self) -> None: + existing = Task( + id="existing-conflict-org", + title="Engineering execution", + project_id="proj1", + session_id="sess-company:wi-existing-conflict-org", + assigned_to="engineer", + org_id="persisted-org", + metadata={ + "execution_mode": "company_mode", + "runtime_model": "multi_team_org", + "work_item_runtime": True, + "work_item_projection_id": "engineering-execute", + "work_item_turn_type": "execute", + "company_profile": "custom", + "organization_id": "persisted-org", + "delegation_seat_id": "seat-engineer", + }, + ) + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + engine.store = SimpleNamespace( + get_runtime_task_for_work_item=AsyncMock(return_value=existing), + save_delegation_work_item=AsyncMock(), + save_task=AsyncMock(), + link_work_item_runtime_task=AsyncMock(return_value=True), + ) + engine.memory = SimpleNamespace(ensure_session=AsyncMock()) + work_item = DelegationWorkItem( + work_item_id="wi-existing-conflict-org", + run_id="run-existing-conflict-org", + cell_id="team::engineering", + team_instance_id="team-instance-1", + role_id="engineer", + seat_id="seat-engineer", + title="Engineering execution", + summary="Implement the requested change.", + kind="execute", + projection_id="engineering-execute", + metadata={"seat_id": "seat-engineer", "team_id": "team::engineering"}, + ) + with self.assertRaises(ServiceError) as ctx: + await engine._ensure_runtime_work_item_task( + work_item=work_item, + parent_session_id="sess-company", + original_message="Build the thing.", + decision=RouterDecision( + mode=ExecutionMode.COMPANY_MODE, + company_profile="custom", + org_id="selected-org", + ), + runtime_topology={ + "final_decider_role_id": "lead", + "seats": [ + { + "seat_id": "seat-engineer", + "team_id": "team::engineering", + "role_id": "engineer", + "employee_assignment": {"employee_id": "eng-1", "name": "Engineer"}, + "metadata": {"role_name": "Engineer"}, + } + ], + }, + delegation_playbook={}, + secretary_context="", + target_output_dir=None, + origin_channel="cli", + origin_chat_id="", + origin_thread_id="", + origin_task_id=None, + attachment_refs=[], + attachment_context="", + force_native_execution=False, + ) + assert ctx.exception.code == "org_id_conflict" + engine.store.save_task.assert_not_awaited() + + async def test_ensure_runtime_work_item_task_repairs_existing_custom_task_from_persisted_durable_org(self) -> None: + existing = Task( + id="existing-durable-org", + title="Engineering execution", + project_id="proj1", + session_id="sess-company:wi-existing-durable-org", + assigned_to="engineer", + org_id="persisted-org", + metadata={ + "execution_mode": "company_mode", + "runtime_model": "multi_team_org", + "work_item_runtime": True, + "work_item_projection_id": "engineering-execute", + "work_item_turn_type": "execute", + "company_profile": "custom", + "organization_id": "stale-org", + "delegation_seat_id": "seat-engineer", + }, + ) + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + engine.store = SimpleNamespace( + get_runtime_task_for_work_item=AsyncMock(return_value=existing), + save_delegation_work_item=AsyncMock(), + save_task=AsyncMock(), + link_work_item_runtime_task=AsyncMock(return_value=True), + ) + engine.memory = SimpleNamespace(ensure_session=AsyncMock()) + work_item = DelegationWorkItem( + work_item_id="wi-existing-durable-org", + run_id="run-existing-durable-org", + cell_id="team::engineering", + team_instance_id="team-instance-1", + role_id="engineer", + seat_id="seat-engineer", + title="Engineering execution", + summary="Implement the requested change.", + kind="execute", + projection_id="engineering-execute", + metadata={"seat_id": "seat-engineer", "team_id": "team::engineering"}, + ) + repaired = await engine._ensure_runtime_work_item_task( + work_item=work_item, + parent_session_id="sess-company", + original_message="Build the thing.", + decision=RouterDecision( + mode=ExecutionMode.COMPANY_MODE, + company_profile="custom", + org_id=None, + ), + runtime_topology={ + "final_decider_role_id": "lead", + "seats": [ + { + "seat_id": "seat-engineer", + "team_id": "team::engineering", + "role_id": "engineer", + "employee_assignment": {"employee_id": "eng-1", "name": "Engineer"}, + "metadata": {"role_name": "Engineer"}, + } + ], + }, + delegation_playbook={}, + secretary_context="", + target_output_dir=None, + origin_channel="cli", + origin_chat_id="", + origin_thread_id="", + origin_task_id=None, + attachment_refs=[], + attachment_context="", + force_native_execution=False, + ) + self.assertEqual(repaired.id, "existing-durable-org") + self.assertEqual(repaired.org_id, "persisted-org") + self.assertEqual(repaired.metadata["org_id"], "persisted-org") + self.assertEqual(repaired.metadata["organization_id"], "persisted-org") + engine.store.save_task.assert_awaited_with(existing) + + async def test_delivery_self_evolution_custom_without_durable_org_fails_closed(self) -> None: + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + waiting_task = Task( + id="waiting-custom", + project_id="proj1", + session_id="sess-delivery", + metadata={ + "execution_mode": "company_mode", + "company_profile": "custom", + "work_item_runtime": True, + "work_item_projection_id": "delivery", + "work_item_turn_type": "deliver", + }, + ) + engine.store = SimpleNamespace(get_task=AsyncMock(return_value=waiting_task)) + engine._mark_company_runtime_checkpoint_status = AsyncMock() + engine._create_company_self_evolution_root_work_item = AsyncMock() + checkpoint = ExecutionCheckpoint( + checkpoint_id="cp-delivery", + project_id="proj1", + session_id="sess-delivery", + task_id="waiting-custom", + checkpoint_type="company_delivery_feedback", + payload={"waiting_task_id": "waiting-custom", "task_ids": ["waiting-custom"]}, + ) + result = await engine._run_company_delivery_self_evolution_consumed( + checkpoint, + action="approve", + ) + assert "durable org identity" in result + engine._mark_company_runtime_checkpoint_status.assert_awaited_once_with( + checkpoint, + status="invalid", + ) + engine._create_company_self_evolution_root_work_item.assert_not_awaited() + + async def test_delivery_self_evolution_custom_uses_durable_org_over_payload_and_config(self) -> None: + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + waiting_task = Task( + id="waiting-custom-2", + project_id="proj1", + session_id="sess-delivery", + org_id="selected-org", + metadata={ + "execution_mode": "company_mode", + "company_profile": "custom", + "work_item_runtime": True, + "work_item_projection_id": "delivery", + "work_item_turn_type": "deliver", + }, + ) + engine.store = SimpleNamespace(get_task=AsyncMock(return_value=waiting_task)) + engine._mark_company_runtime_checkpoint_status = AsyncMock() + engine.org_engine = None + engine._company_followup_target_task = MagicMock( + return_value=SimpleNamespace(assigned_to="lead", metadata={}), + ) + engine.company_executor = SimpleNamespace() + engine._self_evolution_assignments_by_role = MagicMock(return_value={}) + engine._create_company_self_evolution_root_work_item = AsyncMock(return_value=None) + checkpoint = ExecutionCheckpoint( + checkpoint_id="cp-delivery-2", + project_id="proj1", + session_id="sess-delivery", + task_id="waiting-custom-2", + checkpoint_type="company_delivery_feedback", + payload={ + "waiting_task_id": "waiting-custom-2", + "task_ids": ["waiting-custom-2"], + "organization_id": "stale-org", + }, + ) + await engine._run_company_delivery_self_evolution_consumed( + checkpoint, + action="approve", + ) + call = engine._create_company_self_evolution_root_work_item.await_args + assert call is not None + assert call.kwargs["organization_id"] == "selected-org" + + async def test_delivery_self_evolution_corporate_still_uses_config_default(self) -> None: + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + waiting_task = Task( + id="waiting-corporate", + project_id="proj1", + session_id="sess-delivery", + metadata={ + "execution_mode": "company_mode", + "company_profile": "corporate", + "work_item_runtime": True, + "work_item_projection_id": "delivery", + "work_item_turn_type": "deliver", + }, + ) + engine.store = SimpleNamespace(get_task=AsyncMock(return_value=waiting_task)) + engine._mark_company_runtime_checkpoint_status = AsyncMock() + engine.org_engine = None + engine._company_followup_target_task = MagicMock( + return_value=SimpleNamespace(assigned_to="lead", metadata={}), + ) + engine.company_executor = SimpleNamespace() + engine._self_evolution_assignments_by_role = MagicMock(return_value={}) + engine._create_company_self_evolution_root_work_item = AsyncMock(return_value=None) + checkpoint = ExecutionCheckpoint( + checkpoint_id="cp-delivery-corp", + project_id="proj1", + session_id="sess-delivery", + task_id="waiting-corporate", + checkpoint_type="company_delivery_feedback", + payload={"waiting_task_id": "waiting-corporate", "task_ids": ["waiting-corporate"]}, + ) + await engine._run_company_delivery_self_evolution_consumed( + checkpoint, + action="approve", + ) + call = engine._create_company_self_evolution_root_work_item.await_args + assert call is not None + assert call.kwargs["organization_id"] == DEFAULT_ORGANIZATION_ID + + async def test_delivery_self_evolution_metadata_only_legacy_custom_uses_metadata_org(self) -> None: + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + waiting_task = Task( + id="waiting-legacy-metadata", + project_id="proj1", + session_id="sess-delivery", + metadata={ + "work_item_runtime": True, + "work_item_projection_id": "delivery", + "work_item_turn_type": "deliver", + "org_id": "selected-org", + }, + ) + engine.store = SimpleNamespace(get_task=AsyncMock(return_value=waiting_task)) + engine._mark_company_runtime_checkpoint_status = AsyncMock() + engine.org_engine = None + engine._company_followup_target_task = MagicMock( + return_value=SimpleNamespace(assigned_to="lead", metadata={}), + ) + engine.company_executor = SimpleNamespace() + engine._self_evolution_assignments_by_role = MagicMock(return_value={}) + engine._create_company_self_evolution_root_work_item = AsyncMock(return_value=None) + checkpoint = ExecutionCheckpoint( + checkpoint_id="cp-delivery-legacy", + project_id="proj1", + session_id="sess-delivery", + task_id="waiting-legacy-metadata", + checkpoint_type="company_delivery_feedback", + payload={ + "waiting_task_id": "waiting-legacy-metadata", + "task_ids": ["waiting-legacy-metadata"], + }, + ) + await engine._run_company_delivery_self_evolution_consumed( + checkpoint, + action="approve", + ) + call = engine._create_company_self_evolution_root_work_item.await_args + assert call is not None + assert call.kwargs["organization_id"] == "selected-org" + + async def test_delivery_self_evolution_conflicting_org_ids_fail_closed(self) -> None: + engine = OPCEngine(config=OPCConfig(), project_id="proj1") + engine.config.org.company_profile = "custom" + engine.config.org.organization_id = "active-org" + waiting_task = Task( + id="waiting-conflict", + project_id="proj1", + session_id="sess-delivery", + org_id="selected-org", + metadata={ + "work_item_runtime": True, + "work_item_projection_id": "delivery", + "work_item_turn_type": "deliver", + "org_id": "other-org", + }, + ) + engine.store = SimpleNamespace(get_task=AsyncMock(return_value=waiting_task)) + engine._mark_company_runtime_checkpoint_status = AsyncMock() + engine.org_engine = None + engine._company_followup_target_task = MagicMock( + return_value=SimpleNamespace(assigned_to="lead", metadata={}), + ) + engine.company_executor = SimpleNamespace() + engine._self_evolution_assignments_by_role = MagicMock(return_value={}) + engine._create_company_self_evolution_root_work_item = AsyncMock(return_value=None) + checkpoint = ExecutionCheckpoint( + checkpoint_id="cp-delivery-conflict", + project_id="proj1", + session_id="sess-delivery", + task_id="waiting-conflict", + checkpoint_type="company_delivery_feedback", + payload={ + "waiting_task_id": "waiting-conflict", + "task_ids": ["waiting-conflict"], + }, + ) + result = await engine._run_company_delivery_self_evolution_consumed( + checkpoint, + action="approve", + ) + assert "org identity" in result + engine._mark_company_runtime_checkpoint_status.assert_awaited_once_with( + checkpoint, + status="invalid", + ) + engine._create_company_self_evolution_root_work_item.assert_not_awaited() + async def test_company_materialized_work_item_uses_selected_agent_over_template_preference(self) -> None: saved_tasks: list[Task] = []