diff --git a/src/aiu_trace_analyzer/core/engine.py b/src/aiu_trace_analyzer/core/engine.py index 019406c..b3d3fea 100644 --- a/src/aiu_trace_analyzer/core/engine.py +++ b/src/aiu_trace_analyzer/core/engine.py @@ -30,6 +30,7 @@ def run(self) -> int: self.exporter.export(events) # drain the context buffers (if any) + # accumulated warnings ride along as trace_issue meta-events that the exporter captures drain = self.processor.drain() # export any events emitted during drain self.exporter.export(drain) diff --git a/src/aiu_trace_analyzer/core/processing.py b/src/aiu_trace_analyzer/core/processing.py index 7c2008a..b7b47f6 100644 --- a/src/aiu_trace_analyzer/core/processing.py +++ b/src/aiu_trace_analyzer/core/processing.py @@ -5,7 +5,7 @@ import aiu_trace_analyzer.trace_view as aiuev import aiu_trace_analyzer.pipeline.context as procCTX -from aiu_trace_analyzer.types import TraceEvent +from aiu_trace_analyzer.types import DiagnosticEvent, TraceEvent from aiu_trace_analyzer.core.duplicate_hold import IntermediateDuplicateAndHoldContext, duplicate_and_hold from aiu_trace_analyzer.export.exporter import JsonFileTraceExporter from aiu_trace_analyzer.core.stage_profile import StageProfile, StageProfileChecker @@ -81,6 +81,11 @@ def process(self, event: TraceEvent) -> list[aiuev.AbstractEventType]: # turn into a list, pre/post have do be able to expand single events into lists aiulog.log(aiulog.DEBUG, "Processing event:", event) + if isinstance(event, DiagnosticEvent): + output_event_list = self.convert_events([event]) + self.event_count += len(output_event_list) + return output_event_list + event_list = self.pre_process(event) output_event_list = self.convert_events(event_list) @@ -143,4 +148,5 @@ def drain(self) -> list[aiuev.AbstractEventType]: # then process the events that came back using the remaining pre-processing hooks + pipeline for event in pending: next_event_list += self.process(event) + return next_event_list diff --git a/src/aiu_trace_analyzer/export/exporter.py b/src/aiu_trace_analyzer/export/exporter.py index 28d174d..c338c56 100644 --- a/src/aiu_trace_analyzer/export/exporter.py +++ b/src/aiu_trace_analyzer/export/exporter.py @@ -10,6 +10,7 @@ import aiu_trace_analyzer.logger as aiulog import aiu_trace_analyzer.trace_view as tv +from aiu_trace_analyzer.types import TRACE_ISSUE_EVENT_NAME from aiu_trace_analyzer.verification.report import ( VERIFICATION_RESULT_NAME, VERIFICATION_TEST_RESULT_NAME, @@ -77,6 +78,13 @@ def __init__(self, target_uri, timescale="ms", settings=None) -> None: # take (a list) of events and append to the traceview def export(self, data: list[tv.AbstractEventType]): for event in data: + # trace_issue meta-events go into otherData, not into the trace event stream + if event.ph == "M" and event.name == TRACE_ISSUE_EVENT_NAME: + severity = "errors" if event.args["is_error"] else "warnings" + findings = self.traceview.other_data.setdefault(severity, []) + findings.append({"finding": event.args["finding"], + "text": event.args["text"]}) + continue self.traceview.append_trace_event(event.json()) def export_meta(self, meta_data): diff --git a/src/aiu_trace_analyzer/pipeline/context.py b/src/aiu_trace_analyzer/pipeline/context.py index afab6bb..a67246a 100644 --- a/src/aiu_trace_analyzer/pipeline/context.py +++ b/src/aiu_trace_analyzer/pipeline/context.py @@ -1,7 +1,7 @@ # Copyright 2024-2025 IBM Corporation import aiu_trace_analyzer.logger as aiulog -from aiu_trace_analyzer.types import TraceEvent, TraceWarning +from aiu_trace_analyzer.types import DiagnosticEvent, TraceEvent, TraceWarning, TRACE_ISSUE_EVENT_NAME class AbstractContext: @@ -52,6 +52,19 @@ def print_warnings(self) -> None: if w.has_warning(): aiulog.log(aiulog.WARN, w) + def emit_issue_events(self) -> list[TraceEvent]: + ''' + emit each active warning as a meta-event so the exporter can fold it into the output json. + ''' + return [ + DiagnosticEvent({"ph": "M", "ts": 0, "pid": 0, + "name": TRACE_ISSUE_EVENT_NAME, + "args": {"finding": name, + "text": str(w), + "is_error": w.is_error()}}) + for name, w in self.warnings.items() if w.has_warning() + ] + def add_warning(self, warning: TraceWarning): self.warnings[warning.get_name()] = warning @@ -72,13 +85,13 @@ def drain(self) -> list[TraceEvent]: a list of events. Events are drained following the sequence of registered processing functions. ''' - return [] + return self.emit_issue_events() def _emit_verification_events(self) -> list[TraceEvent]: return [ - TraceEvent({"ph": "M", "ts": 0, "pid": 0, - "name": "verification_data", - "args": w.to_verification_event_args()}) + DiagnosticEvent({"ph": "M", "ts": 0, "pid": 0, + "name": "verification_data", + "args": w.to_verification_event_args()}) for w in self.warnings.values() ] @@ -91,10 +104,10 @@ def _get_test_result_status(self) -> str: return "pass" def _emit_test_result_event(self, test_name: str) -> TraceEvent: - return TraceEvent({"ph": "M", "ts": 0, "pid": 0, - "name": "verification_test_result", - "args": {"test": test_name, - "result": self._get_test_result_status()}}) + return DiagnosticEvent({"ph": "M", "ts": 0, "pid": 0, + "name": "verification_test_result", + "args": {"test": test_name, + "result": self._get_test_result_status()}}) class AbstractVerificationContext(AbstractContext): diff --git a/src/aiu_trace_analyzer/types.py b/src/aiu_trace_analyzer/types.py index a981b03..c7589a2 100644 --- a/src/aiu_trace_analyzer/types.py +++ b/src/aiu_trace_analyzer/types.py @@ -11,6 +11,15 @@ class TraceEvent(dict): pass +class DiagnosticEvent(TraceEvent): + pass + + +# name of the meta-event used to carry an accumulated warning through the pipeline +# so the exporter can fold it into the output json instead of losing it in the console +TRACE_ISSUE_EVENT_NAME = "trace_issue" + + class InputDialect: categories = set() dialect_map = {} @@ -292,13 +301,19 @@ def update(self, def has_warning(self) -> bool: return self.occurred + def is_error(self) -> bool: + return self.warn_level == aiulog.ERROR + + def severity(self) -> str: + return "error" if self.is_error() else "warning" + def add_instance(self, data: dict) -> None: self._instances.append(data) def to_verification_event_args(self) -> dict: return { "finding": self.name, - "is_error": self.warn_level == aiulog.ERROR, + "is_error": self.is_error(), "count": self.args_list.get("count", len(self._instances)), "instances": list(self._instances), } diff --git a/tests/aiu_trace_analyzer/core/test_processing_issues.py b/tests/aiu_trace_analyzer/core/test_processing_issues.py new file mode 100644 index 0000000..0306396 --- /dev/null +++ b/tests/aiu_trace_analyzer/core/test_processing_issues.py @@ -0,0 +1,78 @@ +# Copyright 2024-2026 IBM Corporation + +import json + +import pytest + +from aiu_trace_analyzer.core.processing import EventProcessor +from aiu_trace_analyzer.pipeline.context import AbstractContext +from aiu_trace_analyzer.types import TraceWarning, TRACE_ISSUE_EVENT_NAME +from aiu_trace_analyzer.export.exporter import JsonFileTraceExporter + + +@pytest.fixture +def warned_context() -> AbstractContext: + warning = TraceWarning( + name="long_dur", + text="OVC: Detected {d[count]} long event(s).", + data={"count": 0}, + update_fn={"count": int.__add__}, + auto_log=False, + ) + ctx = AbstractContext(warnings=[warning]) + ctx.enable() + return ctx + + +def _processor_with(context: AbstractContext) -> EventProcessor: + proc = EventProcessor() + # register the stage directly to avoid pulling in a full StageProfile for the test + proc.stages.append((lambda event, ctx: [event], context, {})) + return proc + + +def test_drain_emits_active_warning_as_meta_event(warned_context): + warned_context.issue_warning("long_dur", {"count": 3}) + + drained = _processor_with(warned_context).drain() + + issue_events = [e for e in drained if e.name == TRACE_ISSUE_EVENT_NAME] + assert len(issue_events) == 1 + assert issue_events[0].args == {"finding": "long_dur", + "text": "OVC: Detected 3 long event(s).", + "is_error": False} + + +def test_drain_emits_nothing_when_no_warning(warned_context): + drained = _processor_with(warned_context).drain() + + assert [e for e in drained if e.name == TRACE_ISSUE_EVENT_NAME] == [] + + +def test_warning_reaches_exporter_other_data(warned_context): + warned_context.issue_warning("long_dur", {"count": 3}) + + drained = _processor_with(warned_context).drain() + exporter = JsonFileTraceExporter(target_uri="unused.json") + exporter.export(drained) + + output = json.loads(exporter.get_data()) + assert output["otherData"]["warnings"] == [{ + "finding": "long_dur", + "text": "OVC: Detected 3 long event(s).", + }] + assert output["traceEvents"] == [] + + +def test_drain_warning_bypasses_remaining_pipeline_stages(warned_context): + warned_context.issue_warning("long_dur", {"count": 3}) + proc = _processor_with(warned_context) + + def drop_everything(event, ctx): + return [] + + proc.stages.append((drop_everything, None, {})) + + drained = proc.drain() + + assert [e for e in drained if e.name == TRACE_ISSUE_EVENT_NAME] diff --git a/tests/aiu_trace_analyzer/export/test_exporter.py b/tests/aiu_trace_analyzer/export/test_exporter.py new file mode 100644 index 0000000..1126fc5 --- /dev/null +++ b/tests/aiu_trace_analyzer/export/test_exporter.py @@ -0,0 +1,63 @@ +# Copyright 2024-2026 IBM Corporation + +import json + +import pytest + +from aiu_trace_analyzer.types import TRACE_ISSUE_EVENT_NAME +from aiu_trace_analyzer.trace_view import AbstractEventType +from aiu_trace_analyzer.export.exporter import JsonFileTraceExporter + + +@pytest.fixture +def json_exporter() -> JsonFileTraceExporter: + return JsonFileTraceExporter(target_uri="unused.json") + + +def _issue_event(name: str, text: str, is_error: bool = False) -> AbstractEventType: + return AbstractEventType.from_dict({ + "ph": "M", "ts": 0, "pid": 0, + "name": TRACE_ISSUE_EVENT_NAME, + "args": {"finding": name, "text": text, "is_error": is_error}, + }) + + +def _instant_event() -> AbstractEventType: + return AbstractEventType.from_dict({ + "ph": "i", "ts": 1, "pid": 0, "tid": 0, "s": "g", + "name": "regular_event", "args": {}, + }) + + +def test_export_captures_issue_events(json_exporter): + text = "OVC: Detected 3 event(s) with long duration." + json_exporter.export([_issue_event("long_dur", text)]) + + other_data = json.loads(json_exporter.get_data())["otherData"] + assert other_data["warnings"] == [{"finding": "long_dur", "text": text}] + + +def test_export_separates_errors_from_warnings(json_exporter): + json_exporter.export([_issue_event("long_dur", "warn text"), + _issue_event("bad_ts", "error text", is_error=True)]) + + other_data = json.loads(json_exporter.get_data())["otherData"] + assert other_data["warnings"] == [{"finding": "long_dur", "text": "warn text"}] + assert other_data["errors"] == [{"finding": "bad_ts", "text": "error text"}] + + +def test_export_issue_events_do_not_leak_into_trace(json_exporter): + json_exporter.export([_issue_event("long_dur", "text"), _instant_event()]) + + dumped = json.loads(json_exporter.get_data()) + names = [e["name"] for e in dumped["traceEvents"]] + assert TRACE_ISSUE_EVENT_NAME not in names + assert "regular_event" in names + + +def test_export_no_issue_section_when_absent(json_exporter): + json_exporter.export([_instant_event()]) + + other_data = json.loads(json_exporter.get_data())["otherData"] + assert "warnings" not in other_data + assert "errors" not in other_data diff --git a/tests/aiu_trace_analyzer/pipeline/test_context.py b/tests/aiu_trace_analyzer/pipeline/test_context.py index 0c0c4ef..377b5d6 100644 --- a/tests/aiu_trace_analyzer/pipeline/test_context.py +++ b/tests/aiu_trace_analyzer/pipeline/test_context.py @@ -3,7 +3,7 @@ import pytest import math -from aiu_trace_analyzer.types import TraceWarning +from aiu_trace_analyzer.types import TraceWarning, TRACE_ISSUE_EVENT_NAME from aiu_trace_analyzer.pipeline import AbstractContext from aiu_trace_analyzer.pipeline.context import AbstractVerificationContext @@ -135,6 +135,45 @@ def test_issue_warning(abstract_context): assert abstract_context.warnings["pytest"].__str__() == "A Warning with 2 args: 2 and 5.0" +def test_emit_issue_events_none_when_inactive(abstract_context): + abstract_context.warnings["pytest"].auto_log = False # disable auto-output for tests + assert abstract_context.emit_issue_events() == [] + + +def test_emit_issue_events(abstract_context): + abstract_context.warnings["pytest"].auto_log = False # disable auto-output for tests + abstract_context.issue_warning("pytest", {"count": 1, "max": 5.0}) + + events = abstract_context.emit_issue_events() + + assert len(events) == 1 + assert events[0]["ph"] == "M" + assert events[0]["name"] == TRACE_ISSUE_EVENT_NAME + assert events[0]["args"] == {"finding": "pytest", + "text": "A Warning with 2 args: 1 and 5.0", + "is_error": False} + + +def test_emit_issue_events_of_error_warning(): + error = TraceWarning( + name="pytest_err", + text="An Error with {d[count]} occurrence(s)", + data={"count": 0}, + update_fn={"count": int.__add__}, + auto_log=False, + is_error=True, + ) + context = AbstractContext(warnings=[error]) + context.issue_warning("pytest_err", {"count": 1}) + + events = context.emit_issue_events() + + assert len(events) == 1 + assert events[0]["args"] == {"finding": "pytest_err", + "text": "An Error with 1 occurrence(s)", + "is_error": True} + + def test_drain(abstract_context): assert abstract_context.drain() == [] @@ -248,8 +287,15 @@ def test_v1_drain_no_warnings_produces_pass(verif_context_no_warnings): def test_v2_drain_warn_level_warning_produces_warn(verif_context_warn): verif_context_warn.warnings["test_w"].update({"count": 1}) events = verif_context_warn.drain() + issue_events = _find_events(events, TRACE_ISSUE_EVENT_NAME) test_result = _find_events(events, "verification_test_result")[0] assert test_result["args"]["result"] == "warn" + assert len(issue_events) == 1 + assert issue_events[0]["args"] == { + "finding": "test_w", + "text": "Found 1 issues", + "is_error": False, + } data_events = _find_events(events, "verification_data") assert len(data_events) == 1 assert data_events[0]["args"]["is_error"] is False @@ -258,8 +304,15 @@ def test_v2_drain_warn_level_warning_produces_warn(verif_context_warn): def test_v3_drain_error_level_warning_produces_fail(verif_context_error): verif_context_error.warnings["test_err"].update({"count": 1}) events = verif_context_error.drain() + issue_events = _find_events(events, TRACE_ISSUE_EVENT_NAME) test_result = _find_events(events, "verification_test_result")[0] assert test_result["args"]["result"] == "fail" + assert len(issue_events) == 1 + assert issue_events[0]["args"] == { + "finding": "test_err", + "text": "Found 1 errors", + "is_error": True, + } data_events = _find_events(events, "verification_data") assert data_events[0]["args"]["is_error"] is True