Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2215414282 | ||
|
|
bd48aec5aa | ||
|
|
b50b87c32b | ||
|
|
9e7173c2a7 | ||
|
|
70cb1e0789 | ||
|
|
500d2c4c67 | ||
|
|
77be2b39ff | ||
|
|
e2a0434941 | ||
|
|
9840953687 | ||
|
|
e55dd74c19 | ||
|
|
dd8a618719 | ||
|
|
adf454d77d | ||
|
|
303563abf1 | ||
|
|
2a4481f2f9 | ||
|
|
13d06b8a60 | ||
|
|
443a85528b |
@@ -0,0 +1,18 @@
|
||||
MKA_WHISPER_MODEL=/path/to/ggml-large-v3-turbo.bin
|
||||
MKA_WHISPER_EXECUTABLE=whisper-cli
|
||||
MKA_FFMPEG_EXECUTABLE=ffmpeg
|
||||
MKA_PROTOCOL_MODEL=qwen3.8:27b
|
||||
MKA_OLLAMA_ENDPOINT=http://127.0.0.1:11434
|
||||
MKA_DATA_ROOT=data/meetings
|
||||
MKA_GLOSSARY_DATABASE=data/database/glossary.sqlite3
|
||||
MKA_WHISPER_THREADS=auto
|
||||
MKA_DIARIZATION_MODE=auto
|
||||
MKA_DIARIZATION_RUNTIME=native
|
||||
# MKA_DIARIZATION_CONTAINER_IMAGE=
|
||||
# MKA_DIARIZATION_CONTAINER_ARGS=[]
|
||||
|
||||
# Example AMD ROCm container configuration for a compatible workstation:
|
||||
# MKA_DIARIZATION_MODE=gpu
|
||||
# MKA_DIARIZATION_RUNTIME=container
|
||||
# MKA_DIARIZATION_CONTAINER_IMAGE=rocm/pytorch:rocm7.2.1_ubuntu24.04_py3.12_pytorch_release_2.9.1
|
||||
# MKA_DIARIZATION_CONTAINER_ARGS='["--device=/dev/kfd","--device=/dev/dri","--group-add","video"]'
|
||||
|
||||
@@ -34,6 +34,7 @@ dist/
|
||||
*.sqlite3
|
||||
|
||||
# Meeting data
|
||||
data/meetings/*
|
||||
data/recordings/*
|
||||
data/transcripts/*
|
||||
data/database/*
|
||||
@@ -42,3 +43,4 @@ data/database/*
|
||||
!data/recordings/.gitkeep
|
||||
!data/transcripts/.gitkeep
|
||||
!data/database/.gitkeep
|
||||
!data/meetings/.gitkeep
|
||||
|
||||
+21
-4
@@ -2,18 +2,35 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Added
|
||||
|
||||
- Retained live stage and total timings for processing and protocol regeneration.
|
||||
- Immutable diagnostic generation history with atomic latest-result publication.
|
||||
|
||||
- Auto/Fast/Efficient/Powersave protocol profiles with backend-selected threads
|
||||
by default, also applied during mapped-speaker regeneration.
|
||||
|
||||
- Versioned UTF-8 YAML import/export for the SQLite terminology glossary, with
|
||||
stable entry IDs, complete validation, and atomic replacement semantics.
|
||||
- First Streamlit MVP for audio upload, Meeting Context entry, participant
|
||||
management, processing progress and protocol editing.
|
||||
- Application service and Meeting Lab adapter with environment-based runtime
|
||||
configuration.
|
||||
- Per-meeting upload and edited-protocol persistence.
|
||||
- Unit tests for context construction, configuration translation, progress,
|
||||
failure and result handling.
|
||||
- ADR 0011 for the Meeting Lab MVP backend architecture.
|
||||
|
||||
### Changed
|
||||
|
||||
- Record active glossary configuration without rewriting protocol transcript input.
|
||||
|
||||
- Reconciled the documented MVP with the validated Meeting Lab backend.
|
||||
- Selected `whisper.cpp` for transcription and made diarization optional.
|
||||
- Defined the Python API, progress-event and Meeting Context integration
|
||||
boundaries.
|
||||
- Set the user-facing Meeting Assistant GUI as the next milestone.
|
||||
|
||||
### Added
|
||||
|
||||
- ADR 0011 for the Meeting Lab MVP backend architecture.
|
||||
|
||||
## [0.1.0] - 2026-07-14
|
||||
|
||||
### Added
|
||||
|
||||
+80
-7
@@ -1,5 +1,14 @@
|
||||
# Project Knowledge
|
||||
|
||||
## Meeting language
|
||||
|
||||
The GUI's `meeting_language` is passed through `MeetingDetails.language` to both
|
||||
Meeting Context `meeting.language` and `MvpMeetingConfig.language` (Whisper).
|
||||
Meeting Lab derives explicit protocol output language from the persisted context,
|
||||
also during regeneration, and records derived `output_language` provenance in
|
||||
protocol runtime metadata. Missing context language defaults to German. There is
|
||||
no separate protocol-language setting or automatic context/transcript translation.
|
||||
|
||||
## Vision
|
||||
|
||||
Meeting Assistant turns recorded meetings into reviewed, user-facing protocols
|
||||
@@ -15,7 +24,7 @@ are derived, reproducible artifacts and require human review.
|
||||
|
||||
Meeting Lab owns reusable and experimental processing:
|
||||
|
||||
- FFmpeg audio preparation
|
||||
- FFmpeg audio preparation, with optional normalization enabled by default
|
||||
- `whisper.cpp` transcription
|
||||
- optional `pyannote.audio` diarization
|
||||
- protocol generation
|
||||
@@ -36,21 +45,44 @@ Meeting Assistant owns user interaction and product workflow:
|
||||
- audio-file selection
|
||||
- structured Meeting Context editing
|
||||
- participant management
|
||||
- versioned YAML import/export for reusable People lists, using replace semantics
|
||||
- versioned JSON import/export for the complete user-configurable run-input
|
||||
form; source media is represented only by optional filename metadata and must
|
||||
be selected again after import
|
||||
- optional explicit speaker mapping
|
||||
- post-diarization speaker review and protocol-only regeneration from existing
|
||||
run artifacts
|
||||
- pipeline launch and progress display
|
||||
- protocol review and editing
|
||||
- a global SQLite terminology glossary whose active canonical core terms and
|
||||
recognition aliases are rendered through Meeting Context into direct
|
||||
protocol prompts
|
||||
- versioned UTF-8 YAML glossary interchange (`version: 1`) that preserves entry
|
||||
IDs and atomically replaces SQLite state only after complete validation
|
||||
- export and presentation of protocol versions
|
||||
|
||||
Processing logic must not be duplicated in the application.
|
||||
|
||||
Diarization runtime, container image and ordered container arguments are
|
||||
machine configuration. `MKA_DIARIZATION_CONTAINER_ARGS` is a JSON array of
|
||||
strings forwarded unchanged to Meeting Lab; these details are not normal UI
|
||||
controls.
|
||||
|
||||
People-list YAML is a Meeting Assistant application concern and contains only
|
||||
stable IDs, display names, roles, organizations and attendance states. It does
|
||||
not contain meeting metadata or processing settings. Named team or meeting
|
||||
templates may build on this later, but are not part of the current mechanism.
|
||||
|
||||
## Validated MVP Pipeline
|
||||
|
||||
```text
|
||||
Imported audio
|
||||
-> prepare as mono, 16 kHz PCM WAV with FFmpeg
|
||||
-> always prepare as mono, 16 kHz PCM WAV with FFmpeg
|
||||
(optionally normalize loudness; default on)
|
||||
-> transcribe with whisper.cpp and large-v3-turbo
|
||||
-> optionally diarize with pyannote.audio Community-1
|
||||
-> generate a protocol directly from the full transcript
|
||||
-> when diarized, wait for optional speaker review before explicit protocol generation
|
||||
-> otherwise generate a protocol directly from the full transcript
|
||||
-> review and edit by a human
|
||||
```
|
||||
|
||||
@@ -64,17 +96,37 @@ PyTorch/ROCm processed the same duration in approximately 98.6 seconds. These
|
||||
are validation observations, not performance guarantees or hardware
|
||||
requirements. CPU execution remains supported and may be substantially slower.
|
||||
|
||||
Protocol-generation performance is selected through intentionally abstract
|
||||
profiles rather than hardware controls in the GUI. `auto` is the default and
|
||||
delegates thread selection to the inference backend; with Ollama this means
|
||||
omitting `num_thread` entirely. The explicit resource profiles currently map
|
||||
`fast` to 16 Ollama CPU threads, `efficient` to 10, and `powersave` to 4. These
|
||||
concrete mappings belong to the backend/configuration boundary and may evolve
|
||||
independently of the user-facing profile semantics.
|
||||
|
||||
## Meeting Context and Speakers
|
||||
|
||||
`MeetingContext` is a structured domain object containing meeting metadata and
|
||||
participants. It can also contain optional explicit mappings from
|
||||
`MeetingContext` is a structured domain object containing meeting metadata,
|
||||
present participants, and people who were mentioned without attending. The GUI
|
||||
records exactly `present` or `mentioned_only`; missing legacy participant status
|
||||
defaults to `present` in Meeting Lab. It can also contain optional explicit mappings from
|
||||
`SPEAKER_XX` labels to `participant_id` values.
|
||||
|
||||
A speaker mapping is authoritative only when a user explicitly confirms it.
|
||||
A speaker mapping is authoritative only when a user explicitly confirms it and
|
||||
may reference only a present participant.
|
||||
Speakers otherwise remain anonymous. Automatic speaker-name inference is not
|
||||
allowed. The GUI must create and edit Meeting Context; hand-written YAML is not
|
||||
a product requirement.
|
||||
|
||||
Speaker selectors filter participants using a snapshot of all current widget
|
||||
selections, falling back to saved mappings before first interaction. Filtering
|
||||
retains existing assignments; clearing a selection releases the participant.
|
||||
Completeness counts and warnings cover detected speakers only and do not gate
|
||||
protocol generation or change mapping persistence. A diarized run persists the
|
||||
`awaiting_speaker_review` checkpoint after diarization; an empty, partial or
|
||||
complete mapping can then explicitly generate a protocol from the existing
|
||||
artifacts. Failed generation preserves that checkpoint and saved mappings for retry.
|
||||
|
||||
## Progress Contract
|
||||
|
||||
Meeting Lab emits stage-based progress events for:
|
||||
@@ -92,7 +144,7 @@ measurable progress is available.
|
||||
## Protocol Policy
|
||||
|
||||
Direct full-transcript protocol generation is the practical MVP direction.
|
||||
`qwen3.8:27B` has shown strong readability and contextual synthesis;
|
||||
`qwen3.8:27b` has shown strong readability and contextual synthesis;
|
||||
`qwen3.6:35B-A3B` has shown somewhat more conservative behavior in some areas.
|
||||
Experimental dual-model and diarization-assisted hard-fact extraction has not
|
||||
demonstrated reliably better strict attribution accuracy and is not mandatory.
|
||||
@@ -101,6 +153,11 @@ The product should ultimately offer both a detailed contextual protocol and a
|
||||
shorter participant/distribution version. The short version remains follow-up
|
||||
work if it is not available for the first GUI milestone.
|
||||
|
||||
Confirmed meeting-specific corrections should ultimately be reusable by both
|
||||
protocol views. The detailed correction and selective-regeneration decision is
|
||||
recorded in [ADR 0012](docs/adr/0012-post-run-corrections.md); it is later
|
||||
product work, not a current MVP requirement.
|
||||
|
||||
## Design Principles
|
||||
|
||||
- offline-first where practical
|
||||
@@ -111,3 +168,19 @@ work if it is not available for the first GUI milestone.
|
||||
- no automatic identity claims
|
||||
- human review of generated protocols
|
||||
- small modules and simple interfaces
|
||||
|
||||
Active glossary aliases are forwarded to Meeting Lab as provenance metadata.
|
||||
Canonical terminology remains Meeting Context guidance; the exact protocol
|
||||
transcript input and the raw Whisper/diarization artifacts are not rewritten.
|
||||
`glossary_replacements` remains empty.
|
||||
|
||||
Processing and protocol regeneration show measured monotonic stage/total timings,
|
||||
including frozen failure durations. The latest timing display survives ordinary
|
||||
Streamlit reruns. Initial-run durations are saved in `run_metadata.json`;
|
||||
regeneration display timings remain session-local. Worker callbacks capture UI
|
||||
configuration before dispatch and never read Streamlit session state.
|
||||
|
||||
Meeting Lab retains immutable generation records including prompt, transcript
|
||||
input, model response/metadata, glossary configuration, mappings and Meeting
|
||||
Context. Latest protocol/diagnostic/context paths resolve through an atomic
|
||||
`protocol/current` link. Copy whole runs with relative symlinks preserved.
|
||||
|
||||
@@ -8,8 +8,10 @@ local processing backend.
|
||||
|
||||
The first product milestone is a desktop GUI that lets a user:
|
||||
|
||||
- select an existing audio recording
|
||||
- create and edit structured meeting metadata and participants
|
||||
- select an existing WAV, FLAC, or M4A recording
|
||||
- create and edit structured meeting metadata and relevant people, including whether
|
||||
they were present or only mentioned
|
||||
- import and export the reusable People list as versioned YAML
|
||||
- optionally map anonymous `SPEAKER_XX` labels to known participants
|
||||
- start the Meeting Lab processing pipeline
|
||||
- follow stage-based progress
|
||||
@@ -21,12 +23,20 @@ reference recorder, but imported audio is not tied to OBS-specific behavior.
|
||||
|
||||
## Processing Pipeline
|
||||
|
||||
The existing **Meeting language** selector controls both transcription and protocol
|
||||
output: `de` means German for both, and `en` means English for both. The selection
|
||||
is saved in Meeting Context (`meeting.language`) and reused for protocol-only
|
||||
regeneration without retranscription. Older contexts without a language retain
|
||||
German protocol output. Names, speaker mappings and authored Meeting Context are
|
||||
preserved; no transcript translation stage is added.
|
||||
|
||||
```text
|
||||
Audio file
|
||||
-> FFmpeg preparation (mono, 16 kHz PCM WAV)
|
||||
-> FFmpeg preparation (mono, 16 kHz PCM WAV; normalization optional)
|
||||
-> whisper.cpp transcription (large-v3-turbo)
|
||||
-> optional pyannote.audio Community-1 diarization
|
||||
-> direct full-transcript protocol generation
|
||||
-> optional speaker review and explicit protocol generation (when diarized)
|
||||
or direct full-transcript protocol generation (without diarization)
|
||||
-> human review and export
|
||||
```
|
||||
|
||||
@@ -40,6 +50,17 @@ compatibility path. Diarization is optional and produces anonymous speaker
|
||||
labels. A label identifies a participant only when the user explicitly
|
||||
confirms the mapping; automatic speaker-name inference is not allowed.
|
||||
|
||||
Speaker selectors hide participants already assigned to other speakers while keeping
|
||||
the current assignment available. The UI shows detected, assigned and unassigned
|
||||
counts, marks unassigned speakers, and warns when mappings remain incomplete.
|
||||
Anonymous-speaker protocol generation remains available.
|
||||
|
||||
After a diarized run, the result view lists detected `SPEAKER_XX` labels with
|
||||
short transcript excerpts. Confirmed mappings regenerate only the protocol
|
||||
from the existing diarized transcript; audio preparation, Whisper and Pyannote
|
||||
are not rerun. `Unmapped / Unknown` remains valid, and the original anonymous
|
||||
diarized transcript is preserved.
|
||||
|
||||
## Product Outputs
|
||||
|
||||
The product direction includes:
|
||||
@@ -58,8 +79,186 @@ orchestration and protocol-generation logic. Meeting Assistant owns the GUI,
|
||||
context and participant editing, explicit speaker mapping, progress display,
|
||||
protocol editing and export.
|
||||
|
||||
The code currently contains the initial Meeting Assistant project and domain
|
||||
foundation. The GUI has not yet been implemented.
|
||||
Audio normalization is enabled by default and can be disabled in the processing
|
||||
options. This controls loudness normalization only: Meeting Lab still prepares
|
||||
every WAV, FLAC or M4A source as canonical audio before transcription.
|
||||
|
||||
The **Performance profile** selector controls protocol-generation runtime using
|
||||
the abstract Auto (default), Fast, Efficient, and Powersave profiles. Auto lets
|
||||
the inference backend select its own thread configuration; for Ollama, Meeting
|
||||
Assistant intentionally sends no `num_thread` option. Fast, Efficient, and
|
||||
Powersave are explicit resource profiles currently mapped to 16, 10, and 4
|
||||
Ollama CPU threads. These concrete mappings may evolve independently of the UI
|
||||
semantics.
|
||||
|
||||
The People section can export its current entries to a UTF-8 `people.yaml` file
|
||||
and replace them from a previous `.yaml` or `.yml` export. Stable person IDs,
|
||||
names, roles, organizations and attendance states are retained. This is a small
|
||||
reuse mechanism, not a server-side participant library or named meeting-template
|
||||
system.
|
||||
|
||||
## Run the Streamlit MVP
|
||||
|
||||
The development setup expects `meeting-assistant` and `meeting-lab` to be
|
||||
sibling repositories. Meeting Lab currently imports its API through the
|
||||
`src.meeting_lab` package path, so its repository root must be supplied on
|
||||
`PYTHONPATH`. Meeting Assistant does not modify `sys.path` at runtime.
|
||||
|
||||
Create a virtual environment and install Meeting Assistant, including
|
||||
Streamlit:
|
||||
|
||||
```bash
|
||||
python3 -m venv .venv
|
||||
.venv/bin/pip install -e '.[dev]'
|
||||
.venv/bin/pip install -e ../meeting-lab
|
||||
```
|
||||
|
||||
Configure at least the local whisper.cpp model. All supported values are shown
|
||||
in `.env.example`; export them in the shell because the application does not
|
||||
load `.env` files implicitly:
|
||||
|
||||
```bash
|
||||
export MKA_WHISPER_MODEL=/path/to/ggml-large-v3-turbo.bin
|
||||
export MKA_WHISPER_EXECUTABLE=whisper-cli
|
||||
export MKA_FFMPEG_EXECUTABLE=ffmpeg
|
||||
export MKA_PROTOCOL_MODEL=qwen3.8:27b
|
||||
```
|
||||
|
||||
Optional machine-specific settings include:
|
||||
|
||||
- `MKA_DATA_ROOT` (default: `data/meetings`)
|
||||
- `MKA_OLLAMA_ENDPOINT` (default: `http://127.0.0.1:11434`)
|
||||
- `MKA_WHISPER_THREADS` (default: `auto`)
|
||||
- `MKA_FFMPEG_EXECUTABLE` (default: `ffmpeg` found on `PATH`)
|
||||
- `MKA_DIARIZATION_MODE` (`auto`, `cpu`, or `gpu`; default: `auto`)
|
||||
- `MKA_DIARIZATION_RUNTIME` (`native` or `container`; default: `native`)
|
||||
- `MKA_DIARIZATION_CONTAINER_IMAGE` (required for container diarization)
|
||||
- `MKA_DIARIZATION_CONTAINER_ARGS` (default: no extra arguments), encoded as a
|
||||
JSON array of strings so ordering and leading dashes are preserved exactly
|
||||
- `MKA_GLOSSARY_DATABASE` (default: `data/database/glossary.sqlite3`), the local
|
||||
SQLite file used by the global terminology glossary
|
||||
|
||||
## Terminology glossary
|
||||
|
||||
The Streamlit **Terminology glossary** section manages recurring product,
|
||||
material, organization, acronym, and technical names. Store canonical core
|
||||
terms such as `Secugrid HS`, not every compound such as `Secugrid HS Düse`.
|
||||
Aliases help the protocol model recognize transcript variants while retaining
|
||||
the surrounding wording. Only active entries are added to Meeting Context as
|
||||
authoritative terminology for direct protocol generation; inactive entries
|
||||
remain stored but are omitted. The production database starts empty.
|
||||
|
||||
The database defaults to `data/database/glossary.sqlite3`. It is created and
|
||||
bootstrapped automatically and can be moved with `MKA_GLOSSARY_DATABASE`.
|
||||
Glossary integration does not rewrite raw or diarized transcript artifacts.
|
||||
|
||||
Use **Export glossary** to download the complete SQLite-backed glossary as
|
||||
UTF-8 YAML, and **Import glossary** to replace it from a previously exported
|
||||
file. Imports are fully validated before a single SQLite transaction replaces
|
||||
the current glossary, so malformed or conflicting files leave existing data
|
||||
unchanged. Schema version 1 is:
|
||||
|
||||
```yaml
|
||||
version: 1
|
||||
glossary:
|
||||
- id: 1
|
||||
canonical_term: Secugrid HS
|
||||
aliases:
|
||||
- Sikirgut
|
||||
category: product
|
||||
description: Canonical product spelling
|
||||
active: true
|
||||
```
|
||||
|
||||
`id` is the stable SQLite glossary-entry identifier. Categories are `product`,
|
||||
`material`, `organization`, `technical_term`, `acronym`, or `other`.
|
||||
|
||||
## Input configuration import and export
|
||||
|
||||
Use **Export inputs** to save the current Meeting Assistant run-input form as a
|
||||
small, versioned JSON file and **Import inputs** to restore it later. The file
|
||||
contains meeting metadata, participant records, language, date, audio
|
||||
normalization, and diarization choices. It contains form configuration only:
|
||||
generated prompts, protocols, run artifacts, and source-media contents are not
|
||||
included.
|
||||
|
||||
The original source filename may be retained as a reminder, but importing does
|
||||
not restore an upload. Select the audio file explicitly before starting the new
|
||||
run.
|
||||
|
||||
Protocol generation with `qwen3.8:27b` explicitly requests a 32,768-token
|
||||
Ollama context with thinking disabled; Ollama's machine default may otherwise
|
||||
be only 4,096. Keep normal prompt input at approximately 29,000 tokens or less.
|
||||
Although a 31,038-token synthetic prompt passed, larger input is not assumed
|
||||
safe merely because the model advertises a 262,144-token native context.
|
||||
|
||||
For example, a compatible AMD ROCm workstation can configure validated
|
||||
container access without adding controls to the Streamlit UI:
|
||||
|
||||
```bash
|
||||
export MKA_DIARIZATION_MODE=gpu
|
||||
export MKA_DIARIZATION_RUNTIME=container
|
||||
export MKA_DIARIZATION_CONTAINER_IMAGE=rocm/pytorch:rocm7.2.1_ubuntu24.04_py3.12_pytorch_release_2.9.1
|
||||
export MKA_DIARIZATION_CONTAINER_ARGS='["--device=/dev/kfd","--device=/dev/dri","--group-add","video"]'
|
||||
```
|
||||
|
||||
The JSON elements are forwarded as four separate, ordered Meeting Lab
|
||||
container arguments. Runtime and hardware details remain machine-specific
|
||||
environment configuration; the UI continues to expose only the diarization
|
||||
on/off choice.
|
||||
|
||||
From the Meeting Assistant repository, start the UI with:
|
||||
|
||||
```bash
|
||||
PYTHONPATH=src:../meeting-lab .venv/bin/streamlit run src/mka/ui/streamlit_app.py
|
||||
```
|
||||
|
||||
### North launcher
|
||||
|
||||
In the post-diarization Assistant worktree on North, use the checked-in
|
||||
launcher instead of creating another virtual environment:
|
||||
|
||||
```bash
|
||||
./run_north.sh
|
||||
```
|
||||
|
||||
It reuses North's validated Assistant virtual environment and external
|
||||
Whisper/Ollama/Docker-ROCm runtime while loading Meeting Assistant from this
|
||||
worktree and Meeting Lab from
|
||||
`/opt/git-projekts/meeting-lab`. `HF_TOKEN`, when present, is
|
||||
passed through to the disposable diarization container without being stored by
|
||||
the launcher.
|
||||
|
||||
Streamlit opens `http://localhost:8501` by default. Installing the sibling
|
||||
Meeting Lab project supplies its runtime requirements such as PyYAML and
|
||||
Requests. Optional diarization dependencies are needed only when diarization
|
||||
is enabled.
|
||||
|
||||
Uploaded source files are stored under
|
||||
`data/meetings/<meeting-id>/uploads/`. Meeting Lab run artifacts are stored
|
||||
under `data/meetings/<meeting-id>/runs/`. The reviewed protocol is saved as
|
||||
`protocol_edited.md` inside its run directory; the generated `protocol.md`
|
||||
remains unchanged.
|
||||
|
||||
The upload remains in its original format and is passed unchanged to Meeting Lab.
|
||||
Meeting Lab creates the canonical mono 16 kHz signed PCM16 WAV run artifact used by
|
||||
transcription; Meeting Assistant does not duplicate audio conversion.
|
||||
|
||||
See [Architecture](docs/architecture.md), [Project Knowledge](PROJECT_KNOWLEDGE.md),
|
||||
[Roadmap](ROADMAP.md) and [ADR 0011](docs/adr/0011-use-meeting-lab-mvp-backend.md).
|
||||
|
||||
Active glossary aliases are forwarded to Meeting Lab as provenance metadata.
|
||||
Canonical terminology remains Meeting Context guidance; the exact protocol
|
||||
transcript input and the raw Whisper/diarization artifacts are not rewritten.
|
||||
`glossary_replacements` remains empty.
|
||||
|
||||
Processing and protocol regeneration show measured monotonic stage/total timings,
|
||||
including frozen failure durations. The latest timing display survives ordinary
|
||||
Streamlit reruns. Initial-run durations are saved in `run_metadata.json`;
|
||||
regeneration display timings remain session-local. Worker callbacks capture UI
|
||||
configuration before dispatch and never read Streamlit session state.
|
||||
|
||||
Meeting Lab retains immutable generation records including prompt, transcript
|
||||
input, model response/metadata, glossary configuration, mappings and Meeting
|
||||
Context. Latest protocol/diagnostic/context paths resolve through an atomic
|
||||
`protocol/current` link. Copy whole runs with relative symlinks preserved.
|
||||
|
||||
+14
-2
@@ -5,8 +5,10 @@
|
||||
Build the user-facing application over the validated Meeting Lab Python API.
|
||||
|
||||
- select an existing audio file
|
||||
- reliably import WAV, FLAC and M4A through canonical FFmpeg preparation
|
||||
- offer optional audio normalization, enabled by default
|
||||
- create and edit Meeting Context through structured fields
|
||||
- add and manage participants
|
||||
- add and manage present and `mentioned_only` people
|
||||
- optionally map `SPEAKER_XX` labels to participants with explicit confirmation
|
||||
- configure and start `run_mvp_meeting(...)`
|
||||
- display stage-based status and real progress when available
|
||||
@@ -18,6 +20,9 @@ The milestone must keep diarization, GPU acceleration and fixed-duration audio
|
||||
chunking optional. It must not launch the Meeting Lab CLI as a subprocess or
|
||||
duplicate Meeting Lab processing logic.
|
||||
|
||||
Complete real-meeting validation of this vertical slice before expanding the
|
||||
post-run workflow.
|
||||
|
||||
## Following Milestone: Distribution Protocol
|
||||
|
||||
- derive or generate a shorter participant-facing version
|
||||
@@ -28,7 +33,12 @@ duplicate Meeting Lab processing logic.
|
||||
|
||||
## Later Product Work
|
||||
|
||||
- transcript viewing and correction workflows
|
||||
- speaker/name review with explicit human confirmation
|
||||
- deterministic correction of suitable derived artifacts, beginning with a
|
||||
simple search-and-replace path
|
||||
- protocol-only regeneration from existing transcription and diarization plus
|
||||
corrected mappings and Meeting Context
|
||||
- transcript viewing and broader correction workflows
|
||||
- recording and artifact lifecycle management
|
||||
- search, tags, projects and meeting history
|
||||
- richer Markdown, PDF and DOCX export
|
||||
@@ -42,6 +52,8 @@ duplicate Meeting Lab processing logic.
|
||||
- live transcription and real-time summaries
|
||||
- company glossary and custom vocabulary
|
||||
- voice identification with explicit consent and confirmation
|
||||
- automatic name and speaker suggestions only as non-authoritative candidates
|
||||
for later human review
|
||||
- audio cleanup
|
||||
- OCR for shared screens
|
||||
- RAG and knowledge-graph integration
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
# ADR 0012: Preserve Corrections as Post-Run Knowledge
|
||||
|
||||
## Status
|
||||
|
||||
Accepted as a future product direction; not required for the current MVP.
|
||||
|
||||
## Context
|
||||
|
||||
Real-meeting review exposes corrections that should not require repeating
|
||||
expensive processing. These include misspelled or recurring name variants,
|
||||
people mentioned but omitted from the initial Meeting Context, and confirmed
|
||||
`SPEAKER_XX -> participant_id` mappings. A destructive edit would lose useful
|
||||
provenance, while rerunning Whisper or diarization would add cost without
|
||||
improving a deterministic correction.
|
||||
|
||||
The product is also expected to provide both a detailed contextual protocol and
|
||||
a short participant/distribution protocol. Both should eventually use the same
|
||||
confirmed context and corrections.
|
||||
|
||||
## Decision
|
||||
|
||||
Original machine-generated artifacts remain immutable. Corrected or reviewed
|
||||
artifacts are separate derivatives. Human corrections should evolve into
|
||||
structured, meeting-specific knowledge rather than opaque destructive edits;
|
||||
the correction schema is deliberately not defined by this ADR.
|
||||
|
||||
An early correction tool may be simple deterministic search and replace, for
|
||||
example `Grossman` to `Herr Grossmann`. Applying such a confirmed correction to
|
||||
a suitable editable artifact requires neither Whisper, diarization nor an LLM
|
||||
call.
|
||||
|
||||
A later review stage may run after diarization or the complete initial run. It
|
||||
may propose likely person/name matches, allow addition of previously omitted
|
||||
mentioned people, and present anonymous speaker mappings for review. Every
|
||||
suggestion is non-authoritative: anonymous labels remain anonymous until the
|
||||
user explicitly confirms or corrects them, and `mentioned_only` people cannot
|
||||
be mapped as speakers.
|
||||
|
||||
After confirmation, the product should offer two paths:
|
||||
|
||||
1. A fast path applies confirmed speaker, name or text corrections
|
||||
deterministically where that is semantically safe.
|
||||
2. A quality path regenerates protocol output using the existing transcription,
|
||||
existing diarization, corrected Meeting Context and confirmed mappings or
|
||||
corrections. It reruns protocol generation only. Whisper and diarization run
|
||||
again only when separately requested or technically necessary.
|
||||
|
||||
The likely later workflow is therefore:
|
||||
|
||||
```text
|
||||
Audio -> transcription -> optional diarization -> initial protocol
|
||||
-> name/person/speaker review -> human confirmation
|
||||
-> deterministic correction or protocol-only regeneration
|
||||
-> final reviewed detailed and/or distribution protocol
|
||||
```
|
||||
|
||||
The exact ordering may evolve, and this review stage is not mandatory for the
|
||||
current Streamlit MVP.
|
||||
|
||||
## Consequences
|
||||
|
||||
- Human knowledge can improve outputs without unnecessary upstream work.
|
||||
- Provenance is retained because originals and corrected derivatives coexist.
|
||||
- Correction data can later be reused consistently across detailed and short
|
||||
protocol views.
|
||||
- Search/replace, correction storage, review UI, identity suggestions and
|
||||
protocol-only rerun controls remain future implementation work.
|
||||
@@ -0,0 +1,40 @@
|
||||
# Alpha paired validation
|
||||
|
||||
Run the complete Assistant suite with `PYTHONPATH=src:<Lab candidate root>`.
|
||||
`tests/test_alpha_pair.py` exercises the real adapter, context, orchestration,
|
||||
transcript selection and generation persistence with mocked audio preparation,
|
||||
Whisper, diarization and model responses. It covers de/en, all performance
|
||||
profiles, anonymous generation, mapped regeneration, glossary provenance and
|
||||
immutable history. The test skips when Lab is unavailable; a release validation
|
||||
must run it with Lab available and report no skips.
|
||||
|
||||
Run Lab's full `python -m unittest discover -s tests` at its candidate revision.
|
||||
History tests inject write/publication failures and a process interruption.
|
||||
UI tests use the real executor and reject worker access to session state.
|
||||
|
||||
These tests do not measure recognition quality or model language compliance.
|
||||
Before tagging, perform the final local GUI/model smoke with the exact paired
|
||||
revisions and configured model/runtime. Reuse copied historical transcripts
|
||||
where possible; do not rewrite original regression evidence. Preserve the known
|
||||
GTM-Hub human-reference whitespace. Record both candidate SHAs in release notes.
|
||||
|
||||
## Historical Lab fixture dependency
|
||||
|
||||
Four older experiment tests require six ignored files, absent from a fresh Git
|
||||
worktree. Do not add them to the alpha source boundary or silently skip the tests:
|
||||
|
||||
- `artifacts/experiments/evidence_observations_v3/20260819_v3_single_run/h_resulting_action/parsed_observations.json`
|
||||
- `artifacts/experiments/negative_act_form_v0/20260820_qwen35_9b_single_run/na-01/v3_style_input_observations.json`
|
||||
- The same negative-act filename under `na-03`, `na-04`, `na-05`, and `na-06`.
|
||||
|
||||
For the alpha audit, these files were copied from the original Lab worktree to a
|
||||
separate temporary test directory. That directory exposed the candidate's
|
||||
`tests`, `prompts`, `samples`, `src`, and `scripts` as resource links. The full
|
||||
suite ran there with the candidate Lab root on `PYTHONPATH` and the candidate's
|
||||
absolute tests directory supplied to unittest discovery. All 413 tests passed.
|
||||
The bare candidate worktree run instead reports four missing-fixture errors.
|
||||
This is a historical test portability limitation, not a runtime dependency.
|
||||
|
||||
The paired Assistant suite passed all 107 tests with the candidate Lab backend;
|
||||
no Whisper, diarization or Ollama inference was executed. These counts describe
|
||||
the preparation validation and do not replace the final smoke-test record.
|
||||
+15
-2
@@ -27,7 +27,7 @@ subprocess.
|
||||
|
||||
```text
|
||||
Source audio
|
||||
-> FFmpeg normalization/preparation
|
||||
-> FFmpeg preparation (normalization optional, default on)
|
||||
-> mono, 16 kHz PCM WAV
|
||||
-> whisper.cpp transcription with large-v3-turbo
|
||||
-> optional pyannote.audio Community-1 diarization
|
||||
@@ -85,7 +85,11 @@ this context and its mappings.
|
||||
|
||||
Meeting Lab uses FFmpeg to prepare a consistent local-processing input. The
|
||||
current practical target is mono, 16 kHz PCM WAV. The imported source remains a
|
||||
separate source artifact.
|
||||
separate source artifact. Preparation always runs for WAV, FLAC and M4A,
|
||||
regardless of the normalization switch. When enabled, Meeting Lab currently
|
||||
uses `loudnorm=I=-16:LRA=11:TP=-1.5`, an isolated conservative default for
|
||||
speech recordings that may be revisited after empirical comparison. Meeting
|
||||
Assistant passes only an on/off choice and does not own filter parameters.
|
||||
|
||||
### Transcription
|
||||
|
||||
@@ -137,6 +141,15 @@ and reviewed protocols are derived versions. Generated artifacts should retain
|
||||
their input version, backend/model configuration, prompt version and timestamp
|
||||
where practical.
|
||||
|
||||
## Later Post-Run Correction Flow
|
||||
|
||||
Post-run corrections are planned as a separate, non-mandatory workflow after
|
||||
initial protocol generation. As detailed in [ADR 0012](adr/0012-post-run-corrections.md),
|
||||
confirmed name, person and anonymous-speaker corrections should become
|
||||
meeting-specific knowledge. They may then be applied deterministically to safe
|
||||
derived artifacts or used for protocol-only regeneration without needlessly
|
||||
rerunning transcription or diarization.
|
||||
|
||||
## Future Extensions
|
||||
|
||||
- shorter distribution protocols
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
# Meeting Assistant — First functional alpha
|
||||
|
||||
Version: `0.1.0a1`
|
||||
Tag: `v0.1.0-alpha.1`
|
||||
|
||||
This is the first functional end-to-end alpha of the local Meeting Assistant.
|
||||
It provides local meeting configuration and audio workflow, FFmpeg preparation
|
||||
and normalization, Whisper transcription, optional diarization, explicit
|
||||
speaker mapping, anonymous-speaker protocol generation, mapped-speaker
|
||||
regeneration, German and English meeting-language handling, glossary handling
|
||||
with YAML interchange, glossary provenance without transcript mutation,
|
||||
configurable protocol performance profiles, processing timings, generation
|
||||
history and provenance, and atomic publication of complete successful protocol
|
||||
generations. The GTM-Hub qualitative regression case is included.
|
||||
|
||||
The validated candidate revisions for this release were:
|
||||
|
||||
- Assistant: `b50b87c32b3941310c719f459fd45faa87658b68`
|
||||
- Meeting Lab: `d2e3636b949280402728308022c99bdee6ea1066`
|
||||
|
||||
The immutable release revisions are the commits targeted by the paired
|
||||
`v0.1.0-alpha.1` tags in the two repositories.
|
||||
|
||||
## Accepted alpha limitations
|
||||
|
||||
- A local runtime, model, and container setup is required.
|
||||
- Generated protocols require human review, especially responsibility and
|
||||
action attribution; speaker identity requires explicit confirmation.
|
||||
- Diarization can be imperfect. The final native smoke environment did not
|
||||
contain `pyannote.audio==4.0.7` and PyTorch, so real diarization and mapped
|
||||
speaker regeneration were not rerun there; automated tests cover these paths.
|
||||
- Real German and English `qwen3.8:27b` protocol generation was validated.
|
||||
- Context limits can cause attribution fallback or oversized-input rejection.
|
||||
- Historical ignored Lab fixtures are required by some old experiment tests and
|
||||
are not automatically present in a completely fresh Lab worktree.
|
||||
- Generation publication assumes a local POSIX filesystem; power-loss
|
||||
durability is not guaranteed.
|
||||
- Visual redesign is post-alpha, and concise distribution-protocol rendering
|
||||
remains future work.
|
||||
+4
-2
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
||||
|
||||
[project]
|
||||
name = "meeting-assistant"
|
||||
version = "0.1.0"
|
||||
version = "0.1.0a1"
|
||||
description = "Local-first meeting transcription and knowledge capture assistant"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.11"
|
||||
@@ -12,7 +12,9 @@ authors = [
|
||||
{ name = "Martin Tazl" }
|
||||
]
|
||||
dependencies = [
|
||||
"pydantic>=2,<3"
|
||||
"pydantic>=2,<3",
|
||||
"PyYAML>=6,<7",
|
||||
"streamlit>=1.40,<2",
|
||||
]
|
||||
|
||||
[project.optional-dependencies]
|
||||
|
||||
Executable
+57
@@ -0,0 +1,57 @@
|
||||
#!/usr/bin/env bash
|
||||
set -euo pipefail
|
||||
|
||||
assistant_root="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd -P)"
|
||||
assistant_src="$assistant_root/src"
|
||||
lab_root=/opt/git-projekts/meeting-lab
|
||||
north_venv=/opt/git-projekts/meeting-assistant/.venv
|
||||
whisper_executable=/home/martin/whisper.cpp/whisper.cpp/build/bin/whisper-cli
|
||||
whisper_model=/home/martin/whisper.cpp/whisper.cpp/models/ggml-large-v3-turbo.bin
|
||||
|
||||
require_directory() {
|
||||
if [[ ! -d "$1" ]]; then
|
||||
printf 'Required directory is missing: %s\n' "$1" >&2
|
||||
exit 1
|
||||
fi
|
||||
}
|
||||
|
||||
require_file() {
|
||||
if [[ ! -f "$1" ]]; then
|
||||
printf 'Required file is missing: %s\n' "$1" >&2
|
||||
exit 1
|
||||
fi
|
||||
}
|
||||
|
||||
require_executable() {
|
||||
if [[ ! -x "$1" ]]; then
|
||||
printf 'Required executable is missing or not executable: %s\n' "$1" >&2
|
||||
exit 1
|
||||
fi
|
||||
}
|
||||
|
||||
require_directory "$assistant_src"
|
||||
require_file "$assistant_src/mka/ui/streamlit_app.py"
|
||||
require_directory "$lab_root"
|
||||
require_file "$lab_root/src/meeting_lab/__init__.py"
|
||||
require_directory "$north_venv"
|
||||
require_executable "$north_venv/bin/streamlit"
|
||||
require_executable "$whisper_executable"
|
||||
require_file "$whisper_model"
|
||||
if ! command -v docker >/dev/null 2>&1; then
|
||||
printf 'Required executable is not on PATH: docker\n' >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
export MKA_WHISPER_EXECUTABLE="$whisper_executable"
|
||||
export MKA_WHISPER_MODEL="$whisper_model"
|
||||
export MKA_PROTOCOL_MODEL=qwen3.8:27b
|
||||
export MKA_OLLAMA_ENDPOINT=http://127.0.0.1:11434
|
||||
export MKA_DIARIZATION_MODE=auto
|
||||
export MKA_DIARIZATION_RUNTIME=container
|
||||
export MKA_DIARIZATION_CONTAINER_IMAGE=rocm/pytorch:rocm7.2.1_ubuntu24.04_py3.12_pytorch_release_2.9.1
|
||||
export MKA_DIARIZATION_CONTAINER_ARGS='["--device=/dev/kfd","--device=/dev/dri","--group-add","video"]'
|
||||
|
||||
# Keep feature sources ahead of the reused venv's editable installations.
|
||||
export PYTHONPATH="$assistant_src:$lab_root"
|
||||
|
||||
exec "$north_venv/bin/streamlit" run "$assistant_src/mka/ui/streamlit_app.py"
|
||||
@@ -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",
|
||||
]
|
||||
@@ -0,0 +1,106 @@
|
||||
"""Environment-backed configuration for the Meeting Assistant application."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
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
|
||||
glossary_database: Path = Path("data/database/glossary.sqlite3")
|
||||
whisper_executable: str = "whisper-cli"
|
||||
ffmpeg_executable: str = "ffmpeg"
|
||||
protocol_model: str = "qwen3.8:27b"
|
||||
ollama_endpoint: str = "http://127.0.0.1:11434"
|
||||
protocol_num_ctx: int = 32_768
|
||||
protocol_safe_input_token_budget: int = 29_000
|
||||
language: str = "de"
|
||||
threads: str | int = "auto"
|
||||
diarization_mode: str = "auto"
|
||||
diarization_runtime: str = "native"
|
||||
diarization_container_image: str | None = None
|
||||
diarization_container_args: tuple[str, ...] = ()
|
||||
|
||||
@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,
|
||||
glossary_database=Path(
|
||||
os.getenv("MKA_GLOSSARY_DATABASE", "data/database/glossary.sqlite3")
|
||||
),
|
||||
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"),
|
||||
diarization_container_args=_parse_container_args(
|
||||
os.getenv("MKA_DIARIZATION_CONTAINER_ARGS")
|
||||
),
|
||||
)
|
||||
|
||||
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."
|
||||
)
|
||||
if any(
|
||||
not isinstance(argument, str) or not argument
|
||||
for argument in self.diarization_container_args
|
||||
):
|
||||
raise ConfigurationError(
|
||||
"MKA_DIARIZATION_CONTAINER_ARGS must contain only non-empty strings."
|
||||
)
|
||||
|
||||
|
||||
def _parse_container_args(value: str | None) -> tuple[str, ...]:
|
||||
"""Parse an ordered JSON array of opaque container command arguments."""
|
||||
if value is None or not value.strip():
|
||||
return ()
|
||||
try:
|
||||
parsed = json.loads(value)
|
||||
except json.JSONDecodeError as exc:
|
||||
raise ConfigurationError(
|
||||
"MKA_DIARIZATION_CONTAINER_ARGS must be a JSON array of strings."
|
||||
) from exc
|
||||
if not isinstance(parsed, list) or any(
|
||||
not isinstance(argument, str) or not argument for argument in parsed
|
||||
):
|
||||
raise ConfigurationError(
|
||||
"MKA_DIARIZATION_CONTAINER_ARGS must be a JSON array of non-empty strings."
|
||||
)
|
||||
return tuple(parsed)
|
||||
@@ -0,0 +1,369 @@
|
||||
"""SQLite-backed global terminology glossary."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sqlite3
|
||||
from collections.abc import Iterable, Iterator
|
||||
from contextlib import contextmanager
|
||||
from dataclasses import dataclass
|
||||
from datetime import UTC, datetime
|
||||
from pathlib import Path
|
||||
|
||||
GLOSSARY_CATEGORIES = (
|
||||
"product",
|
||||
"material",
|
||||
"organization",
|
||||
"technical_term",
|
||||
"acronym",
|
||||
"other",
|
||||
)
|
||||
|
||||
|
||||
class GlossaryConflictError(ValueError):
|
||||
"""Raised when a canonical term or alias conflicts with existing terminology."""
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GlossaryEntry:
|
||||
id: int
|
||||
canonical_term: str
|
||||
category: str
|
||||
description: str | None
|
||||
is_active: bool
|
||||
aliases: tuple[str, ...]
|
||||
created_at: str
|
||||
updated_at: str
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GlossaryReplacementEntry:
|
||||
"""Validated values used to atomically replace the persisted glossary."""
|
||||
|
||||
id: int
|
||||
canonical_term: str
|
||||
category: str
|
||||
description: str | None
|
||||
is_active: bool
|
||||
aliases: tuple[str, ...]
|
||||
|
||||
|
||||
class GlossaryRepository:
|
||||
"""Small data-access boundary for the local glossary database."""
|
||||
|
||||
def __init__(self, database_path: Path) -> None:
|
||||
self.database_path = Path(database_path)
|
||||
|
||||
def initialize(self) -> None:
|
||||
"""Create the database and current schema when absent."""
|
||||
self.database_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
with self._connect() as connection:
|
||||
connection.executescript(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS glossary_entries (
|
||||
id INTEGER PRIMARY KEY,
|
||||
canonical_term TEXT NOT NULL COLLATE NOCASE UNIQUE,
|
||||
category TEXT NOT NULL CHECK (category IN (
|
||||
'product', 'material', 'organization',
|
||||
'technical_term', 'acronym', 'other'
|
||||
)),
|
||||
description TEXT,
|
||||
is_active INTEGER NOT NULL DEFAULT 1 CHECK (is_active IN (0, 1)),
|
||||
created_at TEXT NOT NULL,
|
||||
updated_at TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS glossary_aliases (
|
||||
id INTEGER PRIMARY KEY,
|
||||
entry_id INTEGER NOT NULL REFERENCES glossary_entries(id)
|
||||
ON DELETE CASCADE,
|
||||
alias TEXT NOT NULL COLLATE NOCASE UNIQUE,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_glossary_aliases_entry_id
|
||||
ON glossary_aliases(entry_id);
|
||||
PRAGMA user_version = 1;
|
||||
"""
|
||||
)
|
||||
|
||||
def create(
|
||||
self,
|
||||
canonical_term: str,
|
||||
category: str,
|
||||
*,
|
||||
aliases: Iterable[str] = (),
|
||||
description: str | None = None,
|
||||
is_active: bool = True,
|
||||
) -> GlossaryEntry:
|
||||
canonical, normalized_aliases = self._validate_values(canonical_term, category, aliases)
|
||||
now = _timestamp()
|
||||
try:
|
||||
with self._connect() as connection:
|
||||
self._ensure_terms_available(connection, canonical, normalized_aliases)
|
||||
cursor = connection.execute(
|
||||
"""INSERT INTO glossary_entries
|
||||
(canonical_term, category, description, is_active, created_at, updated_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?)""",
|
||||
(canonical, category, _optional_text(description), is_active, now, now),
|
||||
)
|
||||
entry_id = int(cursor.lastrowid)
|
||||
connection.executemany(
|
||||
"INSERT INTO glossary_aliases (entry_id, alias, created_at) VALUES (?, ?, ?)",
|
||||
((entry_id, alias, now) for alias in normalized_aliases),
|
||||
)
|
||||
except sqlite3.IntegrityError as exc:
|
||||
raise GlossaryConflictError("Canonical term or alias already exists.") from exc
|
||||
return self.get(entry_id)
|
||||
|
||||
def get(self, entry_id: int) -> GlossaryEntry:
|
||||
with self._connect() as connection:
|
||||
row = connection.execute(
|
||||
"SELECT * FROM glossary_entries WHERE id = ?", (entry_id,)
|
||||
).fetchone()
|
||||
if row is None:
|
||||
raise KeyError(f"Unknown glossary entry: {entry_id}")
|
||||
return self._to_entry(connection, row)
|
||||
|
||||
def list(self, search: str = "", *, active_only: bool = False) -> list[GlossaryEntry]:
|
||||
clauses: list[str] = []
|
||||
parameters: list[object] = []
|
||||
if active_only:
|
||||
clauses.append("entry.is_active = 1")
|
||||
if search.strip():
|
||||
clauses.append(
|
||||
"(entry.canonical_term LIKE ? COLLATE NOCASE "
|
||||
"OR entry.category LIKE ? COLLATE NOCASE "
|
||||
"OR entry.description LIKE ? COLLATE NOCASE "
|
||||
"OR alias.alias LIKE ? COLLATE NOCASE)"
|
||||
)
|
||||
pattern = f"%{search.strip()}%"
|
||||
parameters.extend([pattern] * 4)
|
||||
where = f"WHERE {' AND '.join(clauses)}" if clauses else ""
|
||||
with self._connect() as connection:
|
||||
rows = connection.execute(
|
||||
f"""SELECT DISTINCT entry.* FROM glossary_entries AS entry
|
||||
LEFT JOIN glossary_aliases AS alias ON alias.entry_id = entry.id
|
||||
{where} ORDER BY entry.canonical_term COLLATE NOCASE""", # noqa: S608
|
||||
parameters,
|
||||
).fetchall()
|
||||
return [self._to_entry(connection, row) for row in rows]
|
||||
|
||||
def update(
|
||||
self,
|
||||
entry_id: int,
|
||||
canonical_term: str,
|
||||
category: str,
|
||||
*,
|
||||
aliases: Iterable[str] = (),
|
||||
description: str | None = None,
|
||||
is_active: bool = True,
|
||||
) -> GlossaryEntry:
|
||||
canonical, normalized_aliases = self._validate_values(canonical_term, category, aliases)
|
||||
try:
|
||||
with self._connect() as connection:
|
||||
exists = connection.execute(
|
||||
"SELECT 1 FROM glossary_entries WHERE id = ?", (entry_id,)
|
||||
).fetchone()
|
||||
if exists is None:
|
||||
raise KeyError(f"Unknown glossary entry: {entry_id}")
|
||||
self._ensure_terms_available(
|
||||
connection, canonical, normalized_aliases, excluding_entry_id=entry_id
|
||||
)
|
||||
now = _timestamp()
|
||||
connection.execute(
|
||||
"""UPDATE glossary_entries SET canonical_term = ?, category = ?,
|
||||
description = ?, is_active = ?, updated_at = ? WHERE id = ?""",
|
||||
(
|
||||
canonical,
|
||||
category,
|
||||
_optional_text(description),
|
||||
is_active,
|
||||
now,
|
||||
entry_id,
|
||||
),
|
||||
)
|
||||
connection.execute("DELETE FROM glossary_aliases WHERE entry_id = ?", (entry_id,))
|
||||
connection.executemany(
|
||||
"INSERT INTO glossary_aliases (entry_id, alias, created_at) VALUES (?, ?, ?)",
|
||||
((entry_id, alias, now) for alias in normalized_aliases),
|
||||
)
|
||||
except sqlite3.IntegrityError as exc:
|
||||
raise GlossaryConflictError("Canonical term or alias already exists.") from exc
|
||||
return self.get(entry_id)
|
||||
|
||||
def set_active(self, entry_id: int, is_active: bool) -> GlossaryEntry:
|
||||
entry = self.get(entry_id)
|
||||
return self.update(
|
||||
entry.id,
|
||||
entry.canonical_term,
|
||||
entry.category,
|
||||
aliases=entry.aliases,
|
||||
description=entry.description,
|
||||
is_active=is_active,
|
||||
)
|
||||
|
||||
def delete(self, entry_id: int) -> None:
|
||||
with self._connect() as connection:
|
||||
cursor = connection.execute("DELETE FROM glossary_entries WHERE id = ?", (entry_id,))
|
||||
if cursor.rowcount == 0:
|
||||
raise KeyError(f"Unknown glossary entry: {entry_id}")
|
||||
|
||||
def replace_all(self, entries: Iterable[GlossaryReplacementEntry]) -> None:
|
||||
"""Atomically replace every glossary entry, preserving supplied stable IDs."""
|
||||
replacements = tuple(entries)
|
||||
validated: list[GlossaryReplacementEntry] = []
|
||||
seen_ids: set[int] = set()
|
||||
seen_terms: set[str] = set()
|
||||
for entry in replacements:
|
||||
if type(entry.id) is not int or entry.id <= 0:
|
||||
raise ValueError("Glossary entry IDs must be positive integers.")
|
||||
if entry.id in seen_ids:
|
||||
raise GlossaryConflictError(f"Duplicate glossary entry ID: {entry.id}.")
|
||||
canonical, aliases = self._validate_values(
|
||||
entry.canonical_term, entry.category, entry.aliases
|
||||
)
|
||||
folded_terms = {canonical.casefold(), *(alias.casefold() for alias in aliases)}
|
||||
if seen_terms.intersection(folded_terms):
|
||||
raise GlossaryConflictError("Canonical term or alias already exists.")
|
||||
seen_ids.add(entry.id)
|
||||
seen_terms.update(folded_terms)
|
||||
validated.append(
|
||||
GlossaryReplacementEntry(
|
||||
id=entry.id,
|
||||
canonical_term=canonical,
|
||||
category=entry.category,
|
||||
description=_optional_text(entry.description),
|
||||
is_active=entry.is_active,
|
||||
aliases=aliases,
|
||||
)
|
||||
)
|
||||
|
||||
now = _timestamp()
|
||||
try:
|
||||
with self._connect() as connection:
|
||||
connection.execute("DELETE FROM glossary_entries")
|
||||
connection.executemany(
|
||||
"""INSERT INTO glossary_entries
|
||||
(id, canonical_term, category, description, is_active, created_at, updated_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)""",
|
||||
(
|
||||
(
|
||||
entry.id,
|
||||
entry.canonical_term,
|
||||
entry.category,
|
||||
entry.description,
|
||||
entry.is_active,
|
||||
now,
|
||||
now,
|
||||
)
|
||||
for entry in validated
|
||||
),
|
||||
)
|
||||
connection.executemany(
|
||||
"INSERT INTO glossary_aliases (entry_id, alias, created_at) VALUES (?, ?, ?)",
|
||||
((entry.id, alias, now) for entry in validated for alias in entry.aliases),
|
||||
)
|
||||
except sqlite3.IntegrityError as exc:
|
||||
raise GlossaryConflictError("Canonical term or alias already exists.") from exc
|
||||
|
||||
@contextmanager
|
||||
def _connect(self) -> Iterator[sqlite3.Connection]:
|
||||
connection = sqlite3.connect(self.database_path, timeout=5)
|
||||
try:
|
||||
connection.row_factory = sqlite3.Row
|
||||
connection.execute("PRAGMA foreign_keys = ON")
|
||||
connection.execute("PRAGMA journal_mode = WAL")
|
||||
connection.execute("PRAGMA busy_timeout = 5000")
|
||||
with connection:
|
||||
yield connection
|
||||
finally:
|
||||
connection.close()
|
||||
|
||||
@staticmethod
|
||||
def _validate_values(
|
||||
canonical_term: str, category: str, aliases: Iterable[str]
|
||||
) -> tuple[str, tuple[str, ...]]:
|
||||
canonical = canonical_term.strip()
|
||||
if not canonical:
|
||||
raise ValueError("Canonical term is required.")
|
||||
if category not in GLOSSARY_CATEGORIES:
|
||||
raise ValueError(f"Unsupported glossary category: {category}")
|
||||
normalized_aliases = tuple(
|
||||
dict.fromkeys(alias.strip() for alias in aliases if alias.strip())
|
||||
)
|
||||
folded = [alias.casefold() for alias in normalized_aliases]
|
||||
if len(folded) != len(set(folded)) or canonical.casefold() in folded:
|
||||
raise GlossaryConflictError(
|
||||
"Aliases must be unique and differ from the canonical term."
|
||||
)
|
||||
return canonical, normalized_aliases
|
||||
|
||||
@staticmethod
|
||||
def _ensure_terms_available(
|
||||
connection: sqlite3.Connection,
|
||||
canonical: str,
|
||||
aliases: tuple[str, ...],
|
||||
*,
|
||||
excluding_entry_id: int | None = None,
|
||||
) -> None:
|
||||
terms = (canonical, *aliases)
|
||||
placeholders = ", ".join("?" for _ in terms)
|
||||
exclusion = "AND entry_id != ?" if excluding_entry_id is not None else ""
|
||||
alias_parameters: list[object] = [*terms]
|
||||
if excluding_entry_id is not None:
|
||||
alias_parameters.append(excluding_entry_id)
|
||||
alias_conflict = connection.execute(
|
||||
f"SELECT 1 FROM glossary_aliases WHERE alias IN ({placeholders}) {exclusion} LIMIT 1", # noqa: S608
|
||||
alias_parameters,
|
||||
).fetchone()
|
||||
entry_exclusion = "AND id != ?" if excluding_entry_id is not None else ""
|
||||
entry_parameters: list[object] = [*terms]
|
||||
if excluding_entry_id is not None:
|
||||
entry_parameters.append(excluding_entry_id)
|
||||
canonical_conflict = connection.execute(
|
||||
f"SELECT 1 FROM glossary_entries WHERE canonical_term IN ({placeholders}) " # noqa: S608
|
||||
f"{entry_exclusion} LIMIT 1",
|
||||
entry_parameters,
|
||||
).fetchone()
|
||||
if alias_conflict or canonical_conflict:
|
||||
raise GlossaryConflictError("Canonical term or alias already exists.")
|
||||
|
||||
@staticmethod
|
||||
def _to_entry(connection: sqlite3.Connection, row: sqlite3.Row) -> GlossaryEntry:
|
||||
aliases = connection.execute(
|
||||
"SELECT alias FROM glossary_aliases WHERE entry_id = ? ORDER BY alias COLLATE NOCASE",
|
||||
(row["id"],),
|
||||
).fetchall()
|
||||
return GlossaryEntry(
|
||||
id=row["id"],
|
||||
canonical_term=row["canonical_term"],
|
||||
category=row["category"],
|
||||
description=row["description"],
|
||||
is_active=bool(row["is_active"]),
|
||||
aliases=tuple(alias["alias"] for alias in aliases),
|
||||
created_at=row["created_at"],
|
||||
updated_at=row["updated_at"],
|
||||
)
|
||||
|
||||
|
||||
def render_glossary_terms(entries: Iterable[GlossaryEntry]) -> list[str]:
|
||||
"""Render concise canonical terms with recognition aliases for Meeting Context."""
|
||||
rendered = []
|
||||
for entry in entries:
|
||||
item = entry.canonical_term
|
||||
if entry.aliases:
|
||||
item += f" (aliases: {', '.join(entry.aliases)})"
|
||||
rendered.append(item)
|
||||
return rendered
|
||||
|
||||
|
||||
def glossary_alias_mapping(entries: Iterable[GlossaryEntry]) -> dict[str, str]:
|
||||
"""Return explicit alias configuration for provenance, not substitution."""
|
||||
return {alias: entry.canonical_term for entry in entries for alias in entry.aliases}
|
||||
|
||||
|
||||
def _timestamp() -> str:
|
||||
return datetime.now(UTC).isoformat(timespec="seconds")
|
||||
|
||||
|
||||
def _optional_text(value: str | None) -> str | None:
|
||||
stripped = value.strip() if value else ""
|
||||
return stripped or None
|
||||
@@ -0,0 +1,122 @@
|
||||
"""Versioned UTF-8 YAML import and export for the terminology glossary."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Sequence
|
||||
|
||||
import yaml
|
||||
|
||||
from mka.application.glossary import (
|
||||
GLOSSARY_CATEGORIES,
|
||||
GlossaryEntry,
|
||||
GlossaryReplacementEntry,
|
||||
)
|
||||
|
||||
GLOSSARY_YAML_VERSION = 1
|
||||
|
||||
|
||||
class GlossaryYamlError(ValueError):
|
||||
"""Raised when a glossary YAML document is malformed or invalid."""
|
||||
|
||||
|
||||
def export_glossary_yaml(entries: Sequence[GlossaryEntry]) -> str:
|
||||
"""Serialize all glossary state in a deterministic, versioned document."""
|
||||
document = {
|
||||
"version": GLOSSARY_YAML_VERSION,
|
||||
"glossary": [
|
||||
{
|
||||
"id": entry.id,
|
||||
"canonical_term": entry.canonical_term,
|
||||
"aliases": list(entry.aliases),
|
||||
"category": entry.category,
|
||||
"description": entry.description,
|
||||
"active": entry.is_active,
|
||||
}
|
||||
for entry in entries
|
||||
],
|
||||
}
|
||||
return yaml.safe_dump(document, allow_unicode=True, sort_keys=False, default_flow_style=False)
|
||||
|
||||
|
||||
def import_glossary_yaml(content: str | bytes) -> list[GlossaryReplacementEntry]:
|
||||
"""Parse and completely validate a glossary replacement document."""
|
||||
try:
|
||||
if isinstance(content, bytes):
|
||||
content = content.decode("utf-8")
|
||||
document = yaml.safe_load(content)
|
||||
except (yaml.YAMLError, UnicodeDecodeError) as exc:
|
||||
raise GlossaryYamlError(f"Malformed glossary YAML: {exc}") from exc
|
||||
|
||||
if not isinstance(document, dict):
|
||||
raise GlossaryYamlError("Glossary YAML must contain a top-level mapping.")
|
||||
version = document.get("version")
|
||||
if type(version) is not int or version != GLOSSARY_YAML_VERSION:
|
||||
raise GlossaryYamlError(
|
||||
f"Unsupported glossary YAML version {version!r}; expected {GLOSSARY_YAML_VERSION}."
|
||||
)
|
||||
raw_entries = document.get("glossary")
|
||||
if not isinstance(raw_entries, list):
|
||||
raise GlossaryYamlError("Glossary YAML must contain a top-level 'glossary' list.")
|
||||
|
||||
result: list[GlossaryReplacementEntry] = []
|
||||
seen_ids: set[int] = set()
|
||||
seen_terms: set[str] = set()
|
||||
for index, raw in enumerate(raw_entries, start=1):
|
||||
if not isinstance(raw, dict):
|
||||
raise GlossaryYamlError(f"Glossary entry {index} must be a mapping.")
|
||||
entry_id = raw.get("id")
|
||||
if type(entry_id) is not int or entry_id <= 0:
|
||||
raise GlossaryYamlError(f"Glossary entry {index} requires a positive integer id.")
|
||||
if entry_id in seen_ids:
|
||||
raise GlossaryYamlError(f"Duplicate glossary entry id: {entry_id}.")
|
||||
seen_ids.add(entry_id)
|
||||
|
||||
canonical = _required_text(raw, "canonical_term", index).strip()
|
||||
category = raw.get("category")
|
||||
if category not in GLOSSARY_CATEGORIES:
|
||||
allowed = ", ".join(GLOSSARY_CATEGORIES)
|
||||
raise GlossaryYamlError(
|
||||
f"Glossary entry {index} has invalid category; expected one of: {allowed}."
|
||||
)
|
||||
aliases_value = raw.get("aliases")
|
||||
if not isinstance(aliases_value, list) or any(
|
||||
not isinstance(alias, str) or not alias.strip() for alias in aliases_value
|
||||
):
|
||||
raise GlossaryYamlError(f"Glossary entry {index} aliases must be a list of text.")
|
||||
aliases = tuple(alias.strip() for alias in aliases_value)
|
||||
folded = [canonical.casefold(), *(alias.casefold() for alias in aliases)]
|
||||
if len(folded) != len(set(folded)):
|
||||
raise GlossaryYamlError(
|
||||
f"Glossary entry {index} aliases must be unique and differ from its canonical term."
|
||||
)
|
||||
duplicate = next((term for term in folded if term in seen_terms), None)
|
||||
if duplicate is not None:
|
||||
raise GlossaryYamlError(
|
||||
f"Glossary entry {index} contains a duplicate canonical term or alias."
|
||||
)
|
||||
seen_terms.update(folded)
|
||||
|
||||
description = raw.get("description")
|
||||
if description is not None and not isinstance(description, str):
|
||||
raise GlossaryYamlError(f"Glossary entry {index} description must be text or null.")
|
||||
active = raw.get("active")
|
||||
if type(active) is not bool:
|
||||
raise GlossaryYamlError(f"Glossary entry {index} active must be true or false.")
|
||||
result.append(
|
||||
GlossaryReplacementEntry(
|
||||
id=entry_id,
|
||||
canonical_term=canonical,
|
||||
aliases=aliases,
|
||||
category=category,
|
||||
description=description,
|
||||
is_active=active,
|
||||
)
|
||||
)
|
||||
return result
|
||||
|
||||
|
||||
def _required_text(entry: dict[object, object], field: str, index: int) -> str:
|
||||
value = entry.get(field)
|
||||
if not isinstance(value, str) or not value.strip():
|
||||
raise GlossaryYamlError(f"Glossary entry {index} requires a non-empty {field}.")
|
||||
return value
|
||||
@@ -0,0 +1,561 @@
|
||||
"""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
|
||||
|
||||
import yaml
|
||||
|
||||
from mka.application.config import AppSettings
|
||||
from mka.application.glossary import (
|
||||
GlossaryRepository,
|
||||
glossary_alias_mapping,
|
||||
render_glossary_terms,
|
||||
)
|
||||
from mka.application.performance import DEFAULT_PERFORMANCE_PROFILE, resolve_ollama_num_thread
|
||||
from mka.application.progress_timing import ProcessingTimer
|
||||
|
||||
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: ...
|
||||
|
||||
def regenerate_protocol(
|
||||
self,
|
||||
run_dir: Path,
|
||||
meeting_context: Any,
|
||||
progress_sink: Callable[[Any], None],
|
||||
**options: Any,
|
||||
) -> 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
|
||||
performance_profile: str = DEFAULT_PERFORMANCE_PROFILE
|
||||
|
||||
|
||||
@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
|
||||
speaker_attribution_available: bool | None = None
|
||||
awaiting_speaker_review: bool = False
|
||||
detected_speaker_count: int | None = None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SpeakerReview:
|
||||
speaker_label: str
|
||||
excerpts: tuple[str, ...]
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SpeakerMappingReview:
|
||||
speakers: tuple[SpeakerReview, ...]
|
||||
participants: tuple[tuple[str, str], ...]
|
||||
current_mappings: dict[str, str]
|
||||
|
||||
|
||||
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,
|
||||
glossary: GlossaryRepository | None = None,
|
||||
timer: ProcessingTimer | None = None,
|
||||
) -> None:
|
||||
self.settings = settings
|
||||
self.meeting_lab = meeting_lab
|
||||
self.glossary = glossary or GlossaryRepository(settings.glossary_database)
|
||||
self.glossary.initialize()
|
||||
self.timer = timer or ProcessingTimer()
|
||||
|
||||
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,
|
||||
},
|
||||
}
|
||||
self._apply_glossary(data)
|
||||
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,
|
||||
"protocol_num_ctx": self.settings.protocol_num_ctx,
|
||||
"protocol_safe_input_token_budget": (
|
||||
self.settings.protocol_safe_input_token_budget
|
||||
),
|
||||
"protocol_num_thread": resolve_ollama_num_thread(options.performance_profile),
|
||||
"glossary_aliases": glossary_alias_mapping(self.glossary.list(active_only=True)),
|
||||
"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),
|
||||
"diarization_container_args": self.settings.diarization_container_args,
|
||||
"stop_after_diarization": options.diarization_enabled,
|
||||
}
|
||||
)
|
||||
current_stage: str | None = None
|
||||
self.timer.start_run()
|
||||
|
||||
def relay(event: Any) -> None:
|
||||
nonlocal current_stage
|
||||
if event.stage in STAGES and event.status == "started":
|
||||
current_stage = event.stage
|
||||
self._update_timer(event)
|
||||
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)
|
||||
|
||||
try:
|
||||
result = self.meeting_lab.run(config, context, relay)
|
||||
except Exception:
|
||||
self.timer.finish_run()
|
||||
raise
|
||||
self.timer.finish_run()
|
||||
self._persist_timing(result.run_dir)
|
||||
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)
|
||||
failure_message = failure.get("message") or "Meeting Lab processing failed."
|
||||
failure_type = failure.get("type")
|
||||
if failure_type and not failure_message.startswith(f"{failure_type}:"):
|
||||
failure_message = f"{failure_type}: {failure_message}"
|
||||
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_message,
|
||||
)
|
||||
if result.protocol_path is None:
|
||||
return ProcessingOutcome(
|
||||
succeeded=True,
|
||||
run_dir=result.run_dir,
|
||||
original_protocol=None,
|
||||
protocol_path=None,
|
||||
awaiting_speaker_review=True,
|
||||
detected_speaker_count=self._detected_speaker_count(result.run_dir),
|
||||
)
|
||||
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,
|
||||
speaker_attribution_available=self._speaker_attribution_available(result.run_dir),
|
||||
)
|
||||
|
||||
def load_speaker_mapping_review(
|
||||
self,
|
||||
run_dir: Path,
|
||||
*,
|
||||
excerpts_per_speaker: int = 3,
|
||||
minimum_excerpt_characters: int = 20,
|
||||
) -> SpeakerMappingReview | None:
|
||||
"""Load detected labels and small deterministic identification excerpts."""
|
||||
run_dir = Path(run_dir)
|
||||
transcript_path = run_dir / "diarization" / "transcript_diarized.json"
|
||||
context_path = run_dir / "context" / "meeting_context.yaml"
|
||||
if not transcript_path.is_file() or not context_path.is_file():
|
||||
return None
|
||||
transcript = json.loads(transcript_path.read_text(encoding="utf-8-sig"))
|
||||
context_data = yaml.safe_load(context_path.read_text(encoding="utf-8"))
|
||||
if not isinstance(transcript, dict) or not isinstance(context_data, dict):
|
||||
raise ValueError("Existing run contains malformed speaker review artifacts.")
|
||||
segments = transcript.get("segments")
|
||||
if not isinstance(segments, list):
|
||||
raise ValueError("Diarized transcript has no segments list.")
|
||||
|
||||
texts_by_speaker: dict[str, list[str]] = {}
|
||||
short_by_speaker: dict[str, list[str]] = {}
|
||||
for segment in segments:
|
||||
if not isinstance(segment, dict):
|
||||
continue
|
||||
label = segment.get("speaker_id")
|
||||
text = segment.get("text")
|
||||
if (
|
||||
not isinstance(label, str)
|
||||
or not label.startswith("SPEAKER_")
|
||||
or label == "SPEAKER_UNASSIGNED"
|
||||
or not isinstance(text, str)
|
||||
or not text.strip()
|
||||
):
|
||||
continue
|
||||
normalized = " ".join(text.split())
|
||||
target = (
|
||||
texts_by_speaker
|
||||
if len(normalized) >= minimum_excerpt_characters
|
||||
else short_by_speaker
|
||||
)
|
||||
target.setdefault(label, []).append(normalized)
|
||||
|
||||
labels = sorted(set(texts_by_speaker) | set(short_by_speaker))
|
||||
speakers = []
|
||||
for label in labels:
|
||||
candidates = texts_by_speaker.get(label) or short_by_speaker.get(label, [])
|
||||
speakers.append(SpeakerReview(label, tuple(candidates[:excerpts_per_speaker])))
|
||||
participants = tuple(
|
||||
(participant["participant_id"], participant["display_name"])
|
||||
for participant in context_data.get("participants", [])
|
||||
if isinstance(participant, dict)
|
||||
and isinstance(participant.get("participant_id"), str)
|
||||
and isinstance(participant.get("display_name"), str)
|
||||
)
|
||||
mappings = context_data.get("speaker_mappings")
|
||||
return SpeakerMappingReview(
|
||||
speakers=tuple(speakers),
|
||||
participants=participants,
|
||||
current_mappings=dict(mappings) if isinstance(mappings, dict) else {},
|
||||
)
|
||||
|
||||
def regenerate_protocol(
|
||||
self,
|
||||
run_dir: Path,
|
||||
speaker_mappings: dict[str, str],
|
||||
progress_sink: Callable[[AppProgressEvent], None] | None = None,
|
||||
performance_profile: str = DEFAULT_PERFORMANCE_PROFILE,
|
||||
) -> ProcessingOutcome:
|
||||
"""Persist confirmed mappings and regenerate only the direct protocol."""
|
||||
review = self.load_speaker_mapping_review(run_dir)
|
||||
if review is None:
|
||||
raise ValueError("This run has no diarization artifacts to map.")
|
||||
detected = {speaker.speaker_label for speaker in review.speakers}
|
||||
participant_ids = {participant_id for participant_id, _ in review.participants}
|
||||
unknown_labels = sorted(set(speaker_mappings) - detected)
|
||||
if unknown_labels:
|
||||
raise ValueError(f"Unknown diarization speaker label: {unknown_labels[0]}")
|
||||
unknown_participants = sorted(set(speaker_mappings.values()) - participant_ids)
|
||||
if unknown_participants:
|
||||
raise ValueError(f"Unknown participant ID: {unknown_participants[0]}")
|
||||
assigned_participants = list(speaker_mappings.values())
|
||||
if len(assigned_participants) != len(set(assigned_participants)):
|
||||
raise ValueError("A participant may be assigned to only one speaker label.")
|
||||
|
||||
context_path = Path(run_dir) / "context" / "meeting_context.yaml"
|
||||
context_data = yaml.safe_load(context_path.read_text(encoding="utf-8"))
|
||||
if not isinstance(context_data, dict):
|
||||
raise ValueError("Existing run contains malformed Meeting Context.")
|
||||
context_data["speaker_mappings"] = dict(sorted(speaker_mappings.items()))
|
||||
self._apply_glossary(context_data)
|
||||
context = self.meeting_lab.create_context(context_data)
|
||||
# A checkpoint mapping is user input, not a derived protocol artifact. Save it
|
||||
# before inference so a failed generation remains retryable with the same review.
|
||||
if context_path.is_symlink():
|
||||
context_path.unlink()
|
||||
context_path.write_text(
|
||||
yaml.safe_dump(context.data, allow_unicode=True, sort_keys=False), encoding="utf-8"
|
||||
)
|
||||
|
||||
def relay(event: Any) -> None:
|
||||
self._update_timer(event)
|
||||
if progress_sink is not None:
|
||||
progress_sink(
|
||||
AppProgressEvent(
|
||||
stage=event.stage,
|
||||
status=event.status,
|
||||
elapsed_seconds=event.elapsed_seconds,
|
||||
progress=event.progress,
|
||||
message=event.message,
|
||||
)
|
||||
)
|
||||
|
||||
self.timer.start_run()
|
||||
try:
|
||||
result = self.meeting_lab.regenerate_protocol(
|
||||
Path(run_dir),
|
||||
context,
|
||||
relay,
|
||||
model=self.settings.protocol_model,
|
||||
ollama_endpoint=self.settings.ollama_endpoint,
|
||||
protocol_num_ctx=self.settings.protocol_num_ctx,
|
||||
protocol_safe_input_token_budget=(self.settings.protocol_safe_input_token_budget),
|
||||
protocol_num_thread=resolve_ollama_num_thread(performance_profile),
|
||||
glossary_aliases=glossary_alias_mapping(self.glossary.list(active_only=True)),
|
||||
)
|
||||
except Exception:
|
||||
self.timer.finish_run()
|
||||
raise
|
||||
self.timer.finish_run()
|
||||
self._persist_timing(Path(result.run_dir))
|
||||
protocol_path = Path(result.protocol_path)
|
||||
return ProcessingOutcome(
|
||||
succeeded=True,
|
||||
run_dir=Path(result.run_dir),
|
||||
original_protocol=protocol_path.read_text(encoding="utf-8"),
|
||||
protocol_path=protocol_path,
|
||||
speaker_attribution_available=self._speaker_attribution_available(result.run_dir),
|
||||
)
|
||||
|
||||
def _apply_glossary(self, context_data: dict[str, Any]) -> None:
|
||||
"""Merge active terminology into context without discarding other entities."""
|
||||
glossary_terms = render_glossary_terms(self.glossary.list(active_only=True))
|
||||
known_entities = context_data.setdefault("known_entities", {})
|
||||
if glossary_terms:
|
||||
known_entities["Authoritative terminology"] = glossary_terms
|
||||
else:
|
||||
known_entities.pop("Authoritative terminology", None)
|
||||
rules = context_data.setdefault("context_rules", {})
|
||||
rules["glossary_canonical_spelling"] = (
|
||||
"Use canonical glossary spellings only when the meeting clearly refers to "
|
||||
"those terms; do not invent matches or replace unrelated words."
|
||||
)
|
||||
rules["glossary_core_terms"] = (
|
||||
"Preserve surrounding context and use canonical core terms inside compounds "
|
||||
"where appropriate."
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _speaker_attribution_available(run_dir: Path | None) -> bool | None:
|
||||
if run_dir is None:
|
||||
return None
|
||||
metadata_path = Path(run_dir) / "protocol" / "runtime_metadata.json"
|
||||
if not metadata_path.is_file():
|
||||
return None
|
||||
try:
|
||||
metadata = json.loads(metadata_path.read_text(encoding="utf-8"))
|
||||
except (OSError, json.JSONDecodeError):
|
||||
return None
|
||||
value = metadata.get("speaker_attribution_available")
|
||||
return value if isinstance(value, bool) else None
|
||||
|
||||
@staticmethod
|
||||
def _detected_speaker_count(run_dir: Path | None) -> int | None:
|
||||
if run_dir is None:
|
||||
return None
|
||||
metadata_path = Path(run_dir) / "run_metadata.json"
|
||||
try:
|
||||
metadata = json.loads(metadata_path.read_text(encoding="utf-8"))
|
||||
except (OSError, json.JSONDecodeError):
|
||||
return None
|
||||
value = metadata.get("diarization", {}).get("speaker_count")
|
||||
return value if isinstance(value, int) else None
|
||||
|
||||
@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
|
||||
|
||||
def _update_timer(self, event: Any) -> None:
|
||||
"""Apply a backend progress event to the shared per-run timer."""
|
||||
if event.stage in STAGES and event.status == "started":
|
||||
self.timer.start_stage(event.stage)
|
||||
elif event.stage in STAGES and event.status in {"completed", "skipped"}:
|
||||
self.timer.finish_stage(event.stage)
|
||||
if event.stage in {"completed", "failed"}:
|
||||
self.timer.finish_run()
|
||||
|
||||
def _persist_timing(self, run_dir: Path | None) -> None:
|
||||
"""Merge final monotonic durations into Meeting Lab run metadata."""
|
||||
if run_dir is None:
|
||||
return
|
||||
metadata_path = Path(run_dir) / "run_metadata.json"
|
||||
metadata: dict[str, Any] = {}
|
||||
if metadata_path.is_file():
|
||||
try:
|
||||
existing = json.loads(metadata_path.read_text(encoding="utf-8"))
|
||||
if isinstance(existing, dict):
|
||||
metadata = existing
|
||||
except (OSError, json.JSONDecodeError):
|
||||
return
|
||||
snapshot = self.timer.snapshot()
|
||||
metadata["timing"] = {
|
||||
"stages_seconds": snapshot.stage_durations,
|
||||
"total_seconds": snapshot.total_duration,
|
||||
}
|
||||
metadata_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
try:
|
||||
metadata_path.write_text(
|
||||
json.dumps(metadata, indent=2, sort_keys=True) + "\n", encoding="utf-8"
|
||||
)
|
||||
except OSError:
|
||||
return
|
||||
@@ -0,0 +1,108 @@
|
||||
"""Versioned YAML import and export for reusable People lists."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
from collections.abc import Sequence
|
||||
|
||||
import yaml
|
||||
|
||||
from mka.application.meeting_service import ParticipantInput
|
||||
|
||||
PEOPLE_YAML_VERSION = 1
|
||||
ATTENDANCE_VALUES = frozenset({"present", "mentioned_only"})
|
||||
PARTICIPANT_ID_PATTERN = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.:-]*$")
|
||||
|
||||
|
||||
class PeopleYamlError(ValueError):
|
||||
"""Raised when a reusable People-list document is invalid."""
|
||||
|
||||
|
||||
def export_people_yaml(people: Sequence[ParticipantInput]) -> str:
|
||||
"""Serialize people deterministically without meeting-specific data."""
|
||||
document = {
|
||||
"version": PEOPLE_YAML_VERSION,
|
||||
"people": [
|
||||
{
|
||||
"participant_id": person.participant_id,
|
||||
"display_name": person.display_name,
|
||||
"role": person.role,
|
||||
"organization": person.organization,
|
||||
"attendance_status": person.attendance_status,
|
||||
}
|
||||
for person in people
|
||||
],
|
||||
}
|
||||
return yaml.safe_dump(
|
||||
document,
|
||||
allow_unicode=True,
|
||||
sort_keys=False,
|
||||
default_flow_style=False,
|
||||
)
|
||||
|
||||
|
||||
def import_people_yaml(content: str | bytes) -> list[ParticipantInput]:
|
||||
"""Parse and validate a complete replacement People list."""
|
||||
try:
|
||||
if isinstance(content, bytes):
|
||||
content = content.decode("utf-8")
|
||||
document = yaml.safe_load(content)
|
||||
except (yaml.YAMLError, UnicodeDecodeError) as exc:
|
||||
raise PeopleYamlError(f"Malformed People YAML: {exc}") from exc
|
||||
|
||||
if not isinstance(document, dict):
|
||||
raise PeopleYamlError("People YAML must contain a top-level mapping.")
|
||||
version = document.get("version")
|
||||
if type(version) is not int or version != PEOPLE_YAML_VERSION:
|
||||
raise PeopleYamlError(
|
||||
f"Unsupported People YAML version {version!r}; expected {PEOPLE_YAML_VERSION}."
|
||||
)
|
||||
entries = document.get("people")
|
||||
if not isinstance(entries, list):
|
||||
raise PeopleYamlError("People YAML must contain a top-level 'people' list.")
|
||||
|
||||
people: list[ParticipantInput] = []
|
||||
seen_ids: set[str] = set()
|
||||
for index, entry in enumerate(entries, start=1):
|
||||
if not isinstance(entry, dict):
|
||||
raise PeopleYamlError(f"Person {index} must be a mapping.")
|
||||
participant_id = _required_text(entry, "participant_id", index)
|
||||
if PARTICIPANT_ID_PATTERN.fullmatch(participant_id) is None:
|
||||
raise PeopleYamlError(f"Person {index} has invalid participant_id {participant_id!r}.")
|
||||
if participant_id in seen_ids:
|
||||
raise PeopleYamlError(f"Duplicate participant_id: {participant_id!r}.")
|
||||
seen_ids.add(participant_id)
|
||||
|
||||
display_name = _required_text(entry, "display_name", index)
|
||||
role = _optional_text(entry, "role", index)
|
||||
organization = _optional_text(entry, "organization", index)
|
||||
attendance_status = entry.get("attendance_status")
|
||||
if attendance_status not in ATTENDANCE_VALUES:
|
||||
allowed = ", ".join(sorted(ATTENDANCE_VALUES))
|
||||
raise PeopleYamlError(
|
||||
f"Person {index} has invalid attendance_status; expected one of: {allowed}."
|
||||
)
|
||||
people.append(
|
||||
ParticipantInput(
|
||||
participant_id=participant_id,
|
||||
display_name=display_name,
|
||||
role=role,
|
||||
organization=organization,
|
||||
attendance_status=attendance_status,
|
||||
)
|
||||
)
|
||||
return people
|
||||
|
||||
|
||||
def _required_text(entry: dict[object, object], field: str, index: int) -> str:
|
||||
value = entry.get(field)
|
||||
if not isinstance(value, str) or not value.strip():
|
||||
raise PeopleYamlError(f"Person {index} requires a non-empty {field}.")
|
||||
return value
|
||||
|
||||
|
||||
def _optional_text(entry: dict[object, object], field: str, index: int) -> str:
|
||||
value = entry.get(field, "")
|
||||
if not isinstance(value, str):
|
||||
raise PeopleYamlError(f"Person {index} field {field} must be text.")
|
||||
return value
|
||||
@@ -0,0 +1,23 @@
|
||||
"""Abstract performance profiles resolved to current backend runtime options."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
DEFAULT_PERFORMANCE_PROFILE = "auto"
|
||||
PERFORMANCE_PROFILES = ("auto", "fast", "efficient", "powersave")
|
||||
|
||||
_OLLAMA_THREADS_BY_PROFILE = {
|
||||
"fast": 16,
|
||||
"efficient": 10,
|
||||
"powersave": 4,
|
||||
}
|
||||
|
||||
|
||||
def resolve_ollama_num_thread(profile: str) -> int | None:
|
||||
"""Resolve a profile, leaving automatic thread selection to Ollama for Auto."""
|
||||
if profile == "auto":
|
||||
return None
|
||||
try:
|
||||
return _OLLAMA_THREADS_BY_PROFILE[profile]
|
||||
except KeyError as exc:
|
||||
allowed = ", ".join(PERFORMANCE_PROFILES)
|
||||
raise ValueError(f"Unknown performance profile {profile!r}; expected: {allowed}.") from exc
|
||||
@@ -0,0 +1,80 @@
|
||||
"""Thread-safe monotonic timing state for meeting-processing progress."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from dataclasses import dataclass
|
||||
from threading import Lock
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class TimingSnapshot:
|
||||
"""Presentation-neutral runtime values measured in seconds."""
|
||||
|
||||
stage_durations: dict[str, float]
|
||||
active_stage: str | None
|
||||
total_duration: float
|
||||
running: bool
|
||||
|
||||
|
||||
class ProcessingTimer:
|
||||
"""Track pipeline stages independently from pipeline business logic."""
|
||||
|
||||
def __init__(self, clock: Callable[[], float] = time.perf_counter) -> None:
|
||||
self._clock = clock
|
||||
self._lock = Lock()
|
||||
self._run_started: float | None = None
|
||||
self._run_finished: float | None = None
|
||||
self._stage_started: dict[str, float] = {}
|
||||
self._stage_finished: dict[str, float] = {}
|
||||
self._active_stage: str | None = None
|
||||
|
||||
def start_run(self) -> None:
|
||||
with self._lock:
|
||||
self._run_started = self._clock()
|
||||
self._run_finished = None
|
||||
self._stage_started.clear()
|
||||
self._stage_finished.clear()
|
||||
self._active_stage = None
|
||||
|
||||
def start_stage(self, stage: str) -> None:
|
||||
with self._lock:
|
||||
now = self._clock()
|
||||
if self._active_stage is not None and self._active_stage != stage:
|
||||
self._stage_finished.setdefault(self._active_stage, now)
|
||||
self._stage_started.setdefault(stage, now)
|
||||
self._active_stage = stage
|
||||
|
||||
def finish_stage(self, stage: str) -> None:
|
||||
with self._lock:
|
||||
if stage in self._stage_started:
|
||||
self._stage_finished.setdefault(stage, self._clock())
|
||||
if self._active_stage == stage:
|
||||
self._active_stage = None
|
||||
|
||||
def finish_run(self) -> None:
|
||||
with self._lock:
|
||||
if self._run_finished is not None:
|
||||
return
|
||||
now = self._clock()
|
||||
if self._active_stage is not None:
|
||||
self._stage_finished.setdefault(self._active_stage, now)
|
||||
self._active_stage = None
|
||||
if self._run_started is not None:
|
||||
self._run_finished = now
|
||||
|
||||
def snapshot(self) -> TimingSnapshot:
|
||||
with self._lock:
|
||||
now = self._run_finished if self._run_finished is not None else self._clock()
|
||||
durations = {
|
||||
stage: max(0.0, self._stage_finished.get(stage, now) - started)
|
||||
for stage, started in self._stage_started.items()
|
||||
}
|
||||
total = max(0.0, now - self._run_started) if self._run_started is not None else 0.0
|
||||
return TimingSnapshot(
|
||||
stage_durations=durations,
|
||||
active_stage=self._active_stage,
|
||||
total_duration=total,
|
||||
running=self._run_started is not None and self._run_finished is None,
|
||||
)
|
||||
@@ -0,0 +1,218 @@
|
||||
"""Versioned JSON import/export for the user-facing run input form."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from collections.abc import Sequence
|
||||
from dataclasses import dataclass
|
||||
from datetime import date
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from mka.application.meeting_service import ParticipantInput, stable_id
|
||||
from mka.application.people_yaml import ATTENDANCE_VALUES, PARTICIPANT_ID_PATTERN
|
||||
|
||||
RUN_INPUT_SCHEMA_VERSION = 1
|
||||
SUPPORTED_LANGUAGES = frozenset({"de", "en"})
|
||||
|
||||
|
||||
class RunInputJsonError(ValueError):
|
||||
"""Raised when a run-input JSON document is malformed or unsupported."""
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RunInputState:
|
||||
"""All user-configurable values needed before starting a processing run."""
|
||||
|
||||
title: str
|
||||
description: str
|
||||
language: str
|
||||
meeting_date: date | None
|
||||
participants: tuple[ParticipantInput, ...]
|
||||
audio_normalization: bool
|
||||
diarization_enabled: bool
|
||||
source_file_name: str | None = None
|
||||
|
||||
@classmethod
|
||||
def defaults(cls) -> RunInputState:
|
||||
return cls(
|
||||
title="",
|
||||
description="",
|
||||
language="de",
|
||||
meeting_date=date.today(),
|
||||
participants=(ParticipantInput(participant_id="", display_name=""),),
|
||||
audio_normalization=True,
|
||||
diarization_enabled=False,
|
||||
)
|
||||
|
||||
|
||||
def export_run_inputs(state: RunInputState) -> str:
|
||||
"""Serialize form state as deterministic, human-readable JSON."""
|
||||
source_file_name = _source_file_name(state.source_file_name)
|
||||
document = {
|
||||
"schema_version": RUN_INPUT_SCHEMA_VERSION,
|
||||
"meeting": {
|
||||
"title": state.title,
|
||||
"description": state.description,
|
||||
"language": state.language,
|
||||
"date": state.meeting_date.isoformat() if state.meeting_date else None,
|
||||
"participants": [
|
||||
{
|
||||
"participant_id": person.participant_id,
|
||||
"display_name": person.display_name,
|
||||
"role": person.role,
|
||||
"organization": person.organization,
|
||||
"attendance_status": person.attendance_status,
|
||||
}
|
||||
for person in state.participants
|
||||
],
|
||||
},
|
||||
"processing": {
|
||||
"audio_normalization": state.audio_normalization,
|
||||
"diarization_enabled": state.diarization_enabled,
|
||||
},
|
||||
"source_file_name": source_file_name,
|
||||
}
|
||||
return json.dumps(document, ensure_ascii=False, indent=2) + "\n"
|
||||
|
||||
|
||||
def import_run_inputs(content: str | bytes) -> RunInputState:
|
||||
"""Parse supported fields, applying current defaults to omitted optional fields."""
|
||||
try:
|
||||
if isinstance(content, bytes):
|
||||
content = content.decode("utf-8")
|
||||
document = json.loads(content)
|
||||
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
|
||||
raise RunInputJsonError(f"Malformed input JSON: {exc}") from exc
|
||||
if not isinstance(document, dict):
|
||||
raise RunInputJsonError("Input JSON must contain a top-level object.")
|
||||
version = document.get("schema_version")
|
||||
if type(version) is not int or version != RUN_INPUT_SCHEMA_VERSION:
|
||||
raise RunInputJsonError(
|
||||
f"Unsupported input schema_version {version!r}; expected {RUN_INPUT_SCHEMA_VERSION}."
|
||||
)
|
||||
|
||||
defaults = RunInputState.defaults()
|
||||
meeting = _optional_mapping(document, "meeting")
|
||||
processing = _optional_mapping(document, "processing")
|
||||
title = _optional_text(meeting, "title", defaults.title)
|
||||
description = _optional_text(meeting, "description", defaults.description)
|
||||
language = _optional_text(meeting, "language", defaults.language)
|
||||
if language not in SUPPORTED_LANGUAGES:
|
||||
raise RunInputJsonError(f"Unsupported meeting language: {language!r}.")
|
||||
meeting_date = _meeting_date(meeting, defaults.meeting_date)
|
||||
participants = _participants(meeting, defaults.participants)
|
||||
audio_normalization = _optional_bool(
|
||||
processing, "audio_normalization", defaults.audio_normalization
|
||||
)
|
||||
diarization_enabled = _optional_bool(
|
||||
processing, "diarization_enabled", defaults.diarization_enabled
|
||||
)
|
||||
source_file_name = _source_file_name(document.get("source_file_name"))
|
||||
return RunInputState(
|
||||
title=title,
|
||||
description=description,
|
||||
language=language,
|
||||
meeting_date=meeting_date,
|
||||
participants=participants,
|
||||
audio_normalization=audio_normalization,
|
||||
diarization_enabled=diarization_enabled,
|
||||
source_file_name=source_file_name,
|
||||
)
|
||||
|
||||
|
||||
def run_input_filename(title: str) -> str:
|
||||
"""Build a stable download filename without filesystem-specific characters."""
|
||||
suffix = stable_id(title) if title.strip() else "untitled"
|
||||
return f"meeting-inputs-{suffix}.json"
|
||||
|
||||
|
||||
def _optional_mapping(document: dict[str, Any], field: str) -> dict[str, Any]:
|
||||
value = document.get(field, {})
|
||||
if not isinstance(value, dict):
|
||||
raise RunInputJsonError(f"Field {field!r} must be an object.")
|
||||
return value
|
||||
|
||||
|
||||
def _optional_text(document: dict[str, Any], field: str, default: str) -> str:
|
||||
value = document.get(field, default)
|
||||
if not isinstance(value, str):
|
||||
raise RunInputJsonError(f"Field {field!r} must be text.")
|
||||
return value
|
||||
|
||||
|
||||
def _optional_bool(document: dict[str, Any], field: str, default: bool) -> bool:
|
||||
value = document.get(field, default)
|
||||
if type(value) is not bool:
|
||||
raise RunInputJsonError(f"Field {field!r} must be true or false.")
|
||||
return value
|
||||
|
||||
|
||||
def _meeting_date(document: dict[str, Any], default: date | None) -> date | None:
|
||||
if "date" not in document:
|
||||
return default
|
||||
value = document["date"]
|
||||
if value is None:
|
||||
return None
|
||||
if not isinstance(value, str):
|
||||
raise RunInputJsonError("Field 'date' must be an ISO date or null.")
|
||||
try:
|
||||
return date.fromisoformat(value)
|
||||
except ValueError as exc:
|
||||
raise RunInputJsonError("Field 'date' must be a valid ISO date or null.") from exc
|
||||
|
||||
|
||||
def _participants(
|
||||
document: dict[str, Any], default: Sequence[ParticipantInput]
|
||||
) -> tuple[ParticipantInput, ...]:
|
||||
if "participants" not in document:
|
||||
return tuple(default)
|
||||
entries = document["participants"]
|
||||
if not isinstance(entries, list):
|
||||
raise RunInputJsonError("Field 'participants' must be a list.")
|
||||
people: list[ParticipantInput] = []
|
||||
seen_ids: set[str] = set()
|
||||
for index, entry in enumerate(entries, start=1):
|
||||
if not isinstance(entry, dict):
|
||||
raise RunInputJsonError(f"Participant {index} must be an object.")
|
||||
participant_id = _optional_text(entry, "participant_id", "")
|
||||
display_name = _optional_text(entry, "display_name", "")
|
||||
role = _optional_text(entry, "role", "")
|
||||
organization = _optional_text(entry, "organization", "")
|
||||
attendance = _optional_text(entry, "attendance_status", "present")
|
||||
if bool(participant_id) != bool(display_name):
|
||||
raise RunInputJsonError(
|
||||
f"Participant {index} must provide both participant_id and display_name."
|
||||
)
|
||||
if participant_id and PARTICIPANT_ID_PATTERN.fullmatch(participant_id) is None:
|
||||
raise RunInputJsonError(
|
||||
f"Participant {index} has invalid participant_id {participant_id!r}."
|
||||
)
|
||||
if participant_id in seen_ids:
|
||||
raise RunInputJsonError(f"Duplicate participant_id: {participant_id!r}.")
|
||||
if participant_id:
|
||||
seen_ids.add(participant_id)
|
||||
if attendance not in ATTENDANCE_VALUES:
|
||||
raise RunInputJsonError(
|
||||
f"Participant {index} has invalid attendance_status {attendance!r}."
|
||||
)
|
||||
people.append(
|
||||
ParticipantInput(
|
||||
participant_id=participant_id,
|
||||
display_name=display_name,
|
||||
role=role,
|
||||
organization=organization,
|
||||
attendance_status=attendance,
|
||||
)
|
||||
)
|
||||
return tuple(people)
|
||||
|
||||
|
||||
def _source_file_name(value: Any) -> str | None:
|
||||
if value is None:
|
||||
return None
|
||||
if not isinstance(value, str) or not value.strip():
|
||||
raise RunInputJsonError("Field 'source_file_name' must be non-empty text or null.")
|
||||
if Path(value).name != value:
|
||||
raise RunInputJsonError("Field 'source_file_name' must be a filename, not a path.")
|
||||
return value
|
||||
@@ -0,0 +1 @@
|
||||
"""Adapters for external processing systems."""
|
||||
@@ -0,0 +1,68 @@
|
||||
"""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,
|
||||
regenerate_mvp_protocol,
|
||||
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
|
||||
self._regenerate_mvp_protocol = regenerate_mvp_protocol
|
||||
|
||||
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,
|
||||
)
|
||||
|
||||
def regenerate_protocol(
|
||||
self,
|
||||
run_dir: Any,
|
||||
meeting_context: Any,
|
||||
progress_sink: Callable[[Any], None],
|
||||
**options: Any,
|
||||
) -> Any:
|
||||
"""Regenerate protocol artifacts without rerunning media processing."""
|
||||
return self._regenerate_mvp_protocol(
|
||||
run_dir,
|
||||
meeting_context=meeting_context,
|
||||
progress_sink=progress_sink,
|
||||
**options,
|
||||
)
|
||||
@@ -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)
|
||||
|
||||
@@ -0,0 +1,785 @@
|
||||
"""Streamlit presentation layer for the first Meeting Assistant MVP."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from datetime import date
|
||||
from queue import Empty, Queue
|
||||
from typing import Any
|
||||
from uuid import uuid4
|
||||
|
||||
import streamlit as st
|
||||
|
||||
from mka.application.config import AppSettings, ConfigurationError
|
||||
from mka.application.glossary import (
|
||||
GLOSSARY_CATEGORIES,
|
||||
GlossaryConflictError,
|
||||
GlossaryRepository,
|
||||
)
|
||||
from mka.application.glossary_yaml import (
|
||||
GlossaryYamlError,
|
||||
export_glossary_yaml,
|
||||
import_glossary_yaml,
|
||||
)
|
||||
from mka.application.meeting_service import (
|
||||
STAGES,
|
||||
AppProgressEvent,
|
||||
MeetingDetails,
|
||||
MeetingProcessingService,
|
||||
ParticipantInput,
|
||||
ProcessingOptions,
|
||||
SpeakerMappingReview,
|
||||
stable_id,
|
||||
)
|
||||
from mka.application.people_yaml import (
|
||||
PeopleYamlError,
|
||||
export_people_yaml,
|
||||
import_people_yaml,
|
||||
)
|
||||
from mka.application.performance import DEFAULT_PERFORMANCE_PROFILE, PERFORMANCE_PROFILES
|
||||
from mka.application.progress_timing import TimingSnapshot
|
||||
from mka.application.run_inputs import (
|
||||
RunInputJsonError,
|
||||
RunInputState,
|
||||
export_run_inputs,
|
||||
import_run_inputs,
|
||||
run_input_filename,
|
||||
)
|
||||
from mka.integrations.meeting_lab import (
|
||||
MeetingLabGateway,
|
||||
MeetingLabUnavailableError,
|
||||
)
|
||||
|
||||
STAGE_LABELS = {
|
||||
"preparing": "Preparation",
|
||||
"transcription": "Transcription",
|
||||
"diarization": "Diarization",
|
||||
"protocol_generation": "Protocol generation",
|
||||
}
|
||||
|
||||
|
||||
def _parse_aliases(value: str) -> tuple[str, ...]:
|
||||
"""Parse one alias per line while tolerating comma-separated input."""
|
||||
return tuple(
|
||||
alias.strip() for line in value.splitlines() for alias in line.split(",") if alias.strip()
|
||||
)
|
||||
|
||||
|
||||
def _render_glossary(repository: GlossaryRepository) -> None:
|
||||
"""Render simple global glossary CRUD controls."""
|
||||
with st.expander("Terminology glossary", expanded=False):
|
||||
st.caption(
|
||||
"Store canonical core terms and recognition aliases. Compound phrases are "
|
||||
"composed from meeting context during protocol generation."
|
||||
)
|
||||
import_file = st.file_uploader(
|
||||
"Import glossary",
|
||||
type=["yaml", "yml"],
|
||||
key="glossary_import_file",
|
||||
help="Validate and atomically replace the SQLite glossary from versioned YAML.",
|
||||
)
|
||||
action_columns = st.columns(2)
|
||||
if action_columns[0].button("Import glossary", disabled=import_file is None):
|
||||
try:
|
||||
imported = import_glossary_yaml(import_file.getvalue())
|
||||
repository.replace_all(imported)
|
||||
except (GlossaryYamlError, GlossaryConflictError, ValueError) as exc:
|
||||
st.error(str(exc))
|
||||
else:
|
||||
st.session_state.glossary_import_message = (
|
||||
f"Imported {len(imported)} glossary entries."
|
||||
)
|
||||
st.rerun()
|
||||
action_columns[1].download_button(
|
||||
"Export glossary",
|
||||
data=export_glossary_yaml(repository.list()).encode("utf-8"),
|
||||
file_name="terminology-glossary.yaml",
|
||||
mime="application/yaml",
|
||||
)
|
||||
if message := st.session_state.pop("glossary_import_message", None):
|
||||
st.success(message)
|
||||
|
||||
with st.form("glossary_add"):
|
||||
columns = st.columns(2)
|
||||
canonical = columns[0].text_input("Canonical term")
|
||||
category = columns[1].selectbox("Category", GLOSSARY_CATEGORIES)
|
||||
aliases = st.text_area("Aliases", help="One per line or comma-separated.")
|
||||
description = st.text_area("Optional description")
|
||||
if st.form_submit_button("Add glossary entry"):
|
||||
try:
|
||||
repository.create(
|
||||
canonical,
|
||||
category,
|
||||
aliases=_parse_aliases(aliases),
|
||||
description=description,
|
||||
)
|
||||
except (GlossaryConflictError, ValueError) as exc:
|
||||
st.error(str(exc))
|
||||
else:
|
||||
st.success("Glossary entry added.")
|
||||
st.rerun()
|
||||
|
||||
search = st.text_input("Search glossary")
|
||||
active_only = st.checkbox("Show active entries only")
|
||||
entries = repository.list(search, active_only=active_only)
|
||||
st.caption(f"{len(entries)} glossary entries")
|
||||
for entry in entries:
|
||||
status = "active" if entry.is_active else "inactive"
|
||||
with (
|
||||
st.expander(f"{entry.canonical_term} · {entry.category} · {status}"),
|
||||
st.form(f"glossary_edit_{entry.id}"),
|
||||
):
|
||||
columns = st.columns(2)
|
||||
edited_canonical = columns[0].text_input("Canonical term", entry.canonical_term)
|
||||
edited_category = columns[1].selectbox(
|
||||
"Category",
|
||||
GLOSSARY_CATEGORIES,
|
||||
index=GLOSSARY_CATEGORIES.index(entry.category),
|
||||
)
|
||||
edited_aliases = st.text_area("Aliases", "\n".join(entry.aliases))
|
||||
edited_description = st.text_area("Optional description", entry.description or "")
|
||||
edited_active = st.checkbox("Active", value=entry.is_active)
|
||||
delete_confirmed = st.checkbox(
|
||||
"Permanently delete this entry",
|
||||
help="Prefer clearing Active for normal use.",
|
||||
)
|
||||
action_columns = st.columns(2)
|
||||
save = action_columns[0].form_submit_button("Save changes")
|
||||
delete = action_columns[1].form_submit_button(
|
||||
"Delete", disabled=not delete_confirmed
|
||||
)
|
||||
try:
|
||||
if save:
|
||||
repository.update(
|
||||
entry.id,
|
||||
edited_canonical,
|
||||
edited_category,
|
||||
aliases=_parse_aliases(edited_aliases),
|
||||
description=edited_description,
|
||||
is_active=edited_active,
|
||||
)
|
||||
st.rerun()
|
||||
if delete:
|
||||
repository.delete(entry.id)
|
||||
st.rerun()
|
||||
except (GlossaryConflictError, ValueError) as exc:
|
||||
st.error(str(exc))
|
||||
|
||||
|
||||
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", "")
|
||||
st.session_state.setdefault("source_media_widget_generation", 0)
|
||||
_apply_pending_run_inputs()
|
||||
_apply_pending_edited_protocol()
|
||||
|
||||
|
||||
def _queue_run_inputs(value: RunInputState) -> None:
|
||||
"""Defer form restoration until the beginning of the next Streamlit run."""
|
||||
st.session_state["pending_run_inputs"] = value
|
||||
|
||||
|
||||
def _apply_pending_run_inputs() -> None:
|
||||
"""Restore imported values before any corresponding widget is instantiated."""
|
||||
value = st.session_state.pop("pending_run_inputs", None)
|
||||
if value is None:
|
||||
return
|
||||
st.session_state.update(
|
||||
{
|
||||
"meeting_title": value.title,
|
||||
"meeting_description": value.description,
|
||||
"meeting_language": value.language,
|
||||
"meeting_has_date": value.meeting_date is not None,
|
||||
"meeting_date": value.meeting_date or date.today(),
|
||||
"participants": _people_to_rows(list(value.participants)),
|
||||
"audio_normalization": value.audio_normalization,
|
||||
"diarization_enabled": value.diarization_enabled,
|
||||
"imported_source_file_name": value.source_file_name,
|
||||
"run_input_import_message": (
|
||||
"Meeting inputs imported. Select the source media before processing."
|
||||
),
|
||||
"source_media_widget_generation": (
|
||||
st.session_state.get("source_media_widget_generation", 0) + 1
|
||||
),
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def _apply_pending_edited_protocol() -> None:
|
||||
"""Apply a deferred widget value before the widget is instantiated."""
|
||||
if "pending_edited_protocol" in st.session_state:
|
||||
st.session_state["edited_protocol"] = st.session_state.pop("pending_edited_protocol")
|
||||
|
||||
|
||||
def _queue_edited_protocol(value: str) -> None:
|
||||
"""Defer an edited-protocol widget update until the next Streamlit run."""
|
||||
st.session_state["pending_edited_protocol"] = value
|
||||
|
||||
|
||||
def _people_to_rows(people: list[ParticipantInput]) -> list[dict[str, str]]:
|
||||
"""Create fresh widget rows while preserving reusable person IDs."""
|
||||
return [
|
||||
{
|
||||
"row_id": uuid4().hex,
|
||||
"participant_id": person.participant_id,
|
||||
"display_name": person.display_name,
|
||||
"role": person.role,
|
||||
"organization": person.organization,
|
||||
"attendance_status": person.attendance_status,
|
||||
}
|
||||
for person in people
|
||||
]
|
||||
|
||||
|
||||
def _render_run_input_import() -> None:
|
||||
"""Render import controls before the widgets whose state they restore."""
|
||||
st.subheader("Input configuration")
|
||||
imported_file = st.file_uploader(
|
||||
"Import inputs",
|
||||
type=["json"],
|
||||
key="run_input_import_file",
|
||||
help="Restore form values only; source media is never included.",
|
||||
)
|
||||
if st.button("Import input configuration", disabled=imported_file is None):
|
||||
try:
|
||||
imported = import_run_inputs(imported_file.getvalue())
|
||||
except RunInputJsonError as exc:
|
||||
st.error(str(exc))
|
||||
else:
|
||||
_queue_run_inputs(imported)
|
||||
st.rerun()
|
||||
if message := st.session_state.pop("run_input_import_message", None):
|
||||
st.success(message)
|
||||
|
||||
|
||||
def _render_participants() -> list[ParticipantInput]:
|
||||
st.subheader("People")
|
||||
st.caption("Record whether each relevant person attended or was only mentioned.")
|
||||
imported_file = st.file_uploader(
|
||||
"Import people",
|
||||
type=["yaml", "yml"],
|
||||
help="Replace the current People list from a versioned YAML export.",
|
||||
)
|
||||
if st.button("Import people list", disabled=imported_file is None):
|
||||
try:
|
||||
imported_people = import_people_yaml(imported_file.getvalue())
|
||||
except PeopleYamlError as exc:
|
||||
st.error(str(exc))
|
||||
else:
|
||||
st.session_state.participants = _people_to_rows(imported_people)
|
||||
st.session_state.people_import_message = f"Imported {len(imported_people)} people."
|
||||
st.rerun()
|
||||
if message := st.session_state.pop("people_import_message", None):
|
||||
st.success(message)
|
||||
|
||||
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()
|
||||
action_columns = st.columns(2)
|
||||
if action_columns[0].button("Add person"):
|
||||
rows.append(_new_participant())
|
||||
st.rerun()
|
||||
people = [
|
||||
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
|
||||
]
|
||||
exported_yaml = export_people_yaml(people).encode("utf-8")
|
||||
action_columns[1].download_button(
|
||||
"Export people",
|
||||
data=exported_yaml,
|
||||
file_name="people.yaml",
|
||||
mime="application/yaml",
|
||||
)
|
||||
return people
|
||||
|
||||
|
||||
def _format_duration(seconds: float) -> str:
|
||||
"""Format numeric seconds as MM:SS or H:MM:SS."""
|
||||
whole_seconds = max(0, int(seconds))
|
||||
hours, remainder = divmod(whole_seconds, 3600)
|
||||
minutes, seconds = divmod(remainder, 60)
|
||||
return f"{hours}:{minutes:02d}:{seconds:02d}" if hours else f"{minutes:02d}:{seconds:02d}"
|
||||
|
||||
|
||||
def _render_progress_status(
|
||||
status_box: Any,
|
||||
stage_table: Any,
|
||||
states: dict[str, str],
|
||||
timing: TimingSnapshot,
|
||||
message: str,
|
||||
) -> None:
|
||||
suffix = " …" if timing.running else ""
|
||||
status_box.info(f"{message} — total runtime {_format_duration(timing.total_duration)}{suffix}")
|
||||
stage_table.table(
|
||||
[
|
||||
{
|
||||
"Stage": STAGE_LABELS[stage],
|
||||
"Status": states[stage],
|
||||
"Runtime": (
|
||||
_format_duration(timing.stage_durations[stage])
|
||||
+ (" …" if timing.active_stage == stage else "")
|
||||
if stage in timing.stage_durations
|
||||
else ""
|
||||
),
|
||||
}
|
||||
for stage in STAGES
|
||||
]
|
||||
+ [
|
||||
{
|
||||
"Stage": "Total runtime",
|
||||
"Status": "",
|
||||
"Runtime": _format_duration(timing.total_duration) + suffix,
|
||||
}
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def _apply_progress_event(event: AppProgressEvent, states: dict[str, str]) -> None:
|
||||
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"
|
||||
|
||||
|
||||
def _run_with_live_progress(
|
||||
service: MeetingProcessingService,
|
||||
work: Callable[[Callable[[AppProgressEvent], None]], Any],
|
||||
states: dict[str, str],
|
||||
initial_message: str,
|
||||
) -> tuple[Any, dict[str, str], str, TimingSnapshot]:
|
||||
"""Run pipeline work while rendering the shared timer and progress events."""
|
||||
status_box = st.empty()
|
||||
stage_table = st.empty()
|
||||
progress_slot = st.empty()
|
||||
event_queue: Queue[AppProgressEvent] = Queue()
|
||||
latest_message = initial_message
|
||||
progress_bar = None
|
||||
with ThreadPoolExecutor(max_workers=1) as executor:
|
||||
future = executor.submit(work, event_queue.put)
|
||||
while not future.done():
|
||||
try:
|
||||
while True:
|
||||
event = event_queue.get_nowait()
|
||||
_apply_progress_event(event, states)
|
||||
latest_message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||
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)}: "
|
||||
f"{event.progress:.0%}"
|
||||
),
|
||||
)
|
||||
except Empty:
|
||||
pass
|
||||
_render_progress_status(
|
||||
status_box,
|
||||
stage_table,
|
||||
states,
|
||||
service.timer.snapshot(),
|
||||
latest_message,
|
||||
)
|
||||
time.sleep(0.2)
|
||||
try:
|
||||
result = future.result()
|
||||
except Exception:
|
||||
while not event_queue.empty():
|
||||
event = event_queue.get_nowait()
|
||||
_apply_progress_event(event, states)
|
||||
latest_message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||
service.timer.finish_run()
|
||||
running_stage = next(
|
||||
(stage for stage, status in states.items() if status == "running"),
|
||||
None,
|
||||
)
|
||||
if running_stage is not None:
|
||||
states[running_stage] = "failed"
|
||||
snapshot = service.timer.snapshot()
|
||||
_render_progress_status(status_box, stage_table, states, snapshot, latest_message)
|
||||
st.session_state["processing_progress"] = (states.copy(), latest_message, snapshot)
|
||||
raise
|
||||
while not event_queue.empty():
|
||||
event = event_queue.get_nowait()
|
||||
_apply_progress_event(event, states)
|
||||
latest_message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||
running_stage = next((stage for stage, status in states.items() if status == "running"), None)
|
||||
if running_stage is not None:
|
||||
states[running_stage] = "completed" if getattr(result, "succeeded", True) else "failed"
|
||||
snapshot = service.timer.snapshot()
|
||||
_render_progress_status(status_box, stage_table, states, snapshot, latest_message)
|
||||
st.session_state["processing_progress"] = (states.copy(), latest_message, snapshot)
|
||||
return result, states, latest_message, snapshot
|
||||
|
||||
|
||||
def _remember_regeneration_timing(
|
||||
states: dict[str, str], message: str, timing: TimingSnapshot
|
||||
) -> None:
|
||||
st.session_state["regeneration_progress"] = (states.copy(), message, timing)
|
||||
|
||||
|
||||
def _speaker_options(
|
||||
participants: tuple[str, ...], mappings: dict[str, str | None], speaker_label: str
|
||||
) -> list[str | None]:
|
||||
"""Reserve other speakers' participants while retaining this speaker's mapping."""
|
||||
current = mappings.get(speaker_label)
|
||||
reserved = {value for label, value in mappings.items() if label != speaker_label}
|
||||
options: list[str | None] = [None]
|
||||
options.extend(person for person in participants if person == current or person not in reserved)
|
||||
if current is not None and current not in options:
|
||||
options.append(current)
|
||||
return options
|
||||
|
||||
|
||||
def _mapping_counts(
|
||||
speaker_labels: tuple[str, ...], mappings: dict[str, str | None]
|
||||
) -> tuple[int, int, int]:
|
||||
"""Count only detected speakers, including explicit cleared selections."""
|
||||
detected = len(speaker_labels)
|
||||
assigned = sum(mappings.get(label) is not None for label in speaker_labels)
|
||||
return detected, assigned, detected - assigned
|
||||
|
||||
|
||||
def _render_speaker_mapping(review: SpeakerMappingReview, run_name: str) -> dict[str, str]:
|
||||
st.subheader("Identify diarized speakers")
|
||||
st.caption(
|
||||
"Confirm identities explicitly. Unmapped speakers remain anonymous; "
|
||||
"the diarized source transcript is not modified."
|
||||
)
|
||||
participant_names = dict(review.participants)
|
||||
# Read every widget before rendering so later speakers also reserve their person.
|
||||
keys = {
|
||||
speaker.speaker_label: f"speaker_mapping_{run_name}_{speaker.speaker_label}"
|
||||
for speaker in review.speakers
|
||||
}
|
||||
mappings = {
|
||||
label: st.session_state.get(key, review.current_mappings.get(label))
|
||||
for label, key in keys.items()
|
||||
}
|
||||
detected, assigned, unassigned = _mapping_counts(tuple(keys), mappings)
|
||||
st.markdown(f"**{detected} speakers detected · {assigned} assigned · {unassigned} unassigned**")
|
||||
if unassigned:
|
||||
st.warning(
|
||||
"Some detected speakers have no confirmed participant mapping. "
|
||||
"Check whether a participant is missing or speaker assignment is incomplete. "
|
||||
"You can still generate a protocol with anonymous speakers."
|
||||
)
|
||||
selections: dict[str, str] = {}
|
||||
for speaker in review.speakers:
|
||||
label = speaker.speaker_label
|
||||
current = mappings[label]
|
||||
options = _speaker_options(tuple(participant_names), mappings, label)
|
||||
st.session_state[keys[label]] = current
|
||||
selected = st.selectbox(
|
||||
f"{label} — Unassigned" if current is None else label,
|
||||
options=options,
|
||||
format_func=lambda value, names=participant_names: (
|
||||
"Unmapped / Unknown" if value is None else names.get(value, value)
|
||||
),
|
||||
key=keys[label],
|
||||
)
|
||||
if selected is not None:
|
||||
selections[label] = selected
|
||||
for excerpt in speaker.excerpts:
|
||||
st.caption(f"“{excerpt}”")
|
||||
return selections
|
||||
|
||||
|
||||
def _render_result() -> None:
|
||||
outcome = st.session_state.outcome
|
||||
if outcome is None:
|
||||
return
|
||||
st.divider()
|
||||
st.header("Protocol result")
|
||||
if retained_progress := st.session_state.get("processing_progress"):
|
||||
states, message, timing = retained_progress
|
||||
_render_progress_status(st.empty(), st.empty(), states, timing, message)
|
||||
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
|
||||
|
||||
awaiting_review = getattr(outcome, "awaiting_speaker_review", False)
|
||||
if awaiting_review:
|
||||
count = getattr(outcome, "detected_speaker_count", None)
|
||||
detected = f"{count} speakers detected" if count is not None else "speakers detected"
|
||||
st.info(
|
||||
f"Diarization complete — {detected}. Assign speakers if desired, then generate the protocol."
|
||||
)
|
||||
else:
|
||||
st.success("Processing completed. Review the generated protocol before use.")
|
||||
st.caption(f"Run artifacts: {outcome.run_dir}")
|
||||
if outcome.speaker_attribution_available is False:
|
||||
st.warning(
|
||||
"Speaker attribution was unavailable for this protocol because the "
|
||||
"safe-budget fallback removed diarization labels from the prompt."
|
||||
)
|
||||
if message := st.session_state.pop("speaker_mapping_message", None):
|
||||
st.success(message)
|
||||
if not awaiting_review:
|
||||
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}")
|
||||
|
||||
try:
|
||||
service = MeetingProcessingService(AppSettings.from_environment(), MeetingLabGateway())
|
||||
review = service.load_speaker_mapping_review(outcome.run_dir)
|
||||
except (MeetingLabUnavailableError, OSError, ValueError) as exc:
|
||||
st.warning(f"Speaker mapping is unavailable: {exc}")
|
||||
return
|
||||
if review is None or not review.speakers:
|
||||
return
|
||||
|
||||
selections = _render_speaker_mapping(review, outcome.run_dir.name)
|
||||
|
||||
duplicate_assignments = len(selections.values()) != len(set(selections.values()))
|
||||
if duplicate_assignments:
|
||||
st.error("Each participant can be assigned to only one detected speaker.")
|
||||
if st.button(
|
||||
"Generate protocol" if awaiting_review else "Regenerate protocol with confirmed speakers",
|
||||
disabled=duplicate_assignments,
|
||||
type="primary",
|
||||
):
|
||||
st.session_state.pop("regeneration_progress", None)
|
||||
run_dir = outcome.run_dir
|
||||
mapping_items = tuple(selections.items())
|
||||
selected_profile = st.session_state.get("performance_profile", DEFAULT_PERFORMANCE_PROFILE)
|
||||
states = {stage: "skipped" for stage in STAGES}
|
||||
states["protocol_generation"] = "pending"
|
||||
try:
|
||||
regenerated, states, message, timing = _run_with_live_progress(
|
||||
service,
|
||||
lambda progress_sink: service.regenerate_protocol(
|
||||
run_dir,
|
||||
dict(mapping_items),
|
||||
progress_sink=progress_sink,
|
||||
performance_profile=selected_profile,
|
||||
),
|
||||
states,
|
||||
"Starting protocol generation" if awaiting_review else "Starting protocol regeneration",
|
||||
)
|
||||
except (OSError, RuntimeError, ValueError) as exc:
|
||||
_remember_regeneration_timing(
|
||||
states,
|
||||
"Protocol regeneration failed",
|
||||
service.timer.snapshot(),
|
||||
)
|
||||
st.error(f"Protocol generation failed: {exc}")
|
||||
else:
|
||||
_remember_regeneration_timing(states, message, timing)
|
||||
st.session_state.outcome = regenerated
|
||||
_queue_edited_protocol(regenerated.original_protocol or "")
|
||||
st.session_state.speaker_mapping_message = (
|
||||
"Speaker mappings saved and protocol generated."
|
||||
)
|
||||
st.rerun()
|
||||
|
||||
|
||||
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.")
|
||||
|
||||
settings = AppSettings.from_environment()
|
||||
glossary = GlossaryRepository(settings.glossary_database)
|
||||
glossary.initialize()
|
||||
_render_glossary(glossary)
|
||||
|
||||
_render_run_input_import()
|
||||
|
||||
st.header("Meeting and audio")
|
||||
audio = st.file_uploader(
|
||||
"Audio recording",
|
||||
type=["wav", "flac", "m4a"],
|
||||
key=f"source_media_{st.session_state.source_media_widget_generation}",
|
||||
)
|
||||
imported_source_name = st.session_state.get("imported_source_file_name")
|
||||
if audio is None and imported_source_name:
|
||||
st.caption(
|
||||
f"Previous source filename: {imported_source_name}. "
|
||||
"Select the media file again before processing."
|
||||
)
|
||||
title = st.text_input("Meeting title", key="meeting_title")
|
||||
description = st.text_area("Description / context", height=100, key="meeting_description")
|
||||
metadata_columns = st.columns(3)
|
||||
language = metadata_columns[0].selectbox(
|
||||
"Meeting language",
|
||||
options=["de", "en"],
|
||||
key="meeting_language",
|
||||
help="Language for both transcription and the generated protocol.",
|
||||
)
|
||||
has_date = metadata_columns[1].checkbox(
|
||||
"Meeting date is known", value=True, key="meeting_has_date"
|
||||
)
|
||||
selected_date: date | None = (
|
||||
metadata_columns[2].date_input("Meeting date", key="meeting_date") if has_date else None
|
||||
)
|
||||
|
||||
participants = _render_participants()
|
||||
|
||||
st.header("Processing options")
|
||||
performance_profile = st.selectbox(
|
||||
"Performance profile",
|
||||
options=PERFORMANCE_PROFILES,
|
||||
format_func=str.title,
|
||||
key="performance_profile",
|
||||
help="Controls protocol-generation performance using an abstract runtime profile.",
|
||||
)
|
||||
audio_normalization = st.checkbox(
|
||||
"Audio normalization",
|
||||
value=True,
|
||||
key="audio_normalization",
|
||||
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",
|
||||
key="diarization_enabled",
|
||||
help="Speaker labels remain anonymous; identities are never inferred.",
|
||||
)
|
||||
|
||||
current_inputs = RunInputState(
|
||||
title=title,
|
||||
description=description,
|
||||
language=language,
|
||||
meeting_date=selected_date,
|
||||
participants=tuple(participants),
|
||||
audio_normalization=audio_normalization,
|
||||
diarization_enabled=diarization_enabled,
|
||||
source_file_name=(audio.name if audio is not None else imported_source_name),
|
||||
)
|
||||
st.download_button(
|
||||
"Export inputs",
|
||||
data=export_run_inputs(current_inputs).encode("utf-8"),
|
||||
file_name=run_input_filename(title),
|
||||
mime="application/json",
|
||||
help="Download the current form configuration without source media or run results.",
|
||||
)
|
||||
|
||||
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:
|
||||
try:
|
||||
service = MeetingProcessingService(settings, MeetingLabGateway(), glossary)
|
||||
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.session_state.pop("regeneration_progress", None)
|
||||
st.session_state.pop("processing_progress", None)
|
||||
st.header("Processing status")
|
||||
states = {stage: "pending" for stage in STAGES}
|
||||
if not diarization_enabled:
|
||||
states["diarization"] = "skipped"
|
||||
outcome, _, _, _ = _run_with_live_progress(
|
||||
service,
|
||||
lambda progress_sink: service.process(
|
||||
audio_path,
|
||||
meeting,
|
||||
participants,
|
||||
ProcessingOptions(
|
||||
diarization_enabled=diarization_enabled,
|
||||
audio_normalization=audio_normalization,
|
||||
performance_profile=performance_profile,
|
||||
),
|
||||
progress_sink=progress_sink,
|
||||
),
|
||||
states,
|
||||
"Starting processing",
|
||||
)
|
||||
st.session_state.outcome = outcome
|
||||
st.session_state.edited_protocol = outcome.original_protocol or ""
|
||||
if outcome.succeeded:
|
||||
st.success("Processing completed.")
|
||||
else:
|
||||
st.error(
|
||||
f"Processing failed during {outcome.failed_stage}: "
|
||||
f"{outcome.error_message} 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()
|
||||
@@ -0,0 +1,138 @@
|
||||
"""Paired Assistant/Lab smoke with only expensive media/model boundaries mocked."""
|
||||
|
||||
import json
|
||||
import shutil
|
||||
from dataclasses import replace
|
||||
from unittest.mock import Mock
|
||||
|
||||
import pytest
|
||||
|
||||
from mka.application.meeting_service import MeetingProcessingService, ProcessingOptions
|
||||
from mka.integrations.meeting_lab import MeetingLabGateway
|
||||
from test_meeting_service import make_service, meeting, participants
|
||||
|
||||
mvp = pytest.importorskip("src.meeting_lab.orchestration.mvp")
|
||||
from src.meeting_lab.audio import PreparedAudio # noqa: E402
|
||||
from src.meeting_lab.diarization.backend import DiarizationResult # noqa: E402
|
||||
from src.meeting_lab.llm.ollama import OllamaGeneration # noqa: E402
|
||||
from src.meeting_lab.protocol.generate_direct_protocol import generate_direct_protocol # noqa: E402
|
||||
from src.meeting_lab.transcription.whisper import TranscriptionResult # noqa: E402
|
||||
|
||||
|
||||
def fake_prepare(source, destination, **kwargs):
|
||||
destination.parent.mkdir(parents=True, exist_ok=True)
|
||||
shutil.copyfile(source, destination)
|
||||
return PreparedAudio(source, source.suffix[1:], destination, "ffmpeg", "ffmpeg")
|
||||
|
||||
|
||||
def fake_transcribe(audio, model, output, language, **kwargs):
|
||||
output.mkdir(parents=True, exist_ok=True)
|
||||
raw, transcript, text, metadata = [
|
||||
output / n
|
||||
for n in ("whisper_raw.json", "transcript.json", "transcript.txt", "runtime_metadata.json")
|
||||
]
|
||||
raw.write_text("{}")
|
||||
transcript.write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"text": "Lumini project discussion.",
|
||||
"segments": [{"id": 0, "start": 0, "end": 2, "text": "Lumini project discussion."}],
|
||||
}
|
||||
)
|
||||
)
|
||||
text.write_text("Lumini project discussion.")
|
||||
metadata.write_text("{}")
|
||||
return TranscriptionResult(output, raw, transcript, text, metadata, 0.1)
|
||||
|
||||
|
||||
def fake_diarize(audio, output, mode, **kwargs):
|
||||
output.mkdir(parents=True, exist_ok=True)
|
||||
paths = [
|
||||
output / name
|
||||
for name in (
|
||||
"metadata.json",
|
||||
"diarization.rttm",
|
||||
"exclusive_diarization.rttm",
|
||||
"turns.json",
|
||||
"exclusive_turns.json",
|
||||
)
|
||||
]
|
||||
metadata = {"speaker_count": 1, "runtime_seconds": 0.1}
|
||||
paths[0].write_text(json.dumps(metadata))
|
||||
for p in paths[1:]:
|
||||
p.write_text("[]")
|
||||
paths[-1].write_text(json.dumps([{"start": 0, "end": 2, "speaker_id": "SPEAKER_00"}]))
|
||||
return DiarizationResult(output, *paths, metadata)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("language,expected", [("de", "German"), ("en", "English")])
|
||||
@pytest.mark.parametrize(
|
||||
"profile,threads", [("auto", None), ("fast", 16), ("efficient", 10), ("powersave", 4)]
|
||||
)
|
||||
def test_paired_anonymous_generation_and_mapped_regeneration(
|
||||
tmp_path, monkeypatch, language, expected, profile, threads
|
||||
):
|
||||
template, _ = make_service(tmp_path)
|
||||
service = MeetingProcessingService(template.settings, MeetingLabGateway())
|
||||
service.glossary.create("Luminy", "product", aliases=("Lumini",))
|
||||
calls = []
|
||||
|
||||
def model_call(endpoint, model, prompt, **kwargs):
|
||||
calls.append((prompt, kwargs))
|
||||
return OllamaGeneration(
|
||||
{"response": "# Mock protocol", "done": True}, "# Mock protocol", 0.1
|
||||
)
|
||||
|
||||
def generate(transcript, context, **kwargs):
|
||||
return generate_direct_protocol(
|
||||
transcript, context, **kwargs, model_check=lambda *_: {}, generation_call=model_call
|
||||
)
|
||||
|
||||
prepare = Mock(side_effect=fake_prepare)
|
||||
transcribe = Mock(side_effect=fake_transcribe)
|
||||
diarize = Mock(side_effect=fake_diarize)
|
||||
monkeypatch.setattr(mvp, "prepare_audio", prepare)
|
||||
monkeypatch.setattr(mvp, "transcribe_audio", transcribe)
|
||||
monkeypatch.setattr(mvp, "diarize_audio", diarize)
|
||||
monkeypatch.setattr(mvp, "generate_direct_protocol", generate)
|
||||
audio = tmp_path / "sample.wav"
|
||||
audio.write_bytes(b"mock audio")
|
||||
outcome = service.process(
|
||||
audio,
|
||||
replace(meeting(), language=language),
|
||||
participants(),
|
||||
ProcessingOptions(diarization_enabled=True, performance_profile=profile),
|
||||
)
|
||||
assert outcome.succeeded
|
||||
assert outcome.awaiting_speaker_review
|
||||
assert prepare.call_args.kwargs["normalization_enabled"] is True
|
||||
assert transcribe.call_args.args[3] == language
|
||||
root = outcome.run_dir
|
||||
assert calls == []
|
||||
service.regenerate_protocol(root, {}, performance_profile=profile)
|
||||
first = root / "protocol/generations/001"
|
||||
before = {p.name: p.read_bytes() for p in first.iterdir()}
|
||||
original = (root / "diarization/transcript_diarized.json").read_bytes()
|
||||
first_meta = json.loads((first / "runtime_metadata.json").read_text())
|
||||
assert first_meta["speaker_mapping"] == {}
|
||||
assert first_meta["output_language"] == language
|
||||
assert "SPEAKER_00" in (first / "transcript_input.txt").read_text()
|
||||
assert "timing" in json.loads((root / "run_metadata.json").read_text())
|
||||
service.regenerate_protocol(root, {"SPEAKER_00": "martin"}, performance_profile=profile)
|
||||
assert prepare.call_count == transcribe.call_count == diarize.call_count == 1
|
||||
assert (root / "diarization/transcript_diarized.json").read_bytes() == original
|
||||
assert {p.name: p.read_bytes() for p in first.iterdir()} == before
|
||||
second = root / "protocol/generations/002"
|
||||
metadata = json.loads((second / "runtime_metadata.json").read_text())
|
||||
assert metadata["speaker_mapping"] == {"SPEAKER_00": "martin"}
|
||||
assert metadata["speaker_mapping_names"] == {"SPEAKER_00": "Martin"}
|
||||
assert metadata["output_language"] == language
|
||||
assert metadata["num_thread"] == threads
|
||||
assert metadata["glossary_aliases_configured"]["Lumini"] == "Luminy"
|
||||
assert metadata["glossary_replacements"] == []
|
||||
assert "Lumini" in (second / "transcript_input.txt").read_text()
|
||||
assert (root / "protocol.md").resolve() == second / "protocol.md"
|
||||
for prompt, options in calls:
|
||||
assert f"Write the meeting protocol in {expected}." in prompt
|
||||
assert "Luminy" in prompt
|
||||
assert options["num_thread"] == threads
|
||||
@@ -0,0 +1,76 @@
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from mka.application.config import AppSettings, ConfigurationError
|
||||
|
||||
|
||||
def test_settings_require_whisper_model() -> None:
|
||||
settings = AppSettings(data_root=Path("runs"), whisper_model=None)
|
||||
|
||||
with pytest.raises(ConfigurationError, match="MKA_WHISPER_MODEL"):
|
||||
settings.validate_for_processing()
|
||||
|
||||
|
||||
def test_settings_reject_missing_whisper_model(tmp_path: Path) -> None:
|
||||
settings = AppSettings(
|
||||
data_root=tmp_path / "runs",
|
||||
whisper_model=tmp_path / "missing.bin",
|
||||
)
|
||||
|
||||
with pytest.raises(ConfigurationError, match="does not exist"):
|
||||
settings.validate_for_processing()
|
||||
|
||||
|
||||
def test_settings_accept_valid_native_configuration(tmp_path: Path) -> None:
|
||||
model = tmp_path / "model.bin"
|
||||
model.write_bytes(b"model")
|
||||
settings = AppSettings(data_root=tmp_path / "runs", whisper_model=model)
|
||||
|
||||
settings.validate_for_processing()
|
||||
|
||||
assert settings.diarization_runtime == "native"
|
||||
assert settings.diarization_container_args == ()
|
||||
|
||||
|
||||
def test_environment_defaults_to_no_diarization_container_args(monkeypatch) -> None:
|
||||
monkeypatch.delenv("MKA_DIARIZATION_CONTAINER_ARGS", raising=False)
|
||||
|
||||
settings = AppSettings.from_environment()
|
||||
|
||||
assert settings.diarization_container_args == ()
|
||||
assert settings.glossary_database == Path("data/database/glossary.sqlite3")
|
||||
|
||||
|
||||
def test_environment_configures_glossary_database(monkeypatch, tmp_path: Path) -> None:
|
||||
database = tmp_path / "terms.sqlite3"
|
||||
monkeypatch.setenv("MKA_GLOSSARY_DATABASE", str(database))
|
||||
|
||||
assert AppSettings.from_environment().glossary_database == database
|
||||
|
||||
|
||||
def test_environment_parses_multiple_ordered_container_args(monkeypatch) -> None:
|
||||
monkeypatch.setenv(
|
||||
"MKA_DIARIZATION_CONTAINER_ARGS",
|
||||
'["--device=/dev/kfd", "--device=/dev/dri", "--group-add", "video"]',
|
||||
)
|
||||
|
||||
settings = AppSettings.from_environment()
|
||||
|
||||
assert settings.diarization_container_args == (
|
||||
"--device=/dev/kfd",
|
||||
"--device=/dev/dri",
|
||||
"--group-add",
|
||||
"video",
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"value",
|
||||
["not-json", '"--flag"', '["valid", ""]', '["valid", 1]'],
|
||||
)
|
||||
def test_environment_rejects_invalid_container_args(monkeypatch, value: str) -> None:
|
||||
monkeypatch.setenv("MKA_DIARIZATION_CONTAINER_ARGS", value)
|
||||
|
||||
with pytest.raises(ConfigurationError, match="JSON array"):
|
||||
AppSettings.from_environment()
|
||||
@@ -0,0 +1,78 @@
|
||||
import sqlite3
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from mka.application.glossary import GlossaryConflictError, GlossaryRepository
|
||||
|
||||
|
||||
def repository(tmp_path: Path) -> GlossaryRepository:
|
||||
result = GlossaryRepository(tmp_path / "database" / "glossary.sqlite3")
|
||||
result.initialize()
|
||||
return result
|
||||
|
||||
|
||||
def test_initialization_creates_empty_versioned_database(tmp_path: Path) -> None:
|
||||
glossary = repository(tmp_path)
|
||||
|
||||
assert glossary.database_path.is_file()
|
||||
assert glossary.list() == []
|
||||
with sqlite3.connect(glossary.database_path) as connection:
|
||||
assert connection.execute("PRAGMA user_version").fetchone()[0] == 1
|
||||
tables = {
|
||||
row[0]
|
||||
for row in connection.execute("SELECT name FROM sqlite_master WHERE type = 'table'")
|
||||
}
|
||||
assert {"glossary_entries", "glossary_aliases"} <= tables
|
||||
|
||||
|
||||
def test_create_read_update_and_deactivate_with_multiple_aliases(tmp_path: Path) -> None:
|
||||
glossary = repository(tmp_path)
|
||||
created = glossary.create(
|
||||
"Secugrid HS",
|
||||
"product",
|
||||
aliases=("Sikirgut", "Secugrid H S"),
|
||||
description="Canonical core product name",
|
||||
)
|
||||
|
||||
assert glossary.get(created.id).aliases == ("Secugrid H S", "Sikirgut")
|
||||
assert glossary.list("sikir")[0].canonical_term == "Secugrid HS"
|
||||
|
||||
updated = glossary.update(
|
||||
created.id,
|
||||
"Secugrid HS",
|
||||
"technical_term",
|
||||
aliases=("Sekugrid HS",),
|
||||
description="Updated",
|
||||
is_active=False,
|
||||
)
|
||||
|
||||
assert updated.category == "technical_term"
|
||||
assert updated.aliases == ("Sekugrid HS",)
|
||||
assert not updated.is_active
|
||||
assert glossary.list(active_only=True) == []
|
||||
|
||||
|
||||
def test_terms_and_aliases_are_unique_case_insensitively_across_entries(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
glossary = repository(tmp_path)
|
||||
glossary.create("PBAT", "acronym", aliases=("Polybutylene adipate terephthalate",))
|
||||
|
||||
with pytest.raises(GlossaryConflictError):
|
||||
glossary.create("pbat", "material")
|
||||
with pytest.raises(GlossaryConflictError):
|
||||
glossary.create("Other", "other", aliases=("PBAT",))
|
||||
with pytest.raises(GlossaryConflictError):
|
||||
glossary.create("Polybutylene adipate terephthalate", "material")
|
||||
|
||||
|
||||
def test_delete_removes_entry_and_aliases(tmp_path: Path) -> None:
|
||||
glossary = repository(tmp_path)
|
||||
entry = glossary.create("Luminy", "product", aliases=("Lumini",))
|
||||
|
||||
glossary.delete(entry.id)
|
||||
|
||||
assert glossary.list() == []
|
||||
with sqlite3.connect(glossary.database_path) as connection:
|
||||
assert connection.execute("SELECT COUNT(*) FROM glossary_aliases").fetchone()[0] == 0
|
||||
@@ -0,0 +1,147 @@
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from mka.application.glossary import GlossaryRepository
|
||||
from mka.application.glossary_yaml import (
|
||||
GlossaryYamlError,
|
||||
export_glossary_yaml,
|
||||
import_glossary_yaml,
|
||||
)
|
||||
|
||||
|
||||
def repository(tmp_path: Path) -> GlossaryRepository:
|
||||
result = GlossaryRepository(tmp_path / "glossary.sqlite3")
|
||||
result.initialize()
|
||||
return result
|
||||
|
||||
|
||||
def test_round_trip_preserves_ids_aliases_and_active_state(tmp_path: Path) -> None:
|
||||
source = repository(tmp_path / "source")
|
||||
active = source.create(
|
||||
"Größenmaß",
|
||||
"technical_term",
|
||||
aliases=("Groessenmass", "Größen-Maß"),
|
||||
description="Canonical UTF-8 spelling",
|
||||
)
|
||||
inactive = source.create("Legacy Product", "product", is_active=False)
|
||||
|
||||
encoded = export_glossary_yaml(source.list()).encode("utf-8")
|
||||
target = repository(tmp_path / "target")
|
||||
target.create("Will be replaced", "other")
|
||||
target.replace_all(import_glossary_yaml(encoded))
|
||||
|
||||
entries = target.list()
|
||||
assert [entry.id for entry in entries] == [active.id, inactive.id]
|
||||
assert entries[0].canonical_term == "Größenmaß"
|
||||
assert entries[0].aliases == ("Groessenmass", "Größen-Maß")
|
||||
assert entries[0].description == "Canonical UTF-8 spelling"
|
||||
assert entries[0].is_active is True
|
||||
assert entries[1].is_active is False
|
||||
|
||||
|
||||
def test_malformed_yaml_does_not_replace_existing_glossary(tmp_path: Path) -> None:
|
||||
glossary = repository(tmp_path)
|
||||
existing = glossary.create("PBAT", "acronym")
|
||||
|
||||
with pytest.raises(GlossaryYamlError, match="Malformed glossary YAML"):
|
||||
import_glossary_yaml(b"version: 1\nglossary: [")
|
||||
|
||||
assert glossary.list() == [existing]
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"entries, message",
|
||||
[
|
||||
(
|
||||
[
|
||||
{
|
||||
"id": 1,
|
||||
"canonical_term": "PBAT",
|
||||
"aliases": [],
|
||||
"category": "acronym",
|
||||
"description": None,
|
||||
"active": True,
|
||||
},
|
||||
{
|
||||
"id": 2,
|
||||
"canonical_term": "pbat",
|
||||
"aliases": [],
|
||||
"category": "material",
|
||||
"description": None,
|
||||
"active": True,
|
||||
},
|
||||
],
|
||||
"duplicate canonical term or alias",
|
||||
),
|
||||
(
|
||||
[
|
||||
{
|
||||
"id": 1,
|
||||
"canonical_term": "Secugrid",
|
||||
"aliases": ["Sikirgut", "SIKIRGUT"],
|
||||
"category": "product",
|
||||
"description": None,
|
||||
"active": True,
|
||||
},
|
||||
],
|
||||
"aliases must be unique",
|
||||
),
|
||||
(
|
||||
[
|
||||
{
|
||||
"id": 1,
|
||||
"canonical_term": "Secugrid",
|
||||
"aliases": ["PBAT"],
|
||||
"category": "product",
|
||||
"description": None,
|
||||
"active": True,
|
||||
},
|
||||
{
|
||||
"id": 2,
|
||||
"canonical_term": "PBAT",
|
||||
"aliases": [],
|
||||
"category": "acronym",
|
||||
"description": None,
|
||||
"active": False,
|
||||
},
|
||||
],
|
||||
"duplicate canonical term or alias",
|
||||
),
|
||||
],
|
||||
)
|
||||
def test_duplicate_terms_and_aliases_are_rejected(
|
||||
entries: list[dict[str, object]], message: str
|
||||
) -> None:
|
||||
import yaml
|
||||
|
||||
content = yaml.safe_dump({"version": 1, "glossary": entries})
|
||||
|
||||
with pytest.raises(GlossaryYamlError, match=message):
|
||||
import_glossary_yaml(content)
|
||||
|
||||
|
||||
def test_validation_failure_leaves_database_unchanged(tmp_path: Path) -> None:
|
||||
glossary = repository(tmp_path)
|
||||
existing = glossary.create("Existing", "other", aliases=("Existing alias",))
|
||||
invalid = b"""version: 1
|
||||
glossary:
|
||||
- id: 10
|
||||
canonical_term: Duplicate
|
||||
aliases: []
|
||||
category: other
|
||||
description: null
|
||||
active: true
|
||||
- id: 11
|
||||
canonical_term: duplicate
|
||||
aliases: []
|
||||
category: other
|
||||
description: null
|
||||
active: false
|
||||
"""
|
||||
|
||||
with pytest.raises(GlossaryYamlError):
|
||||
imported = import_glossary_yaml(invalid)
|
||||
glossary.replace_all(imported)
|
||||
|
||||
assert glossary.list() == [existing]
|
||||
@@ -0,0 +1,88 @@
|
||||
"""Exercise the real executor and Streamlit reruns without model calls."""
|
||||
|
||||
from pathlib import Path
|
||||
from threading import current_thread
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
import streamlit
|
||||
from streamlit.testing.v1 import AppTest
|
||||
|
||||
from mka.application.meeting_service import AppProgressEvent, SpeakerMappingReview, SpeakerReview
|
||||
from mka.application.progress_timing import ProcessingTimer
|
||||
from mka.ui import streamlit_app as ui
|
||||
|
||||
|
||||
class GuardedStreamlit:
|
||||
"""Fail if a worker touches the UI session proxy, including callback reads."""
|
||||
|
||||
def __getattr__(self, name):
|
||||
if name == "session_state" and current_thread().name.startswith("ThreadPoolExecutor"):
|
||||
raise AssertionError("Worker accessed Streamlit session state")
|
||||
return getattr(streamlit, name)
|
||||
|
||||
|
||||
@pytest.mark.parametrize("fails", [False, True])
|
||||
@pytest.mark.parametrize("mapped", [False, True])
|
||||
def test_regeneration_captures_ui_values_before_real_worker_and_retains_timing(
|
||||
monkeypatch, fails, mapped
|
||||
):
|
||||
outcome = SimpleNamespace(
|
||||
succeeded=True,
|
||||
run_dir=Path("/tmp/alpha-mapping-test"),
|
||||
speaker_attribution_available=True,
|
||||
original_protocol="Existing protocol",
|
||||
)
|
||||
review = SpeakerMappingReview((SpeakerReview("SPEAKER_00", ()),), (("a", "A"),), {})
|
||||
calls = []
|
||||
|
||||
class Service:
|
||||
def __init__(self, *args):
|
||||
self.timer = ProcessingTimer()
|
||||
|
||||
def load_speaker_mapping_review(self, run_dir):
|
||||
return review
|
||||
|
||||
def regenerate_protocol(self, run_dir, selections, *, progress_sink, performance_profile):
|
||||
assert current_thread().name.startswith("ThreadPoolExecutor")
|
||||
calls.append((run_dir, selections, performance_profile))
|
||||
self.timer.start_run()
|
||||
self.timer.start_stage("protocol_generation")
|
||||
progress_sink(AppProgressEvent("protocol_generation", "started", 0.0))
|
||||
if fails:
|
||||
raise RuntimeError("model unavailable")
|
||||
self.timer.finish_stage("protocol_generation")
|
||||
self.timer.finish_run()
|
||||
progress_sink(AppProgressEvent("protocol_generation", "completed", 0.0))
|
||||
return outcome
|
||||
|
||||
monkeypatch.setattr(ui, "st", GuardedStreamlit())
|
||||
monkeypatch.setattr(ui, "MeetingProcessingService", Service)
|
||||
monkeypatch.setattr(ui.AppSettings, "from_environment", lambda: None)
|
||||
monkeypatch.setattr(ui, "MeetingLabGateway", lambda: None)
|
||||
|
||||
def result_app(outcome):
|
||||
import streamlit as st
|
||||
|
||||
from mka.ui.streamlit_app import _render_result
|
||||
|
||||
st.session_state.outcome = outcome
|
||||
st.session_state.performance_profile = "fast"
|
||||
_render_result()
|
||||
|
||||
app = AppTest.from_function(result_app, args=(outcome,)).run()
|
||||
if mapped:
|
||||
app.selectbox[0].select("a").run()
|
||||
button = next(button for button in app.button if "Regenerate" in button.label)
|
||||
assert not button.disabled
|
||||
button.click().run()
|
||||
assert not app.exception
|
||||
assert calls == [(outcome.run_dir, {"SPEAKER_00": "a"} if mapped else {}, "fast")]
|
||||
states, message, timing = app.session_state["processing_progress"]
|
||||
assert states["protocol_generation"] == ("failed" if fails else "completed")
|
||||
assert not timing.running
|
||||
assert timing.total_duration >= timing.stage_durations["protocol_generation"] >= 0
|
||||
app.run()
|
||||
assert not app.exception
|
||||
assert app.session_state["processing_progress"][2] == timing
|
||||
assert any("total runtime" in item.value for item in app.info)
|
||||
@@ -0,0 +1,826 @@
|
||||
import json
|
||||
from dataclasses import dataclass, replace
|
||||
from datetime import date
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
import yaml
|
||||
|
||||
from mka.application.config import AppSettings
|
||||
from mka.application.meeting_service import (
|
||||
MeetingDetails,
|
||||
MeetingProcessingService,
|
||||
ParticipantInput,
|
||||
ProcessingOptions,
|
||||
)
|
||||
from mka.application.progress_timing import ProcessingTimer
|
||||
|
||||
|
||||
@dataclass
|
||||
class FakeContext:
|
||||
data: dict[str, Any]
|
||||
|
||||
@property
|
||||
def meeting_id(self) -> str:
|
||||
return self.data["meeting"]["meeting_id"]
|
||||
|
||||
|
||||
class FakeMeetingLab:
|
||||
def __init__(self, run_dir: Path) -> None:
|
||||
self.run_dir = run_dir
|
||||
self.context_data: dict[str, Any] | None = None
|
||||
self.config_values: dict[str, Any] | None = None
|
||||
self.fail = False
|
||||
self.regeneration: dict[str, Any] | None = None
|
||||
|
||||
def create_context(self, data: dict[str, Any]) -> FakeContext:
|
||||
self.context_data = data
|
||||
if not data["meeting"]["title"]:
|
||||
raise ValueError("meeting.title must be present and non-empty")
|
||||
participant_ids = [item["participant_id"] for item in data["participants"]]
|
||||
if len(participant_ids) != len(set(participant_ids)):
|
||||
raise ValueError("Duplicate participant_id")
|
||||
return FakeContext(data)
|
||||
|
||||
def create_config(self, values: dict[str, Any]) -> dict[str, Any]:
|
||||
self.config_values = values
|
||||
return values
|
||||
|
||||
def run(self, config: Any, meeting_context: Any, progress_sink: Any) -> Any:
|
||||
self.run_dir.mkdir(parents=True, exist_ok=True)
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="preparing",
|
||||
status="started",
|
||||
elapsed_seconds=0.1,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="preparing",
|
||||
status="completed",
|
||||
elapsed_seconds=0.2,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="transcription",
|
||||
status="started",
|
||||
elapsed_seconds=0.2,
|
||||
progress=0.25,
|
||||
message="transcribing",
|
||||
)
|
||||
)
|
||||
if self.fail:
|
||||
(self.run_dir / "run_metadata.json").write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"failure": {
|
||||
"stage": "whisper",
|
||||
"type": "TranscriptionError",
|
||||
"message": "model failed",
|
||||
}
|
||||
}
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="failed",
|
||||
status="failed",
|
||||
elapsed_seconds=0.3,
|
||||
progress=None,
|
||||
message="transcription: model failed",
|
||||
)
|
||||
)
|
||||
return SimpleNamespace(exit_code=2, run_dir=self.run_dir, protocol_path=None)
|
||||
if config["stop_after_diarization"]:
|
||||
write_speaker_review_artifacts(self.run_dir)
|
||||
(self.run_dir / "run_metadata.json").write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"status": "awaiting_speaker_review",
|
||||
"diarization": {"speaker_count": 2},
|
||||
}
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="diarization", status="completed", elapsed_seconds=0.3,
|
||||
progress=None, message=None,
|
||||
)
|
||||
)
|
||||
return SimpleNamespace(exit_code=0, run_dir=self.run_dir, protocol_path=None)
|
||||
protocol = self.run_dir / "protocol.md"
|
||||
protocol.write_text("# Generated protocol\n", encoding="utf-8")
|
||||
return SimpleNamespace(exit_code=0, run_dir=self.run_dir, protocol_path=protocol)
|
||||
|
||||
def regenerate_protocol(
|
||||
self,
|
||||
run_dir: Path,
|
||||
meeting_context: Any,
|
||||
progress_sink: Any,
|
||||
**options: Any,
|
||||
) -> Any:
|
||||
self.regeneration = {
|
||||
"run_dir": run_dir,
|
||||
"meeting_context": meeting_context,
|
||||
"options": options,
|
||||
}
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="protocol_generation",
|
||||
status="started",
|
||||
elapsed_seconds=0.0,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
protocol = Path(run_dir) / "protocol.md"
|
||||
protocol.write_text("# Regenerated protocol\n", encoding="utf-8")
|
||||
protocol_dir = Path(run_dir) / "protocol"
|
||||
protocol_dir.mkdir(exist_ok=True)
|
||||
(protocol_dir / "runtime_metadata.json").write_text(
|
||||
json.dumps({"speaker_attribution_available": True}), encoding="utf-8"
|
||||
)
|
||||
return SimpleNamespace(exit_code=0, run_dir=run_dir, protocol_path=protocol)
|
||||
|
||||
|
||||
def make_service(tmp_path: Path) -> tuple[MeetingProcessingService, FakeMeetingLab]:
|
||||
model = tmp_path / "model.bin"
|
||||
model.write_bytes(b"model")
|
||||
gateway = FakeMeetingLab(tmp_path / "backend-run")
|
||||
settings = AppSettings(
|
||||
data_root=tmp_path / "meetings",
|
||||
whisper_model=model,
|
||||
glossary_database=tmp_path / "glossary.sqlite3",
|
||||
whisper_executable="/opt/whisper-cli",
|
||||
protocol_model="test:model",
|
||||
diarization_mode="gpu",
|
||||
)
|
||||
return MeetingProcessingService(settings, gateway), gateway
|
||||
|
||||
|
||||
def meeting() -> MeetingDetails:
|
||||
return MeetingDetails(
|
||||
title="Architecture Review",
|
||||
language="de",
|
||||
meeting_date=date(2026, 8, 23),
|
||||
description="Review the MVP.",
|
||||
)
|
||||
|
||||
|
||||
def participants() -> list[ParticipantInput]:
|
||||
return [
|
||||
ParticipantInput(
|
||||
participant_id="martin",
|
||||
display_name="Martin",
|
||||
role="Project lead",
|
||||
organization="Engineering",
|
||||
),
|
||||
ParticipantInput(
|
||||
participant_id="alex",
|
||||
display_name="Alex",
|
||||
organization="Engineering",
|
||||
),
|
||||
]
|
||||
|
||||
|
||||
def write_speaker_review_artifacts(run_dir: Path) -> Path:
|
||||
diarization_dir = run_dir / "diarization"
|
||||
context_dir = run_dir / "context"
|
||||
diarization_dir.mkdir(parents=True)
|
||||
context_dir.mkdir()
|
||||
transcript_path = diarization_dir / "transcript_diarized.json"
|
||||
transcript_path.write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"speaker_labels_anonymous": True,
|
||||
"segments": [
|
||||
{
|
||||
"speaker_id": "SPEAKER_01",
|
||||
"text": "I will prepare all raw materials before Wednesday.",
|
||||
},
|
||||
{"speaker_id": "SPEAKER_00", "text": "Yes."},
|
||||
{
|
||||
"speaker_id": "SPEAKER_00",
|
||||
"text": "We will run the production trial on Wednesday.",
|
||||
},
|
||||
{
|
||||
"speaker_id": "SPEAKER_00",
|
||||
"text": "The trial requires the complete production team.",
|
||||
},
|
||||
{"speaker_id": None, "text": "Unassigned text."},
|
||||
],
|
||||
}
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
(context_dir / "meeting_context.yaml").write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"schema_version": "1",
|
||||
"meeting": {
|
||||
"meeting_id": "speaker-review",
|
||||
"title": "Speaker review",
|
||||
"language": "en",
|
||||
},
|
||||
"participants": [
|
||||
{"participant_id": "martin", "display_name": "Martin"},
|
||||
{"participant_id": "anna", "display_name": "Anna"},
|
||||
],
|
||||
"speaker_mappings": {},
|
||||
"mentioned_people": [],
|
||||
"organization": {"departments": []},
|
||||
"known_entities": {},
|
||||
}
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
return transcript_path
|
||||
|
||||
|
||||
def test_build_context_uses_actual_v1_shape(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
|
||||
context = service.build_context(meeting(), participants())
|
||||
|
||||
assert context.meeting_id == "architecture-review"
|
||||
assert gateway.context_data is not None
|
||||
assert gateway.context_data["meeting"]["date"] == "2026-08-23"
|
||||
assert gateway.context_data["meeting"]["notes"] == "Review the MVP."
|
||||
assert gateway.context_data["participants"][0] == {
|
||||
"participant_id": "martin",
|
||||
"display_name": "Martin",
|
||||
"aliases": [],
|
||||
"role": "Project lead",
|
||||
"department": "engineering",
|
||||
"attendance_status": "present",
|
||||
"notes": None,
|
||||
}
|
||||
assert gateway.context_data["organization"]["departments"] == [
|
||||
{"id": "engineering", "name": "Engineering", "aliases": []}
|
||||
]
|
||||
|
||||
|
||||
def test_build_context_preserves_explicit_speaker_mapping(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
|
||||
service.build_context(meeting(), participants(), {"SPEAKER_00": "martin"})
|
||||
|
||||
assert gateway.context_data is not None
|
||||
assert gateway.context_data["speaker_mappings"] == {"SPEAKER_00": "martin"}
|
||||
|
||||
|
||||
def test_build_context_includes_only_active_authoritative_glossary_terms(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
service.glossary.create("Secugrid HS", "product", aliases=("Sikirgut", "Secugrid H S"))
|
||||
inactive = service.glossary.create("Old Name", "other")
|
||||
service.glossary.set_active(inactive.id, False)
|
||||
|
||||
service.build_context(meeting(), participants())
|
||||
|
||||
assert gateway.context_data is not None
|
||||
assert gateway.context_data["known_entities"] == {
|
||||
"Authoritative terminology": ["Secugrid HS (aliases: Secugrid H S, Sikirgut)"]
|
||||
}
|
||||
rules = gateway.context_data["context_rules"]
|
||||
assert "do not invent matches" in rules["glossary_canonical_spelling"]
|
||||
assert "inside compounds" in rules["glossary_core_terms"]
|
||||
assert "Old Name" not in str(gateway.context_data)
|
||||
|
||||
|
||||
def test_glossary_is_rendered_into_meeting_lab_protocol_context(tmp_path: Path) -> None:
|
||||
meeting_context = pytest.importorskip("src.meeting_lab.models.meeting_context")
|
||||
service, gateway = make_service(tmp_path)
|
||||
service.glossary.create("PBAT", "acronym", aliases=("P B A T",))
|
||||
context = service.build_context(meeting(), participants())
|
||||
|
||||
real_context = meeting_context.create_meeting_context(context.data)
|
||||
prompt_context = meeting_context.render_meeting_context_for_prompt(real_context)
|
||||
|
||||
assert "Authoritative terminology: PBAT (aliases: P B A T)" in prompt_context
|
||||
assert "Use canonical glossary spellings" in prompt_context
|
||||
assert "use canonical core terms inside compounds" in prompt_context
|
||||
|
||||
|
||||
def test_build_context_translates_mentioned_only_person(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
people = participants() + [
|
||||
ParticipantInput(
|
||||
participant_id="sam",
|
||||
display_name="Sam",
|
||||
attendance_status="mentioned_only",
|
||||
)
|
||||
]
|
||||
|
||||
service.build_context(meeting(), people)
|
||||
|
||||
assert gateway.context_data is not None
|
||||
assert [item["participant_id"] for item in gateway.context_data["participants"]] == [
|
||||
"martin",
|
||||
"alex",
|
||||
]
|
||||
assert gateway.context_data["mentioned_people"] == [
|
||||
{
|
||||
"person_id": "sam",
|
||||
"display_name": "Sam",
|
||||
"aliases": [],
|
||||
"role": None,
|
||||
"department": None,
|
||||
"attendance_status": "mentioned_only",
|
||||
"notes": None,
|
||||
}
|
||||
]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("language", ["de", "en"])
|
||||
def test_process_translates_configuration_and_disables_diarization(
|
||||
tmp_path: Path,
|
||||
language: str,
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(
|
||||
audio,
|
||||
replace(meeting(), language=language),
|
||||
participants(),
|
||||
ProcessingOptions(diarization_enabled=False),
|
||||
)
|
||||
|
||||
assert outcome.succeeded
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["language"] == language
|
||||
assert service.build_context(replace(meeting(), language=language), participants()).data[
|
||||
"meeting"
|
||||
]["language"] == language
|
||||
assert gateway.config_values["diarization"] == "off"
|
||||
assert gateway.config_values["model"] == "test:model"
|
||||
assert gateway.config_values["protocol_num_ctx"] == 32_768
|
||||
assert gateway.config_values["protocol_safe_input_token_budget"] == 29_000
|
||||
assert gateway.config_values["whisper_executable"] == "/opt/whisper-cli"
|
||||
assert gateway.config_values["ffmpeg_executable"] == "ffmpeg"
|
||||
assert gateway.config_values["audio_normalization"] is True
|
||||
assert gateway.config_values["diarization_runtime"] == "native"
|
||||
assert gateway.config_values["diarization_container_args"] == ()
|
||||
assert gateway.config_values["output_root"] == (
|
||||
tmp_path / "meetings" / "architecture-review" / "runs"
|
||||
)
|
||||
|
||||
|
||||
def test_process_propagates_disabled_audio_normalization(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.m4a"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(
|
||||
audio,
|
||||
meeting(),
|
||||
participants(),
|
||||
ProcessingOptions(audio_normalization=False),
|
||||
)
|
||||
|
||||
assert outcome.succeeded
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["audio_normalization"] is False
|
||||
|
||||
|
||||
def test_processing_options_default_to_audio_normalization_on() -> None:
|
||||
assert ProcessingOptions().audio_normalization is True
|
||||
|
||||
|
||||
def test_process_propagates_enabled_diarization_and_progress(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.flac"
|
||||
audio.write_bytes(b"audio")
|
||||
events = []
|
||||
|
||||
outcome = service.process(
|
||||
audio,
|
||||
meeting(),
|
||||
participants(),
|
||||
ProcessingOptions(diarization_enabled=True),
|
||||
progress_sink=events.append,
|
||||
)
|
||||
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["diarization"] == "gpu"
|
||||
assert gateway.config_values["stop_after_diarization"] is True
|
||||
assert outcome.awaiting_speaker_review is True
|
||||
assert outcome.detected_speaker_count == 2
|
||||
assert [(event.stage, event.status) for event in events[:2]] == [
|
||||
("preparing", "started"),
|
||||
("preparing", "completed"),
|
||||
]
|
||||
assert events[2].progress == 0.25
|
||||
|
||||
|
||||
def test_diarization_checkpoint_defers_protocol_but_disabled_run_does_not(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
checkpoint = service.process(
|
||||
audio, meeting(), participants(), ProcessingOptions(diarization_enabled=True)
|
||||
)
|
||||
assert checkpoint.succeeded
|
||||
assert checkpoint.awaiting_speaker_review
|
||||
assert checkpoint.original_protocol is None
|
||||
assert checkpoint.protocol_path is None
|
||||
|
||||
normal_root = tmp_path / "without-diarization"
|
||||
normal_root.mkdir()
|
||||
normal_service, normal_gateway = make_service(normal_root)
|
||||
normal_audio = normal_root / "meeting.wav"
|
||||
normal_audio.write_bytes(b"audio")
|
||||
normal = normal_service.process(
|
||||
normal_audio, meeting(), participants(), ProcessingOptions(diarization_enabled=False)
|
||||
)
|
||||
assert normal.succeeded
|
||||
assert not normal.awaiting_speaker_review
|
||||
assert normal.protocol_path is not None
|
||||
assert normal_gateway.config_values["stop_after_diarization"] is False
|
||||
|
||||
|
||||
def test_process_propagates_ordered_diarization_container_args(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
service.settings = replace(
|
||||
service.settings,
|
||||
diarization_runtime="container",
|
||||
diarization_container_image="runtime/image:tag",
|
||||
diarization_container_args=(
|
||||
"--network=host",
|
||||
"--label",
|
||||
"meeting-test",
|
||||
),
|
||||
)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(
|
||||
audio,
|
||||
meeting(),
|
||||
participants(),
|
||||
ProcessingOptions(diarization_enabled=True),
|
||||
)
|
||||
|
||||
assert outcome.succeeded
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["diarization_container_args"] == (
|
||||
"--network=host",
|
||||
"--label",
|
||||
"meeting-test",
|
||||
)
|
||||
|
||||
|
||||
def test_result_and_user_edit_are_preserved_separately(tmp_path: Path) -> None:
|
||||
service, _ = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
edited_path = service.save_edited_protocol(outcome.run_dir, "# Reviewed\n")
|
||||
|
||||
assert outcome.original_protocol == "# Generated protocol\n"
|
||||
assert outcome.protocol_path.read_text(encoding="utf-8") == "# Generated protocol\n"
|
||||
assert edited_path.read_text(encoding="utf-8") == "# Reviewed\n"
|
||||
|
||||
|
||||
def test_failure_reports_stage_and_preserves_run_dir(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
gateway.fail = True
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
assert not outcome.succeeded
|
||||
assert outcome.failed_stage == "transcription"
|
||||
assert outcome.error_message == "TranscriptionError: model failed"
|
||||
assert outcome.run_dir == gateway.run_dir
|
||||
assert (gateway.run_dir / "run_metadata.json").is_file()
|
||||
|
||||
|
||||
def test_speaker_review_lists_detected_labels_participants_and_excerpts(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
source = write_speaker_review_artifacts(gateway.run_dir)
|
||||
|
||||
review = service.load_speaker_mapping_review(gateway.run_dir, excerpts_per_speaker=2)
|
||||
|
||||
assert review is not None
|
||||
assert [speaker.speaker_label for speaker in review.speakers] == [
|
||||
"SPEAKER_00",
|
||||
"SPEAKER_01",
|
||||
]
|
||||
assert review.speakers[0].excerpts == (
|
||||
"We will run the production trial on Wednesday.",
|
||||
"The trial requires the complete production team.",
|
||||
)
|
||||
assert review.participants == (("martin", "Martin"), ("anna", "Anna"))
|
||||
assert review.current_mappings == {}
|
||||
assert "SPEAKER_00" in source.read_text(encoding="utf-8")
|
||||
|
||||
|
||||
def test_protocol_only_regeneration_persists_mappings_without_rewriting_transcript(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
service.glossary.create("ENLYZE", "organization", aliases=("Enlyse",))
|
||||
source = write_speaker_review_artifacts(gateway.run_dir)
|
||||
original_source = source.read_bytes()
|
||||
events = []
|
||||
|
||||
outcome = service.regenerate_protocol(
|
||||
gateway.run_dir,
|
||||
{"SPEAKER_00": "martin"},
|
||||
progress_sink=events.append,
|
||||
)
|
||||
|
||||
assert outcome.succeeded
|
||||
assert outcome.original_protocol == "# Regenerated protocol\n"
|
||||
assert outcome.speaker_attribution_available is True
|
||||
assert gateway.context_data is not None
|
||||
assert gateway.context_data["speaker_mappings"] == {"SPEAKER_00": "martin"}
|
||||
assert gateway.context_data["known_entities"] == {
|
||||
"Authoritative terminology": ["ENLYZE (aliases: Enlyse)"]
|
||||
}
|
||||
assert gateway.regeneration is not None
|
||||
assert gateway.regeneration["meeting_context"].data["speaker_mappings"] == {
|
||||
"SPEAKER_00": "martin"
|
||||
}
|
||||
assert gateway.regeneration["options"]["protocol_num_ctx"] == 32_768
|
||||
assert gateway.regeneration["options"]["glossary_aliases"] == {"Enlyse": "ENLYZE"}
|
||||
assert events[0].stage == "protocol_generation"
|
||||
assert source.read_bytes() == original_source
|
||||
|
||||
|
||||
def test_protocol_regeneration_allows_unmapped_and_rejects_duplicate_participant(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
|
||||
outcome = service.regenerate_protocol(gateway.run_dir, {})
|
||||
|
||||
assert outcome.succeeded
|
||||
assert gateway.context_data is not None
|
||||
assert gateway.context_data["speaker_mappings"] == {}
|
||||
|
||||
with pytest.raises(ValueError, match="only one speaker"):
|
||||
service.regenerate_protocol(
|
||||
gateway.run_dir,
|
||||
{"SPEAKER_00": "martin", "SPEAKER_01": "martin"},
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"mappings",
|
||||
[
|
||||
{"SPEAKER_00": "martin"},
|
||||
{"SPEAKER_00": "martin", "SPEAKER_01": "anna"},
|
||||
],
|
||||
)
|
||||
def test_protocol_generation_uses_partial_and_complete_confirmed_mappings(
|
||||
tmp_path: Path, mappings: dict[str, str]
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
|
||||
outcome = service.regenerate_protocol(gateway.run_dir, mappings)
|
||||
|
||||
assert outcome.succeeded
|
||||
assert gateway.regeneration is not None
|
||||
assert gateway.regeneration["meeting_context"].data["speaker_mappings"] == mappings
|
||||
|
||||
|
||||
def test_failed_protocol_generation_keeps_checkpoint_mapping_for_retry(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
source = write_speaker_review_artifacts(gateway.run_dir)
|
||||
original_source = source.read_bytes()
|
||||
original_regeneration = gateway.regenerate_protocol
|
||||
|
||||
def fail_regeneration(*args, **kwargs):
|
||||
raise RuntimeError("model unavailable")
|
||||
|
||||
gateway.regenerate_protocol = fail_regeneration
|
||||
with pytest.raises(RuntimeError, match="model unavailable"):
|
||||
service.regenerate_protocol(gateway.run_dir, {"SPEAKER_00": "martin"})
|
||||
|
||||
persisted = yaml.safe_load(
|
||||
(gateway.run_dir / "context" / "meeting_context.yaml").read_text(encoding="utf-8")
|
||||
)
|
||||
assert persisted["speaker_mappings"] == {"SPEAKER_00": "martin"}
|
||||
gateway.regenerate_protocol = original_regeneration
|
||||
retry = service.regenerate_protocol(gateway.run_dir, {"SPEAKER_00": "martin"})
|
||||
assert retry.succeeded
|
||||
assert source.read_bytes() == original_source
|
||||
|
||||
|
||||
def test_fallback_attribution_loss_is_read_from_runtime_metadata(tmp_path: Path) -> None:
|
||||
run_dir = tmp_path / "run"
|
||||
protocol_dir = run_dir / "protocol"
|
||||
protocol_dir.mkdir(parents=True)
|
||||
(protocol_dir / "runtime_metadata.json").write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"speaker_attribution_available": False,
|
||||
"speaker_attribution_loss_reason": "plain_transcript_fallback",
|
||||
}
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
assert MeetingProcessingService._speaker_attribution_available(run_dir) is False
|
||||
|
||||
|
||||
def test_uploaded_source_is_preserved_in_meeting_directory(tmp_path: Path) -> None:
|
||||
service, _ = make_service(tmp_path)
|
||||
source = SimpleNamespace(getbuffer=lambda: b"source audio")
|
||||
|
||||
destination = service.preserve_upload("meeting-1", "../unsafe.wav", source)
|
||||
|
||||
assert destination.parent == tmp_path / "meetings" / "meeting-1" / "uploads"
|
||||
assert destination.name.endswith("_unsafe.wav")
|
||||
assert destination.read_bytes() == b"source audio"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("profile", "threads"),
|
||||
[("fast", 16), ("efficient", 10), ("powersave", 4)],
|
||||
)
|
||||
def test_process_resolves_profile_for_protocol_generation(
|
||||
tmp_path: Path, profile: str, threads: int
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
service.process(
|
||||
audio,
|
||||
meeting(),
|
||||
participants(),
|
||||
ProcessingOptions(performance_profile=profile),
|
||||
)
|
||||
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["protocol_num_thread"] == threads
|
||||
|
||||
|
||||
def test_process_defaults_to_backend_thread_selection(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["protocol_num_thread"] is None
|
||||
|
||||
|
||||
def test_protocol_regeneration_propagates_explicit_profile(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
|
||||
service.regenerate_protocol(
|
||||
gateway.run_dir,
|
||||
{},
|
||||
performance_profile="fast",
|
||||
)
|
||||
|
||||
assert gateway.regeneration is not None
|
||||
assert gateway.regeneration["options"]["protocol_num_thread"] == 16
|
||||
|
||||
|
||||
def test_process_passes_only_active_explicit_glossary_aliases(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
service.glossary.create("Luminy", "product", aliases=("Lumini",))
|
||||
inactive = service.glossary.create("PBAT", "acronym", aliases=("PBRT",))
|
||||
service.glossary.set_active(inactive.id, False)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["glossary_aliases"] == {"Lumini": "Luminy"}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("fails", [False, True])
|
||||
def test_process_persists_completed_and_failed_timings(tmp_path: Path, fails: bool) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
gateway.fail = fails
|
||||
values = iter([0.0, 1.0, 9.0, 10.0, 20.0])
|
||||
service.timer = ProcessingTimer(lambda: next(values))
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
metadata = json.loads((gateway.run_dir / "run_metadata.json").read_text(encoding="utf-8"))
|
||||
assert metadata["timing"] == {
|
||||
"stages_seconds": {"preparing": 8.0, "transcription": 10.0},
|
||||
"total_seconds": 20.0,
|
||||
}
|
||||
assert outcome.succeeded is not fails
|
||||
if fails:
|
||||
assert metadata["failure"]["type"] == "TranscriptionError"
|
||||
|
||||
|
||||
def test_regeneration_starts_fresh_timer_and_retains_final_duration(
|
||||
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
clock = SimpleNamespace(value=10.0)
|
||||
service.timer = ProcessingTimer(lambda: clock.value)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
clock.value = 100.0
|
||||
observed = []
|
||||
|
||||
def regenerate(run_dir: Path, context: Any, progress_sink: Any, **options: Any) -> Any:
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="protocol_generation",
|
||||
status="started",
|
||||
elapsed_seconds=0.0,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
clock.value = 105.0
|
||||
observed.append(service.timer.snapshot())
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="protocol_generation",
|
||||
status="completed",
|
||||
elapsed_seconds=5.0,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
clock.value = 109.0
|
||||
protocol = run_dir / "protocol.md"
|
||||
protocol.write_text("# Regenerated protocol\n", encoding="utf-8")
|
||||
return SimpleNamespace(run_dir=run_dir, protocol_path=protocol)
|
||||
|
||||
monkeypatch.setattr(gateway, "regenerate_protocol", regenerate)
|
||||
|
||||
service.regenerate_protocol(gateway.run_dir, {"SPEAKER_00": "martin"})
|
||||
final = service.timer.snapshot()
|
||||
|
||||
assert observed[0].active_stage == "protocol_generation"
|
||||
assert observed[0].stage_durations == {"protocol_generation": 5.0}
|
||||
assert observed[0].total_duration == 5.0
|
||||
assert observed[0].running is True
|
||||
assert final.stage_durations == {"protocol_generation": 5.0}
|
||||
assert final.total_duration == 9.0
|
||||
assert final.running is False
|
||||
|
||||
|
||||
def test_failed_regeneration_freezes_active_and_total_durations(
|
||||
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
clock = SimpleNamespace(value=20.0)
|
||||
service.timer = ProcessingTimer(lambda: clock.value)
|
||||
|
||||
def fail_regeneration(run_dir: Path, context: Any, progress_sink: Any, **options: Any) -> Any:
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="protocol_generation",
|
||||
status="started",
|
||||
elapsed_seconds=0.0,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
clock.value = 27.0
|
||||
raise RuntimeError("generation failed")
|
||||
|
||||
monkeypatch.setattr(gateway, "regenerate_protocol", fail_regeneration)
|
||||
|
||||
with pytest.raises(RuntimeError, match="generation failed"):
|
||||
service.regenerate_protocol(gateway.run_dir, {})
|
||||
clock.value = 99.0
|
||||
final = service.timer.snapshot()
|
||||
|
||||
assert final.stage_durations == {"protocol_generation": 7.0}
|
||||
assert final.total_duration == 7.0
|
||||
assert final.running is False
|
||||
@@ -0,0 +1,127 @@
|
||||
import pytest
|
||||
|
||||
from mka.application.meeting_service import ParticipantInput
|
||||
from mka.application.people_yaml import (
|
||||
PeopleYamlError,
|
||||
export_people_yaml,
|
||||
import_people_yaml,
|
||||
)
|
||||
|
||||
|
||||
def sample_people() -> list[ParticipantInput]:
|
||||
return [
|
||||
ParticipantInput(
|
||||
participant_id="participant-35c1528c",
|
||||
display_name="Martin Tazl",
|
||||
role="Beirat",
|
||||
organization="Verwaltungsbeirat",
|
||||
attendance_status="present",
|
||||
),
|
||||
ParticipantInput(
|
||||
participant_id="participant-56aae811",
|
||||
display_name="Norbert Hebbelmann",
|
||||
role="Gast",
|
||||
organization="Eigentümergemeinschaft",
|
||||
attendance_status="mentioned_only",
|
||||
),
|
||||
]
|
||||
|
||||
|
||||
def test_export_import_round_trip_preserves_all_people_fields() -> None:
|
||||
people = sample_people()
|
||||
|
||||
exported = export_people_yaml(people)
|
||||
|
||||
assert exported.startswith("version: 1\npeople:\n")
|
||||
assert import_people_yaml(exported) == people
|
||||
assert export_people_yaml(import_people_yaml(exported)) == exported
|
||||
|
||||
|
||||
def test_import_preserves_stable_ids_roles_organizations_and_attendance() -> None:
|
||||
imported = import_people_yaml(export_people_yaml(sample_people()))
|
||||
|
||||
assert [person.participant_id for person in imported] == [
|
||||
"participant-35c1528c",
|
||||
"participant-56aae811",
|
||||
]
|
||||
assert imported[0].role == "Beirat"
|
||||
assert imported[1].organization == "Eigentümergemeinschaft"
|
||||
assert [person.attendance_status for person in imported] == [
|
||||
"present",
|
||||
"mentioned_only",
|
||||
]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("content", ["people: [", b"\xff\xfe"])
|
||||
def test_malformed_yaml_is_rejected(content: str | bytes) -> None:
|
||||
with pytest.raises(PeopleYamlError, match="Malformed People YAML"):
|
||||
import_people_yaml(content)
|
||||
|
||||
|
||||
def test_unsupported_version_is_rejected() -> None:
|
||||
with pytest.raises(PeopleYamlError, match="Unsupported People YAML version"):
|
||||
import_people_yaml("version: 2\npeople: []\n")
|
||||
|
||||
|
||||
def test_missing_id_is_rejected() -> None:
|
||||
content = """\
|
||||
version: 1
|
||||
people:
|
||||
- display_name: Martin Tazl
|
||||
attendance_status: present
|
||||
"""
|
||||
|
||||
with pytest.raises(PeopleYamlError, match="participant_id"):
|
||||
import_people_yaml(content)
|
||||
|
||||
|
||||
def test_invalid_id_is_rejected_without_generating_a_replacement() -> None:
|
||||
content = """\
|
||||
version: 1
|
||||
people:
|
||||
- participant_id: invalid id
|
||||
display_name: Martin
|
||||
attendance_status: present
|
||||
"""
|
||||
|
||||
with pytest.raises(PeopleYamlError, match="invalid participant_id"):
|
||||
import_people_yaml(content)
|
||||
|
||||
|
||||
def test_duplicate_id_is_rejected() -> None:
|
||||
content = """\
|
||||
version: 1
|
||||
people:
|
||||
- participant_id: martin
|
||||
display_name: Martin
|
||||
attendance_status: present
|
||||
- participant_id: martin
|
||||
display_name: Martin Duplicate
|
||||
attendance_status: mentioned_only
|
||||
"""
|
||||
|
||||
with pytest.raises(PeopleYamlError, match="Duplicate participant_id"):
|
||||
import_people_yaml(content)
|
||||
|
||||
|
||||
def test_invalid_attendance_is_rejected() -> None:
|
||||
content = """\
|
||||
version: 1
|
||||
people:
|
||||
- participant_id: martin
|
||||
display_name: Martin
|
||||
attendance_status: absent
|
||||
"""
|
||||
|
||||
with pytest.raises(PeopleYamlError, match="invalid attendance_status"):
|
||||
import_people_yaml(content)
|
||||
|
||||
|
||||
def test_failed_import_does_not_mutate_existing_people() -> None:
|
||||
existing = sample_people()
|
||||
snapshot = list(existing)
|
||||
|
||||
with pytest.raises(PeopleYamlError):
|
||||
import_people_yaml("version: 1\npeople: not-a-list\n")
|
||||
|
||||
assert existing == snapshot
|
||||
@@ -0,0 +1,21 @@
|
||||
import pytest
|
||||
|
||||
from mka.application.meeting_service import ProcessingOptions
|
||||
from mka.application.performance import PERFORMANCE_PROFILES, resolve_ollama_num_thread
|
||||
|
||||
|
||||
def test_default_profile_is_auto() -> None:
|
||||
assert ProcessingOptions().performance_profile == "auto"
|
||||
assert PERFORMANCE_PROFILES[0] == "auto"
|
||||
|
||||
|
||||
def test_auto_profile_has_no_ollama_thread_override() -> None:
|
||||
assert resolve_ollama_num_thread("auto") is None
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("profile", "threads"),
|
||||
[("fast", 16), ("efficient", 10), ("powersave", 4)],
|
||||
)
|
||||
def test_profiles_resolve_to_current_ollama_thread_counts(profile: str, threads: int) -> None:
|
||||
assert resolve_ollama_num_thread(profile) == threads
|
||||
@@ -0,0 +1,45 @@
|
||||
from mka.application.progress_timing import ProcessingTimer
|
||||
|
||||
|
||||
class FakeClock:
|
||||
def __init__(self, value: float = 0.0) -> None:
|
||||
self.value = value
|
||||
|
||||
def __call__(self) -> float:
|
||||
return self.value
|
||||
|
||||
|
||||
def test_completed_duration_is_retained_while_active_stage_increases() -> None:
|
||||
clock = FakeClock(10.0)
|
||||
timer = ProcessingTimer(clock)
|
||||
timer.start_run()
|
||||
timer.start_stage("preparing")
|
||||
clock.value = 18.0
|
||||
timer.finish_stage("preparing")
|
||||
timer.start_stage("transcription")
|
||||
|
||||
clock.value = 20.0
|
||||
first = timer.snapshot()
|
||||
clock.value = 25.5
|
||||
second = timer.snapshot()
|
||||
|
||||
assert first.stage_durations == {"preparing": 8.0, "transcription": 2.0}
|
||||
assert second.stage_durations == {"preparing": 8.0, "transcription": 7.5}
|
||||
assert second.total_duration == 15.5
|
||||
|
||||
|
||||
def test_failed_active_stage_and_total_duration_are_frozen() -> None:
|
||||
clock = FakeClock(2.0)
|
||||
timer = ProcessingTimer(clock)
|
||||
timer.start_run()
|
||||
clock.value = 5.0
|
||||
timer.start_stage("transcription")
|
||||
clock.value = 14.0
|
||||
timer.finish_run()
|
||||
clock.value = 99.0
|
||||
|
||||
snapshot = timer.snapshot()
|
||||
|
||||
assert snapshot.stage_durations["transcription"] == 9.0
|
||||
assert snapshot.total_duration == 12.0
|
||||
assert snapshot.running is False
|
||||
@@ -0,0 +1,136 @@
|
||||
import json
|
||||
from datetime import date
|
||||
|
||||
import pytest
|
||||
|
||||
from mka.application.meeting_service import ParticipantInput
|
||||
from mka.application.run_inputs import (
|
||||
RunInputJsonError,
|
||||
RunInputState,
|
||||
export_run_inputs,
|
||||
import_run_inputs,
|
||||
run_input_filename,
|
||||
)
|
||||
|
||||
|
||||
def populated_state() -> RunInputState:
|
||||
return RunInputState(
|
||||
title="F&E Technik Abteilungs Jour Fixe",
|
||||
description="Review the production trial.",
|
||||
language="de",
|
||||
meeting_date=date(2026, 8, 25),
|
||||
participants=(
|
||||
ParticipantInput(
|
||||
participant_id="martin",
|
||||
display_name="Martin Tazl",
|
||||
role="Project lead",
|
||||
organization="Engineering",
|
||||
),
|
||||
ParticipantInput(
|
||||
participant_id="alex",
|
||||
display_name="Alexander Funk",
|
||||
attendance_status="mentioned_only",
|
||||
),
|
||||
),
|
||||
audio_normalization=False,
|
||||
diarization_enabled=True,
|
||||
source_file_name="2026-08-25_jour_fixe.wav",
|
||||
)
|
||||
|
||||
|
||||
def test_current_form_state_serializes_with_schema_and_all_supported_values() -> None:
|
||||
document = json.loads(export_run_inputs(populated_state()))
|
||||
|
||||
assert document["schema_version"] == 1
|
||||
assert document["meeting"]["title"] == "F&E Technik Abteilungs Jour Fixe"
|
||||
assert document["meeting"]["description"] == "Review the production trial."
|
||||
assert document["meeting"]["language"] == "de"
|
||||
assert document["meeting"]["date"] == "2026-08-25"
|
||||
assert document["meeting"]["participants"][1] == {
|
||||
"participant_id": "alex",
|
||||
"display_name": "Alexander Funk",
|
||||
"role": "",
|
||||
"organization": "",
|
||||
"attendance_status": "mentioned_only",
|
||||
}
|
||||
assert document["processing"] == {
|
||||
"audio_normalization": False,
|
||||
"diarization_enabled": True,
|
||||
}
|
||||
|
||||
|
||||
def test_export_import_round_trip_restores_form_and_participants() -> None:
|
||||
original = populated_state()
|
||||
|
||||
restored = import_run_inputs(export_run_inputs(original))
|
||||
|
||||
assert restored == original
|
||||
assert restored.participants[0].participant_id == "martin"
|
||||
assert restored.participants[0].display_name == "Martin Tazl"
|
||||
|
||||
|
||||
def test_missing_optional_fields_use_current_defaults() -> None:
|
||||
restored = import_run_inputs('{"schema_version": 1}')
|
||||
|
||||
defaults = RunInputState.defaults()
|
||||
assert restored == defaults
|
||||
|
||||
|
||||
def test_unknown_safe_fields_are_ignored() -> None:
|
||||
restored = import_run_inputs(
|
||||
json.dumps(
|
||||
{
|
||||
"schema_version": 1,
|
||||
"meeting": {"title": "Known", "future_field": {"value": 1}},
|
||||
"processing": {"future_toggle": True},
|
||||
"future_section": [1, 2, 3],
|
||||
}
|
||||
)
|
||||
)
|
||||
|
||||
assert restored.title == "Known"
|
||||
|
||||
|
||||
@pytest.mark.parametrize("content", ["not-json", "[]", b"\xff"])
|
||||
def test_malformed_json_is_rejected(content: str | bytes) -> None:
|
||||
with pytest.raises(RunInputJsonError, match="Malformed|top-level"):
|
||||
import_run_inputs(content)
|
||||
|
||||
|
||||
def test_unsupported_future_schema_is_rejected() -> None:
|
||||
with pytest.raises(RunInputJsonError, match="Unsupported.*schema_version"):
|
||||
import_run_inputs('{"schema_version": 2}')
|
||||
|
||||
|
||||
def test_media_contents_are_never_serialized() -> None:
|
||||
exported = export_run_inputs(populated_state())
|
||||
|
||||
assert "2026-08-25_jour_fixe.wav" in exported
|
||||
assert "audio_bytes" not in exported
|
||||
assert "base64" not in exported
|
||||
|
||||
|
||||
def test_source_filename_must_not_be_a_machine_specific_path() -> None:
|
||||
with pytest.raises(RunInputJsonError, match="filename, not a path"):
|
||||
import_run_inputs('{"schema_version": 1, "source_file_name": "/tmp/meeting.wav"}')
|
||||
|
||||
|
||||
def test_invalid_or_duplicate_participants_are_not_silently_reinterpreted() -> None:
|
||||
document = {
|
||||
"schema_version": 1,
|
||||
"meeting": {
|
||||
"participants": [
|
||||
{"participant_id": "martin", "display_name": "Martin"},
|
||||
{"participant_id": "martin", "display_name": "Someone else"},
|
||||
]
|
||||
},
|
||||
}
|
||||
|
||||
with pytest.raises(RunInputJsonError, match="Duplicate participant_id"):
|
||||
import_run_inputs(json.dumps(document))
|
||||
|
||||
|
||||
def test_export_filename_is_readable_and_machine_independent() -> None:
|
||||
assert run_input_filename("F&E Technik Jour Fixe") == (
|
||||
"meeting-inputs-f-e-technik-jour-fixe.json"
|
||||
)
|
||||
@@ -0,0 +1,134 @@
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
from streamlit.testing.v1 import AppTest
|
||||
|
||||
from mka.ui.streamlit_app import _mapping_counts, _speaker_options
|
||||
|
||||
|
||||
@pytest.mark.parametrize("label", ["SPEAKER_00", "SPEAKER_01", "SPEAKER_02"])
|
||||
def test_unmapped_options_include_all_participants(label):
|
||||
assert _speaker_options(("a", "b"), {}, label) == [None, "a", "b"]
|
||||
|
||||
|
||||
def test_options_reserve_other_assignments_and_preserve_current():
|
||||
mappings = {"SPEAKER_00": "a"}
|
||||
assert _speaker_options(("a", "b"), mappings, "SPEAKER_00") == [None, "a", "b"]
|
||||
for label in ("SPEAKER_01", "SPEAKER_02"):
|
||||
assert _speaker_options(("a", "b"), mappings, label) == [None, "b"]
|
||||
mappings["SPEAKER_00"] = "b"
|
||||
assert _speaker_options(("a", "b"), mappings, "SPEAKER_01") == [None, "a"]
|
||||
mappings["SPEAKER_00"] = None
|
||||
assert _speaker_options(("a", "b"), mappings, "SPEAKER_01") == [None, "a", "b"]
|
||||
|
||||
|
||||
def test_existing_duplicate_or_missing_participant_is_not_removed():
|
||||
assert _speaker_options(("b",), {"s0": "a", "s1": "a"}, "s0") == [None, "b", "a"]
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("mappings", "expected"),
|
||||
[
|
||||
({}, (2, 0, 2)),
|
||||
({"s0": "a", "other": "b"}, (2, 1, 1)),
|
||||
({"s0": "a", "s1": "b"}, (2, 2, 0)),
|
||||
({"s0": None}, (2, 0, 2)),
|
||||
],
|
||||
)
|
||||
def test_counts_only_include_detected_speakers(mappings, expected):
|
||||
assert _mapping_counts(("s0", "s1"), mappings) == expected
|
||||
|
||||
|
||||
def mapping_app():
|
||||
from mka.application.meeting_service import SpeakerMappingReview, SpeakerReview
|
||||
from mka.ui.streamlit_app import _render_speaker_mapping
|
||||
|
||||
_render_speaker_mapping(
|
||||
SpeakerMappingReview(
|
||||
tuple(SpeakerReview(f"SPEAKER_0{i}", ()) for i in range(3)),
|
||||
(("a", "Participant A"), ("b", "Participant B"), ("c", "Participant C")),
|
||||
{"SPEAKER_00": "a"},
|
||||
),
|
||||
"test",
|
||||
)
|
||||
|
||||
|
||||
def test_widget_reruns_filter_release_and_update_status():
|
||||
app = AppTest.from_function(mapping_app).run()
|
||||
assert not app.exception
|
||||
assert app.selectbox[0].value == "a"
|
||||
assert "Participant A" not in app.selectbox[1].options
|
||||
assert len(app.warning) == 1
|
||||
assert "3 speakers detected · 1 assigned · 2 unassigned" in app.markdown[0].value
|
||||
assert "Unassigned" in app.selectbox[1].label
|
||||
|
||||
app.selectbox[2].select("c").run()
|
||||
assert app.selectbox[0].value == "a"
|
||||
assert "Participant C" not in app.selectbox[0].options
|
||||
app.selectbox[1].select("b").run()
|
||||
assert not app.warning
|
||||
assert "3 assigned · 0 unassigned" in app.markdown[0].value
|
||||
app.selectbox[0].select(None).run()
|
||||
assert app.warning
|
||||
assert "Participant A" in app.selectbox[1].options
|
||||
assert app.selectbox[1].value == "b"
|
||||
app.selectbox[1].select("a").run()
|
||||
assert "Participant B" in app.selectbox[0].options
|
||||
assert "Participant A" not in app.selectbox[0].options
|
||||
assert app.selectbox[0].value is None
|
||||
assert not app.exception
|
||||
|
||||
|
||||
def test_checkpoint_allows_anonymous_protocol_generation(monkeypatch):
|
||||
from mka.application.meeting_service import SpeakerMappingReview, SpeakerReview
|
||||
from mka.ui import streamlit_app as ui
|
||||
|
||||
outcome = SimpleNamespace(
|
||||
succeeded=True,
|
||||
run_dir=Path("/tmp/mapping-test"),
|
||||
speaker_attribution_available=True,
|
||||
original_protocol="Anonymous protocol",
|
||||
awaiting_speaker_review=True,
|
||||
detected_speaker_count=1,
|
||||
)
|
||||
review = SpeakerMappingReview((SpeakerReview("SPEAKER_00", ()),), (("a", "A"),), {})
|
||||
calls = []
|
||||
|
||||
class Service:
|
||||
def __init__(self, *args):
|
||||
pass
|
||||
|
||||
def load_speaker_mapping_review(self, run_dir):
|
||||
return review
|
||||
|
||||
def regenerate_protocol(self, run_dir, selections, **kwargs):
|
||||
calls.append(selections)
|
||||
return outcome
|
||||
|
||||
monkeypatch.setattr(ui, "MeetingProcessingService", Service)
|
||||
monkeypatch.setattr(ui.AppSettings, "from_environment", lambda: None)
|
||||
monkeypatch.setattr(ui, "MeetingLabGateway", lambda: None)
|
||||
monkeypatch.setattr(ui, "_remember_regeneration_timing", lambda *args: None, raising=False)
|
||||
monkeypatch.setattr(
|
||||
ui,
|
||||
"_run_with_live_progress",
|
||||
lambda service, action, states, message: (action(None), states, message, None),
|
||||
raising=False,
|
||||
)
|
||||
|
||||
def result_app(outcome):
|
||||
import streamlit as st
|
||||
|
||||
from mka.ui.streamlit_app import _render_result
|
||||
|
||||
st.session_state.outcome = outcome
|
||||
_render_result()
|
||||
|
||||
app = AppTest.from_function(result_app, args=(outcome,)).run()
|
||||
generate = next(button for button in app.button if button.label == "Generate protocol")
|
||||
assert not generate.disabled
|
||||
assert app.warning
|
||||
generate.click().run()
|
||||
assert calls == [{}]
|
||||
assert not app.exception
|
||||
@@ -0,0 +1,95 @@
|
||||
from datetime import date
|
||||
|
||||
from mka.application.meeting_service import ParticipantInput
|
||||
from mka.application.progress_timing import TimingSnapshot
|
||||
from mka.application.run_inputs import RunInputState
|
||||
from mka.ui import streamlit_app
|
||||
|
||||
|
||||
def test_alias_input_accepts_lines_and_commas() -> None:
|
||||
assert streamlit_app._parse_aliases("Sikirgut\nSekugrid HS, Secugrid H S") == (
|
||||
"Sikirgut",
|
||||
"Sekugrid HS",
|
||||
"Secugrid H S",
|
||||
)
|
||||
|
||||
|
||||
def test_regenerated_protocol_widget_value_is_deferred_until_next_run(
|
||||
monkeypatch,
|
||||
) -> None:
|
||||
state = {"edited_protocol": "old protocol"}
|
||||
monkeypatch.setattr(streamlit_app.st, "session_state", state)
|
||||
|
||||
streamlit_app._queue_edited_protocol("regenerated protocol")
|
||||
|
||||
assert state["edited_protocol"] == "old protocol"
|
||||
assert state["pending_edited_protocol"] == "regenerated protocol"
|
||||
|
||||
streamlit_app._apply_pending_edited_protocol()
|
||||
|
||||
assert state["edited_protocol"] == "regenerated protocol"
|
||||
assert "pending_edited_protocol" not in state
|
||||
|
||||
|
||||
def test_imported_inputs_are_applied_via_pending_state_before_widgets(
|
||||
monkeypatch,
|
||||
) -> None:
|
||||
existing_upload = object()
|
||||
state = {"source_media_0": existing_upload, "source_media_widget_generation": 0}
|
||||
monkeypatch.setattr(streamlit_app.st, "session_state", state)
|
||||
imported = RunInputState(
|
||||
title="Imported meeting",
|
||||
description="Imported context",
|
||||
language="en",
|
||||
meeting_date=date(2026, 8, 25),
|
||||
participants=(ParticipantInput("martin", "Martin"),),
|
||||
audio_normalization=False,
|
||||
diarization_enabled=True,
|
||||
source_file_name="meeting.wav",
|
||||
)
|
||||
|
||||
streamlit_app._queue_run_inputs(imported)
|
||||
|
||||
assert "meeting_title" not in state
|
||||
assert state["pending_run_inputs"] == imported
|
||||
|
||||
streamlit_app._apply_pending_run_inputs()
|
||||
|
||||
assert state["meeting_title"] == "Imported meeting"
|
||||
assert state["meeting_description"] == "Imported context"
|
||||
assert state["meeting_language"] == "en"
|
||||
assert state["meeting_has_date"] is True
|
||||
assert state["participants"][0]["participant_id"] == "martin"
|
||||
assert state["audio_normalization"] is False
|
||||
assert state["diarization_enabled"] is True
|
||||
assert state["imported_source_file_name"] == "meeting.wav"
|
||||
assert state["source_media_widget_generation"] == 1
|
||||
assert "source_media_1" not in state
|
||||
assert "pending_run_inputs" not in state
|
||||
assert state["source_media_0"] is existing_upload
|
||||
|
||||
|
||||
def test_duration_formatting_covers_seconds_minutes_and_hours() -> None:
|
||||
assert streamlit_app._format_duration(8.9) == "00:08"
|
||||
assert streamlit_app._format_duration(7 * 60 + 18) == "07:18"
|
||||
assert streamlit_app._format_duration(3600 + 3 * 60 + 42) == "1:03:42"
|
||||
|
||||
|
||||
def test_completed_regeneration_timing_survives_frontend_rerun(monkeypatch) -> None:
|
||||
state = {}
|
||||
monkeypatch.setattr(streamlit_app.st, "session_state", state)
|
||||
timing = TimingSnapshot(
|
||||
stage_durations={"protocol_generation": 4.5},
|
||||
active_stage=None,
|
||||
total_duration=4.5,
|
||||
running=False,
|
||||
)
|
||||
states = {stage: "skipped" for stage in streamlit_app.STAGES}
|
||||
states["protocol_generation"] = "completed"
|
||||
|
||||
streamlit_app._remember_regeneration_timing(states, "Protocol regenerated", timing)
|
||||
|
||||
saved_states, saved_message, saved_timing = state["regeneration_progress"]
|
||||
assert saved_states["protocol_generation"] == "completed"
|
||||
assert saved_message == "Protocol regenerated"
|
||||
assert saved_timing == timing
|
||||
Reference in New Issue
Block a user