205 lines
7.3 KiB
Python
205 lines
7.3 KiB
Python
from __future__ import annotations
|
|
|
|
import uuid
|
|
from datetime import date
|
|
|
|
from fastapi import APIRouter, Depends, Header, Query, Response
|
|
from sqlalchemy import select
|
|
from sqlalchemy.orm import Session
|
|
|
|
from app.api.deps import McpClientContext, get_db, require_mcp_service_token
|
|
from app.core.config import get_settings
|
|
from app.core.errors import AppError
|
|
from app.models.booking import Booking
|
|
from app.models.data_quality import DataQualityIssue
|
|
from app.models.vehicle import Vehicle
|
|
from app.schemas import (
|
|
AttentionVehicleOut,
|
|
McpKnowledgeSearchRequest,
|
|
McpVehicleDetailOut,
|
|
OperationsSummaryOut,
|
|
)
|
|
from app.services.audit import record_audit_event
|
|
from app.services.knowledge import GroundedAnswer, get_knowledge_provider
|
|
from app.services.operations import compute_metrics, list_attention_vehicles
|
|
|
|
router = APIRouter(prefix="/api/v1/integrations/mcp", tags=["mcp"])
|
|
settings = get_settings()
|
|
|
|
|
|
def get_correlation_id(
|
|
x_correlation_id: str | None = Header(default=None, alias="X-Correlation-Id"),
|
|
) -> str:
|
|
"""Preserve the Hub's own inbound correlation ID through MCP client -> Hub -> Fleet
|
|
Ops -> RAGcore -> Fleet Ops Audit; only mint a fresh one when none was supplied or
|
|
it isn't a valid UUID (per the task's own correlation-propagation contract)."""
|
|
if x_correlation_id:
|
|
try:
|
|
return str(uuid.UUID(x_correlation_id))
|
|
except ValueError:
|
|
pass
|
|
return str(uuid.uuid4())
|
|
|
|
|
|
def _audit_service_request(
|
|
db: Session,
|
|
*,
|
|
reported_client_id: str,
|
|
tool: str,
|
|
status_label: str,
|
|
correlation_id: str,
|
|
metadata: dict[str, object] | None = None,
|
|
) -> None:
|
|
record_audit_event(
|
|
db,
|
|
actor_type="service",
|
|
# The shared service token authenticates the Hub, not the caller identity that
|
|
# the Hub reports in a header. Keep attribution authoritative and retain the
|
|
# reported value only as explicitly non-authenticated diagnostic metadata.
|
|
actor_label="itworx-mcp-hub",
|
|
action="mcp_tool_request",
|
|
entity_type="mcp_tool",
|
|
correlation_id=uuid.UUID(correlation_id),
|
|
metadata={
|
|
"tool": tool,
|
|
"status": status_label,
|
|
"reported_client_id": reported_client_id,
|
|
**(metadata or {}),
|
|
},
|
|
)
|
|
db.commit()
|
|
|
|
|
|
def _set_trace_headers(response: Response, correlation_id: str, tenant: str) -> None:
|
|
response.headers["X-Correlation-Id"] = correlation_id
|
|
response.headers["X-Tenant-Id"] = tenant
|
|
response.headers["Cache-Control"] = "no-store"
|
|
|
|
|
|
@router.get("/operations-summary", response_model=OperationsSummaryOut)
|
|
def operations_summary(
|
|
response: Response,
|
|
db: Session = Depends(get_db),
|
|
client: McpClientContext = Depends(require_mcp_service_token),
|
|
correlation_id: str = Depends(get_correlation_id),
|
|
) -> OperationsSummaryOut:
|
|
metrics = compute_metrics(db)
|
|
_set_trace_headers(response, correlation_id, client.tenant)
|
|
_audit_service_request(
|
|
db,
|
|
reported_client_id=client.reported_client_id,
|
|
tool="fleet_ops_get_operations_summary",
|
|
status_label="ok",
|
|
correlation_id=correlation_id,
|
|
)
|
|
return OperationsSummaryOut(tenant=client.tenant, metrics=metrics)
|
|
|
|
|
|
@router.get("/attention-vehicles", response_model=list[AttentionVehicleOut])
|
|
def attention_vehicles(
|
|
response: Response,
|
|
minimum_severity: str = Query(default="medium", pattern="^(low|medium|high)$"),
|
|
date_filter: date | None = Query(default=None, alias="date"),
|
|
limit: int = Query(default=20, ge=1, le=50),
|
|
db: Session = Depends(get_db),
|
|
client: McpClientContext = Depends(require_mcp_service_token),
|
|
correlation_id: str = Depends(get_correlation_id),
|
|
) -> list[AttentionVehicleOut]:
|
|
results = list_attention_vehicles(
|
|
db, minimum_severity=minimum_severity, on_or_before=date_filter, limit=limit
|
|
)
|
|
_set_trace_headers(response, correlation_id, client.tenant)
|
|
_audit_service_request(
|
|
db,
|
|
reported_client_id=client.reported_client_id,
|
|
tool="fleet_ops_list_attention_vehicles",
|
|
status_label="ok",
|
|
correlation_id=correlation_id,
|
|
)
|
|
return [AttentionVehicleOut(**r) for r in results]
|
|
|
|
|
|
@router.get("/vehicles/{vehicle_ref}", response_model=McpVehicleDetailOut)
|
|
def vehicle_details(
|
|
vehicle_ref: str,
|
|
response: Response,
|
|
db: Session = Depends(get_db),
|
|
client: McpClientContext = Depends(require_mcp_service_token),
|
|
correlation_id: str = Depends(get_correlation_id),
|
|
) -> McpVehicleDetailOut:
|
|
_set_trace_headers(response, correlation_id, client.tenant)
|
|
vehicle = db.scalar(select(Vehicle).where(Vehicle.public_ref == vehicle_ref))
|
|
if vehicle is None:
|
|
_audit_service_request(
|
|
db,
|
|
reported_client_id=client.reported_client_id,
|
|
tool="fleet_ops_get_vehicle_details",
|
|
status_label="not_found",
|
|
correlation_id=correlation_id,
|
|
)
|
|
raise AppError("VEHICLE_NOT_FOUND", "Vehicle not found.", status_code=404)
|
|
|
|
open_issue_count = len(
|
|
db.scalars(
|
|
select(DataQualityIssue.id).where(
|
|
DataQualityIssue.entity_type == "vehicle",
|
|
DataQualityIssue.entity_id == vehicle.id,
|
|
DataQualityIssue.status == "open",
|
|
)
|
|
).all()
|
|
)
|
|
current_booking = db.scalar(
|
|
select(Booking).where(Booking.vehicle_id == vehicle.id, Booking.status == "active")
|
|
)
|
|
|
|
_audit_service_request(
|
|
db,
|
|
reported_client_id=client.reported_client_id,
|
|
tool="fleet_ops_get_vehicle_details",
|
|
status_label="ok",
|
|
correlation_id=correlation_id,
|
|
)
|
|
return McpVehicleDetailOut(
|
|
public_ref=vehicle.public_ref,
|
|
make=vehicle.make,
|
|
model=vehicle.model,
|
|
model_year=vehicle.model_year,
|
|
location=vehicle.location,
|
|
operational_status=vehicle.operational_status,
|
|
odometer_km=vehicle.odometer_km,
|
|
next_service_km=vehicle.next_service_km,
|
|
open_quality_issue_count=open_issue_count,
|
|
current_booking_ref=current_booking.public_ref if current_booking else None,
|
|
)
|
|
|
|
|
|
@router.post("/search-knowledge", response_model=GroundedAnswer)
|
|
def search_knowledge(
|
|
body: McpKnowledgeSearchRequest,
|
|
response: Response,
|
|
db: Session = Depends(get_db),
|
|
client: McpClientContext = Depends(require_mcp_service_token),
|
|
correlation_id: str = Depends(get_correlation_id),
|
|
) -> GroundedAnswer:
|
|
provider = get_knowledge_provider()
|
|
answer = provider.ask(body.question, correlation_id, language=body.locale)
|
|
source_count_available = len(answer.sources)
|
|
answer.sources = answer.sources[: body.max_sources]
|
|
_set_trace_headers(response, correlation_id, client.tenant)
|
|
response.headers["X-Sources-Available"] = str(source_count_available)
|
|
response.headers["X-Sources-Returned"] = str(len(answer.sources))
|
|
_audit_service_request(
|
|
db,
|
|
reported_client_id=client.reported_client_id,
|
|
tool="fleet_ops_search_knowledge",
|
|
status_label=answer.evidence_state,
|
|
correlation_id=correlation_id,
|
|
metadata={
|
|
"tenant": client.tenant,
|
|
"locale": body.locale,
|
|
"sources_available": source_count_available,
|
|
"sources_returned": len(answer.sources),
|
|
},
|
|
)
|
|
return answer
|