Skip to content

Commit 990762a

Browse files
authored
Merge pull request #3 from APIForge-Organisation/dev
feat: Phase 1 complete — DRIFT detection, unit tests, CI fix
2 parents 0e18ee2 + 648a3a3 commit 990762a

10 files changed

Lines changed: 709 additions & 7 deletions

File tree

.github/workflows/ci.yml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ name: CI
22

33
on:
44
push:
5-
branches: [dev]
5+
branches: [main, dev]
66
pull_request:
77
branches: [main, dev]
88

CHANGELOG.md

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,20 @@ Format: [Keep a Changelog](https://keepachangelog.com/en/1.1.0/) — versioning
66

77
---
88

9+
## [1.0.3] — 2026-05-15
10+
11+
### Added
12+
13+
- `DRIFT` insight type: detects progressive latency degradation using ordinary least squares over the last 30 days — emitted when slope ≥ 5ms/day over 7+ data points, with a 30-day projection
14+
- `DRIFT` filter chip added to the dashboard Insights view
15+
- 60 unit tests covering `aggregator`, `database`, `insights` and `middleware`
16+
17+
### Fixed
18+
19+
- CI badge was pointing to `main` with no workflow run — CI now triggers on both `main` and `dev`
20+
21+
---
22+
923
## [1.0.0] — 2026-05-15
1024

1125
### Fixed

apiforgepy/dashboard.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1250,8 +1250,9 @@
12501250
12511251
const types = [
12521252
{id:'ALL',label:'All'},{id:'PERF',label:'Performance'},
1253-
{id:'ANOMALY',label:'Anomaly'},{id:'DEAD',label:'Dead'},
1254-
{id:'UNTRACKED',label:'Untracked'},{id:'OK',label:'OK'},
1253+
{id:'DRIFT',label:'Drift'},{id:'ANOMALY',label:'Anomaly'},
1254+
{id:'DEAD',label:'Dead'},{id:'UNTRACKED',label:'Untracked'},
1255+
{id:'OK',label:'OK'},
12551256
];
12561257
const filtered = INSIGHTS.filter(i =>
12571258
(typeFilter === 'ALL' || i.type === typeFilter) &&

apiforgepy/database.py

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -254,6 +254,21 @@ def get_releases(self) -> list[dict]:
254254
""").fetchall()
255255
return [dict(r) for r in rows]
256256

257+
def get_drift_data(self) -> list[dict]:
258+
"""Returns one row per (route, method, day) over the last 30 days for drift detection."""
259+
since_30d = _now_sec() - 30 * 86_400
260+
rows = self._conn.execute("""
261+
SELECT
262+
route, method,
263+
CAST(bucket_ts / 86400 AS INTEGER) as day_bucket,
264+
AVG(lat_p90) as p90
265+
FROM api_metrics
266+
WHERE bucket_ts >= ? AND lat_p90 IS NOT NULL
267+
GROUP BY route, method, day_bucket
268+
ORDER BY route, method, day_bucket
269+
""", (since_30d,)).fetchall()
270+
return [dict(r) for r in rows]
271+
257272
def get_global_time_series(self, hours: int = 24) -> list[dict]:
258273
since = _now_sec() - hours * 3600
259274
rows = self._conn.execute("""

apiforgepy/insights.py

Lines changed: 67 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,11 @@
11
import math
22
import time
33

4-
DEAD_ENDPOINT_DAYS = 21
5-
REGRESSION_THRESHOLD = 0.20
6-
ANOMALY_Z_THRESHOLD = 2.5
4+
DEAD_ENDPOINT_DAYS = 21
5+
REGRESSION_THRESHOLD = 0.20
6+
ANOMALY_Z_THRESHOLD = 2.5
7+
DRIFT_SLOPE_THRESHOLD = 5.0 # ms/day above which progressive drift is reported
8+
DRIFT_MIN_DAYS = 7 # minimum number of daily data points required
79

810

911
def get_insights(db) -> list[dict]:
@@ -13,6 +15,7 @@ def get_insights(db) -> list[dict]:
1315
_detect_dead_endpoints,
1416
_detect_release_regressions,
1517
_detect_untracked_routes,
18+
_detect_drift,
1619
):
1720
try:
1821
insights.extend(fn(db))
@@ -160,6 +163,67 @@ def _detect_release_regressions(db) -> list[dict]:
160163
return insights
161164

162165

166+
def _detect_drift(db) -> list[dict]:
167+
rows = db.get_drift_data()
168+
if not rows:
169+
return []
170+
171+
# Group daily P90 samples by endpoint
172+
by_endpoint: dict[str, dict] = {}
173+
for row in rows:
174+
key = f"{row['method']}|{row['route']}"
175+
if key not in by_endpoint:
176+
by_endpoint[key] = {"method": row["method"], "route": row["route"], "points": []}
177+
by_endpoint[key]["points"].append({"x": row["day_bucket"], "y": row["p90"]})
178+
179+
insights = []
180+
for ep in by_endpoint.values():
181+
points = ep["points"]
182+
if len(points) < DRIFT_MIN_DAYS:
183+
continue
184+
185+
# Ordinary least squares on (day_index, p90) pairs
186+
x0 = points[0]["x"]
187+
xs = [p["x"] - x0 for p in points]
188+
ys = [p["y"] for p in points]
189+
n = len(xs)
190+
sum_x = sum(xs)
191+
sum_y = sum(ys)
192+
sum_xy = sum(xs[i] * ys[i] for i in range(n))
193+
sum_x2 = sum(x * x for x in xs)
194+
denom = n * sum_x2 - sum_x ** 2
195+
if denom == 0:
196+
continue
197+
198+
slope = (n * sum_xy - sum_x * sum_y) / denom
199+
if slope < DRIFT_SLOPE_THRESHOLD:
200+
continue
201+
202+
observed_days = xs[-1]
203+
projection_30 = round(slope * 30)
204+
method, route = ep["method"], ep["route"]
205+
day_str = "day" if observed_days == 1 else "days"
206+
207+
insights.append({
208+
"type": "DRIFT",
209+
"severity": "warning",
210+
"route": route,
211+
"method": method,
212+
"message": (
213+
f"`{method} {route}` has been progressively degrading for "
214+
f"{observed_days} {day_str}: +{slope:.1f}ms/day. "
215+
f"30-day projection: +{projection_30}ms."
216+
),
217+
"data": {
218+
"slope_ms_per_day": slope,
219+
"observed_days": observed_days,
220+
"projection_30d_ms": projection_30,
221+
},
222+
})
223+
224+
return insights
225+
226+
163227
def _detect_untracked_routes(db) -> list[dict]:
164228
return [
165229
{

pyproject.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
44

55
[project]
66
name = "apiforgepy"
7-
version = "1.0.2"
7+
version = "1.0.3"
88

99
description = "API observability & intelligence for FastAPI/Starlette — local-first, privacy-first"
1010
readme = "README.md"

tests/test_aggregator.py

Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
1+
import pytest
2+
from apiforgepy.aggregator import Aggregator
3+
4+
5+
def make_transport():
6+
class Spy:
7+
def __init__(self):
8+
self.calls = []
9+
def write(self, rows):
10+
self.calls.append(rows)
11+
return Spy()
12+
13+
14+
def base_event(**overrides):
15+
defaults = dict(
16+
route="/test", method="GET", status=200,
17+
duration_ms=100.0, timestamp="2026-01-01T00:00:00Z",
18+
env="test", release=None, service="svc", response_size=None,
19+
)
20+
return {**defaults, **overrides}
21+
22+
23+
class TestRecord:
24+
def test_accumulates_durations(self):
25+
t = make_transport()
26+
agg = Aggregator(t, flush_interval_ms=999_999_000)
27+
agg.record(base_event(duration_ms=10))
28+
agg.record(base_event(duration_ms=20))
29+
key = next(iter(agg._buffer))
30+
assert len(agg._buffer[key]["durations"]) == 2
31+
32+
def test_increments_2xx_counter(self):
33+
t = make_transport()
34+
agg = Aggregator(t, flush_interval_ms=999_999_000)
35+
agg.record(base_event(status=200))
36+
agg.record(base_event(status=204))
37+
key = next(iter(agg._buffer))
38+
assert agg._buffer[key]["status_2xx"] == 2
39+
40+
def test_increments_4xx_counter(self):
41+
t = make_transport()
42+
agg = Aggregator(t, flush_interval_ms=999_999_000)
43+
agg.record(base_event(status=404))
44+
agg.record(base_event(status=429))
45+
key = next(iter(agg._buffer))
46+
assert agg._buffer[key]["status_4xx"] == 2
47+
assert agg._buffer[key]["status_2xx"] == 0
48+
49+
def test_increments_5xx_counter(self):
50+
t = make_transport()
51+
agg = Aggregator(t, flush_interval_ms=999_999_000)
52+
agg.record(base_event(status=500))
53+
key = next(iter(agg._buffer))
54+
assert agg._buffer[key]["status_5xx"] == 1
55+
56+
def test_separate_buckets_per_route(self):
57+
t = make_transport()
58+
agg = Aggregator(t, flush_interval_ms=999_999_000)
59+
agg.record(base_event(route="/a"))
60+
agg.record(base_event(route="/b"))
61+
assert len(agg._buffer) == 2
62+
63+
def test_release_separates_bucket_key(self):
64+
t = make_transport()
65+
agg = Aggregator(t, flush_interval_ms=999_999_000)
66+
agg.record(base_event(release="v1"))
67+
agg.record(base_event(release="v2"))
68+
assert len(agg._buffer) == 2
69+
70+
71+
class TestFlush:
72+
def test_sends_rows_and_clears_buffer(self):
73+
t = make_transport()
74+
agg = Aggregator(t, flush_interval_ms=999_999_000)
75+
agg.record(base_event())
76+
agg._flush()
77+
assert len(t.calls) == 1
78+
assert len(t.calls[0]) == 1
79+
assert len(agg._buffer) == 0
80+
81+
def test_noop_when_buffer_empty(self):
82+
t = make_transport()
83+
agg = Aggregator(t, flush_interval_ms=999_999_000)
84+
agg._flush()
85+
assert len(t.calls) == 0
86+
87+
def test_computes_percentiles(self):
88+
t = make_transport()
89+
agg = Aggregator(t, flush_interval_ms=999_999_000)
90+
for i in range(1, 11):
91+
agg.record(base_event(duration_ms=float(i * 10)))
92+
agg._flush()
93+
row = t.calls[0][0]
94+
assert 50 <= row["lat_p50"] <= 60
95+
assert 90 <= row["lat_p90"] <= 100
96+
assert row["lat_p99"] >= 90
97+
98+
def test_correct_lat_min_max(self):
99+
t = make_transport()
100+
agg = Aggregator(t, flush_interval_ms=999_999_000)
101+
agg.record(base_event(duration_ms=5.0))
102+
agg.record(base_event(duration_ms=95.0))
103+
agg._flush()
104+
row = t.calls[0][0]
105+
assert row["lat_min"] == pytest.approx(5.0)
106+
assert row["lat_max"] == pytest.approx(95.0)
107+
108+
def test_total_calls_matches_records(self):
109+
t = make_transport()
110+
agg = Aggregator(t, flush_interval_ms=999_999_000)
111+
for _ in range(7):
112+
agg.record(base_event())
113+
agg._flush()
114+
assert t.calls[0][0]["total_calls"] == 7
115+
116+
def test_stop_flushes_buffer(self):
117+
t = make_transport()
118+
agg = Aggregator(t, flush_interval_ms=999_999_000)
119+
agg.record(base_event())
120+
agg.stop()
121+
assert len(t.calls) == 1, "stop() must flush remaining events"

0 commit comments

Comments
 (0)