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
6 changes: 5 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,11 @@ Then run:
flowsense analyze <dag_id>
```

The CLI report includes a DAG summary and separate tables for task drift,
handoff drift, change points, trends, propagation paths, and diagnostics.
Results are ordered by severity or subject so repeated analyses remain easy to
compare.

## Library API

FlowSense can also be used as a Python library through its supported top-level
Expand Down Expand Up @@ -283,7 +288,6 @@ The current implementation should be considered experimental and is not yet inte
Planned areas include:

- DAG-level analysis models
- richer CLI reporting
- broader Airflow compatibility testing

## License
Expand Down
91 changes: 2 additions & 89 deletions src/flowsense/cli/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,9 @@

import typer
from rich.console import Console
from rich.table import Table

from flowsense.application import analyze_dag
from flowsense.cli.report import render_analysis
from flowsense.domain import AnalysisPolicy, MappedTaskAggregation
from flowsense.infrastructure.airflow import AirflowApiError, AirflowClient

Expand Down Expand Up @@ -63,91 +63,4 @@ def analyze(
console.print(f"[bold red]Airflow request failed:[/bold red] {exc}")
raise typer.Exit(code=1) from exc

console.print(f"\n[bold]FlowSense Analysis — {analysis.dag_id}[/bold]\n")

table = Table()

table.add_column("Task")
table.add_column("Baseline")
table.add_column("Current")
table.add_column("Deviation")
table.add_column("Z-Score")
table.add_column("Severity")
table.add_column("Impact")

for task_id, result in analysis.drift_results.items():
impact = analysis.task_impacts.get(task_id)
impact_label = impact.classification if impact else "-"

table.add_row(
task_id,
f"{result.baseline:.2f}s",
f"{result.current:.2f}s",
f"{result.deviation_percent:+.1f}%",
f"{result.robust_z_score:.2f}",
result.severity,
impact_label,
)

console.print(table)

console.print(f"\nOverall Severity: [bold]{analysis.overall_severity}[/bold]")

if analysis.change_point_results or analysis.handoff_change_point_results:
console.print("\n[bold]Change Points[/bold]\n")

for result in [
*analysis.change_point_results.values(),
*analysis.handoff_change_point_results.values(),
]:
change = (
f"{result.change_percent:+.1f}%"
if result.change_percent is not None
else "n/a"
)
console.print(
f"{result.subject_id}: {result.direction} at observation "
f"{result.change_index + 1} ({change}, score={result.score:.2f})"
)

if analysis.trend_results or analysis.handoff_trend_results:
console.print("\n[bold]Trends[/bold]\n")

for result in [
*analysis.trend_results.values(),
*analysis.handoff_trend_results.values(),
]:
change = (
f"{result.change_percent:+.1f}%"
if result.change_percent is not None
else "n/a"
)
console.print(
f"{result.subject_id}: {result.direction} "
f"({result.slope_per_observation:+.2f}/run, {change}, "
f"score={result.score:.2f})"
)

if analysis.primary_origin:
console.print(f"Primary Origin: [bold]{analysis.primary_origin.task_id}[/bold]")
console.print(f"Reason: [bold]{analysis.primary_origin.classification}[/bold]")
console.print(f"Severity: [bold]{analysis.primary_origin.severity}[/bold]")
console.print(
f"Propagation Score: {analysis.primary_origin.propagation_score:.2f}"
)

if analysis.propagation_results:
console.print("\n[bold]Propagation Analysis[/bold]\n")

for result in analysis.propagation_results:
console.print(f"Origin: {result.origin_task}")
console.print(f"Path: {' -> '.join(result.path)}")
console.print(f"Propagation Score: {result.propagation_score:.2f}")

if analysis.diagnostics:
console.print("\n[bold yellow]Diagnostics[/bold yellow]\n")

for diagnostic in analysis.diagnostics:
console.print(
f"[{diagnostic.code}] {diagnostic.subject_id}: {diagnostic.message}"
)
render_analysis(console, analysis)
218 changes: 218 additions & 0 deletions src/flowsense/cli/report.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,218 @@
from __future__ import annotations

from rich.console import Console
from rich.panel import Panel
from rich.table import Table
from rich.text import Text

from flowsense.domain import DAGAnalysis, Severity
from flowsense.domain.enums import SEVERITY_SCORE

_SEVERITY_STYLES = {
Severity.NORMAL: "green",
Severity.MEDIUM: "yellow",
Severity.HIGH: "bright_red",
Severity.CRITICAL: "bold red",
}


def _severity_text(severity: Severity) -> Text:
return Text(str(severity), style=_SEVERITY_STYLES[severity])


def _percent(value: float | None) -> str:
return f"{value:+.1f}%" if value is not None else "n/a"


def _render_summary(console: Console, analysis: DAGAnalysis) -> None:
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))

if analysis.primary_origin is not None:
summary.add_row("Primary origin", analysis.primary_origin.task_id)
summary.add_row("Classification", str(analysis.primary_origin.classification))
summary.add_row(
"Propagation score",
f"{analysis.primary_origin.propagation_score:.2f}",
)

console.print(Panel(summary, title="FlowSense Analysis", expand=False))


def _render_task_drift(console: Console, analysis: DAGAnalysis) -> None:
table = Table(title="Task Drift")
table.add_column("Task")
table.add_column("Baseline", justify="right")
table.add_column("Current", justify="right")
table.add_column("Deviation", justify="right")
table.add_column("Z-Score", justify="right")
table.add_column("Severity")
table.add_column("Impact")

ordered_results = sorted(
analysis.drift_results.items(),
key=lambda item: (-SEVERITY_SCORE[item[1].severity], item[0]),
)

for task_id, result in ordered_results:
impact = analysis.task_impacts.get(task_id)
table.add_row(
task_id,
f"{result.baseline:.2f}s",
f"{result.current:.2f}s",
f"{result.deviation_percent:+.1f}%",
f"{result.robust_z_score:.2f}",
_severity_text(result.severity),
str(impact.classification) if impact else "-",
)

console.print(table)


def _render_handoff_drift(console: Console, analysis: DAGAnalysis) -> None:
if not analysis.handoff_drift_results:
return

table = Table(title="Handoff Drift")
table.add_column("Edge")
table.add_column("Baseline", justify="right")
table.add_column("Current", justify="right")
table.add_column("Deviation", justify="right")
table.add_column("Z-Score", justify="right")
table.add_column("Severity")

ordered_results = sorted(
analysis.handoff_drift_results.items(),
key=lambda item: (-SEVERITY_SCORE[item[1].severity], item[0]),
)

for (upstream, downstream), result in ordered_results:
table.add_row(
f"{upstream} -> {downstream}",
f"{result.baseline:.2f}s",
f"{result.current:.2f}s",
f"{result.deviation_percent:+.1f}%",
f"{result.robust_z_score:.2f}",
_severity_text(result.severity),
)

console.print(table)


def _render_change_points(console: Console, analysis: DAGAnalysis) -> None:
results = [
*analysis.change_point_results.values(),
*analysis.handoff_change_point_results.values(),
]
if not results:
return

table = Table(title="Change Points")
table.add_column("Subject")
table.add_column("Direction")
table.add_column("Observation", justify="right")
table.add_column("Before", justify="right")
table.add_column("After", justify="right")
table.add_column("Change", justify="right")
table.add_column("Score", justify="right")

for result in sorted(results, key=lambda item: item.subject_id):
table.add_row(
result.subject_id,
str(result.direction),
str(result.change_index + 1),
f"{result.before_median:.2f}s",
f"{result.after_median:.2f}s",
_percent(result.change_percent),
f"{result.score:.2f}",
)

console.print(table)


def _render_trends(console: Console, analysis: DAGAnalysis) -> None:
results = [
*analysis.trend_results.values(),
*analysis.handoff_trend_results.values(),
]
if not results:
return

table = Table(title="Trends")
table.add_column("Subject")
table.add_column("Direction")
table.add_column("Slope / run", justify="right")
table.add_column("Est. change", justify="right")
table.add_column("Change", justify="right")
table.add_column("Consistency", justify="right")
table.add_column("Score", justify="right")

for result in sorted(results, key=lambda item: item.subject_id):
table.add_row(
result.subject_id,
str(result.direction),
f"{result.slope_per_observation:+.2f}s",
f"{result.estimated_change:+.2f}s",
_percent(result.change_percent),
f"{result.directional_consistency:.0%}",
f"{result.score:.2f}",
)

console.print(table)


def _render_propagation(console: Console, analysis: DAGAnalysis) -> None:
if not analysis.propagation_results:
return

table = Table(title="Propagation")
table.add_column("Origin")
table.add_column("Path")
table.add_column("Affected", justify="right")
table.add_column("Score", justify="right")

for result in sorted(
analysis.propagation_results,
key=lambda item: (item.origin_task, item.path),
):
table.add_row(
result.origin_task,
" -> ".join(result.path),
str(len(result.affected_tasks)),
f"{result.propagation_score:.2f}",
)

console.print(table)


def _render_diagnostics(console: Console, analysis: DAGAnalysis) -> None:
if not analysis.diagnostics:
return

table = Table(title="Diagnostics", title_style="bold yellow")
table.add_column("Code", style="yellow")
table.add_column("Subject")
table.add_column("Message")

for diagnostic in sorted(
analysis.diagnostics,
key=lambda item: (item.code, item.subject_id),
):
table.add_row(diagnostic.code, diagnostic.subject_id, diagnostic.message)

console.print(table)


def render_analysis(console: Console, analysis: DAGAnalysis) -> None:
"""Render a complete human-readable analysis report."""
_render_summary(console, analysis)
_render_task_drift(console, analysis)
_render_handoff_drift(console, analysis)
_render_change_points(console, analysis)
_render_trends(console, analysis)
_render_propagation(console, analysis)
_render_diagnostics(console, analysis)
Loading
Loading