Author SHA1 Message Date
admin 2215414282 Add post-diarization speaker review workflow 2026-09-17 19:16:57 +02:00
admin bd48aec5aa Release Meeting Assistant v0.1.0-alpha.1 2026-09-12 12:37:44 +02:00
admin b50b87c32b Document alpha validation and historical fixture dependency 2026-09-12 11:59:00 +02:00
admin 9e7173c2a7 Validate paired alpha workflow without inference 2026-09-12 11:57:08 +02:00
admin 70cb1e0789 Retain live processing timings across UI reruns 2026-09-12 11:55:19 +02:00
admin 500d2c4c67 Forward glossary provenance without transcript substitution 2026-09-12 11:51:38 +02:00
admin 77be2b39ff Add protocol performance profiles 2026-09-12 11:50:30 +02:00
admin e2a0434941 Add versioned glossary YAML import and export 2026-09-12 11:23:18 +02:00
admin 9840953687 Improve speaker mapping UX 2026-09-12 11:22:09 +02:00
admin e55dd74c19 Treat language selection as meeting language 2026-09-11 10:25:45 +02:00
admin dd8a618719 feat: add speaker mapping and workflow improvements 2026-08-25 15:38:58 +02:00
admin adf454d77d feat: configure 32k context for protocol generation 2026-08-25 15:31:33 +02:00
admin 303563abf1 Add container diarization configuration 2026-08-24 15:59:18 +02:00
admin 2a4481f2f9 Add reusable people list import and export 2026-08-24 14:26:30 +02:00
admin 13d06b8a60 Add Streamlit meeting assistant MVP 2026-08-24 10:28:44 +02:00
admin 443a85528b Merge branch 'docs/reconcile-meeting-lab-architecture' 2026-08-23 20:19:55 +02:00
admin 490d557cb3 Reconcile Meeting Assistant architecture with Meeting Lab 2026-08-23 20:19:47 +02:00
admin e028f1c82c Document development workflow 2026-08-23 19:56:15 +02:00
47 changed files with 5444 additions and 586 deletions
+18
View File
@@ -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"]'
+2
View File
@@ -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
+31
View File
@@ -1,5 +1,36 @@
# Changelog
## [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.
## [0.1.0] - 2026-07-14
### Added
+38
View File
@@ -0,0 +1,38 @@
# Contributing
1. Branch Strategy
2. Development Workflow
3. Commit Messages
4. Definition of Done
5. Coding Standards
6. Documentation Requirements
7. Testing
8. Creating ADRs
## Development Workflow
See ADR 0010 for the overall development workflow.
---
## Commit Messages
Commit messages should be:
- short
- descriptive
- written in the imperative mood
Examples:
- Add meeting import
- Implement transcript persistence
- Refactor storage layer
- Fix timestamp serialization
Avoid messages such as:
- Update
- Changes
- Fix
- Miscellaneous
+176 -153
View File
@@ -1,163 +1,186 @@
# 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
The project creates a complete pipeline that transforms spoken meetings into structured organizational knowledge.
Meeting Assistant turns recorded meetings into reviewed, user-facing protocols
and, over time, searchable organizational knowledge. It is the application
layer over the reusable Meeting Lab processing backend.
The transcript is the canonical source.
The source audio and canonical transcript are preserved. Generated protocols
are derived, reproducible artifacts and require human review.
AI-generated summaries are reproducible artefacts.
## System Boundary
---
### Meeting Lab
Meeting Lab owns reusable and experimental processing:
- FFmpeg audio preparation, with optional normalization enabled by default
- `whisper.cpp` transcription
- optional `pyannote.audio` diarization
- protocol generation
- pipeline orchestration and progress events
Its reusable integration surface is the Python API:
- `MvpMeetingConfig`
- `MvpRunResult`
- `run_mvp_meeting(...)`
The CLI is a thin adapter, not the Meeting Assistant integration boundary.
### Meeting Assistant
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
-> 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
-> 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
```
Mandatory fixed five-minute chunking is not part of the productive MVP path.
Chunking or segmentation used internally by an engine does not shape the
Meeting Assistant storage or domain model.
On the North reference machine, `whisper.cpp` with Vulkan on an AMD RX 9070
transcribed approximately 93.5 minutes in about 100 seconds. Pyannote through
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,
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 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:
- `preparing`
- `transcription`
- `diarization`
- `protocol_generation`
- `completed`
- `failed`
The GUI consumes this event interface. Percentages are shown only when real
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.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.
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.
Replaceable AI components.
Vendor independence.
Modular architecture.
Small focused modules.
Simple interfaces.
---
## Core Pipeline
Recorder
↓
Transcription
↓
Speaker Identification
↓
LLM Analysis
↓
Knowledge Extraction
↓
Storage
↓
Export
---
## AI Components
Recorder
Responsible only for recording.
No AI.
---
Transcription
Responsible only for speech-to-text.
No summarization.
---
Speaker Identification
Responsible only for identifying speakers.
Must not modify transcript text.
---
LLM Analysis
Responsible for:
- Summary
- Decisions
- Action Items
- Risks
- Questions
- Knowledge extraction
---
Storage
Stores:
Audio
Transcript
Metadata
AI Results
Knowledge Objects
---
## Architectural Rule
Every module has exactly one responsibility.
---
## Transcript Policy
Never modify the original transcript.
Corrections must create a new derived version.
---
## AI Output Policy
Every AI output should be reproducible.
Prompt version should be stored.
Model should be stored.
Timestamp should be stored.
---
## Long-Term Goal
Every meeting becomes searchable organizational knowledge.
No information should be lost after the meeting.
---
## Engineering Principles
The project follows an architecture-first development approach.
Before implementing a feature:
- define the domain model
- define module boundaries
- document architectural decisions
Implementation is intentionally delayed until the architecture is considered sufficiently stable.
- offline-first where practical
- explicit application/backend boundaries
- replaceable AI components
- CPU-compatible operation with optional acceleration
- immutable source artifacts and traceable derived artifacts
- 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.
+225 -87
View File
@@ -1,126 +1,264 @@
# Meeting Knowledge Assistant
# Meeting Assistant
## Vision
Meeting Assistant is the user-facing application for turning recorded meetings
into reviewed, distributable protocols. It uses Meeting Lab as its reusable
local processing backend.
The Meeting Knowledge Assistant (MKA) is an AI-assisted desktop application that records meetings, creates high-quality transcripts and transforms them into structured knowledge.
## Current MVP Direction
Unlike traditional meeting assistants, MKA does not aim to replace human interaction during meetings.
Its goal is to create an accurate, searchable and long-term knowledge base from spoken communication.
The first product milestone is a desktop GUI that lets a user:
The project follows an "AI-first" architecture:
- 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
- review and edit the generated protocol
- export the reviewed result
Audio
→ Transcription
→ Speaker Identification
→ AI Analysis
→ Structured Knowledge
→ Search
The MVP does not include integrated recording. OBS Studio remains a recommended
reference recorder, but imported audio is not tied to OBS-specific behavior.
---
## Processing Pipeline
## Core Features
- Import of existing meeting recordings
- OBS Studio as the recommended reference recorder
- High-quality local transcription
- Speaker diarization
- Editable working transcripts
- AI-generated meeting minutes
- Executive summaries
- Action items
- Decision tracking
- Knowledge extraction
- Local file-based storage
- Export to Markdown, PDF and DOCX
---
## Recording Workflow
The MVP does not include its own audio recorder.
OBS Studio is the recommended reference recorder for online meetings. It can
capture system audio and microphone audio without joining the meeting as a bot.
Initial workflow:
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; normalization optional)
-> whisper.cpp transcription (large-v3-turbo)
-> optional pyannote.audio Community-1 diarization
-> optional speaker review and explicit protocol generation (when diarized)
or direct full-transcript protocol generation (without diarization)
-> human review and export
```
Online Meeting
Meeting Assistant calls Meeting Lab's reusable Python API
(`MvpMeetingConfig`, `MvpRunResult` and `run_mvp_meeting(...)`). It does not run
the Meeting Lab CLI as a subprocess. Fixed five-minute audio chunking is not a
required part of this workflow.
↓
GPU acceleration is supported where available, but CPU execution remains a
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.
OBS Studio
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.
WAV Recording
## Product Outputs
↓
The product direction includes:
Meeting Assistant Import
- a detailed, contextual protocol for review and continued work
- a shorter version suitable for participants and distribution
↓
The detailed direct-protocol flow is the practical MVP direction. The shorter
version may be implemented after the initial GUI. Human review remains part of
the workflow for all generated protocols.
Transcription
## Project Boundary
↓
Meeting Lab owns reusable audio preparation, transcription, diarization,
orchestration and protocol-generation logic. Meeting Assistant owns the GUI,
context and participant editing, explicit speaker mapping, progress display,
protocol editing and export.
Speaker Diarization
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.
Analysis and Export
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
## Design Goals
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.
- High transcription quality
- Modular architecture
- Offline-first whenever practical
- Replaceable AI components
- Vendor independence
- Long-term maintainability
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
```
## Philosophy
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:
The transcript is the primary asset.
```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
```
Everything else (summaries, reports, action items, knowledge extraction)
can always be regenerated using better AI models in the future.
Optional machine-specific settings include:
Therefore:
- `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
Audio
→ Transcript
is considered immutable.
## Terminology glossary
AI output is considered reproducible.
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.
## Planned Architecture
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:
Recorder
↓
Transcription
↓
Speaker Identification
↓
LLM Processing
↓
Knowledge Database
↓
Export
```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`.
## Current Status
## Input configuration import and export
Project planning.
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.
No implementation has started yet.
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.
+49 -99
View File
@@ -1,112 +1,62 @@
# Roadmap
## Phase 1: Minimum Viable Processing Pipeline
## Next Milestone: Meeting Assistant MVP GUI
- Import existing WAV recordings
- Validate recording metadata
- Create meeting directory
- Create canonical transcript
- Save transcript as JSON
- Export transcript as Markdown
- Archive confirmed recordings as FLAC
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 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
- handle completed and failed runs clearly
- display and edit the detailed generated protocol
- export the reviewed protocol
## Phase 2
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.
Speaker diarization
Complete real-meeting validation of this vertical slice before expanding the
post-run workflow.
- Speaker detection
- Speaker naming
- Timeline view
## Following Milestone: Distribution Protocol
---
- derive or generate a shorter participant-facing version
- allow review and editing before distribution
- export both detailed and short versions
- retain provenance linking each version to its transcript, configuration,
prompt and model where practical
## Phase 3
## Later Product Work
Meeting intelligence
- 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
- optional SQLite indexing and full-text search
- integrated recording controls
- calendar and conferencing integrations
- semantic search and organizational knowledge features
- Executive Summary
- Action Items
- Decisions
- Risks
- Open Questions
---
## Phase 4
Knowledge system
- SQLite database
- Full text search
- Semantic search
- Tags
- Projects
- Participants
---
## Phase 5
Desktop application
- Recording UI
- Transcript viewer
- Search
- Export
---
## Phase 6
Integrations
- Outlook
- Teams
- Google Calendar
- Local LLM
- OpenAI
- Ollama
---
## Future Phase: Integrated Recording
- Capture microphone audio
- Capture system audio
- Support separate audio channels
- Provide recording controls in the desktop application
- Replace or complement the external OBS workflow
- Live transcription
- Real-time summaries
- Company glossary
- Custom vocabulary
- Automatic project detection
- Meeting templates
- Voice identification
- Audio cleanup
## Research and Optional Extensions
- 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
- Timeline with bookmarks
- RAG knowledge integration
- Multi-language meetings
- Translation
- Automatic follow-up generation
- RAG and knowledge-graph integration
- multi-language meetings and translation
- dual-model or evidence pipelines only if later validation demonstrates a
material reliability benefit
View File
+7 -2
View File
@@ -2,7 +2,7 @@
## Status
Accepted
Accepted, amended by [ADR 0011](0011-use-meeting-lab-mvp-backend.md)
## Context
@@ -29,4 +29,9 @@ The transcript is treated as the canonical source. AI-generated summaries and kn
- The system must preserve original recordings and original transcripts.
- AI outputs must be reproducible where practical.
- Modules should remain replaceable.
- The project should avoid unnecessary vendor lock-in.
- The project should avoid unnecessary vendor lock-in.
ADR 0011 narrows the practical MVP to a user-facing application over the
reusable Meeting Lab backend. Direct protocol generation is the validated MVP
path; a separate knowledge-extraction stage remains a possible future
extension rather than an MVP requirement.
+1 -1
View File
@@ -2,7 +2,7 @@
## Status
Accepted
Superseded by [ADR 0011](0011-use-meeting-lab-mvp-backend.md)
## Context
+9 -1
View File
@@ -2,7 +2,7 @@
## Status
Accepted
Accepted, amended by [ADR 0011](0011-use-meeting-lab-mvp-backend.md)
## Context
@@ -23,3 +23,11 @@ Speaker diarization will be treated as a separate pipeline step after transcript
- Human correction of speaker names should be supported later.
- The diarization module must hide the concrete engine behind an internal interface.
- Model name, engine version and timestamp should be stored with every diarization result.
## Amendment
ADR 0011 makes diarization optional for the MVP and selects `pyannote.audio`
Community-1 as the current preferred backend. CPU execution remains supported,
with GPU acceleration used when available. Speaker labels remain anonymous
unless a user explicitly confirms a `SPEAKER_XX` to participant mapping in the
Meeting Context. Automatic speaker-name inference is not allowed.
+3 -3
View File
@@ -2,7 +2,7 @@
## Status
Accepted
Accepted, amended by [ADR 0011](0011-use-meeting-lab-mvp-backend.md)
## Context
@@ -18,8 +18,8 @@ removing temporary source files without risking data loss.
Recordings shall initially be captured in a processing-friendly lossless
format, normally WAV.
After transcription and diarization have completed, the recording may be
converted to an archive format.
After transcription and, when enabled, diarization have completed, the
recording may be converted to an archive format.
The preferred archive formats are:
@@ -2,7 +2,7 @@
## Status
Accepted
Accepted, amended by [ADR 0011](0011-use-meeting-lab-mvp-backend.md)
## Context
@@ -13,8 +13,9 @@ OBS Studio already provides stable, configurable and cross-platform audio
capture. Reimplementing audio capture during the initial project phase would
add substantial complexity without directly improving transcription quality.
The primary value of the Meeting Assistant lies in transcription, speaker
diarization, analysis, structured knowledge extraction and export.
The primary value of the Meeting Assistant lies in transcription, optional
speaker diarization, protocol generation, review and export. ADR 0011 defines
the validated processing boundary and current MVP protocol path.
## Decision
+55
View File
@@ -0,0 +1,55 @@
# ADR 0010: Development Workflow
## Status
Accepted
## Context
The project is expected to evolve over an extended period and will be developed with the assistance of AI coding tools.
To maintain code quality, architectural consistency and a comprehensible project history, a common development workflow is required.
## Decision
The project follows a feature-branch based development workflow.
The `main` branch shall always represent a stable and working state.
Development of new functionality shall take place on dedicated feature branches.
Examples:
- feature/domain-model
- feature/storage
- feature/import
- feature/transcription
- feature/diarization
- feature/export
- feature/ui
Feature branches shall be merged into `main` only after:
- implementation is complete
- documentation has been updated where required
- all tests pass
- static analysis completes without errors
Architectural changes shall be documented using Architecture Decision Records (ADRs).
Project versioning follows Semantic Versioning.
Version numbers shall be maintained in:
- `pyproject.toml`
- `CHANGELOG.md`
Each released version shall be tagged in Git.
## Consequences
- Development history remains easy to understand.
- Experimental work is isolated from the stable branch.
- Architectural decisions remain documented.
- Releases are reproducible through Git tags.
@@ -0,0 +1,94 @@
# ADR 0011: Use Meeting Lab as the MVP Processing Backend
## Status
Accepted
Supersedes ADR 0002 and amends ADRs 0001 and 0003.
## Context
Meeting Lab experiments have validated a practical local pipeline for turning
a recorded meeting into an editable protocol. Earlier Meeting Assistant plans
selected `faster-whisper`, treated diarization as required and described a
separate knowledge-extraction stage as part of the primary pipeline. Work on an
old feature branch also proposed mandatory fixed five-minute audio chunks and
chunk-oriented storage.
The validated backend now has a reusable Python API, optional diarization and a
direct full-transcript protocol path. The product boundary must keep backend
processing in Meeting Lab and user interaction in Meeting Assistant.
## Decision
Meeting Lab is the reusable processing backend and experimental engine.
Meeting Assistant is the user-facing application and calls Meeting Lab through
its Python API rather than launching its CLI as a subprocess.
The MVP processing flow is:
```text
Source audio
-> FFmpeg preparation (mono, 16 kHz PCM WAV)
-> whisper.cpp transcription (large-v3-turbo)
-> optional pyannote.audio Community-1 diarization
-> direct full-transcript protocol generation
-> human review and editing
```
The productive path does not require fixed five-minute audio chunks. Internal
streaming or segmentation remains an implementation detail of Meeting Lab and
must not define Meeting Assistant's domain or storage model.
Meeting Assistant integrates with these public Meeting Lab interfaces:
- `MvpMeetingConfig`
- `MvpRunResult`
- `run_mvp_meeting(...)`
- stage-based progress events: `preparing`, `transcription`, `diarization`,
`protocol_generation`, `completed` and `failed`
Progress percentages are displayed only when the backend reports real
measurable progress.
`MeetingContext` is the structured input containing meeting metadata and
participants. It may contain explicitly confirmed `SPEAKER_XX -> participant_id`
mappings. Such mappings are authoritative only when confirmed by the user.
The system must not infer speaker names automatically.
The preferred local engines are:
- `whisper.cpp` with `large-v3-turbo` for transcription
- `pyannote.audio` Community-1 for optional diarization
GPU acceleration is optional. Vulkan transcription and PyTorch/ROCm
diarization are validated acceleration paths on North's AMD RX 9070. CPU
execution remains the compatibility fallback.
The direct full-transcript protocol path is the practical MVP direction.
`qwen3.8:27B` has produced strong readability and contextual synthesis, while
`qwen3.6:35B-A3B` has sometimes behaved more conservatively. Dual-model and
diarization-assisted hard-fact experiments have not established sufficiently
reliable attribution gains to justify mandatory MVP complexity. Model choice
remains backend configuration, and human review is required.
Meeting Assistant owns the GUI, structured context editor, participant and
speaker-mapping controls, progress presentation, protocol editor and export.
Meeting Lab owns audio preparation, transcription, diarization, orchestration
and protocol-generation logic.
The product should ultimately provide a detailed contextual protocol and a
shorter participant/distribution version. The shorter version may follow the
first GUI milestone.
## Consequences
- ADR 0002's `faster-whisper` selection is no longer current.
- Fixed-duration chunking and chunk-centric storage are not MVP requirements.
- Diarization and GPU acceleration remain optional.
- The Meeting Lab CLI remains useful for command-line operation but is not the
application integration boundary.
- Meeting Context is created through the future GUI; users are not expected to
hand-write YAML.
- Backend experiments can evolve without moving processing logic into the GUI.
- Protocols remain AI-assisted outputs that require human review.
+67
View File
@@ -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.
+40
View File
@@ -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.
+138 -233
View File
@@ -2,254 +2,159 @@
## Purpose
The Meeting Knowledge Assistant transforms recorded meetings into structured organizational knowledge through a modular processing pipeline.
Meeting Assistant is the user-facing application for preparing meeting
context, running the Meeting Lab pipeline and reviewing its results. Meeting
Lab is the reusable processing backend and experimental engine.
---
# Design Principles
- Architecture first
- Offline first where practical
- Immutable source data
- Replaceable AI components
- Small, focused modules
- Explicit interfaces
- Reproducible AI outputs
- External recording and internal processing are separate responsibilities.
- The processing pipeline must not depend on OBS-specific metadata or behavior.
- Imported recordings must be treated like recordings from any other supported source.
---
## High-Level Pipeline
External Recorder
↓
Audio Import
↓
Recording Validation
↓
Meeting Storage
↓
Transcription
↓
Speaker Diarization
↓
Working Transcript
↓
LLM Analysis
↓
Knowledge Extraction
↓
Export
---
## Core Pipeline
External Recording
↓
Audio Import and Validation
↓
Transcription
↓
Speaker Diarization
↓
Working Transcript
↓
LLM Analysis
↓
Knowledge Extraction
↓
Storage and Export
# Domain Model
Meeting
│
├── Recording
├── Transcript
│ ├── Canonical
│ └── Working
├── Speakers
├── Participants
├── AI Artifacts
├── Knowledge Objects
├── Attachments
└── Exports
---
# Module Responsibilities
## Input and Ingest
Responsible for importing existing recordings into the application.
Initial reference source:
- OBS Studio
Initial supported formats:
- WAV
- FLAC
Responsibilities:
- validate the input file
- collect technical metadata
- calculate checksums
- copy or move the recording into the meeting directory
- create the initial Recording domain object
The ingest component must not perform transcription or modify the audio content.
## Recorder
The recorder module is reserved for a future integrated recording implementation.
It is not required for the MVP.
The initial application workflow uses externally created recordings, with OBS
Studio as the recommended reference recorder.
---
## Transcription
Responsible only for speech-to-text conversion.
Input:
Recording
Output:
Canonical Transcript
---
## Diarization
Responsible only for speaker identification.
Input:
Recording + Canonical Transcript
Output:
Working Transcript
---
## LLM
Responsible for semantic analysis.
Input:
Transcript
Output:
AI Artifacts
---
## Knowledge Extraction
Responsible for creating structured knowledge.
Input:
AI Artifacts
Output:
Knowledge Objects
---
## Export
Responsible for creating user-facing documents.
Input:
Knowledge Objects
Output:
Markdown
PDF
DOCX
---
## Artifact Lifecycle
## System Boundary
```text
Imported Recording
↓
Validated Source Recording
↓
Canonical Transcript
↓
Working Transcript
↓
AI Artifacts
↓
Knowledge Objects
↓
Exports
Meeting Assistant Meeting Lab
----------------- -----------
Audio selection ------> Audio preparation
Meeting Context editor ------> Transcription
Participant management ------> Optional diarization
Explicit speaker mapping ------> Protocol generation
Progress presentation <------ Progress events
Protocol editor and export <------ Run result and artifacts
```
---
Meeting Assistant calls the Meeting Lab Python API directly. The Meeting Lab
CLI is a thin adapter over that API and must not be launched as an application
subprocess.
## Initial Recording Strategy
## MVP Processing Flow
The MVP does not implement platform-specific audio capture.
```text
Source audio
-> FFmpeg preparation (normalization optional, default on)
-> mono, 16 kHz PCM WAV
-> whisper.cpp transcription with large-v3-turbo
-> optional pyannote.audio Community-1 diarization
-> direct full-transcript protocol generation
-> human review and editing
-> export
```
OBS Studio is the recommended reference recorder for online meetings.
The MVP does not require fixed five-minute audio chunks, semantic chunking, a
separate evidence-extraction pipeline or multiple LLMs. Any internal
segmentation remains a Meeting Lab implementation detail.
The application initially processes existing WAV or FLAC recordings. Integrated
## Backend Integration
recording remains a future extension and must not be required by transcription,
The reusable Meeting Lab interface consists of:
diarization or analysis modules.
- `MvpMeetingConfig`, the run configuration
- `run_mvp_meeting(...)`, the orchestration entry point
- `MvpRunResult`, the completed run result
- stage-based progress events
---
The known stages are:
# Storage Strategy
```text
preparing
transcription
diarization
protocol_generation
completed
failed
```
The domain model is independent of the storage backend.
The GUI shows the current stage. It shows a percentage only when the event
contains real measurable progress; stage changes must not be presented as
invented percentages.
Possible implementations:
## Meeting Context
- File System
- SQLite
- PostgreSQL
- Cloud Storage
`MeetingContext` is structured domain input rather than an informal prompt or a
file users must author manually. It contains meeting metadata and participants
and can include optional mappings:
---
```text
SPEAKER_XX -> participant_id
```
# Future Extensions
Mappings are authoritative only after explicit user confirmation. Diarization
labels otherwise remain anonymous, and the application must not infer speaker
names automatically. The Meeting Assistant GUI owns creation and editing of
this context and its mappings.
- Live transcription
- Video processing
- OCR
- Semantic search
- Knowledge graph
- Company glossary
- Multi-language meetings
- Local LLM support
## Processing Components
### Audio Preparation
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. 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
`whisper.cpp` with `large-v3-turbo` is the preferred local backend. Vulkan is a
validated acceleration path on North's AMD RX 9070, while CPU execution remains
a compatibility fallback. Engine and model details should be retained with the
result for traceability.
### Diarization
`pyannote.audio` Community-1 is preferred when diarization is enabled. It can
run on CPU and may use GPU acceleration such as PyTorch/ROCm where available.
Diarization is optional and its anonymous labels do not modify the canonical
transcript or assert participant identity.
### Protocol Generation
The direct full-transcript path is the practical MVP. Current model experiments
favor `qwen3.8:27B` for readability and contextual synthesis, while
`qwen3.6:35B-A3B` has shown more conservative behavior in some areas. Neither a
specific dual-model arrangement nor an evidence pipeline is an application
requirement. Generated protocols require human review.
## Application Responsibilities
The first GUI milestone provides:
- audio-file selection
- structured Meeting Context and participant editing
- optional explicit speaker mapping
- pipeline start and configuration
- progress and failure presentation
- detailed protocol display and editing
- export of the reviewed result
The product direction also includes a shorter participant/distribution
protocol. It may be delivered after the first GUI milestone.
## Artifact and Storage Principles
The existing Meeting domain remains the application-level container for source
recordings, canonical and working transcripts, context, generated artifacts and
exports. Storage remains independent of processing segmentation. In particular,
chunks are not primary Meeting Assistant domain objects or the required unit of
MVP persistence.
Source audio and the canonical transcription output are preserved. Corrections
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
- integrated recording
- searchable meeting history and optional database indexes
- live transcription
- OCR and video processing
- semantic search and organizational knowledge features
+17 -1
View File
@@ -2,7 +2,23 @@ Transcript
Original speech-to-text output.
Diarization
Assignment of transcript segments to speakers.
Assignment of transcript segments to anonymous speaker labels. It does not by
itself identify participants.
Meeting Context
Structured meeting metadata, participants and optional explicitly confirmed
speaker-to-participant mappings supplied to the processing pipeline.
Speaker Mapping
An explicitly confirmed association from a diarization label such as
`SPEAKER_00` to a participant ID. It must not be inferred automatically.
Detailed Protocol
Contextual protocol generated from the full transcript for human review and
continued work.
Distribution Protocol
Shorter, participant-facing version of a reviewed meeting protocol.
Knowledge Object
Structured information extracted from meetings.
+39
View File
@@ -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
View File
@@ -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
View File
@@ -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"
+17
View File
@@ -0,0 +1,17 @@
"""Application services for Meeting Assistant workflows."""
from mka.application.meeting_service import (
MeetingDetails,
MeetingProcessingService,
ParticipantInput,
ProcessingOptions,
ProcessingOutcome,
)
__all__ = [
"MeetingDetails",
"MeetingProcessingService",
"ParticipantInput",
"ProcessingOptions",
"ProcessingOutcome",
]
+106
View File
@@ -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)
+369
View File
@@ -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
+122
View File
@@ -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
+561
View File
@@ -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
+108
View File
@@ -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
+23
View File
@@ -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
+80
View File
@@ -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,
)
+218
View File
@@ -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
+1
View File
@@ -0,0 +1 @@
"""Adapters for external processing systems."""
+68
View File
@@ -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,
)
+1 -1
View File
@@ -22,4 +22,4 @@ class DomainModel(BaseModel):
id: UUID = Field(default_factory=uuid4)
created_at: datetime = Field(default_factory=utc_now)
updated_at: datetime = Field(default_factory=utc_now)
updated_at: datetime = Field(default_factory=utc_now)
+785
View File
@@ -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()
+138
View File
@@ -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
+76
View File
@@ -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()
+78
View File
@@ -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
+147
View File
@@ -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]
+88
View File
@@ -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)
+826
View File
@@ -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
+127
View File
@@ -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
+21
View File
@@ -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
+45
View File
@@ -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
+136
View File
@@ -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"
)
+134
View File
@@ -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
+95
View File
@@ -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