diff --git a/README.md b/README.md index bee7f90..e6e7167 100644 --- a/README.md +++ b/README.md @@ -786,4 +786,5 @@ OpenOPC is built with gratitude for several open-source projects that helped sha - [BloopAI/vibe-kanban](https://github.com/BloopAI/vibe-kanban) for inspiration around kanban-centered agent work management and task visibility. - [msitarzewski/agency-agents](https://github.com/msitarzewski/agency-agents) for the talent-template foundation. All talent templates included in this repository are imported from `agency-agents`. - [HKUDS/nanobot](https://github.com/HKUDS/nanobot) for inspiration around skill-oriented agent design and `SKILL.md`-style organization. +- [pixel-agents-hq/pixel-agents](https://github.com/pixel-agents-hq/pixel-agents) for inspiration around the animated pixel-art office visualization of agent activity. diff --git a/opc/engine.py b/opc/engine.py index 6e69856..364e9b0 100644 --- a/opc/engine.py +++ b/opc/engine.py @@ -1028,7 +1028,7 @@ class OPCEngine: try: data = json.loads(path.read_text(encoding="utf-8")) except Exception: - logger.debug("Failed to load project company staffing defaults", exc_info=True) + logger.opt(exception=True).debug("Failed to load project company staffing defaults") return {} scopes = dict(data.get("scopes", {}) or {}) if isinstance(data, dict) else {} return dict(scopes.get(self._company_staffing_scope_key(decision, company_profile=company_profile), {}) or {}) @@ -1549,7 +1549,7 @@ class OPCEngine: try: tasks = await self.store.get_tasks(project_id=self.project_id or "default") except Exception: - logger.debug("failed to inspect company tasks before recording top-level reply", exc_info=True) + logger.opt(exception=True).debug("failed to inspect company tasks before recording top-level reply") return marker_fallback session_key = str(session_id or "").strip() for task in tasks: @@ -5379,10 +5379,9 @@ class OPCEngine: claimed_by_seat_id="", ) except Exception: - logger.debug( + logger.opt(exception=True).debug( "company runtime suspend: hold/release failed for %s", work_item_id, - exc_info=True, ) else: task.metadata["dispatch_hold"] = "company_runtime_suspended" @@ -5402,7 +5401,7 @@ class OPCEngine: } await self.store.save_external_session(latest_session) except Exception: - logger.debug("company runtime suspend: external session status update failed", exc_info=True) + logger.opt(exception=True).debug("company runtime suspend: external session status update failed") fresh = await self.store.get_task(task.id) target = fresh or task @@ -5434,7 +5433,7 @@ class OPCEngine: }, ) except Exception: - logger.debug("company runtime suspend: role session idle update failed", exc_info=True) + logger.opt(exception=True).debug("company runtime suspend: role session idle update failed") return affected async def suspend_company_runtime( @@ -5618,7 +5617,7 @@ class OPCEngine: if work_item_id: work_item_by_id[work_item_id] = item except Exception: - logger.debug("company runtime resume: failed to load run work items", exc_info=True) + logger.opt(exception=True).debug("company runtime resume: failed to load run work items") for task in tasks: if resume_task_ids is not None and task.id not in resume_task_ids: refreshed.append(task) @@ -5759,7 +5758,7 @@ class OPCEngine: claimed_by_seat_id="", ) except Exception: - logger.debug("company runtime resume: phase restore/hold clear failed", exc_info=True) + logger.opt(exception=True).debug("company runtime resume: phase restore/hold clear failed") try: await update_work_item( work_item_id, @@ -5772,7 +5771,7 @@ class OPCEngine: claimed_by_seat_id="", ) except Exception: - logger.debug("company runtime resume: fallback hold clear failed", exc_info=True) + logger.opt(exception=True).debug("company runtime resume: fallback hold clear failed") role_session_id = str(task.metadata.get("delegation_role_session_id", "") or "").strip() if role_session_id and callable(update_role_session): @@ -5789,7 +5788,7 @@ class OPCEngine: }, ) except Exception: - logger.debug("company runtime resume: role session update failed", exc_info=True) + logger.opt(exception=True).debug("company runtime resume: role session update failed") await self.store.save_task(task) fresh = await self.store.get_task(task.id) refreshed.append(fresh or task) @@ -6303,10 +6302,9 @@ class OPCEngine: await self._save_company_feedback_followup_checkpoint(task, tasks, plan) restored_checkpoint = True except Exception: - logger.debug( + logger.opt(exception=True).debug( "Best-effort restore of human feedback checkpoint failed for task {}", task.id, - exc_info=True, ) if restored_checkpoint: task.metadata = dict(task.metadata or {}) @@ -7072,10 +7070,9 @@ class OPCEngine: metadata_updates=work_item_updates, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( "failed to update closed delivery review work item metadata for {}", work_item_id, - exc_info=True, ) except TypeError: try: @@ -7086,16 +7083,14 @@ class OPCEngine: metadata_updates=work_item_updates, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( "failed to approve closed delivery review work item for {}", work_item_id, - exc_info=True, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( "failed to approve closed delivery review work item for {}", work_item_id, - exc_info=True, ) async def _terminalize_company_delivery_feedback_checkpoint( @@ -7210,7 +7205,7 @@ class OPCEngine: checkpoint_types=["company_delivery_feedback"], ) except Exception: - logger.debug("failed to inspect pending delivery feedback checkpoints", exc_info=True) + logger.opt(exception=True).debug("failed to inspect pending delivery feedback checkpoints") pending = [] pending_task_ids = { str( @@ -7227,10 +7222,9 @@ class OPCEngine: try: await self._save_company_feedback_followup_checkpoint(task, tasks, plan) except Exception: - logger.debug( + logger.opt(exception=True).debug( "failed to restore missing delivery feedback checkpoint for task {}", task.id, - exc_info=True, ) @staticmethod @@ -8293,7 +8287,7 @@ class OPCEngine: metadata_updates={"prompt_contract": target_contract}, ) except Exception: - logger.debug("Best-effort external target prompt_contract update failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort external target prompt_contract update failed") else: target_contract = prompt_contract_from_work_item( SimpleNamespace( @@ -8335,7 +8329,7 @@ class OPCEngine: if updated is not None: work_metadata = dict(getattr(updated, "metadata", {}) or {}) except Exception: - logger.debug("Best-effort external prompt_contract update failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort external prompt_contract update failed") work_metadata = {**work_metadata, **metadata_updates} else: work_metadata = {**work_metadata, **metadata_updates} @@ -9230,7 +9224,7 @@ class OPCEngine: try: tasks = await self.store.get_tasks(project_id=self.project_id or "default") except Exception: - logger.debug("failed to load tasks while resolving company parent session", exc_info=True) + logger.opt(exception=True).debug("failed to load tasks while resolving company parent session") return "" for task in tasks: task_session_id = str(getattr(task, "session_id", "") or "").strip() @@ -9313,9 +9307,9 @@ class OPCEngine: if checkpoint is not None: return checkpoint except Exception: - logger.debug("direct checkpoint lookup failed", exc_info=True) + logger.opt(exception=True).debug("direct checkpoint lookup failed") except Exception: - logger.debug("direct checkpoint lookup failed", exc_info=True) + logger.opt(exception=True).debug("direct checkpoint lookup failed") listing_getter = getattr(self.store, "get_execution_checkpoints", None) if not callable(listing_getter): diff --git a/opc/layer2_organization/company_mode.py b/opc/layer2_organization/company_mode.py index 91b58db..68a3c34 100644 --- a/opc/layer2_organization/company_mode.py +++ b/opc/layer2_organization/company_mode.py @@ -3314,7 +3314,7 @@ class CompanyWorkItemExecutor: try: await self.store.save_task(task) except Exception: - logger.debug("Best-effort runtime Task projection sync failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort runtime Task projection sync failed") async def _refresh_ready_work_items( self, @@ -3372,11 +3372,11 @@ class CompanyWorkItemExecutor: try: self._signal_dispatcher_wake() except Exception: - logger.debug("Best-effort dispatcher wake after dependency release failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort dispatcher wake after dependency release failed") try: await self._notify_kanban_changed() except Exception: - logger.debug("Best-effort kanban notify after dependency release failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort kanban notify after dependency release failed") run_id = str(work_items[0].run_id or "").strip() return await self.store.list_delegation_work_items(run_id) @@ -4384,7 +4384,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("company runtime cancellation: failed session idle reset", exc_info=True) + logger.opt(exception=True).debug("company runtime cancellation: failed session idle reset") for claimed_task in claimed_tasks: if claimed_task.status in {TaskStatus.DONE, TaskStatus.FAILED, TaskStatus.CANCELLED}: continue @@ -4432,17 +4432,15 @@ class CompanyWorkItemExecutor: claimed_task.metadata["dispatch_hold"] = "company_runtime_suspended" claimed_task.metadata["suspended_phase"] = phase_value except Exception: - logger.debug( + logger.opt(exception=True).debug( "company runtime cancellation: failed suspend hold release", - exc_info=True, ) if self.save_task and self._store_is_ready(self.store): try: await self.save_task(claimed_task) except Exception: - logger.debug( + logger.opt(exception=True).debug( "company runtime cancellation: failed suspended task save", - exc_info=True, ) raise finally: @@ -4664,7 +4662,7 @@ class CompanyWorkItemExecutor: try: work_item = await self.store.get_delegation_work_item(work_item_id) except Exception: - logger.debug("work-item revision stale guard: failed to load work item", exc_info=True) + logger.opt(exception=True).debug("work-item revision stale guard: failed to load work item") return None if work_item is None: return None @@ -5079,9 +5077,8 @@ class CompanyWorkItemExecutor: task.metadata["suspended_phase"] = phase_value task.status = task_status_for_phase(phase) if isinstance(phase, Phase) else task.status except Exception: - logger.debug( + logger.opt(exception=True).debug( "company runtime cancellation: failed to apply suspend hold", - exc_info=True, ) if self.save_task and self._store_is_ready(self.store): await self.save_task(task) @@ -5611,7 +5608,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("Best-effort close of orphan report card failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort close of orphan report card failed") return None # The report turn's prose IS the handoff. Try a structured parse @@ -5646,7 +5643,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("Best-effort close of orphan report card failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort close of orphan report card failed") return None parent_metadata = dict(getattr(parent_item, "metadata", {}) or {}) @@ -5667,7 +5664,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("Best-effort close of non-reviewable report card failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort close of non-reviewable report card failed") await self._record_work_item_runtime_diagnostic( code="report_parent_not_reviewable", severity="info", @@ -5761,9 +5758,8 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( "report_done: failed to stamp worker_report on review card", - exc_info=True, ) else: await self._record_work_item_runtime_diagnostic( @@ -6508,7 +6504,7 @@ class CompanyWorkItemExecutor: ) worker_metadata = {**worker_metadata, "prompt_contract": target_prompt_contract} except Exception: - logger.debug("Best-effort target prompt_contract snapshot update failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort target prompt_contract snapshot update failed") review_prompt_contract = make_prompt_contract( task_brief=( "Review the completed child deliverable and decide whether to " @@ -6547,7 +6543,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("Best-effort in-flight review refresh failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort in-flight review refresh failed") return existing_card attempt_no = self._next_review_attempt(worker_metadata) @@ -6614,7 +6610,7 @@ class CompanyWorkItemExecutor: try: await self.store.save_delegation_work_item(review_work_item) except Exception: - logger.debug("Best-effort review work-item create failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort review work-item create failed") return None # Persist the attempt counter on the worker so future calls can # locate the current review without scanning. @@ -6624,7 +6620,7 @@ class CompanyWorkItemExecutor: metadata_updates={"review_attempt_count": attempt_no}, ) except Exception: - logger.debug("Best-effort review_attempt_count update failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort review_attempt_count update failed") return review_work_item async def _ensure_report_work_item_for_work_item( @@ -6705,7 +6701,7 @@ class CompanyWorkItemExecutor: ) worker_metadata = {**worker_metadata, "prompt_contract": target_prompt_contract} except Exception: - logger.debug("Best-effort target prompt_contract update before report failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort target prompt_contract update before report failed") if worker_item is not None and not is_manager_reviewable_turn(worker_item): await self._record_work_item_runtime_diagnostic( code="report_parent_not_reviewable", @@ -6750,7 +6746,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("Best-effort in-flight report refresh failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort in-flight report refresh failed") return existing_card attempt_no = self._next_report_attempt(worker_metadata) @@ -6817,7 +6813,7 @@ class CompanyWorkItemExecutor: try: await self.store.save_delegation_work_item(report_work_item) except Exception: - logger.debug("Best-effort report work-item create failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort report work-item create failed") return None try: await self.store.update_delegation_work_item( @@ -6825,7 +6821,7 @@ class CompanyWorkItemExecutor: metadata_updates={"report_attempt_count": attempt_no}, ) except Exception: - logger.debug("Best-effort report_attempt_count update failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort report_attempt_count update failed") return report_work_item async def _finalize_review_work_item(self, review_task: Task) -> None: @@ -6964,8 +6960,8 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug( - "Best-effort close of unparseable review card failed", exc_info=True + logger.opt(exception=True).debug( + "Best-effort close of unparseable review card failed" ) await self._ack_lifecycle_inbox_for_review( review_task=review_task, @@ -7094,7 +7090,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("Best-effort review work-item finalization failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort review work-item finalization failed") await self._ack_lifecycle_inbox_for_review( review_task=review_task, review_work_item_id=review_work_item_id, @@ -7108,7 +7104,7 @@ class CompanyWorkItemExecutor: try: self._signal_dispatcher_wake() except Exception: - logger.debug("_signal_dispatcher_wake failed", exc_info=True) + logger.opt(exception=True).debug("_signal_dispatcher_wake failed") if child_phase == Phase.APPROVED: await self._refresh_delegation_dependents(review_task) await self._notify_kanban_changed() @@ -7276,7 +7272,7 @@ class CompanyWorkItemExecutor: task=review_task, ) except Exception: - logger.debug("Best-effort lifecycle inbox cleanup after review failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort lifecycle inbox cleanup after review failed") @staticmethod def _resolve_max_review_reworks( @@ -7414,8 +7410,8 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug( - "verdict-parse-retry: closing prior review card failed", exc_info=True + logger.opt(exception=True).debug( + "verdict-parse-retry: closing prior review card failed" ) completion_report = str( @@ -7454,14 +7450,14 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug( - "verdict-parse-retry: extending new summary failed", exc_info=True + logger.opt(exception=True).debug( + "verdict-parse-retry: extending new summary failed" ) try: self._signal_dispatcher_wake() except Exception: - logger.debug("verdict-parse-retry: dispatcher wake failed", exc_info=True) + logger.opt(exception=True).debug("verdict-parse-retry: dispatcher wake failed") logger.info( f"verdict-parse-retry spawned: child={target_work_item_id} " @@ -7516,7 +7512,7 @@ class CompanyWorkItemExecutor: }, ) except Exception: - logger.debug("Best-effort review work-item close failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort review work-item close failed") @staticmethod def _review_feedback_version(metadata: dict[str, Any] | None) -> int: @@ -9022,7 +9018,7 @@ class CompanyWorkItemExecutor: delivery_policy["requires_user_feedback"] = True break except Exception: - logger.debug("Failed to read delivery policy from intake work-item plan", exc_info=True) + logger.opt(exception=True).debug("Failed to read delivery policy from intake work-item plan") # Owner-facing synthetic delivery cards are the stable handoff point # for follow-up directives. A projection-level false must not suppress # the human review card; review closure is an explicit runtime tool. @@ -11031,7 +11027,7 @@ class CompanyWorkItemExecutor: self._signal_dispatcher_wake() await self._notify_kanban_changed() except Exception: - logger.debug("Best-effort follow-up dependency frontier refresh failed", exc_info=True) + logger.opt(exception=True).debug("Best-effort follow-up dependency frontier refresh failed") return created_dependency_ids def _capture_environment_manifest(self, task: Task, result: TaskResult) -> None: @@ -12714,7 +12710,7 @@ class CompanyWorkItemExecutor: try: run = await self.store.get_delegation_run(run_id) except Exception: - logger.debug("Failed to load delegation run for owner review lifecycle update", exc_info=True) + logger.opt(exception=True).debug("Failed to load delegation run for owner review lifecycle update") return if run is None: return @@ -12731,7 +12727,7 @@ class CompanyWorkItemExecutor: try: await self.store.save_delegation_run(run) except Exception: - logger.debug("Failed to save delegation run owner review lifecycle update", exc_info=True) + logger.opt(exception=True).debug("Failed to save delegation run owner review lifecycle update") async def _finalize_completed_work_item(self, task: Task) -> None: if self._is_authoritative_delivery_work_item(task): diff --git a/opc/layer2_organization/company_runtime.py b/opc/layer2_organization/company_runtime.py index 66793f1..b18e36e 100644 --- a/opc/layer2_organization/company_runtime.py +++ b/opc/layer2_organization/company_runtime.py @@ -967,7 +967,7 @@ class CompanyRuntime: }, ) except Exception: - logger.debug("company runtime resume reset: role session persist failed", exc_info=True) + logger.opt(exception=True).debug("company runtime resume reset: role session persist failed") def _role_session_id(self, task: Task, *, role_id: str) -> str: explicit = str((task.metadata or {}).get("delegation_role_session_id", "") or "").strip() diff --git a/opc/layer2_organization/recruiter.py b/opc/layer2_organization/recruiter.py index 3672b33..3427025 100644 --- a/opc/layer2_organization/recruiter.py +++ b/opc/layer2_organization/recruiter.py @@ -404,7 +404,7 @@ class CompanyRecruiter: try: final_decider_role_id = str(getter() or "").strip() except Exception: - logger.debug("failed to resolve final decider role for recruiter payload", exc_info=True) + logger.opt(exception=True).debug("failed to resolve final decider role for recruiter payload") return { "final_decider_role_id": final_decider_role_id, diff --git a/opc/layer3_agent/external_broker.py b/opc/layer3_agent/external_broker.py index 89d33cf..4d60ee3 100644 --- a/opc/layer3_agent/external_broker.py +++ b/opc/layer3_agent/external_broker.py @@ -429,10 +429,9 @@ class ExternalAgentBroker: role_session_id, adapter.agent_type ) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"PR6 role adapter-state read failed " f"sid={role_session_id} agent={adapter.agent_type}", - exc_info=True, ) entry = None if isinstance(entry, dict): @@ -452,9 +451,8 @@ class ExternalAgentBroker: opc_session_id=role_session_id, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"External resume restore: get_external_session by role failed for {adapter.agent_type}/{role_session_id}", - exc_info=True, ) prior = None if prior is None: @@ -465,9 +463,8 @@ class ExternalAgentBroker: task_id=task.id, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"External resume restore: get_external_session failed for {adapter.agent_type}/{task.id}", - exc_info=True, ) return if prior is None: @@ -1387,7 +1384,7 @@ class ExternalAgentBroker: ) return str(path) except Exception: - logger.debug("ExternalAgentBroker: failed to write raw external log", exc_info=True) + logger.opt(exception=True).debug("ExternalAgentBroker: failed to write raw external log") return "" def _enrich_structured_result_artifacts( @@ -2231,8 +2228,7 @@ class ExternalAgentBroker: token_record, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"PR6 role adapter-state write failed " f"sid={role_session_id} agent={adapter.agent_type}", - exc_info=True, ) diff --git a/opc/layer4_tools/registry.py b/opc/layer4_tools/registry.py index d473f4b..8e17d48 100644 --- a/opc/layer4_tools/registry.py +++ b/opc/layer4_tools/registry.py @@ -157,7 +157,16 @@ class ToolRegistry: result = await tool.func(**call_args) output = {"result": result, "success": True} except Exception as e: - logger.error(f"Tool {name} failed ({type(e).__name__}): {e}", exc_info=True) + # Loguru has no stdlib-style ``exc_info`` kwarg: extra kwargs are + # format() arguments, which forces str.format() on the message — an + # error message containing ``{...}`` (e.g. a JSON error body) then + # raises KeyError FROM the logging call, escaping this handler and + # killing the caller instead of returning the error output below. + # Positional formatting keeps brace-containing values inert, and + # opt(exception=True) is the loguru way to log the traceback. + logger.opt(exception=True).error( + "Tool {} failed ({}): {}", name, type(e).__name__, e + ) output = { "error": str(e), "traceback": traceback.format_exc(), diff --git a/opc/plugins/office_ui/services/context.py b/opc/plugins/office_ui/services/context.py index a7f8259..1b91f2e 100644 --- a/opc/plugins/office_ui/services/context.py +++ b/opc/plugins/office_ui/services/context.py @@ -174,7 +174,7 @@ class OfficeServiceContext: try: wire(engine) except Exception: - logger.debug("Failed to wire service project engine callbacks", exc_info=True) + logger.opt(exception=True).debug("Failed to wire service project engine callbacks") return engine async def activate_project(self, project_id: str) -> Any: diff --git a/opc/plugins/office_ui/services/project.py b/opc/plugins/office_ui/services/project.py index fc94a0e..b2f9d66 100644 --- a/opc/plugins/office_ui/services/project.py +++ b/opc/plugins/office_ui/services/project.py @@ -77,7 +77,7 @@ class ProjectService: if asyncio.iscoroutine(maybe): await maybe except Exception: - logger.debug(f"Failed to close project store for {project_id}", exc_info=True) + logger.opt(exception=True).debug(f"Failed to close project store for {project_id}") async def list(self, *, active_project_id: str | None = None) -> ServiceResult: active = active_project_id or self.context.active_engine_project_id() @@ -247,7 +247,7 @@ class ProjectService: try: await active_engine.store.close() except Exception: - logger.debug("Failed to close active project store before delete", exc_info=True) + logger.opt(exception=True).debug("Failed to close active project store before delete") shutil.rmtree(str(projects_dir), ignore_errors=True) workplace = self.context.project_workplace(project_id) @@ -263,7 +263,7 @@ class ProjectService: if asyncio.iscoroutine(maybe): await maybe except Exception: - logger.debug(f"memory.delete_project failed for {project_id}", exc_info=True) + logger.opt(exception=True).debug(f"memory.delete_project failed for {project_id}") events = [ServiceEvent("project_deleted", {"project_id": project_id})] payload: dict[str, Any] = {"project_id": project_id, "deleted_channels": deleted_channels} diff --git a/opc/plugins/office_ui/services/session.py b/opc/plugins/office_ui/services/session.py index 5ee8769..aa4dfc4 100644 --- a/opc/plugins/office_ui/services/session.py +++ b/opc/plugins/office_ui/services/session.py @@ -267,7 +267,7 @@ class SessionService: try: await store.save_task(task) except Exception: - logger.debug("failed to mark company runtime stop state", exc_info=True) + logger.opt(exception=True).debug("failed to mark company runtime stop state") async def _clear_company_runtime_stop_state(self, *, engine: Any, task_ids: list[str]) -> None: store = getattr(engine, "store", None) @@ -297,7 +297,7 @@ class SessionService: try: await store.save_task(task) except Exception: - logger.debug("failed to clear company runtime stop state", exc_info=True) + logger.opt(exception=True).debug("failed to clear company runtime stop state") def _normalize_requested_config( self, @@ -445,7 +445,7 @@ class SessionService: ) events.append(ServiceEvent("collab_sync_push", collab)) except Exception: - logger.warning("create_session collab_sync build failed", exc_info=True) + logger.opt(exception=True).warning("create_session collab_sync build failed") return ServiceResult(session_payload, events) def _session_metadata( @@ -606,7 +606,7 @@ class SessionService: if int(message_count or 0) > 0: return "message_history" except Exception: - logger.debug("Failed to inspect session message count for config lock", exc_info=True) + logger.opt(exception=True).debug("Failed to inspect session message count for config lock") status = getattr(getattr(task, "status", None), "value", getattr(task, "status", None)) status_value = str(status or "").strip().lower() if status_value and status_value != "pending": @@ -958,7 +958,7 @@ class SessionService: stop_intent_id=stop_intent_id, ) except Exception: - logger.warning("suspend_company_runtime failed during service stop", exc_info=True) + logger.opt(exception=True).warning("suspend_company_runtime failed during service stop") if suspended is not None: for candidate in list(suspended.get("task_ids", []) or []): candidate_id = str(candidate or "").strip() @@ -1000,7 +1000,7 @@ class SessionService: }, ) except Exception: - logger.debug("failed to insert company runtime stop system message", exc_info=True) + logger.opt(exception=True).debug("failed to insert company runtime stop system message") payload = { **default_payload, "status": "suspended", diff --git a/opc/plugins/office_ui/snapshot_builder.py b/opc/plugins/office_ui/snapshot_builder.py index 06f03b4..c4d289a 100644 --- a/opc/plugins/office_ui/snapshot_builder.py +++ b/opc/plugins/office_ui/snapshot_builder.py @@ -2630,7 +2630,7 @@ async def _build_company_runtime_control_by_task( if sid and sid not in checkpoints_by_session: checkpoints_by_session[sid] = checkpoint except Exception: - logger.debug("snapshot: failed to load company runtime checkpoints", exc_info=True) + logger.opt(exception=True).debug("snapshot: failed to load company runtime checkpoints") result: dict[str, dict[str, Any]] = {} for parent_session_id, group in tasks_by_parent_session.items(): @@ -2892,7 +2892,7 @@ async def build_project_index_sync( try: tasks = await engine.store.get_tasks(project_id=project_id) except Exception: - logger.warning("Failed to load tasks for project index", exc_info=True) + logger.opt(exception=True).warning("Failed to load tasks for project index") existing_task_ids = { str(getattr(task, "id", "") or "").strip() @@ -3286,7 +3286,7 @@ async def build_collab_sync( event_adapter=event_adapter, ) except Exception: - logger.warning("Failed to build company-mode kanban projection", exc_info=True) + logger.opt(exception=True).warning("Failed to build company-mode kanban projection") # Per-session DelegationWorkItem rollups produced alongside the company # kanban projection. Used by the work-item-driven Execution Progress # panel — keyed by ``task.session_id`` (matches the formatted_sessions @@ -3308,7 +3308,7 @@ async def build_collab_sync( if backfilled: logger.info(f"Reconciled {backfilled} total messages from engine → ChatStore") except Exception: - logger.warning("Session reconciliation failed (non-fatal)", exc_info=True) + logger.opt(exception=True).warning("Session reconciliation failed (non-fatal)") existing_task_ids = { str(getattr(task, "id", "") or "").strip() diff --git a/opc/plugins/office_ui/ws_handler.py b/opc/plugins/office_ui/ws_handler.py index 238b473..c9c57e4 100644 --- a/opc/plugins/office_ui/ws_handler.py +++ b/opc/plugins/office_ui/ws_handler.py @@ -703,9 +703,8 @@ class WSHandler: if not already_subscribed: event_bus.subscribe_all(runtime_event_callback) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"Failed to wire UI callbacks for project engine {getattr(engine, 'project_id', None)!r}", - exc_info=True, ) @staticmethod @@ -823,7 +822,7 @@ class WSHandler: try: summary = await store.reset_orphan_running_tasks(lease_seconds=lease_seconds) except Exception: - logger.warning("heal_orphan_tasks_on_boot: reset_orphan_running_tasks failed", exc_info=True) + logger.opt(exception=True).warning("heal_orphan_tasks_on_boot: reset_orphan_running_tasks failed") return reset_count = summary.get("statuses_reset", 0) locks_cleared = summary.get("locks_cleared", 0) @@ -861,9 +860,8 @@ class WSHandler: try: still_ours = await store.renew_task_lock(task_id) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"Heartbeat (initial) renew_task_lock failed for {task_id}", - exc_info=True, ) still_ours = True if not still_ours: @@ -873,9 +871,8 @@ class WSHandler: try: still_ours = await store.renew_task_lock(task_id) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"Heartbeat renew_task_lock failed for {task_id}", - exc_info=True, ) continue if not still_ours: @@ -3165,7 +3162,7 @@ class WSHandler: except Exception: pass else: - logger.error(f"WS handler error for {msg_type}: {type(e).__name__}: {e!r}", exc_info=True) + logger.opt(exception=True).error(f"WS handler error for {msg_type}: {type(e).__name__}: {e!r}") try: await self._send_ack(ws, ok=False, error=str(e) or type(e).__name__, action=msg_type) except Exception: @@ -3238,16 +3235,14 @@ class WSHandler: except asyncio.CancelledError: raise except Exception as exc: - logger.warning( + logger.opt(exception=True).warning( f"Project index sent, but snapshot refresh failed for {project_id}: {type(exc).__name__}: {exc!r}", - exc_info=True, ) except asyncio.CancelledError: raise except Exception as exc: - logger.error( + logger.opt(exception=True).error( f"Failed to build project index for {project_id}: {type(exc).__name__}: {exc!r}", - exc_info=True, ) if send_error_ack and self._client_active_project_id(ws) == project_id: await self._send_ack( @@ -3288,7 +3283,7 @@ class WSHandler: except asyncio.CancelledError: raise except Exception: - logger.warning("Initial websocket collab_sync push failed", exc_info=True) + logger.opt(exception=True).warning("Initial websocket collab_sync push failed") try: org_info = await self._build_org_info_payload() if self._client_active_project_id(ws) == project_id: @@ -3296,7 +3291,7 @@ class WSHandler: except asyncio.CancelledError: raise except Exception: - logger.warning("Initial websocket org_info push failed", exc_info=True) + logger.opt(exception=True).warning("Initial websocket org_info push failed") async def _handle_collab_sync(self, ws: Any, data: dict) -> None: engine, project_id = await self._engine_for_request(data) @@ -3457,7 +3452,7 @@ class WSHandler: event_adapter=self.event_adapter, ) except Exception: - logger.warning("Failed to load company kanban projection for kanban_switch_view", exc_info=True) + logger.opt(exception=True).warning("Failed to load company kanban projection for kanban_switch_view") if self._exec_mode in {"company", "org", "custom"} and not company_columns: company_boards = [{ "board_id": project_id, @@ -4127,7 +4122,7 @@ class WSHandler: if not bg_task.done(): bg_task.cancel() except Exception: - logger.debug(f"Failed to cancel background task for {task_id}", exc_info=True) + logger.opt(exception=True).debug(f"Failed to cancel background task for {task_id}") @staticmethod def _is_ws_disconnect_error(exc: BaseException) -> bool: @@ -4375,7 +4370,7 @@ class WSHandler: return exc = task.exception() if exc is not None: - logger.error(f"Background task failed: {exc}", exc_info=exc) + logger.opt(exception=exc).error(f"Background task failed: {exc}") # Notify frontend for session tasks so UI doesn't stay stuck on "thinking" context = self._task_bg_context.get(task) or {} task_id = self._find_task_id_for_bg_task(task) @@ -4604,7 +4599,7 @@ class WSHandler: name=f"task-heartbeat:{task_id}", ) except Exception: - logger.debug("failed to mark run_task company runtime running", exc_info=True) + logger.opt(exception=True).debug("failed to mark run_task company runtime running") response = await engine.process_message( content, project_id=pid, @@ -4642,9 +4637,9 @@ class WSHandler: idle_target = await self._resolve_company_runtime_target(task_id, engine=engine) await self._set_company_runtime_control(idle_target or company_runtime_target, state="idle") except Exception: - logger.debug("failed to mark run_task company runtime idle", exc_info=True) + logger.opt(exception=True).debug("failed to mark run_task company runtime idle") except Exception as e: - logger.error(f"Task execution error: {e}", exc_info=True) + logger.opt(exception=True).error(f"Task execution error: {e}") # Broadcast: task failed (stays in in-progress, user can retry) if task_id: await self.broadcast({"type": "board_task_status_changed", "payload": { @@ -4655,7 +4650,7 @@ class WSHandler: failed_target = await self._resolve_company_runtime_target(task_id, engine=engine) await self._set_company_runtime_control(failed_target or company_runtime_target, state="idle") except Exception: - logger.debug("failed to clear run_task company runtime after error", exc_info=True) + logger.opt(exception=True).debug("failed to clear run_task company runtime after error") err_meta: dict[str, Any] = {} if task_id: err_meta["task_id"] = task_id @@ -4685,7 +4680,7 @@ class WSHandler: ) await self.broadcast({"type": "collab_sync_push", "payload": collab}) except Exception: - logger.warning("Post-run collab_sync broadcast failed (non-fatal)", exc_info=True) + logger.opt(exception=True).warning("Post-run collab_sync broadcast failed (non-fatal)") async def _mirror_agent_message( self, @@ -4846,9 +4841,8 @@ class WSHandler: project_id=project_id, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"Failed to update human escalation checkpoint status for {normalized_escalation_id}", - exc_info=True, ) return None if updated is not None: @@ -4908,9 +4902,8 @@ class WSHandler: project_id=project_id, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"Failed to load unresolved human escalation cards for {channel_id}", - exc_info=True, ) return [] @@ -4955,9 +4948,8 @@ class WSHandler: try: cards = await getter(channel_id, project_id=project_id) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"Failed to load unresolved execution checkpoint cards for {channel_id}", - exc_info=True, ) return [] @@ -4990,9 +4982,8 @@ class WSHandler: project_id=project_id, ) except Exception: - logger.debug( + logger.opt(exception=True).debug( f"Failed to reconcile execution checkpoint card {checkpoint_id}", - exc_info=True, ) continue if updated is not None: @@ -5102,7 +5093,7 @@ class WSHandler: project_id=project_id, ) except Exception: - logger.debug(f"session_detail: transcript page load failed for {task_id}", exc_info=True) + logger.opt(exception=True).debug(f"session_detail: transcript page load failed for {task_id}") transcript_total_count = 0 transcript_has_more = False @@ -5183,14 +5174,14 @@ class WSHandler: try: project_tasks = await store.get_tasks(project_id=project_id) except Exception: - logger.debug("session_detail: failed to load project tasks for runtime control", exc_info=True) + logger.opt(exception=True).debug("session_detail: failed to load project tasks for runtime control") project_tasks = [task] try: runtime_control_meta = ( await _build_company_runtime_control_by_task(run_engine, project_tasks, project_id) ).get(task_id, {}) except Exception: - logger.debug("session_detail: failed to build runtime control payload", exc_info=True) + logger.opt(exception=True).debug("session_detail: failed to build runtime control payload") runtime_control_meta = {} runtime_meta = dict(task_meta.get("runtime_v2", {}) or {}) member_session_meta = dict(task_meta.get("member_session_state", {}) or {}) @@ -6022,7 +6013,7 @@ class WSHandler: try: await store.save_task(task) except Exception: - logger.debug("failed to mark company runtime stop state", exc_info=True) + logger.opt(exception=True).debug("failed to mark company runtime stop state") async def _clear_company_runtime_stop_state( self, @@ -6060,7 +6051,7 @@ class WSHandler: try: await store.save_task(task) except Exception: - logger.debug("failed to clear company runtime stop state", exc_info=True) + logger.opt(exception=True).debug("failed to clear company runtime stop state") async def _finalize_company_runtime_stop(self, target: dict[str, Any], *, stop_intent_id: str) -> None: runtime_engine = target.get("engine") or self.engine @@ -6081,7 +6072,7 @@ class WSHandler: stop_intent_id=stop_intent_id, ) except Exception: - logger.warning(f"suspend_company_runtime failed for {origin_task_id}", exc_info=True) + logger.opt(exception=True).warning(f"suspend_company_runtime failed for {origin_task_id}") if suspended is not None: for candidate in list(suspended.get("task_ids", []) or []): @@ -6101,7 +6092,7 @@ class WSHandler: stop_intent_id=stop_intent_id, ) except Exception: - logger.debug("failed to mark company runtime fully suspended", exc_info=True) + logger.opt(exception=True).debug("failed to mark company runtime fully suspended") try: await self._set_company_runtime_control( target, @@ -6110,7 +6101,7 @@ class WSHandler: stop_intent_id=stop_intent_id, ) except Exception: - logger.debug("failed to broadcast company runtime suspended state", exc_info=True) + logger.opt(exception=True).debug("failed to broadcast company runtime suspended state") channel_id = f"session:{origin_task_id or target.get('parent_task_id', '')}" pid = self._normalize_project_id(getattr(runtime_engine, "project_id", None)) try: @@ -6130,12 +6121,12 @@ class WSHandler: ) await self.broadcast({"type": "session_message", "payload": msg}) except Exception: - logger.warning(f"Failed to insert suspend message for {origin_task_id}", exc_info=True) + logger.opt(exception=True).warning(f"Failed to insert suspend message for {origin_task_id}") else: try: await self._clear_company_runtime_stop_state(target, stop_intent_id=stop_intent_id) except Exception: - logger.debug("failed to clear company runtime stop state after suspend failure", exc_info=True) + logger.opt(exception=True).debug("failed to clear company runtime stop state after suspend failure") try: await self._set_company_runtime_control( target, @@ -6160,7 +6151,7 @@ class WSHandler: try: task = await store.get_task(task_id) except Exception: - logger.debug(f"Failed to load task for stop: {task_id}", exc_info=True) + logger.opt(exception=True).debug(f"Failed to load task for stop: {task_id}") if task is None: await self._send_ack(ws, ok=False, error="task_not_found", project_id=run_project_id, task_id=task_id) return @@ -6170,7 +6161,7 @@ class WSHandler: try: target = await self._resolve_company_runtime_target(task_id, engine=run_engine) except Exception: - logger.warning(f"failed to resolve company runtime stop target for {task_id}", exc_info=True) + logger.opt(exception=True).warning(f"failed to resolve company runtime stop target for {task_id}") target = None if target is not None: parent_session_id = str(target.get("parent_session_id", "") or "").strip() @@ -6178,7 +6169,7 @@ class WSHandler: try: existing_checkpoint = await run_engine.get_pending_company_runtime_suspend_checkpoint(parent_session_id) except Exception: - logger.debug("failed to check existing company suspend checkpoint", exc_info=True) + logger.opt(exception=True).debug("failed to check existing company suspend checkpoint") existing_intent = self._company_stop_intents.get(parent_session_id) existing_finalizer = self._company_stop_finalize_tasks.get(parent_session_id) if existing_checkpoint is not None: @@ -6227,7 +6218,7 @@ class WSHandler: try: all_task_ids = await self._cancel_task_tree(task_id, preserve_history=True, store=store) except Exception: - logger.warning(f"_cancel_task_tree failed for {task_id}", exc_info=True) + logger.opt(exception=True).warning(f"_cancel_task_tree failed for {task_id}") all_task_ids = [task_id] self._stop_requested_task_ids.update(all_task_ids) @@ -6257,7 +6248,7 @@ class WSHandler: stop_payload["agent_id"] = resolved_agent await self.broadcast({"type": "agent_runtime_update", "payload": stop_payload}) except Exception: - logger.warning(f"Failed to broadcast stop status for {tid}", exc_info=True) + logger.opt(exception=True).warning(f"Failed to broadcast stop status for {tid}") # Insert system message (only for the primary task — children are internal) channel_id = f"session:{task_id}" @@ -6273,7 +6264,7 @@ class WSHandler: ) await self.broadcast({"type": "session_message", "payload": msg}) except Exception: - logger.warning(f"Failed to insert stop message for {task_id}", exc_info=True) + logger.opt(exception=True).warning(f"Failed to insert stop message for {task_id}") await self._send_ack(ws, ok=True) async def _handle_session_resume(self, ws: Any, data: dict) -> None: @@ -6298,7 +6289,7 @@ class WSHandler: try: task = await run_engine.store.get_task(task_id) except Exception: - logger.warning(f"session_resume: get_task failed for {task_id}", exc_info=True) + logger.opt(exception=True).warning(f"session_resume: get_task failed for {task_id}") if task is None: await self._send_ack(ws, ok=False, error="task_not_found", project_id=run_project_id, task_id=task_id) return @@ -6311,7 +6302,7 @@ class WSHandler: try: target = await self._resolve_company_runtime_target(task_id, engine=run_engine) except Exception: - logger.warning(f"session_resume: failed to resolve company runtime target for {task_id}", exc_info=True) + logger.opt(exception=True).warning(f"session_resume: failed to resolve company runtime target for {task_id}") target = None if target is not None: parent_session_id = str(target.get("parent_session_id", "") or "").strip() @@ -6325,12 +6316,12 @@ class WSHandler: await self._send_ack(ws, ok=False, error="stop_finalize_in_progress") return except Exception: - logger.debug("session_resume: stop finalizer ended with error", exc_info=True) + logger.opt(exception=True).debug("session_resume: stop finalizer ended with error") self._company_stop_intents.pop(parent_session_id, None) try: await self._set_company_runtime_control(target, state="resuming") except Exception: - logger.debug("session_resume: failed to broadcast resuming state", exc_info=True) + logger.opt(exception=True).debug("session_resume: failed to broadcast resuming state") if not session_id: await self._send_ack(ws, ok=False, error="session_not_found") return @@ -6373,7 +6364,7 @@ class WSHandler: except TypeError: checkpoints = await getter(project_id) except Exception: - logger.debug("failed to list checkpoints for explicit reply routing", exc_info=True) + logger.opt(exception=True).debug("failed to list checkpoints for explicit reply routing") checkpoints = [] for checkpoint in list(checkpoints or []): if str(getattr(checkpoint, "checkpoint_id", "") or "").strip() != normalized_checkpoint_id: @@ -6391,7 +6382,7 @@ class WSHandler: maybe_checkpoint = direct_lookup(normalized_checkpoint_id) checkpoint = await maybe_checkpoint if inspect.isawaitable(maybe_checkpoint) else maybe_checkpoint except Exception: - logger.debug("failed to load checkpoint by id for explicit reply routing", exc_info=True) + logger.opt(exception=True).debug("failed to load checkpoint by id for explicit reply routing") checkpoint = None if checkpoint is not None: if str(getattr(checkpoint, "checkpoint_id", "") or "").strip() != normalized_checkpoint_id: @@ -6428,7 +6419,7 @@ class WSHandler: try: target = await self._resolve_company_runtime_target(candidate_task_id, engine=engine) except Exception: - logger.debug("failed to resolve company runtime target for delivery feedback", exc_info=True) + logger.opt(exception=True).debug("failed to resolve company runtime target for delivery feedback") target = None if not target: continue @@ -6511,7 +6502,7 @@ class WSHandler: if visible is True: return True except Exception: - logger.debug("failed to evaluate delivery feedback checkpoint visibility", exc_info=True) + logger.opt(exception=True).debug("failed to evaluate delivery feedback checkpoint visibility") checkpoint_session_id = str(getattr(checkpoint, "session_id", "") or "").strip() if checkpoint_session_id == requested_session_id: @@ -6539,7 +6530,7 @@ class WSHandler: try: waiting_task = await get_task(waiting_task_id) except Exception: - logger.debug("failed to load waiting task for delivery feedback visibility", exc_info=True) + logger.opt(exception=True).debug("failed to load waiting task for delivery feedback visibility") waiting_task = None if waiting_task is not None: waiting_session_id = str(getattr(waiting_task, "session_id", "") or "").strip() @@ -6690,10 +6681,9 @@ class WSHandler: created_at=self._checkpoint_created_timestamp(checkpoint), ) except Exception: - logger.debug( + logger.opt(exception=True).debug( "failed to insert terminal synthetic checkpoint card; retrying status update: checkpoint_id={}", normalized_checkpoint_id, - exc_info=True, ) updated = await self.chat_store.update_checkpoint_status( normalized_checkpoint_id, @@ -6736,7 +6726,7 @@ class WSHandler: checkpoint_types=["company_delivery_feedback"], ) except Exception: - logger.debug("failed to list pending delivery feedback checkpoints", exc_info=True) + logger.opt(exception=True).debug("failed to list pending delivery feedback checkpoints") return [] superseded_ids: list[str] = [] @@ -6802,7 +6792,7 @@ class WSHandler: waiting_task.status = TaskStatus.DONE await store.save_task(waiting_task) except Exception: - logger.debug("failed to mark delivery feedback checkpoint superseded", exc_info=True) + logger.opt(exception=True).debug("failed to mark delivery feedback checkpoint superseded") continue superseded_ids.append(checkpoint_id) @@ -7014,7 +7004,7 @@ class WSHandler: try: waiting_task = await store.get_task(waiting_task_id) except Exception: - logger.debug("failed to load delivery feedback waiting task", exc_info=True) + logger.opt(exception=True).debug("failed to load delivery feedback waiting task") target = await self._company_delivery_feedback_parent_target( task_id=task_id, waiting_task_id=waiting_task_id, @@ -7082,7 +7072,7 @@ class WSHandler: try: parent_task = await run_engine.store.get_task(parent_task_id) except Exception: - logger.debug("failed to load parent task for delivery feedback reply", exc_info=True) + logger.opt(exception=True).debug("failed to load parent task for delivery feedback reply") 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 = "" @@ -7163,7 +7153,7 @@ class WSHandler: ) await self.broadcast({"type": "session_message", "payload": msg}) except Exception: - logger.debug("failed to write delivery feedback reply error", exc_info=True) + logger.opt(exception=True).debug("failed to write delivery feedback reply error") finally: if self._chat_store_is_ready(self.chat_store): await self._flush_progress(parent_task_id, project_id=pid) @@ -7207,7 +7197,7 @@ class WSHandler: try: target = await self._resolve_company_runtime_target(task_id, engine=run_engine) except Exception: - logger.debug("failed to resolve company suspend reply target", exc_info=True) + logger.opt(exception=True).debug("failed to resolve company suspend reply target") return False if target is None: return False @@ -7231,7 +7221,7 @@ class WSHandler: await self.broadcast({"type": "session_message", "payload": helper}) return True except Exception: - logger.debug("company stop finalizer failed before follow-up routing", exc_info=True) + logger.opt(exception=True).debug("company stop finalizer failed before follow-up routing") checkpoint = None get_checkpoint = getattr(run_engine, "get_active_company_runtime_suspend_checkpoint", None) @@ -7239,7 +7229,7 @@ class WSHandler: try: checkpoint = await get_checkpoint(parent_session_id) except Exception: - logger.debug("failed to load active company suspend checkpoint", exc_info=True) + logger.opt(exception=True).debug("failed to load active company suspend checkpoint") if checkpoint is None: return False if str(getattr(checkpoint, "status", "") or "").strip() != "pending": @@ -7296,14 +7286,14 @@ class WSHandler: try: await self._set_company_runtime_control(target, state="resuming") except Exception: - logger.debug("failed to broadcast company suspend reply routing state", exc_info=True) + logger.opt(exception=True).debug("failed to broadcast company suspend reply routing state") parent_task = None if self._store_is_ready(run_engine.store): try: parent_task = await run_engine.store.get_task(parent_task_id) except Exception: - logger.debug("failed to load parent task for company suspend reply", exc_info=True) + logger.opt(exception=True).debug("failed to load parent task for company suspend reply") 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 = "" @@ -7373,9 +7363,8 @@ class WSHandler: chat_exc, ) else: - logger.debug( + logger.opt(exception=True).debug( "Failed to write company suspend reply error message", - exc_info=True, ) finally: if self._chat_store_is_ready(self.chat_store): @@ -7416,7 +7405,7 @@ class WSHandler: inferred_ws = str(Path(comms_root).parent) return _comms.resolve_layout(inferred_ws, project_id, session_id).root except Exception: - logger.debug(f"_resolve_task_comms_dir failed for task {getattr(task, 'id', '?')}", exc_info=True) + logger.opt(exception=True).debug(f"_resolve_task_comms_dir failed for task {getattr(task, 'id', '?')}") return None async def _handle_session_delete(self, ws: Any, data: dict) -> None: @@ -7534,7 +7523,7 @@ class WSHandler: metadata=assistant_turn_meta or None, ) except Exception as e: - logger.error(f"Dispatcher error, falling back to engine: {e}", exc_info=True) + logger.opt(exception=True).error(f"Dispatcher error, falling back to engine: {e}") process_kwargs = { "session_id": session_id, "attachment_refs": attachment_refs, @@ -7568,7 +7557,7 @@ class WSHandler: try: checkpoint = await direct_lookup(normalized_checkpoint_id) except Exception: - logger.debug("failed to load checkpoint by id from engine", exc_info=True) + logger.opt(exception=True).debug("failed to load checkpoint by id from engine") if checkpoint is None: store = getattr(engine, "store", None) getter = getattr(store, "get_execution_checkpoints", None) @@ -7689,7 +7678,7 @@ class WSHandler: broadcast_update=False, ) except Exception: - logger.debug("failed to persist checkpoint card terminal state", exc_info=True) + logger.opt(exception=True).debug("failed to persist checkpoint card terminal state") return None async def _process_session_message( @@ -7776,7 +7765,7 @@ class WSHandler: company_runtime_target = await self._resolve_company_runtime_target(task_id, engine=engine) await self._set_company_runtime_control(company_runtime_target, state="running") except Exception: - logger.debug("failed to mark company session runtime running", exc_info=True) + logger.opt(exception=True).debug("failed to mark company session runtime running") try: engine_mode, company_profile = self._resolve_engine_mode( session_exec_mode, @@ -7856,7 +7845,7 @@ class WSHandler: idle_target = await self._resolve_company_runtime_target(task_id, engine=engine) await self._set_company_runtime_control(idle_target or company_runtime_target, state="idle") except Exception: - logger.debug("failed to mark company session runtime idle", exc_info=True) + logger.opt(exception=True).debug("failed to mark company session runtime idle") except asyncio.CancelledError: cancelled_ids = await self._mark_task_tree_cancelled_if_active(task_id, store=engine.store) for cancelled_id in cancelled_ids: @@ -7891,7 +7880,7 @@ class WSHandler: failed_target = await self._resolve_company_runtime_target(task_id, engine=engine) await self._set_company_runtime_control(failed_target or company_runtime_target, state="idle") except Exception: - logger.debug("failed to clear company session runtime after error", exc_info=True) + logger.opt(exception=True).debug("failed to clear company session runtime after error") finally: # Stop heartbeat before anything else so we don't keep bumping # execution_locked_at for a task we've just finished handling. @@ -8497,7 +8486,7 @@ class WSHandler: ) await self.broadcast({"type": "session_message", "payload": msg}) except Exception as e: - logger.error(f"Secretary processing error: {e}", exc_info=True) + logger.opt(exception=True).error(f"Secretary processing error: {e}") msg = await self.chat_store.insert_message( channel_id=secretary_channel, sender="system", @@ -8547,9 +8536,8 @@ class WSHandler: await self._send_service_error(ws, exc, action="switch_project") return except Exception as exc: - logger.error( + logger.opt(exception=True).error( f"Failed to prepare project switch for {new_id}: {type(exc).__name__}: {exc!r}", - exc_info=True, ) await self._send_ack( ws, @@ -8831,7 +8819,7 @@ class WSHandler: list(getattr(org, "employees", []) or []), ) except Exception: - logger.debug("Failed to write employee registry for active org", exc_info=True) + logger.opt(exception=True).debug("Failed to write employee registry for active org") payload = build_org_config_payload_from_config( config, organization_id=organization_id, @@ -9016,7 +9004,7 @@ class WSHandler: list(validated_config.org.employees), ) except Exception: - logger.debug("Failed to load employee registry for applied org config", exc_info=True) + logger.opt(exception=True).debug("Failed to load employee registry for applied org config") roles_before = {r.id for r in existing.org.roles} roles_after = {r.id for r in validated_config.org.roles} employees_before = len(existing.org.employees)