diff --git a/.flow/tasks/fn-16.17.json b/.flow/tasks/fn-16.17.json new file mode 100644 index 000000000..459dec8de --- /dev/null +++ b/.flow/tasks/fn-16.17.json @@ -0,0 +1,22 @@ +{ + "assignee": "bordumbb@gmail.com", + "claim_note": "", + "claimed_at": "2026-09-27T17:28:56.373909Z", + "created_at": "2026-09-27T17:26:27.656485Z", + "depends_on": [], + "epic": "fn-16", + "evidence": { + "commits": [], + "prs": [], + "tests": [ + "pytest tests/integration/test_investigation_completed_at.py tests/integration/api/test_dashboard.py -m integration", + "pnpm test" + ] + }, + "id": "fn-16.17", + "priority": null, + "spec_path": ".flow/tasks/fn-16.17.md", + "status": "done", + "title": "Dashboard stats on the unified schema (completed_at + real numbers)", + "updated_at": "2026-09-27T17:58:48.568254Z" +} diff --git a/.flow/tasks/fn-16.17.md b/.flow/tasks/fn-16.17.md new file mode 100644 index 000000000..98a491d37 --- /dev/null +++ b/.flow/tasks/fn-16.17.md @@ -0,0 +1,22 @@ +# fn-16.17 Dashboard stats on the unified schema (completed_at + real numbers) + +## Description +TBD + +## Acceptance +## Acceptance +- Migration 036 adds investigations.completed_at, set by a trigger the first time outcome is set (re-runnable, no BEGIN/COMMIT) +- get_dashboard_stats: active = status active, completed today = completed_at today; dashboard recent list reads alert JSONB (dataset_ids[0], metric display_name/anomaly_type, severity) and COALESCE(outcome status) +- frontend fetchDashboardStats maps the snake_case API and no longer falls back to mock numbers; the page shows no fake numbers on error +- Integration tests on the migrated schema; vitest for the frontend + + +## Done summary +Migration 036 adds investigations.completed_at, stamped by a trigger the first time outcome is set (cleared if outcome is cleared; re-runnable; backfills created_at for existing completions). +get_dashboard_stats: active = status active, completed today = completed_at >= CURRENT_DATE. list_investigations now returns summaries from the alert JSONB (dataset_ids[0], display_name/anomaly_type, severity, COALESCE outcome status), "unknown" for non-AnomalyAlert replays. +Frontend: fetchDashboardStats used /dashboard/stats (404: missing /api/v1) and fell back to mock numbers; it now uses the generated client and maps snake_case. New DashboardStatsCards shows "—" and an alert on error. +Tests: tests/integration/test_investigation_completed_at.py, tests/integration/api/test_dashboard.py, src/lib/api/dashboard.test.ts, src/features/dashboard/dashboard-stats.test.tsx (all RED first). +## Evidence +- Commits: +- Tests: pytest tests/integration/test_investigation_completed_at.py tests/integration/api/test_dashboard.py -m integration, pnpm test +- PRs: diff --git a/.flow/tasks/fn-16.18.json b/.flow/tasks/fn-16.18.json new file mode 100644 index 000000000..80fd33c2e --- /dev/null +++ b/.flow/tasks/fn-16.18.json @@ -0,0 +1,21 @@ +{ + "assignee": "bordumbb@gmail.com", + "claim_note": "", + "claimed_at": "2026-09-27T17:58:48.835363Z", + "created_at": "2026-09-27T17:26:27.924541Z", + "depends_on": [], + "epic": "fn-16", + "evidence": { + "commits": [], + "prs": [], + "tests": [ + "pytest tests/integration/core/test_permission_service.py -m integration" + ] + }, + "id": "fn-16.18", + "priority": null, + "spec_path": ".flow/tasks/fn-16.18.md", + "status": "done", + "title": "RBAC datasource grants match alert datasource_id", + "updated_at": "2026-09-27T18:04:18.537780Z" +} diff --git a/.flow/tasks/fn-16.18.md b/.flow/tasks/fn-16.18.md new file mode 100644 index 000000000..35763cb55 --- /dev/null +++ b/.flow/tasks/fn-16.18.md @@ -0,0 +1,18 @@ +# fn-16.18 RBAC datasource grants match alert datasource_id + +## Description +TBD + +## Acceptance +## Acceptance +- PermissionService datasource grants (user and team) match investigations via alert->>datasource_id (4 places) +- Integration tests: can_access_investigation + get_accessible_investigation_ids + + +## Done summary +PermissionService matched datasource grants on the dropped investigations.data_source_id (4 places; the whole EXISTS/UNION query failed, so every access check raised). Now compares pg.data_source_id::text with i.alert->>datasource_id. +Tests: tests/integration/core/test_permission_service.py, user and team grants, can_access_investigation + get_accessible_investigation_ids (RED: column i.data_source_id does not exist). +## Evidence +- Commits: +- Tests: pytest tests/integration/core/test_permission_service.py -m integration +- PRs: diff --git a/.flow/tasks/fn-16.19.json b/.flow/tasks/fn-16.19.json new file mode 100644 index 000000000..90bc01cb7 --- /dev/null +++ b/.flow/tasks/fn-16.19.json @@ -0,0 +1,22 @@ +{ + "assignee": "bordumbb@gmail.com", + "claim_note": "", + "claimed_at": "2026-09-27T18:04:18.795162Z", + "created_at": "2026-09-27T17:26:28.178232Z", + "depends_on": [], + "epic": "fn-16", + "evidence": { + "commits": [], + "prs": [], + "tests": [ + "pytest tests/integration/test_fix_feedback_export_context.py -m integration", + "pytest tests/unit/services/test_feedback.py" + ] + }, + "id": "fn-16.19", + "priority": null, + "spec_path": ".flow/tasks/fn-16.19.md", + "status": "done", + "title": "Fix feedback export: issue link from issue_investigation_runs", + "updated_at": "2026-09-27T18:11:55.210616Z" +} diff --git a/.flow/tasks/fn-16.19.md b/.flow/tasks/fn-16.19.md new file mode 100644 index 000000000..42d26424f --- /dev/null +++ b/.flow/tasks/fn-16.19.md @@ -0,0 +1,19 @@ +# fn-16.19 Fix feedback export: issue link from issue_investigation_runs + +## Description +TBD + +## Acceptance +## Acceptance +- FixFeedbackService export context reads the issue link from issue_investigation_runs (not investigations.issue_id); no "None" issue ids +- Real-schema test of the context lookup (fix_feedback/fix_executions have no migrations; flagged) + + +## Done summary +FixFeedbackService export read issue_id from investigations (dropped column). The per-record lookup is now investigation_context(): created_at from investigations plus the latest issue_investigation_runs issue, omitted when none (previously would have been "None"). +Note: fix_feedback/fix_executions have no migrations and neither FixFeedbackService nor FixExecutionService is constructed anywhere, so the export itself cannot run; flagged. +Tests: tests/integration/test_fix_feedback_export_context.py (2 cases); unit test_feedback.py still passes. +## Evidence +- Commits: +- Tests: pytest tests/integration/test_fix_feedback_export_context.py -m integration, pytest tests/unit/services/test_feedback.py +- PRs: diff --git a/.flow/tasks/fn-16.20.json b/.flow/tasks/fn-16.20.json new file mode 100644 index 000000000..d2a09e108 --- /dev/null +++ b/.flow/tasks/fn-16.20.json @@ -0,0 +1,21 @@ +{ + "assignee": "bordumbb@gmail.com", + "claim_note": "", + "claimed_at": "2026-09-27T18:11:55.473365Z", + "created_at": "2026-09-27T17:26:28.429990Z", + "depends_on": [], + "epic": "fn-16", + "evidence": { + "commits": [], + "prs": [], + "tests": [ + "pytest dataing-ee/tests/integration/api/test_runbook_from_issue.py -m integration" + ] + }, + "id": "fn-16.20", + "priority": null, + "spec_path": ".flow/tasks/fn-16.20.md", + "status": "done", + "title": "EE runbook generation from the unified schema", + "updated_at": "2026-09-27T19:32:58.954009Z" +} diff --git a/.flow/tasks/fn-16.20.md b/.flow/tasks/fn-16.20.md new file mode 100644 index 000000000..95fa798c9 --- /dev/null +++ b/.flow/tasks/fn-16.20.md @@ -0,0 +1,20 @@ +# fn-16.20 EE runbook generation from the unified schema + +## Description +TBD + +## Acceptance +## Acceptance +- RunbookGenerator reads resolution_note, issue_labels, and the latest linked investigation outcome via issue_investigation_runs +- POST /runbooks/from-issue records created_from_investigation_id from the linked run +- EE integration test on the migrated schema + + +## Done summary +RunbookGenerator read issues.resolution/metadata and investigations.issue_id/synthesis/metadata (all dropped). It now reads resolution_note, labels from issue_labels, and the outcome of the latest issue_investigation_runs investigation, and returns that investigation_id; POST /runbooks/from-issue uses it (its own investigations.issue_id lookup is gone). +_row_to_response decodes JSONB text (AppDatabase has no JSONB codec on main), without which every runbook-returning route failed pydantic list validation. +Tests: dataing-ee/tests/integration/api/test_runbook_from_issue.py (latest run wins; issue without investigation). RED: column inv.issue_id does not exist. +## Evidence +- Commits: +- Tests: pytest dataing-ee/tests/integration/api/test_runbook_from_issue.py -m integration +- PRs: diff --git a/.flow/tasks/fn-16.21.json b/.flow/tasks/fn-16.21.json new file mode 100644 index 000000000..9e7979c0b --- /dev/null +++ b/.flow/tasks/fn-16.21.json @@ -0,0 +1,14 @@ +{ + "assignee": null, + "claim_note": "", + "claimed_at": null, + "created_at": "2026-09-27T17:26:28.678421Z", + "depends_on": [], + "epic": "fn-16", + "id": "fn-16.21", + "priority": null, + "spec_path": ".flow/tasks/fn-16.21.md", + "status": "todo", + "title": "Share the issue AnomalyAlert helper between CE and EE", + "updated_at": "2026-09-27T17:26:28.678692Z" +} diff --git a/.flow/tasks/fn-16.21.md b/.flow/tasks/fn-16.21.md new file mode 100644 index 000000000..a7b26f06d --- /dev/null +++ b/.flow/tasks/fn-16.21.md @@ -0,0 +1,18 @@ +# fn-16.21 Share the issue AnomalyAlert helper between CE and EE + +## Description +TBD + +## Acceptance +## Acceptance +- One CE helper builds the issue AnomalyAlert for both the CE spawn route (#159) and the EE spawn_investigation action (#183) +- Deferred until busy-cartwright (fn-59) lands: the issues.py import lines are next to its uncommitted edits + + +## Done summary +TBD + +## Evidence +- Commits: +- Tests: +- PRs: diff --git a/frontend/app/src/features/dashboard/dashboard-page.tsx b/frontend/app/src/features/dashboard/dashboard-page.tsx index 42d9433ee..0661519d8 100644 --- a/frontend/app/src/features/dashboard/dashboard-page.tsx +++ b/frontend/app/src/features/dashboard/dashboard-page.tsx @@ -1,17 +1,21 @@ import { useQuery } from "@tanstack/react-query"; import { Link } from "react-router-dom"; -import { Search, Database, CheckCircle2, Plus } from "lucide-react"; +import { Plus } from "lucide-react"; import { Card, CardContent, CardHeader, CardTitle } from "@/components/ui/Card"; import { Button } from "@/components/ui/Button"; -import { Skeleton } from "@/components/ui/skeleton"; import { fetchDashboardStats } from "@/lib/api/dashboard"; +import { DashboardStatsCards } from "./dashboard-stats"; import { RecentInvestigations } from "./recent-investigations"; import { PageHeader } from "@/components/shared/page-header"; import { useRole } from "@/lib/auth"; export function DashboardPage() { - const { data: stats, isLoading } = useQuery({ + const { + data: stats, + isLoading, + isError, + } = useQuery({ queryKey: ["dashboard-stats"], queryFn: fetchDashboardStats, }); @@ -33,60 +37,11 @@ export function DashboardPage() { } /> - {/* Stats Grid */} -
- - - - Active Investigations - - - - - {isLoading ? ( - - ) : ( -
- {stats?.activeInvestigations ?? 0} -
- )} -
-
- - - - - Completed Today - - - - - {isLoading ? ( - - ) : ( -
- {stats?.completedToday ?? 0} -
- )} -
-
- - - - Data Sources - - - - {isLoading ? ( - - ) : ( -
- {stats?.dataSources ?? 0} -
- )} -
-
-
+ {/* Recent Investigations */} diff --git a/frontend/app/src/features/dashboard/dashboard-stats.test.tsx b/frontend/app/src/features/dashboard/dashboard-stats.test.tsx new file mode 100644 index 000000000..3ba5a8154 --- /dev/null +++ b/frontend/app/src/features/dashboard/dashboard-stats.test.tsx @@ -0,0 +1,43 @@ +import { render, screen } from "@testing-library/react"; +import { describe, expect, it } from "vitest"; + +import { DashboardStatsCards } from "./dashboard-stats"; + +describe("DashboardStatsCards", () => { + it("shows each stat", () => { + render( + , + ); + + expect(screen.getByLabelText("Active Investigations")).toHaveTextContent( + "4", + ); + expect(screen.getByLabelText("Completed Today")).toHaveTextContent("1"); + expect(screen.getByLabelText("Data Sources")).toHaveTextContent("2"); + }); + + it("shows no numbers when the stats failed to load", () => { + render( + , + ); + + for (const title of [ + "Active Investigations", + "Completed Today", + "Data Sources", + ]) { + expect(screen.getByLabelText(title)).toHaveTextContent("—"); + } + expect(screen.getByRole("alert")).toHaveTextContent( + "Couldn't load dashboard stats.", + ); + }); +}); diff --git a/frontend/app/src/features/dashboard/dashboard-stats.tsx b/frontend/app/src/features/dashboard/dashboard-stats.tsx new file mode 100644 index 000000000..9a3177e79 --- /dev/null +++ b/frontend/app/src/features/dashboard/dashboard-stats.tsx @@ -0,0 +1,52 @@ +import { CheckCircle2, Database, Search } from "lucide-react"; + +import { Card, CardContent, CardHeader, CardTitle } from "@/components/ui/Card"; +import { Skeleton } from "@/components/ui/skeleton"; +import type { DashboardStats } from "@/lib/api/dashboard"; + +const STAT_CARDS = [ + { key: "activeInvestigations", title: "Active Investigations", Icon: Search }, + { key: "completedToday", title: "Completed Today", Icon: CheckCircle2 }, + { key: "dataSources", title: "Data Sources", Icon: Database }, +] as const; + +interface DashboardStatsCardsProps { + stats: DashboardStats | undefined; + isLoading: boolean; + isError: boolean; +} + +export function DashboardStatsCards({ + stats, + isLoading, + isError, +}: DashboardStatsCardsProps) { + return ( +
+
+ {STAT_CARDS.map(({ key, title, Icon }) => ( + + + {title} + + + + {isLoading ? ( + + ) : ( +
+ {stats ? stats[key] : "—"} +
+ )} +
+
+ ))} +
+ {isError && ( +

+ Couldn't load dashboard stats. +

+ )} +
+ ); +} diff --git a/frontend/app/src/lib/api/dashboard.test.ts b/frontend/app/src/lib/api/dashboard.test.ts new file mode 100644 index 000000000..1fdc12bee --- /dev/null +++ b/frontend/app/src/lib/api/dashboard.test.ts @@ -0,0 +1,43 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; + +import { fetchDashboardStats } from "./dashboard"; + +function respondWith(status: number, body: unknown) { + const fetchMock = vi.fn(async () => ({ + ok: status < 400, + status, + json: async () => body, + })); + vi.stubGlobal("fetch", fetchMock); + return fetchMock; +} + +afterEach(() => { + vi.unstubAllGlobals(); +}); + +describe("fetchDashboardStats", () => { + it("reads the stats from GET /api/v1/dashboard/stats", async () => { + const fetchMock = respondWith(200, { + active_investigations: 4, + completed_today: 1, + data_sources: 2, + }); + + await expect(fetchDashboardStats()).resolves.toEqual({ + activeInvestigations: 4, + completedToday: 1, + dataSources: 2, + }); + expect(fetchMock).toHaveBeenCalledWith( + expect.stringMatching(/\/api\/v1\/dashboard\/stats$/), + expect.objectContaining({ method: "GET" }), + ); + }); + + it("fails instead of inventing numbers when the API errors", async () => { + respondWith(500, { detail: "database unavailable" }); + + await expect(fetchDashboardStats()).rejects.toThrow("database unavailable"); + }); +}); diff --git a/frontend/app/src/lib/api/dashboard.ts b/frontend/app/src/lib/api/dashboard.ts index 0302eb445..4327dd0bb 100644 --- a/frontend/app/src/lib/api/dashboard.ts +++ b/frontend/app/src/lib/api/dashboard.ts @@ -1,4 +1,4 @@ -import customInstance from "./client"; +import { getStatsApiV1DashboardStatsGet } from "./generated/dashboard/dashboard"; export interface DashboardStats { activeInvestigations: number; @@ -7,17 +7,10 @@ export interface DashboardStats { } export async function fetchDashboardStats(): Promise { - try { - return await customInstance({ - url: "/dashboard/stats", - method: "GET", - }); - } catch { - // Return mock data if endpoint doesn't exist - return { - activeInvestigations: 3, - completedToday: 7, - dataSources: 2, - }; - } + const stats = await getStatsApiV1DashboardStatsGet(); + return { + activeInvestigations: stats.active_investigations, + completedToday: stats.completed_today, + dataSources: stats.data_sources, + }; } diff --git a/python-packages/dataing-ee/src/dataing_ee/core/runbook/generator.py b/python-packages/dataing-ee/src/dataing_ee/core/runbook/generator.py index 8d68939e7..60ad7f1e4 100644 --- a/python-packages/dataing-ee/src/dataing_ee/core/runbook/generator.py +++ b/python-packages/dataing-ee/src/dataing_ee/core/runbook/generator.py @@ -2,6 +2,7 @@ from __future__ import annotations +import json import logging from dataclasses import dataclass from typing import Any @@ -26,6 +27,15 @@ class GeneratedRunbook: prevention_notes: str | None dataset_id: str | None labels: list[str] + investigation_id: UUID | None = None + + +def _json_object(value: Any) -> dict[str, Any]: + """A JSONB object column value as a dict (AppDatabase returns JSONB as text).""" + if isinstance(value, str): + value = json.loads(value) + result: dict[str, Any] = value or {} + return result class RunbookGenerator: @@ -47,25 +57,28 @@ async def generate_from_issue( Returns: GeneratedRunbook or None if issue not suitable """ - # Fetch issue with investigation + # Fetch the issue with the outcome of its most recently spawned investigation issue = await db.fetch_one( """ SELECT - i.id, i.title, i.description, - i.status, - i.resolution, + i.resolution_note, i.dataset_id, - i.metadata, - inv.id as investigation_id, - inv.synthesis, - inv.metadata as inv_metadata + ARRAY( + SELECT label FROM issue_labels WHERE issue_id = i.id ORDER BY label + ) AS labels, + run.investigation_id, + inv.outcome FROM issues i - LEFT JOIN investigations inv ON inv.issue_id = i.id + LEFT JOIN LATERAL ( + SELECT investigation_id FROM issue_investigation_runs + WHERE issue_id = i.id + ORDER BY created_at DESC + LIMIT 1 + ) run ON true + LEFT JOIN investigations inv ON inv.id = run.investigation_id WHERE i.id = $1 AND i.tenant_id = $2 - ORDER BY inv.created_at DESC - LIMIT 1 """, issue_id, tenant_id, @@ -78,10 +91,10 @@ async def generate_from_issue( # Extract data title = issue.get("title", "Untitled Issue") description = issue.get("description", "") - resolution = issue.get("resolution", "") - synthesis = issue.get("synthesis") or {} + resolution = issue.get("resolution_note") or "" + synthesis = _json_object(issue.get("outcome")) dataset_id = issue.get("dataset_id") - metadata = issue.get("metadata") or {} + metadata = {"labels": issue.get("labels") or []} # Build runbook content symptoms = self._extract_symptoms(description, synthesis) @@ -117,6 +130,7 @@ async def generate_from_issue( prevention_notes=prevention_notes, dataset_id=dataset_id, labels=labels, + investigation_id=issue.get("investigation_id"), ) def _extract_symptoms( diff --git a/python-packages/dataing-ee/src/dataing_ee/entrypoints/api/routes/runbooks.py b/python-packages/dataing-ee/src/dataing_ee/entrypoints/api/routes/runbooks.py index 9981efe56..d70f7f040 100644 --- a/python-packages/dataing-ee/src/dataing_ee/entrypoints/api/routes/runbooks.py +++ b/python-packages/dataing-ee/src/dataing_ee/entrypoints/api/routes/runbooks.py @@ -2,6 +2,7 @@ from __future__ import annotations +import json import logging from datetime import datetime from typing import Annotated, Any @@ -488,13 +489,6 @@ async def generate_runbook_from_issue( detail="Could not generate runbook from issue", ) - # Get investigation ID if exists - inv_row = await db.fetch_one( - "SELECT id FROM investigations WHERE issue_id = $1 ORDER BY created_at DESC LIMIT 1", - issue_id, - ) - investigation_id = inv_row["id"] if inv_row else None - # Insert runbook row = await db.fetch_one( """ @@ -521,7 +515,7 @@ async def generate_runbook_from_issue( generated.prevention_notes, body.publish, issue_id, - investigation_id, + generated.investigation_id, auth.user_id, ) @@ -725,6 +719,14 @@ def _usefulness_score(helpful: int, not_helpful: int) -> float: return max(0.0, 0.1 * helpful - 0.05 * not_helpful) +def _json_list(value: Any) -> list[dict[str, Any]]: + """A JSONB array column value as a list (AppDatabase returns JSONB as text).""" + if isinstance(value, str): + value = json.loads(value) + result: list[dict[str, Any]] = value or [] + return result + + def _row_to_response(row: dict[str, Any]) -> RunbookResponse: """Convert database row to response model.""" return RunbookResponse( @@ -735,10 +737,10 @@ def _row_to_response(row: dict[str, Any]) -> RunbookResponse: summary=row.get("summary"), dataset_id=row.get("dataset_id"), labels=row.get("labels") or [], - symptoms=row.get("symptoms") or [], + symptoms=_json_list(row.get("symptoms")), root_cause=row.get("root_cause"), - verification_steps=row.get("verification_steps") or [], - fix_steps=row.get("fix_steps") or [], + verification_steps=_json_list(row.get("verification_steps")), + fix_steps=_json_list(row.get("fix_steps")), prevention_notes=row.get("prevention_notes"), is_published=row["is_published"], view_count=row["view_count"], diff --git a/python-packages/dataing-ee/tests/integration/api/test_runbook_from_issue.py b/python-packages/dataing-ee/tests/integration/api/test_runbook_from_issue.py new file mode 100644 index 000000000..71ddab248 --- /dev/null +++ b/python-packages/dataing-ee/tests/integration/api/test_runbook_from_issue.py @@ -0,0 +1,157 @@ +"""Integration tests for POST /runbooks/from-issue/{issue_id} on the migrated schema.""" + +from __future__ import annotations + +import json +from datetime import UTC, datetime, timedelta +from typing import Any +from uuid import UUID, uuid4 + +import httpx +import pytest +from dataing_ee.entrypoints.api.routes.runbooks import router +from fastapi import FastAPI + +from dataing.adapters.db.app_db import AppDatabase +from dataing.entrypoints.api.deps import get_app_db +from dataing.entrypoints.api.middleware.auth import ApiKeyContext, verify_api_key + +pytestmark = pytest.mark.integration + +RESOLUTION_NOTE = "Re-ran the orders load after the upstream export landed" +OUTCOME = { + "status": "completed", + "root_cause": "The 02:00 orders load ran before the upstream export finished", + "recommendations": [ + { + "description": "Gate the orders load on the upstream export sensor", + "type": "preventive", + } + ], + "tags": ["late-data"], +} + + +async def _create_tenant(db: AppDatabase) -> UUID: + tenant_id = uuid4() + await db.execute( + "INSERT INTO tenants (id, name, slug) VALUES ($1, $2, $3)", + tenant_id, + "Test Tenant", + f"test-{tenant_id.hex[:12]}", + ) + return tenant_id + + +async def _create_resolved_issue(db: AppDatabase, tenant_id: UUID) -> UUID: + issue_id = uuid4() + await db.execute( + """INSERT INTO issues + (id, tenant_id, number, title, description, status, dataset_id, resolution_note) + VALUES ($1, $2, 1, $3, $4, 'resolved', 'public.orders', $5)""", + issue_id, + tenant_id, + "Null spike in orders.customer_id", + "customer_id is null on a quarter of today's orders", + RESOLUTION_NOTE, + ) + for label in ("orders", "null-rate"): + await db.execute( + "INSERT INTO issue_labels (issue_id, label) VALUES ($1, $2)", issue_id, label + ) + return issue_id + + +async def _investigate( + db: AppDatabase, + tenant_id: UUID, + issue_id: UUID, + outcome: dict[str, Any], + *, + spawned_at: datetime, +) -> UUID: + """Store an investigation spawned from the issue, with its final outcome.""" + investigation_id = uuid4() + await db.execute( + "INSERT INTO investigations (id, tenant_id, alert, outcome) VALUES ($1, $2, $3, $4)", + investigation_id, + tenant_id, + json.dumps({"dataset_ids": ["public.orders"]}), + json.dumps(outcome), + ) + await db.execute( + """INSERT INTO issue_investigation_runs + (issue_id, investigation_id, trigger_type, created_at) + VALUES ($1, $2, 'human', $3)""", + issue_id, + investigation_id, + spawned_at, + ) + return investigation_id + + +async def _generate(db: AppDatabase, tenant_id: UUID, issue_id: UUID) -> dict[str, Any]: + app = FastAPI() + app.include_router(router, prefix="/api/v1") + app.dependency_overrides[get_app_db] = lambda: db + app.dependency_overrides[verify_api_key] = lambda: ApiKeyContext( + key_id=uuid4(), + tenant_id=tenant_id, + tenant_slug="test", + tenant_name="Test Tenant", + user_id=None, + scopes=["read", "write"], + ) + transport = httpx.ASGITransport(app=app) + async with httpx.AsyncClient(transport=transport, base_url="http://test") as client: + response = await client.post( + f"/api/v1/runbooks/from-issue/{issue_id}", json={"publish": False} + ) + assert response.status_code == 201, response.text + runbook: dict[str, Any] = response.json() + return runbook + + +async def test_runbook_is_built_from_the_issue_and_its_latest_investigation( + migrated_db: AppDatabase, +) -> None: + """The issue's resolution and labels, and the latest run's outcome, fill the runbook.""" + tenant_id = await _create_tenant(migrated_db) + issue_id = await _create_resolved_issue(migrated_db, tenant_id) + now = datetime.now(UTC) + await _investigate( + migrated_db, + tenant_id, + issue_id, + {"status": "failed", "root_cause": "inconclusive"}, + spawned_at=now - timedelta(hours=1), + ) + investigation_id = await _investigate(migrated_db, tenant_id, issue_id, OUTCOME, spawned_at=now) + + runbook = await _generate(migrated_db, tenant_id, issue_id) + + assert runbook["created_from_issue_id"] == str(issue_id) + assert runbook["created_from_investigation_id"] == str(investigation_id) + assert runbook["dataset_id"] == "public.orders" + assert runbook["root_cause"] == OUTCOME["root_cause"] + assert [step["description"] for step in runbook["fix_steps"]] == [ + RESOLUTION_NOTE, + "Gate the orders load on the upstream export sensor", + ] + assert runbook["prevention_notes"] == "Gate the orders load on the upstream export sensor" + assert sorted(runbook["labels"]) == ["late-data", "null-rate", "orders"] + + +async def test_runbook_without_an_investigation_uses_the_issue_alone( + migrated_db: AppDatabase, +) -> None: + """An issue resolved without an investigation still yields a runbook.""" + tenant_id = await _create_tenant(migrated_db) + issue_id = await _create_resolved_issue(migrated_db, tenant_id) + + runbook = await _generate(migrated_db, tenant_id, issue_id) + + assert runbook["created_from_investigation_id"] is None + assert runbook["root_cause"] is None + assert [step["description"] for step in runbook["fix_steps"]] == [RESOLUTION_NOTE] + assert sorted(runbook["labels"]) == ["null-rate", "orders"] diff --git a/python-packages/dataing/migrations/036_investigation_completed_at.sql b/python-packages/dataing/migrations/036_investigation_completed_at.sql new file mode 100644 index 000000000..798aff181 --- /dev/null +++ b/python-packages/dataing/migrations/036_investigation_completed_at.sql @@ -0,0 +1,32 @@ +-- Record when an investigation completes. +-- investigations.status is generated from outcome, but nothing records when the +-- outcome was set. A trigger stamps completed_at the first time outcome becomes +-- non-null (and clears it if the outcome is cleared), so writers only set outcome. +-- Re-runnable: every statement is idempotent. + +ALTER TABLE investigations ADD COLUMN IF NOT EXISTS completed_at TIMESTAMPTZ; + +COMMENT ON COLUMN investigations.completed_at IS + 'When outcome was first set (NULL while the investigation is active)'; + +CREATE OR REPLACE FUNCTION set_investigation_completed_at() RETURNS trigger AS $$ +BEGIN + IF NEW.outcome IS NULL THEN + NEW.completed_at := NULL; + ELSIF NEW.completed_at IS NULL THEN + NEW.completed_at := NOW(); + END IF; + RETURN NEW; +END; +$$ LANGUAGE plpgsql; + +DROP TRIGGER IF EXISTS investigations_set_completed_at ON investigations; +CREATE TRIGGER investigations_set_completed_at + BEFORE INSERT OR UPDATE OF outcome ON investigations + FOR EACH ROW EXECUTE FUNCTION set_investigation_completed_at(); + +-- Investigations completed before this migration: the completion time is unknown, +-- so use the creation time. +UPDATE investigations +SET completed_at = created_at +WHERE outcome IS NOT NULL AND completed_at IS NULL; diff --git a/python-packages/dataing/src/dataing/adapters/db/app_db.py b/python-packages/dataing/src/dataing/adapters/db/app_db.py index a469a0846..a1b727fcc 100644 --- a/python-packages/dataing/src/dataing/adapters/db/app_db.py +++ b/python-packages/dataing/src/dataing/adapters/db/app_db.py @@ -572,24 +572,24 @@ async def get_investigation( async def list_investigations( self, tenant_id: UUID, - status: str | None = None, limit: int = 50, offset: int = 0, ) -> list[dict[str, Any]]: - """List investigations for a tenant.""" - if status: - return await self.fetch_all( - """SELECT * FROM investigations - WHERE tenant_id = $1 AND status = $2 - ORDER BY created_at DESC - LIMIT $3 OFFSET $4""", - tenant_id, - status, - limit, - offset, - ) + """List investigation summaries for a tenant, newest first. + + Summaries come from the alert JSONB: the primary dataset, the metric's + display name (else the anomaly type) and the severity. Alerts that are not + an AnomalyAlert, such as imported replays, report "unknown". + """ return await self.fetch_all( - """SELECT * FROM investigations + """SELECT id, + COALESCE(alert->'dataset_ids'->>0, 'unknown') AS dataset_id, + COALESCE(NULLIF(alert #>> '{metric_spec,display_name}', ''), + alert->>'anomaly_type', 'unknown') AS metric_name, + COALESCE(outcome->>'status', status) AS status, + alert->>'severity' AS severity, + created_at + FROM investigations WHERE tenant_id = $1 ORDER BY created_at DESC LIMIT $2 OFFSET $3""", @@ -793,18 +793,17 @@ async def get_monthly_usage( # Dashboard stats async def get_dashboard_stats(self, tenant_id: UUID) -> dict[str, Any]: """Get dashboard statistics for a tenant.""" - # Active investigations + # Active investigations (no outcome yet) active_result = await self.fetch_one( """SELECT COUNT(*) as count FROM investigations - WHERE tenant_id = $1 AND status IN ('pending', 'in_progress')""", + WHERE tenant_id = $1 AND status = 'active'""", tenant_id, ) - # Completed today + # Completed today (outcome set today, whatever it says) completed_result = await self.fetch_one( """SELECT COUNT(*) as count FROM investigations - WHERE tenant_id = $1 AND status = 'completed' - AND completed_at >= CURRENT_DATE""", + WHERE tenant_id = $1 AND completed_at >= CURRENT_DATE""", tenant_id, ) diff --git a/python-packages/dataing/src/dataing/core/rbac/permission_service.py b/python-packages/dataing/src/dataing/core/rbac/permission_service.py index af241fece..1c85a4191 100644 --- a/python-packages/dataing/src/dataing/core/rbac/permission_service.py +++ b/python-packages/dataing/src/dataing/core/rbac/permission_service.py @@ -77,7 +77,7 @@ async def can_access_investigation(self, user_id: UUID, investigation_id: UUID) -- Datasource-based grant (user) SELECT 1 FROM permission_grants pg - JOIN investigations i ON pg.data_source_id = i.data_source_id + JOIN investigations i ON pg.data_source_id::text = i.alert->>'datasource_id' WHERE pg.user_id = $1 AND i.id = $2 UNION ALL @@ -102,7 +102,7 @@ async def can_access_investigation(self, user_id: UUID, investigation_id: UUID) -- Team grants (datasource-based) SELECT 1 FROM permission_grants pg JOIN team_members tm ON pg.team_id = tm.team_id - JOIN investigations i ON pg.data_source_id = i.data_source_id + JOIN investigations i ON pg.data_source_id::text = i.alert->>'datasource_id' WHERE tm.user_id = $1 AND i.id = $2 ) """, @@ -158,7 +158,7 @@ async def get_accessible_investigation_ids( -- Datasource grant (user) OR EXISTS ( SELECT 1 FROM permission_grants pg - WHERE pg.user_id = $1 AND pg.data_source_id = i.data_source_id + WHERE pg.user_id = $1 AND pg.data_source_id::text = i.alert->>'datasource_id' ) -- Team grants @@ -172,7 +172,7 @@ async def get_accessible_investigation_ids( SELECT tag_id FROM investigation_tags WHERE investigation_id = i.id ) - OR pg.data_source_id = i.data_source_id + OR pg.data_source_id::text = i.alert->>'datasource_id' ) ) ) diff --git a/python-packages/dataing/src/dataing/entrypoints/api/routes/dashboard.py b/python-packages/dataing/src/dataing/entrypoints/api/routes/dashboard.py index e4e781579..235a9c7c6 100644 --- a/python-packages/dataing/src/dataing/entrypoints/api/routes/dashboard.py +++ b/python-packages/dataing/src/dataing/entrypoints/api/routes/dashboard.py @@ -69,7 +69,7 @@ async def get_dashboard( dataset_id=inv["dataset_id"], metric_name=inv["metric_name"], status=inv["status"], - severity=inv.get("severity"), + severity=inv["severity"], created_at=inv["created_at"], ) for inv in recent diff --git a/python-packages/dataing/src/dataing/services/feedback.py b/python-packages/dataing/src/dataing/services/feedback.py index 7692d3791..f35c14f63 100644 --- a/python-packages/dataing/src/dataing/services/feedback.py +++ b/python-packages/dataing/src/dataing/services/feedback.py @@ -518,22 +518,7 @@ async def export_feedback_for_finetuning( context: dict[str, Any] = {} if include_context: - # Get investigation context - inv_row = await self.db.fetch_one( - """ - SELECT - issue_id, - created_at as investigation_created_at - FROM investigations - WHERE id = $1 - """, - row["investigation_id"], - ) - if inv_row: - context["issue_id"] = str(inv_row["issue_id"]) - context["investigation_created_at"] = inv_row[ - "investigation_created_at" - ].isoformat() + context = await self.investigation_context(row["investigation_id"]) records.append( FeedbackExportRecord( @@ -557,6 +542,44 @@ async def export_feedback_for_finetuning( return records + async def investigation_context(self, investigation_id: UUID) -> dict[str, Any]: + """Context an exported feedback record carries about its investigation. + + Names the issue the investigation was most recently spawned from, if any. + + Args: + investigation_id: The investigation the feedback is about. + + Returns: + investigation_created_at, plus issue_id when an issue spawned it. + """ + if not self.db: + return {} + + row = await self.db.fetch_one( + """ + SELECT i.created_at AS investigation_created_at, run.issue_id + FROM investigations i + LEFT JOIN LATERAL ( + SELECT issue_id FROM issue_investigation_runs + WHERE investigation_id = i.id + ORDER BY created_at DESC + LIMIT 1 + ) run ON true + WHERE i.id = $1 + """, + investigation_id, + ) + if row is None: + return {} + + context: dict[str, Any] = { + "investigation_created_at": row["investigation_created_at"].isoformat() + } + if row["issue_id"] is not None: + context["issue_id"] = str(row["issue_id"]) + return context + def get_feedback(self, feedback_id: UUID) -> FixFeedback | None: """Get feedback by ID from cache. diff --git a/python-packages/dataing/tests/integration/api/test_dashboard.py b/python-packages/dataing/tests/integration/api/test_dashboard.py new file mode 100644 index 000000000..74cf63480 --- /dev/null +++ b/python-packages/dataing/tests/integration/api/test_dashboard.py @@ -0,0 +1,216 @@ +"""Integration tests for the /dashboard routes on the migrated schema.""" + +from __future__ import annotations + +import json +from datetime import UTC, datetime, timedelta +from typing import Any +from uuid import UUID, uuid4 + +import httpx +import pytest +from fastapi import FastAPI + +from dataing.adapters.db.app_db import AppDatabase +from dataing.core.domain_types import AnomalyAlert, MetricSpec +from dataing.entrypoints.api.deps import get_app_db +from dataing.entrypoints.api.middleware.auth import ApiKeyContext, verify_api_key +from dataing.entrypoints.api.routes.dashboard import router + +pytestmark = pytest.mark.integration + +CREATED_AT = datetime(2026, 9, 1, 12, 0, tzinfo=UTC) + + +async def _create_tenant(db: AppDatabase) -> UUID: + tenant_id = uuid4() + await db.execute( + "INSERT INTO tenants (id, name, slug) VALUES ($1, $2, $3)", + tenant_id, + "Test Tenant", + f"test-{tenant_id.hex[:12]}", + ) + return tenant_id + + +async def _create_datasource(db: AppDatabase, tenant_id: UUID, *, is_active: bool = True) -> UUID: + datasource_id = uuid4() + await db.execute( + """INSERT INTO data_sources + (id, tenant_id, name, type, connection_config_encrypted, is_active) + VALUES ($1, $2, $3, 'postgresql', 'unused', $4)""", + datasource_id, + tenant_id, + f"warehouse-{datasource_id.hex[:8]}", + is_active, + ) + return datasource_id + + +def _alert( + datasource_id: UUID, + *, + dataset_ids: list[str] | None = None, + display_name: str = "null_rate on customer_id", + severity: str = "high", +) -> dict[str, Any]: + """An alert as POST /investigations stores it.""" + alert = AnomalyAlert( + dataset_ids=dataset_ids or ["public.orders"], + metric_spec=MetricSpec( + metric_type="column", + expression="customer_id", + display_name=display_name, + columns_referenced=["customer_id"], + ), + anomaly_type="null_rate", + expected_value=0.01, + actual_value=0.25, + deviation_pct=2400.0, + anomaly_date="2026-09-01", + severity=severity, + ) + return {**alert.model_dump(mode="json"), "datasource_id": str(datasource_id)} + + +async def _create_investigation( + db: AppDatabase, + tenant_id: UUID, + alert: dict[str, Any], + *, + created_at: datetime = CREATED_AT, +) -> UUID: + investigation_id = uuid4() + await db.execute( + "INSERT INTO investigations (id, tenant_id, alert, created_at) VALUES ($1, $2, $3, $4)", + investigation_id, + tenant_id, + json.dumps(alert), + created_at, + ) + return investigation_id + + +async def _complete( + db: AppDatabase, investigation_id: UUID, *, status: str = "completed", days_ago: int = 0 +) -> None: + """Record an outcome, as if the investigation finished `days_ago` days ago.""" + await db.execute( + "UPDATE investigations SET outcome = $2 WHERE id = $1", + investigation_id, + json.dumps({"status": status}), + ) + if days_ago: + await db.execute( + "UPDATE investigations SET completed_at = NOW() - make_interval(days => $2) " + "WHERE id = $1", + investigation_id, + days_ago, + ) + + +async def _get(db: AppDatabase, tenant_id: UUID, path: str) -> dict[str, Any]: + app = FastAPI() + app.include_router(router, prefix="/api/v1") + app.dependency_overrides[get_app_db] = lambda: db + app.dependency_overrides[verify_api_key] = lambda: ApiKeyContext( + key_id=uuid4(), + tenant_id=tenant_id, + tenant_slug="test", + tenant_name="Test Tenant", + user_id=None, + scopes=["read"], + ) + transport = httpx.ASGITransport(app=app) + async with httpx.AsyncClient(transport=transport, base_url="http://test") as client: + response = await client.get(f"/api/v1/dashboard{path}") + assert response.status_code == 200, response.text + body: dict[str, Any] = response.json() + return body + + +async def test_stats_count_the_tenants_active_and_completed_today( + migrated_db: AppDatabase, +) -> None: + """Active means no outcome yet; completed today means the outcome was set today.""" + tenant_id = await _create_tenant(migrated_db) + datasource_id = await _create_datasource(migrated_db, tenant_id) + await _create_datasource(migrated_db, tenant_id) + await _create_datasource(migrated_db, tenant_id, is_active=False) + alert = _alert(datasource_id) + for _ in range(2): + await _create_investigation(migrated_db, tenant_id, alert) + await _complete(migrated_db, await _create_investigation(migrated_db, tenant_id, alert)) + await _complete( + migrated_db, + await _create_investigation(migrated_db, tenant_id, alert), + status="failed", + ) + await _complete( + migrated_db, await _create_investigation(migrated_db, tenant_id, alert), days_ago=2 + ) + other_tenant_id = await _create_tenant(migrated_db) + await _create_investigation(migrated_db, other_tenant_id, alert) + await _complete(migrated_db, await _create_investigation(migrated_db, other_tenant_id, alert)) + + stats = await _get(migrated_db, tenant_id, "/stats") + + assert stats == {"active_investigations": 2, "completed_today": 2, "data_sources": 2} + + +async def test_dashboard_lists_recent_investigations_from_their_alerts( + migrated_db: AppDatabase, +) -> None: + """Recent investigations are summarised from the alert JSONB, newest first.""" + tenant_id = await _create_tenant(migrated_db) + datasource_id = await _create_datasource(migrated_db, tenant_id) + active = await _create_investigation( + migrated_db, + tenant_id, + _alert(datasource_id, dataset_ids=["public.orders", "public.customers"]), + ) + failed = await _create_investigation( + migrated_db, + tenant_id, + _alert(datasource_id, display_name="", severity="low"), + created_at=CREATED_AT + timedelta(hours=1), + ) + await _complete(migrated_db, failed, status="failed") + # POST /investigations/import stores replays with an alert that is not an AnomalyAlert + replay = await _create_investigation( + migrated_db, + tenant_id, + {"replay_of": str(uuid4()), "is_replay": True}, + created_at=CREATED_AT + timedelta(hours=2), + ) + await _complete(migrated_db, replay) + + body = await _get(migrated_db, tenant_id, "/") + + assert body["recent_investigations"] == [ + { + "id": str(replay), + "dataset_id": "unknown", + "metric_name": "unknown", + "status": "completed", + "severity": None, + "created_at": (CREATED_AT + timedelta(hours=2)).isoformat().replace("+00:00", "Z"), + }, + { + "id": str(failed), + "dataset_id": "public.orders", + "metric_name": "null_rate", + "status": "failed", + "severity": "low", + "created_at": (CREATED_AT + timedelta(hours=1)).isoformat().replace("+00:00", "Z"), + }, + { + "id": str(active), + "dataset_id": "public.orders", + "metric_name": "null_rate on customer_id", + "status": "active", + "severity": "high", + "created_at": CREATED_AT.isoformat().replace("+00:00", "Z"), + }, + ] + assert body["stats"] == {"active_investigations": 1, "completed_today": 2, "data_sources": 1} diff --git a/python-packages/dataing/tests/integration/core/test_permission_service.py b/python-packages/dataing/tests/integration/core/test_permission_service.py new file mode 100644 index 000000000..94fae44b8 --- /dev/null +++ b/python-packages/dataing/tests/integration/core/test_permission_service.py @@ -0,0 +1,131 @@ +"""Integration tests for PermissionService datasource grants on the migrated schema.""" + +from __future__ import annotations + +import json +from uuid import UUID, uuid4 + +import pytest + +from dataing.adapters.db.app_db import AppDatabase +from dataing.core.domain_types import AnomalyAlert, MetricSpec +from dataing.core.rbac.permission_service import PermissionService + +pytestmark = pytest.mark.integration + + +async def _create_org(db: AppDatabase) -> UUID: + """An org is a tenant and an organization with the same id (the JWT org_id).""" + org_id = uuid4() + slug = f"test-{org_id.hex[:12]}" + await db.execute( + "INSERT INTO tenants (id, name, slug) VALUES ($1, 'Test Org', $2)", org_id, slug + ) + await db.execute( + "INSERT INTO organizations (id, name, slug) VALUES ($1, 'Test Org', $2)", org_id, slug + ) + return org_id + + +async def _create_member(db: AppDatabase, org_id: UUID) -> UUID: + user_id = uuid4() + await db.execute( + "INSERT INTO users (id, email) VALUES ($1, $2)", user_id, f"{user_id.hex[:12]}@example.com" + ) + await db.execute( + "INSERT INTO org_memberships (user_id, org_id, role) VALUES ($1, $2, 'member')", + user_id, + org_id, + ) + return user_id + + +async def _create_datasource(db: AppDatabase, org_id: UUID) -> UUID: + datasource_id = uuid4() + await db.execute( + """INSERT INTO data_sources (id, tenant_id, name, type, connection_config_encrypted) + VALUES ($1, $2, $3, 'postgresql', 'unused')""", + datasource_id, + org_id, + f"warehouse-{datasource_id.hex[:8]}", + ) + return datasource_id + + +async def _create_investigation(db: AppDatabase, org_id: UUID, datasource_id: UUID) -> UUID: + """Store an investigation the way POST /investigations does.""" + alert = AnomalyAlert( + dataset_ids=["public.orders"], + metric_spec=MetricSpec.from_column("customer_id"), + anomaly_type="null_rate", + expected_value=0.01, + actual_value=0.25, + deviation_pct=2400.0, + anomaly_date="2026-09-01", + severity="high", + ) + investigation_id = uuid4() + await db.execute( + "INSERT INTO investigations (id, tenant_id, alert) VALUES ($1, $2, $3)", + investigation_id, + org_id, + json.dumps({**alert.model_dump(mode="json"), "datasource_id": str(datasource_id)}), + ) + return investigation_id + + +async def _assert_access_follows_datasource( + db: AppDatabase, org_id: UUID, user_id: UUID, visible: UUID, hidden: UUID +) -> None: + async with db.acquire() as conn: + service = PermissionService(conn) + assert await service.can_access_investigation(user_id, visible) + assert not await service.can_access_investigation(user_id, hidden) + assert await service.get_accessible_investigation_ids(user_id, org_id) == [visible] + + +async def test_user_datasource_grant_opens_that_datasources_investigations( + migrated_db: AppDatabase, +) -> None: + """A member granted a datasource sees the investigations that ran against it.""" + org_id = await _create_org(migrated_db) + user_id = await _create_member(migrated_db, org_id) + granted = await _create_datasource(migrated_db, org_id) + other = await _create_datasource(migrated_db, org_id) + visible = await _create_investigation(migrated_db, org_id, granted) + hidden = await _create_investigation(migrated_db, org_id, other) + await migrated_db.execute( + "INSERT INTO permission_grants (org_id, user_id, data_source_id) VALUES ($1, $2, $3)", + org_id, + user_id, + granted, + ) + + await _assert_access_follows_datasource(migrated_db, org_id, user_id, visible, hidden) + + +async def test_team_datasource_grant_opens_that_datasources_investigations( + migrated_db: AppDatabase, +) -> None: + """A member of a team granted a datasource sees the investigations that ran against it.""" + org_id = await _create_org(migrated_db) + user_id = await _create_member(migrated_db, org_id) + granted = await _create_datasource(migrated_db, org_id) + other = await _create_datasource(migrated_db, org_id) + visible = await _create_investigation(migrated_db, org_id, granted) + hidden = await _create_investigation(migrated_db, org_id, other) + team_id = uuid4() + await migrated_db.execute( + "INSERT INTO teams (id, org_id, name) VALUES ($1, $2, 'Data Platform')", team_id, org_id + ) + await migrated_db.execute( + "INSERT INTO team_members (team_id, user_id) VALUES ($1, $2)", team_id, user_id + ) + await migrated_db.execute( + "INSERT INTO permission_grants (org_id, team_id, data_source_id) VALUES ($1, $2, $3)", + org_id, + team_id, + granted, + ) + + await _assert_access_follows_datasource(migrated_db, org_id, user_id, visible, hidden) diff --git a/python-packages/dataing/tests/integration/test_fix_feedback_export_context.py b/python-packages/dataing/tests/integration/test_fix_feedback_export_context.py new file mode 100644 index 000000000..6ffad54b0 --- /dev/null +++ b/python-packages/dataing/tests/integration/test_fix_feedback_export_context.py @@ -0,0 +1,80 @@ +"""Integration tests for the investigation context of fix feedback exports.""" + +from __future__ import annotations + +import json +from datetime import UTC, datetime +from uuid import UUID, uuid4 + +import pytest + +from dataing.adapters.db.app_db import AppDatabase +from dataing.services.feedback import FixFeedbackService + +pytestmark = pytest.mark.integration + +CREATED_AT = datetime(2026, 9, 1, 12, 0, tzinfo=UTC) + + +async def _create_investigation(db: AppDatabase) -> tuple[UUID, UUID]: + tenant_id = uuid4() + await db.execute( + "INSERT INTO tenants (id, name, slug) VALUES ($1, $2, $3)", + tenant_id, + "Test Tenant", + f"test-{tenant_id.hex[:12]}", + ) + investigation_id = uuid4() + await db.execute( + "INSERT INTO investigations (id, tenant_id, alert, created_at) VALUES ($1, $2, $3, $4)", + investigation_id, + tenant_id, + json.dumps({"dataset_ids": ["public.orders"]}), + CREATED_AT, + ) + return tenant_id, investigation_id + + +async def _spawn_from_issue( + db: AppDatabase, tenant_id: UUID, investigation_id: UUID, number: int +) -> UUID: + issue_id = uuid4() + await db.execute( + "INSERT INTO issues (id, tenant_id, number, title) VALUES ($1, $2, $3, 'Null spike')", + issue_id, + tenant_id, + number, + ) + await db.execute( + """INSERT INTO issue_investigation_runs (issue_id, investigation_id, trigger_type) + VALUES ($1, $2, 'human')""", + issue_id, + investigation_id, + ) + return issue_id + + +async def test_context_names_the_issue_the_investigation_was_spawned_from( + migrated_db: AppDatabase, +) -> None: + """The issue comes from issue_investigation_runs.""" + tenant_id, investigation_id = await _create_investigation(migrated_db) + issue_id = await _spawn_from_issue(migrated_db, tenant_id, investigation_id, number=1) + + context = await FixFeedbackService(db=migrated_db).investigation_context(investigation_id) + + assert context == { + "issue_id": str(issue_id), + "investigation_created_at": CREATED_AT.isoformat(), + } + + +async def test_context_has_no_issue_for_an_investigation_started_directly( + migrated_db: AppDatabase, +) -> None: + """Investigations started from an alert have no issue to name.""" + _, investigation_id = await _create_investigation(migrated_db) + + context = await FixFeedbackService(db=migrated_db).investigation_context(investigation_id) + + assert context == {"investigation_created_at": CREATED_AT.isoformat()} diff --git a/python-packages/dataing/tests/integration/test_investigation_completed_at.py b/python-packages/dataing/tests/integration/test_investigation_completed_at.py new file mode 100644 index 000000000..94f86c5ce --- /dev/null +++ b/python-packages/dataing/tests/integration/test_investigation_completed_at.py @@ -0,0 +1,78 @@ +"""Integration tests for investigations.completed_at on the migrated schema.""" + +from __future__ import annotations + +import json +from uuid import UUID, uuid4 + +import pytest + +from dataing.adapters.db.app_db import AppDatabase + +pytestmark = pytest.mark.integration + + +async def _create_investigation(db: AppDatabase) -> UUID: + tenant_id = uuid4() + await db.execute( + "INSERT INTO tenants (id, name, slug) VALUES ($1, $2, $3)", + tenant_id, + "Test Tenant", + f"test-{tenant_id.hex[:12]}", + ) + investigation_id = uuid4() + await db.execute( + "INSERT INTO investigations (id, tenant_id, alert) VALUES ($1, $2, $3)", + investigation_id, + tenant_id, + json.dumps({"dataset_ids": ["public.orders"]}), + ) + return investigation_id + + +async def _completion(db: AppDatabase, investigation_id: UUID) -> tuple[str, object]: + row = await db.fetch_one( + "SELECT status, completed_at FROM investigations WHERE id = $1", investigation_id + ) + assert row is not None + return row["status"], row["completed_at"] + + +async def test_completed_at_is_stamped_when_the_outcome_is_first_set( + migrated_db: AppDatabase, +) -> None: + """Writers only set outcome; the database records when that first happened.""" + investigation_id = await _create_investigation(migrated_db) + assert await _completion(migrated_db, investigation_id) == ("active", None) + + await migrated_db.execute( + "UPDATE investigations SET outcome = $2 WHERE id = $1", + investigation_id, + json.dumps({"status": "completed"}), + ) + status, completed_at = await _completion(migrated_db, investigation_id) + assert status == "completed" + assert completed_at is not None + + await migrated_db.execute( + "UPDATE investigations SET outcome = $2 WHERE id = $1", + investigation_id, + json.dumps({"status": "completed", "root_cause": "late upstream load"}), + ) + assert await _completion(migrated_db, investigation_id) == ("completed", completed_at) + + +async def test_clearing_the_outcome_reopens_the_investigation(migrated_db: AppDatabase) -> None: + """An investigation without an outcome is active and has no completion time.""" + investigation_id = await _create_investigation(migrated_db) + await migrated_db.execute( + "UPDATE investigations SET outcome = $2 WHERE id = $1", + investigation_id, + json.dumps({"status": "completed"}), + ) + + await migrated_db.execute( + "UPDATE investigations SET outcome = NULL WHERE id = $1", investigation_id + ) + + assert await _completion(migrated_db, investigation_id) == ("active", None)