The original docs described two n8n workflows but the repository only ever shipped one (return-processing); the sketched second workflow (knowledge sync) depends on RAGcore, which isn't connected here, so it stays deferred. Add POST /api/v1/integrations/n8n/scheduled-scan (X-Service-Token protected, same pattern as the return callback), calling the same run_scan() the manual "Run quality scan" UI action uses and recording a service-actor data_quality_scan_run audit event. run_scan() already only creates an issue for a condition without one open, so overlapping triggers do no duplicate domain work. n8n/mobilityops-scheduled-quality-scan.json (hourly schedule + manual test trigger, both feeding the same HTTP call) ships "active": false so it can't fire anywhere until deliberately published. Verified live against the local n8n instance via the Manual test trigger: full green execution, and the resulting data_quality_scan_run audit event (actor_type=service, actor_label="n8n scheduled scan") confirms the real round trip, not just a contract test. deploy/unraid/setup-scheduled-scan.sh mirrors the existing return-workflow publish script for the shared Unraid n8n.
92 lines
3.3 KiB
Python
92 lines
3.3 KiB
Python
from __future__ import annotations
|
|
|
|
import uuid
|
|
from datetime import UTC, datetime
|
|
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 ScanResultOut
|
|
from app.services.audit import record_audit_event
|
|
from app.services.data_quality import run_scan
|
|
|
|
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)
|