From ec36ee83aa7d2c57017a14ef7f42b7ff2cdc859f Mon Sep 17 00:00:00 2001 From: bordumb Date: Sun, 27 Sep 2026 21:02:26 +0100 Subject: [PATCH] fix(api): finish moving investigation queries to the unified schema MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Follow-up to #183. These call sites still read columns that 013_unified_investigation.sql dropped and failed at runtime: - Dashboard: get_dashboard_stats filtered on completed_at and on statuses the generated status column never takes. Migration 036 adds investigations.completed_at, stamped by a trigger the first time outcome is set; active = no outcome, completed today = completed_at today. list_investigations now returns summaries from the alert JSONB (primary dataset, metric display name or anomaly type, severity), "unknown" for imported replays. The frontend called /dashboard/stats without /api/v1, read camelCase keys the API never sends and fell back to mock numbers; it now uses the generated client and shows "—" and an error on failure. - RBAC: PermissionService matched datasource grants on investigations.data_source_id, so every access check raised. It now compares with alert->>'datasource_id'. - Fix feedback export: the investigation context read investigations.issue_id; it now names the issue from issue_investigation_runs. - EE runbooks: generation from an issue read issues.resolution/metadata and investigations.issue_id/synthesis; it now uses resolution_note, issue_labels and the latest linked investigation's outcome, and the route records that investigation. Runbook responses also decode JSONB text, which failed validation on every runbook route. Integration tests run against the migrated schema via migrated_db. Co-Authored-By: Claude Opus 5.5 Signed-off-by: Claude --- .flow/tasks/fn-16.17.json | 22 ++ .flow/tasks/fn-16.17.md | 22 ++ .flow/tasks/fn-16.18.json | 21 ++ .flow/tasks/fn-16.18.md | 18 ++ .flow/tasks/fn-16.19.json | 22 ++ .flow/tasks/fn-16.19.md | 19 ++ .flow/tasks/fn-16.20.json | 21 ++ .flow/tasks/fn-16.20.md | 20 ++ .flow/tasks/fn-16.21.json | 14 ++ .flow/tasks/fn-16.21.md | 18 ++ .../src/features/dashboard/dashboard-page.tsx | 69 +----- .../dashboard/dashboard-stats.test.tsx | 43 ++++ .../features/dashboard/dashboard-stats.tsx | 52 +++++ frontend/app/src/lib/api/dashboard.test.ts | 43 ++++ frontend/app/src/lib/api/dashboard.ts | 21 +- .../src/dataing_ee/core/runbook/generator.py | 42 ++-- .../entrypoints/api/routes/runbooks.py | 24 +- .../api/test_runbook_from_issue.py | 157 +++++++++++++ .../036_investigation_completed_at.sql | 32 +++ .../dataing/src/dataing/adapters/db/app_db.py | 37 ++- .../dataing/core/rbac/permission_service.py | 8 +- .../entrypoints/api/routes/dashboard.py | 2 +- .../dataing/src/dataing/services/feedback.py | 55 +++-- .../tests/integration/api/test_dashboard.py | 216 ++++++++++++++++++ .../core/test_permission_service.py | 131 +++++++++++ .../test_fix_feedback_export_context.py | 80 +++++++ .../test_investigation_completed_at.py | 78 +++++++ 27 files changed, 1151 insertions(+), 136 deletions(-) create mode 100644 .flow/tasks/fn-16.17.json create mode 100644 .flow/tasks/fn-16.17.md create mode 100644 .flow/tasks/fn-16.18.json create mode 100644 .flow/tasks/fn-16.18.md create mode 100644 .flow/tasks/fn-16.19.json create mode 100644 .flow/tasks/fn-16.19.md create mode 100644 .flow/tasks/fn-16.20.json create mode 100644 .flow/tasks/fn-16.20.md create mode 100644 .flow/tasks/fn-16.21.json create mode 100644 .flow/tasks/fn-16.21.md create mode 100644 frontend/app/src/features/dashboard/dashboard-stats.test.tsx create mode 100644 frontend/app/src/features/dashboard/dashboard-stats.tsx create mode 100644 frontend/app/src/lib/api/dashboard.test.ts create mode 100644 python-packages/dataing-ee/tests/integration/api/test_runbook_from_issue.py create mode 100644 python-packages/dataing/migrations/036_investigation_completed_at.sql create mode 100644 python-packages/dataing/tests/integration/api/test_dashboard.py create mode 100644 python-packages/dataing/tests/integration/core/test_permission_service.py create mode 100644 python-packages/dataing/tests/integration/test_fix_feedback_export_context.py create mode 100644 python-packages/dataing/tests/integration/test_investigation_completed_at.py 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)