Files
impactflow_discovery/app/routers/discovery.py
T
Claude 866bd60225 Persist discovery answers to CSV for reprocessing and recovery
Save each discovery conversation's prompts and answers to durable CSV
files (per-conversation + append-only master log) on both save and
completion, so answers survive an extraction error, can be re-fed to the
AI, and can be reviewed/resumed by the user.

- app/services/answer_store.py: canonical prompt list + atomic CSV writes,
  master append, and read-back helpers (DB stays system of record; CSV
  failures are logged, never fatal).
- discovery router: write CSV on /respond and /complete; new endpoints
  GET /answers, GET /answers.csv, POST /reprocess (shared extraction
  helper; locked profiles return 409).
- discovery.html: prefill/resume from saved answers after an error and a
  "Re-run analysis" button wired to /reprocess.
- scripts/reprocess_csv.py: offline CLI to re-run extraction from a CSV
  (print or --write-db).
- QUESTIONS_DIR / QUESTIONS_MASTER_CSV config, .gitignore, README, tests.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Dg6XWUwprmP5QCL18HxssY
2026-06-19 01:06:40 +00:00

627 lines
21 KiB
Python

"""Discovery API routes: start a conversation, save responses, generate and
confirm a profile.
All routes are user-scoped via the dual-auth dependency. The MCP server uses
the X-API-Key header and operates under the synthetic admin user; the
browser app uses a Google-issued JWT.
"""
import json
import os
import uuid
from datetime import datetime, timezone
from fastapi import APIRouter, Depends, HTTPException
from fastapi.responses import Response
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app import schemas
from app.auth import get_current_user
from app.database import get_db
from app.models import (
DiscoveryConversation,
DiscoveryProfile,
ProfileRevision,
ReflectionMessage,
User,
)
from app.services import answer_store
from app.services.answer_store import DISCOVERY_PROMPTS
from app.services.extractor import DiscoveryExtractionError, DiscoveryExtractor
from app.services.reflector import (
EDITABLE_FIELDS,
ReflectionCoach,
ReflectionError,
)
from app.services.profile_history import goal_history
router = APIRouter(prefix="/discovery", tags=["discovery"])
def _now() -> datetime:
return datetime.now(timezone.utc)
def _profile_fields(profile: DiscoveryProfile) -> dict:
"""The seven editable prose fields as a plain dict (for snapshots and the
reflector's profile context)."""
return {f: getattr(profile, f) for f in EDITABLE_FIELDS}
def _record_revision(
db: AsyncSession,
profile: DiscoveryProfile,
source: str,
note: str | None = None,
) -> None:
"""Snapshot the profile's editable prose into profile_revision. Caller
commits."""
db.add(
ProfileRevision(
id=str(uuid.uuid4()),
profile_id=profile.id,
user_id=profile.user_id,
source=source,
fields_json=json.dumps(_profile_fields(profile)),
note=note,
created_at=_now(),
)
)
async def _generate_profile(
db: AsyncSession,
conversation: DiscoveryConversation,
source: str,
) -> tuple[DiscoveryProfile, str | None]:
"""Run the extractor over a conversation's answers and stage a new profile
(plus a revision snapshot) on the session. The caller commits.
Shared by ``/complete`` (source="extraction") and ``/reprocess``
(source="reprocess"). Raises HTTPException(400) when there is nothing to
analyze and HTTPException(502) on an extractor failure.
"""
responses = {
p["key"]: getattr(conversation, p["column"], None) or ""
for p in DISCOVERY_PROMPTS
}
if not any(text.strip() for text in responses.values()):
raise HTTPException(
status_code=400, detail="No responses available to analyze"
)
api_key = os.getenv("ANTHROPIC_API_KEY")
model = os.getenv("ANTHROPIC_MODEL", "claude-sonnet-4-6")
try:
extractor = DiscoveryExtractor(api_key=api_key, model=model)
data = await extractor.extract(responses)
except DiscoveryExtractionError as exc:
raise HTTPException(status_code=502, detail=str(exc)) from exc
profile = DiscoveryProfile(
id=str(uuid.uuid4()),
user_id=conversation.user_id,
conversation_id=conversation.id,
generated_at=_now(),
triad=data.get("triad"),
probable_type=_as_int(data.get("probable_type")),
wing=_as_int(data.get("wing")),
instinctual_variant=data.get("instinctual_variant"),
instinctual_stack=data.get("instinctual_stack"),
love_summary=data.get("love_summary"),
strength_summary=data.get("strength_summary"),
mission_summary=data.get("mission_summary"),
vocation_summary=data.get("vocation_summary"),
overlap_narrative=data.get("overlap_narrative"),
short_term_goals=data.get("short_term_goals"),
long_term_goals=data.get("long_term_goals"),
confidence_json=json.dumps(data.get("confidence", {})),
locked=False,
)
db.add(profile)
_record_revision(db, profile, source=source)
return profile, data.get("extraction_notes")
def _to_profile_response(
profile: DiscoveryProfile, extraction_notes: str | None = None
) -> schemas.ProfileResponse:
"""Build a ProfileResponse from a stored profile row."""
confidence = None
if profile.confidence_json:
try:
confidence = schemas.Confidence(**json.loads(profile.confidence_json))
except (json.JSONDecodeError, TypeError, ValueError):
confidence = None
return schemas.ProfileResponse(
id=profile.id,
user_id=profile.user_id,
conversation_id=profile.conversation_id,
generated_at=profile.generated_at,
triad=profile.triad,
probable_type=profile.probable_type,
wing=profile.wing,
instinctual_variant=profile.instinctual_variant,
instinctual_stack=profile.instinctual_stack,
love_summary=profile.love_summary,
strength_summary=profile.strength_summary,
mission_summary=profile.mission_summary,
vocation_summary=profile.vocation_summary,
overlap_narrative=profile.overlap_narrative,
short_term_goals=profile.short_term_goals,
long_term_goals=profile.long_term_goals,
confidence=confidence,
locked=profile.locked,
extraction_notes=extraction_notes,
)
async def _latest_profile(
db: AsyncSession, user_id: str
) -> DiscoveryProfile | None:
stmt = (
select(DiscoveryProfile)
.where(DiscoveryProfile.user_id == user_id)
.order_by(DiscoveryProfile.generated_at.desc())
)
result = await db.execute(stmt)
return result.scalars().first()
async def _owned_conversation(
db: AsyncSession, conversation_id: str, user: User
) -> DiscoveryConversation:
"""Fetch a conversation and assert the caller owns it (admins bypass).
Returns 404 for both 'not found' and 'not yours' so the existence of
other users' conversations isn't leaked."""
conversation = await db.get(DiscoveryConversation, conversation_id)
if conversation is None:
raise HTTPException(status_code=404, detail="Conversation not found")
if conversation.user_id != user.id and user.role != "admin":
raise HTTPException(status_code=404, detail="Conversation not found")
return conversation
@router.post("/start", response_model=schemas.StartResponse)
async def start_conversation(
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
conversation = DiscoveryConversation(
id=str(uuid.uuid4()),
user_id=user.id,
started_at=_now(),
)
db.add(conversation)
await db.commit()
return schemas.StartResponse(conversation_id=conversation.id)
@router.put(
"/{conversation_id}/respond", response_model=schemas.RespondResponse
)
async def save_responses(
conversation_id: str,
payload: schemas.RespondRequest,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
conversation = await _owned_conversation(db, conversation_id, user)
conversation.prompt_alive = payload.prompt_alive
conversation.prompt_friction = payload.prompt_friction
conversation.prompt_pull = payload.prompt_pull
conversation.prompt_recognition = payload.prompt_recognition
conversation.prompt_future = payload.prompt_future
conversation.prompt_goals_short = payload.prompt_goals_short
conversation.prompt_goals_long = payload.prompt_goals_long
await db.commit()
# Durable CSV copy, written before extraction runs so the answers survive
# an extraction error and can be reviewed or re-processed later.
answer_store.save_answers(conversation, user, status="responses_saved")
return schemas.RespondResponse(
conversation_id=conversation_id, status="responses_saved"
)
@router.post(
"/{conversation_id}/complete", response_model=schemas.ProfileResponse
)
async def complete_conversation(
conversation_id: str,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
conversation = await _owned_conversation(db, conversation_id, user)
profile, extraction_notes = await _generate_profile(
db, conversation, source="extraction"
)
conversation.completed_at = _now()
await db.commit()
# Refresh the durable CSV copy now that the conversation is complete.
answer_store.save_answers(conversation, user, status="completed")
return _to_profile_response(profile, extraction_notes=extraction_notes)
@router.get(
"/{conversation_id}/answers", response_model=schemas.AnswersResponse
)
async def get_answers(
conversation_id: str,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""The saved prompts and answers for a conversation, so the person can
review or resume from their original responses (e.g. after an error)."""
conversation = await _owned_conversation(db, conversation_id, user)
answers = [
schemas.AnswerItem(
prompt_key=p["key"],
prompt_title=p["title"],
answer=getattr(conversation, p["column"], None) or "",
)
for p in DISCOVERY_PROMPTS
]
return schemas.AnswersResponse(
conversation_id=conversation.id,
started_at=conversation.started_at,
completed_at=conversation.completed_at,
answers=answers,
)
@router.get("/{conversation_id}/answers.csv")
async def download_answers_csv(
conversation_id: str,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""Download a conversation's answers as a CSV file (built from the DB so it
works even if the on-disk copy was never written)."""
conversation = await _owned_conversation(db, conversation_id, user)
status = "completed" if conversation.completed_at else "responses_saved"
csv_text = answer_store.conversation_csv_text(conversation, user, status)
filename = f"discovery-{conversation_id}.csv"
return Response(
content=csv_text,
media_type="text/csv",
headers={
"Content-Disposition": f'attachment; filename="{filename}"'
},
)
@router.post(
"/{conversation_id}/reprocess", response_model=schemas.ProfileResponse
)
async def reprocess_conversation(
conversation_id: str,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""Re-run AI extraction over a conversation's saved answers, producing a
fresh profile. Used to recover from an extraction error or to regenerate a
profile after the answers were re-fed. The latest profile must be unlocked.
"""
conversation = await _owned_conversation(db, conversation_id, user)
existing = await _latest_profile(db, user.id)
if existing is not None and existing.locked:
raise HTTPException(
status_code=409,
detail="Profile is affirmed and locked; it cannot be reprocessed.",
)
profile, extraction_notes = await _generate_profile(
db, conversation, source="reprocess"
)
if conversation.completed_at is None:
conversation.completed_at = _now()
await db.commit()
answer_store.save_answers(conversation, user, status="completed")
return _to_profile_response(profile, extraction_notes=extraction_notes)
@router.get("/profile/me", response_model=schemas.ProfileResponse)
async def get_my_profile(
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
profile = await _latest_profile(db, user.id)
if profile is None:
raise HTTPException(status_code=404, detail="No profile for this user")
return _to_profile_response(profile)
@router.patch("/profile/me", response_model=schemas.ProfileResponse)
async def update_my_profile(
payload: schemas.ProfileUpdate,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""Edit the prose of the latest profile. The person owns their words, so
they can revise any summary, the narrative, or their goals — but only
while the profile is unlocked. Affirming (locking) makes it final."""
profile = await _latest_profile(db, user.id)
if profile is None:
raise HTTPException(status_code=404, detail="No profile for this user")
if profile.locked:
raise HTTPException(
status_code=409,
detail="Profile is locked; it can no longer be edited.",
)
updates = payload.model_dump(exclude_unset=True)
if not updates:
raise HTTPException(status_code=400, detail="No fields to update")
for field, value in updates.items():
setattr(profile, field, value)
_record_revision(db, profile, source="manual_edit", note="manual edit")
await db.commit()
await db.refresh(profile)
return _to_profile_response(profile)
@router.put(
"/profile/me/confirm", response_model=schemas.ConfirmResponse
)
async def confirm_my_profile(
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
profile = await _latest_profile(db, user.id)
if profile is None:
raise HTTPException(status_code=404, detail="No profile for this user")
profile.locked = True
await db.commit()
return schemas.ConfirmResponse(status="locked")
# -- Phase 2: AI coach reflection loop ---------------------------------------
async def _reflection_history(
db: AsyncSession, profile_id: str
) -> list[ReflectionMessage]:
stmt = (
select(ReflectionMessage)
.where(ReflectionMessage.profile_id == profile_id)
.order_by(ReflectionMessage.sequence)
)
return list((await db.execute(stmt)).scalars().all())
def _msg_out(m: ReflectionMessage) -> schemas.ReflectionMessageOut:
return schemas.ReflectionMessageOut(
role=m.role,
content=m.content,
sequence=m.sequence,
created_at=m.created_at,
)
@router.get(
"/profile/me/reflection", response_model=schemas.ReflectionThreadResponse
)
async def get_reflection(
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""The reflection dialogue so far for the user's latest profile."""
profile = await _latest_profile(db, user.id)
if profile is None:
raise HTTPException(status_code=404, detail="No profile for this user")
history = await _reflection_history(db, profile.id)
return schemas.ReflectionThreadResponse(
messages=[_msg_out(m) for m in history],
profile=_to_profile_response(profile),
)
@router.post(
"/profile/me/reflect", response_model=schemas.ReflectTurnResponse
)
async def reflect_on_profile(
payload: schemas.ReflectRequest,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""Advance the AI-coach reflection loop by one turn.
An empty message starts the loop (the coach's opening reflection); a
non-empty message is recorded as the person's turn before the coach
replies. When the person's input implies a correction, the coach proposes
revisions which are applied to the profile (mirror, not compass) and
snapshotted. Affirming is the separate ``/confirm`` lock.
"""
profile = await _latest_profile(db, user.id)
if profile is None:
raise HTTPException(status_code=404, detail="No profile for this user")
if profile.locked:
raise HTTPException(
status_code=409,
detail="Profile is affirmed and locked; reflection is closed.",
)
history = await _reflection_history(db, profile.id)
next_seq = (history[-1].sequence + 1) if history else 0
person_text = payload.message.strip()
if person_text:
db.add(
ReflectionMessage(
id=str(uuid.uuid4()),
profile_id=profile.id,
user_id=profile.user_id,
role="person",
content=person_text,
sequence=next_seq,
created_at=_now(),
)
)
next_seq += 1
elif history:
# No new message and the loop has already opened — nothing to do.
raise HTTPException(
status_code=400, detail="Provide a message to continue reflecting."
)
coach_history = [
{"role": m.role, "content": m.content} for m in history
]
if person_text:
coach_history.append({"role": "person", "content": person_text})
api_key = os.getenv("ANTHROPIC_API_KEY")
model = os.getenv("ANTHROPIC_MODEL", "claude-sonnet-4-6")
try:
coach = ReflectionCoach(api_key=api_key, model=model)
result = await coach.reflect(
_profile_fields(profile) | {"triad": profile.triad},
coach_history,
focus=payload.focus,
)
except ReflectionError as exc:
raise HTTPException(status_code=502, detail=str(exc)) from exc
revised = False
revisions = result.get("revisions")
if revisions:
for field, value in revisions.items():
if field in EDITABLE_FIELDS:
setattr(profile, field, value)
revised = True
_record_revision(
db,
profile,
source="reflection",
note=result.get("revision_note") or "reflection revision",
)
coach_msg = ReflectionMessage(
id=str(uuid.uuid4()),
profile_id=profile.id,
user_id=profile.user_id,
role="coach",
content=result["message"],
sequence=next_seq,
created_at=_now(),
)
db.add(coach_msg)
await db.commit()
await db.refresh(profile)
await db.refresh(coach_msg)
return schemas.ReflectTurnResponse(
message=_msg_out(coach_msg),
profile=_to_profile_response(profile),
revised=revised,
revision_note=result.get("revision_note") if revised else None,
)
@router.get(
"/profile/me/revisions",
response_model=list[schemas.ProfileRevisionOut],
)
async def get_profile_revisions(
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""The profile's edit/iteration history, newest first."""
profile = await _latest_profile(db, user.id)
if profile is None:
raise HTTPException(status_code=404, detail="No profile for this user")
stmt = (
select(ProfileRevision)
.where(ProfileRevision.profile_id == profile.id)
.order_by(ProfileRevision.created_at.desc())
)
rows = (await db.execute(stmt)).scalars().all()
out = []
for r in rows:
try:
fields = json.loads(r.fields_json)
except (json.JSONDecodeError, TypeError):
fields = {}
out.append(
schemas.ProfileRevisionOut(
id=r.id,
source=r.source,
fields=fields,
note=r.note,
created_at=r.created_at,
)
)
return out
@router.get(
"/profile/me/goal-history", response_model=schemas.GoalHistoryOut
)
async def get_goal_history(
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
"""How the person's near/long-term goals evolved over time, derived from
the profile revision snapshots (Phase 5)."""
profile = await _latest_profile(db, user.id)
if profile is None:
raise HTTPException(status_code=404, detail="No profile for this user")
rev_stmt = (
select(ProfileRevision)
.where(ProfileRevision.profile_id == profile.id)
.order_by(ProfileRevision.created_at)
)
rows = (await db.execute(rev_stmt)).scalars().all()
revisions = []
for r in rows:
try:
fields = json.loads(r.fields_json)
except (json.JSONDecodeError, TypeError):
fields = {}
revisions.append(
{"fields": fields, "source": r.source, "created_at": r.created_at}
)
timelines = goal_history(revisions)
return schemas.GoalHistoryOut(
short_term_goals=[
schemas.GoalHistoryEntry(**e) for e in timelines["short_term_goals"]
],
long_term_goals=[
schemas.GoalHistoryEntry(**e) for e in timelines["long_term_goals"]
],
)
@router.get(
"/conversation/{conversation_id}",
response_model=schemas.ConversationResponse,
)
async def get_conversation(
conversation_id: str,
db: AsyncSession = Depends(get_db),
user: User = Depends(get_current_user),
):
conversation = await _owned_conversation(db, conversation_id, user)
return schemas.ConversationResponse.model_validate(conversation)
def _as_int(value) -> int | None:
"""Coerce the model's numeric fields to int, tolerating strings/None."""
if value is None:
return None
try:
return int(value)
except (TypeError, ValueError):
return None