n8n: surface real per-workflow evidence on the integration status page
Fleet Ops integration status no longer depends only on a config boolean or the most recent outbox event: N8nIntegrationStatus now reports per-canonical-workflow evidence (last successful outbox delivery for the return workflow, latest service-triggered data_quality_scan_run for the scan workflow, latest n8n_workflow_failure_registered for the error handler, and "not built" for the still-blocked RAGcore sync), plus an error-handler summary (total failures registered, latest failure + which workflow). Automation page renders this as a localized workflow table (EN/NL/FR) with technical workflow names tucked under a "Technical details" disclosure, matching the existing progressive-disclosure pattern. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Sonnet 5
parent
e39c0a1dd6
commit
4049c0c6b1
@@ -223,6 +223,18 @@ class SearchResponse(BaseModel):
|
||||
results: list[SearchResultItem]
|
||||
|
||||
|
||||
class N8nWorkflowEvidence(BaseModel):
|
||||
name: str
|
||||
built: bool
|
||||
last_seen_at: datetime | None
|
||||
|
||||
|
||||
class N8nErrorHandlerStatus(BaseModel):
|
||||
total_failures_registered: int
|
||||
latest_failure_at: datetime | None
|
||||
latest_failure_workflow: str | None
|
||||
|
||||
|
||||
class N8nIntegrationStatus(BaseModel):
|
||||
configured: bool
|
||||
dispatch_enabled: bool
|
||||
@@ -233,6 +245,10 @@ class N8nIntegrationStatus(BaseModel):
|
||||
succeeded: int
|
||||
latest_success_at: datetime | None
|
||||
latest_failure_at: datetime | None
|
||||
expected_workflow_count: int
|
||||
known_workflow_count: int
|
||||
workflows: list[N8nWorkflowEvidence]
|
||||
error_handler: N8nErrorHandlerStatus
|
||||
|
||||
|
||||
class McpHubIntegrationStatus(BaseModel):
|
||||
|
||||
@@ -6,11 +6,21 @@ from sqlalchemy import func, select
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from app.core.config import get_settings
|
||||
from app.models.audit import AuditEvent
|
||||
from app.models.outbox import OutboxEvent
|
||||
from app.schemas import N8nIntegrationStatus
|
||||
from app.schemas import N8nErrorHandlerStatus, N8nIntegrationStatus, N8nWorkflowEvidence
|
||||
|
||||
settings = get_settings()
|
||||
|
||||
# The 4 canonical Fleet Ops n8n workflows (see n8n/workflows/MANIFEST.md). Workflow 3
|
||||
# (RAGcore Procedure Sync) is not built yet, so it always reports no evidence.
|
||||
_CANONICAL_WORKFLOWS = (
|
||||
"Fleet Ops — Vehicle Return Orchestration",
|
||||
"Fleet Ops — Scheduled Data Quality Scan",
|
||||
"Fleet Ops — RAGcore Procedure Sync",
|
||||
"Fleet Ops — Workflow Error Handler",
|
||||
)
|
||||
|
||||
|
||||
def derive_n8n_status(db: Session) -> N8nIntegrationStatus:
|
||||
counts: dict[str, int] = dict(
|
||||
@@ -42,6 +52,52 @@ def derive_n8n_status(db: Session) -> N8nIntegrationStatus:
|
||||
else:
|
||||
state = "no_evidence"
|
||||
|
||||
# Scheduled scan evidence: only service-triggered runs count as n8n evidence, not
|
||||
# runs an operator triggered manually from the Data Quality page.
|
||||
latest_scan_at = db.scalar(
|
||||
select(func.max(AuditEvent.occurred_at)).where(
|
||||
AuditEvent.action == "data_quality_scan_run",
|
||||
AuditEvent.actor_type == "service",
|
||||
)
|
||||
)
|
||||
|
||||
# Error handler evidence: registrations posted by the "Fleet Ops — Workflow Error
|
||||
# Handler" n8n workflow itself, which also doubles as proof that workflow is wired
|
||||
# up and firing correctly.
|
||||
total_failures_registered = (
|
||||
db.scalar(
|
||||
select(func.count(AuditEvent.id)).where(
|
||||
AuditEvent.action == "n8n_workflow_failure_registered"
|
||||
)
|
||||
)
|
||||
or 0
|
||||
)
|
||||
latest_failure_row = db.execute(
|
||||
select(AuditEvent.occurred_at, AuditEvent.after_json)
|
||||
.where(AuditEvent.action == "n8n_workflow_failure_registered")
|
||||
.order_by(AuditEvent.occurred_at.desc())
|
||||
.limit(1)
|
||||
).first()
|
||||
latest_handler_failure_at = latest_failure_row[0] if latest_failure_row else None
|
||||
latest_handler_failure_workflow = (
|
||||
(latest_failure_row[1] or {}).get("workflow_name") if latest_failure_row else None
|
||||
)
|
||||
|
||||
evidence_by_workflow = {
|
||||
"Fleet Ops — Vehicle Return Orchestration": latest_success_at,
|
||||
"Fleet Ops — Scheduled Data Quality Scan": latest_scan_at,
|
||||
"Fleet Ops — RAGcore Procedure Sync": None,
|
||||
"Fleet Ops — Workflow Error Handler": latest_handler_failure_at,
|
||||
}
|
||||
workflows = [
|
||||
N8nWorkflowEvidence(
|
||||
name=name,
|
||||
built=name != "Fleet Ops — RAGcore Procedure Sync",
|
||||
last_seen_at=evidence_by_workflow[name],
|
||||
)
|
||||
for name in _CANONICAL_WORKFLOWS
|
||||
]
|
||||
|
||||
return N8nIntegrationStatus(
|
||||
configured=bool(settings.n8n_webhook_url),
|
||||
dispatch_enabled=settings.n8n_dispatch_enabled,
|
||||
@@ -52,4 +108,12 @@ def derive_n8n_status(db: Session) -> N8nIntegrationStatus:
|
||||
succeeded=succeeded,
|
||||
latest_success_at=latest_success_at,
|
||||
latest_failure_at=latest_failure_at,
|
||||
expected_workflow_count=len(_CANONICAL_WORKFLOWS),
|
||||
known_workflow_count=sum(1 for w in workflows if w.last_seen_at is not None),
|
||||
workflows=workflows,
|
||||
error_handler=N8nErrorHandlerStatus(
|
||||
total_failures_registered=total_failures_registered,
|
||||
latest_failure_at=latest_handler_failure_at,
|
||||
latest_failure_workflow=latest_handler_failure_workflow,
|
||||
),
|
||||
)
|
||||
|
||||
@@ -39,3 +39,82 @@ def test_integration_status_is_operational_once_all_failed_events_resolved(ops_c
|
||||
body = response.json()["n8n"]
|
||||
assert body["failed"] == 0
|
||||
assert body["state"] == "operational"
|
||||
|
||||
|
||||
def test_integration_status_lists_all_four_canonical_workflows(ops_client):
|
||||
body = ops_client.get("/api/v1/integrations/status").json()["n8n"]
|
||||
assert body["expected_workflow_count"] == 4
|
||||
names = {w["name"] for w in body["workflows"]}
|
||||
assert names == {
|
||||
"Fleet Ops — Vehicle Return Orchestration",
|
||||
"Fleet Ops — Scheduled Data Quality Scan",
|
||||
"Fleet Ops — RAGcore Procedure Sync",
|
||||
"Fleet Ops — Workflow Error Handler",
|
||||
}
|
||||
ragcore_sync = next(w for w in body["workflows"] if "RAGcore" in w["name"])
|
||||
assert ragcore_sync["built"] is False
|
||||
assert ragcore_sync["last_seen_at"] is None
|
||||
|
||||
|
||||
def test_integration_status_scheduled_scan_evidence_only_counts_service_runs(client, ops_client):
|
||||
from app.core.config import get_settings
|
||||
|
||||
settings = get_settings()
|
||||
|
||||
before = ops_client.get("/api/v1/integrations/status").json()["n8n"]
|
||||
scan_workflow = next(
|
||||
w for w in before["workflows"] if w["name"].endswith("Scheduled Data Quality Scan")
|
||||
)
|
||||
assert scan_workflow["last_seen_at"] is None
|
||||
|
||||
scan = client.post(
|
||||
"/api/v1/integrations/n8n/scheduled-scan",
|
||||
headers={"X-Service-Token": settings.n8n_callback_token},
|
||||
)
|
||||
assert scan.status_code == 200
|
||||
|
||||
after = ops_client.get("/api/v1/integrations/status").json()["n8n"]
|
||||
scan_workflow = next(
|
||||
w for w in after["workflows"] if w["name"].endswith("Scheduled Data Quality Scan")
|
||||
)
|
||||
assert scan_workflow["last_seen_at"] is not None
|
||||
assert after["known_workflow_count"] > before["known_workflow_count"]
|
||||
|
||||
|
||||
def test_integration_status_reflects_error_handler_registrations(client, ops_client):
|
||||
import uuid
|
||||
|
||||
from app.core.config import get_settings
|
||||
|
||||
settings = get_settings()
|
||||
before = ops_client.get("/api/v1/integrations/status").json()["n8n"]
|
||||
|
||||
execution_id = str(uuid.uuid4())
|
||||
report = client.post(
|
||||
"/api/v1/integrations/n8n/workflow-error",
|
||||
json={
|
||||
"workflow_id": "mobilityops-return-processing",
|
||||
"workflow_name": "Fleet Ops — Vehicle Return Orchestration",
|
||||
"execution_id": execution_id,
|
||||
"failed_at": "2026-08-04T10:15:00Z",
|
||||
"error_category": "httpError",
|
||||
"error_summary": "Simulated failure for status test",
|
||||
"trigger_context": "webhook",
|
||||
"attempt": 1,
|
||||
},
|
||||
headers={"X-Service-Token": settings.n8n_callback_token},
|
||||
)
|
||||
assert report.status_code == 200
|
||||
|
||||
after = ops_client.get("/api/v1/integrations/status").json()["n8n"]
|
||||
assert (
|
||||
after["error_handler"]["total_failures_registered"]
|
||||
== before["error_handler"]["total_failures_registered"] + 1
|
||||
)
|
||||
assert after["error_handler"]["latest_failure_workflow"] == (
|
||||
"Fleet Ops — Vehicle Return Orchestration"
|
||||
)
|
||||
handler_workflow = next(
|
||||
w for w in after["workflows"] if w["name"].endswith("Workflow Error Handler")
|
||||
)
|
||||
assert handler_workflow["last_seen_at"] is not None
|
||||
|
||||
Reference in New Issue
Block a user