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
27 changes: 27 additions & 0 deletions src/scherlok/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -1341,6 +1341,23 @@ def _table_health(
}


def _relative_age(detected_at: str, now: datetime | None = None) -> str:
"""Short relative age ("2h ago") for a stored ISO timestamp."""
past = datetime.fromisoformat(detected_at)
if past.tzinfo is None:
past = past.replace(tzinfo=timezone.utc)
seconds = max(0, int(((now or datetime.now(timezone.utc)) - past).total_seconds()))
minutes, _ = divmod(seconds, 60)
if minutes < 1:
return "just now"
if minutes < 60:
return f"{minutes}m ago"
hours, _ = divmod(minutes, 60)
if hours < 24:
return f"{hours}h ago"
return f"{hours // 24}d ago"


@app.command()
def status(
output: str = typer.Option(
Expand Down Expand Up @@ -1374,6 +1391,7 @@ def status(
tables = connector.list_tables()

records: list[dict] = []
latest = store.get_latest_anomaly_per_table()
for table in tables:
vol = store.get_latest_profile(table, "volume")
sch = store.get_latest_profile(table, "schema")
Expand All @@ -1383,6 +1401,7 @@ def status(
"columns": len(sch["columns"]) if sch else None,
"status": _table_health(connector, store, table, vol, sch),
"last_profiled": vol.get("timestamp") if vol else None,
"last_anomaly": latest.get(table),
})

if json_mode:
Expand All @@ -1394,13 +1413,21 @@ def status(
tbl.add_column("Columns", justify="right")
tbl.add_column("Last Profiled")
tbl.add_column("Status")
tbl.add_column("Last Anomaly")
for r in records:
anomaly = r["last_anomaly"]
last = (
f"{anomaly['type']} · {_relative_age(anomaly['detected_at'])}"
if anomaly
else ""
)
tbl.add_row(
r["table"],
str(r["rows"]) if r["rows"] is not None else "—",
str(r["columns"]) if r["columns"] is not None else "—",
r["last_profiled"] or "—",
_STATUS_RICH[r["status"]],
last,
)
console.print(tbl)
finally:
Expand Down
23 changes: 23 additions & 0 deletions src/scherlok/store/sqlite.py
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,29 @@ def get_anomaly_history(self, days: int = 30) -> list[dict[str, Any]]:
for row in rows
]

def get_latest_anomaly_per_table(self) -> dict[str, dict[str, str]]:
"""Latest anomaly for every table in one query.

Latest means highest detected_at, with the anomaly id as the
deterministic tie-breaker when timestamps match. No history
window: a table that fired long ago still reports it.
"""
rows = self._conn.execute(
"SELECT table_name, anomaly_type, severity, detected_at FROM ("
"SELECT table_name, anomaly_type, severity, detected_at, "
"ROW_NUMBER() OVER (PARTITION BY table_name "
"ORDER BY detected_at DESC, id DESC) AS rn FROM anomalies"
") WHERE rn = 1",
).fetchall()
return {
row["table_name"]: {
"type": row["anomaly_type"],
"severity": row["severity"],
"detected_at": row["detected_at"],
}
for row in rows
}

def close(self) -> None:
"""Close the database connection."""
self._conn.close()
116 changes: 116 additions & 0 deletions tests/test_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,7 @@ def test_produces_valid_json_array(self):
"volume": {"row_count": 100, "timestamp": "2026-08-31T12:00:00+00:00"},
"schema": {"columns": [{"name": "id"}]},
}.get(pt)
store.get_latest_anomaly_per_table.return_value = {}

result = runner.invoke(app, ["status", "--output", "json"])

Expand All @@ -185,6 +186,7 @@ def test_produces_valid_json_array(self):
"columns": 1,
"status": "healthy",
"last_profiled": "2026-08-31T12:00:00+00:00",
"last_anomaly": None,
}

def test_unprofiled_table_has_null_fields(self):
Expand All @@ -197,6 +199,7 @@ def test_unprofiled_table_has_null_fields(self):
store = MagicMock()
mock_store_cls.return_value = store
store.get_latest_profile.return_value = None
store.get_latest_anomaly_per_table.return_value = {}

result = runner.invoke(app, ["status", "--output", "json"])

Expand All @@ -205,6 +208,7 @@ def test_unprofiled_table_has_null_fields(self):
assert data[0]["columns"] is None
assert data[0]["last_profiled"] is None
assert data[0]["status"] == "unknown"
assert data[0]["last_anomaly"] is None

def test_multiple_tables(self):
connector = _mock_connector(tables=["orders", "users"])
Expand All @@ -219,6 +223,7 @@ def test_multiple_tables(self):
"volume": {"row_count": 50, "timestamp": "2026-08-30T00:00:00+00:00"},
"schema": {"columns": [{"name": "id"}, {"name": "name"}]},
}.get(pt)
store.get_latest_anomaly_per_table.return_value = {}

result = runner.invoke(app, ["status", "--output", "json"])

Expand All @@ -242,19 +247,130 @@ def test_text_mode_still_works(self):
"volume": {"row_count": 100, "timestamp": "2026-08-31T12:00:00+00:00"},
"schema": {"columns": [{"name": "id"}]},
}.get(pt)
store.get_latest_anomaly_per_table.return_value = {}

result = runner.invoke(app, ["status"])

assert result.exit_code == 0
output = ANSI_RE.sub("", result.output)
assert "Table Health" in output
assert "users" in output
assert "Last Anomaly" in output

def test_json_includes_last_anomaly(self):
connector = _mock_connector()
with (
patch("scherlok.cli._get_connector_or_exit", return_value=connector),
patch("scherlok.cli.ProfileStore") as mock_store_cls,
patch("scherlok.cli._table_health", return_value="critical"),
):
store = MagicMock()
mock_store_cls.return_value = store
store.get_latest_profile.side_effect = lambda _t, pt: {
"volume": {"row_count": 100, "timestamp": "2026-08-31T12:00:00+00:00"},
"schema": {"columns": [{"name": "id"}]},
}.get(pt)
store.get_latest_anomaly_per_table.return_value = {
"users": {
"type": "volume_drop",
"severity": "CRITICAL",
"detected_at": "2026-08-31T10:00:00+00:00",
}
}

result = runner.invoke(app, ["status", "--output", "json"])

assert result.exit_code == 0
data = json.loads(result.output)
assert data[0]["last_anomaly"] == {
"type": "volume_drop",
"severity": "CRITICAL",
"detected_at": "2026-08-31T10:00:00+00:00",
}
assert data[0]["table"] == "users"
assert data[0]["status"] == "critical"

def test_text_shows_last_anomaly_column(self):
from datetime import datetime, timedelta, timezone

detected_at = (datetime.now(timezone.utc) - timedelta(hours=2)).isoformat()
connector = _mock_connector()
with (
patch("scherlok.cli._get_connector_or_exit", return_value=connector),
patch("scherlok.cli.ProfileStore") as mock_store_cls,
patch("scherlok.cli._table_health", return_value="critical"),
):
store = MagicMock()
mock_store_cls.return_value = store
store.get_latest_profile.side_effect = lambda _t, pt: {
"volume": {"row_count": 100, "timestamp": "2026-08-31T12:00:00+00:00"},
"schema": {"columns": [{"name": "id"}]},
}.get(pt)
store.get_latest_anomaly_per_table.return_value = {
"users": {
"type": "volume_drop",
"severity": "CRITICAL",
"detected_at": detected_at,
}
}

result = runner.invoke(
app, ["status"], env={"NO_COLOR": "1", "COLUMNS": "200"}
)

assert result.exit_code == 0
output = ANSI_RE.sub("", result.output)
assert "Last Anomaly" in output
assert "volume_drop · 2h ago" in output

def test_text_leaves_last_anomaly_empty_without_anomalies(self):
connector = _mock_connector()
with (
patch("scherlok.cli._get_connector_or_exit", return_value=connector),
patch("scherlok.cli.ProfileStore") as mock_store_cls,
patch("scherlok.cli._table_health", return_value="healthy"),
):
store = MagicMock()
mock_store_cls.return_value = store
store.get_latest_profile.side_effect = lambda _t, pt: {
"volume": {"row_count": 100, "timestamp": "2026-08-31T12:00:00+00:00"},
"schema": {"columns": [{"name": "id"}]},
}.get(pt)
store.get_latest_anomaly_per_table.return_value = {}

result = runner.invoke(app, ["status"])

assert result.exit_code == 0
output = ANSI_RE.sub("", result.output)
assert "Last Anomaly" in output
assert "ago" not in output

def test_invalid_output_value_exits_1(self):
result = runner.invoke(app, ["status", "--output", "xml"])
assert result.exit_code == 1


def test_relative_age_formats_minutes_hours_days():
from datetime import datetime, timezone

from scherlok.cli import _relative_age

now = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
assert _relative_age("2026-09-01T11:59:30+00:00", now) == "just now"
assert _relative_age("2026-09-01T11:15:00+00:00", now) == "45m ago"
assert _relative_age("2026-09-01T10:00:00+00:00", now) == "2h ago"
assert _relative_age("2026-08-29T12:00:00+00:00", now) == "3d ago"


def test_relative_age_clamps_future_timestamps():
from datetime import datetime, timezone

from scherlok.cli import _relative_age

now = datetime(2026, 9, 1, 12, 0, 0, tzinfo=timezone.utc)
assert _relative_age("2026-09-01T13:00:00+00:00", now) == "just now"


class TestHistoryJson:
def test_produces_valid_json_array(self):
with patch("scherlok.cli.ProfileStore") as mock_store_cls:
Expand Down
85 changes: 85 additions & 0 deletions tests/test_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -99,3 +99,88 @@ def test_anomaly_history_empty(self):
history = store.get_anomaly_history(days=30)
assert history == []
store.close()

def test_latest_anomaly_per_table(self):
store = _temp_store()
store.save_anomalies([
{
"table": "users",
"type": "volume_drop",
"severity": Severity.CRITICAL,
"message": "Row count dropped 60%",
},
{
"table": "orders",
"type": "schema_drift",
"severity": Severity.WARNING,
"message": "Column added: email",
},
])
store.save_anomalies([
{
"table": "users",
"type": "freshness_gap",
"severity": Severity.WARNING,
"message": "Table is stale",
},
])
latest = store.get_latest_anomaly_per_table()
assert latest["users"]["type"] == "freshness_gap"
assert latest["users"]["severity"] == "WARNING"
assert latest["users"]["detected_at"]
assert latest["orders"] == {
"type": "schema_drift",
"severity": "WARNING",
"detected_at": latest["orders"]["detected_at"],
}
assert "payments" not in latest
store.close()

def test_latest_anomaly_per_table_tie_breaks_on_id(self):
store = _temp_store()
detected_at = "2026-09-01T00:00:00+00:00"
for anomaly_type, severity, message in [
("volume_drop", "CRITICAL", "first"),
("schema_drift", "WARNING", "second"),
]:
store._conn.execute(
"INSERT INTO anomalies "
"(table_name, anomaly_type, severity, message, detected_at) "
"VALUES (?, ?, ?, ?, ?)",
("users", anomaly_type, severity, message, detected_at),
)
store._conn.commit()
latest = store.get_latest_anomaly_per_table()
assert latest["users"]["type"] == "schema_drift"
store.close()

def test_latest_anomaly_per_table_empty(self):
store = _temp_store()
assert store.get_latest_anomaly_per_table() == {}
store.close()

def test_latest_anomaly_per_table_issues_one_query(self):
store = _temp_store()
store.save_anomalies([
{
"table": "users",
"type": "volume_drop",
"severity": Severity.CRITICAL,
"message": "Row count dropped 60%",
},
{
"table": "orders",
"type": "schema_drift",
"severity": Severity.WARNING,
"message": "Column added: email",
},
])
seen: list[str] = []
store._conn.set_trace_callback(seen.append)
try:
store.get_latest_anomaly_per_table()
finally:
store._conn.set_trace_callback(None)
selects = [q for q in seen if q.lstrip().upper().startswith("SELECT")]
assert len(selects) == 1
store.close()
Loading