Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,5 +18,12 @@
- Do not reintroduce floating unpinned `ruff`/`ty` in CI without lockfile pins
- Prefer project-scoped changes; no drive-by refactors outside the issue

## MVP foundation (issues #75–#81)
- Sample policy fixtures: [`docs/mvp/fixtures/`](docs/mvp/fixtures/) — Markdown with H1 + `## Controls`; covered by `backend/tests/test_mvp_fixtures.py`
- Verdict events: `backend/app/services/verdict_log.py` — validate returns `verdict_event` with required keys (`action`, `verdict`, `repo`, `timestamp`, `policy_id`/`policyId`, `validation_id`, `rule_id`)
- Policies API: `GET /policies`, `GET /policies/{id}`; `POST /ingest` returns `policy_id` and versions by content hash
- Backend tests: `cd backend && uv run pytest` (pyproject sets `pythonpath = ["."]`)
- Prefer TDD on MVP slices; update the matching `docs/handoffs/YYYY-MM-DD-*.md` when status changes

## Speckit
Feature plans live under `specs/` and `.specify/`. Constitution: `docs/constitution.md` → `.specify/memory/constitution.md`.
63 changes: 56 additions & 7 deletions backend/app/api/ingest.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,17 +2,58 @@
import tempfile

from fastapi import APIRouter, BackgroundTasks, Depends, File, HTTPException, UploadFile
from sqlmodel import Session, select
from sqlmodel import Session, col, select

from ..core.db import get_session
from ..core.hashing import calculate_sha256
from ..models.policy import Policy
from ..models.task import IngestionTask, TaskStatus
from ..services.pipeline import start_ingestion_pipeline
from ..worker.tasks import run_background_task

router = APIRouter()


def _upsert_policy(
session: Session,
*,
name: str,
file_hash: str,
task_id,
source_url: str | None = None,
) -> Policy:
"""Create or reuse a Policy version row for this ingest."""
existing_hash = session.exec(
select(Policy).where(Policy.hash == file_hash)
).first()
if existing_hash:
return existing_hash

previous = session.exec(
select(Policy)
.where(Policy.name == name, Policy.is_current.is_(True))
.order_by(col(Policy.version).desc())
).first()

next_version = (previous.version + 1) if previous else 1
if previous:
previous.is_current = False
session.add(previous)

policy = Policy(
name=name,
source_url=source_url,
version=next_version,
hash=file_hash,
task_id=task_id,
is_current=True,
)
session.add(policy)
session.commit()
session.refresh(policy)
return policy


@router.post("/ingest", status_code=202)
async def ingest_document(
background_tasks: BackgroundTasks,
Expand All @@ -21,16 +62,21 @@ async def ingest_document(
):
content = await file.read()
file_hash = calculate_sha256(content)
name = file.filename or "untitled"

# Check for existing task (idempotency)
existing_task = session.exec(
select(IngestionTask).where(IngestionTask.source_hash == file_hash)
).first()
if existing_task:
policy = _upsert_policy(
session, name=name, file_hash=file_hash, task_id=existing_task.id
)
return {
"task_id": existing_task.id,
"status": existing_task.status,
"message": "File already processed or in progress.",
"policy_id": str(policy.id),
}

# Create new task
Expand All @@ -39,9 +85,11 @@ async def ingest_document(
session.commit()
session.refresh(task)

policy = _upsert_policy(session, name=name, file_hash=file_hash, task_id=task.id)

# Save file temporarily for processing
temp_dir = tempfile.mkdtemp(prefix=f"ingest_{task.id}_")
temp_path = os.path.join(temp_dir, file.filename)
temp_path = os.path.join(temp_dir, name)
with open(temp_path, "wb") as f:
f.write(content)

Expand All @@ -50,11 +98,12 @@ async def ingest_document(
run_background_task, start_ingestion_pipeline, str(task.id), temp_path
)

task_id = str(task.id)
task_status = task.status
task_progress = task.progress_pct

return {"task_id": task_id, "status": task_status, "progress_pct": task_progress}
return {
"task_id": str(task.id),
"status": task.status,
"progress_pct": task.progress_pct,
"policy_id": str(policy.id),
}


@router.get("/ingest/{task_id}")
Expand Down
43 changes: 43 additions & 0 deletions backend/app/api/policies.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
from uuid import UUID

from fastapi import APIRouter, Depends, HTTPException
from sqlmodel import Session, col, select

from ..core.db import get_session
from ..models.policy import Policy

router = APIRouter()


def _serialize(policy: Policy) -> dict:
return {
"id": str(policy.id),
"name": policy.name,
"source_url": policy.source_url,
"version": policy.version,
"hash": policy.hash,
"task_id": str(policy.task_id) if policy.task_id else None,
"is_current": policy.is_current,
"created_at": policy.created_at.isoformat(),
}


@router.get("/policies")
async def list_policies(
current_only: bool = True,
session: Session = Depends(get_session),
):
statement = select(Policy)
if current_only:
statement = statement.where(Policy.is_current.is_(True))
statement = statement.order_by(col(Policy.name), col(Policy.version).desc())
policies = session.exec(statement).all()
return [_serialize(p) for p in policies]


@router.get("/policies/{policy_id}")
async def get_policy(policy_id: UUID, session: Session = Depends(get_session)):
policy = session.get(Policy, policy_id)
if not policy:
raise HTTPException(status_code=404, detail="Policy not found")
return _serialize(policy)
70 changes: 52 additions & 18 deletions backend/app/api/validate.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,12 @@
from pydantic import BaseModel
from sqlmodel import Session, col, select

from ..core.config import settings
from ..core.db import get_session
from ..models.rule import GovernanceRule
from ..services.logger import ActivityLog, log_activity
from ..services.logger import ActivityLog, engine as activity_engine, log_activity
from ..services.validator import validator
from ..services.verdict_log import build_verdict_event, emit_verdict_log

router = APIRouter()

Expand All @@ -17,6 +19,8 @@ class ValidateRequest(BaseModel):
code_snippet: str
rule_id: str
context: str | None = None
repo: str | None = None
policy_id: str | None = None


class HITLAction(BaseModel):
Expand All @@ -40,19 +44,44 @@ async def submit_validation(
# Perform validation
result = await validator.validate_risk(request.code_snippet, rule_content)

# Log activity
status = "complete" if result["verdict"] != "MID" else "pending"

# Pre-allocate validation id via activity log, then attach structured verdict
details = {
"request": request.model_dump(),
"result": result,
}
log = log_activity(
action="risk_validation",
status="complete" if result["verdict"] != "MID" else "pending",
details={
"request": request.model_dump(),
"result": result,
},
status=status,
details=details,
)

verdict_event = build_verdict_event(
action="risk_validation",
verdict=result["verdict"],
repo=request.repo or getattr(settings, "PROJECT_NAME", "") or "",
policy_id=request.policy_id,
validation_id=str(log.id),
rule_id=request.rule_id,
timestamp=log.timestamp if log.timestamp.tzinfo else log.timestamp.replace(tzinfo=None),
)
emit_verdict_log(verdict_event)

# Persist structured event on the activity row (same engine as log_activity)
with Session(activity_engine) as s:
row = s.get(ActivityLog, log.id)
if row:
merged = dict(row.details or {})
merged["verdict_event"] = verdict_event
row.details = merged
s.add(row)
s.commit()

return {
"validation_id": str(log.id),
"status": "processing" if result["verdict"] == "MID" else "complete",
"verdict_event": verdict_event,
}


Expand All @@ -63,14 +92,18 @@ async def get_validation(id: UUID, session: Session = Depends(get_session)):
raise HTTPException(status_code=404, detail="Validation result not found")

result = log.details.get("result", {})
return {
verdict_event = log.details.get("verdict_event")
payload = {
"validation_id": str(log.id),
"verdict": result.get("verdict", "HIGH"),
"reasoning": result.get("reasoning", "No reasoning provided"),
"activity_logged": True,
"created_at": log.timestamp.isoformat(),
"status": log.status,
}
if verdict_event:
payload["verdict_event"] = verdict_event
return payload


@router.patch("/validate/{id}")
Expand Down Expand Up @@ -118,15 +151,16 @@ async def list_validations(
if verdict and current_verdict != verdict.upper():
continue

results.append(
{
"validation_id": str(log.id),
"verdict": current_verdict,
"reasoning": res.get("reasoning"),
"activity_logged": True,
"created_at": log.timestamp.isoformat(),
"status": log.status,
}
)
item = {
"validation_id": str(log.id),
"verdict": current_verdict,
"reasoning": res.get("reasoning"),
"activity_logged": True,
"created_at": log.timestamp.isoformat(),
"status": log.status,
}
if "verdict_event" in (log.details or {}):
item["verdict_event"] = log.details["verdict_event"]
results.append(item)

return results
17 changes: 17 additions & 0 deletions backend/app/models/policy.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
"""Policy metadata + version history (#75)."""

from datetime import UTC, datetime
from uuid import UUID, uuid4

from sqlmodel import Field, SQLModel


class Policy(SQLModel, table=True):
id: UUID = Field(default_factory=uuid4, primary_key=True)
name: str = Field(index=True)
source_url: str | None = None
version: int = Field(default=1)
hash: str = Field(index=True)
task_id: UUID | None = Field(default=None, foreign_key="ingestiontask.id")
is_current: bool = Field(default=True, index=True)
created_at: datetime = Field(default_factory=lambda: datetime.now(UTC))
3 changes: 3 additions & 0 deletions backend/app/services/logger.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,14 @@ def log_activity(
):
"""
Persist activity log to database.
Returns an expunged instance safe to read after the session closes.
"""
with Session(engine) as session:
log = ActivityLog(
action=action, details=details, task_id=task_id, status=status
)
session.add(log)
session.commit()
session.refresh(log)
session.expunge(log)
return log
56 changes: 56 additions & 0 deletions backend/app/services/verdict_log.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
"""Canonical structured verdict event schema (#76). Loki shipping is #77."""

from __future__ import annotations

from datetime import UTC, datetime
from typing import Any

from app.core.logging import logger

REQUIRED_VERDICT_KEYS = (
"action",
"verdict",
"repo",
"timestamp",
"policy_id",
"policyId",
"validation_id",
"rule_id",
)


def build_verdict_event(
*,
action: str,
verdict: str,
repo: str | None = None,
policy_id: str | None = None,
validation_id: str | None = None,
rule_id: str | None = None,
timestamp: datetime | None = None,
extra: dict[str, Any] | None = None,
) -> dict[str, Any]:
"""Build a JSON-serializable verdict event with required keys."""
ts = timestamp or datetime.now(UTC)
pid = policy_id
event: dict[str, Any] = {
"action": action,
"verdict": verdict,
"repo": repo or "",
"timestamp": ts.isoformat(),
"policy_id": pid,
"policyId": pid,
"validation_id": validation_id,
"rule_id": rule_id,
}
if extra:
event.update(extra)
return event


def emit_verdict_log(event: dict[str, Any]) -> None:
"""Emit structured verdict to application logger (stdout JSON-ish)."""
missing = [k for k in ("action", "verdict", "timestamp", "validation_id") if k not in event]
if missing:
logger.warning("verdict_event missing keys: %s", missing)
logger.info("verdict_event %s", event)
3 changes: 2 additions & 1 deletion backend/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware

from app.api import activity, ingest, manifest, validate
from app.api import activity, ingest, manifest, policies, validate
from app.core.config import settings
from app.core.db import init_db
from app.core.logging import logger
Expand Down Expand Up @@ -36,6 +36,7 @@ async def lifespan(app: FastAPI):

# Include routers
app.include_router(ingest.router, tags=["Ingestion"])
app.include_router(policies.router, tags=["Policies"])
app.include_router(manifest.router, tags=["Manifest"])
app.include_router(validate.router, tags=["Validation"])
app.include_router(activity.router, tags=["Activity"])
Expand Down
Loading
Loading