Add Streamlit meeting assistant MVP

This commit is contained in:
2026-08-24 10:28:44 +02:00
parent 443a85528b
commit 13d06b8a60
19 changed files with 1188 additions and 20 deletions
+17
View File
@@ -0,0 +1,17 @@
"""Application services for Meeting Assistant workflows."""
from mka.application.meeting_service import (
MeetingDetails,
MeetingProcessingService,
ParticipantInput,
ProcessingOptions,
ProcessingOutcome,
)
__all__ = [
"MeetingDetails",
"MeetingProcessingService",
"ParticipantInput",
"ProcessingOptions",
"ProcessingOutcome",
]
+69
View File
@@ -0,0 +1,69 @@
"""Environment-backed configuration for the Meeting Assistant application."""
from __future__ import annotations
import os
from dataclasses import dataclass
from pathlib import Path
class ConfigurationError(ValueError):
"""Raised when required processing configuration is unavailable."""
@dataclass(frozen=True)
class AppSettings:
"""Machine-specific Meeting Lab settings kept outside the UI."""
data_root: Path
whisper_model: Path | None
whisper_executable: str = "whisper-cli"
ffmpeg_executable: str = "ffmpeg"
protocol_model: str = "qwen3.8:27B"
ollama_endpoint: str = "http://127.0.0.1:11434"
language: str = "de"
threads: str | int = "auto"
diarization_mode: str = "auto"
diarization_runtime: str = "native"
diarization_container_image: str | None = None
@classmethod
def from_environment(cls) -> AppSettings:
"""Load settings from MKA_* environment variables."""
whisper_model = os.getenv("MKA_WHISPER_MODEL")
threads = os.getenv("MKA_WHISPER_THREADS", "auto")
parsed_threads: str | int = int(threads) if threads.isdigit() else threads
return cls(
data_root=Path(os.getenv("MKA_DATA_ROOT", "data/meetings")),
whisper_model=Path(whisper_model) if whisper_model else None,
whisper_executable=os.getenv("MKA_WHISPER_EXECUTABLE", "whisper-cli"),
ffmpeg_executable=os.getenv("MKA_FFMPEG_EXECUTABLE", "ffmpeg"),
protocol_model=os.getenv("MKA_PROTOCOL_MODEL", "qwen3.8:27B"),
ollama_endpoint=os.getenv("MKA_OLLAMA_ENDPOINT", "http://127.0.0.1:11434"),
language=os.getenv("MKA_LANGUAGE", "de"),
threads=parsed_threads,
diarization_mode=os.getenv("MKA_DIARIZATION_MODE", "auto"),
diarization_runtime=os.getenv("MKA_DIARIZATION_RUNTIME", "native"),
diarization_container_image=os.getenv("MKA_DIARIZATION_CONTAINER_IMAGE"),
)
def validate_for_processing(self) -> None:
"""Validate settings needed before a Meeting Lab run starts."""
if self.whisper_model is None:
raise ConfigurationError("MKA_WHISPER_MODEL is not configured.")
if not self.whisper_model.is_file():
raise ConfigurationError(
f"Configured Whisper model does not exist: {self.whisper_model}"
)
if not self.whisper_executable.strip():
raise ConfigurationError("MKA_WHISPER_EXECUTABLE must not be empty.")
if not self.ffmpeg_executable.strip():
raise ConfigurationError("MKA_FFMPEG_EXECUTABLE must not be empty.")
if self.diarization_mode not in {"auto", "cpu", "gpu"}:
raise ConfigurationError("MKA_DIARIZATION_MODE must be one of: auto, cpu, gpu.")
if self.diarization_runtime not in {"native", "container"}:
raise ConfigurationError("MKA_DIARIZATION_RUNTIME must be native or container.")
if self.diarization_runtime == "container" and not self.diarization_container_image:
raise ConfigurationError(
"MKA_DIARIZATION_CONTAINER_IMAGE is required for container runtime."
)
+267
View File
@@ -0,0 +1,267 @@
"""Use-case service for processing one uploaded meeting recording."""
from __future__ import annotations
import json
import re
import shutil
from collections.abc import Callable
from dataclasses import dataclass
from datetime import date
from pathlib import Path
from typing import Any, Protocol
from uuid import uuid4
from mka.application.config import AppSettings
STAGES = ("preparing", "transcription", "diarization", "protocol_generation")
class MeetingLabPort(Protocol):
"""Operations the application requires from Meeting Lab."""
def create_context(self, data: dict[str, Any]) -> Any: ...
def create_config(self, values: dict[str, Any]) -> Any: ...
def run(
self,
config: Any,
meeting_context: Any,
progress_sink: Callable[[Any], None],
) -> Any: ...
@dataclass(frozen=True)
class MeetingDetails:
title: str
language: str = "de"
meeting_date: date | None = None
description: str = ""
meeting_id: str | None = None
@dataclass(frozen=True)
class ParticipantInput:
participant_id: str
display_name: str
role: str = ""
organization: str = ""
attendance_status: str = "present"
@dataclass(frozen=True)
class ProcessingOptions:
diarization_enabled: bool = False
audio_normalization: bool = True
@dataclass(frozen=True)
class AppProgressEvent:
stage: str
status: str
elapsed_seconds: float
progress: float | None = None
message: str | None = None
@dataclass(frozen=True)
class ProcessingOutcome:
succeeded: bool
run_dir: Path | None
original_protocol: str | None
protocol_path: Path | None
failed_stage: str | None = None
error_message: str | None = None
def stable_id(value: str) -> str:
"""Return a schema-safe ID based on user text, with a random fallback."""
normalized = re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-")
return normalized or f"participant-{uuid4().hex[:8]}"
class MeetingProcessingService:
"""Translate UI input into Meeting Lab calls and persisted artifacts."""
def __init__(self, settings: AppSettings, meeting_lab: MeetingLabPort) -> None:
self.settings = settings
self.meeting_lab = meeting_lab
def build_context(
self,
meeting: MeetingDetails,
participants: list[ParticipantInput],
speaker_mappings: dict[str, str] | None = None,
) -> Any:
"""Build and validate Meeting Context V1 through Meeting Lab."""
meeting_id = meeting.meeting_id or stable_id(meeting.title)
organizations = sorted(
{item.organization.strip() for item in participants if item.organization.strip()}
)
departments = [
{"id": stable_id(name), "name": name, "aliases": []} for name in organizations
]
department_ids = {item["name"]: item["id"] for item in departments}
data = {
"schema_version": "1",
"meeting": {
"meeting_id": meeting_id,
"title": meeting.title.strip(),
"language": meeting.language,
"date": meeting.meeting_date.isoformat() if meeting.meeting_date else None,
"objective": "",
"notes": meeting.description.strip(),
},
"participants": [
{
"participant_id": item.participant_id.strip(),
"display_name": item.display_name.strip(),
"aliases": [],
"role": item.role.strip() or None,
"department": department_ids.get(item.organization.strip()),
"attendance_status": "present",
"notes": None,
}
for item in participants
if item.attendance_status == "present"
],
"speaker_mappings": dict(speaker_mappings or {}),
"mentioned_people": [
{
"person_id": item.participant_id.strip(),
"display_name": item.display_name.strip(),
"aliases": [],
"role": item.role.strip() or None,
"department": department_ids.get(item.organization.strip()),
"attendance_status": "mentioned_only",
"notes": None,
}
for item in participants
if item.attendance_status == "mentioned_only"
],
"organization": {
"name": None,
"departments": departments,
"abbreviations": {},
},
"known_entities": {},
"context_rules": {
"participant_list_is_authoritative": True,
"do_not_infer_roles": True,
"do_not_infer_departments": True,
"do_not_infer_responsibilities": True,
"mentioned_people_are_not_participants": True,
},
}
return self.meeting_lab.create_context(data)
def preserve_upload(self, meeting_id: str, filename: str, source: Any) -> Path:
"""Persist an uploaded source before processing starts."""
safe_filename = Path(filename).name
upload_dir = self.settings.data_root / stable_id(meeting_id) / "uploads"
upload_dir.mkdir(parents=True, exist_ok=True)
destination = upload_dir / f"{uuid4().hex[:8]}_{safe_filename}"
with destination.open("wb") as target:
if hasattr(source, "getbuffer"):
target.write(source.getbuffer())
else:
shutil.copyfileobj(source, target)
return destination
def process(
self,
audio_file: Path,
meeting: MeetingDetails,
participants: list[ParticipantInput],
options: ProcessingOptions,
progress_sink: Callable[[AppProgressEvent], None] | None = None,
speaker_mappings: dict[str, str] | None = None,
) -> ProcessingOutcome:
"""Validate input, run Meeting Lab and expose a UI-oriented result."""
self.settings.validate_for_processing()
context = self.build_context(meeting, participants, speaker_mappings)
meeting_id = context.meeting_id
output_root = self.settings.data_root / stable_id(meeting_id) / "runs"
config = self.meeting_lab.create_config(
{
"audio_file": audio_file,
"whisper_model": self.settings.whisper_model,
"whisper_executable": self.settings.whisper_executable,
"ffmpeg_executable": self.settings.ffmpeg_executable,
"audio_normalization": options.audio_normalization,
"output_root": output_root,
"language": meeting.language or self.settings.language,
"threads": self.settings.threads,
"model": self.settings.protocol_model,
"ollama_endpoint": self.settings.ollama_endpoint,
"diarization": (
self.settings.diarization_mode if options.diarization_enabled else "off"
),
"diarization_runtime": self.settings.diarization_runtime,
"diarization_container_image": (self.settings.diarization_container_image),
}
)
current_stage: str | None = None
def relay(event: Any) -> None:
nonlocal current_stage
if event.stage in STAGES and event.status == "started":
current_stage = event.stage
app_event = AppProgressEvent(
stage=event.stage,
status=event.status,
elapsed_seconds=event.elapsed_seconds,
progress=event.progress,
message=event.message,
)
if progress_sink is not None:
progress_sink(app_event)
result = self.meeting_lab.run(config, context, relay)
if result.exit_code != 0:
failure = self._read_failure(result.run_dir)
backend_stage = failure.get("stage")
failed_stage = {
"validation": "preparing",
"setup": "preparing",
"audio_preparation": "preparing",
"whisper": "transcription",
"transcript_validation": "transcription",
"protocol": "protocol_generation",
}.get(backend_stage, backend_stage)
return ProcessingOutcome(
succeeded=False,
run_dir=result.run_dir,
original_protocol=None,
protocol_path=None,
failed_stage=failed_stage or current_stage or "preparing",
error_message=failure.get("message") or "Meeting Lab processing failed.",
)
protocol_path = Path(result.protocol_path)
return ProcessingOutcome(
succeeded=True,
run_dir=result.run_dir,
original_protocol=protocol_path.read_text(encoding="utf-8"),
protocol_path=protocol_path,
)
@staticmethod
def _read_failure(run_dir: Path | None) -> dict[str, str]:
if run_dir is None:
return {}
metadata_path = Path(run_dir) / "run_metadata.json"
if not metadata_path.is_file():
return {}
try:
failure = json.loads(metadata_path.read_text(encoding="utf-8")).get("failure")
except (OSError, json.JSONDecodeError):
return {}
return failure if isinstance(failure, dict) else {}
@staticmethod
def save_edited_protocol(run_dir: Path, text: str) -> Path:
"""Persist user edits beside, never over, the generated protocol."""
destination = Path(run_dir) / "protocol_edited.md"
destination.write_text(text, encoding="utf-8")
return destination
+1
View File
@@ -0,0 +1 @@
"""Adapters for external processing systems."""
+51
View File
@@ -0,0 +1,51 @@
"""Narrow adapter around the reusable Meeting Lab Python API."""
from __future__ import annotations
from collections.abc import Callable, Mapping
from typing import Any
class MeetingLabUnavailableError(RuntimeError):
"""Raised when the Meeting Lab package cannot be imported."""
class MeetingLabGateway:
"""Load and delegate to Meeting Lab without coupling Streamlit to it."""
def __init__(self) -> None:
try:
from src.meeting_lab.models.meeting_context import create_meeting_context
from src.meeting_lab.orchestration.mvp import (
MvpMeetingConfig,
run_mvp_meeting,
)
except ImportError as exc:
raise MeetingLabUnavailableError(
"Meeting Lab is unavailable. Add the Meeting Lab repository root "
"to PYTHONPATH as described in README.md."
) from exc
self._config_type = MvpMeetingConfig
self._create_context = create_meeting_context
self._run_mvp_meeting = run_mvp_meeting
def create_context(self, data: dict[str, Any]) -> Any:
"""Validate structured context through Meeting Lab's domain boundary."""
return self._create_context(data)
def create_config(self, values: Mapping[str, Any]) -> Any:
"""Create the concrete Meeting Lab run configuration."""
return self._config_type(**values)
def run(
self,
config: Any,
meeting_context: Any,
progress_sink: Callable[[Any], None],
) -> Any:
"""Run Meeting Lab synchronously and relay progress callbacks."""
return self._run_mvp_meeting(
config,
meeting_context=meeting_context,
progress_sink=progress_sink,
)
+1 -1
View File
@@ -22,4 +22,4 @@ class DomainModel(BaseModel):
id: UUID = Field(default_factory=uuid4)
created_at: datetime = Field(default_factory=utc_now)
updated_at: datetime = Field(default_factory=utc_now)
updated_at: datetime = Field(default_factory=utc_now)
+253
View File
@@ -0,0 +1,253 @@
"""Streamlit presentation layer for the first Meeting Assistant MVP."""
from __future__ import annotations
from datetime import date
from typing import Any
from uuid import uuid4
import streamlit as st
from mka.application.config import AppSettings, ConfigurationError
from mka.application.meeting_service import (
STAGES,
AppProgressEvent,
MeetingDetails,
MeetingProcessingService,
ParticipantInput,
ProcessingOptions,
stable_id,
)
from mka.integrations.meeting_lab import (
MeetingLabGateway,
MeetingLabUnavailableError,
)
STAGE_LABELS = {
"preparing": "Preparation",
"transcription": "Transcription",
"diarization": "Diarization",
"protocol_generation": "Protocol generation",
}
def _new_participant() -> dict[str, str]:
return {
"row_id": uuid4().hex,
"participant_id": "",
"display_name": "",
"role": "",
"organization": "",
"attendance_status": "present",
}
def _initialize_state() -> None:
st.session_state.setdefault("participants", [_new_participant()])
st.session_state.setdefault("outcome", None)
st.session_state.setdefault("edited_protocol", "")
def _render_participants() -> list[ParticipantInput]:
st.subheader("People")
st.caption("Record whether each relevant person attended or was only mentioned.")
rows = st.session_state.participants
remove_index: int | None = None
for index, row in enumerate(rows):
row_id = row["row_id"]
columns = st.columns([2, 2, 2, 2, 2, 0.6])
row["display_name"] = columns[0].text_input(
"Name", value=row["display_name"], key=f"name_{row_id}"
)
suggested_id = row["participant_id"] or stable_id(row["display_name"])
row["participant_id"] = columns[1].text_input(
"Participant ID", value=suggested_id, key=f"id_{row_id}"
)
row["role"] = columns[2].text_input("Role", value=row["role"], key=f"role_{row_id}")
row["organization"] = columns[3].text_input(
"Organization / department",
value=row["organization"],
key=f"organization_{row_id}",
)
row["attendance_status"] = columns[4].selectbox(
"Attendance",
options=["present", "mentioned_only"],
format_func=lambda value: {
"present": "Present / participated",
"mentioned_only": "Mentioned, but not present",
}[value],
index=0 if row["attendance_status"] == "present" else 1,
key=f"attendance_{row_id}",
)
if columns[5].button("Remove", key=f"remove_{row_id}"):
remove_index = index
if remove_index is not None:
rows.pop(remove_index)
st.rerun()
if st.button("Add person"):
rows.append(_new_participant())
st.rerun()
return [
ParticipantInput(
participant_id=row["participant_id"],
display_name=row["display_name"],
role=row["role"],
organization=row["organization"],
attendance_status=row["attendance_status"],
)
for row in rows
]
def _progress_callback(
status_box: Any,
stage_table: Any,
progress_slot: Any,
states: dict[str, str],
) -> Any:
progress_bar = None
def update(event: AppProgressEvent) -> None:
nonlocal progress_bar
if event.stage in states:
states[event.stage] = "running" if event.status == "started" else event.status
if event.stage == "failed":
running = next(
(stage for stage, status in states.items() if status == "running"),
None,
)
if running:
states[running] = "failed"
elapsed = f"{event.elapsed_seconds:.1f} s"
message = event.message or STAGE_LABELS.get(event.stage, event.stage)
status_box.info(f"{message} — elapsed {elapsed}")
stage_table.table(
[{"Stage": STAGE_LABELS[stage], "Status": states[stage]} for stage in STAGES]
)
if event.progress is not None:
if progress_bar is None:
progress_bar = progress_slot.progress(0.0)
progress_bar.progress(
min(max(event.progress, 0.0), 1.0),
text=f"{STAGE_LABELS.get(event.stage, event.stage)}: {event.progress:.0%}",
)
return update
def _render_result() -> None:
outcome = st.session_state.outcome
if outcome is None:
return
st.divider()
st.header("Protocol result")
if not outcome.succeeded:
st.error(f"Processing failed during {outcome.failed_stage}: {outcome.error_message}")
if outcome.run_dir:
st.code(str(outcome.run_dir))
st.caption("Intermediate artifacts and run metadata were preserved here.")
return
st.success("Processing completed. Review the generated protocol before use.")
st.caption(f"Run artifacts: {outcome.run_dir}")
with st.expander("Original generated protocol", expanded=False):
st.markdown(outcome.original_protocol or "")
edited = st.text_area(
"Editable protocol",
key="edited_protocol",
height=500,
help="The original protocol.md remains unchanged.",
)
if st.button("Save edited protocol", type="primary"):
path = MeetingProcessingService.save_edited_protocol(outcome.run_dir, edited)
st.success(f"Saved edited protocol to {path}")
def main() -> None:
st.set_page_config(page_title="Meeting Assistant", layout="wide")
_initialize_state()
st.title("Meeting Assistant")
st.caption("Create meeting context, run Meeting Lab, and review the protocol.")
st.header("Meeting and audio")
audio = st.file_uploader("Audio recording", type=["wav", "flac", "m4a"])
title = st.text_input("Meeting title")
description = st.text_area("Description / context", height=100)
metadata_columns = st.columns(3)
language = metadata_columns[0].selectbox("Meeting language", options=["de", "en"])
has_date = metadata_columns[1].checkbox("Meeting date is known", value=True)
selected_date: date | None = (
metadata_columns[2].date_input("Meeting date") if has_date else None
)
participants = _render_participants()
st.header("Processing options")
audio_normalization = st.checkbox(
"Audio normalization",
value=True,
help=(
"Normalize speech loudness during preparation. WAV, FLAC, and M4A "
"are always converted to the canonical Meeting Lab audio format, "
"regardless of this setting."
),
)
diarization_enabled = st.checkbox(
"Enable speaker diarization",
help="Speaker labels remain anonymous; identities are never inferred.",
)
if st.button("Start processing", type="primary", disabled=audio is None):
if not title.strip():
st.error("Meeting title is required.")
elif any(
not item.display_name.strip() or not item.participant_id.strip()
for item in participants
):
st.error("Every person row needs a name and participant ID.")
else:
settings = AppSettings.from_environment()
try:
service = MeetingProcessingService(settings, MeetingLabGateway())
meeting = MeetingDetails(
title=title,
language=language,
meeting_date=selected_date,
description=description,
)
meeting_id = stable_id(title)
audio_path = service.preserve_upload(meeting_id, audio.name, audio)
st.header("Processing status")
status_box = st.empty()
stage_table = st.empty()
progress_slot = st.empty()
states = {stage: "pending" for stage in STAGES}
if not diarization_enabled:
states["diarization"] = "skipped"
callback = _progress_callback(status_box, stage_table, progress_slot, states)
outcome = service.process(
audio_path,
meeting,
participants,
ProcessingOptions(
diarization_enabled=diarization_enabled,
audio_normalization=audio_normalization,
),
progress_sink=callback,
)
st.session_state.outcome = outcome
st.session_state.edited_protocol = outcome.original_protocol or ""
if outcome.succeeded:
status_box.success("Processing completed.")
else:
status_box.error("Processing failed. Artifacts were preserved.")
except (ConfigurationError, MeetingLabUnavailableError, ValueError) as exc:
st.error(str(exc))
except Exception as exc:
st.error(f"Unable to process meeting: {exc}")
_render_result()
if __name__ == "__main__":
main()