Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 5 additions & 8 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
2 changes: 2 additions & 0 deletions src/flowsense/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
ChangeDirection,
ChangePointResult,
DAGAnalysis,
DAGAnalysisSummary,
DriftResult,
FlowSenseError,
ImpactClassification,
Expand Down Expand Up @@ -38,6 +39,7 @@
"ChangeDirection",
"ChangePointResult",
"DAGAnalysis",
"DAGAnalysisSummary",
"DAGDataSource",
"DriftResult",
"FlowSenseError",
Expand Down
13 changes: 13 additions & 0 deletions src/flowsense/cli/report.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
2 changes: 2 additions & 0 deletions src/flowsense/domain/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
AnalysisDiagnostic,
ChangePointResult,
DAGAnalysis,
DAGAnalysisSummary,
DriftResult,
PropagationResult,
RootCauseResult,
Expand All @@ -30,6 +31,7 @@
"ChangeDirection",
"ChangePointResult",
"DAGAnalysis",
"DAGAnalysisSummary",
"DriftResult",
"FlowSenseError",
"ImpactClassification",
Expand Down
87 changes: 87 additions & 0 deletions src/flowsense/domain/results.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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),
)
17 changes: 17 additions & 0 deletions src/flowsense/mcp/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
3 changes: 3 additions & 0 deletions tests/test_cli_report.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
134 changes: 134 additions & 0 deletions tests/test_dag_summary.py
Original file line number Diff line number Diff line change
@@ -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
15 changes: 15 additions & 0 deletions tests/test_mcp_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"]

Expand Down
1 change: 1 addition & 0 deletions tests/test_public_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ def test_top_level_api_declares_supported_exports() -> None:
"ChangeDirection",
"ChangePointResult",
"DAGAnalysis",
"DAGAnalysisSummary",
"DAGDataSource",
"Severity",
"MappedTaskAggregation",
Expand Down
Loading