feat(transcription): add Supabase schema and API endpoints for CIS

This commit is contained in:
2026-05-20 21:03:00 +00:00
parent d68b63cedb
commit a746aed937
10 changed files with 818 additions and 17 deletions
+11 -1
View File
@@ -161,4 +161,14 @@ async def process_worker_timetable(file_content, user_node_data, worker_node_dat
raise
finally:
logging.info(f"Closing driver for {worker_node_data['worker_db_name']}")
driver.close_driver(neo_driver)
driver.close_driver(neo_driver)
@router.get("/current-period")
async def get_current_period(user_id: str = ""):
# Phase 1: return stub — TODO: implement Neo4j query in Phase 2
return {
"period_id": None,
"event_type": None,
"event_label": None,
"start_time": None,
"end_time": None,
}
+63
View File
@@ -0,0 +1,63 @@
"""Canvas events router — batch write and query canvas event logs."""
from fastapi import APIRouter, Depends, HTTPException, Query
from typing import List, Optional
from datetime import datetime
from modules.auth.supabase_bearer import SupabaseBearer
from modules.transcription.models import CanvasEventCreate, CanvasEventResponse
router = APIRouter()
def get_supabase_client():
"""Get Supabase service role client."""
from modules.database.supabase.utils.client import SupabaseServiceRoleClient
return SupabaseServiceRoleClient()
def get_user_id(credentials=Depends(SupabaseBearer())) -> str:
"""Extract user_id from Supabase JWT token."""
return credentials.get("sub", credentials.get("user_id", ""))
@router.post("/canvas-events")
async def batch_write_canvas_events(
events: List[CanvasEventCreate],
user_id: str = Depends(get_user_id),
):
"""Batch write canvas events."""
supabase = get_supabase_client()
# Filter events to only this user's
user_events = [e for e in events if e.user_id == user_id]
if not user_events:
return {"message": "No events to write", "count": 0}
event_data = [e.model_dump() for e in user_events]
result = supabase.supabase.table("canvas_events").insert(event_data).execute()
return {"message": f"Wrote {len(event_data)} events", "count": len(event_data)}
@router.get("/canvas-events", response_model=List[CanvasEventResponse])
async def get_canvas_events(
session_id: Optional[str] = Query(None),
user_id: str = Depends(get_user_id),
limit: int = Query(100, ge=1, le=1000),
):
"""Get canvas events for a session or user."""
supabase = get_supabase_client()
query = supabase.supabase.table("canvas_events").select("*").eq("user_id", user_id)
if session_id:
query = query.eq("session_id", session_id)
query = query.order("timestamp", desc=True).limit(limit)
result = query.execute()
return result.data
+102
View File
@@ -0,0 +1,102 @@
"""Keyword watches router — CRUD for keyword watch rules and events."""
from fastapi import APIRouter, Depends, HTTPException
from typing import List
from modules.auth.supabase_bearer import SupabaseBearer
from modules.transcription.models import (
KeywordWatchCreate,
KeywordWatchResponse,
KeywordEventCreate,
KeywordEventResponse,
)
router = APIRouter()
def get_supabase_client():
"""Get Supabase service role client."""
from modules.database.supabase.utils.client import SupabaseServiceRoleClient
return SupabaseServiceRoleClient()
def get_user_id(credentials=Depends(SupabaseBearer())) -> str:
"""Extract user_id from Supabase JWT token."""
return credentials.get("sub", credentials.get("user_id", ""))
@router.get("/keywords", response_model=List[KeywordWatchResponse])
async def list_keyword_watches(
user_id: str = Depends(get_user_id),
):
"""List user's keyword watches."""
supabase = get_supabase_client()
result = supabase.supabase.table("keyword_watches").select("*").eq("user_id", user_id).execute()
return result.data
@router.post("/keywords", response_model=KeywordWatchResponse)
async def create_keyword_watch(
watch: KeywordWatchCreate,
user_id: str = Depends(get_user_id),
):
"""Create a keyword watch."""
supabase = get_supabase_client()
data = {
"user_id": user_id,
"keyword": watch.keyword.lower(), # Store lowercase for case-insensitive matching
"match_type": watch.match_type,
"action": watch.action,
}
result = supabase.supabase.table("keyword_watches").insert(data).execute()
if not result.data:
raise HTTPException(status_code=500, detail="Failed to create keyword watch")
return result.data[0]
@router.delete("/keywords/{watch_id}")
async def delete_keyword_watch(
watch_id: str,
user_id: str = Depends(get_user_id),
):
"""Delete a keyword watch."""
supabase = get_supabase_client()
# Verify ownership
existing = supabase.supabase.table("keyword_watches").select("*").eq("id", watch_id).eq("user_id", user_id).execute()
if not existing.data:
raise HTTPException(status_code=404, detail="Keyword watch not found")
supabase.supabase.table("keyword_watches").delete().eq("id", watch_id).execute()
return {"message": "Keyword watch deleted"}
@router.post("/keywords/events")
async def log_keyword_event(
event: KeywordEventCreate,
user_id: str = Depends(get_user_id),
):
"""Log a keyword event (triggered when a watch matches)."""
supabase = get_supabase_client()
# Verify session ownership
session_check = supabase.supabase.table("transcription_sessions").select("id").eq("id", event.session_id).eq("user_id", user_id).execute()
if not session_check.data:
raise HTTPException(status_code=404, detail="Session not found")
data = event.model_dump()
result = supabase.supabase.table("keyword_events").insert(data).execute()
if not result.data:
raise HTTPException(status_code=500, detail="Failed to log keyword event")
return result.data[0]
+293
View File
@@ -0,0 +1,293 @@
"""Transcription sessions router — CRUD endpoints for transcription sessions and segments."""
from fastapi import APIRouter, Depends, HTTPException, Query
from typing import Optional, List
from datetime import datetime
from modules.auth.supabase_bearer import SupabaseBearer
from modules.transcription.models import (
TranscriptionSessionCreate,
TranscriptionSessionUpdate,
TranscriptionSessionResponse,
SessionListResponse,
TranscriptionSegmentCreate,
TranscriptionSegmentResponse,
SummaryGenerateRequest,
SummaryResponse,
ExportFormat,
)
router = APIRouter()
def get_supabase_client():
"""Get Supabase service role client."""
from modules.database.supabase.utils.client import SupabaseServiceRoleClient
return SupabaseServiceRoleClient()
def get_user_id(credentials=Depends(SupabaseBearer())) -> str:
"""Extract user_id from Supabase JWT token."""
return credentials.get("sub", credentials.get("user_id", ""))
@router.post("/sessions", response_model=TranscriptionSessionResponse)
async def create_session(
session_data: TranscriptionSessionCreate,
user_id: str = Depends(get_user_id),
):
"""Create a new transcription session."""
supabase = get_supabase_client()
data = {
"user_id": user_id,
"title": session_data.title,
"canvas_type": session_data.canvas_type,
}
result = supabase.supabase.table("transcription_sessions").insert(data).execute()
if not result.data:
raise HTTPException(status_code=500, detail="Failed to create session")
return result.data[0]
@router.patch("/sessions/{session_id}", response_model=TranscriptionSessionResponse)
async def update_session(
session_id: str,
update_data: TranscriptionSessionUpdate,
user_id: str = Depends(get_user_id),
):
"""Update a transcription session (end, tag, title)."""
supabase = get_supabase_client()
# Verify ownership
existing = supabase.supabase.table("transcription_sessions").select("*").eq("id", session_id).eq("user_id", user_id).execute()
if not existing.data:
raise HTTPException(status_code=404, detail="Session not found")
# Build update dict (only non-None fields)
updates = {k: v for k, v in update_data.model_dump().items() if v is not None}
updates["updated_at"] = datetime.utcnow().isoformat()
result = supabase.supabase.table("transcription_sessions").update(updates).eq("id", session_id).execute()
if not result.data:
raise HTTPException(status_code=500, detail="Failed to update session")
return result.data[0]
@router.get("/sessions", response_model=SessionListResponse)
async def list_sessions(
user_id: str = Depends(get_user_id),
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
timetable_period_id: Optional[str] = None,
):
"""List transcription sessions for the current user (paginated)."""
supabase = get_supabase_client()
query = supabase.supabase.table("transcription_sessions").select("*", count="exact").eq("user_id", user_id)
if timetable_period_id:
query = query.eq("timetable_period_id", timetable_period_id)
query = query.order("started_at", desc=True).range((page - 1) * page_size, page * page_size - 1)
result = query.execute()
return SessionListResponse(
sessions=result.data,
total=result.count or 0,
page=page,
page_size=page_size,
)
@router.get("/sessions/{session_id}", response_model=dict)
async def get_session(
session_id: str,
user_id: str = Depends(get_user_id),
):
"""Get a session with its segments and summaries."""
supabase = get_supabase_client()
# Get session
session_result = supabase.supabase.table("transcription_sessions").select("*").eq("id", session_id).eq("user_id", user_id).execute()
if not session_result.data:
raise HTTPException(status_code=404, detail="Session not found")
# Get segments
segments_result = supabase.supabase.table("transcription_segments").select("*").eq("session_id", session_id).order("sequence_index").execute()
# Get summaries
summaries_result = supabase.supabase.table("transcription_summaries").select("*").eq("session_id", session_id).execute()
return {
"session": session_result.data[0],
"segments": segments_result.data,
"summaries": summaries_result.data,
}
@router.delete("/sessions/{session_id}")
async def delete_session(
session_id: str,
user_id: str = Depends(get_user_id),
):
"""Soft delete a transcription session."""
supabase = get_supabase_client()
# Verify ownership
existing = supabase.supabase.table("transcription_sessions").select("*").eq("id", session_id).eq("user_id", user_id).execute()
if not existing.data:
raise HTTPException(status_code=404, detail="Session not found")
# Soft delete: set ended_at and mark metadata
result = supabase.supabase.table("transcription_sessions").update({
"ended_at": datetime.utcnow().isoformat(),
"metadata": {"deleted": True},
}).eq("id", session_id).execute()
return {"message": "Session deleted"}
@router.post("/sessions/{session_id}/segments")
async def upsert_segments(
session_id: str,
segments: List[TranscriptionSegmentCreate],
user_id: str = Depends(get_user_id),
):
"""Batch upsert segments for a session."""
supabase = get_supabase_client()
# Verify session exists and user owns it
session_check = supabase.supabase.table("transcription_sessions").select("id").eq("id", session_id).eq("user_id", user_id).execute()
if not session_check.data:
raise HTTPException(status_code=404, detail="Session not found")
# Batch insert segments
segment_data = [s.model_dump() for s in segments]
if segment_data:
result = supabase.supabase.table("transcription_segments").insert(segment_data).execute()
# Update segment count on session
supabase.supabase.table("transcription_sessions").update({
"segment_count": len(segment_data),
}).eq("id", session_id).execute()
return {"message": f"Upserted {len(segment_data)} segments", "count": len(segment_data)}
@router.get("/sessions/{session_id}/segments", response_model=List[TranscriptionSegmentResponse])
async def list_segments(
session_id: str,
user_id: str = Depends(get_user_id),
):
"""List all segments for a session."""
supabase = get_supabase_client()
# Verify ownership
session_check = supabase.supabase.table("transcription_sessions").select("id").eq("id", session_id).eq("user_id", user_id).execute()
if not session_check.data:
raise HTTPException(status_code=404, detail="Session not found")
result = supabase.supabase.table("transcription_segments").select("*").eq("session_id", session_id).order("sequence_index").execute()
return result.data
@router.post("/sessions/{session_id}/summaries", response_model=SummaryResponse)
async def generate_summary(
session_id: str,
summary_request: SummaryGenerateRequest,
user_id: str = Depends(get_user_id),
):
"""Generate a summary for a session (Phase 1 stub)."""
supabase = get_supabase_client()
# Verify session exists and user owns it
session_check = supabase.supabase.table("transcription_sessions").select("id").eq("id", session_id).eq("user_id", user_id).execute()
if not session_check.data:
raise HTTPException(status_code=404, detail="Session not found")
# Phase 1 stub: TODO implement LLM call in Phase 3
content = "[TODO: Generate summary via LLM — provider={}, model={}]".format(
summary_request.provider, summary_request.model
)
# Save summary to database
summary_data = {
"session_id": session_id,
"user_id": user_id,
"summary_type": summary_request.summary_type,
"content": content,
"llm_provider": summary_request.provider,
"llm_model": summary_request.model,
}
result = supabase.supabase.table("transcription_summaries").insert(summary_data).execute()
if not result.data:
raise HTTPException(status_code=500, detail="Failed to save summary")
return result.data[0]
@router.get("/sessions/{session_id}/summaries", response_model=List[SummaryResponse])
async def list_summaries(
session_id: str,
user_id: str = Depends(get_user_id),
):
"""List summaries for a session."""
supabase = get_supabase_client()
# Verify ownership
session_check = supabase.supabase.table("transcription_sessions").select("id").eq("id", session_id).eq("user_id", user_id).execute()
if not session_check.data:
raise HTTPException(status_code=404, detail="Session not found")
result = supabase.supabase.table("transcription_summaries").select("*").eq("session_id", session_id).execute()
return result.data
@router.post("/sessions/{session_id}/export")
async def export_session(
session_id: str,
export_format: ExportFormat,
user_id: str = Depends(get_user_id),
):
"""Export session as SRT, TXT, or JSON (Phase 1 stub)."""
supabase = get_supabase_client()
# Verify ownership
session_check = supabase.supabase.table("transcription_sessions").select("id").eq("id", session_id).eq("user_id", user_id).execute()
if not session_check.data:
raise HTTPException(status_code=404, detail="Session not found")
# Get segments
segments_result = supabase.supabase.table("transcription_segments").select("*").eq("session_id", session_id).order("sequence_index").execute()
segments = segments_result.data
if export_format.format == "srt":
# Phase 1 stub — implement in Phase 3
return {"format": "srt", "content": "[TODO: Generate SRT from segments]"}
elif export_format.format == "txt":
text = "\n".join(s["text"] for s in segments)
return {"format": "txt", "content": text}
elif export_format.format == "json":
return {"format": "json", "content": {"segments": segments}}
else:
raise HTTPException(status_code=400, detail=f"Unsupported format: {export_format.format}")
+28 -16
View File
@@ -1,24 +1,31 @@
from fastapi import APIRouter, Request
"""DEPRECATED - utterance.py
This module is deprecated. All transcription data now goes through
the new sessions router at /transcribe/sessions.
The old endpoints are kept for backwards compatibility but will be
removed in a future release. New integrations should use:
POST /transcribe/sessions - create session
POST /transcribe/sessions/{id}/segments - batch upsert segments
"""
from fastapi import APIRouter, Request, HTTPException
import os
import queue
from dotenv import load_dotenv
import json
load_dotenv()
router = APIRouter(deprecated=True)
router = APIRouter()
@router.post("/handle_whisper_live_eos_utterance/{user_id}")
async def handle_whisper_live_eos_utterance(user_id: str, request: Request):
"""DEPRECATED - Use POST /transcribe/sessions instead."""
data = await request.json()
utterance = data.get("utterance")
print(f"Utterance: {utterance}")
start = data.get("start")
end = data.get("end")
print(f"Start: {start}")
print(f"End: {end}")
eos = data.get("eos")
print(f"Eos: {eos}")
# Log deprecation warning
print(f"[DEPRECATED] /handle_whisper_live_eos_utterance called for user {user_id}. "
f"Please migrate to POST /transcribe/sessions/{{id}}/segments")
# Keep writing to flat file for backwards compat
user_dir = f"../../data/users/{user_id}/transcripts"
if not os.path.exists(user_dir):
os.makedirs(user_dir)
@@ -27,16 +34,21 @@ async def handle_whisper_live_eos_utterance(user_id: str, request: Request):
with open(log_file, "a") as f:
f.write(json.dumps(data) + "\n")
return {"message": "Utterance logged successfully"}
return {"message": "Utterance logged (deprecated - migrate to sessions API)", "_deprecated": True}
@router.get("/get_utterances/{user_id}")
async def get_utterances(user_id: str):
"""DEPRECATED - Use GET /transcribe/sessions/{id}/segments instead."""
print(f"[DEPRECATED] /get_utterances called for user {user_id}. "
f"Please migrate to GET /transcribe/sessions/{{id}}/segments")
user_dir = f"../../data/users/{user_id}/transcripts"
log_file = os.path.join(user_dir, "utterances.log")
if not os.path.exists(log_file):
return {"utterances": []}
return {"utterances": [], "_deprecated": True}
with open(log_file, "r") as f:
utterances = [json.loads(line) for line in f]
return {"utterances": utterances}
return {"utterances": utterances, "_deprecated": True}