Compare commits
9
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bd48aec5aa | ||
|
|
b50b87c32b | ||
|
|
9e7173c2a7 | ||
|
|
70cb1e0789 | ||
|
|
500d2c4c67 | ||
|
|
77be2b39ff | ||
|
|
e2a0434941 | ||
|
|
9840953687 | ||
|
|
e55dd74c19 |
@@ -4,6 +4,14 @@
|
||||
|
||||
### Added
|
||||
|
||||
- Retained live stage and total timings for processing and protocol regeneration.
|
||||
- Immutable diagnostic generation history with atomic latest-result publication.
|
||||
|
||||
- Auto/Fast/Efficient/Powersave protocol profiles with backend-selected threads
|
||||
by default, also applied during mapped-speaker regeneration.
|
||||
|
||||
- Versioned UTF-8 YAML import/export for the SQLite terminology glossary, with
|
||||
stable entry IDs, complete validation, and atomic replacement semantics.
|
||||
- First Streamlit MVP for audio upload, Meeting Context entry, participant
|
||||
management, processing progress and protocol editing.
|
||||
- Application service and Meeting Lab adapter with environment-based runtime
|
||||
@@ -15,6 +23,8 @@
|
||||
|
||||
### Changed
|
||||
|
||||
- Record active glossary configuration without rewriting protocol transcript input.
|
||||
|
||||
- Reconciled the documented MVP with the validated Meeting Lab backend.
|
||||
- Selected `whisper.cpp` for transcription and made diarization optional.
|
||||
- Defined the Python API, progress-event and Meeting Context integration
|
||||
|
||||
@@ -1,5 +1,14 @@
|
||||
# Project Knowledge
|
||||
|
||||
## Meeting language
|
||||
|
||||
The GUI's `meeting_language` is passed through `MeetingDetails.language` to both
|
||||
Meeting Context `meeting.language` and `MvpMeetingConfig.language` (Whisper).
|
||||
Meeting Lab derives explicit protocol output language from the persisted context,
|
||||
also during regeneration, and records derived `output_language` provenance in
|
||||
protocol runtime metadata. Missing context language defaults to German. There is
|
||||
no separate protocol-language setting or automatic context/transcript translation.
|
||||
|
||||
## Vision
|
||||
|
||||
Meeting Assistant turns recorded meetings into reviewed, user-facing protocols
|
||||
@@ -48,6 +57,8 @@ Meeting Assistant owns user interaction and product workflow:
|
||||
- a global SQLite terminology glossary whose active canonical core terms and
|
||||
recognition aliases are rendered through Meeting Context into direct
|
||||
protocol prompts
|
||||
- versioned UTF-8 YAML glossary interchange (`version: 1`) that preserves entry
|
||||
IDs and atomically replaces SQLite state only after complete validation
|
||||
- export and presentation of protocol versions
|
||||
|
||||
Processing logic must not be duplicated in the application.
|
||||
@@ -84,6 +95,14 @@ PyTorch/ROCm processed the same duration in approximately 98.6 seconds. These
|
||||
are validation observations, not performance guarantees or hardware
|
||||
requirements. CPU execution remains supported and may be substantially slower.
|
||||
|
||||
Protocol-generation performance is selected through intentionally abstract
|
||||
profiles rather than hardware controls in the GUI. `auto` is the default and
|
||||
delegates thread selection to the inference backend; with Ollama this means
|
||||
omitting `num_thread` entirely. The explicit resource profiles currently map
|
||||
`fast` to 16 Ollama CPU threads, `efficient` to 10, and `powersave` to 4. These
|
||||
concrete mappings belong to the backend/configuration boundary and may evolve
|
||||
independently of the user-facing profile semantics.
|
||||
|
||||
## Meeting Context and Speakers
|
||||
|
||||
`MeetingContext` is a structured domain object containing meeting metadata,
|
||||
@@ -98,6 +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
|
||||
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
|
||||
|
||||
Meeting Lab emits stage-based progress events for:
|
||||
@@ -139,3 +164,19 @@ product work, not a current MVP requirement.
|
||||
- no automatic identity claims
|
||||
- human review of generated protocols
|
||||
- small modules and simple interfaces
|
||||
|
||||
Active glossary aliases are forwarded to Meeting Lab as provenance metadata.
|
||||
Canonical terminology remains Meeting Context guidance; the exact protocol
|
||||
transcript input and the raw Whisper/diarization artifacts are not rewritten.
|
||||
`glossary_replacements` remains empty.
|
||||
|
||||
Processing and protocol regeneration show measured monotonic stage/total timings,
|
||||
including frozen failure durations. The latest timing display survives ordinary
|
||||
Streamlit reruns. Initial-run durations are saved in `run_metadata.json`;
|
||||
regeneration display timings remain session-local. Worker callbacks capture UI
|
||||
configuration before dispatch and never read Streamlit session state.
|
||||
|
||||
Meeting Lab retains immutable generation records including prompt, transcript
|
||||
input, model response/metadata, glossary configuration, mappings and Meeting
|
||||
Context. Latest protocol/diagnostic/context paths resolve through an atomic
|
||||
`protocol/current` link. Copy whole runs with relative symlinks preserved.
|
||||
|
||||
@@ -23,6 +23,13 @@ reference recorder, but imported audio is not tied to OBS-specific behavior.
|
||||
|
||||
## Processing Pipeline
|
||||
|
||||
The existing **Meeting language** selector controls both transcription and protocol
|
||||
output: `de` means German for both, and `en` means English for both. The selection
|
||||
is saved in Meeting Context (`meeting.language`) and reused for protocol-only
|
||||
regeneration without retranscription. Older contexts without a language retain
|
||||
German protocol output. Names, speaker mappings and authored Meeting Context are
|
||||
preserved; no transcript translation stage is added.
|
||||
|
||||
```text
|
||||
Audio file
|
||||
-> FFmpeg preparation (mono, 16 kHz PCM WAV; normalization optional)
|
||||
@@ -42,6 +49,11 @@ compatibility path. Diarization is optional and produces anonymous speaker
|
||||
labels. A label identifies a participant only when the user explicitly
|
||||
confirms the mapping; automatic speaker-name inference is not allowed.
|
||||
|
||||
Speaker selectors hide participants already assigned to other speakers while keeping
|
||||
the current assignment available. The UI shows detected, assigned and unassigned
|
||||
counts, marks unassigned speakers, and warns when mappings remain incomplete.
|
||||
Anonymous-speaker protocol generation remains available.
|
||||
|
||||
After a diarized run, the result view lists detected `SPEAKER_XX` labels with
|
||||
short transcript excerpts. Confirmed mappings regenerate only the protocol
|
||||
from the existing diarized transcript; audio preparation, Whisper and Pyannote
|
||||
@@ -70,6 +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
|
||||
every WAV, FLAC or M4A source as canonical audio before transcription.
|
||||
|
||||
The **Performance profile** selector controls protocol-generation runtime using
|
||||
the abstract Auto (default), Fast, Efficient, and Powersave profiles. Auto lets
|
||||
the inference backend select its own thread configuration; for Ollama, Meeting
|
||||
Assistant intentionally sends no `num_thread` option. Fast, Efficient, and
|
||||
Powersave are explicit resource profiles currently mapped to 16, 10, and 4
|
||||
Ollama CPU threads. These concrete mappings may evolve independently of the UI
|
||||
semantics.
|
||||
|
||||
The People section can export its current entries to a UTF-8 `people.yaml` file
|
||||
and replace them from a previous `.yaml` or `.yml` export. Stable person IDs,
|
||||
names, roles, organizations and attendance states are retained. This is a small
|
||||
@@ -131,6 +151,27 @@ The database defaults to `data/database/glossary.sqlite3`. It is created and
|
||||
bootstrapped automatically and can be moved with `MKA_GLOSSARY_DATABASE`.
|
||||
Glossary integration does not rewrite raw or diarized transcript artifacts.
|
||||
|
||||
Use **Export glossary** to download the complete SQLite-backed glossary as
|
||||
UTF-8 YAML, and **Import glossary** to replace it from a previously exported
|
||||
file. Imports are fully validated before a single SQLite transaction replaces
|
||||
the current glossary, so malformed or conflicting files leave existing data
|
||||
unchanged. Schema version 1 is:
|
||||
|
||||
```yaml
|
||||
version: 1
|
||||
glossary:
|
||||
- id: 1
|
||||
canonical_term: Secugrid HS
|
||||
aliases:
|
||||
- Sikirgut
|
||||
category: product
|
||||
description: Canonical product spelling
|
||||
active: true
|
||||
```
|
||||
|
||||
`id` is the stable SQLite glossary-entry identifier. Categories are `product`,
|
||||
`material`, `organization`, `technical_term`, `acronym`, or `other`.
|
||||
|
||||
## Input configuration import and export
|
||||
|
||||
Use **Export inputs** to save the current Meeting Assistant run-input form as a
|
||||
@@ -188,3 +229,19 @@ transcription; Meeting Assistant does not duplicate audio conversion.
|
||||
|
||||
See [Architecture](docs/architecture.md), [Project Knowledge](PROJECT_KNOWLEDGE.md),
|
||||
[Roadmap](ROADMAP.md) and [ADR 0011](docs/adr/0011-use-meeting-lab-mvp-backend.md).
|
||||
|
||||
Active glossary aliases are forwarded to Meeting Lab as provenance metadata.
|
||||
Canonical terminology remains Meeting Context guidance; the exact protocol
|
||||
transcript input and the raw Whisper/diarization artifacts are not rewritten.
|
||||
`glossary_replacements` remains empty.
|
||||
|
||||
Processing and protocol regeneration show measured monotonic stage/total timings,
|
||||
including frozen failure durations. The latest timing display survives ordinary
|
||||
Streamlit reruns. Initial-run durations are saved in `run_metadata.json`;
|
||||
regeneration display timings remain session-local. Worker callbacks capture UI
|
||||
configuration before dispatch and never read Streamlit session state.
|
||||
|
||||
Meeting Lab retains immutable generation records including prompt, transcript
|
||||
input, model response/metadata, glossary configuration, mappings and Meeting
|
||||
Context. Latest protocol/diagnostic/context paths resolve through an atomic
|
||||
`protocol/current` link. Copy whole runs with relative symlinks preserved.
|
||||
|
||||
@@ -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]
|
||||
name = "meeting-assistant"
|
||||
version = "0.1.0"
|
||||
version = "0.1.0a1"
|
||||
description = "Local-first meeting transcription and knowledge capture assistant"
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.11"
|
||||
|
||||
@@ -35,6 +35,18 @@ class GlossaryEntry:
|
||||
updated_at: str
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class GlossaryReplacementEntry:
|
||||
"""Validated values used to atomically replace the persisted glossary."""
|
||||
|
||||
id: int
|
||||
canonical_term: str
|
||||
category: str
|
||||
description: str | None
|
||||
is_active: bool
|
||||
aliases: tuple[str, ...]
|
||||
|
||||
|
||||
class GlossaryRepository:
|
||||
"""Small data-access boundary for the local glossary database."""
|
||||
|
||||
@@ -194,6 +206,64 @@ class GlossaryRepository:
|
||||
if cursor.rowcount == 0:
|
||||
raise KeyError(f"Unknown glossary entry: {entry_id}")
|
||||
|
||||
def replace_all(self, entries: Iterable[GlossaryReplacementEntry]) -> None:
|
||||
"""Atomically replace every glossary entry, preserving supplied stable IDs."""
|
||||
replacements = tuple(entries)
|
||||
validated: list[GlossaryReplacementEntry] = []
|
||||
seen_ids: set[int] = set()
|
||||
seen_terms: set[str] = set()
|
||||
for entry in replacements:
|
||||
if type(entry.id) is not int or entry.id <= 0:
|
||||
raise ValueError("Glossary entry IDs must be positive integers.")
|
||||
if entry.id in seen_ids:
|
||||
raise GlossaryConflictError(f"Duplicate glossary entry ID: {entry.id}.")
|
||||
canonical, aliases = self._validate_values(
|
||||
entry.canonical_term, entry.category, entry.aliases
|
||||
)
|
||||
folded_terms = {canonical.casefold(), *(alias.casefold() for alias in aliases)}
|
||||
if seen_terms.intersection(folded_terms):
|
||||
raise GlossaryConflictError("Canonical term or alias already exists.")
|
||||
seen_ids.add(entry.id)
|
||||
seen_terms.update(folded_terms)
|
||||
validated.append(
|
||||
GlossaryReplacementEntry(
|
||||
id=entry.id,
|
||||
canonical_term=canonical,
|
||||
category=entry.category,
|
||||
description=_optional_text(entry.description),
|
||||
is_active=entry.is_active,
|
||||
aliases=aliases,
|
||||
)
|
||||
)
|
||||
|
||||
now = _timestamp()
|
||||
try:
|
||||
with self._connect() as connection:
|
||||
connection.execute("DELETE FROM glossary_entries")
|
||||
connection.executemany(
|
||||
"""INSERT INTO glossary_entries
|
||||
(id, canonical_term, category, description, is_active, created_at, updated_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?)""",
|
||||
(
|
||||
(
|
||||
entry.id,
|
||||
entry.canonical_term,
|
||||
entry.category,
|
||||
entry.description,
|
||||
entry.is_active,
|
||||
now,
|
||||
now,
|
||||
)
|
||||
for entry in validated
|
||||
),
|
||||
)
|
||||
connection.executemany(
|
||||
"INSERT INTO glossary_aliases (entry_id, alias, created_at) VALUES (?, ?, ?)",
|
||||
((entry.id, alias, now) for entry in validated for alias in entry.aliases),
|
||||
)
|
||||
except sqlite3.IntegrityError as exc:
|
||||
raise GlossaryConflictError("Canonical term or alias already exists.") from exc
|
||||
|
||||
@contextmanager
|
||||
def _connect(self) -> Iterator[sqlite3.Connection]:
|
||||
connection = sqlite3.connect(self.database_path, timeout=5)
|
||||
@@ -285,6 +355,11 @@ def render_glossary_terms(entries: Iterable[GlossaryEntry]) -> list[str]:
|
||||
return rendered
|
||||
|
||||
|
||||
def glossary_alias_mapping(entries: Iterable[GlossaryEntry]) -> dict[str, str]:
|
||||
"""Return explicit alias configuration for provenance, not substitution."""
|
||||
return {alias: entry.canonical_term for entry in entries for alias in entry.aliases}
|
||||
|
||||
|
||||
def _timestamp() -> str:
|
||||
return datetime.now(UTC).isoformat(timespec="seconds")
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
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")
|
||||
|
||||
@@ -65,6 +71,7 @@ class ParticipantInput:
|
||||
class ProcessingOptions:
|
||||
diarization_enabled: bool = False
|
||||
audio_normalization: bool = True
|
||||
performance_profile: str = DEFAULT_PERFORMANCE_PROFILE
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
@@ -114,11 +121,13 @@ class MeetingProcessingService:
|
||||
settings: AppSettings,
|
||||
meeting_lab: MeetingLabPort,
|
||||
glossary: GlossaryRepository | None = None,
|
||||
timer: ProcessingTimer | None = None,
|
||||
) -> None:
|
||||
self.settings = settings
|
||||
self.meeting_lab = meeting_lab
|
||||
self.glossary = glossary or GlossaryRepository(settings.glossary_database)
|
||||
self.glossary.initialize()
|
||||
self.timer = timer or ProcessingTimer()
|
||||
|
||||
def build_context(
|
||||
self,
|
||||
@@ -232,6 +241,8 @@ class MeetingProcessingService:
|
||||
"protocol_safe_input_token_budget": (
|
||||
self.settings.protocol_safe_input_token_budget
|
||||
),
|
||||
"protocol_num_thread": resolve_ollama_num_thread(options.performance_profile),
|
||||
"glossary_aliases": glossary_alias_mapping(self.glossary.list(active_only=True)),
|
||||
"diarization": (
|
||||
self.settings.diarization_mode if options.diarization_enabled else "off"
|
||||
),
|
||||
@@ -241,11 +252,13 @@ class MeetingProcessingService:
|
||||
}
|
||||
)
|
||||
current_stage: str | None = None
|
||||
self.timer.start_run()
|
||||
|
||||
def relay(event: Any) -> None:
|
||||
nonlocal current_stage
|
||||
if event.stage in STAGES and event.status == "started":
|
||||
current_stage = event.stage
|
||||
self._update_timer(event)
|
||||
app_event = AppProgressEvent(
|
||||
stage=event.stage,
|
||||
status=event.status,
|
||||
@@ -256,7 +269,13 @@ class MeetingProcessingService:
|
||||
if progress_sink is not None:
|
||||
progress_sink(app_event)
|
||||
|
||||
result = self.meeting_lab.run(config, context, relay)
|
||||
try:
|
||||
result = self.meeting_lab.run(config, context, relay)
|
||||
except Exception:
|
||||
self.timer.finish_run()
|
||||
raise
|
||||
self.timer.finish_run()
|
||||
self._persist_timing(result.run_dir)
|
||||
if result.exit_code != 0:
|
||||
failure = self._read_failure(result.run_dir)
|
||||
backend_stage = failure.get("stage")
|
||||
@@ -357,6 +376,7 @@ class MeetingProcessingService:
|
||||
run_dir: Path,
|
||||
speaker_mappings: dict[str, str],
|
||||
progress_sink: Callable[[AppProgressEvent], None] | None = None,
|
||||
performance_profile: str = DEFAULT_PERFORMANCE_PROFILE,
|
||||
) -> ProcessingOutcome:
|
||||
"""Persist confirmed mappings and regenerate only the direct protocol."""
|
||||
review = self.load_speaker_mapping_review(run_dir)
|
||||
@@ -383,6 +403,7 @@ class MeetingProcessingService:
|
||||
context = self.meeting_lab.create_context(context_data)
|
||||
|
||||
def relay(event: Any) -> None:
|
||||
self._update_timer(event)
|
||||
if progress_sink is not None:
|
||||
progress_sink(
|
||||
AppProgressEvent(
|
||||
@@ -394,15 +415,23 @@ class MeetingProcessingService:
|
||||
)
|
||||
)
|
||||
|
||||
result = self.meeting_lab.regenerate_protocol(
|
||||
Path(run_dir),
|
||||
context,
|
||||
relay,
|
||||
model=self.settings.protocol_model,
|
||||
ollama_endpoint=self.settings.ollama_endpoint,
|
||||
protocol_num_ctx=self.settings.protocol_num_ctx,
|
||||
protocol_safe_input_token_budget=(self.settings.protocol_safe_input_token_budget),
|
||||
)
|
||||
self.timer.start_run()
|
||||
try:
|
||||
result = self.meeting_lab.regenerate_protocol(
|
||||
Path(run_dir),
|
||||
context,
|
||||
relay,
|
||||
model=self.settings.protocol_model,
|
||||
ollama_endpoint=self.settings.ollama_endpoint,
|
||||
protocol_num_ctx=self.settings.protocol_num_ctx,
|
||||
protocol_safe_input_token_budget=(self.settings.protocol_safe_input_token_budget),
|
||||
protocol_num_thread=resolve_ollama_num_thread(performance_profile),
|
||||
glossary_aliases=glossary_alias_mapping(self.glossary.list(active_only=True)),
|
||||
)
|
||||
except Exception:
|
||||
self.timer.finish_run()
|
||||
raise
|
||||
self.timer.finish_run()
|
||||
protocol_path = Path(result.protocol_path)
|
||||
return ProcessingOutcome(
|
||||
succeeded=True,
|
||||
@@ -463,3 +492,38 @@ class MeetingProcessingService:
|
||||
destination = Path(run_dir) / "protocol_edited.md"
|
||||
destination.write_text(text, encoding="utf-8")
|
||||
return destination
|
||||
|
||||
def _update_timer(self, event: Any) -> None:
|
||||
"""Apply a backend progress event to the shared per-run timer."""
|
||||
if event.stage in STAGES and event.status == "started":
|
||||
self.timer.start_stage(event.stage)
|
||||
elif event.stage in STAGES and event.status in {"completed", "skipped"}:
|
||||
self.timer.finish_stage(event.stage)
|
||||
if event.stage in {"completed", "failed"}:
|
||||
self.timer.finish_run()
|
||||
|
||||
def _persist_timing(self, run_dir: Path | None) -> None:
|
||||
"""Merge final monotonic durations into Meeting Lab run metadata."""
|
||||
if run_dir is None:
|
||||
return
|
||||
metadata_path = Path(run_dir) / "run_metadata.json"
|
||||
metadata: dict[str, Any] = {}
|
||||
if metadata_path.is_file():
|
||||
try:
|
||||
existing = json.loads(metadata_path.read_text(encoding="utf-8"))
|
||||
if isinstance(existing, dict):
|
||||
metadata = existing
|
||||
except (OSError, json.JSONDecodeError):
|
||||
return
|
||||
snapshot = self.timer.snapshot()
|
||||
metadata["timing"] = {
|
||||
"stages_seconds": snapshot.stage_durations,
|
||||
"total_seconds": snapshot.total_duration,
|
||||
}
|
||||
metadata_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
try:
|
||||
metadata_path.write_text(
|
||||
json.dumps(metadata, indent=2, sort_keys=True) + "\n", encoding="utf-8"
|
||||
)
|
||||
except OSError:
|
||||
return
|
||||
|
||||
@@ -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
|
||||
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from datetime import date
|
||||
from queue import Empty, Queue
|
||||
from typing import Any
|
||||
from uuid import uuid4
|
||||
|
||||
@@ -14,6 +18,11 @@ from mka.application.glossary import (
|
||||
GlossaryConflictError,
|
||||
GlossaryRepository,
|
||||
)
|
||||
from mka.application.glossary_yaml import (
|
||||
GlossaryYamlError,
|
||||
export_glossary_yaml,
|
||||
import_glossary_yaml,
|
||||
)
|
||||
from mka.application.meeting_service import (
|
||||
STAGES,
|
||||
AppProgressEvent,
|
||||
@@ -21,6 +30,7 @@ from mka.application.meeting_service import (
|
||||
MeetingProcessingService,
|
||||
ParticipantInput,
|
||||
ProcessingOptions,
|
||||
SpeakerMappingReview,
|
||||
stable_id,
|
||||
)
|
||||
from mka.application.people_yaml import (
|
||||
@@ -28,6 +38,8 @@ from mka.application.people_yaml import (
|
||||
export_people_yaml,
|
||||
import_people_yaml,
|
||||
)
|
||||
from mka.application.performance import DEFAULT_PERFORMANCE_PROFILE, PERFORMANCE_PROFILES
|
||||
from mka.application.progress_timing import TimingSnapshot
|
||||
from mka.application.run_inputs import (
|
||||
RunInputJsonError,
|
||||
RunInputState,
|
||||
@@ -62,6 +74,33 @@ def _render_glossary(repository: GlossaryRepository) -> None:
|
||||
"Store canonical core terms and recognition aliases. Compound phrases are "
|
||||
"composed from meeting context during protocol generation."
|
||||
)
|
||||
import_file = st.file_uploader(
|
||||
"Import glossary",
|
||||
type=["yaml", "yml"],
|
||||
key="glossary_import_file",
|
||||
help="Validate and atomically replace the SQLite glossary from versioned YAML.",
|
||||
)
|
||||
action_columns = st.columns(2)
|
||||
if action_columns[0].button("Import glossary", disabled=import_file is None):
|
||||
try:
|
||||
imported = import_glossary_yaml(import_file.getvalue())
|
||||
repository.replace_all(imported)
|
||||
except (GlossaryYamlError, GlossaryConflictError, ValueError) as exc:
|
||||
st.error(str(exc))
|
||||
else:
|
||||
st.session_state.glossary_import_message = (
|
||||
f"Imported {len(imported)} glossary entries."
|
||||
)
|
||||
st.rerun()
|
||||
action_columns[1].download_button(
|
||||
"Export glossary",
|
||||
data=export_glossary_yaml(repository.list()).encode("utf-8"),
|
||||
file_name="terminology-glossary.yaml",
|
||||
mime="application/yaml",
|
||||
)
|
||||
if message := st.session_state.pop("glossary_import_message", None):
|
||||
st.success(message)
|
||||
|
||||
with st.form("glossary_add"):
|
||||
columns = st.columns(2)
|
||||
canonical = columns[0].text_input("Canonical term")
|
||||
@@ -304,40 +343,199 @@ def _render_participants() -> list[ParticipantInput]:
|
||||
return people
|
||||
|
||||
|
||||
def _progress_callback(
|
||||
def _format_duration(seconds: float) -> str:
|
||||
"""Format numeric seconds as MM:SS or H:MM:SS."""
|
||||
whole_seconds = max(0, int(seconds))
|
||||
hours, remainder = divmod(whole_seconds, 3600)
|
||||
minutes, seconds = divmod(remainder, 60)
|
||||
return f"{hours}:{minutes:02d}:{seconds:02d}" if hours else f"{minutes:02d}:{seconds:02d}"
|
||||
|
||||
|
||||
def _render_progress_status(
|
||||
status_box: Any,
|
||||
stage_table: Any,
|
||||
progress_slot: Any,
|
||||
states: dict[str, str],
|
||||
) -> Any:
|
||||
progress_bar = None
|
||||
timing: TimingSnapshot,
|
||||
message: str,
|
||||
) -> None:
|
||||
suffix = " …" if timing.running else ""
|
||||
status_box.info(f"{message} — total runtime {_format_duration(timing.total_duration)}{suffix}")
|
||||
stage_table.table(
|
||||
[
|
||||
{
|
||||
"Stage": STAGE_LABELS[stage],
|
||||
"Status": states[stage],
|
||||
"Runtime": (
|
||||
_format_duration(timing.stage_durations[stage])
|
||||
+ (" …" if timing.active_stage == stage else "")
|
||||
if stage in timing.stage_durations
|
||||
else ""
|
||||
),
|
||||
}
|
||||
for stage in STAGES
|
||||
]
|
||||
+ [
|
||||
{
|
||||
"Stage": "Total runtime",
|
||||
"Status": "",
|
||||
"Runtime": _format_duration(timing.total_duration) + suffix,
|
||||
}
|
||||
]
|
||||
)
|
||||
|
||||
def update(event: AppProgressEvent) -> None:
|
||||
nonlocal progress_bar
|
||||
if event.stage in states:
|
||||
states[event.stage] = "running" if event.status == "started" else event.status
|
||||
if event.stage == "failed":
|
||||
running = next(
|
||||
|
||||
def _apply_progress_event(event: AppProgressEvent, states: dict[str, str]) -> None:
|
||||
if event.stage in states:
|
||||
states[event.stage] = "running" if event.status == "started" else event.status
|
||||
if event.stage == "failed":
|
||||
running = next((stage for stage, status in states.items() if status == "running"), None)
|
||||
if running:
|
||||
states[running] = "failed"
|
||||
|
||||
|
||||
def _run_with_live_progress(
|
||||
service: MeetingProcessingService,
|
||||
work: Callable[[Callable[[AppProgressEvent], None]], Any],
|
||||
states: dict[str, str],
|
||||
initial_message: str,
|
||||
) -> tuple[Any, dict[str, str], str, TimingSnapshot]:
|
||||
"""Run pipeline work while rendering the shared timer and progress events."""
|
||||
status_box = st.empty()
|
||||
stage_table = st.empty()
|
||||
progress_slot = st.empty()
|
||||
event_queue: Queue[AppProgressEvent] = Queue()
|
||||
latest_message = initial_message
|
||||
progress_bar = None
|
||||
with ThreadPoolExecutor(max_workers=1) as executor:
|
||||
future = executor.submit(work, event_queue.put)
|
||||
while not future.done():
|
||||
try:
|
||||
while True:
|
||||
event = event_queue.get_nowait()
|
||||
_apply_progress_event(event, states)
|
||||
latest_message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||
if event.progress is not None:
|
||||
if progress_bar is None:
|
||||
progress_bar = progress_slot.progress(0.0)
|
||||
progress_bar.progress(
|
||||
min(max(event.progress, 0.0), 1.0),
|
||||
text=(
|
||||
f"{STAGE_LABELS.get(event.stage, event.stage)}: "
|
||||
f"{event.progress:.0%}"
|
||||
),
|
||||
)
|
||||
except Empty:
|
||||
pass
|
||||
_render_progress_status(
|
||||
status_box,
|
||||
stage_table,
|
||||
states,
|
||||
service.timer.snapshot(),
|
||||
latest_message,
|
||||
)
|
||||
time.sleep(0.2)
|
||||
try:
|
||||
result = future.result()
|
||||
except Exception:
|
||||
while not event_queue.empty():
|
||||
event = event_queue.get_nowait()
|
||||
_apply_progress_event(event, states)
|
||||
latest_message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||
service.timer.finish_run()
|
||||
running_stage = next(
|
||||
(stage for stage, status in states.items() if status == "running"),
|
||||
None,
|
||||
)
|
||||
if running:
|
||||
states[running] = "failed"
|
||||
elapsed = f"{event.elapsed_seconds:.1f} s"
|
||||
message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||
status_box.info(f"{message} — elapsed {elapsed}")
|
||||
stage_table.table(
|
||||
[{"Stage": STAGE_LABELS[stage], "Status": states[stage]} for stage in STAGES]
|
||||
)
|
||||
if event.progress is not None:
|
||||
if progress_bar is None:
|
||||
progress_bar = progress_slot.progress(0.0)
|
||||
progress_bar.progress(
|
||||
min(max(event.progress, 0.0), 1.0),
|
||||
text=f"{STAGE_LABELS.get(event.stage, event.stage)}: {event.progress:.0%}",
|
||||
)
|
||||
if running_stage is not None:
|
||||
states[running_stage] = "failed"
|
||||
snapshot = service.timer.snapshot()
|
||||
_render_progress_status(status_box, stage_table, states, snapshot, latest_message)
|
||||
st.session_state["processing_progress"] = (states.copy(), latest_message, snapshot)
|
||||
raise
|
||||
while not event_queue.empty():
|
||||
event = event_queue.get_nowait()
|
||||
_apply_progress_event(event, states)
|
||||
latest_message = event.message or STAGE_LABELS.get(event.stage, event.stage)
|
||||
running_stage = next((stage for stage, status in states.items() if status == "running"), None)
|
||||
if running_stage is not None:
|
||||
states[running_stage] = "completed" if getattr(result, "succeeded", True) else "failed"
|
||||
snapshot = service.timer.snapshot()
|
||||
_render_progress_status(status_box, stage_table, states, snapshot, latest_message)
|
||||
st.session_state["processing_progress"] = (states.copy(), latest_message, snapshot)
|
||||
return result, states, latest_message, snapshot
|
||||
|
||||
return update
|
||||
|
||||
def _remember_regeneration_timing(
|
||||
states: dict[str, str], message: str, timing: TimingSnapshot
|
||||
) -> None:
|
||||
st.session_state["regeneration_progress"] = (states.copy(), message, timing)
|
||||
|
||||
|
||||
def _speaker_options(
|
||||
participants: tuple[str, ...], mappings: dict[str, str | None], speaker_label: str
|
||||
) -> list[str | None]:
|
||||
"""Reserve other speakers' participants while retaining this speaker's mapping."""
|
||||
current = mappings.get(speaker_label)
|
||||
reserved = {value for label, value in mappings.items() if label != speaker_label}
|
||||
options: list[str | None] = [None]
|
||||
options.extend(person for person in participants if person == current or person not in reserved)
|
||||
if current is not None and current not in options:
|
||||
options.append(current)
|
||||
return options
|
||||
|
||||
|
||||
def _mapping_counts(
|
||||
speaker_labels: tuple[str, ...], mappings: dict[str, str | None]
|
||||
) -> tuple[int, int, int]:
|
||||
"""Count only detected speakers, including explicit cleared selections."""
|
||||
detected = len(speaker_labels)
|
||||
assigned = sum(mappings.get(label) is not None for label in speaker_labels)
|
||||
return detected, assigned, detected - assigned
|
||||
|
||||
|
||||
def _render_speaker_mapping(review: SpeakerMappingReview, run_name: str) -> dict[str, str]:
|
||||
st.subheader("Identify diarized speakers")
|
||||
st.caption(
|
||||
"Confirm identities explicitly. Unmapped speakers remain anonymous; "
|
||||
"the diarized source transcript is not modified."
|
||||
)
|
||||
participant_names = dict(review.participants)
|
||||
# Read every widget before rendering so later speakers also reserve their person.
|
||||
keys = {
|
||||
speaker.speaker_label: f"speaker_mapping_{run_name}_{speaker.speaker_label}"
|
||||
for speaker in review.speakers
|
||||
}
|
||||
mappings = {
|
||||
label: st.session_state.get(key, review.current_mappings.get(label))
|
||||
for label, key in keys.items()
|
||||
}
|
||||
detected, assigned, unassigned = _mapping_counts(tuple(keys), mappings)
|
||||
st.markdown(f"**{detected} speakers detected · {assigned} assigned · {unassigned} unassigned**")
|
||||
if unassigned:
|
||||
st.warning(
|
||||
"Some detected speakers have no confirmed participant mapping. "
|
||||
"Check whether a participant is missing or speaker assignment is incomplete. "
|
||||
"You can still generate a protocol with anonymous speakers."
|
||||
)
|
||||
selections: dict[str, str] = {}
|
||||
for speaker in review.speakers:
|
||||
label = speaker.speaker_label
|
||||
current = mappings[label]
|
||||
options = _speaker_options(tuple(participant_names), mappings, label)
|
||||
st.session_state[keys[label]] = current
|
||||
selected = st.selectbox(
|
||||
f"{label} — Unassigned" if current is None else label,
|
||||
options=options,
|
||||
format_func=lambda value, names=participant_names: (
|
||||
"Unmapped / Unknown" if value is None else names.get(value, value)
|
||||
),
|
||||
key=keys[label],
|
||||
)
|
||||
if selected is not None:
|
||||
selections[label] = selected
|
||||
for excerpt in speaker.excerpts:
|
||||
st.caption(f"“{excerpt}”")
|
||||
return selections
|
||||
|
||||
|
||||
def _render_result() -> None:
|
||||
@@ -346,6 +544,9 @@ def _render_result() -> None:
|
||||
return
|
||||
st.divider()
|
||||
st.header("Protocol result")
|
||||
if retained_progress := st.session_state.get("processing_progress"):
|
||||
states, message, timing = retained_progress
|
||||
_render_progress_status(st.empty(), st.empty(), states, timing, message)
|
||||
if not outcome.succeeded:
|
||||
st.error(f"Processing failed during {outcome.failed_stage}: {outcome.error_message}")
|
||||
if outcome.run_dir:
|
||||
@@ -383,29 +584,7 @@ def _render_result() -> None:
|
||||
if review is None or not review.speakers:
|
||||
return
|
||||
|
||||
st.subheader("Identify diarized speakers")
|
||||
st.caption(
|
||||
"Confirm identities explicitly. Unmapped speakers remain anonymous; "
|
||||
"the diarized source transcript is not modified."
|
||||
)
|
||||
participant_names = dict(review.participants)
|
||||
options = [None, *participant_names]
|
||||
selections: dict[str, str] = {}
|
||||
for speaker in review.speakers:
|
||||
current = review.current_mappings.get(speaker.speaker_label)
|
||||
selected = st.selectbox(
|
||||
speaker.speaker_label,
|
||||
options=options,
|
||||
index=options.index(current) if current in options else 0,
|
||||
format_func=lambda value, names=participant_names: (
|
||||
"Unmapped / Unknown" if value is None else names[value]
|
||||
),
|
||||
key=f"speaker_mapping_{outcome.run_dir.name}_{speaker.speaker_label}",
|
||||
)
|
||||
if selected is not None:
|
||||
selections[speaker.speaker_label] = selected
|
||||
for excerpt in speaker.excerpts:
|
||||
st.caption(f"“{excerpt}”")
|
||||
selections = _render_speaker_mapping(review, outcome.run_dir.name)
|
||||
|
||||
duplicate_assignments = len(selections.values()) != len(set(selections.values()))
|
||||
if duplicate_assignments:
|
||||
@@ -415,12 +594,33 @@ def _render_result() -> None:
|
||||
disabled=duplicate_assignments,
|
||||
type="primary",
|
||||
):
|
||||
st.session_state.pop("regeneration_progress", None)
|
||||
run_dir = outcome.run_dir
|
||||
mapping_items = tuple(selections.items())
|
||||
selected_profile = st.session_state.get("performance_profile", DEFAULT_PERFORMANCE_PROFILE)
|
||||
states = {stage: "skipped" for stage in STAGES}
|
||||
states["protocol_generation"] = "pending"
|
||||
try:
|
||||
with st.spinner("Regenerating protocol without rerunning audio processing..."):
|
||||
regenerated = service.regenerate_protocol(outcome.run_dir, selections)
|
||||
regenerated, states, message, timing = _run_with_live_progress(
|
||||
service,
|
||||
lambda progress_sink: service.regenerate_protocol(
|
||||
run_dir,
|
||||
dict(mapping_items),
|
||||
progress_sink=progress_sink,
|
||||
performance_profile=selected_profile,
|
||||
),
|
||||
states,
|
||||
"Starting protocol regeneration",
|
||||
)
|
||||
except (OSError, RuntimeError, ValueError) as exc:
|
||||
_remember_regeneration_timing(
|
||||
states,
|
||||
"Protocol regeneration failed",
|
||||
service.timer.snapshot(),
|
||||
)
|
||||
st.error(f"Protocol regeneration failed: {exc}")
|
||||
else:
|
||||
_remember_regeneration_timing(states, message, timing)
|
||||
st.session_state.outcome = regenerated
|
||||
_queue_edited_protocol(regenerated.original_protocol or "")
|
||||
st.session_state.speaker_mapping_message = (
|
||||
@@ -458,7 +658,10 @@ def main() -> None:
|
||||
description = st.text_area("Description / context", height=100, key="meeting_description")
|
||||
metadata_columns = st.columns(3)
|
||||
language = metadata_columns[0].selectbox(
|
||||
"Meeting language", options=["de", "en"], key="meeting_language"
|
||||
"Meeting language",
|
||||
options=["de", "en"],
|
||||
key="meeting_language",
|
||||
help="Language for both transcription and the generated protocol.",
|
||||
)
|
||||
has_date = metadata_columns[1].checkbox(
|
||||
"Meeting date is known", value=True, key="meeting_has_date"
|
||||
@@ -470,6 +673,13 @@ def main() -> None:
|
||||
participants = _render_participants()
|
||||
|
||||
st.header("Processing options")
|
||||
performance_profile = st.selectbox(
|
||||
"Performance profile",
|
||||
options=PERFORMANCE_PROFILES,
|
||||
format_func=str.title,
|
||||
key="performance_profile",
|
||||
help="Controls protocol-generation performance using an abstract runtime profile.",
|
||||
)
|
||||
audio_normalization = st.checkbox(
|
||||
"Audio normalization",
|
||||
value=True,
|
||||
@@ -523,30 +733,34 @@ def main() -> None:
|
||||
)
|
||||
meeting_id = stable_id(title)
|
||||
audio_path = service.preserve_upload(meeting_id, audio.name, audio)
|
||||
st.session_state.pop("regeneration_progress", None)
|
||||
st.session_state.pop("processing_progress", None)
|
||||
st.header("Processing status")
|
||||
status_box = st.empty()
|
||||
stage_table = st.empty()
|
||||
progress_slot = st.empty()
|
||||
states = {stage: "pending" for stage in STAGES}
|
||||
if not diarization_enabled:
|
||||
states["diarization"] = "skipped"
|
||||
callback = _progress_callback(status_box, stage_table, progress_slot, states)
|
||||
outcome = service.process(
|
||||
audio_path,
|
||||
meeting,
|
||||
participants,
|
||||
ProcessingOptions(
|
||||
diarization_enabled=diarization_enabled,
|
||||
audio_normalization=audio_normalization,
|
||||
outcome, _, _, _ = _run_with_live_progress(
|
||||
service,
|
||||
lambda progress_sink: service.process(
|
||||
audio_path,
|
||||
meeting,
|
||||
participants,
|
||||
ProcessingOptions(
|
||||
diarization_enabled=diarization_enabled,
|
||||
audio_normalization=audio_normalization,
|
||||
performance_profile=performance_profile,
|
||||
),
|
||||
progress_sink=progress_sink,
|
||||
),
|
||||
progress_sink=callback,
|
||||
states,
|
||||
"Starting processing",
|
||||
)
|
||||
st.session_state.outcome = outcome
|
||||
st.session_state.edited_protocol = outcome.original_protocol or ""
|
||||
if outcome.succeeded:
|
||||
status_box.success("Processing completed.")
|
||||
st.success("Processing completed.")
|
||||
else:
|
||||
status_box.error(
|
||||
st.error(
|
||||
f"Processing failed during {outcome.failed_stage}: "
|
||||
f"{outcome.error_message} Artifacts were preserved."
|
||||
)
|
||||
|
||||
@@ -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,
|
||||
ProcessingOptions,
|
||||
)
|
||||
from mka.application.progress_timing import ProcessingTimer
|
||||
|
||||
|
||||
@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(
|
||||
tmp_path: Path,
|
||||
language: str,
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(
|
||||
audio, meeting(), participants(), ProcessingOptions(diarization_enabled=False)
|
||||
audio,
|
||||
replace(meeting(), language=language),
|
||||
participants(),
|
||||
ProcessingOptions(diarization_enabled=False),
|
||||
)
|
||||
|
||||
assert outcome.succeeded
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["language"] == language
|
||||
assert service.build_context(replace(meeting(), language=language), participants()).data[
|
||||
"meeting"
|
||||
]["language"] == language
|
||||
assert gateway.config_values["diarization"] == "off"
|
||||
assert gateway.config_values["model"] == "test:model"
|
||||
assert gateway.config_values["protocol_num_ctx"] == 32_768
|
||||
@@ -503,6 +513,7 @@ def test_protocol_only_regeneration_persists_mappings_without_rewriting_transcri
|
||||
"SPEAKER_00": "martin"
|
||||
}
|
||||
assert gateway.regeneration["options"]["protocol_num_ctx"] == 32_768
|
||||
assert gateway.regeneration["options"]["glossary_aliases"] == {"Enlyse": "ENLYZE"}
|
||||
assert events[0].stage == "protocol_generation"
|
||||
assert source.read_bytes() == original_source
|
||||
|
||||
@@ -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.name.endswith("_unsafe.wav")
|
||||
assert destination.read_bytes() == b"source audio"
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("profile", "threads"),
|
||||
[("fast", 16), ("efficient", 10), ("powersave", 4)],
|
||||
)
|
||||
def test_process_resolves_profile_for_protocol_generation(
|
||||
tmp_path: Path, profile: str, threads: int
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
service.process(
|
||||
audio,
|
||||
meeting(),
|
||||
participants(),
|
||||
ProcessingOptions(performance_profile=profile),
|
||||
)
|
||||
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["protocol_num_thread"] == threads
|
||||
|
||||
|
||||
def test_process_defaults_to_backend_thread_selection(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["protocol_num_thread"] is None
|
||||
|
||||
|
||||
def test_protocol_regeneration_propagates_explicit_profile(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
|
||||
service.regenerate_protocol(
|
||||
gateway.run_dir,
|
||||
{},
|
||||
performance_profile="fast",
|
||||
)
|
||||
|
||||
assert gateway.regeneration is not None
|
||||
assert gateway.regeneration["options"]["protocol_num_thread"] == 16
|
||||
|
||||
|
||||
def test_process_passes_only_active_explicit_glossary_aliases(tmp_path: Path) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
service.glossary.create("Luminy", "product", aliases=("Lumini",))
|
||||
inactive = service.glossary.create("PBAT", "acronym", aliases=("PBRT",))
|
||||
service.glossary.set_active(inactive.id, False)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
assert gateway.config_values is not None
|
||||
assert gateway.config_values["glossary_aliases"] == {"Lumini": "Luminy"}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("fails", [False, True])
|
||||
def test_process_persists_completed_and_failed_timings(tmp_path: Path, fails: bool) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
gateway.fail = fails
|
||||
values = iter([0.0, 1.0, 9.0, 10.0, 20.0])
|
||||
service.timer = ProcessingTimer(lambda: next(values))
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
|
||||
outcome = service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
metadata = json.loads((gateway.run_dir / "run_metadata.json").read_text(encoding="utf-8"))
|
||||
assert metadata["timing"] == {
|
||||
"stages_seconds": {"preparing": 8.0, "transcription": 10.0},
|
||||
"total_seconds": 20.0,
|
||||
}
|
||||
assert outcome.succeeded is not fails
|
||||
if fails:
|
||||
assert metadata["failure"]["type"] == "TranscriptionError"
|
||||
|
||||
|
||||
def test_regeneration_starts_fresh_timer_and_retains_final_duration(
|
||||
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
clock = SimpleNamespace(value=10.0)
|
||||
service.timer = ProcessingTimer(lambda: clock.value)
|
||||
audio = tmp_path / "meeting.wav"
|
||||
audio.write_bytes(b"audio")
|
||||
service.process(audio, meeting(), participants(), ProcessingOptions())
|
||||
|
||||
clock.value = 100.0
|
||||
observed = []
|
||||
|
||||
def regenerate(run_dir: Path, context: Any, progress_sink: Any, **options: Any) -> Any:
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="protocol_generation",
|
||||
status="started",
|
||||
elapsed_seconds=0.0,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
clock.value = 105.0
|
||||
observed.append(service.timer.snapshot())
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="protocol_generation",
|
||||
status="completed",
|
||||
elapsed_seconds=5.0,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
clock.value = 109.0
|
||||
protocol = run_dir / "protocol.md"
|
||||
protocol.write_text("# Regenerated protocol\n", encoding="utf-8")
|
||||
return SimpleNamespace(run_dir=run_dir, protocol_path=protocol)
|
||||
|
||||
monkeypatch.setattr(gateway, "regenerate_protocol", regenerate)
|
||||
|
||||
service.regenerate_protocol(gateway.run_dir, {"SPEAKER_00": "martin"})
|
||||
final = service.timer.snapshot()
|
||||
|
||||
assert observed[0].active_stage == "protocol_generation"
|
||||
assert observed[0].stage_durations == {"protocol_generation": 5.0}
|
||||
assert observed[0].total_duration == 5.0
|
||||
assert observed[0].running is True
|
||||
assert final.stage_durations == {"protocol_generation": 5.0}
|
||||
assert final.total_duration == 9.0
|
||||
assert final.running is False
|
||||
|
||||
|
||||
def test_failed_regeneration_freezes_active_and_total_durations(
|
||||
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||
) -> None:
|
||||
service, gateway = make_service(tmp_path)
|
||||
write_speaker_review_artifacts(gateway.run_dir)
|
||||
clock = SimpleNamespace(value=20.0)
|
||||
service.timer = ProcessingTimer(lambda: clock.value)
|
||||
|
||||
def fail_regeneration(run_dir: Path, context: Any, progress_sink: Any, **options: Any) -> Any:
|
||||
progress_sink(
|
||||
SimpleNamespace(
|
||||
stage="protocol_generation",
|
||||
status="started",
|
||||
elapsed_seconds=0.0,
|
||||
progress=None,
|
||||
message=None,
|
||||
)
|
||||
)
|
||||
clock.value = 27.0
|
||||
raise RuntimeError("generation failed")
|
||||
|
||||
monkeypatch.setattr(gateway, "regenerate_protocol", fail_regeneration)
|
||||
|
||||
with pytest.raises(RuntimeError, match="generation failed"):
|
||||
service.regenerate_protocol(gateway.run_dir, {})
|
||||
clock.value = 99.0
|
||||
final = service.timer.snapshot()
|
||||
|
||||
assert final.stage_durations == {"protocol_generation": 7.0}
|
||||
assert final.total_duration == 7.0
|
||||
assert final.running is False
|
||||
|
||||
@@ -0,0 +1,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 mka.application.meeting_service import ParticipantInput
|
||||
from mka.application.progress_timing import TimingSnapshot
|
||||
from mka.application.run_inputs import RunInputState
|
||||
from mka.ui import streamlit_app
|
||||
|
||||
@@ -66,3 +67,29 @@ def test_imported_inputs_are_applied_via_pending_state_before_widgets(
|
||||
assert "source_media_1" not in state
|
||||
assert "pending_run_inputs" not in state
|
||||
assert state["source_media_0"] is existing_upload
|
||||
|
||||
|
||||
def test_duration_formatting_covers_seconds_minutes_and_hours() -> None:
|
||||
assert streamlit_app._format_duration(8.9) == "00:08"
|
||||
assert streamlit_app._format_duration(7 * 60 + 18) == "07:18"
|
||||
assert streamlit_app._format_duration(3600 + 3 * 60 + 42) == "1:03:42"
|
||||
|
||||
|
||||
def test_completed_regeneration_timing_survives_frontend_rerun(monkeypatch) -> None:
|
||||
state = {}
|
||||
monkeypatch.setattr(streamlit_app.st, "session_state", state)
|
||||
timing = TimingSnapshot(
|
||||
stage_durations={"protocol_generation": 4.5},
|
||||
active_stage=None,
|
||||
total_duration=4.5,
|
||||
running=False,
|
||||
)
|
||||
states = {stage: "skipped" for stage in streamlit_app.STAGES}
|
||||
states["protocol_generation"] = "completed"
|
||||
|
||||
streamlit_app._remember_regeneration_timing(states, "Protocol regenerated", timing)
|
||||
|
||||
saved_states, saved_message, saved_timing = state["regeneration_progress"]
|
||||
assert saved_states["protocol_generation"] == "completed"
|
||||
assert saved_message == "Protocol regenerated"
|
||||
assert saved_timing == timing
|
||||
|
||||
Reference in New Issue
Block a user