Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bd48aec5aa | ||
|
|
b50b87c32b | ||
|
|
9e7173c2a7 | ||
|
|
70cb1e0789 | ||
|
|
500d2c4c67 | ||
|
|
77be2b39ff | ||
|
|
e2a0434941 | ||
|
|
9840953687 | ||
|
|
e55dd74c19 |
@@ -4,6 +4,14 @@
|
|||||||
|
|
||||||
### Added
|
### 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
|
- First Streamlit MVP for audio upload, Meeting Context entry, participant
|
||||||
management, processing progress and protocol editing.
|
management, processing progress and protocol editing.
|
||||||
- Application service and Meeting Lab adapter with environment-based runtime
|
- Application service and Meeting Lab adapter with environment-based runtime
|
||||||
@@ -15,6 +23,8 @@
|
|||||||
|
|
||||||
### Changed
|
### Changed
|
||||||
|
|
||||||
|
- Record active glossary configuration without rewriting protocol transcript input.
|
||||||
|
|
||||||
- Reconciled the documented MVP with the validated Meeting Lab backend.
|
- Reconciled the documented MVP with the validated Meeting Lab backend.
|
||||||
- Selected `whisper.cpp` for transcription and made diarization optional.
|
- Selected `whisper.cpp` for transcription and made diarization optional.
|
||||||
- Defined the Python API, progress-event and Meeting Context integration
|
- Defined the Python API, progress-event and Meeting Context integration
|
||||||
|
|||||||
@@ -1,5 +1,14 @@
|
|||||||
# Project Knowledge
|
# 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
|
## Vision
|
||||||
|
|
||||||
Meeting Assistant turns recorded meetings into reviewed, user-facing protocols
|
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
|
- a global SQLite terminology glossary whose active canonical core terms and
|
||||||
recognition aliases are rendered through Meeting Context into direct
|
recognition aliases are rendered through Meeting Context into direct
|
||||||
protocol prompts
|
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
|
- export and presentation of protocol versions
|
||||||
|
|
||||||
Processing logic must not be duplicated in the application.
|
Processing logic must not be duplicated in the application.
|
||||||
@@ -84,6 +95,14 @@ PyTorch/ROCm processed the same duration in approximately 98.6 seconds. These
|
|||||||
are validation observations, not performance guarantees or hardware
|
are validation observations, not performance guarantees or hardware
|
||||||
requirements. CPU execution remains supported and may be substantially slower.
|
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
|
## Meeting Context and Speakers
|
||||||
|
|
||||||
`MeetingContext` is a structured domain object containing meeting metadata,
|
`MeetingContext` is a structured domain object containing meeting metadata,
|
||||||
@@ -98,6 +117,12 @@ Speakers otherwise remain anonymous. Automatic speaker-name inference is not
|
|||||||
allowed. The GUI must create and edit Meeting Context; hand-written YAML is not
|
allowed. The GUI must create and edit Meeting Context; hand-written YAML is not
|
||||||
a product requirement.
|
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 regeneration or change mapping persistence.
|
||||||
|
|
||||||
## Progress Contract
|
## Progress Contract
|
||||||
|
|
||||||
Meeting Lab emits stage-based progress events for:
|
Meeting Lab emits stage-based progress events for:
|
||||||
@@ -139,3 +164,19 @@ product work, not a current MVP requirement.
|
|||||||
- no automatic identity claims
|
- no automatic identity claims
|
||||||
- human review of generated protocols
|
- human review of generated protocols
|
||||||
- small modules and simple interfaces
|
- 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.
|
||||||
|
|||||||
@@ -23,6 +23,13 @@ reference recorder, but imported audio is not tied to OBS-specific behavior.
|
|||||||
|
|
||||||
## Processing Pipeline
|
## 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
|
```text
|
||||||
Audio file
|
Audio file
|
||||||
-> FFmpeg preparation (mono, 16 kHz PCM WAV; normalization optional)
|
-> FFmpeg preparation (mono, 16 kHz PCM WAV; normalization optional)
|
||||||
@@ -42,6 +49,11 @@ compatibility path. Diarization is optional and produces anonymous speaker
|
|||||||
labels. A label identifies a participant only when the user explicitly
|
labels. A label identifies a participant only when the user explicitly
|
||||||
confirms the mapping; automatic speaker-name inference is not allowed.
|
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
|
After a diarized run, the result view lists detected `SPEAKER_XX` labels with
|
||||||
short transcript excerpts. Confirmed mappings regenerate only the protocol
|
short transcript excerpts. Confirmed mappings regenerate only the protocol
|
||||||
from the existing diarized transcript; audio preparation, Whisper and Pyannote
|
from the existing diarized transcript; audio preparation, Whisper and Pyannote
|
||||||
@@ -70,6 +82,14 @@ Audio normalization is enabled by default and can be disabled in the processing
|
|||||||
options. This controls loudness normalization only: Meeting Lab still prepares
|
options. This controls loudness normalization only: Meeting Lab still prepares
|
||||||
every WAV, FLAC or M4A source as canonical audio before transcription.
|
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
|
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,
|
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
|
names, roles, organizations and attendance states are retained. This is a small
|
||||||
@@ -131,6 +151,27 @@ The database defaults to `data/database/glossary.sqlite3`. It is created and
|
|||||||
bootstrapped automatically and can be moved with `MKA_GLOSSARY_DATABASE`.
|
bootstrapped automatically and can be moved with `MKA_GLOSSARY_DATABASE`.
|
||||||
Glossary integration does not rewrite raw or diarized transcript artifacts.
|
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
|
## Input configuration import and export
|
||||||
|
|
||||||
Use **Export inputs** to save the current Meeting Assistant run-input form as a
|
Use **Export inputs** to save the current Meeting Assistant run-input form as a
|
||||||
@@ -188,3 +229,19 @@ transcription; Meeting Assistant does not duplicate audio conversion.
|
|||||||
|
|
||||||
See [Architecture](docs/architecture.md), [Project Knowledge](PROJECT_KNOWLEDGE.md),
|
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).
|
[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.
|
||||||
|
|||||||
@@ -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.
|
||||||
@@ -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
@@ -4,7 +4,7 @@ build-backend = "hatchling.build"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "meeting-assistant"
|
name = "meeting-assistant"
|
||||||
version = "0.1.0"
|
version = "0.1.0a1"
|
||||||
description = "Local-first meeting transcription and knowledge capture assistant"
|
description = "Local-first meeting transcription and knowledge capture assistant"
|
||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
|
|||||||
@@ -35,6 +35,18 @@ class GlossaryEntry:
|
|||||||
updated_at: str
|
updated_at: str
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class GlossaryReplacementEntry:
|
||||||
|
"""Validated values used to atomically replace the persisted glossary."""
|
||||||
|
|
||||||
|
id: int
|
||||||
|
canonical_term: str
|
||||||
|
category: str
|
||||||
|
description: str | None
|
||||||
|
is_active: bool
|
||||||
|
aliases: tuple[str, ...]
|
||||||
|
|
||||||
|
|
||||||
class GlossaryRepository:
|
class GlossaryRepository:
|
||||||
"""Small data-access boundary for the local glossary database."""
|
"""Small data-access boundary for the local glossary database."""
|
||||||
|
|
||||||
@@ -194,6 +206,64 @@ class GlossaryRepository:
|
|||||||
if cursor.rowcount == 0:
|
if cursor.rowcount == 0:
|
||||||
raise KeyError(f"Unknown glossary entry: {entry_id}")
|
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
|
@contextmanager
|
||||||
def _connect(self) -> Iterator[sqlite3.Connection]:
|
def _connect(self) -> Iterator[sqlite3.Connection]:
|
||||||
connection = sqlite3.connect(self.database_path, timeout=5)
|
connection = sqlite3.connect(self.database_path, timeout=5)
|
||||||
@@ -285,6 +355,11 @@ def render_glossary_terms(entries: Iterable[GlossaryEntry]) -> list[str]:
|
|||||||
return rendered
|
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:
|
def _timestamp() -> str:
|
||||||
return datetime.now(UTC).isoformat(timespec="seconds")
|
return datetime.now(UTC).isoformat(timespec="seconds")
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -15,7 +15,13 @@ from uuid import uuid4
|
|||||||
import yaml
|
import yaml
|
||||||
|
|
||||||
from mka.application.config import AppSettings
|
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")
|
STAGES = ("preparing", "transcription", "diarization", "protocol_generation")
|
||||||
|
|
||||||
@@ -65,6 +71,7 @@ class ParticipantInput:
|
|||||||
class ProcessingOptions:
|
class ProcessingOptions:
|
||||||
diarization_enabled: bool = False
|
diarization_enabled: bool = False
|
||||||
audio_normalization: bool = True
|
audio_normalization: bool = True
|
||||||
|
performance_profile: str = DEFAULT_PERFORMANCE_PROFILE
|
||||||
|
|
||||||
|
|
||||||
@dataclass(frozen=True)
|
@dataclass(frozen=True)
|
||||||
@@ -114,11 +121,13 @@ class MeetingProcessingService:
|
|||||||
settings: AppSettings,
|
settings: AppSettings,
|
||||||
meeting_lab: MeetingLabPort,
|
meeting_lab: MeetingLabPort,
|
||||||
glossary: GlossaryRepository | None = None,
|
glossary: GlossaryRepository | None = None,
|
||||||
|
timer: ProcessingTimer | None = None,
|
||||||
) -> None:
|
) -> None:
|
||||||
self.settings = settings
|
self.settings = settings
|
||||||
self.meeting_lab = meeting_lab
|
self.meeting_lab = meeting_lab
|
||||||
self.glossary = glossary or GlossaryRepository(settings.glossary_database)
|
self.glossary = glossary or GlossaryRepository(settings.glossary_database)
|
||||||
self.glossary.initialize()
|
self.glossary.initialize()
|
||||||
|
self.timer = timer or ProcessingTimer()
|
||||||
|
|
||||||
def build_context(
|
def build_context(
|
||||||
self,
|
self,
|
||||||
@@ -232,6 +241,8 @@ class MeetingProcessingService:
|
|||||||
"protocol_safe_input_token_budget": (
|
"protocol_safe_input_token_budget": (
|
||||||
self.settings.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": (
|
"diarization": (
|
||||||
self.settings.diarization_mode if options.diarization_enabled else "off"
|
self.settings.diarization_mode if options.diarization_enabled else "off"
|
||||||
),
|
),
|
||||||
@@ -241,11 +252,13 @@ class MeetingProcessingService:
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
current_stage: str | None = None
|
current_stage: str | None = None
|
||||||
|
self.timer.start_run()
|
||||||
|
|
||||||
def relay(event: Any) -> None:
|
def relay(event: Any) -> None:
|
||||||
nonlocal current_stage
|
nonlocal current_stage
|
||||||
if event.stage in STAGES and event.status == "started":
|
if event.stage in STAGES and event.status == "started":
|
||||||
current_stage = event.stage
|
current_stage = event.stage
|
||||||
|
self._update_timer(event)
|
||||||
app_event = AppProgressEvent(
|
app_event = AppProgressEvent(
|
||||||
stage=event.stage,
|
stage=event.stage,
|
||||||
status=event.status,
|
status=event.status,
|
||||||
@@ -256,7 +269,13 @@ class MeetingProcessingService:
|
|||||||
if progress_sink is not None:
|
if progress_sink is not None:
|
||||||
progress_sink(app_event)
|
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:
|
if result.exit_code != 0:
|
||||||
failure = self._read_failure(result.run_dir)
|
failure = self._read_failure(result.run_dir)
|
||||||
backend_stage = failure.get("stage")
|
backend_stage = failure.get("stage")
|
||||||
@@ -357,6 +376,7 @@ class MeetingProcessingService:
|
|||||||
run_dir: Path,
|
run_dir: Path,
|
||||||
speaker_mappings: dict[str, str],
|
speaker_mappings: dict[str, str],
|
||||||
progress_sink: Callable[[AppProgressEvent], None] | None = None,
|
progress_sink: Callable[[AppProgressEvent], None] | None = None,
|
||||||
|
performance_profile: str = DEFAULT_PERFORMANCE_PROFILE,
|
||||||
) -> ProcessingOutcome:
|
) -> ProcessingOutcome:
|
||||||
"""Persist confirmed mappings and regenerate only the direct protocol."""
|
"""Persist confirmed mappings and regenerate only the direct protocol."""
|
||||||
review = self.load_speaker_mapping_review(run_dir)
|
review = self.load_speaker_mapping_review(run_dir)
|
||||||
@@ -383,6 +403,7 @@ class MeetingProcessingService:
|
|||||||
context = self.meeting_lab.create_context(context_data)
|
context = self.meeting_lab.create_context(context_data)
|
||||||
|
|
||||||
def relay(event: Any) -> None:
|
def relay(event: Any) -> None:
|
||||||
|
self._update_timer(event)
|
||||||
if progress_sink is not None:
|
if progress_sink is not None:
|
||||||
progress_sink(
|
progress_sink(
|
||||||
AppProgressEvent(
|
AppProgressEvent(
|
||||||
@@ -394,15 +415,23 @@ class MeetingProcessingService:
|
|||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
result = self.meeting_lab.regenerate_protocol(
|
self.timer.start_run()
|
||||||
Path(run_dir),
|
try:
|
||||||
context,
|
result = self.meeting_lab.regenerate_protocol(
|
||||||
relay,
|
Path(run_dir),
|
||||||
model=self.settings.protocol_model,
|
context,
|
||||||
ollama_endpoint=self.settings.ollama_endpoint,
|
relay,
|
||||||
protocol_num_ctx=self.settings.protocol_num_ctx,
|
model=self.settings.protocol_model,
|
||||||
protocol_safe_input_token_budget=(self.settings.protocol_safe_input_token_budget),
|
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()
|
||||||
protocol_path = Path(result.protocol_path)
|
protocol_path = Path(result.protocol_path)
|
||||||
return ProcessingOutcome(
|
return ProcessingOutcome(
|
||||||
succeeded=True,
|
succeeded=True,
|
||||||
@@ -463,3 +492,38 @@ class MeetingProcessingService:
|
|||||||
destination = Path(run_dir) / "protocol_edited.md"
|
destination = Path(run_dir) / "protocol_edited.md"
|
||||||
destination.write_text(text, encoding="utf-8")
|
destination.write_text(text, encoding="utf-8")
|
||||||
return destination
|
return destination
|
||||||
|
|
||||||
|
def _update_timer(self, event: Any) -> None:
|
||||||
|
"""Apply a backend progress event to the shared per-run timer."""
|
||||||
|
if event.stage in STAGES and event.status == "started":
|
||||||
|
self.timer.start_stage(event.stage)
|
||||||
|
elif event.stage in STAGES and event.status in {"completed", "skipped"}:
|
||||||
|
self.timer.finish_stage(event.stage)
|
||||||
|
if event.stage in {"completed", "failed"}:
|
||||||
|
self.timer.finish_run()
|
||||||
|
|
||||||
|
def _persist_timing(self, run_dir: Path | None) -> None:
|
||||||
|
"""Merge final monotonic durations into Meeting Lab run metadata."""
|
||||||
|
if run_dir is None:
|
||||||
|
return
|
||||||
|
metadata_path = Path(run_dir) / "run_metadata.json"
|
||||||
|
metadata: dict[str, Any] = {}
|
||||||
|
if metadata_path.is_file():
|
||||||
|
try:
|
||||||
|
existing = json.loads(metadata_path.read_text(encoding="utf-8"))
|
||||||
|
if isinstance(existing, dict):
|
||||||
|
metadata = existing
|
||||||
|
except (OSError, json.JSONDecodeError):
|
||||||
|
return
|
||||||
|
snapshot = self.timer.snapshot()
|
||||||
|
metadata["timing"] = {
|
||||||
|
"stages_seconds": snapshot.stage_durations,
|
||||||
|
"total_seconds": snapshot.total_duration,
|
||||||
|
}
|
||||||
|
metadata_path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
try:
|
||||||
|
metadata_path.write_text(
|
||||||
|
json.dumps(metadata, indent=2, sort_keys=True) + "\n", encoding="utf-8"
|
||||||
|
)
|
||||||
|
except OSError:
|
||||||
|
return
|
||||||
|
|||||||
@@ -0,0 +1,23 @@
|
|||||||
|
"""Abstract performance profiles resolved to current backend runtime options."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
DEFAULT_PERFORMANCE_PROFILE = "auto"
|
||||||
|
PERFORMANCE_PROFILES = ("auto", "fast", "efficient", "powersave")
|
||||||
|
|
||||||
|
_OLLAMA_THREADS_BY_PROFILE = {
|
||||||
|
"fast": 16,
|
||||||
|
"efficient": 10,
|
||||||
|
"powersave": 4,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def resolve_ollama_num_thread(profile: str) -> int | None:
|
||||||
|
"""Resolve a profile, leaving automatic thread selection to Ollama for Auto."""
|
||||||
|
if profile == "auto":
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
return _OLLAMA_THREADS_BY_PROFILE[profile]
|
||||||
|
except KeyError as exc:
|
||||||
|
allowed = ", ".join(PERFORMANCE_PROFILES)
|
||||||
|
raise ValueError(f"Unknown performance profile {profile!r}; expected: {allowed}.") from exc
|
||||||
@@ -0,0 +1,80 @@
|
|||||||
|
"""Thread-safe monotonic timing state for meeting-processing progress."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import time
|
||||||
|
from collections.abc import Callable
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from threading import Lock
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class TimingSnapshot:
|
||||||
|
"""Presentation-neutral runtime values measured in seconds."""
|
||||||
|
|
||||||
|
stage_durations: dict[str, float]
|
||||||
|
active_stage: str | None
|
||||||
|
total_duration: float
|
||||||
|
running: bool
|
||||||
|
|
||||||
|
|
||||||
|
class ProcessingTimer:
|
||||||
|
"""Track pipeline stages independently from pipeline business logic."""
|
||||||
|
|
||||||
|
def __init__(self, clock: Callable[[], float] = time.perf_counter) -> None:
|
||||||
|
self._clock = clock
|
||||||
|
self._lock = Lock()
|
||||||
|
self._run_started: float | None = None
|
||||||
|
self._run_finished: float | None = None
|
||||||
|
self._stage_started: dict[str, float] = {}
|
||||||
|
self._stage_finished: dict[str, float] = {}
|
||||||
|
self._active_stage: str | None = None
|
||||||
|
|
||||||
|
def start_run(self) -> None:
|
||||||
|
with self._lock:
|
||||||
|
self._run_started = self._clock()
|
||||||
|
self._run_finished = None
|
||||||
|
self._stage_started.clear()
|
||||||
|
self._stage_finished.clear()
|
||||||
|
self._active_stage = None
|
||||||
|
|
||||||
|
def start_stage(self, stage: str) -> None:
|
||||||
|
with self._lock:
|
||||||
|
now = self._clock()
|
||||||
|
if self._active_stage is not None and self._active_stage != stage:
|
||||||
|
self._stage_finished.setdefault(self._active_stage, now)
|
||||||
|
self._stage_started.setdefault(stage, now)
|
||||||
|
self._active_stage = stage
|
||||||
|
|
||||||
|
def finish_stage(self, stage: str) -> None:
|
||||||
|
with self._lock:
|
||||||
|
if stage in self._stage_started:
|
||||||
|
self._stage_finished.setdefault(stage, self._clock())
|
||||||
|
if self._active_stage == stage:
|
||||||
|
self._active_stage = None
|
||||||
|
|
||||||
|
def finish_run(self) -> None:
|
||||||
|
with self._lock:
|
||||||
|
if self._run_finished is not None:
|
||||||
|
return
|
||||||
|
now = self._clock()
|
||||||
|
if self._active_stage is not None:
|
||||||
|
self._stage_finished.setdefault(self._active_stage, now)
|
||||||
|
self._active_stage = None
|
||||||
|
if self._run_started is not None:
|
||||||
|
self._run_finished = now
|
||||||
|
|
||||||
|
def snapshot(self) -> TimingSnapshot:
|
||||||
|
with self._lock:
|
||||||
|
now = self._run_finished if self._run_finished is not None else self._clock()
|
||||||
|
durations = {
|
||||||
|
stage: max(0.0, self._stage_finished.get(stage, now) - started)
|
||||||
|
for stage, started in self._stage_started.items()
|
||||||
|
}
|
||||||
|
total = max(0.0, now - self._run_started) if self._run_started is not None else 0.0
|
||||||
|
return TimingSnapshot(
|
||||||
|
stage_durations=durations,
|
||||||
|
active_stage=self._active_stage,
|
||||||
|
total_duration=total,
|
||||||
|
running=self._run_started is not None and self._run_finished is None,
|
||||||
|
)
|
||||||
+280
-66
@@ -2,7 +2,11 @@
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import time
|
||||||
|
from collections.abc import Callable
|
||||||
|
from concurrent.futures import ThreadPoolExecutor
|
||||||
from datetime import date
|
from datetime import date
|
||||||
|
from queue import Empty, Queue
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from uuid import uuid4
|
from uuid import uuid4
|
||||||
|
|
||||||
@@ -14,6 +18,11 @@ from mka.application.glossary import (
|
|||||||
GlossaryConflictError,
|
GlossaryConflictError,
|
||||||
GlossaryRepository,
|
GlossaryRepository,
|
||||||
)
|
)
|
||||||
|
from mka.application.glossary_yaml import (
|
||||||
|
GlossaryYamlError,
|
||||||
|
export_glossary_yaml,
|
||||||
|
import_glossary_yaml,
|
||||||
|
)
|
||||||
from mka.application.meeting_service import (
|
from mka.application.meeting_service import (
|
||||||
STAGES,
|
STAGES,
|
||||||
AppProgressEvent,
|
AppProgressEvent,
|
||||||
@@ -21,6 +30,7 @@ from mka.application.meeting_service import (
|
|||||||
MeetingProcessingService,
|
MeetingProcessingService,
|
||||||
ParticipantInput,
|
ParticipantInput,
|
||||||
ProcessingOptions,
|
ProcessingOptions,
|
||||||
|
SpeakerMappingReview,
|
||||||
stable_id,
|
stable_id,
|
||||||
)
|
)
|
||||||
from mka.application.people_yaml import (
|
from mka.application.people_yaml import (
|
||||||
@@ -28,6 +38,8 @@ from mka.application.people_yaml import (
|
|||||||
export_people_yaml,
|
export_people_yaml,
|
||||||
import_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 (
|
from mka.application.run_inputs import (
|
||||||
RunInputJsonError,
|
RunInputJsonError,
|
||||||
RunInputState,
|
RunInputState,
|
||||||
@@ -62,6 +74,33 @@ def _render_glossary(repository: GlossaryRepository) -> None:
|
|||||||
"Store canonical core terms and recognition aliases. Compound phrases are "
|
"Store canonical core terms and recognition aliases. Compound phrases are "
|
||||||
"composed from meeting context during protocol generation."
|
"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"):
|
with st.form("glossary_add"):
|
||||||
columns = st.columns(2)
|
columns = st.columns(2)
|
||||||
canonical = columns[0].text_input("Canonical term")
|
canonical = columns[0].text_input("Canonical term")
|
||||||
@@ -304,40 +343,199 @@ def _render_participants() -> list[ParticipantInput]:
|
|||||||
return people
|
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,
|
status_box: Any,
|
||||||
stage_table: Any,
|
stage_table: Any,
|
||||||
progress_slot: Any,
|
|
||||||
states: dict[str, str],
|
states: dict[str, str],
|
||||||
) -> Any:
|
timing: TimingSnapshot,
|
||||||
progress_bar = None
|
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
|
def _apply_progress_event(event: AppProgressEvent, states: dict[str, str]) -> None:
|
||||||
if event.stage in states:
|
if event.stage in states:
|
||||||
states[event.stage] = "running" if event.status == "started" else event.status
|
states[event.stage] = "running" if event.status == "started" else event.status
|
||||||
if event.stage == "failed":
|
if event.stage == "failed":
|
||||||
running = next(
|
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"),
|
(stage for stage, status in states.items() if status == "running"),
|
||||||
None,
|
None,
|
||||||
)
|
)
|
||||||
if running:
|
if running_stage is not None:
|
||||||
states[running] = "failed"
|
states[running_stage] = "failed"
|
||||||
elapsed = f"{event.elapsed_seconds:.1f} s"
|
snapshot = service.timer.snapshot()
|
||||||
message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
_render_progress_status(status_box, stage_table, states, snapshot, latest_message)
|
||||||
status_box.info(f"{message} — elapsed {elapsed}")
|
st.session_state["processing_progress"] = (states.copy(), latest_message, snapshot)
|
||||||
stage_table.table(
|
raise
|
||||||
[{"Stage": STAGE_LABELS[stage], "Status": states[stage]} for stage in STAGES]
|
while not event_queue.empty():
|
||||||
)
|
event = event_queue.get_nowait()
|
||||||
if event.progress is not None:
|
_apply_progress_event(event, states)
|
||||||
if progress_bar is None:
|
latest_message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||||
progress_bar = progress_slot.progress(0.0)
|
running_stage = next((stage for stage, status in states.items() if status == "running"), None)
|
||||||
progress_bar.progress(
|
if running_stage is not None:
|
||||||
min(max(event.progress, 0.0), 1.0),
|
states[running_stage] = "completed" if getattr(result, "succeeded", True) else "failed"
|
||||||
text=f"{STAGE_LABELS.get(event.stage, event.stage)}: {event.progress:.0%}",
|
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:
|
def _render_result() -> None:
|
||||||
@@ -346,6 +544,9 @@ def _render_result() -> None:
|
|||||||
return
|
return
|
||||||
st.divider()
|
st.divider()
|
||||||
st.header("Protocol result")
|
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:
|
if not outcome.succeeded:
|
||||||
st.error(f"Processing failed during {outcome.failed_stage}: {outcome.error_message}")
|
st.error(f"Processing failed during {outcome.failed_stage}: {outcome.error_message}")
|
||||||
if outcome.run_dir:
|
if outcome.run_dir:
|
||||||
@@ -383,29 +584,7 @@ def _render_result() -> None:
|
|||||||
if review is None or not review.speakers:
|
if review is None or not review.speakers:
|
||||||
return
|
return
|
||||||
|
|
||||||
st.subheader("Identify diarized speakers")
|
selections = _render_speaker_mapping(review, outcome.run_dir.name)
|
||||||
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}”")
|
|
||||||
|
|
||||||
duplicate_assignments = len(selections.values()) != len(set(selections.values()))
|
duplicate_assignments = len(selections.values()) != len(set(selections.values()))
|
||||||
if duplicate_assignments:
|
if duplicate_assignments:
|
||||||
@@ -415,12 +594,33 @@ def _render_result() -> None:
|
|||||||
disabled=duplicate_assignments,
|
disabled=duplicate_assignments,
|
||||||
type="primary",
|
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:
|
try:
|
||||||
with st.spinner("Regenerating protocol without rerunning audio processing..."):
|
regenerated, states, message, timing = _run_with_live_progress(
|
||||||
regenerated = service.regenerate_protocol(outcome.run_dir, selections)
|
service,
|
||||||
|
lambda progress_sink: service.regenerate_protocol(
|
||||||
|
run_dir,
|
||||||
|
dict(mapping_items),
|
||||||
|
progress_sink=progress_sink,
|
||||||
|
performance_profile=selected_profile,
|
||||||
|
),
|
||||||
|
states,
|
||||||
|
"Starting protocol regeneration",
|
||||||
|
)
|
||||||
except (OSError, RuntimeError, ValueError) as exc:
|
except (OSError, RuntimeError, ValueError) as exc:
|
||||||
|
_remember_regeneration_timing(
|
||||||
|
states,
|
||||||
|
"Protocol regeneration failed",
|
||||||
|
service.timer.snapshot(),
|
||||||
|
)
|
||||||
st.error(f"Protocol regeneration failed: {exc}")
|
st.error(f"Protocol regeneration failed: {exc}")
|
||||||
else:
|
else:
|
||||||
|
_remember_regeneration_timing(states, message, timing)
|
||||||
st.session_state.outcome = regenerated
|
st.session_state.outcome = regenerated
|
||||||
_queue_edited_protocol(regenerated.original_protocol or "")
|
_queue_edited_protocol(regenerated.original_protocol or "")
|
||||||
st.session_state.speaker_mapping_message = (
|
st.session_state.speaker_mapping_message = (
|
||||||
@@ -458,7 +658,10 @@ def main() -> None:
|
|||||||
description = st.text_area("Description / context", height=100, key="meeting_description")
|
description = st.text_area("Description / context", height=100, key="meeting_description")
|
||||||
metadata_columns = st.columns(3)
|
metadata_columns = st.columns(3)
|
||||||
language = metadata_columns[0].selectbox(
|
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(
|
has_date = metadata_columns[1].checkbox(
|
||||||
"Meeting date is known", value=True, key="meeting_has_date"
|
"Meeting date is known", value=True, key="meeting_has_date"
|
||||||
@@ -470,6 +673,13 @@ def main() -> None:
|
|||||||
participants = _render_participants()
|
participants = _render_participants()
|
||||||
|
|
||||||
st.header("Processing options")
|
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 = st.checkbox(
|
||||||
"Audio normalization",
|
"Audio normalization",
|
||||||
value=True,
|
value=True,
|
||||||
@@ -523,30 +733,34 @@ def main() -> None:
|
|||||||
)
|
)
|
||||||
meeting_id = stable_id(title)
|
meeting_id = stable_id(title)
|
||||||
audio_path = service.preserve_upload(meeting_id, audio.name, audio)
|
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")
|
st.header("Processing status")
|
||||||
status_box = st.empty()
|
|
||||||
stage_table = st.empty()
|
|
||||||
progress_slot = st.empty()
|
|
||||||
states = {stage: "pending" for stage in STAGES}
|
states = {stage: "pending" for stage in STAGES}
|
||||||
if not diarization_enabled:
|
if not diarization_enabled:
|
||||||
states["diarization"] = "skipped"
|
states["diarization"] = "skipped"
|
||||||
callback = _progress_callback(status_box, stage_table, progress_slot, states)
|
outcome, _, _, _ = _run_with_live_progress(
|
||||||
outcome = service.process(
|
service,
|
||||||
audio_path,
|
lambda progress_sink: service.process(
|
||||||
meeting,
|
audio_path,
|
||||||
participants,
|
meeting,
|
||||||
ProcessingOptions(
|
participants,
|
||||||
diarization_enabled=diarization_enabled,
|
ProcessingOptions(
|
||||||
audio_normalization=audio_normalization,
|
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.outcome = outcome
|
||||||
st.session_state.edited_protocol = outcome.original_protocol or ""
|
st.session_state.edited_protocol = outcome.original_protocol or ""
|
||||||
if outcome.succeeded:
|
if outcome.succeeded:
|
||||||
status_box.success("Processing completed.")
|
st.success("Processing completed.")
|
||||||
else:
|
else:
|
||||||
status_box.error(
|
st.error(
|
||||||
f"Processing failed during {outcome.failed_stage}: "
|
f"Processing failed during {outcome.failed_stage}: "
|
||||||
f"{outcome.error_message} Artifacts were preserved."
|
f"{outcome.error_message} Artifacts were preserved."
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -0,0 +1,135 @@
|
|||||||
|
"""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 prepare.call_args.kwargs["normalization_enabled"] is True
|
||||||
|
assert transcribe.call_args.args[3] == language
|
||||||
|
root = outcome.run_dir
|
||||||
|
first = root / "protocol/generations/001"
|
||||||
|
before = {p.name: p.read_bytes() for p in first.iterdir()}
|
||||||
|
original = (root / "diarization/transcript_diarized.json").read_bytes()
|
||||||
|
first_meta = json.loads((first / "runtime_metadata.json").read_text())
|
||||||
|
assert first_meta["speaker_mapping"] == {}
|
||||||
|
assert first_meta["output_language"] == language
|
||||||
|
assert "SPEAKER_00" in (first / "transcript_input.txt").read_text()
|
||||||
|
assert "timing" in json.loads((root / "run_metadata.json").read_text())
|
||||||
|
service.regenerate_protocol(root, {"SPEAKER_00": "martin"}, performance_profile=profile)
|
||||||
|
assert prepare.call_count == transcribe.call_count == diarize.call_count == 1
|
||||||
|
assert (root / "diarization/transcript_diarized.json").read_bytes() == original
|
||||||
|
assert {p.name: p.read_bytes() for p in first.iterdir()} == before
|
||||||
|
second = root / "protocol/generations/002"
|
||||||
|
metadata = json.loads((second / "runtime_metadata.json").read_text())
|
||||||
|
assert metadata["speaker_mapping"] == {"SPEAKER_00": "martin"}
|
||||||
|
assert metadata["speaker_mapping_names"] == {"SPEAKER_00": "Martin"}
|
||||||
|
assert metadata["output_language"] == language
|
||||||
|
assert metadata["num_thread"] == threads
|
||||||
|
assert metadata["glossary_aliases_configured"]["Lumini"] == "Luminy"
|
||||||
|
assert metadata["glossary_replacements"] == []
|
||||||
|
assert "Lumini" in (second / "transcript_input.txt").read_text()
|
||||||
|
assert (root / "protocol.md").resolve() == second / "protocol.md"
|
||||||
|
for prompt, options in calls:
|
||||||
|
assert f"Write the meeting protocol in {expected}." in prompt
|
||||||
|
assert "Luminy" in prompt
|
||||||
|
assert options["num_thread"] == threads
|
||||||
@@ -0,0 +1,147 @@
|
|||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from mka.application.glossary import GlossaryRepository
|
||||||
|
from mka.application.glossary_yaml import (
|
||||||
|
GlossaryYamlError,
|
||||||
|
export_glossary_yaml,
|
||||||
|
import_glossary_yaml,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def repository(tmp_path: Path) -> GlossaryRepository:
|
||||||
|
result = GlossaryRepository(tmp_path / "glossary.sqlite3")
|
||||||
|
result.initialize()
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
def test_round_trip_preserves_ids_aliases_and_active_state(tmp_path: Path) -> None:
|
||||||
|
source = repository(tmp_path / "source")
|
||||||
|
active = source.create(
|
||||||
|
"Größenmaß",
|
||||||
|
"technical_term",
|
||||||
|
aliases=("Groessenmass", "Größen-Maß"),
|
||||||
|
description="Canonical UTF-8 spelling",
|
||||||
|
)
|
||||||
|
inactive = source.create("Legacy Product", "product", is_active=False)
|
||||||
|
|
||||||
|
encoded = export_glossary_yaml(source.list()).encode("utf-8")
|
||||||
|
target = repository(tmp_path / "target")
|
||||||
|
target.create("Will be replaced", "other")
|
||||||
|
target.replace_all(import_glossary_yaml(encoded))
|
||||||
|
|
||||||
|
entries = target.list()
|
||||||
|
assert [entry.id for entry in entries] == [active.id, inactive.id]
|
||||||
|
assert entries[0].canonical_term == "Größenmaß"
|
||||||
|
assert entries[0].aliases == ("Groessenmass", "Größen-Maß")
|
||||||
|
assert entries[0].description == "Canonical UTF-8 spelling"
|
||||||
|
assert entries[0].is_active is True
|
||||||
|
assert entries[1].is_active is False
|
||||||
|
|
||||||
|
|
||||||
|
def test_malformed_yaml_does_not_replace_existing_glossary(tmp_path: Path) -> None:
|
||||||
|
glossary = repository(tmp_path)
|
||||||
|
existing = glossary.create("PBAT", "acronym")
|
||||||
|
|
||||||
|
with pytest.raises(GlossaryYamlError, match="Malformed glossary YAML"):
|
||||||
|
import_glossary_yaml(b"version: 1\nglossary: [")
|
||||||
|
|
||||||
|
assert glossary.list() == [existing]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"entries, message",
|
||||||
|
[
|
||||||
|
(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"id": 1,
|
||||||
|
"canonical_term": "PBAT",
|
||||||
|
"aliases": [],
|
||||||
|
"category": "acronym",
|
||||||
|
"description": None,
|
||||||
|
"active": True,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 2,
|
||||||
|
"canonical_term": "pbat",
|
||||||
|
"aliases": [],
|
||||||
|
"category": "material",
|
||||||
|
"description": None,
|
||||||
|
"active": True,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
"duplicate canonical term or alias",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"id": 1,
|
||||||
|
"canonical_term": "Secugrid",
|
||||||
|
"aliases": ["Sikirgut", "SIKIRGUT"],
|
||||||
|
"category": "product",
|
||||||
|
"description": None,
|
||||||
|
"active": True,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
"aliases must be unique",
|
||||||
|
),
|
||||||
|
(
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"id": 1,
|
||||||
|
"canonical_term": "Secugrid",
|
||||||
|
"aliases": ["PBAT"],
|
||||||
|
"category": "product",
|
||||||
|
"description": None,
|
||||||
|
"active": True,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": 2,
|
||||||
|
"canonical_term": "PBAT",
|
||||||
|
"aliases": [],
|
||||||
|
"category": "acronym",
|
||||||
|
"description": None,
|
||||||
|
"active": False,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
"duplicate canonical term or alias",
|
||||||
|
),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_duplicate_terms_and_aliases_are_rejected(
|
||||||
|
entries: list[dict[str, object]], message: str
|
||||||
|
) -> None:
|
||||||
|
import yaml
|
||||||
|
|
||||||
|
content = yaml.safe_dump({"version": 1, "glossary": entries})
|
||||||
|
|
||||||
|
with pytest.raises(GlossaryYamlError, match=message):
|
||||||
|
import_glossary_yaml(content)
|
||||||
|
|
||||||
|
|
||||||
|
def test_validation_failure_leaves_database_unchanged(tmp_path: Path) -> None:
|
||||||
|
glossary = repository(tmp_path)
|
||||||
|
existing = glossary.create("Existing", "other", aliases=("Existing alias",))
|
||||||
|
invalid = b"""version: 1
|
||||||
|
glossary:
|
||||||
|
- id: 10
|
||||||
|
canonical_term: Duplicate
|
||||||
|
aliases: []
|
||||||
|
category: other
|
||||||
|
description: null
|
||||||
|
active: true
|
||||||
|
- id: 11
|
||||||
|
canonical_term: duplicate
|
||||||
|
aliases: []
|
||||||
|
category: other
|
||||||
|
description: null
|
||||||
|
active: false
|
||||||
|
"""
|
||||||
|
|
||||||
|
with pytest.raises(GlossaryYamlError):
|
||||||
|
imported = import_glossary_yaml(invalid)
|
||||||
|
glossary.replace_all(imported)
|
||||||
|
|
||||||
|
assert glossary.list() == [existing]
|
||||||
@@ -0,0 +1,88 @@
|
|||||||
|
"""Exercise the real executor and Streamlit reruns without model calls."""
|
||||||
|
|
||||||
|
from pathlib import Path
|
||||||
|
from threading import current_thread
|
||||||
|
from types import SimpleNamespace
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
import streamlit
|
||||||
|
from streamlit.testing.v1 import AppTest
|
||||||
|
|
||||||
|
from mka.application.meeting_service import AppProgressEvent, SpeakerMappingReview, SpeakerReview
|
||||||
|
from mka.application.progress_timing import ProcessingTimer
|
||||||
|
from mka.ui import streamlit_app as ui
|
||||||
|
|
||||||
|
|
||||||
|
class GuardedStreamlit:
|
||||||
|
"""Fail if a worker touches the UI session proxy, including callback reads."""
|
||||||
|
|
||||||
|
def __getattr__(self, name):
|
||||||
|
if name == "session_state" and current_thread().name.startswith("ThreadPoolExecutor"):
|
||||||
|
raise AssertionError("Worker accessed Streamlit session state")
|
||||||
|
return getattr(streamlit, name)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("fails", [False, True])
|
||||||
|
@pytest.mark.parametrize("mapped", [False, True])
|
||||||
|
def test_regeneration_captures_ui_values_before_real_worker_and_retains_timing(
|
||||||
|
monkeypatch, fails, mapped
|
||||||
|
):
|
||||||
|
outcome = SimpleNamespace(
|
||||||
|
succeeded=True,
|
||||||
|
run_dir=Path("/tmp/alpha-mapping-test"),
|
||||||
|
speaker_attribution_available=True,
|
||||||
|
original_protocol="Existing protocol",
|
||||||
|
)
|
||||||
|
review = SpeakerMappingReview((SpeakerReview("SPEAKER_00", ()),), (("a", "A"),), {})
|
||||||
|
calls = []
|
||||||
|
|
||||||
|
class Service:
|
||||||
|
def __init__(self, *args):
|
||||||
|
self.timer = ProcessingTimer()
|
||||||
|
|
||||||
|
def load_speaker_mapping_review(self, run_dir):
|
||||||
|
return review
|
||||||
|
|
||||||
|
def regenerate_protocol(self, run_dir, selections, *, progress_sink, performance_profile):
|
||||||
|
assert current_thread().name.startswith("ThreadPoolExecutor")
|
||||||
|
calls.append((run_dir, selections, performance_profile))
|
||||||
|
self.timer.start_run()
|
||||||
|
self.timer.start_stage("protocol_generation")
|
||||||
|
progress_sink(AppProgressEvent("protocol_generation", "started", 0.0))
|
||||||
|
if fails:
|
||||||
|
raise RuntimeError("model unavailable")
|
||||||
|
self.timer.finish_stage("protocol_generation")
|
||||||
|
self.timer.finish_run()
|
||||||
|
progress_sink(AppProgressEvent("protocol_generation", "completed", 0.0))
|
||||||
|
return outcome
|
||||||
|
|
||||||
|
monkeypatch.setattr(ui, "st", GuardedStreamlit())
|
||||||
|
monkeypatch.setattr(ui, "MeetingProcessingService", Service)
|
||||||
|
monkeypatch.setattr(ui.AppSettings, "from_environment", lambda: None)
|
||||||
|
monkeypatch.setattr(ui, "MeetingLabGateway", lambda: None)
|
||||||
|
|
||||||
|
def result_app(outcome):
|
||||||
|
import streamlit as st
|
||||||
|
|
||||||
|
from mka.ui.streamlit_app import _render_result
|
||||||
|
|
||||||
|
st.session_state.outcome = outcome
|
||||||
|
st.session_state.performance_profile = "fast"
|
||||||
|
_render_result()
|
||||||
|
|
||||||
|
app = AppTest.from_function(result_app, args=(outcome,)).run()
|
||||||
|
if mapped:
|
||||||
|
app.selectbox[0].select("a").run()
|
||||||
|
button = next(button for button in app.button if "Regenerate" in button.label)
|
||||||
|
assert not button.disabled
|
||||||
|
button.click().run()
|
||||||
|
assert not app.exception
|
||||||
|
assert calls == [(outcome.run_dir, {"SPEAKER_00": "a"} if mapped else {}, "fast")]
|
||||||
|
states, message, timing = app.session_state["processing_progress"]
|
||||||
|
assert states["protocol_generation"] == ("failed" if fails else "completed")
|
||||||
|
assert not timing.running
|
||||||
|
assert timing.total_duration >= timing.stage_durations["protocol_generation"] >= 0
|
||||||
|
app.run()
|
||||||
|
assert not app.exception
|
||||||
|
assert app.session_state["processing_progress"][2] == timing
|
||||||
|
assert any("total runtime" in item.value for item in app.info)
|
||||||
@@ -14,6 +14,7 @@ from mka.application.meeting_service import (
|
|||||||
ParticipantInput,
|
ParticipantInput,
|
||||||
ProcessingOptions,
|
ProcessingOptions,
|
||||||
)
|
)
|
||||||
|
from mka.application.progress_timing import ProcessingTimer
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
@@ -323,19 +324,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(
|
def test_process_translates_configuration_and_disables_diarization(
|
||||||
tmp_path: Path,
|
tmp_path: Path,
|
||||||
|
language: str,
|
||||||
) -> None:
|
) -> None:
|
||||||
service, gateway = make_service(tmp_path)
|
service, gateway = make_service(tmp_path)
|
||||||
audio = tmp_path / "meeting.wav"
|
audio = tmp_path / "meeting.wav"
|
||||||
audio.write_bytes(b"audio")
|
audio.write_bytes(b"audio")
|
||||||
|
|
||||||
outcome = service.process(
|
outcome = service.process(
|
||||||
audio, meeting(), participants(), ProcessingOptions(diarization_enabled=False)
|
audio,
|
||||||
|
replace(meeting(), language=language),
|
||||||
|
participants(),
|
||||||
|
ProcessingOptions(diarization_enabled=False),
|
||||||
)
|
)
|
||||||
|
|
||||||
assert outcome.succeeded
|
assert outcome.succeeded
|
||||||
assert gateway.config_values is not None
|
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["diarization"] == "off"
|
||||||
assert gateway.config_values["model"] == "test:model"
|
assert gateway.config_values["model"] == "test:model"
|
||||||
assert gateway.config_values["protocol_num_ctx"] == 32_768
|
assert gateway.config_values["protocol_num_ctx"] == 32_768
|
||||||
@@ -503,6 +513,7 @@ def test_protocol_only_regeneration_persists_mappings_without_rewriting_transcri
|
|||||||
"SPEAKER_00": "martin"
|
"SPEAKER_00": "martin"
|
||||||
}
|
}
|
||||||
assert gateway.regeneration["options"]["protocol_num_ctx"] == 32_768
|
assert gateway.regeneration["options"]["protocol_num_ctx"] == 32_768
|
||||||
|
assert gateway.regeneration["options"]["glossary_aliases"] == {"Enlyse": "ENLYZE"}
|
||||||
assert events[0].stage == "protocol_generation"
|
assert events[0].stage == "protocol_generation"
|
||||||
assert source.read_bytes() == original_source
|
assert source.read_bytes() == original_source
|
||||||
|
|
||||||
@@ -552,3 +563,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.parent == tmp_path / "meetings" / "meeting-1" / "uploads"
|
||||||
assert destination.name.endswith("_unsafe.wav")
|
assert destination.name.endswith("_unsafe.wav")
|
||||||
assert destination.read_bytes() == b"source audio"
|
assert destination.read_bytes() == b"source audio"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("profile", "threads"),
|
||||||
|
[("fast", 16), ("efficient", 10), ("powersave", 4)],
|
||||||
|
)
|
||||||
|
def test_process_resolves_profile_for_protocol_generation(
|
||||||
|
tmp_path: Path, profile: str, threads: int
|
||||||
|
) -> None:
|
||||||
|
service, gateway = make_service(tmp_path)
|
||||||
|
audio = tmp_path / "meeting.wav"
|
||||||
|
audio.write_bytes(b"audio")
|
||||||
|
|
||||||
|
service.process(
|
||||||
|
audio,
|
||||||
|
meeting(),
|
||||||
|
participants(),
|
||||||
|
ProcessingOptions(performance_profile=profile),
|
||||||
|
)
|
||||||
|
|
||||||
|
assert gateway.config_values is not None
|
||||||
|
assert gateway.config_values["protocol_num_thread"] == threads
|
||||||
|
|
||||||
|
|
||||||
|
def test_process_defaults_to_backend_thread_selection(tmp_path: Path) -> None:
|
||||||
|
service, gateway = make_service(tmp_path)
|
||||||
|
audio = tmp_path / "meeting.wav"
|
||||||
|
audio.write_bytes(b"audio")
|
||||||
|
|
||||||
|
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||||
|
|
||||||
|
assert gateway.config_values is not None
|
||||||
|
assert gateway.config_values["protocol_num_thread"] is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_protocol_regeneration_propagates_explicit_profile(tmp_path: Path) -> None:
|
||||||
|
service, gateway = make_service(tmp_path)
|
||||||
|
write_speaker_review_artifacts(gateway.run_dir)
|
||||||
|
|
||||||
|
service.regenerate_protocol(
|
||||||
|
gateway.run_dir,
|
||||||
|
{},
|
||||||
|
performance_profile="fast",
|
||||||
|
)
|
||||||
|
|
||||||
|
assert gateway.regeneration is not None
|
||||||
|
assert gateway.regeneration["options"]["protocol_num_thread"] == 16
|
||||||
|
|
||||||
|
|
||||||
|
def test_process_passes_only_active_explicit_glossary_aliases(tmp_path: Path) -> None:
|
||||||
|
service, gateway = make_service(tmp_path)
|
||||||
|
service.glossary.create("Luminy", "product", aliases=("Lumini",))
|
||||||
|
inactive = service.glossary.create("PBAT", "acronym", aliases=("PBRT",))
|
||||||
|
service.glossary.set_active(inactive.id, False)
|
||||||
|
audio = tmp_path / "meeting.wav"
|
||||||
|
audio.write_bytes(b"audio")
|
||||||
|
|
||||||
|
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||||
|
|
||||||
|
assert gateway.config_values is not None
|
||||||
|
assert gateway.config_values["glossary_aliases"] == {"Lumini": "Luminy"}
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("fails", [False, True])
|
||||||
|
def test_process_persists_completed_and_failed_timings(tmp_path: Path, fails: bool) -> None:
|
||||||
|
service, gateway = make_service(tmp_path)
|
||||||
|
gateway.fail = fails
|
||||||
|
values = iter([0.0, 1.0, 9.0, 10.0, 20.0])
|
||||||
|
service.timer = ProcessingTimer(lambda: next(values))
|
||||||
|
audio = tmp_path / "meeting.wav"
|
||||||
|
audio.write_bytes(b"audio")
|
||||||
|
|
||||||
|
outcome = service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||||
|
|
||||||
|
metadata = json.loads((gateway.run_dir / "run_metadata.json").read_text(encoding="utf-8"))
|
||||||
|
assert metadata["timing"] == {
|
||||||
|
"stages_seconds": {"preparing": 8.0, "transcription": 10.0},
|
||||||
|
"total_seconds": 20.0,
|
||||||
|
}
|
||||||
|
assert outcome.succeeded is not fails
|
||||||
|
if fails:
|
||||||
|
assert metadata["failure"]["type"] == "TranscriptionError"
|
||||||
|
|
||||||
|
|
||||||
|
def test_regeneration_starts_fresh_timer_and_retains_final_duration(
|
||||||
|
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||||
|
) -> None:
|
||||||
|
service, gateway = make_service(tmp_path)
|
||||||
|
write_speaker_review_artifacts(gateway.run_dir)
|
||||||
|
clock = SimpleNamespace(value=10.0)
|
||||||
|
service.timer = ProcessingTimer(lambda: clock.value)
|
||||||
|
audio = tmp_path / "meeting.wav"
|
||||||
|
audio.write_bytes(b"audio")
|
||||||
|
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||||
|
|
||||||
|
clock.value = 100.0
|
||||||
|
observed = []
|
||||||
|
|
||||||
|
def regenerate(run_dir: Path, context: Any, progress_sink: Any, **options: Any) -> Any:
|
||||||
|
progress_sink(
|
||||||
|
SimpleNamespace(
|
||||||
|
stage="protocol_generation",
|
||||||
|
status="started",
|
||||||
|
elapsed_seconds=0.0,
|
||||||
|
progress=None,
|
||||||
|
message=None,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
clock.value = 105.0
|
||||||
|
observed.append(service.timer.snapshot())
|
||||||
|
progress_sink(
|
||||||
|
SimpleNamespace(
|
||||||
|
stage="protocol_generation",
|
||||||
|
status="completed",
|
||||||
|
elapsed_seconds=5.0,
|
||||||
|
progress=None,
|
||||||
|
message=None,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
clock.value = 109.0
|
||||||
|
protocol = run_dir / "protocol.md"
|
||||||
|
protocol.write_text("# Regenerated protocol\n", encoding="utf-8")
|
||||||
|
return SimpleNamespace(run_dir=run_dir, protocol_path=protocol)
|
||||||
|
|
||||||
|
monkeypatch.setattr(gateway, "regenerate_protocol", regenerate)
|
||||||
|
|
||||||
|
service.regenerate_protocol(gateway.run_dir, {"SPEAKER_00": "martin"})
|
||||||
|
final = service.timer.snapshot()
|
||||||
|
|
||||||
|
assert observed[0].active_stage == "protocol_generation"
|
||||||
|
assert observed[0].stage_durations == {"protocol_generation": 5.0}
|
||||||
|
assert observed[0].total_duration == 5.0
|
||||||
|
assert observed[0].running is True
|
||||||
|
assert final.stage_durations == {"protocol_generation": 5.0}
|
||||||
|
assert final.total_duration == 9.0
|
||||||
|
assert final.running is False
|
||||||
|
|
||||||
|
|
||||||
|
def test_failed_regeneration_freezes_active_and_total_durations(
|
||||||
|
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||||
|
) -> None:
|
||||||
|
service, gateway = make_service(tmp_path)
|
||||||
|
write_speaker_review_artifacts(gateway.run_dir)
|
||||||
|
clock = SimpleNamespace(value=20.0)
|
||||||
|
service.timer = ProcessingTimer(lambda: clock.value)
|
||||||
|
|
||||||
|
def fail_regeneration(run_dir: Path, context: Any, progress_sink: Any, **options: Any) -> Any:
|
||||||
|
progress_sink(
|
||||||
|
SimpleNamespace(
|
||||||
|
stage="protocol_generation",
|
||||||
|
status="started",
|
||||||
|
elapsed_seconds=0.0,
|
||||||
|
progress=None,
|
||||||
|
message=None,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
clock.value = 27.0
|
||||||
|
raise RuntimeError("generation failed")
|
||||||
|
|
||||||
|
monkeypatch.setattr(gateway, "regenerate_protocol", fail_regeneration)
|
||||||
|
|
||||||
|
with pytest.raises(RuntimeError, match="generation failed"):
|
||||||
|
service.regenerate_protocol(gateway.run_dir, {})
|
||||||
|
clock.value = 99.0
|
||||||
|
final = service.timer.snapshot()
|
||||||
|
|
||||||
|
assert final.stage_durations == {"protocol_generation": 7.0}
|
||||||
|
assert final.total_duration == 7.0
|
||||||
|
assert final.running is False
|
||||||
|
|||||||
@@ -0,0 +1,21 @@
|
|||||||
|
import pytest
|
||||||
|
|
||||||
|
from mka.application.meeting_service import ProcessingOptions
|
||||||
|
from mka.application.performance import PERFORMANCE_PROFILES, resolve_ollama_num_thread
|
||||||
|
|
||||||
|
|
||||||
|
def test_default_profile_is_auto() -> None:
|
||||||
|
assert ProcessingOptions().performance_profile == "auto"
|
||||||
|
assert PERFORMANCE_PROFILES[0] == "auto"
|
||||||
|
|
||||||
|
|
||||||
|
def test_auto_profile_has_no_ollama_thread_override() -> None:
|
||||||
|
assert resolve_ollama_num_thread("auto") is None
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("profile", "threads"),
|
||||||
|
[("fast", 16), ("efficient", 10), ("powersave", 4)],
|
||||||
|
)
|
||||||
|
def test_profiles_resolve_to_current_ollama_thread_counts(profile: str, threads: int) -> None:
|
||||||
|
assert resolve_ollama_num_thread(profile) == threads
|
||||||
@@ -0,0 +1,45 @@
|
|||||||
|
from mka.application.progress_timing import ProcessingTimer
|
||||||
|
|
||||||
|
|
||||||
|
class FakeClock:
|
||||||
|
def __init__(self, value: float = 0.0) -> None:
|
||||||
|
self.value = value
|
||||||
|
|
||||||
|
def __call__(self) -> float:
|
||||||
|
return self.value
|
||||||
|
|
||||||
|
|
||||||
|
def test_completed_duration_is_retained_while_active_stage_increases() -> None:
|
||||||
|
clock = FakeClock(10.0)
|
||||||
|
timer = ProcessingTimer(clock)
|
||||||
|
timer.start_run()
|
||||||
|
timer.start_stage("preparing")
|
||||||
|
clock.value = 18.0
|
||||||
|
timer.finish_stage("preparing")
|
||||||
|
timer.start_stage("transcription")
|
||||||
|
|
||||||
|
clock.value = 20.0
|
||||||
|
first = timer.snapshot()
|
||||||
|
clock.value = 25.5
|
||||||
|
second = timer.snapshot()
|
||||||
|
|
||||||
|
assert first.stage_durations == {"preparing": 8.0, "transcription": 2.0}
|
||||||
|
assert second.stage_durations == {"preparing": 8.0, "transcription": 7.5}
|
||||||
|
assert second.total_duration == 15.5
|
||||||
|
|
||||||
|
|
||||||
|
def test_failed_active_stage_and_total_duration_are_frozen() -> None:
|
||||||
|
clock = FakeClock(2.0)
|
||||||
|
timer = ProcessingTimer(clock)
|
||||||
|
timer.start_run()
|
||||||
|
clock.value = 5.0
|
||||||
|
timer.start_stage("transcription")
|
||||||
|
clock.value = 14.0
|
||||||
|
timer.finish_run()
|
||||||
|
clock.value = 99.0
|
||||||
|
|
||||||
|
snapshot = timer.snapshot()
|
||||||
|
|
||||||
|
assert snapshot.stage_durations["transcription"] == 9.0
|
||||||
|
assert snapshot.total_duration == 12.0
|
||||||
|
assert snapshot.running is False
|
||||||
@@ -0,0 +1,132 @@
|
|||||||
|
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_anonymous_regeneration_remains_allowed(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",
|
||||||
|
)
|
||||||
|
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()
|
||||||
|
regenerate = app.button[1]
|
||||||
|
assert not regenerate.disabled
|
||||||
|
assert app.warning
|
||||||
|
regenerate.click().run()
|
||||||
|
assert calls == [{}]
|
||||||
|
assert not app.exception
|
||||||
@@ -1,6 +1,7 @@
|
|||||||
from datetime import date
|
from datetime import date
|
||||||
|
|
||||||
from mka.application.meeting_service import ParticipantInput
|
from mka.application.meeting_service import ParticipantInput
|
||||||
|
from mka.application.progress_timing import TimingSnapshot
|
||||||
from mka.application.run_inputs import RunInputState
|
from mka.application.run_inputs import RunInputState
|
||||||
from mka.ui import streamlit_app
|
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 "source_media_1" not in state
|
||||||
assert "pending_run_inputs" not in state
|
assert "pending_run_inputs" not in state
|
||||||
assert state["source_media_0"] is existing_upload
|
assert state["source_media_0"] is existing_upload
|
||||||
|
|
||||||
|
|
||||||
|
def test_duration_formatting_covers_seconds_minutes_and_hours() -> None:
|
||||||
|
assert streamlit_app._format_duration(8.9) == "00:08"
|
||||||
|
assert streamlit_app._format_duration(7 * 60 + 18) == "07:18"
|
||||||
|
assert streamlit_app._format_duration(3600 + 3 * 60 + 42) == "1:03:42"
|
||||||
|
|
||||||
|
|
||||||
|
def test_completed_regeneration_timing_survives_frontend_rerun(monkeypatch) -> None:
|
||||||
|
state = {}
|
||||||
|
monkeypatch.setattr(streamlit_app.st, "session_state", state)
|
||||||
|
timing = TimingSnapshot(
|
||||||
|
stage_durations={"protocol_generation": 4.5},
|
||||||
|
active_stage=None,
|
||||||
|
total_duration=4.5,
|
||||||
|
running=False,
|
||||||
|
)
|
||||||
|
states = {stage: "skipped" for stage in streamlit_app.STAGES}
|
||||||
|
states["protocol_generation"] = "completed"
|
||||||
|
|
||||||
|
streamlit_app._remember_regeneration_timing(states, "Protocol regenerated", timing)
|
||||||
|
|
||||||
|
saved_states, saved_message, saved_timing = state["regeneration_progress"]
|
||||||
|
assert saved_states["protocol_generation"] == "completed"
|
||||||
|
assert saved_message == "Protocol regenerated"
|
||||||
|
assert saved_timing == timing
|
||||||
|
|||||||
Reference in New Issue
Block a user