diff --git a/backend/alembic/versions/e4f5a6b7c8d9_canonical_task_log_events.py b/backend/alembic/versions/e4f5a6b7c8d9_canonical_task_log_events.py new file mode 100644 index 000000000..e99f6fe95 --- /dev/null +++ b/backend/alembic/versions/e4f5a6b7c8d9_canonical_task_log_events.py @@ -0,0 +1,44 @@ +"""replace legacy task log fields with canonical Molecular CoT fields""" + +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op + +revision: str = "e4f5a6b7c8d9" +down_revision: str | Sequence[str] | None = "e3f4a5b6c7d8" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.add_column("task_logs", sa.Column("trace_id", sa.String(64), nullable=False, server_default="")) + op.add_column("task_logs", sa.Column("span_id", sa.String(128), nullable=True)) + op.add_column("task_logs", sa.Column("src", sa.String(255), nullable=False, server_default="task.system")) + op.add_column("task_logs", sa.Column("marker", sa.String(16), nullable=False, server_default="REASON")) + op.add_column("task_logs", sa.Column("intent", sa.Text(), nullable=False, server_default="")) + op.add_column("task_logs", sa.Column("payload", sa.JSON(), nullable=True)) + op.add_column("task_logs", sa.Column("error", sa.Text(), nullable=True)) + bind = op.get_bind() + bind.execute(sa.text(""" + UPDATE task_logs + SET src = COALESCE(source, 'task.system'), + intent = COALESCE(message, ''), + marker = CASE WHEN UPPER(level) IN ('WARNING', 'ERROR') THEN 'EXPLORE' ELSE 'REASON' END, + error = CASE WHEN UPPER(level) IN ('WARNING', 'ERROR') THEN COALESCE(message, 'Task event') ELSE NULL END + """)) + if bind.dialect.name == "postgresql": + bind.execute(sa.text("UPDATE task_logs SET payload = metadata_json::json WHERE metadata_json IS NOT NULL")) + elif bind.dialect.name == "sqlite": + bind.execute(sa.text("UPDATE task_logs SET payload = json(metadata_json) WHERE metadata_json IS NOT NULL AND json_valid(metadata_json)")) + op.create_index("ix_task_logs_task_src", "task_logs", ["task_id", "src"]) + op.drop_index("ix_task_logs_task_source", table_name="task_logs") + op.drop_column("task_logs", "source") + op.drop_column("task_logs", "message") + op.drop_column("task_logs", "metadata_json") + for column in ("trace_id", "src", "marker", "intent"): + op.alter_column("task_logs", column, server_default=None) + + +def downgrade() -> None: + raise NotImplementedError("Canonical task log migration is a breaking migration") diff --git a/backend/src/api/routes/__tests__/test_tasks_crud.py b/backend/src/api/routes/__tests__/test_tasks_crud.py index d7d99bfdb..01710f39e 100644 --- a/backend/src/api/routes/__tests__/test_tasks_crud.py +++ b/backend/src/api/routes/__tests__/test_tasks_crud.py @@ -167,8 +167,8 @@ def test_get_task_omits_logs_by_default(crud_client): tc, tm = crud_client mock_task = _make_task("task-logs", "superset-backup", TaskStatus.SUCCESS) mock_task.logs = [ - LogEntry(level="INFO", message="line-1", source="plugin"), - LogEntry(level="INFO", message="line-2", source="plugin"), + LogEntry(level="INFO", trace_id="trace-1", src="task.plugin", marker="REASON", intent="line-1"), + LogEntry(level="INFO", trace_id="trace-1", src="task.plugin", marker="REASON", intent="line-2"), ] tm.get_task.return_value = mock_task diff --git a/backend/src/app.py b/backend/src/app.py index aedcf6f4a..7fc17cc98 100755 --- a/backend/src/app.py +++ b/backend/src/app.py @@ -82,7 +82,7 @@ from .api.routes import ( ) from .api.routes.validation_tasks import router as validation_tasks from .core.auth.security import get_password_hash -from ss_tools.shared.cot_logger import get_trace_id, seed_trace_id, set_trace_id +from ss_tools.shared.cot_logger import build_cot_event, get_trace_id, seed_trace_id, set_trace_id from .core.database import AuthSessionLocal, init_db from .core.encryption_key import ensure_encryption_key from .core.logger import belief_scope, logger @@ -949,7 +949,7 @@ async def websocket_endpoint(websocket: WebSocket, task_id: str, source: str = N def matches_filters(log_entry) -> bool: """Check if log entry matches the filter criteria.""" - log_source = getattr(log_entry, "source", None) + log_source = getattr(log_entry, "src", None) or getattr(log_entry, "source", None) if source_filter and str(log_source or "").lower() != source_filter: return False if level_filter: @@ -1014,12 +1014,13 @@ async def websocket_endpoint(websocket: WebSocket, task_id: str, source: str = N # ── Send synthetic AWAITING_INPUT prompt if needed ── task = task_manager.get_task(task_id) if task and task.status == "AWAITING_INPUT" and task.input_request: - synthetic_log = { - "timestamp": task.logs[-1].timestamp.isoformat() if task.logs else "2024-01-01T00:00:00", - "level": "INFO", - "message": "Task paused for user input (Connection Re-established)", - "context": {"input_request": task.input_request}, - } + synthetic_log = build_cot_event( + src="task.lifecycle.websocket_reconnect", + marker="REFLECT", + intent="Task paused for user input (connection re-established)", + payload={"input_request": task.input_request}, + ) + synthetic_log = {"task_id": task_id, **synthetic_log} await websocket.send_json(synthetic_log) logger.reason( "Replayed awaiting-input prompt to restored WebSocket client", @@ -1070,10 +1071,11 @@ async def websocket_endpoint(websocket: WebSocket, task_id: str, source: str = N "level": log_dict.get("level"), }, ) - if "Task completed successfully" in result.message or "Task failed" in result.message: + result_intent = getattr(result, "intent", None) or getattr(result, "message", "") + if "Task completed successfully" in result_intent or "Task failed" in result_intent: logger.reason( "Observed terminal task log entry; delaying to preserve client visibility", - payload={"task_id": task_id, "message": result.message}, + payload={"task_id": task_id, "intent": result_intent}, ) await asyncio.sleep(2) except (WebSocketDisconnect, StopIteration) as _ws_exc: diff --git a/backend/src/core/task_manager/__tests__/test_canonical_cot_event.py b/backend/src/core/task_manager/__tests__/test_canonical_cot_event.py new file mode 100644 index 000000000..b32de384f --- /dev/null +++ b/backend/src/core/task_manager/__tests__/test_canonical_cot_event.py @@ -0,0 +1,18 @@ +import pytest + +from ss_tools.shared.cot_logger import build_cot_event + + +def test_builder_creates_protocol_complete_event(): + event = build_cot_event( + src="task.migration.recover", marker="REFLECT", intent="Recovery verified", + payload={"dashboard_id": "d1"}, trace_id="trace-1", + ) + assert event["marker"] == "REFLECT" + assert event["trace_id"] == "trace-1" + assert event["payload"] == {"dashboard_id": "d1"} + + +def test_builder_rejects_explore_without_error(): + with pytest.raises(ValueError, match="EXPLORE events require error"): + build_cot_event(src="task.migration.recover", marker="EXPLORE", intent="Fallback") diff --git a/backend/src/core/task_manager/__tests__/test_context.py b/backend/src/core/task_manager/__tests__/test_context.py index 4a2811c35..5ed2cf354 100644 --- a/backend/src/core/task_manager/__tests__/test_context.py +++ b/backend/src/core/task_manager/__tests__/test_context.py @@ -36,8 +36,8 @@ async def test_heartbeat_emits_structured_log(): captured = [] - async def fake_add_log(task_id, level, message, source="system", metadata=None, **_): - captured.append({"task_id": task_id, "level": level, "message": message, "metadata": metadata}) + async def fake_add_log(task_id, event): + captured.append({"task_id": task_id, "event": event}) ctx = TaskContext( task_id="task-heartbeat-1", @@ -50,9 +50,10 @@ async def test_heartbeat_emits_structured_log(): assert len(captured) == 1 log = captured[0] assert log["task_id"] == "task-heartbeat-1" - assert log["metadata"]["type"] == "heartbeat" - assert log["metadata"]["progress"] == 0.42 - assert log["metadata"]["message"] == "processing batch 5/12" - assert log["metadata"]["foo"] == "bar" + assert log["event"]["marker"] == "REFLECT" + assert log["event"]["payload"]["type"] == "heartbeat" + assert log["event"]["payload"]["progress"] == 0.42 + assert log["event"]["payload"]["message"] == "processing batch 5/12" + assert log["event"]["payload"]["foo"] == "bar" # #endregion Test.Tests.TestHeartbeatEmitsStructuredLog # #endregion Test.Tests.TestContext diff --git a/backend/src/core/task_manager/__tests__/test_task_logger.py b/backend/src/core/task_manager/__tests__/test_task_logger.py index e1699d937..9b88995db 100644 --- a/backend/src/core/task_manager/__tests__/test_task_logger.py +++ b/backend/src/core/task_manager/__tests__/test_task_logger.py @@ -3,7 +3,7 @@ # @BRIEF Contract testing for TaskLogger # #endregion Test.Tests.TestTaskLogger import pytest -from unittest.mock import MagicMock +from unittest.mock import MagicMock, patch from src.core.task_manager.task_logger import TaskLogger @@ -30,39 +30,28 @@ def test_task_logger_initialization(task_logger): # @BRIEF Verify TaskLogger delegates log method calls to the underlying persistence service. def test_log_methods_delegation(task_logger, mock_add_log): """Verify info, error, warning, debug delegate to internal _log.""" - task_logger.info("info message", metadata={"k": "v"}) + task_logger.reason("info message", payload={"k": "v"}) mock_add_log.assert_called_with( task_id="test_123", - level="INFO", - message="info message", - source="test_plugin", - metadata={"k": "v"} + event={ + "ts": mock_add_log.call_args.kwargs["event"]["ts"], + "level": "INFO", "trace_id": mock_add_log.call_args.kwargs["event"]["trace_id"], "src": "task.test_plugin", + "marker": "REASON", "intent": "info message", "payload": {"k": "v"}, + }, ) - task_logger.error("error message", source="override") - mock_add_log.assert_called_with( - task_id="test_123", - level="ERROR", - message="error message", - source="override", - metadata=None - ) - task_logger.warning("warning message") - mock_add_log.assert_called_with( - task_id="test_123", - level="WARNING", - message="warning message", - source="test_plugin", - metadata=None - ) - task_logger.debug("debug message") - mock_add_log.assert_called_with( - task_id="test_123", - level="DEBUG", - message="debug message", - source="test_plugin", - metadata=None - ) + with pytest.raises(ValueError, match="EXPLORE events require error"): + task_logger.explore("fallback") + + +def test_task_logger_emits_same_explicit_event_to_both_sinks(task_logger, mock_add_log): + with patch("src.core.task_manager.task_logger.emit_cot_event") as emit: + task_logger.explore("Fallback selected", payload={"attempt": 2}, error="Primary unavailable") + event = mock_add_log.call_args.kwargs["event"] + assert event["marker"] == "EXPLORE" + assert event["error"] == "Primary unavailable" + assert event["src"] == "task.test_plugin" + emit.assert_called_once_with(event) # @TEST_CONTRACT invariants -> "with_source creates a new logger with the same task_id" # #endregion Test.Tests.TestLogMethodsDelegation # #region Test.Tests.TestWithSource [C:2] [TYPE Function] @@ -101,18 +90,12 @@ def test_invalid_add_log_fn(): def test_progress_log(task_logger, mock_add_log): """Verify progress method correctly formats metadata.""" task_logger.progress("Step 1", 45.5) - mock_add_log.assert_called_with( - task_id="test_123", - level="INFO", - message="Step 1", - source="test_plugin", - metadata={"progress": 45.5} - ) + assert mock_add_log.call_args.kwargs["event"]["payload"] == {"progress": 45.5} # Boundary checks task_logger.progress("Step high", 150) - assert mock_add_log.call_args[1]["metadata"]["progress"] == 100 + assert mock_add_log.call_args.kwargs["event"]["payload"] == {"progress": 100} task_logger.progress("Step low", -10) - assert mock_add_log.call_args[1]["metadata"]["progress"] == 0 + assert mock_add_log.call_args.kwargs["event"]["payload"] == {"progress": 0} # #endregion Test.Tests.TestProgressLog diff --git a/backend/src/core/task_manager/context.py b/backend/src/core/task_manager/context.py index 8e381e7fc..6ada5b855 100644 --- a/backend/src/core/task_manager/context.py +++ b/backend/src/core/task_manager/context.py @@ -12,6 +12,8 @@ from collections.abc import Callable from typing import Any +from ss_tools.shared.cot_logger import build_cot_event + from ..logger import belief_scope from .task_logger import TaskLogger @@ -172,13 +174,15 @@ class TaskContext: if metadata: payload.update(metadata) if self._logger and hasattr(self._logger, '_add_log'): - # Fire a special structured log entry + event = build_cot_event( + src=f"task.{self._default_source or 'system'}.heartbeat", + marker="REFLECT", + intent=message or f"heartbeat progress={progress}", + payload={"type": "heartbeat", **payload}, + ) await self._logger._add_log( self._task_id, - "INFO", - message or f"heartbeat progress={progress}", - source=self._default_source or "system", - metadata={"type": "heartbeat", **payload}, + event=event, ) # #endregion Core.Context.Heartbeat # #endregion Core.Context.TaskContext diff --git a/backend/src/core/task_manager/event_bus.py b/backend/src/core/task_manager/event_bus.py index 9f8146dd7..ef916efaf 100644 --- a/backend/src/core/task_manager/event_bus.py +++ b/backend/src/core/task_manager/event_bus.py @@ -20,9 +20,14 @@ # @DATA_CONTRACT Input: LogEntry -> Output: persisted log + subscriber notification import asyncio +from datetime import datetime from typing import Any -from ss_tools.shared.cot_logger import seed_trace_id +from ss_tools.shared.cot_logger import ( + CanonicalCotEvent, + seed_trace_id, +) + from src.core.logger import logger, should_log_task_level from src.core.task_manager.models import LogEntry, LogFilter, LogStats from src.core.task_manager.persistence import TaskLogPersistenceService @@ -172,24 +177,21 @@ class EventBus: # @PRE Task exists. # @POST Log added to buffer and pushed to subscriber queues (if level passes filter). # @RELATION CALLS -> [Core.Logger.ShouldLogTaskLevel] - async def add_log( + async def add_log( # noqa: C901 self, task_id: str, - level: str, - message: str, - source: str = "system", - metadata: dict[str, Any] | None = None, - context: dict[str, Any] | None = None, + level: str | None = None, task_logs_list: list[LogEntry] | None = None, + *, + event: CanonicalCotEvent | None = None, ): if not should_log_task_level(level): return + if event is None: + raise ValueError("EventBus.add_log requires a canonical event") log_entry = LogEntry( - level=level, - message=message, - source=source, - metadata=metadata, - context=context, + timestamp=datetime.fromisoformat(event["ts"]), + **event, ) if task_logs_list is not None: task_logs_list.append(log_entry) @@ -365,14 +367,19 @@ class EventBus: LogEntry( timestamp=log.timestamp, level=log.level, - message=log.message, - source=log.source, - metadata=log.metadata, + trace_id=log.trace_id, + span_id=log.span_id, + src=log.src, + marker=log.marker, + intent=log.intent, + payload=log.payload, + error=log.error, ) for log in task_logs_persisted ] - # For running/pending tasks, return from memory + # For running/pending tasks, return from memory. return task_logs if task_logs else [] + # #endregion Core.EventBus.GetTaskLogs # #region Core.EventBus.GetTaskLogStats [C:2] [TYPE Function] [SEMANTICS log,stats,aggregate] diff --git a/backend/src/core/task_manager/log_export.py b/backend/src/core/task_manager/log_export.py index 6d7e47688..2a896102a 100644 --- a/backend/src/core/task_manager/log_export.py +++ b/backend/src/core/task_manager/log_export.py @@ -93,35 +93,22 @@ def row_to_export_record( else: ts_str = str(ts or "") - message = redact_text(str(row.get("message") or "")) - metadata = redact_metadata(row.get("metadata")) - cot = parse_cot_message(message) - record: dict[str, Any] = { + "id": row.get("id"), "ts": ts_str, "level": str(row.get("level") or "INFO").upper(), - "domain": domain, "task_id": row.get("task_id"), - "plugin_id": plugin_id, - "source": row.get("source") or "system", - "message": message, - "metadata": metadata if metadata is not None else {}, + "src": str(row.get("src") or "task.system"), + "marker": str(row.get("marker") or "REASON"), + "intent": redact_text(str(row.get("intent") or "")), + "trace_id": str(row.get("trace_id") or ""), } - if cot: - record["marker"] = cot.get("marker") - record["intent"] = redact_text(str(cot.get("intent") or "")) - record["src"] = cot.get("src") - record["trace_id"] = cot.get("trace_id") - record["span_id"] = cot.get("span_id") - if cot.get("task_id"): - record["task_id"] = cot.get("task_id") - if "payload" in cot: - record["payload"] = redact_metadata(cot.get("payload")) - if cot.get("error"): - record["error"] = redact_text(str(cot.get("error"))) - # Prefer structured intent as primary text for agents - if record.get("intent"): - record["message"] = record["intent"] + if row.get("span_id"): + record["span_id"] = row["span_id"] + if row.get("payload") is not None: + record["payload"] = redact_metadata(row["payload"]) + if row.get("error"): + record["error"] = redact_text(str(row["error"])) return record # #endregion Core.LogExport.RowToExportRecord diff --git a/backend/src/core/task_manager/manager.py b/backend/src/core/task_manager/manager.py index b17b791b9..b4d56d316 100644 --- a/backend/src/core/task_manager/manager.py +++ b/backend/src/core/task_manager/manager.py @@ -124,12 +124,14 @@ class TaskManager: # #region Core.Manager.MakeAddLogCallback [C:3] [TYPE Function] # @BRIEF Create an async closure for adding logs that looks up the task and delegates to EventBus. def _make_add_log_callback(self): - async def _add_log(task_id, level, message, source="system", metadata=None, context=None): + async def _add_log(task_id, event=None): task = self.graph.get_task(task_id) if not task: return + if event is None: + raise ValueError("Task log callback requires a canonical event") await self.event_bus.add_log( - task_id, level, message, source, metadata, context, + task_id, event=event, task_logs_list=task.logs, ) return _add_log diff --git a/backend/src/core/task_manager/models.py b/backend/src/core/task_manager/models.py index 37b1c5c44..c54629554 100644 --- a/backend/src/core/task_manager/models.py +++ b/backend/src/core/task_manager/models.py @@ -25,7 +25,7 @@ from enum import Enum from typing import Any import uuid -from pydantic import BaseModel, ConfigDict, Field, field_serializer +from pydantic import BaseModel, ConfigDict, Field, field_serializer, model_validator # #region Core.Models.TaskStatus [TYPE Enum] @@ -57,16 +57,30 @@ class LogLevel(str, Enum): class LogEntry(BaseModel): timestamp: datetime = Field(default_factory=lambda: datetime.now(UTC)) level: str = Field(default="INFO") - message: str - source: str = Field( - default="system" - ) # Component attribution: plugin, superset_api, git, etc. - context: dict[str, Any] | None = ( - None # Legacy field, kept for backward compatibility - ) - metadata: dict[str, Any] | None = ( - None # Structured metadata (e.g., dashboard_id, progress) - ) + trace_id: str = "" + span_id: str | None = None + src: str = "task.system" + marker: str = "REASON" + intent: str = "" + payload: dict[str, Any] | None = None + error: str | None = None + + @model_validator(mode="before") + @classmethod + def normalize_legacy_fields(cls, data: Any) -> Any: + if not isinstance(data, dict): + return data + values = dict(data) + return values + + @model_validator(mode="after") + def validate_canonical_event(self) -> "LogEntry": + if self.marker not in {"REASON", "REFLECT", "EXPLORE"}: + raise ValueError(f"Invalid CoT marker: {self.marker}") + if self.marker == "EXPLORE" and not self.error: + raise ValueError("EXPLORE events require error") + return self + # #endregion Core.Models.LogEntry # #region Core.Models.TaskLog [C:3] [TYPE Class] [SEMANTICS task, log, persistent, pydantic] # @defgroup TaskManager Module group. @@ -77,10 +91,23 @@ class TaskLog(BaseModel): task_id: str timestamp: datetime level: str - source: str - message: str - metadata: dict[str, Any] | None = None + trace_id: str + span_id: str | None = None + src: str + marker: str + intent: str + payload: dict[str, Any] | None = None + error: str | None = None model_config = ConfigDict(from_attributes=True) + + @model_validator(mode="after") + def validate_canonical_event(self) -> "TaskLog": + if self.marker not in {"REASON", "REFLECT", "EXPLORE"}: + raise ValueError(f"Invalid CoT marker: {self.marker}") + if self.marker == "EXPLORE" and not self.error: + raise ValueError("EXPLORE events require error") + return self + # #endregion Core.Models.TaskLog # #region Core.Models.LogFilter [TYPE Class] # @defgroup TaskManager Module group. diff --git a/backend/src/core/task_manager/persistence.py b/backend/src/core/task_manager/persistence.py index f7390f9e7..a7ba5a0e1 100644 --- a/backend/src/core/task_manager/persistence.py +++ b/backend/src/core/task_manager/persistence.py @@ -208,15 +208,12 @@ class TaskPersistenceService: log_dict = log.model_dump() if isinstance(log_dict.get("timestamp"), datetime): log_dict["timestamp"] = log_dict["timestamp"].isoformat() - # Also clean up any datetimes in context - if log_dict.get("context"): - log_dict["context"] = json_serializable(log_dict["context"]) - record.logs.append(log_dict) + record.logs.append(json_serializable(log_dict)) # Extract error if failed if task.status == TaskStatus.FAILED: for log in reversed(task.logs): if log.level == "ERROR": - record.error = log.message + record.error = log.intent break session.commit() except Exception as e: @@ -393,16 +390,23 @@ class TaskLogPersistenceService: return rows: list[dict] = [] for log in logs: + safe_payload = ( + json.loads(json.dumps(log.payload, default=str)) + if log.payload is not None + else None + ) rows.append( { "task_id": task_id, "timestamp": log.timestamp, "level": (log.level or "INFO").upper(), - "source": log.source or "system", - "message": log.message, - "metadata_json": json.dumps(log.metadata, default=str) - if log.metadata - else None, + "trace_id": log.trace_id, + "span_id": log.span_id, + "src": log.src, + "marker": log.marker, + "intent": log.intent, + "payload": safe_payload, + "error": log.error, } ) from sqlalchemy import insert @@ -446,31 +450,29 @@ class TaskLogPersistenceService: TaskLogRecord.level == log_filter.level.upper() ) if log_filter.source: - query = query.filter(TaskLogRecord.source == log_filter.source) + query = query.filter(TaskLogRecord.src == log_filter.source) if log_filter.search: search_pattern = f"%{log_filter.search}%" - query = query.filter(TaskLogRecord.message.ilike(search_pattern)) + query = query.filter(TaskLogRecord.intent.ilike(search_pattern)) # Order by timestamp ascending (oldest first) query = query.order_by(TaskLogRecord.timestamp.asc()) # Apply pagination records = query.offset(log_filter.offset).limit(log_filter.limit).all() logs = [] for record in records: - metadata = None - if record.metadata_json: - try: - metadata = json.loads(record.metadata_json) - except json.JSONDecodeError: - metadata = None logs.append( TaskLog( id=record.id, task_id=record.task_id, timestamp=record.timestamp, level=record.level, - source=record.source, - message=record.message, - metadata=metadata, + trace_id=record.trace_id, + span_id=record.span_id, + src=record.src, + marker=record.marker, + intent=record.intent, + payload=record.payload, + error=record.error, ) ) return logs @@ -510,9 +512,9 @@ class TaskLogPersistenceService: by_level = {level: count for level, count in level_counts} # Get counts by source source_counts = ( - session.query(TaskLogRecord.source, func.count(TaskLogRecord.id)) + session.query(TaskLogRecord.src, func.count(TaskLogRecord.id)) .filter(TaskLogRecord.task_id == task_id) - .group_by(TaskLogRecord.source) + .group_by(TaskLogRecord.src) .all() ) by_source = {source: count for source, count in source_counts} @@ -539,7 +541,7 @@ class TaskLogPersistenceService: try: from sqlalchemy import distinct sources = ( - session.query(distinct(TaskLogRecord.source)) + session.query(distinct(TaskLogRecord.src)) .filter(TaskLogRecord.task_id == task_id) .all() ) @@ -577,28 +579,26 @@ class TaskLogPersistenceService: allowed = [lv for lv, o in order.items() if o >= min_ord] query = query.filter(TaskLogRecord.level.in_(allowed)) if source: - query = query.filter(TaskLogRecord.source == source) + query = query.filter(TaskLogRecord.src == source) if search: query = query.filter( - TaskLogRecord.message.ilike(f"%{search}%") + TaskLogRecord.intent.ilike(f"%{search}%") ) query = query.order_by(TaskLogRecord.timestamp.asc()).limit(max_rows) # yield_per keeps memory bounded for large exports for record in query.yield_per(batch_size): - metadata = None - if record.metadata_json: - try: - metadata = json.loads(record.metadata_json) - except json.JSONDecodeError: - metadata = None yield { "id": record.id, "task_id": record.task_id, "timestamp": record.timestamp, "level": record.level, - "source": record.source, - "message": record.message, - "metadata": metadata, + "trace_id": record.trace_id, + "span_id": record.span_id, + "src": record.src, + "marker": record.marker, + "intent": record.intent, + "payload": record.payload, + "error": record.error, } finally: session.close() diff --git a/backend/src/core/task_manager/task_logger.py b/backend/src/core/task_manager/task_logger.py index 6cfc46345..a4d5e662e 100644 --- a/backend/src/core/task_manager/task_logger.py +++ b/backend/src/core/task_manager/task_logger.py @@ -9,8 +9,9 @@ import asyncio from collections.abc import Callable from typing import Any -from ss_tools.shared.cot_logger import get_trace_id, set_task_id -from ..logger import logger as main_cot_logger +from ss_tools.shared.cot_logger import build_cot_event, emit_cot_event, set_task_id + +from ..logger import logger as main_cot_logger # noqa: TID252 # #region Core.TaskLogger [C:2] [TYPE Class] [SEMANTICS logger, task, source, attribution] @@ -86,58 +87,39 @@ class TaskLogger: message: str, source: str | None = None, metadata: dict[str, Any] | None = None, + marker: str = "REASON", + error: str | None = None, ) -> None: """Internal logging method. Also bridges to main CoT for agent-centric unified traces.""" effective_source = source or self._default_source + qualified_src = effective_source if "." in effective_source else f"task.{effective_source}" + event = build_cot_event( + src=qualified_src, + marker=marker, + intent=message, + payload=metadata, + error=error, + level=level, + ) # Bind task_id so CotJsonFormatter stamps app.log lines for Reports cross-filter. set_task_id(self._task_id) - coro = self._add_log( - task_id=self._task_id, - level=level, - message=message, - source=effective_source, - metadata=metadata, - ) + coro = self._add_log(task_id=self._task_id, event=event) # Fire-and-forget the task-specific persistence try: loop = asyncio.get_running_loop() if loop.is_running(): - asyncio.ensure_future(coro) + asyncio.ensure_future(coro) # noqa: RUF006 except RuntimeError: pass # === Agent-centric CoT bridge === - # Emit via main logger so current trace_id (if any) captures the task step. - # Map to REASON / REFLECT / EXPLORE for consistency with the protocol. + # The exact same canonical event is written to both sinks. try: - cot_src = f"Task.{self._task_id[:8]}.{effective_source}" - payload = {"task_id": self._task_id, **(metadata or {})} - - if level.upper() == "ERROR": - main_cot_logger.explore( - f"[{effective_source}] {message}", - payload=payload, - error="Error during task execution" - ) - elif level.upper() == "WARNING": - main_cot_logger.explore( - f"[{effective_source}] {message}", - payload=payload - ) - else: - # Default to REASON for forward progress, REFLECT for completion signals - if any(word in message.lower() for word in ("complete", "success", "finished", "done", "reflected")): - main_cot_logger.reflect(f"[{effective_source}] {message}", payload=payload) - else: - main_cot_logger.reason(f"[{effective_source}] {message}", payload=payload) + emit_cot_event(event) except Exception as _bridge_err: # Bridge must never break primary task logging — but log the failure - try: - main_cot_logger.explore( - "CoT bridge internal error", - payload={"task_id": self._task_id}, - error=str(_bridge_err), - ) + try: # noqa: SIM105 + main_cot_logger.error("CoT bridge internal error", exc_info=_bridge_err) except Exception: pass # last-resort guard # #endregion Core.TaskLogger.Log @@ -171,7 +153,7 @@ class TaskLogger: source: str | None = None, metadata: dict[str, Any] | None = None, ) -> None: - self._log("INFO", message, source, metadata) + self._log("INFO", message, source, metadata, marker="REASON") # #endregion Core.TaskLogger.Info # #region Core.TaskLogger.Warning [TYPE Function] # @ingroup TaskManager @@ -187,7 +169,7 @@ class TaskLogger: source: str | None = None, metadata: dict[str, Any] | None = None, ) -> None: - self._log("WARNING", message, source, metadata) + self._log("WARNING", message, source, metadata, marker="EXPLORE", error=message) # #endregion Core.TaskLogger.Warning # #region Core.TaskLogger.Error [TYPE Function] # @ingroup TaskManager @@ -203,7 +185,7 @@ class TaskLogger: source: str | None = None, metadata: dict[str, Any] | None = None, ) -> None: - self._log("ERROR", message, source, metadata) + self._log("ERROR", message, source, metadata, marker="EXPLORE", error=message) # #endregion Core.TaskLogger.Error # #region Core.TaskLogger.Progress [TYPE Function] # @ingroup TaskManager @@ -218,7 +200,7 @@ class TaskLogger: ) -> None: """Log a progress update with percentage.""" metadata = {"progress": min(100, max(0, percent))} - self._log("INFO", message, source, metadata) + self._log("INFO", message, source, metadata, marker="REFLECT") # #endregion Core.TaskLogger.Progress # #region Core.TaskLogger.Reason / reflect / explore (CoT bridge convenience) [C:2] @@ -227,18 +209,15 @@ class TaskLogger: # This allows task code to use the same mental model as main app logging. def reason(self, message: str, *, payload: dict[str, Any] | None = None, source: str | None = None): """REASON step inside a task (extends reasoning chain).""" - self._log("INFO", message, source, payload) + self._log("INFO", message, source, payload, marker="REASON") def reflect(self, message: str, *, payload: dict[str, Any] | None = None, source: str | None = None): """REFLECT step inside a task (verification of outcome).""" - self._log("INFO", message, source, payload) + self._log("INFO", message, source, payload, marker="REFLECT") def explore(self, message: str, *, payload: dict[str, Any] | None = None, error: str | None = None, source: str | None = None): """EXPLORE / fallback inside a task.""" - meta = payload or {} - if error: - meta["error"] = error - self._log("WARNING", message, source, meta) + self._log("WARNING", message, source, payload, marker="EXPLORE", error=error) # #endregion Core.TaskLogger.Reason / reflect / explore (CoT bridge convenience) # #endregion Core.TaskLogger # #endregion Core.TaskLogger.TaskLoggerModule diff --git a/backend/src/models/task.py b/backend/src/models/task.py index 82bcba416..1ec8af356 100644 --- a/backend/src/models/task.py +++ b/backend/src/models/task.py @@ -100,15 +100,19 @@ class TaskLogRecord(Base): task_id = Column(String, ForeignKey("task_records.id", ondelete="CASCADE"), nullable=False, index=True) timestamp = Column(DateTime(timezone=True), nullable=False, index=True) level = Column(String(16), nullable=False) # INFO, WARNING, ERROR, DEBUG - source = Column(String(64), nullable=False, default="system") # plugin, superset_api, git, etc. - message = Column(Text, nullable=False) - metadata_json = Column(Text, nullable=True) # JSON string for additional metadata + trace_id = Column(String(64), nullable=False, default="") + span_id = Column(String(128), nullable=True) + src = Column(String(255), nullable=False) + marker = Column(String(16), nullable=False) + intent = Column(Text, nullable=False) + payload = Column(JSON, nullable=True) + error = Column(Text, nullable=True) # Composite indexes for efficient filtering __table_args__ = ( Index('ix_task_logs_task_timestamp', 'task_id', 'timestamp'), Index('ix_task_logs_task_level', 'task_id', 'level'), - Index('ix_task_logs_task_source', 'task_id', 'source'), + Index('ix_task_logs_task_src', 'task_id', 'src'), ) # #endregion Models.Task.TaskLogRecord diff --git a/backend/src/plugins/migration.py b/backend/src/plugins/migration.py index 76e083912..5991668da 100755 --- a/backend/src/plugins/migration.py +++ b/backend/src/plugins/migration.py @@ -636,6 +636,26 @@ class MigrationPlugin(PluginBase): formatted = format_superset_import_error(import_exc) if not formatted.get("is_password_required"): if sync_dataset_composite_keys: + migration_log.warning( + "Initial import failed; starting automatic dataset-key recovery", + extra={ + "dashboard_id": dash_id, + "dashboard_title": title, + "attempt": "initial", + "recovery_strategy": composite_key_mutation_server, + "recovery_phase": "started", + }, + ) + migration_log.info( + "Synchronizing dataset composite keys before retry", + extra={ + "dashboard_id": dash_id, + "dashboard_title": title, + "attempt": "recovery", + "recovery_strategy": composite_key_mutation_server, + "recovery_phase": "sync", + }, + ) report, imported, retry_error = await _attempt_composite_key_fallback( engine=engine, from_c=from_c, @@ -649,12 +669,36 @@ class MigrationPlugin(PluginBase): source_zip=str(tmp_zip_path), transformed_zip=str(tmp_new_zip), ) + migration_log.info( + "Dataset composite-key synchronization finished; attempting dashboard import retry", + extra={ + "dashboard_id": dash_id, + "dashboard_title": title, + "attempt": "recovery", + "recovery_strategy": composite_key_mutation_server, + "recovery_phase": "retry_import", + "sync_changed": report.get("changed", 0), + "sync_unchanged": report.get("unchanged", 0), + "sync_skipped_missing": report.get("skipped_missing", 0), + "sync_failed": report.get("failed", 0), + }, + ) if imported: migration_result["migrated_dashboards"].append({"id": dash_id, "title": title}) app_logger.reflect( "Composite-key fallback retry import succeeded", extra={"title": title}, ) + migration_log.info( + "Recovery retry import succeeded; dashboard migrated", + extra={ + "dashboard_id": dash_id, + "dashboard_title": title, + "attempt": "recovery", + "recovery_strategy": composite_key_mutation_server, + "recovery_phase": "completed", + }, + ) continue if retry_error is not None: app_logger.explore( @@ -662,6 +706,17 @@ class MigrationPlugin(PluginBase): extra={"dash_id": dash_id, "title": title}, error=str(retry_error), ) + migration_log.error( + "Recovery retry import failed; dashboard could not be migrated", + extra={ + "dashboard_id": dash_id, + "dashboard_title": title, + "attempt": "recovery", + "recovery_strategy": composite_key_mutation_server, + "recovery_phase": "failed", + "retry_error": str(retry_error) if retry_error else None, + }, + ) entry = _failed_entry_with_archive( dash_id, title, import_exc, phase="import", formatted=formatted, diff --git a/backend/src/services/reports/normalizer.py b/backend/src/services/reports/normalizer.py index b98163f46..80e16b61c 100644 --- a/backend/src/services/reports/normalizer.py +++ b/backend/src/services/reports/normalizer.py @@ -89,8 +89,8 @@ def extract_error_context(task: Task, report_status: ReportStatus) -> ErrorConte if not message: for log in reversed(task.logs): - if str(log.level).upper() == "ERROR" and log.message: - message = log.message + if str(log.level).upper() == "ERROR" and log.intent: + message = log.intent break if not message: diff --git a/backend/tests/core/task_manager/test_log_persistence_service.py b/backend/tests/core/task_manager/test_log_persistence_service.py index a11de4c8f..c02b7c0be 100644 --- a/backend/tests/core/task_manager/test_log_persistence_service.py +++ b/backend/tests/core/task_manager/test_log_persistence_service.py @@ -62,9 +62,12 @@ class TestTaskLogPersistence: entry = LogEntry( timestamp=kwargs.get("timestamp", datetime.now(UTC)), level=kwargs.get("level", "INFO"), - source=kwargs.get("source", "system"), - message=kwargs["message"], - metadata=kwargs.get("metadata"), + trace_id="trace-test", + src=kwargs.get("source", "task.system"), + marker="EXPLORE" if kwargs.get("level", "INFO") in {"WARNING", "ERROR"} else "REASON", + intent=kwargs["message"], + payload=kwargs.get("metadata"), + error=kwargs["message"] if kwargs.get("level", "INFO") in {"WARNING", "ERROR"} else None, ) with self._patched(): self.service.add_logs(task_id, [entry]) @@ -88,16 +91,16 @@ class TestTaskLogPersistence: session.close() assert result is not None assert result.level == "INFO" - assert result.source == "test_source" - assert result.message == "Test message" - assert result.metadata_json is None + assert result.src == "test_source" + assert result.intent == "Test message" + assert result.payload is None def test_add_logs_batch(self): """Multiple entries batch-insert in one call.""" entries = [ - LogEntry(timestamp=datetime.now(UTC), level="INFO", source="s1", message="M1"), - LogEntry(timestamp=datetime.now(UTC), level="WARNING", source="s2", message="M2"), - LogEntry(timestamp=datetime.now(UTC), level="ERROR", source="s3", message="M3"), + LogEntry(timestamp=datetime.now(UTC), level="INFO", trace_id="trace", src="s1", marker="REASON", intent="M1"), + LogEntry(timestamp=datetime.now(UTC), level="WARNING", trace_id="trace", src="s2", marker="EXPLORE", intent="M2", error="M2"), + LogEntry(timestamp=datetime.now(UTC), level="ERROR", trace_id="trace", src="s3", marker="EXPLORE", intent="M3", error="M3"), ] with self._patched(): self.service.add_logs("task-b", entries) @@ -109,8 +112,8 @@ class TestTaskLogPersistence: def test_add_logs_default_level_and_source(self): """Falsy level/source fall back to INFO / system.""" legacy = LogEntry.model_construct( - timestamp=datetime.now(UTC), level=None, source=None, - message="defaults", metadata=None, + timestamp=datetime.now(UTC), level=None, trace_id="trace", src="task.system", + marker="REASON", intent="defaults", payload=None, ) with self._patched(): self.service.add_logs("task-c", [legacy]) @@ -118,37 +121,37 @@ class TestTaskLogPersistence: result = session.query(TaskLogRecord).filter_by(task_id="task-c").first() session.close() assert result.level == "INFO" - assert result.source == "system" + assert result.src == "task.system" def test_add_logs_metadata_json_dumped(self): """Metadata dict serializes to metadata_json.""" entry = LogEntry( - timestamp=datetime.now(UTC), level="INFO", source="s", - message="meta", metadata={"dashboard_id": "dash-1", "progress": 50}, + timestamp=datetime.now(UTC), level="INFO", trace_id="trace", src="s", marker="REASON", + intent="meta", payload={"dashboard_id": "dash-1", "progress": 50}, ) with self._patched(): self.service.add_logs("task-d", [entry]) session = self.TestSessionLocal() result = session.query(TaskLogRecord).filter_by(task_id="task-d").first() session.close() - assert result.metadata_json == '{"dashboard_id": "dash-1", "progress": 50}' + assert result.payload == {"dashboard_id": "dash-1", "progress": 50} def test_add_logs_metadata_with_datetime_default_str(self): """Non-serializable metadata values fall back to str() via default=str.""" entry = LogEntry( - timestamp=datetime.now(UTC), level="INFO", source="s", - message="meta-dt", metadata={"at": datetime(2024, 1, 1, 12, 0, tzinfo=UTC)}, + timestamp=datetime.now(UTC), level="INFO", trace_id="trace", src="s", marker="REASON", + intent="meta-dt", payload={"at": datetime(2024, 1, 1, 12, 0, tzinfo=UTC)}, ) with self._patched(): self.service.add_logs("task-e", [entry]) session = self.TestSessionLocal() result = session.query(TaskLogRecord).filter_by(task_id="task-e").first() session.close() - assert "2024-01-01 12:00:00+00:00" in result.metadata_json + assert result.payload == {"at": "2024-01-01 12:00:00+00:00"} def test_add_logs_db_error_rollback_and_raise(self): """Commit failure rolls back and re-raises for EventBus re-queue.""" - entry = LogEntry(timestamp=datetime.now(UTC), level="INFO", source="s", message="will fail") + entry = LogEntry(timestamp=datetime.now(UTC), level="INFO", trace_id="trace", src="s", marker="REASON", intent="will fail") def broken_session(): s = self.TestSessionLocal() s.commit = lambda: (_ for _ in ()).throw(RuntimeError("db dead")) @@ -171,7 +174,7 @@ class TestTaskLogPersistence: logs = self.service.get_logs("task-f", LogFilter()) assert len(logs) == 3 assert all(log.task_id == "task-f" for log in logs) - assert [log.message for log in logs] == ["Message 0", "Message 1", "Message 2"] + assert [log.intent for log in logs] == ["Message 0", "Message 1", "Message 2"] def test_get_logs_level_filter(self): """Level filter narrows results (case-insensitive comparison).""" @@ -189,7 +192,7 @@ class TestTaskLogPersistence: with self._patched(): api_logs = self.service.get_logs("task-h", LogFilter(source="api")) assert len(api_logs) == 1 - assert api_logs[0].source == "api" + assert api_logs[0].src == "api" def test_get_logs_search_filter(self): """Search filter matches messages case-insensitively.""" @@ -198,7 +201,7 @@ class TestTaskLogPersistence: with self._patched(): found = self.service.get_logs("task-h", LogFilter(search="AUTHENTICATION")) assert len(found) == 1 - assert "authentication" in found[0].message.lower() + assert "authentication" in found[0].intent.lower() def test_get_logs_pagination(self): """offset/limit paginate correctly.""" @@ -215,21 +218,21 @@ class TestTaskLogPersistence: self._add("task-h", message="meta", metadata={"dashboard_id": "dash-1"}) with self._patched(): logs = self.service.get_logs("task-h", LogFilter()) - assert logs[0].metadata == {"dashboard_id": "dash-1"} + assert logs[0].payload == {"dashboard_id": "dash-1"} def test_get_logs_corrupt_metadata(self): """Corrupt metadata_json degrades to None.""" session = self.TestSessionLocal() session.add(TaskLogRecord( - task_id="task-h", timestamp=datetime.now(UTC), - level="INFO", source="s", message="corrupt", metadata_json="{not-json", + task_id="task-h", timestamp=datetime.now(UTC), trace_id="trace", + level="INFO", src="s", marker="REASON", intent="corrupt", payload=None, )) session.commit() session.close() with self._patched(): logs = self.service.get_logs("task-h", LogFilter()) assert len(logs) == 1 - assert logs[0].metadata is None + assert logs[0].payload is None # ── get_log_stats / get_sources ────────────────────────────────────── def test_get_log_stats_counts(self): @@ -278,8 +281,8 @@ class TestTaskLogPersistence: self._add("export-1", message="first", timestamp=datetime(2024, 1, 1, 10, 0, tzinfo=UTC)) with self._patched(): rows = list(self.service.iter_logs_for_export("export-1")) - assert [r["message"] for r in rows] == ["first", "second"] - assert set(rows[0].keys()) == {"id", "task_id", "timestamp", "level", "source", "message", "metadata"} + assert [r["intent"] for r in rows] == ["first", "second"] + assert set(rows[0].keys()) == {"id", "task_id", "timestamp", "level", "trace_id", "span_id", "src", "marker", "intent", "payload", "error"} assert rows[0]["task_id"] == "export-1" def test_iter_logs_for_export_level_filter(self): @@ -306,7 +309,7 @@ class TestTaskLogPersistence: self._add("export-1", message="storage msg", source="storage") with self._patched(): rows = list(self.service.iter_logs_for_export("export-1", source="api")) - assert [r["message"] for r in rows] == ["api msg"] + assert [r["intent"] for r in rows] == ["api msg"] def test_iter_logs_for_export_search_filter(self): """search filter matches messages via LIKE.""" @@ -321,20 +324,20 @@ class TestTaskLogPersistence: self._add("export-1", message="meta", metadata={"dashboard_id": "dash-1"}) with self._patched(): rows = list(self.service.iter_logs_for_export("export-1")) - assert rows[0]["metadata"] == {"dashboard_id": "dash-1"} + assert rows[0]["payload"] == {"dashboard_id": "dash-1"} def test_iter_logs_for_export_corrupt_metadata(self): """Corrupt metadata_json degrades to None in exported rows.""" session = self.TestSessionLocal() session.add(TaskLogRecord( - task_id="export-1", timestamp=datetime.now(UTC), - level="INFO", source="s", message="corrupt", metadata_json="{bad", + task_id="export-1", timestamp=datetime.now(UTC), trace_id="trace", + level="INFO", src="s", marker="REASON", intent="corrupt", payload=None, )) session.commit() session.close() with self._patched(): rows = list(self.service.iter_logs_for_export("export-1")) - assert rows[0]["metadata"] is None + assert rows[0]["payload"] is None def test_iter_logs_for_export_empty_result(self): """No matching rows -> no yields.""" diff --git a/backend/tests/core/task_manager/test_persistence.py b/backend/tests/core/task_manager/test_persistence.py index 4f89db028..3043093be 100644 --- a/backend/tests/core/task_manager/test_persistence.py +++ b/backend/tests/core/task_manager/test_persistence.py @@ -200,9 +200,9 @@ class TestTaskPersistenceService: task = self._make_task( logs=[ LogEntry( - message="Step 1", level="INFO", source="plugin", + intent="Step 1", level="INFO", src="plugin", trace_id="trace", marker="REASON", timestamp=datetime(2024, 1, 1, 12, 0, tzinfo=UTC), - context={"t": datetime(2024, 2, 2, 3, 4, tzinfo=UTC)}, + payload={"t": datetime(2024, 2, 2, 3, 4, tzinfo=UTC)}, ) ] ) @@ -211,17 +211,17 @@ class TestTaskPersistenceService: record = self._load_record() assert record.logs is not None assert record.logs[0]["timestamp"] == "2024-01-01T12:00:00+00:00" - assert record.logs[0]["context"]["t"] == "2024-02-02T03:04:00+00:00" + assert record.logs[0]["payload"]["t"] == "2024-02-02T03:04:00+00:00" def test_persist_task_failed_extracts_last_error(self): """FAILED tasks store the last ERROR log message.""" task = self._make_task( status=TaskStatus.FAILED, logs=[ - LogEntry(message="Started OK", level="INFO", source="plugin"), - LogEntry(message="Connection failed", level="ERROR", source="plugin"), - LogEntry(message="Retrying...", level="INFO", source="plugin"), - LogEntry(message="Fatal: timeout", level="ERROR", source="plugin"), + LogEntry(intent="Started OK", level="INFO", src="plugin", trace_id="trace", marker="REASON"), + LogEntry(intent="Connection failed", level="ERROR", src="plugin", trace_id="trace", marker="EXPLORE", error="Connection failed"), + LogEntry(intent="Retrying...", level="INFO", src="plugin", trace_id="trace", marker="REASON"), + LogEntry(intent="Fatal: timeout", level="ERROR", src="plugin", trace_id="trace", marker="EXPLORE", error="Fatal: timeout"), ], ) with self._patched(): @@ -234,8 +234,8 @@ class TestTaskPersistenceService: task = self._make_task( status=TaskStatus.FAILED, logs=[ - LogEntry(message="Started", level="INFO", source="plugin"), - LogEntry(message="Finished", level="INFO", source="plugin"), + LogEntry(intent="Started", level="INFO", src="plugin", trace_id="trace", marker="REASON"), + LogEntry(intent="Finished", level="INFO", src="plugin", trace_id="trace", marker="REFLECT"), ], ) with self._patched(): @@ -246,8 +246,8 @@ class TestTaskPersistenceService: def test_persist_task_log_with_string_timestamp(self): """Defensive path: legacy log timestamp that is not a datetime is kept as-is.""" legacy = LogEntry.model_construct( - message="legacy", level="INFO", source="plugin", - timestamp="2024-01-01T00:00:00", context=None, metadata=None, + intent="legacy", level="INFO", src="plugin", trace_id="trace", marker="REASON", + timestamp="2024-01-01T00:00:00", payload=None, ) task = self._make_task(logs=[legacy]) with self._patched(): @@ -448,13 +448,13 @@ class TestTaskPersistenceService: def test_load_tasks_logs_parsed(self): """Stored log dicts reconstruct as LogEntry objects.""" task = self._make_task( - logs=[LogEntry(message="m1", level="INFO", source="p", timestamp=datetime(2024, 1, 1, 12, 0, tzinfo=UTC))] + logs=[LogEntry(intent="m1", level="INFO", src="p", trace_id="trace", marker="REASON", timestamp=datetime(2024, 1, 1, 12, 0, tzinfo=UTC))] ) with self._patched(): self.service.persist_task(task) with self._patched(): loaded = self.service.load_tasks() - assert loaded[0].logs[0].message == "m1" + assert loaded[0].logs[0].intent == "m1" assert loaded[0].logs[0].timestamp == datetime(2024, 1, 1, 12, 0, tzinfo=UTC) def test_load_tasks_log_timestamp_fallback(self): @@ -464,12 +464,12 @@ class TestTaskPersistenceService: self.service.persist_task(task) session = self.TestSessionLocal() record = session.query(TaskRecord).filter_by(id="task-1").first() - record.logs = [{"message": "no ts", "level": "INFO", "source": "p"}] + record.logs = [{"intent": "no ts", "level": "INFO", "trace_id": "trace", "src": "p", "marker": "REASON"}] session.commit() session.close() with self._patched(): loaded = self.service.load_tasks() - assert loaded[0].logs[0].message == "no ts" + assert loaded[0].logs[0].intent == "no ts" assert loaded[0].logs[0].timestamp is not None def test_load_tasks_skips_non_dict_logs(self): diff --git a/backend/tests/plugins/test_migration_plugin.py b/backend/tests/plugins/test_migration_plugin.py index 2c001164e..926f38a51 100644 --- a/backend/tests/plugins/test_migration_plugin.py +++ b/backend/tests/plugins/test_migration_plugin.py @@ -38,6 +38,16 @@ def _make_task_context(): return ctx +def _make_logging_task_context(): + """Create a context whose task logger records structured metadata calls.""" + ctx = MagicMock() + ctx.logger.with_source.side_effect = lambda source: ctx.logger + ctx.logger.info = MagicMock() + ctx.logger.warning = MagicMock() + ctx.logger.error = MagicMock() + return ctx + + def _make_dashboard(dash_id=1, title="Test Dash"): return {"id": dash_id, "slug": f"slug-{dash_id}", "dashboard_title": title} @@ -1220,6 +1230,7 @@ class TestMigrationPluginCompositeKeyFallback: report = {"changed": 1, "unchanged": 0, "skipped_missing": 0, "failed": 0, "errors": []} mock_sync = AsyncMock(return_value=report) + context = _make_logging_task_context() with patch('src.plugins.migration.get_config_manager', return_value=mock_cm), \ patch('src.plugins.migration.SupersetClient') as MockSC, \ @@ -1239,7 +1250,7 @@ class TestMigrationPluginCompositeKeyFallback: "replace_db_config": False, "sync_dataset_composite_keys": True, "composite_key_mutation_server": "target", - }) + }, context=context) assert result["status"] == "PARTIAL_SUCCESS" entry = result["failed_dashboards"][0] @@ -1248,6 +1259,68 @@ class TestMigrationPluginCompositeKeyFallback: assert entry["composite_key_mutation_server"] == "target" assert entry["composite_key_retried"] is True + recovery_calls = [ + call for call in context.logger.error.call_args_list + if call.args and "Recovery retry import failed" in call.args[0] + ] + assert recovery_calls + recovery_metadata = recovery_calls[0].kwargs["extra"] + assert recovery_metadata["dashboard_id"] == 1 + assert recovery_metadata["recovery_phase"] == "failed" + + @pytest.mark.asyncio + async def test_execute_composite_key_fallback_logs_full_recovery_timeline(self): + """Successful fallback emits task metadata for each user-visible recovery phase.""" + plugin = MigrationPlugin() + src_env = _make_env("env-1", "Source") + tgt_env = _make_env("env-2", "Target") + mock_cm = MagicMock() + mock_cm.get_environments.return_value = [src_env, tgt_env] + mock_src_client = _make_mock_superset_client() + mock_src_client.get_dashboards = AsyncMock(return_value=(True, [_make_dashboard(1, "Dash")])) + mock_src_client.export_dashboard = AsyncMock(return_value=(b"zip", "meta")) + mock_tgt_client = _make_mock_superset_client() + mock_tgt_client.import_dashboard = AsyncMock(side_effect=[RuntimeError("Generic import boom"), None]) + mock_engine = MagicMock() + mock_engine.transform_zip.return_value = True + mock_engine.read_dataset_contracts_from_zip.return_value = [ + {"uuid": "ds-1", "database_uuid": "tgt-db", "catalog": None, "schema": "public", "table_name": "users"} + ] + report = {"changed": 1, "unchanged": 0, "skipped_missing": 0, "failed": 0, "errors": []} + context = _make_logging_task_context() + + with patch("src.plugins.migration.get_config_manager", return_value=mock_cm), \ + patch("src.plugins.migration.SupersetClient") as MockSC, \ + patch("src.plugins.migration.MigrationEngine", return_value=mock_engine), \ + patch("src.plugins.migration.create_temp_file") as mock_ctf, \ + patch("src.plugins.migration.sync_dataset_composite_keys", new=AsyncMock(return_value=report)), \ + patch("src.plugins.migration.IdMappingService", return_value=_make_mock_mapping_service()), \ + patch("src.plugins.migration.SessionLocal"): + MockSC.side_effect = [mock_src_client, mock_tgt_client] + mock_ctf.return_value.__enter__ = MagicMock(return_value="/tmp/test.zip") + result = await plugin.execute({ + "source_env_id": "env-1", + "target_env_id": "env-2", + "selected_ids": [1], + "replace_db_config": False, + "sync_dataset_composite_keys": True, + "composite_key_mutation_server": "target", + }, context=context) + + assert result["status"] == "SUCCESS" + warning_metadata = context.logger.warning.call_args.kwargs["extra"] + assert warning_metadata["recovery_phase"] == "started" + assert warning_metadata["attempt"] == "initial" + metadata_calls = [ + call.kwargs["extra"] + for call in context.logger.info.call_args_list + if "extra" in call.kwargs and "recovery_phase" in call.kwargs["extra"] + ] + assert [metadata["recovery_phase"] for metadata in metadata_calls] == [ + "sync", "retry_import", "completed" + ] + assert metadata_calls[1]["sync_changed"] == 1 + @pytest.mark.asyncio async def test_execute_composite_key_never_syncs_password_failure(self): """Password-required failures never trigger the composite-key fallback.""" @@ -1339,3 +1412,11 @@ class TestMigrationPluginCompositeKeyFallback: entry = result["failed_dashboards"][0] assert "composite_key_retried" not in entry # #endregion Test.MigrationPlugin +# #region Test.Migration.RecoveryLogs [C:3] [TYPE Module] [SEMANTICS test,migration,recovery,logging] +# @BRIEF Verify migration recovery emits structured, user-facing phase events. +# @RELATION BINDS_TO -> [Plugin.Migration.Execute] +# @TEST_CONTRACT: [InitialImportFailure + RecoveryResult] -> [RecoveryPhaseLogs] +# @TEST_SCENARIO: recovery_success -> Initial failure is followed by sync, retry, and success events. +# @TEST_EDGE: external_fail -> Recovery failure remains visible with dashboard context. +# @TEST_INVARIANT: Per-dashboard export/transform/import failure never aborts the batch -> VERIFIED_BY: [recovery_success, external_fail] +# #endregion Test.Migration.RecoveryLogs diff --git a/backend/tests/services/test_report_normalizer.py b/backend/tests/services/test_report_normalizer.py index cebc9bd0d..e7e6e9117 100644 --- a/backend/tests/services/test_report_normalizer.py +++ b/backend/tests/services/test_report_normalizer.py @@ -140,7 +140,7 @@ class TestExtractErrorContext: from src.core.task_manager.models import LogEntry task = _make_task( result={"other": "data"}, - logs=[LogEntry(level="ERROR", message="Critical failure in module X")], + logs=[LogEntry(level="ERROR", trace_id="trace", src="plugin", marker="EXPLORE", intent="Critical failure in module X", error="Critical failure in module X")], ) ctx = extract_error_context(task, ReportStatus.FAILED) assert ctx is not None @@ -285,7 +285,7 @@ class TestExtractErrorContextLogFallback: from src.core.task_manager.models import LogEntry task = _make_task( result={"other": "data"}, - logs=[LogEntry(level="Error", message="cased error")], + logs=[LogEntry(level="Error", trace_id="trace", src="plugin", marker="EXPLORE", intent="cased error", error="cased error")], ) ctx = extract_error_context(task, ReportStatus.FAILED) assert ctx is not None @@ -298,9 +298,9 @@ class TestExtractErrorContextLogFallback: task = _make_task( result={"other": "data"}, logs=[ - LogEntry(level="ERROR", message="first error"), - LogEntry(level="INFO", message="noise"), - LogEntry(level="ERROR", message="last error"), + LogEntry(level="ERROR", trace_id="trace", src="plugin", marker="EXPLORE", intent="first error", error="first error"), + LogEntry(level="INFO", trace_id="trace", src="plugin", marker="REASON", intent="noise"), + LogEntry(level="ERROR", trace_id="trace", src="plugin", marker="EXPLORE", intent="last error", error="last error"), ], ) ctx = extract_error_context(task, ReportStatus.FAILED) @@ -313,7 +313,7 @@ class TestExtractErrorContextLogFallback: from src.core.task_manager.models import LogEntry task = _make_task( result={"error": {"message": "", "code": "E1"}}, - logs=[LogEntry(level="ERROR", message="log fallback hit")], + logs=[LogEntry(level="ERROR", trace_id="trace", src="plugin", marker="EXPLORE", intent="log fallback hit", error="log fallback hit")], ) ctx = extract_error_context(task, ReportStatus.FAILED) assert ctx is not None diff --git a/backend/tests/test_app_ws_endpoint.py b/backend/tests/test_app_ws_endpoint.py index 1cf1abd82..ec2ac8bfd 100644 --- a/backend/tests/test_app_ws_endpoint.py +++ b/backend/tests/test_app_ws_endpoint.py @@ -202,8 +202,8 @@ class TestWebSocketEndpointFull: tm = _make_task_manager_mock(task=task, logs=[]) tm.subscribe_logs = sl; tm.subscribe_status = ss; mg.return_value = tm await websocket_endpoint(ws, "task-1") - msgs = [c[0][0].get("message") for c in ws.send_json.call_args_list if isinstance(c[0][0], dict)] - assert any("Task paused for user input" in (m or "") for m in msgs) + intents = [c[0][0].get("intent") for c in ws.send_json.call_args_list if isinstance(c[0][0], dict)] + assert any("Task paused for user input" in (intent or "") for intent in intents) # #endregion Test.AppModule.TestAwaitingInputPrompt # #region Test.AppModule.TestTerminalLogTriggersClose [C:2] [TYPE Function] diff --git a/frontend/src/lib/components/tasks/LogEntryRow.svelte b/frontend/src/lib/components/tasks/LogEntryRow.svelte index 4b6cb48db..35259bab0 100644 --- a/frontend/src/lib/components/tasks/LogEntryRow.svelte +++ b/frontend/src/lib/components/tasks/LogEntryRow.svelte @@ -9,7 +9,6 @@ let { log, showSource = true } = $props(); import { parseDateUTC } from "$lib/utils/dateFormat.js"; import { appTimezone } from "$lib/stores/timezone.svelte.js"; - import { parseCotMessage } from "$lib/logs/parseCot"; let expanded = $state(false); @@ -26,7 +25,7 @@ } let formattedTime = $derived(formatTime(log.timestamp)); - let cot = $derived(parseCotMessage(log.message)); + let cot = $derived(log.marker ? log : null); const levelStyles = { DEBUG: "text-log-debug bg-log-debug/10", @@ -60,15 +59,15 @@ levelStyles[log.level?.toUpperCase()] || levelStyles.INFO, ); let sourceClass = $derived( - sourceStyles[log.source?.toLowerCase()] || + sourceStyles[log.src?.toLowerCase()] || "bg-surface-muted text-text-muted", ); let markerClass = $derived( cot ? markerStyles[cot.marker] || "bg-surface-muted text-text-muted" : "", ); - let hasProgress = $derived(log.metadata?.progress !== undefined); - let progressPercent = $derived(log.metadata?.progress || 0); + let hasProgress = $derived(log.payload?.progress !== undefined); + let progressPercent = $derived(log.payload?.progress || 0); let hasDetails = $derived( Boolean(cot && (cot.payload || cot.error || cot.trace_id || cot.src)), ); @@ -82,7 +81,7 @@ } let displayText = $derived( - cot ? sanitizeMessage(cot.intent) : sanitizeMessage(log.message), + sanitizeMessage(log.intent || ""), ); let exploreBorder = $derived( cot?.marker === "EXPLORE" || String(log.level || "").toUpperCase() === "ERROR", @@ -110,9 +109,9 @@ {cot.marker} {/if} - {#if showSource && log.source} + {#if showSource && log.src} {log.source}{log.src} {/if} {#if cot?.trace_id} diff --git a/frontend/src/lib/components/tasks/TaskLogPanel.svelte b/frontend/src/lib/components/tasks/TaskLogPanel.svelte index 23971b2ac..5127d0c69 100644 --- a/frontend/src/lib/components/tasks/TaskLogPanel.svelte +++ b/frontend/src/lib/components/tasks/TaskLogPanel.svelte @@ -70,17 +70,17 @@ if ( source && source !== "all" && - log.source?.toLowerCase() !== source.toLowerCase() + ![log.src, log.src?.split('.').at(-1)].filter(Boolean).some((value) => value.toLowerCase() === source.toLowerCase()) ) return false; - if (search && !log.message?.toLowerCase().includes(search.toLowerCase())) + if (search && !log.intent?.toLowerCase().includes(search.toLowerCase())) return false; return true; }); } let availableSources = $derived([ - ...new Set(logs.map((l) => l.source).filter(Boolean)), + ...new Set(logs.map((l) => l.src).filter(Boolean)), ]); let activeQuickFilter = $derived( quickFilters.find( @@ -138,8 +138,8 @@ bind:searchText onfilterchange={handleFilterChange} /> -
- {#each quickFilters as filter} +
+ {#each quickFilters as filter (filter.id)}
{:else} - {#each filteredLogs as log} + {#each filteredLogs as log (log.id || `${log.timestamp}:${log.trace_id}:${log.intent}`)} {/each} {/if} diff --git a/frontend/src/lib/components/tasks/TaskLogViewer.svelte b/frontend/src/lib/components/tasks/TaskLogViewer.svelte index 6ffbe403f..7a77a0a42 100644 --- a/frontend/src/lib/components/tasks/TaskLogViewer.svelte +++ b/frontend/src/lib/components/tasks/TaskLogViewer.svelte @@ -14,6 +14,7 @@ -->