11 Commits
21 changed files with 2155 additions and 113 deletions
+10
View File
@@ -4,6 +4,14 @@
### 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
@@ -15,6 +23,8 @@
### 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
+46 -1
View File
@@ -1,5 +1,14 @@
# Project Knowledge
## Meeting language
The GUI's `meeting_language` is passed through `MeetingDetails.language` to both
Meeting Context `meeting.language` and `MvpMeetingConfig.language` (Whisper).
Meeting Lab derives explicit protocol output language from the persisted context,
also during regeneration, and records derived `output_language` provenance in
protocol runtime metadata. Missing context language defaults to German. There is
no separate protocol-language setting or automatic context/transcript translation.
## Vision
Meeting Assistant turns recorded meetings into reviewed, user-facing protocols
@@ -48,6 +57,8 @@ Meeting Assistant owns user interaction and product workflow:
- 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.
@@ -70,7 +81,8 @@ Imported audio
(optionally normalize loudness; default on)
-> transcribe with whisper.cpp and large-v3-turbo
-> optionally diarize with pyannote.audio Community-1
-> generate a protocol directly from the full transcript
-> when diarized, wait for optional speaker review before explicit protocol generation
-> otherwise generate a protocol directly from the full transcript
-> review and edit by a human
```
@@ -84,6 +96,14 @@ 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,
@@ -98,6 +118,15 @@ 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:
@@ -139,3 +168,19 @@ product work, not a current MVP requirement.
- 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.
+75 -1
View File
@@ -23,12 +23,20 @@ reference recorder, but imported audio is not tied to OBS-specific behavior.
## Processing Pipeline
The existing **Meeting language** selector controls both transcription and protocol
output: `de` means German for both, and `en` means English for both. The selection
is saved in Meeting Context (`meeting.language`) and reused for protocol-only
regeneration without retranscription. Older contexts without a language retain
German protocol output. Names, speaker mappings and authored Meeting Context are
preserved; no transcript translation stage is added.
```text
Audio file
-> FFmpeg preparation (mono, 16 kHz PCM WAV; normalization optional)
-> whisper.cpp transcription (large-v3-turbo)
-> optional pyannote.audio Community-1 diarization
-> direct full-transcript protocol generation
-> optional speaker review and explicit protocol generation (when diarized)
or direct full-transcript protocol generation (without diarization)
-> human review and export
```
@@ -42,6 +50,11 @@ compatibility path. Diarization is optional and produces anonymous speaker
labels. A label identifies a participant only when the user explicitly
confirms the mapping; automatic speaker-name inference is not allowed.
Speaker selectors hide participants already assigned to other speakers while keeping
the current assignment available. The UI shows detected, assigned and unassigned
counts, marks unassigned speakers, and warns when mappings remain incomplete.
Anonymous-speaker protocol generation remains available.
After a diarized run, the result view lists detected `SPEAKER_XX` labels with
short transcript excerpts. Confirmed mappings regenerate only the protocol
from the existing diarized transcript; audio preparation, Whisper and Pyannote
@@ -70,6 +83,14 @@ Audio normalization is enabled by default and can be disabled in the processing
options. This controls loudness normalization only: Meeting Lab still prepares
every WAV, FLAC or M4A source as canonical audio before transcription.
The **Performance profile** selector controls protocol-generation runtime using
the abstract Auto (default), Fast, Efficient, and Powersave profiles. Auto lets
the inference backend select its own thread configuration; for Ollama, Meeting
Assistant intentionally sends no `num_thread` option. Fast, Efficient, and
Powersave are explicit resource profiles currently mapped to 16, 10, and 4
Ollama CPU threads. These concrete mappings may evolve independently of the UI
semantics.
The People section can export its current entries to a UTF-8 `people.yaml` file
and replace them from a previous `.yaml` or `.yml` export. Stable person IDs,
names, roles, organizations and attendance states are retained. This is a small
@@ -131,6 +152,27 @@ The database defaults to `data/database/glossary.sqlite3`. It is created and
bootstrapped automatically and can be moved with `MKA_GLOSSARY_DATABASE`.
Glossary integration does not rewrite raw or diarized transcript artifacts.
Use **Export glossary** to download the complete SQLite-backed glossary as
UTF-8 YAML, and **Import glossary** to replace it from a previously exported
file. Imports are fully validated before a single SQLite transaction replaces
the current glossary, so malformed or conflicting files leave existing data
unchanged. Schema version 1 is:
```yaml
version: 1
glossary:
- id: 1
canonical_term: Secugrid HS
aliases:
- Sikirgut
category: product
description: Canonical product spelling
active: true
```
`id` is the stable SQLite glossary-entry identifier. Categories are `product`,
`material`, `organization`, `technical_term`, `acronym`, or `other`.
## Input configuration import and export
Use **Export inputs** to save the current Meeting Assistant run-input form as a
@@ -171,6 +213,22 @@ From the Meeting Assistant repository, start the UI with:
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
@@ -188,3 +246,19 @@ 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.
+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.
+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.
+1 -1
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"
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"
+75
View File
@@ -35,6 +35,18 @@ class GlossaryEntry:
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."""
@@ -194,6 +206,64 @@ class GlossaryRepository:
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)
@@ -285,6 +355,11 @@ def render_glossary_terms(entries: Iterable[GlossaryEntry]) -> list[str]:
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")
+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
+314 -27
View File
@@ -15,9 +15,46 @@ from uuid import uuid4
import yaml
from mka.application.config import AppSettings
from mka.application.glossary import GlossaryRepository, render_glossary_terms
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")
_EXCERPT_TOKEN_PATTERN = re.compile(r"[^\W\d_]+|\d+", re.UNICODE)
_EXCERPT_STOP_WORDS = frozenset(
{
"a", "an", "and", "are", "das", "der", "die", "ein", "eine", "for", "i",
"ich", "im", "in", "is", "it", "ja", "mit", "of", "or", "the", "to", "und",
"we", "wir", "with",
}
)
_ACKNOWLEDGEMENT_WORDS = frozenset(
{
"absolutely", "alles", "danke", "genau", "gut", "ja", "klar", "okay", "ok",
"perfekt", "right", "sure", "thanks", "verstanden", "yes",
}
)
_FIRST_PERSON_WORDS = frozenset(
{"i", "ich", "me", "mein", "meine", "my", "unser", "unsere", "we", "wir"}
)
_ACTION_WORDS = frozenset(
{
"agree", "approve", "believe", "can", "decide", "habe", "kann", "möchte", "plan",
"recommend", "think", "übernehme", "werde", "will",
}
)
_DATE_WORDS = frozenset(
{
"april", "august", "december", "dienstag", "donnerstag", "february", "freitag",
"friday", "january", "juli", "june", "mai", "march", "monday", "mittwoch", "montag",
"november", "october", "saturday", "september", "sonntag", "sunday", "thursday",
"tuesday", "wednesday",
}
)
class MeetingLabPort(Protocol):
@@ -65,6 +102,7 @@ class ParticipantInput:
class ProcessingOptions:
diarization_enabled: bool = False
audio_normalization: bool = True
performance_profile: str = DEFAULT_PERFORMANCE_PROFILE
@dataclass(frozen=True)
@@ -85,6 +123,8 @@ class ProcessingOutcome:
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)
@@ -100,6 +140,89 @@ class SpeakerMappingReview:
current_mappings: dict[str, str]
@dataclass(frozen=True)
class _ExcerptCandidate:
speaker_label: str
text: str
segment_index: int
tokens: frozenset[str]
lexical_tokens: frozenset[str]
previous_other_speaker_tokens: frozenset[str]
def _excerpt_tokens(text: str) -> tuple[str, ...]:
return tuple(token.lower() for token in _EXCERPT_TOKEN_PATTERN.findall(text))
def _excerpt_score(candidate: _ExcerptCandidate, distinctive_token_count: int) -> float:
"""Score an excerpt using deterministic lexical signals only."""
word_count = len(candidate.tokens)
lexical_count = len(candidate.lexical_tokens)
score = min(word_count, 24) * 0.3 + min(lexical_count, 16) * 0.6
if word_count < 5:
score -= 8.0
elif word_count < 9:
score -= 3.0
if lexical_count <= 2:
score -= 4.0
if candidate.tokens and candidate.tokens <= _ACKNOWLEDGEMENT_WORDS:
score -= 14.0
if candidate.text.rstrip().endswith("?"):
score -= 5.0
if any(token.isdigit() for token in candidate.tokens):
score += 2.5
if candidate.tokens & _DATE_WORDS:
score += 2.0
if candidate.tokens & _FIRST_PERSON_WORDS:
score += 1.5
if candidate.tokens & _ACTION_WORDS:
score += 2.0
proper_name_count = len(re.findall(r"\b[A-ZÄÖÜ][a-zäöüß]{2,}\b", candidate.text))
score += min(proper_name_count, 3) * 0.75
score += min(sum(len(token) >= 10 for token in candidate.lexical_tokens), 4) * 0.4
score += min(distinctive_token_count, 6) * 0.65
if candidate.previous_other_speaker_tokens and candidate.lexical_tokens:
overlap = len(candidate.lexical_tokens & candidate.previous_other_speaker_tokens)
if overlap / len(candidate.lexical_tokens) >= 0.7:
score -= 5.0
return score
def _segment_bucket(segment_index: int, segment_count: int) -> int:
return min(5, segment_index * 6 // max(segment_count, 1))
def _selection_score(
candidate: _ExcerptCandidate,
base_score: float,
selected: list[_ExcerptCandidate],
selected_buckets: set[int],
segment_count: int,
) -> float:
"""Favor excerpts from new meeting positions and avoid near duplicates."""
score = base_score
bucket = _segment_bucket(candidate.segment_index, segment_count)
if bucket not in selected_buckets:
score += 2.5
if not selected or not candidate.lexical_tokens:
return score
similarities = [
len(candidate.lexical_tokens & item.lexical_tokens)
/ len(candidate.lexical_tokens | item.lexical_tokens)
for item in selected
if item.lexical_tokens
]
maximum_similarity = max(similarities, default=0.0)
if maximum_similarity >= 0.72:
score -= 18.0
elif maximum_similarity >= 0.55:
score -= 6.0
nearest_position = min(abs(candidate.segment_index - item.segment_index) for item in selected)
if nearest_position / max(segment_count - 1, 1) < 0.12:
score -= 1.5
return score
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("-")
@@ -114,11 +237,13 @@ class MeetingProcessingService:
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,
@@ -232,20 +357,25 @@ class MeetingProcessingService:
"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,
@@ -256,7 +386,13 @@ class MeetingProcessingService:
if progress_sink is not None:
progress_sink(app_event)
result = self.meeting_lab.run(config, context, relay)
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")
@@ -280,6 +416,15 @@ class MeetingProcessingService:
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,
@@ -293,7 +438,7 @@ class MeetingProcessingService:
self,
run_dir: Path,
*,
excerpts_per_speaker: int = 3,
excerpts_per_speaker: int = 6,
minimum_excerpt_characters: int = 20,
) -> SpeakerMappingReview | None:
"""Load detected labels and small deterministic identification excerpts."""
@@ -310,34 +455,67 @@ class MeetingProcessingService:
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:
candidates_by_speaker: dict[str, list[_ExcerptCandidate]] = {}
short_candidates_by_speaker: dict[str, list[_ExcerptCandidate]] = {}
previous_label: str | None = None
previous_tokens: frozenset[str] = frozenset()
for segment_index, segment in enumerate(segments):
if not isinstance(segment, dict):
continue
label = segment.get("speaker_id")
text = segment.get("text")
normalized = " ".join(text.split()) if isinstance(text, str) else ""
tokens = frozenset(_excerpt_tokens(normalized))
if (
not isinstance(label, str)
or not label.startswith("SPEAKER_")
or label == "SPEAKER_UNASSIGNED"
or not isinstance(text, str)
or not text.strip()
or not normalized
):
if normalized:
previous_label = label if isinstance(label, str) else None
previous_tokens = tokens
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))
candidate = _ExcerptCandidate(
speaker_label=label,
text=normalized,
segment_index=segment_index,
tokens=tokens,
lexical_tokens=tokens - _EXCERPT_STOP_WORDS,
previous_other_speaker_tokens=(
previous_tokens
if previous_label is not None and previous_label != label
else frozenset()
),
)
target = (
candidates_by_speaker
if len(normalized) >= minimum_excerpt_characters
else short_candidates_by_speaker
)
target.setdefault(label, []).append(candidate)
previous_label = label
previous_tokens = tokens
labels = sorted(set(candidates_by_speaker) | set(short_candidates_by_speaker))
token_speakers: dict[str, set[str]] = {}
for candidate_group in (*candidates_by_speaker.values(), *short_candidates_by_speaker.values()):
for candidate in candidate_group:
for token in candidate.lexical_tokens:
token_speakers.setdefault(token, set()).add(candidate.speaker_label)
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])))
candidates = candidates_by_speaker.get(label) or short_candidates_by_speaker.get(
label, []
)
excerpts = self._select_identifying_excerpts(
candidates,
excerpts_per_speaker=excerpts_per_speaker,
segment_count=len(segments),
token_speakers=token_speakers,
)
speakers.append(SpeakerReview(label, excerpts))
participants = tuple(
(participant["participant_id"], participant["display_name"])
for participant in context_data.get("participants", [])
@@ -352,11 +530,56 @@ class MeetingProcessingService:
current_mappings=dict(mappings) if isinstance(mappings, dict) else {},
)
@staticmethod
def _select_identifying_excerpts(
candidates: list[_ExcerptCandidate],
*,
excerpts_per_speaker: int,
segment_count: int,
token_speakers: dict[str, set[str]],
) -> tuple[str, ...]:
"""Choose informative, distinct excerpts while retaining transcript order."""
if len(candidates) <= excerpts_per_speaker:
return tuple(candidate.text for candidate in candidates)
base_scores = {
candidate: _excerpt_score(
candidate,
sum(
len(token_speakers.get(token, set())) == 1
for token in candidate.lexical_tokens
),
)
for candidate in candidates
}
selected: list[_ExcerptCandidate] = []
selected_buckets: set[int] = set()
while len(selected) < excerpts_per_speaker:
choice = max(
(candidate for candidate in candidates if candidate not in selected),
key=lambda candidate: (
_selection_score(
candidate,
base_scores[candidate],
selected,
selected_buckets,
segment_count,
),
-candidate.segment_index,
),
)
selected.append(choice)
selected_buckets.add(_segment_bucket(choice.segment_index, segment_count))
return tuple(
candidate.text for candidate in sorted(selected, key=lambda item: item.segment_index)
)
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)
@@ -381,8 +604,16 @@ class MeetingProcessingService:
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(
@@ -394,15 +625,24 @@ class MeetingProcessingService:
)
)
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),
)
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,
@@ -444,6 +684,18 @@ class MeetingProcessingService:
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:
@@ -463,3 +715,38 @@ class MeetingProcessingService:
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
+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,
)
+304 -81
View File
@@ -2,7 +2,11 @@
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
@@ -14,6 +18,11 @@ from mka.application.glossary import (
GlossaryConflictError,
GlossaryRepository,
)
from mka.application.glossary_yaml import (
GlossaryYamlError,
export_glossary_yaml,
import_glossary_yaml,
)
from mka.application.meeting_service import (
STAGES,
AppProgressEvent,
@@ -21,6 +30,7 @@ from mka.application.meeting_service import (
MeetingProcessingService,
ParticipantInput,
ProcessingOptions,
SpeakerMappingReview,
stable_id,
)
from mka.application.people_yaml import (
@@ -28,6 +38,8 @@ from mka.application.people_yaml import (
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,
@@ -62,6 +74,33 @@ def _render_glossary(repository: GlossaryRepository) -> None:
"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")
@@ -304,40 +343,199 @@ def _render_participants() -> list[ParticipantInput]:
return people
def _progress_callback(
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,
progress_slot: Any,
states: dict[str, str],
) -> Any:
progress_bar = None
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 update(event: AppProgressEvent) -> None:
nonlocal progress_bar
if event.stage in states:
states[event.stage] = "running" if event.status == "started" else event.status
if event.stage == "failed":
running = next(
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:
states[running] = "failed"
elapsed = f"{event.elapsed_seconds:.1f} s"
message = event.message or STAGE_LABELS.get(event.stage, event.stage)
status_box.info(f"{message} — elapsed {elapsed}")
stage_table.table(
[{"Stage": STAGE_LABELS[stage], "Status": states[stage]} for stage in STAGES]
)
if event.progress is not None:
if progress_bar is None:
progress_bar = progress_slot.progress(0.0)
progress_bar.progress(
min(max(event.progress, 0.0), 1.0),
text=f"{STAGE_LABELS.get(event.stage, event.stage)}: {event.progress:.0%}",
)
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
return update
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:
@@ -346,6 +544,9 @@ def _render_result() -> 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:
@@ -353,7 +554,15 @@ def _render_result() -> None:
st.caption("Intermediate artifacts and run metadata were preserved here.")
return
st.success("Processing completed. Review the generated protocol before use.")
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(
@@ -362,17 +571,18 @@ def _render_result() -> None:
)
if message := st.session_state.pop("speaker_mapping_message", None):
st.success(message)
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}")
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())
@@ -383,48 +593,47 @@ def _render_result() -> None:
if review is None or not review.speakers:
return
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)
options = [None, *participant_names]
selections: dict[str, str] = {}
for speaker in review.speakers:
current = review.current_mappings.get(speaker.speaker_label)
selected = st.selectbox(
speaker.speaker_label,
options=options,
index=options.index(current) if current in options else 0,
format_func=lambda value, names=participant_names: (
"Unmapped / Unknown" if value is None else names[value]
),
key=f"speaker_mapping_{outcome.run_dir.name}_{speaker.speaker_label}",
)
if selected is not None:
selections[speaker.speaker_label] = selected
for excerpt in speaker.excerpts:
st.caption(f"“{excerpt}”")
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(
"Regenerate protocol with confirmed speakers",
"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:
with st.spinner("Regenerating protocol without rerunning audio processing..."):
regenerated = service.regenerate_protocol(outcome.run_dir, selections)
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:
st.error(f"Protocol regeneration failed: {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 regenerated."
"Speaker mappings saved and protocol generated."
)
st.rerun()
@@ -458,7 +667,10 @@ def main() -> None:
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"
"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"
@@ -470,6 +682,13 @@ def main() -> 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,
@@ -523,30 +742,34 @@ def main() -> None:
)
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")
status_box = st.empty()
stage_table = st.empty()
progress_slot = st.empty()
states = {stage: "pending" for stage in STAGES}
if not diarization_enabled:
states["diarization"] = "skipped"
callback = _progress_callback(status_box, stage_table, progress_slot, states)
outcome = service.process(
audio_path,
meeting,
participants,
ProcessingOptions(
diarization_enabled=diarization_enabled,
audio_normalization=audio_normalization,
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,
),
progress_sink=callback,
states,
"Starting processing",
)
st.session_state.outcome = outcome
st.session_state.edited_protocol = outcome.original_protocol or ""
if outcome.succeeded:
status_box.success("Processing completed.")
st.success("Processing completed.")
else:
status_box.error(
st.error(
f"Processing failed during {outcome.failed_stage}: "
f"{outcome.error_message} Artifacts were preserved."
)
+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
+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)
+369 -2
View File
@@ -6,6 +6,7 @@ from types import SimpleNamespace
from typing import Any
import pytest
import yaml
from mka.application.config import AppSettings
from mka.application.meeting_service import (
@@ -14,6 +15,7 @@ from mka.application.meeting_service import (
ParticipantInput,
ProcessingOptions,
)
from mka.application.progress_timing import ProcessingTimer
@dataclass
@@ -98,6 +100,24 @@ class FakeMeetingLab:
)
)
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)
@@ -227,6 +247,18 @@ def write_speaker_review_artifacts(run_dir: Path) -> Path:
return transcript_path
def write_custom_speaker_review_artifacts(
run_dir: Path, segments: list[dict[str, str | None]]
) -> Path:
"""Write custom diarized segments using the standard review context fixture."""
transcript_path = write_speaker_review_artifacts(run_dir)
transcript_path.write_text(
json.dumps({"speaker_labels_anonymous": True, "segments": segments}),
encoding="utf-8",
)
return transcript_path
def test_build_context_uses_actual_v1_shape(tmp_path: Path) -> None:
service, gateway = make_service(tmp_path)
@@ -323,19 +355,28 @@ def test_build_context_translates_mentioned_only_person(tmp_path: Path) -> 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, meeting(), participants(), ProcessingOptions(diarization_enabled=False)
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
@@ -377,7 +418,7 @@ def test_process_propagates_enabled_diarization_and_progress(tmp_path: Path) ->
audio.write_bytes(b"audio")
events = []
service.process(
outcome = service.process(
audio,
meeting(),
participants(),
@@ -387,6 +428,9 @@ def test_process_propagates_enabled_diarization_and_progress(tmp_path: Path) ->
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"),
@@ -394,6 +438,33 @@ def test_process_propagates_enabled_diarization_and_progress(tmp_path: Path) ->
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(
@@ -475,6 +546,89 @@ def test_speaker_review_lists_detected_labels_participants_and_excerpts(
assert "SPEAKER_00" in source.read_text(encoding="utf-8")
def test_speaker_review_prefers_content_rich_excerpts_over_acknowledgements_and_questions(
tmp_path: Path,
) -> None:
service, gateway = make_service(tmp_path)
acknowledgement = "Okay, that sounds good."
question = "Could we continue now?"
rich_excerpts = [
"I will complete the Atlas deployment in Berlin on 14 September with the database team.",
"I recommend moving the Orion product launch because the supplier test is incomplete.",
"We decided that I will own the Kubernetes migration for the payments service.",
"My experience with the Frankfurt factory rollout suggests a two-week validation period.",
"I will present the security audit results to the Apollo project steering group.",
"We should reserve twelve hours for the API integration and production monitoring.",
]
write_custom_speaker_review_artifacts(
gateway.run_dir,
[{"speaker_id": "SPEAKER_00", "text": text} for text in [acknowledgement, question, *rich_excerpts]],
)
review = service.load_speaker_mapping_review(gateway.run_dir)
assert review is not None
excerpts = review.speakers[0].excerpts
assert len(excerpts) == 6
assert acknowledgement not in excerpts
assert question not in excerpts
assert set(excerpts) == set(rich_excerpts)
def test_speaker_review_avoids_near_duplicate_excerpts(tmp_path: Path) -> None:
service, gateway = make_service(tmp_path)
primary = "I will lead the Atlas API migration for the Berlin platform team next Monday."
duplicate = "I will lead the Atlas API migration for the Berlin platform group next Monday."
alternatives = [
"I recommend documenting the supplier quality decision in the project register.",
"My team will complete the production calibration before the factory trial.",
"I will present the customer research findings at the Munich planning workshop.",
"Our database monitoring threshold needs approval from the operations group.",
"I will coordinate the Orion release checklist with the support organization.",
]
write_custom_speaker_review_artifacts(
gateway.run_dir,
[{"speaker_id": "SPEAKER_00", "text": text} for text in [primary, duplicate, *alternatives]],
)
review = service.load_speaker_mapping_review(gateway.run_dir)
assert review is not None
excerpts = review.speakers[0].excerpts
assert primary in excerpts
assert duplicate not in excerpts
assert set(excerpts) == {primary, *alternatives}
def test_speaker_review_spreads_selected_excerpts_across_meeting_positions(
tmp_path: Path,
) -> None:
service, gateway = make_service(tmp_path)
positions = {
0: "I will lead the Atlas deployment review with the engineering team marker alpha.",
1: "I will lead the Atlas deployment review with the engineering team marker bravo.",
6: "I will lead the Atlas deployment review with the engineering team marker charlie.",
7: "I will lead the Atlas deployment review with the engineering team marker delta.",
12: "I will lead the Atlas deployment review with the engineering team marker echo.",
13: "I will lead the Atlas deployment review with the engineering team marker foxtrot.",
18: "I will lead the Atlas deployment review with the engineering team marker golf.",
19: "I will lead the Atlas deployment review with the engineering team marker hotel.",
}
segments = [
{
"speaker_id": "SPEAKER_00" if index in positions else "SPEAKER_01",
"text": positions.get(index, "Acknowledged."),
}
for index in range(20)
]
write_custom_speaker_review_artifacts(gateway.run_dir, segments)
review = service.load_speaker_mapping_review(gateway.run_dir)
assert review is not None
assert review.speakers[0].excerpts == tuple(positions[index] for index in (0, 1, 6, 7, 12, 18))
def test_protocol_only_regeneration_persists_mappings_without_rewriting_transcript(
tmp_path: Path,
) -> None:
@@ -503,6 +657,7 @@ def test_protocol_only_regeneration_persists_mappings_without_rewriting_transcri
"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
@@ -526,6 +681,49 @@ def test_protocol_regeneration_allows_unmapped_and_rejects_duplicate_participant
)
@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"
@@ -552,3 +750,172 @@ def test_uploaded_source_is_preserved_in_meeting_directory(tmp_path: Path) -> No
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
+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
+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
+27
View File
@@ -1,6 +1,7 @@
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
@@ -66,3 +67,29 @@ def test_imported_inputs_are_applied_via_pending_state_before_widgets(
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