specimba's picture
download
raw
6.84 kB
"""FastAPI surface for the canonical Nexus Governor."""
from __future__ import annotations
import os
from typing import Any, Dict, Optional
from fastapi import FastAPI, HTTPException
from fastapi.middleware.cors import CORSMiddleware
from pydantic import BaseModel, Field
from nexus_os.api.service import GovernanceAPIState
class SkillProposalRequest(BaseModel):
"""Incoming GSPP proposal request."""
proposal_id: str = Field(..., min_length=1)
model_id: str = Field(..., min_length=1)
skill: Dict[str, Any] = Field(default_factory=dict)
rationale: str = Field(..., min_length=1)
timestamp: Optional[str] = None
class ApprovalRequest(BaseModel):
"""Manual approval/rejection request for held proposals."""
decision: str = Field(..., pattern="^(approve|reject|hold)$")
reviewer: str = Field(default="speci")
reason: str = Field(default="")
class TaskHeartbeatRequest(BaseModel):
"""Governed agent task heartbeat payload."""
task_id: str = Field(..., min_length=1)
agent_id: str = Field(..., min_length=1)
trace_id: str = Field(..., min_length=1)
progress: int = Field(default=0, ge=0, le=100)
status: str = Field(default="heartbeat", min_length=1)
details: Dict[str, Any] = Field(default_factory=dict)
timestamp: Optional[str] = None
class TaskResultRequest(BaseModel):
"""Governed task completion/failure payload."""
task_id: str = Field(..., min_length=1)
agent_id: str = Field(..., min_length=1)
trace_id: str = Field(..., min_length=1)
outcome: str = Field(..., min_length=1)
summary: str = Field(default="")
artifacts: Dict[str, Any] = Field(default_factory=dict)
timestamp: Optional[str] = None
class RouteLogRequest(BaseModel):
"""Route attempt logging request for 7352 governance."""
route_class: str = Field(..., min_length=1)
provider_id: str = Field(..., min_length=1)
model_family: str = Field(..., min_length=1)
latency_ms: float = Field(..., ge=0)
success: bool
failure_code: Optional[str] = None
fallback_count: int = Field(default=0, ge=0)
stream_disconnect: bool = False
timestamp: Optional[str] = None
metadata: Dict[str, Any] = Field(default_factory=dict)
def create_app(db_path: Optional[str] = None) -> FastAPI:
"""Create the FastAPI governance app."""
state = GovernanceAPIState(
db_path=db_path or os.getenv("NEXUS_API_DB_PATH", "nexus_api.db")
)
api = FastAPI(title="Nexus OS Governance API", version="3.0.0")
api.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_methods=["*"],
allow_headers=["*"],
)
@api.get("/health")
def health() -> Dict[str, Any]:
return {
"status": "operational",
"service": "nexus-governance-api",
"proposals": len(state.proposals),
}
@api.post("/skills/propose")
def propose_skill(request: SkillProposalRequest) -> Dict[str, Any]:
try:
return state.propose_skill(
proposal_id=request.proposal_id,
model_id=request.model_id,
skill=request.skill,
rationale=request.rationale,
timestamp=request.timestamp,
)
except ValueError as exc:
raise HTTPException(status_code=409, detail=str(exc)) from exc
@api.get("/skills/status/{proposal_id}")
def proposal_status(proposal_id: str) -> Dict[str, Any]:
record = state.get_proposal(proposal_id)
if record is None:
raise HTTPException(status_code=404, detail="Proposal not found")
return record
@api.get("/governance/proposals")
def list_proposals() -> Dict[str, Any]:
return state.list_proposals()
@api.post("/governance/approve/{proposal_id}")
def approve_proposal(
proposal_id: str, request: ApprovalRequest
) -> Dict[str, Any]:
try:
record = state.review_proposal(
proposal_id=proposal_id,
decision=request.decision,
reviewer=request.reviewer,
reason=request.reason,
)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
if record is None:
raise HTTPException(status_code=404, detail="Proposal not found")
return record
@api.post("/governance/log")
def log_route_attempt(request: RouteLogRequest) -> Dict[str, Any]:
"""Log a route attempt for governance audit and learning."""
return state.record_route_log(
route_class=request.route_class,
provider_id=request.provider_id,
model_family=request.model_family,
latency_ms=request.latency_ms,
success=request.success,
failure_code=request.failure_code,
fallback_count=request.fallback_count,
stream_disconnect=request.stream_disconnect,
timestamp=request.timestamp,
metadata=request.metadata,
)
@api.get("/governance/routes")
def list_route_logs(
provider_id: Optional[str] = None,
route_class: Optional[str] = None,
limit: int = 100,
) -> Dict[str, Any]:
"""List route logs for analysis."""
return state.list_route_logs(
provider_id=provider_id,
route_class=route_class,
limit=limit,
)
@api.get("/dashboard/stats")
def dashboard_stats() -> Dict[str, Any]:
return state.dashboard_stats()
@api.post("/tasks/heartbeat")
def task_heartbeat(request: TaskHeartbeatRequest) -> Dict[str, Any]:
return state.record_task_heartbeat(
task_id=request.task_id,
agent_id=request.agent_id,
trace_id=request.trace_id,
progress=request.progress,
status=request.status,
details=request.details,
timestamp=request.timestamp,
)
@api.post("/tasks/result")
def task_result(request: TaskResultRequest) -> Dict[str, Any]:
return state.record_task_result(
task_id=request.task_id,
agent_id=request.agent_id,
trace_id=request.trace_id,
outcome=request.outcome,
summary=request.summary,
artifacts=request.artifacts,
timestamp=request.timestamp,
)
@api.get("/tasks/status/{task_id}")
def task_status(task_id: str) -> Dict[str, Any]:
record = state.get_task_status(task_id)
if record is None:
raise HTTPException(status_code=404, detail="Task not found")
return record
@api.on_event("shutdown")
def close_db() -> None:
state.close()
return api
app = create_app()

Xet Storage Details

Size:
6.84 kB
·
Xet hash:
86dae58dccf92ba35b055fc81ab03b1a25a7817b8e459eff4e841f41f6ccac25

Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.