agent-tina / storage.py
KPrashanth's picture
Deploy Agent Tina
9d7051c verified
Raw
History Blame Contribute Delete
7.48 kB
from __future__ import annotations
import re
import secrets
import shutil
from pathlib import Path
from threading import RLock
from fastapi import HTTPException
from huggingface_hub import HfApi, hf_hub_download
from config import HF_DATASET_REPO, HF_TOKEN, MEETINGS_DIR
from models import Meeting, MeetingCreate, Segment, utc_now_iso
STORAGE_LOCK = RLock()
MEETING_ID_PATTERN = re.compile(r"^[A-Za-z0-9_-]{1,80}$")
def ensure_storage() -> None:
MEETINGS_DIR.mkdir(parents=True, exist_ok=True)
def slugify(value: str) -> str:
slug = re.sub(r"[^a-zA-Z0-9]+", "-", value.strip().lower()).strip("-")
return slug[:32] or "meeting"
def meeting_path(meeting_id: str) -> Path:
return meeting_dir(meeting_id) / "metadata.json"
def meeting_dir(meeting_id: str) -> Path:
if not MEETING_ID_PATTERN.fullmatch(meeting_id):
raise HTTPException(status_code=404, detail="Meeting not found.")
return MEETINGS_DIR / meeting_id
def create_meeting_id(title: str) -> str:
ensure_storage()
base = slugify(title)
for _ in range(20):
suffix = secrets.token_urlsafe(6)
meeting_id = f"{base}-{suffix}"
if not meeting_path(meeting_id).exists():
return meeting_id
return f"{base}-{secrets.token_hex(4)}"
def create_meeting(payload: MeetingCreate) -> Meeting:
ensure_storage()
meeting = Meeting(
meeting_id=create_meeting_id(payload.title),
host_key=secrets.token_urlsafe(18),
title=payload.title.strip(),
meeting_type=payload.meeting_type,
context=payload.context.strip(),
known_terms=payload.known_terms.strip(),
)
save_meeting(meeting)
return meeting
def load_meeting(meeting_id: str) -> Meeting:
path = meeting_path(meeting_id)
with STORAGE_LOCK:
if not path.exists():
restore_meeting_from_dataset(meeting_id)
if not path.exists():
raise HTTPException(status_code=404, detail="Meeting not found.")
return Meeting.model_validate_json(path.read_text(encoding="utf-8"))
def save_meeting(meeting: Meeting) -> None:
with STORAGE_LOCK:
write_meeting(meeting)
sync_meeting_to_dataset(meeting.meeting_id)
def write_meeting(meeting: Meeting) -> None:
ensure_storage()
directory = meeting_dir(meeting.meeting_id)
directory.mkdir(parents=True, exist_ok=True)
meeting.updated_at = utc_now_iso()
path = meeting_path(meeting.meeting_id)
temp_path = path.with_suffix(".json.tmp")
temp_path.write_text(
meeting.model_dump_json(indent=2),
encoding="utf-8",
)
temp_path.replace(path)
def add_segment(meeting: Meeting, segment: Segment) -> Meeting:
with STORAGE_LOCK:
latest = load_meeting(meeting.meeting_id)
if not any(item.segment_id == segment.segment_id for item in latest.segments):
latest.segments.append(segment)
write_meeting(latest)
sync_meeting_to_dataset(meeting.meeting_id)
return latest
def update_segment(meeting_id: str, segment: Segment) -> Meeting:
with STORAGE_LOCK:
meeting = load_meeting(meeting_id)
for index, existing in enumerate(meeting.segments):
if existing.segment_id == segment.segment_id:
meeting.segments[index] = segment
write_meeting(meeting)
break
else:
meeting.segments.append(segment)
write_meeting(meeting)
sync_meeting_to_dataset(meeting_id)
return meeting
def save_segment_audio(meeting_id: str, segment_id: str, source_path: str, suffix: str) -> Path:
recordings_dir = meeting_dir(meeting_id) / "recordings"
recordings_dir.mkdir(parents=True, exist_ok=True)
safe_suffix = suffix.lower() if suffix and re.fullmatch(r"\.[a-zA-Z0-9]{1,10}", suffix) else ".audio"
destination = recordings_dir / f"{segment_id}{safe_suffix}"
shutil.copyfile(source_path, destination)
return destination
def save_outputs(meeting: Meeting) -> None:
directory = meeting_dir(meeting.meeting_id)
directory.mkdir(parents=True, exist_ok=True)
if meeting.outputs.raw_transcript_md:
(directory / "raw_transcript.md").write_text(meeting.outputs.raw_transcript_md, encoding="utf-8")
if meeting.outputs.corrected_transcript_md:
(directory / "corrected_transcript.md").write_text(meeting.outputs.corrected_transcript_md, encoding="utf-8")
if meeting.outputs.minutes_md:
(directory / "minutes.md").write_text(meeting.outputs.minutes_md, encoding="utf-8")
sync_meeting_to_dataset(meeting.meeting_id)
def require_host(meeting: Meeting, host_key: str) -> None:
if not host_key or not secrets.compare_digest(host_key, meeting.host_key):
raise HTTPException(status_code=403, detail="Invalid host key.")
def safe_download_path(meeting_id: str, file_name: str) -> Path:
allowed = {"raw_transcript.md", "corrected_transcript.md", "minutes.md", "metadata.json"}
if file_name not in allowed:
raise HTTPException(status_code=404, detail="File not found.")
path = meeting_dir(meeting_id) / file_name
if not path.exists():
restore_dataset_file(meeting_id, file_name)
if not path.exists():
raise HTTPException(status_code=404, detail="File not generated yet.")
return path
def safe_recording_path(meeting_id: str, file_name: str) -> Path:
if Path(file_name).name != file_name:
raise HTTPException(status_code=404, detail="Recording not found.")
path = meeting_dir(meeting_id) / "recordings" / file_name
if not path.exists():
restore_dataset_file(meeting_id, f"recordings/{file_name}")
if not path.exists() or not path.is_file():
raise HTTPException(status_code=404, detail="Recording not found.")
return path
def sync_meeting_to_dataset(meeting_id: str) -> None:
if not HF_TOKEN or not HF_DATASET_REPO:
return
directory = meeting_dir(meeting_id)
if not directory.exists():
return
try:
api = HfApi(token=HF_TOKEN)
api.create_repo(repo_id=HF_DATASET_REPO, repo_type="dataset", private=True, exist_ok=True)
api.upload_folder(
repo_id=HF_DATASET_REPO,
repo_type="dataset",
folder_path=str(directory),
path_in_repo=meeting_id,
commit_message=f"Update Agent Tina meeting {meeting_id}",
)
except Exception as exc:
print(f"Warning: failed to sync meeting {meeting_id} to dataset: {exc}")
def restore_meeting_from_dataset(meeting_id: str) -> None:
restore_dataset_file(meeting_id, "metadata.json")
def restore_dataset_file(meeting_id: str, relative_path: str) -> None:
if not HF_TOKEN or not HF_DATASET_REPO:
return
directory = meeting_dir(meeting_id)
relative = Path(relative_path)
if not relative.parts or ".." in relative.parts or relative.is_absolute():
return
repo_path = f"{meeting_id}/{relative.as_posix()}"
try:
source = Path(
hf_hub_download(
repo_id=HF_DATASET_REPO,
filename=repo_path,
repo_type="dataset",
token=HF_TOKEN,
)
)
destination = directory / relative
destination.parent.mkdir(parents=True, exist_ok=True)
shutil.copyfile(source, destination)
except Exception as exc:
print(f"Warning: failed to restore {repo_path} from dataset: {exc}")