Fresh hub entity per the founder-reviewed decision (not a suggestions rename-bridge): kind: intake per canon/standards/work-record-types_v0.1.md, lifecycle open -> vetted -> routed -> closed(promoted|declined|absorbed). - api/models/base.py::new_uuid7 -- dependency-free RFC 9562 UUIDv7 generator (48-bit ms timestamp, version/variant bits, random remainder); existing tables keep new_uuid (UUIDv4) unchanged, this is opt-in for new work-record entities per the identity-layering canon - api/models/intake.py: Intake + IntakeNote ORM models, mirroring Decision's shape (topic/workplan/repo scope, lane, status, outcome, promoted_to back-link); CHECK constraints enforce scope-required, closed-requires-outcome, promoted-requires-promoted_to at the DB level - migrations/a7c3e9f1b4d2: intakes + intake_notes tables, 3 enum types - api/routers/intake.py: list/create/get/patch + /route + /close + /notes actions, mirroring decisions.py's pattern (409 on invalid transitions, progress event on close) - api/schemas/intake.py: Pydantic create/update/route/close/note schemas - mcp_server/server.py: create_intake, list_intakes, route_intake, close_intake tool wrappers - tests/test_intake.py: 12 tests against the real Postgres test DB (create/list/scope-validation, full lifecycle incl. 409s and the promoted-requires-promoted_to constraint, notes, UUIDv7 verification) Verified live against the running dev API + DB (not just pytest): applied the migration, restarted the MCP server, and ran a full create -> route -> close cycle over the real REST endpoints. No regressions: full existing suite (test_routers_core, test_suggestions, test_mcp_smoke, test_mcp_write_tools, test_mcp_registration, test_consistency_check, test_consistency_sweep) all green. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
143 lines
5.4 KiB
Python
143 lines
5.4 KiB
Python
import hashlib
|
|
import os
|
|
import time
|
|
from contextlib import asynccontextmanager
|
|
|
|
from fastapi import FastAPI
|
|
from fastapi.middleware.cors import CORSMiddleware
|
|
from starlette.middleware.base import BaseHTTPMiddleware
|
|
from starlette.requests import Request
|
|
from starlette.responses import Response as StarletteResponse
|
|
|
|
from api.database import engine
|
|
from api.events import shutdown_publisher
|
|
from api.services.write_idempotency import WriteIdempotencyMiddleware
|
|
from api.routers import decisions, extension_points, intake, progress, state, suggestions, tasks, technical_debt, topics, workstreams, workstream_dependencies
|
|
from api.routers import domains, repos, contributions, sbom, policy, domain_goals, repo_goals, messages, capability_requests, tpsc, services
|
|
from api.routers import token_events
|
|
from api.routers import interface_changes
|
|
from api.routers import flows
|
|
from api.routers import recently_on_scope
|
|
from api.routers import consistency_sweep
|
|
from api.routers import reconciliation
|
|
from api.routers import execution
|
|
from api.routers import fabric
|
|
from api.routers import legacy_meter
|
|
|
|
|
|
class ETagMiddleware(BaseHTTPMiddleware):
|
|
"""Add ETag + conditional-GET (304) support to all JSON GET responses."""
|
|
|
|
async def dispatch(self, request: Request, call_next):
|
|
started = time.perf_counter()
|
|
response = await call_next(request)
|
|
if request.method != "GET":
|
|
response.headers["X-StateHub-Elapsed-Ms"] = f"{(time.perf_counter() - started) * 1000:.1f}"
|
|
return response
|
|
if "application/json" not in response.headers.get("content-type", ""):
|
|
response.headers["X-StateHub-Elapsed-Ms"] = f"{(time.perf_counter() - started) * 1000:.1f}"
|
|
return response
|
|
|
|
body_parts = []
|
|
async for chunk in response.body_iterator:
|
|
body_parts.append(chunk)
|
|
body = b"".join(body_parts)
|
|
elapsed_ms = f"{(time.perf_counter() - started) * 1000:.1f}"
|
|
|
|
etag = '"' + hashlib.md5(body, usedforsecurity=False).hexdigest() + '"'
|
|
if request.headers.get("if-none-match") == etag:
|
|
return StarletteResponse(
|
|
status_code=304,
|
|
headers={
|
|
"ETag": etag,
|
|
"Cache-Control": "no-cache",
|
|
"X-StateHub-Elapsed-Ms": elapsed_ms,
|
|
"X-StateHub-Response-Bytes": "0",
|
|
},
|
|
)
|
|
|
|
headers = {k: v for k, v in response.headers.items() if k.lower() != "content-length"}
|
|
headers["ETag"] = etag
|
|
headers["X-StateHub-Elapsed-Ms"] = elapsed_ms
|
|
headers["X-StateHub-Response-Bytes"] = str(len(body))
|
|
if not any(k.lower() == "cache-control" for k in headers):
|
|
headers["Cache-Control"] = "no-cache"
|
|
return StarletteResponse(
|
|
content=body,
|
|
status_code=response.status_code,
|
|
headers=headers,
|
|
media_type=response.media_type,
|
|
)
|
|
|
|
|
|
@asynccontextmanager
|
|
async def lifespan(app: FastAPI):
|
|
yield
|
|
await shutdown_publisher()
|
|
await engine.dispose()
|
|
|
|
|
|
app = FastAPI(
|
|
title="Custodian State Hub",
|
|
description="Local-first state API for the Custodian agent system.",
|
|
version="0.6.0",
|
|
lifespan=lifespan,
|
|
)
|
|
|
|
_default_dashboard_origins = [
|
|
*(f"http://localhost:{port}" for port in range(3000, 3006)),
|
|
*(f"http://127.0.0.1:{port}" for port in range(3000, 3006)),
|
|
*(f"http://[::1]:{port}" for port in range(3000, 3006)),
|
|
]
|
|
_cors_env = os.getenv("CORS_ORIGINS", ",".join(_default_dashboard_origins))
|
|
_cors_origins = [o.strip() for o in _cors_env.split(",") if o.strip()]
|
|
|
|
app.add_middleware(WriteIdempotencyMiddleware)
|
|
app.add_middleware(ETagMiddleware)
|
|
app.add_middleware(
|
|
CORSMiddleware,
|
|
allow_origins=_cors_origins,
|
|
allow_methods=["GET", "POST", "PATCH", "DELETE", "PUT"],
|
|
allow_headers=["Content-Type", "If-None-Match", "Idempotency-Key", "X-StateHub-Source-Agent", "X-StateHub-Source-Host"],
|
|
expose_headers=["ETag", "X-StateHub-Elapsed-Ms", "X-StateHub-Response-Bytes", "X-StateHub-Cache", "X-StateHub-Idempotency-Replay"],
|
|
)
|
|
|
|
app.include_router(domains.router)
|
|
app.include_router(recently_on_scope.hourly_router)
|
|
app.include_router(recently_on_scope.router)
|
|
app.include_router(consistency_sweep.router)
|
|
app.include_router(repos.router)
|
|
app.include_router(topics.router)
|
|
app.include_router(workstreams.router)
|
|
app.include_router(workstreams.workplan_router)
|
|
app.include_router(workstream_dependencies.router)
|
|
app.include_router(workstream_dependencies.workplan_router)
|
|
app.include_router(tasks.router)
|
|
app.include_router(decisions.router)
|
|
app.include_router(intake.router)
|
|
app.include_router(extension_points.router)
|
|
app.include_router(technical_debt.router)
|
|
app.include_router(progress.router)
|
|
app.include_router(domain_goals.router)
|
|
app.include_router(repo_goals.router)
|
|
app.include_router(contributions.router)
|
|
app.include_router(sbom.router)
|
|
app.include_router(messages.router)
|
|
app.include_router(capability_requests.router)
|
|
app.include_router(suggestions.router)
|
|
app.include_router(tpsc.router)
|
|
app.include_router(services.router)
|
|
app.include_router(token_events.router)
|
|
app.include_router(interface_changes.router)
|
|
app.include_router(flows.router)
|
|
app.include_router(reconciliation.router)
|
|
app.include_router(execution.router)
|
|
app.include_router(fabric.router)
|
|
app.include_router(legacy_meter.router)
|
|
app.include_router(state.router)
|
|
app.include_router(policy.router)
|
|
|
|
|
|
@app.get("/", include_in_schema=False)
|
|
async def root():
|
|
return {"service": "dev-hub", "docs": "/docs"}
|