From c2064cfc8d231c00a9d2ac8ac804db5dcc8eca9e Mon Sep 17 00:00:00 2001 From: omercengiz Date: Fri, 4 Sep 2026 23:41:59 +0300 Subject: [PATCH] feat: add DAG analysis summary --- README.md | 13 ++- src/flowsense/__init__.py | 2 + src/flowsense/cli/report.py | 13 +++ src/flowsense/domain/__init__.py | 2 + src/flowsense/domain/results.py | 87 ++++++++++++++++++++ src/flowsense/mcp/server.py | 17 ++++ tests/test_cli_report.py | 3 + tests/test_dag_summary.py | 134 +++++++++++++++++++++++++++++++ tests/test_mcp_server.py | 15 ++++ tests/test_public_api.py | 1 + 10 files changed, 279 insertions(+), 8 deletions(-) create mode 100644 tests/test_dag_summary.py diff --git a/README.md b/README.md index dca3c2c..2142f6c 100644 --- a/README.md +++ b/README.md @@ -184,6 +184,11 @@ policy = AnalysisPolicy( run. Mapped task durations can be aggregated with `MAX`, `MEAN`, or `SUM`. The same policy options are available through the CLI and MCP tool. +Every `DAGAnalysis` exposes a derived `summary` with task-analysis coverage, +severity distribution, anomalous task and handoff counts, uniquely affected +tasks, structural signals, and diagnostics. The same DAG-level summary is +included in CLI and MCP output. + Custom data sources can implement the `DAGDataSource` protocol and be passed to `analyze_dag`. Names exported directly from `flowsense` form the supported public API. Imports from internal packages such as `flowsense.engine` should be treated @@ -290,14 +295,6 @@ FlowSense is currently in early development. The current implementation should be considered experimental and is not yet intended for production use. -## Roadmap - -Planned areas include: - -- DAG-level analysis models -- richer CLI reporting -- broader Airflow compatibility testing - ## License Licensed under the Apache License 2.0. diff --git a/src/flowsense/__init__.py b/src/flowsense/__init__.py index 7a9365a..3776ee3 100644 --- a/src/flowsense/__init__.py +++ b/src/flowsense/__init__.py @@ -8,6 +8,7 @@ ChangeDirection, ChangePointResult, DAGAnalysis, + DAGAnalysisSummary, DriftResult, FlowSenseError, ImpactClassification, @@ -38,6 +39,7 @@ "ChangeDirection", "ChangePointResult", "DAGAnalysis", + "DAGAnalysisSummary", "DAGDataSource", "DriftResult", "FlowSenseError", diff --git a/src/flowsense/cli/report.py b/src/flowsense/cli/report.py index bf20750..bac8061 100644 --- a/src/flowsense/cli/report.py +++ b/src/flowsense/cli/report.py @@ -25,12 +25,25 @@ def _percent(value: float | None) -> str: def _render_summary(console: Console, analysis: DAGAnalysis) -> None: + dag_summary = analysis.summary summary = Table.grid(padding=(0, 2)) summary.add_column(style="bold") summary.add_column() summary.add_row("DAG", analysis.dag_id) summary.add_row("Runs analyzed", str(analysis.runs_analyzed)) summary.add_row("Overall severity", _severity_text(analysis.overall_severity)) + summary.add_row( + "Task coverage", + f"{dag_summary.analyzed_tasks}/{dag_summary.total_tasks} " + f"({dag_summary.analysis_coverage_percent:.1f}%)", + ) + summary.add_row("Anomalous tasks", str(dag_summary.anomalous_tasks)) + summary.add_row("Anomalous handoffs", str(dag_summary.anomalous_handoffs)) + summary.add_row("Affected tasks", str(dag_summary.affected_tasks)) + summary.add_row( + "Structural signals", + f"{dag_summary.change_points} change points, {dag_summary.trends} trends", + ) if analysis.primary_origin is not None: summary.add_row("Primary origin", analysis.primary_origin.task_id) diff --git a/src/flowsense/domain/__init__.py b/src/flowsense/domain/__init__.py index be6ec75..b6c1655 100644 --- a/src/flowsense/domain/__init__.py +++ b/src/flowsense/domain/__init__.py @@ -16,6 +16,7 @@ AnalysisDiagnostic, ChangePointResult, DAGAnalysis, + DAGAnalysisSummary, DriftResult, PropagationResult, RootCauseResult, @@ -30,6 +31,7 @@ "ChangeDirection", "ChangePointResult", "DAGAnalysis", + "DAGAnalysisSummary", "DriftResult", "FlowSenseError", "ImpactClassification", diff --git a/src/flowsense/domain/results.py b/src/flowsense/domain/results.py index 4ea9914..3c4d572 100644 --- a/src/flowsense/domain/results.py +++ b/src/flowsense/domain/results.py @@ -76,6 +76,23 @@ class AnalysisDiagnostic: message: str +@dataclass(frozen=True) +class DAGAnalysisSummary: + total_tasks: int + analyzed_tasks: int + analysis_coverage_percent: float + normal_tasks: int + medium_tasks: int + high_tasks: int + critical_tasks: int + anomalous_tasks: int + anomalous_handoffs: int + affected_tasks: int + change_points: int + trends: int + diagnostics: int + + @dataclass class DAGAnalysis: dag_id: str @@ -97,3 +114,73 @@ class DAGAnalysis: handoff_trend_results: dict[tuple[str, str], TrendResult] = field( default_factory=dict ) + + @property + def summary(self) -> DAGAnalysisSummary: + task_ids = { + *self.dependencies, + *(task for tasks in self.dependencies.values() for task in tasks), + *self.drift_results, + *self.task_impacts, + *self.change_point_results, + *self.trend_results, + *( + task + for edge in ( + *self.handoff_drift_results, + *self.handoff_change_point_results, + *self.handoff_trend_results, + ) + for task in edge + ), + *( + task + for result in self.propagation_results + for task in (result.origin_task, *result.affected_tasks) + ), + } + severity_counts = { + severity: sum( + result.severity == severity for result in self.drift_results.values() + ) + for severity in Severity + } + analyzed_tasks = len(self.drift_results) + total_tasks = len(task_ids) + anomalous_severities = { + Severity.MEDIUM, + Severity.HIGH, + Severity.CRITICAL, + } + + return DAGAnalysisSummary( + total_tasks=total_tasks, + analyzed_tasks=analyzed_tasks, + analysis_coverage_percent=( + analyzed_tasks / total_tasks * 100 if total_tasks else 0.0 + ), + normal_tasks=severity_counts[Severity.NORMAL], + medium_tasks=severity_counts[Severity.MEDIUM], + high_tasks=severity_counts[Severity.HIGH], + critical_tasks=severity_counts[Severity.CRITICAL], + anomalous_tasks=sum( + result.severity in anomalous_severities + for result in self.drift_results.values() + ), + anomalous_handoffs=sum( + result.severity in anomalous_severities + for result in self.handoff_drift_results.values() + ), + affected_tasks=len( + { + task + for result in self.propagation_results + for task in result.affected_tasks + } + ), + change_points=( + len(self.change_point_results) + len(self.handoff_change_point_results) + ), + trends=len(self.trend_results) + len(self.handoff_trend_results), + diagnostics=len(self.diagnostics), + ) diff --git a/src/flowsense/mcp/server.py b/src/flowsense/mcp/server.py index 6a02922..94d8d84 100644 --- a/src/flowsense/mcp/server.py +++ b/src/flowsense/mcp/server.py @@ -14,10 +14,27 @@ def serialize_analysis(analysis: DAGAnalysis) -> dict: + summary = analysis.summary + return { "dag_id": analysis.dag_id, "runs_analyzed": analysis.runs_analyzed, "overall_severity": analysis.overall_severity, + "summary": { + "total_tasks": summary.total_tasks, + "analyzed_tasks": summary.analyzed_tasks, + "analysis_coverage_percent": summary.analysis_coverage_percent, + "normal_tasks": summary.normal_tasks, + "medium_tasks": summary.medium_tasks, + "high_tasks": summary.high_tasks, + "critical_tasks": summary.critical_tasks, + "anomalous_tasks": summary.anomalous_tasks, + "anomalous_handoffs": summary.anomalous_handoffs, + "affected_tasks": summary.affected_tasks, + "change_points": summary.change_points, + "trends": summary.trends, + "diagnostics": summary.diagnostics, + }, "policy": { "minimum_history": analysis.policy.minimum_history, "baseline_window": analysis.policy.baseline_window, diff --git a/tests/test_cli_report.py b/tests/test_cli_report.py index 8182754..348f303 100644 --- a/tests/test_cli_report.py +++ b/tests/test_cli_report.py @@ -101,6 +101,9 @@ def test_renders_complete_analysis_report() -> None: report = output.getvalue() assert "FlowSense Analysis" in report + assert "Task coverage" in report + assert "1/3 (33.3%)" in report + assert "Structural signals" in report assert "Task Drift" in report assert "Handoff Drift" in report assert "Change Points" in report diff --git a/tests/test_dag_summary.py b/tests/test_dag_summary.py new file mode 100644 index 0000000..82535ee --- /dev/null +++ b/tests/test_dag_summary.py @@ -0,0 +1,134 @@ +from flowsense import ( + AnalysisDiagnostic, + ChangePointResult, + DAGAnalysis, + DriftResult, + PropagationResult, + TrendResult, +) + + +def _drift(task_id: str, severity: str) -> DriftResult: + return DriftResult( + task_id=task_id, + baseline=1.0, + current=2.0, + mad=0.1, + robust_z_score=3.0, + deviation_percent=100.0, + severity=severity, + ) + + +def test_summarizes_dag_analysis_results() -> None: + analysis = DAGAnalysis( + dag_id="demo", + runs_analyzed=8, + overall_severity="CRITICAL", + primary_origin=None, + drift_results={ + "extract": _drift("extract", "NORMAL"), + "transform": _drift("transform", "CRITICAL"), + "load": _drift("load", "MEDIUM"), + }, + handoff_drift_results={ + ("extract", "transform"): _drift("extract->transform", "HIGH"), + ("transform", "load"): _drift("transform->load", "NORMAL"), + }, + task_impacts={}, + propagation_results=[ + PropagationResult( + origin_task="transform", + affected_tasks=["load", "notify"], + path=["transform", "load", "notify"], + propagation_score=0.7, + ), + PropagationResult( + origin_task="transform", + affected_tasks=["load"], + path=["transform", "load"], + propagation_score=0.8, + ), + ], + dependencies={ + "extract": ["transform"], + "transform": ["load"], + "load": ["notify"], + "notify": [], + }, + diagnostics=[ + AnalysisDiagnostic( + code="INSUFFICIENT_TASK_HISTORY", + subject_id="notify", + message="Not enough observations.", + ) + ], + change_point_results={ + "transform": ChangePointResult( + subject_id="transform", + change_index=4, + before_median=1.0, + after_median=2.0, + change_percent=100.0, + score=5.0, + direction="INCREASE", + ) + }, + handoff_change_point_results={ + ("extract", "transform"): ChangePointResult( + subject_id="extract->transform", + change_index=4, + before_median=1.0, + after_median=2.0, + change_percent=100.0, + score=5.0, + direction="INCREASE", + ) + }, + trend_results={ + "load": TrendResult( + subject_id="load", + direction="INCREASING", + slope_per_observation=0.2, + estimated_change=1.4, + change_percent=70.0, + score=4.0, + directional_consistency=0.8, + observations=8, + ) + }, + ) + + summary = analysis.summary + + assert summary.total_tasks == 4 + assert summary.analyzed_tasks == 3 + assert summary.analysis_coverage_percent == 75.0 + assert summary.normal_tasks == 1 + assert summary.medium_tasks == 1 + assert summary.high_tasks == 0 + assert summary.critical_tasks == 1 + assert summary.anomalous_tasks == 2 + assert summary.anomalous_handoffs == 1 + assert summary.affected_tasks == 2 + assert summary.change_points == 2 + assert summary.trends == 1 + assert summary.diagnostics == 1 + + +def test_empty_analysis_has_zero_summary() -> None: + analysis = DAGAnalysis( + dag_id="empty", + runs_analyzed=0, + overall_severity="NORMAL", + primary_origin=None, + drift_results={}, + handoff_drift_results={}, + task_impacts={}, + propagation_results=[], + dependencies={}, + ) + + assert analysis.summary.total_tasks == 0 + assert analysis.summary.analyzed_tasks == 0 + assert analysis.summary.analysis_coverage_percent == 0.0 diff --git a/tests/test_mcp_server.py b/tests/test_mcp_server.py index 19a0760..39d1356 100644 --- a/tests/test_mcp_server.py +++ b/tests/test_mcp_server.py @@ -156,6 +156,21 @@ def test_serialize_analysis() -> None: assert result["dag_id"] == "demo" assert result["runs_analyzed"] == 5 assert result["overall_severity"] == "CRITICAL" + assert result["summary"] == { + "total_tasks": 2, + "analyzed_tasks": 1, + "analysis_coverage_percent": 50.0, + "normal_tasks": 0, + "medium_tasks": 0, + "high_tasks": 0, + "critical_tasks": 1, + "anomalous_tasks": 1, + "anomalous_handoffs": 1, + "affected_tasks": 0, + "change_points": 2, + "trends": 2, + "diagnostics": 1, + } primary_origin = result["primary_origin"] diff --git a/tests/test_public_api.py b/tests/test_public_api.py index 7a42d27..1a8d575 100644 --- a/tests/test_public_api.py +++ b/tests/test_public_api.py @@ -27,6 +27,7 @@ def test_top_level_api_declares_supported_exports() -> None: "ChangeDirection", "ChangePointResult", "DAGAnalysis", + "DAGAnalysisSummary", "DAGDataSource", "Severity", "MappedTaskAggregation",