Files
MobilityOps/backend/app/api/routers/integrations.py
T
NuklearRabbit 0da5251524 n8n: add backend endpoints for the RAGcore Procedure Sync workflow
GET /api/v1/integrations/n8n/procedures lists every procedure Markdown
file Fleet Ops ships (all languages) with a stable per-document id and
content hash, ready for workflow 3 to push into RAGcore. POST
.../procedures-sync-result records the sync outcome as an idempotent
audit event, matching the existing return-callback/workflow-error
pattern. Extracted frontmatter parsing out of the demo knowledge
provider into a shared module so both read the same source of truth.
2026-08-04 19:47:39 +02:00

223 lines
8.4 KiB
Python

from __future__ import annotations
import uuid
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
from fastapi import APIRouter, Depends, Header
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.api.deps import get_db
from app.core.config import get_settings
from app.core.errors import AppError
from app.models.audit import AuditEvent
from app.models.outbox import OutboxEvent
from app.schemas import (
ProcedureDocumentOut,
ProcedureListOut,
ProcedureSyncResultIn,
ProcedureSyncResultResult,
ScanResultOut,
WorkflowErrorReportIn,
WorkflowErrorReportResult,
)
from app.services.audit import record_audit_event
from app.services.data_quality import run_scan
from app.services.knowledge.procedures import iter_procedure_documents
router = APIRouter(prefix="/api/v1/integrations/n8n", tags=["integrations"])
settings = get_settings()
@router.post("/return-callback")
def return_callback(
body: dict[str, Any],
idempotency_key: str = Header(..., alias="Idempotency-Key"),
service_token: str = Header(..., alias="X-Service-Token"),
db: Session = Depends(get_db),
) -> dict:
if service_token != settings.n8n_callback_token:
raise AppError("UNAUTHORIZED_SERVICE", "Invalid service token.", status_code=401)
try:
event_id = uuid.UUID(idempotency_key)
except ValueError as exc:
raise AppError(
"INVALID_IDEMPOTENCY_KEY", "Idempotency-Key must be the event's UUID.", status_code=422
) from exc
event = db.scalar(select(OutboxEvent).where(OutboxEvent.event_id == event_id))
if event is None:
raise AppError("EVENT_NOT_FOUND", "No outbox event matches this event ID.", status_code=404)
# Idempotent by event ID: n8n or our own dispatcher may redeliver the same event
# (e.g. a lost response after a timeout), so this callback must not double-record.
already_recorded = (
db.scalar(
select(AuditEvent.id).where(
AuditEvent.action == "n8n_return_followup_recorded",
AuditEvent.metadata_json["event_id"].astext == str(event_id),
)
)
is not None
)
if not already_recorded:
record_audit_event(
db,
actor_type="service",
actor_label="n8n",
action="n8n_return_followup_recorded",
entity_type="booking",
correlation_id=uuid.UUID(body.get("correlation_id"))
if body.get("correlation_id")
else None,
after={"follow_up": body.get("follow_up"), "summary": body.get("summary")},
metadata={"event_id": str(event_id)},
)
db.commit()
return {
"status": "recorded",
"event_id": str(event_id),
"occurred_at": datetime.now(UTC).isoformat(),
}
@router.post("/scheduled-scan", response_model=ScanResultOut)
def scheduled_scan(
service_token: str = Header(..., alias="X-Service-Token"),
db: Session = Depends(get_db),
) -> ScanResultOut:
"""Triggered by the scheduled n8n quality-scan workflow. Narrow, read-mostly, and
safe to call repeatedly: run_scan() only ever creates an issue for a condition that
doesn't already have one open, so a duplicate or overlapping trigger does no
duplicate domain work -- it just reports zero new issues for anything already known."""
if service_token != settings.n8n_callback_token:
raise AppError("UNAUTHORIZED_SERVICE", "Invalid service token.", status_code=401)
result = run_scan(db, actor_label="n8n scheduled scan", actor_type="service")
return ScanResultOut(created=result.created)
@router.post("/workflow-error", response_model=WorkflowErrorReportResult)
def workflow_error(
body: WorkflowErrorReportIn,
service_token: str = Header(..., alias="X-Service-Token"),
db: Session = Depends(get_db),
) -> WorkflowErrorReportResult:
"""Receives a bounded, secret-free failure report from the central n8n "Fleet Ops --
Workflow Error Handler" workflow, which is attached as the Error Workflow on every
other Fleet Ops n8n workflow. Idempotent on execution_id: n8n may redeliver the same
error report (e.g. after a timed-out response), so this must not double-record."""
if service_token != settings.n8n_callback_token:
raise AppError("UNAUTHORIZED_SERVICE", "Invalid service token.", status_code=401)
already_recorded = (
db.scalar(
select(AuditEvent.id).where(
AuditEvent.action == "n8n_workflow_failure_registered",
AuditEvent.metadata_json["execution_id"].astext == body.execution_id,
)
)
is not None
)
if not already_recorded:
correlation_id: uuid.UUID | None = None
if body.correlation_id:
try:
correlation_id = uuid.UUID(body.correlation_id)
except ValueError:
correlation_id = None
record_audit_event(
db,
actor_type="service",
actor_label="n8n error handler",
action="n8n_workflow_failure_registered",
entity_type="automation",
correlation_id=correlation_id,
after={
"workflow_id": body.workflow_id,
"workflow_name": body.workflow_name,
"error_category": body.error_category,
"error_summary": body.error_summary,
"trigger_context": body.trigger_context,
"attempt": body.attempt,
"retry_action": body.retry_action,
"failed_at": body.failed_at.isoformat(),
},
metadata={"execution_id": body.execution_id},
)
db.commit()
return WorkflowErrorReportResult(
status="already_registered" if already_recorded else "registered",
execution_id=body.execution_id,
occurred_at=datetime.now(UTC),
)
@router.get("/procedures", response_model=ProcedureListOut)
def list_procedures(service_token: str = Header(..., alias="X-Service-Token")) -> ProcedureListOut:
"""Read-only source list for the RAGcore Procedure Sync workflow: every procedure
Markdown file Fleet Ops ships, across every supported language, with a stable
per-document id (source_id) and a content hash so the caller can detect changes
without re-fetching content it already has."""
if service_token != settings.n8n_callback_token:
raise AppError("UNAUTHORIZED_SERVICE", "Invalid service token.", status_code=401)
documents = [
ProcedureDocumentOut(
id=doc.source_id,
language=doc.language,
document_id=doc.document_id,
title=doc.title,
version=doc.version,
content=doc.content,
content_hash=doc.content_hash,
)
for doc in iter_procedure_documents(Path(settings.knowledge_dir))
]
return ProcedureListOut(documents=documents)
@router.post("/procedures-sync-result", response_model=ProcedureSyncResultResult)
def procedures_sync_result(
body: ProcedureSyncResultIn,
service_token: str = Header(..., alias="X-Service-Token"),
db: Session = Depends(get_db),
) -> ProcedureSyncResultResult:
"""Receives a summary (counts only, no document content) from the n8n "Fleet Ops --
RAGcore Procedure Sync" workflow once it finishes uploading procedures to RAGcore.
Idempotent on execution_id, matching the workflow-error and return-callback pattern."""
if service_token != settings.n8n_callback_token:
raise AppError("UNAUTHORIZED_SERVICE", "Invalid service token.", status_code=401)
already_recorded = (
db.scalar(
select(AuditEvent.id).where(
AuditEvent.action == "n8n_procedures_synced",
AuditEvent.metadata_json["execution_id"].astext == body.execution_id,
)
)
is not None
)
if not already_recorded:
record_audit_event(
db,
actor_type="service",
actor_label="n8n procedure sync",
action="n8n_procedures_synced",
entity_type="automation",
after={"synced": body.synced, "failed": body.failed},
metadata={"execution_id": body.execution_id},
)
db.commit()
return ProcedureSyncResultResult(
status="already_registered" if already_recorded else "registered",
execution_id=body.execution_id,
occurred_at=datetime.now(UTC),
)