2 Commits
Author SHA1 Message Date
admin 9ce4af69be Improve speaker-review excerpts 2026-09-17 19:42:35 +02:00
admin 2215414282 Add post-diarization speaker review workflow 2026-09-17 19:16:57 +02:00
8 changed files with 542 additions and 40 deletions
+6 -2
View File
@@ -81,7 +81,8 @@ Imported audio
(optionally normalize loudness; default on) (optionally normalize loudness; default on)
-> transcribe with whisper.cpp and large-v3-turbo -> transcribe with whisper.cpp and large-v3-turbo
-> optionally diarize with pyannote.audio Community-1 -> optionally diarize with pyannote.audio Community-1
-> generate a protocol directly from the full transcript -> when diarized, wait for optional speaker review before explicit protocol generation
-> otherwise generate a protocol directly from the full transcript
-> review and edit by a human -> review and edit by a human
``` ```
@@ -121,7 +122,10 @@ Speaker selectors filter participants using a snapshot of all current widget
selections, falling back to saved mappings before first interaction. Filtering selections, falling back to saved mappings before first interaction. Filtering
retains existing assignments; clearing a selection releases the participant. retains existing assignments; clearing a selection releases the participant.
Completeness counts and warnings cover detected speakers only and do not gate Completeness counts and warnings cover detected speakers only and do not gate
protocol regeneration or change mapping persistence. protocol generation or change mapping persistence. A diarized run persists the
`awaiting_speaker_review` checkpoint after diarization; an empty, partial or
complete mapping can then explicitly generate a protocol from the existing
artifacts. Failed generation preserves that checkpoint and saved mappings for retry.
## Progress Contract ## Progress Contract
+18 -1
View File
@@ -35,7 +35,8 @@ Audio file
-> FFmpeg preparation (mono, 16 kHz PCM WAV; normalization optional) -> FFmpeg preparation (mono, 16 kHz PCM WAV; normalization optional)
-> whisper.cpp transcription (large-v3-turbo) -> whisper.cpp transcription (large-v3-turbo)
-> optional pyannote.audio Community-1 diarization -> optional pyannote.audio Community-1 diarization
-> direct full-transcript protocol generation -> optional speaker review and explicit protocol generation (when diarized)
or direct full-transcript protocol generation (without diarization)
-> human review and export -> human review and export
``` ```
@@ -212,6 +213,22 @@ From the Meeting Assistant repository, start the UI with:
PYTHONPATH=src:../meeting-lab .venv/bin/streamlit run src/mka/ui/streamlit_app.py PYTHONPATH=src:../meeting-lab .venv/bin/streamlit run src/mka/ui/streamlit_app.py
``` ```
### North launcher
In the post-diarization Assistant worktree on North, use the checked-in
launcher instead of creating another virtual environment:
```bash
./run_north.sh
```
It reuses North's validated Assistant virtual environment and external
Whisper/Ollama/Docker-ROCm runtime while loading Meeting Assistant from this
worktree and Meeting Lab from
`/opt/git-projekts/meeting-lab`. `HF_TOKEN`, when present, is
passed through to the disposable diarization container without being stored by
the launcher.
Streamlit opens `http://localhost:8501` by default. Installing the sibling Streamlit opens `http://localhost:8501` by default. Installing the sibling
Meeting Lab project supplies its runtime requirements such as PyYAML and Meeting Lab project supplies its runtime requirements such as PyYAML and
Requests. Optional diarization dependencies are needed only when diarization Requests. Optional diarization dependencies are needed only when diarization
Executable
+57
View File
@@ -0,0 +1,57 @@
#!/usr/bin/env bash
set -euo pipefail
assistant_root="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd -P)"
assistant_src="$assistant_root/src"
lab_root=/opt/git-projekts/meeting-lab
north_venv=/opt/git-projekts/meeting-assistant/.venv
whisper_executable=/home/martin/whisper.cpp/whisper.cpp/build/bin/whisper-cli
whisper_model=/home/martin/whisper.cpp/whisper.cpp/models/ggml-large-v3-turbo.bin
require_directory() {
if [[ ! -d "$1" ]]; then
printf 'Required directory is missing: %s\n' "$1" >&2
exit 1
fi
}
require_file() {
if [[ ! -f "$1" ]]; then
printf 'Required file is missing: %s\n' "$1" >&2
exit 1
fi
}
require_executable() {
if [[ ! -x "$1" ]]; then
printf 'Required executable is missing or not executable: %s\n' "$1" >&2
exit 1
fi
}
require_directory "$assistant_src"
require_file "$assistant_src/mka/ui/streamlit_app.py"
require_directory "$lab_root"
require_file "$lab_root/src/meeting_lab/__init__.py"
require_directory "$north_venv"
require_executable "$north_venv/bin/streamlit"
require_executable "$whisper_executable"
require_file "$whisper_model"
if ! command -v docker >/dev/null 2>&1; then
printf 'Required executable is not on PATH: docker\n' >&2
exit 1
fi
export MKA_WHISPER_EXECUTABLE="$whisper_executable"
export MKA_WHISPER_MODEL="$whisper_model"
export MKA_PROTOCOL_MODEL=qwen3.8:27b
export MKA_OLLAMA_ENDPOINT=http://127.0.0.1:11434
export MKA_DIARIZATION_MODE=auto
export MKA_DIARIZATION_RUNTIME=container
export MKA_DIARIZATION_CONTAINER_IMAGE=rocm/pytorch:rocm7.2.1_ubuntu24.04_py3.12_pytorch_release_2.9.1
export MKA_DIARIZATION_CONTAINER_ARGS='["--device=/dev/kfd","--device=/dev/dri","--group-add","video"]'
# Keep feature sources ahead of the reused venv's editable installations.
export PYTHONPATH="$assistant_src:$lab_root"
exec "$north_venv/bin/streamlit" run "$assistant_src/mka/ui/streamlit_app.py"
+239 -16
View File
@@ -24,6 +24,37 @@ from mka.application.performance import DEFAULT_PERFORMANCE_PROFILE, resolve_oll
from mka.application.progress_timing import ProcessingTimer from mka.application.progress_timing import ProcessingTimer
STAGES = ("preparing", "transcription", "diarization", "protocol_generation") STAGES = ("preparing", "transcription", "diarization", "protocol_generation")
_EXCERPT_TOKEN_PATTERN = re.compile(r"[^\W\d_]+|\d+", re.UNICODE)
_EXCERPT_STOP_WORDS = frozenset(
{
"a", "an", "and", "are", "das", "der", "die", "ein", "eine", "for", "i",
"ich", "im", "in", "is", "it", "ja", "mit", "of", "or", "the", "to", "und",
"we", "wir", "with",
}
)
_ACKNOWLEDGEMENT_WORDS = frozenset(
{
"absolutely", "alles", "danke", "genau", "gut", "ja", "klar", "okay", "ok",
"perfekt", "right", "sure", "thanks", "verstanden", "yes",
}
)
_FIRST_PERSON_WORDS = frozenset(
{"i", "ich", "me", "mein", "meine", "my", "unser", "unsere", "we", "wir"}
)
_ACTION_WORDS = frozenset(
{
"agree", "approve", "believe", "can", "decide", "habe", "kann", "möchte", "plan",
"recommend", "think", "übernehme", "werde", "will",
}
)
_DATE_WORDS = frozenset(
{
"april", "august", "december", "dienstag", "donnerstag", "february", "freitag",
"friday", "january", "juli", "june", "mai", "march", "monday", "mittwoch", "montag",
"november", "october", "saturday", "september", "sonntag", "sunday", "thursday",
"tuesday", "wednesday",
}
)
class MeetingLabPort(Protocol): class MeetingLabPort(Protocol):
@@ -92,6 +123,8 @@ class ProcessingOutcome:
failed_stage: str | None = None failed_stage: str | None = None
error_message: str | None = None error_message: str | None = None
speaker_attribution_available: bool | None = None speaker_attribution_available: bool | None = None
awaiting_speaker_review: bool = False
detected_speaker_count: int | None = None
@dataclass(frozen=True) @dataclass(frozen=True)
@@ -107,6 +140,89 @@ class SpeakerMappingReview:
current_mappings: dict[str, str] current_mappings: dict[str, str]
@dataclass(frozen=True)
class _ExcerptCandidate:
speaker_label: str
text: str
segment_index: int
tokens: frozenset[str]
lexical_tokens: frozenset[str]
previous_other_speaker_tokens: frozenset[str]
def _excerpt_tokens(text: str) -> tuple[str, ...]:
return tuple(token.lower() for token in _EXCERPT_TOKEN_PATTERN.findall(text))
def _excerpt_score(candidate: _ExcerptCandidate, distinctive_token_count: int) -> float:
"""Score an excerpt using deterministic lexical signals only."""
word_count = len(candidate.tokens)
lexical_count = len(candidate.lexical_tokens)
score = min(word_count, 24) * 0.3 + min(lexical_count, 16) * 0.6
if word_count < 5:
score -= 8.0
elif word_count < 9:
score -= 3.0
if lexical_count <= 2:
score -= 4.0
if candidate.tokens and candidate.tokens <= _ACKNOWLEDGEMENT_WORDS:
score -= 14.0
if candidate.text.rstrip().endswith("?"):
score -= 5.0
if any(token.isdigit() for token in candidate.tokens):
score += 2.5
if candidate.tokens & _DATE_WORDS:
score += 2.0
if candidate.tokens & _FIRST_PERSON_WORDS:
score += 1.5
if candidate.tokens & _ACTION_WORDS:
score += 2.0
proper_name_count = len(re.findall(r"\b[A-ZÄÖÜ][a-zäöüß]{2,}\b", candidate.text))
score += min(proper_name_count, 3) * 0.75
score += min(sum(len(token) >= 10 for token in candidate.lexical_tokens), 4) * 0.4
score += min(distinctive_token_count, 6) * 0.65
if candidate.previous_other_speaker_tokens and candidate.lexical_tokens:
overlap = len(candidate.lexical_tokens & candidate.previous_other_speaker_tokens)
if overlap / len(candidate.lexical_tokens) >= 0.7:
score -= 5.0
return score
def _segment_bucket(segment_index: int, segment_count: int) -> int:
return min(5, segment_index * 6 // max(segment_count, 1))
def _selection_score(
candidate: _ExcerptCandidate,
base_score: float,
selected: list[_ExcerptCandidate],
selected_buckets: set[int],
segment_count: int,
) -> float:
"""Favor excerpts from new meeting positions and avoid near duplicates."""
score = base_score
bucket = _segment_bucket(candidate.segment_index, segment_count)
if bucket not in selected_buckets:
score += 2.5
if not selected or not candidate.lexical_tokens:
return score
similarities = [
len(candidate.lexical_tokens & item.lexical_tokens)
/ len(candidate.lexical_tokens | item.lexical_tokens)
for item in selected
if item.lexical_tokens
]
maximum_similarity = max(similarities, default=0.0)
if maximum_similarity >= 0.72:
score -= 18.0
elif maximum_similarity >= 0.55:
score -= 6.0
nearest_position = min(abs(candidate.segment_index - item.segment_index) for item in selected)
if nearest_position / max(segment_count - 1, 1) < 0.12:
score -= 1.5
return score
def stable_id(value: str) -> str: def stable_id(value: str) -> str:
"""Return a schema-safe ID based on user text, with a random fallback.""" """Return a schema-safe ID based on user text, with a random fallback."""
normalized = re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") normalized = re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-")
@@ -249,6 +365,7 @@ class MeetingProcessingService:
"diarization_runtime": self.settings.diarization_runtime, "diarization_runtime": self.settings.diarization_runtime,
"diarization_container_image": (self.settings.diarization_container_image), "diarization_container_image": (self.settings.diarization_container_image),
"diarization_container_args": self.settings.diarization_container_args, "diarization_container_args": self.settings.diarization_container_args,
"stop_after_diarization": options.diarization_enabled,
} }
) )
current_stage: str | None = None current_stage: str | None = None
@@ -299,6 +416,15 @@ class MeetingProcessingService:
failed_stage=failed_stage or current_stage or "preparing", failed_stage=failed_stage or current_stage or "preparing",
error_message=failure_message, error_message=failure_message,
) )
if result.protocol_path is None:
return ProcessingOutcome(
succeeded=True,
run_dir=result.run_dir,
original_protocol=None,
protocol_path=None,
awaiting_speaker_review=True,
detected_speaker_count=self._detected_speaker_count(result.run_dir),
)
protocol_path = Path(result.protocol_path) protocol_path = Path(result.protocol_path)
return ProcessingOutcome( return ProcessingOutcome(
succeeded=True, succeeded=True,
@@ -312,7 +438,7 @@ class MeetingProcessingService:
self, self,
run_dir: Path, run_dir: Path,
*, *,
excerpts_per_speaker: int = 3, excerpts_per_speaker: int = 6,
minimum_excerpt_characters: int = 20, minimum_excerpt_characters: int = 20,
) -> SpeakerMappingReview | None: ) -> SpeakerMappingReview | None:
"""Load detected labels and small deterministic identification excerpts.""" """Load detected labels and small deterministic identification excerpts."""
@@ -329,34 +455,67 @@ class MeetingProcessingService:
if not isinstance(segments, list): if not isinstance(segments, list):
raise ValueError("Diarized transcript has no segments list.") raise ValueError("Diarized transcript has no segments list.")
texts_by_speaker: dict[str, list[str]] = {} candidates_by_speaker: dict[str, list[_ExcerptCandidate]] = {}
short_by_speaker: dict[str, list[str]] = {} short_candidates_by_speaker: dict[str, list[_ExcerptCandidate]] = {}
for segment in segments: previous_label: str | None = None
previous_tokens: frozenset[str] = frozenset()
for segment_index, segment in enumerate(segments):
if not isinstance(segment, dict): if not isinstance(segment, dict):
continue continue
label = segment.get("speaker_id") label = segment.get("speaker_id")
text = segment.get("text") text = segment.get("text")
normalized = " ".join(text.split()) if isinstance(text, str) else ""
tokens = frozenset(_excerpt_tokens(normalized))
if ( if (
not isinstance(label, str) not isinstance(label, str)
or not label.startswith("SPEAKER_") or not label.startswith("SPEAKER_")
or label == "SPEAKER_UNASSIGNED" or label == "SPEAKER_UNASSIGNED"
or not isinstance(text, str) or not normalized
or not text.strip()
): ):
if normalized:
previous_label = label if isinstance(label, str) else None
previous_tokens = tokens
continue continue
normalized = " ".join(text.split())
target = (
texts_by_speaker
if len(normalized) >= minimum_excerpt_characters
else short_by_speaker
)
target.setdefault(label, []).append(normalized)
labels = sorted(set(texts_by_speaker) | set(short_by_speaker)) candidate = _ExcerptCandidate(
speaker_label=label,
text=normalized,
segment_index=segment_index,
tokens=tokens,
lexical_tokens=tokens - _EXCERPT_STOP_WORDS,
previous_other_speaker_tokens=(
previous_tokens
if previous_label is not None and previous_label != label
else frozenset()
),
)
target = (
candidates_by_speaker
if len(normalized) >= minimum_excerpt_characters
else short_candidates_by_speaker
)
target.setdefault(label, []).append(candidate)
previous_label = label
previous_tokens = tokens
labels = sorted(set(candidates_by_speaker) | set(short_candidates_by_speaker))
token_speakers: dict[str, set[str]] = {}
for candidate_group in (*candidates_by_speaker.values(), *short_candidates_by_speaker.values()):
for candidate in candidate_group:
for token in candidate.lexical_tokens:
token_speakers.setdefault(token, set()).add(candidate.speaker_label)
speakers = [] speakers = []
for label in labels: for label in labels:
candidates = texts_by_speaker.get(label) or short_by_speaker.get(label, []) candidates = candidates_by_speaker.get(label) or short_candidates_by_speaker.get(
speakers.append(SpeakerReview(label, tuple(candidates[:excerpts_per_speaker]))) label, []
)
excerpts = self._select_identifying_excerpts(
candidates,
excerpts_per_speaker=excerpts_per_speaker,
segment_count=len(segments),
token_speakers=token_speakers,
)
speakers.append(SpeakerReview(label, excerpts))
participants = tuple( participants = tuple(
(participant["participant_id"], participant["display_name"]) (participant["participant_id"], participant["display_name"])
for participant in context_data.get("participants", []) for participant in context_data.get("participants", [])
@@ -371,6 +530,50 @@ class MeetingProcessingService:
current_mappings=dict(mappings) if isinstance(mappings, dict) else {}, current_mappings=dict(mappings) if isinstance(mappings, dict) else {},
) )
@staticmethod
def _select_identifying_excerpts(
candidates: list[_ExcerptCandidate],
*,
excerpts_per_speaker: int,
segment_count: int,
token_speakers: dict[str, set[str]],
) -> tuple[str, ...]:
"""Choose informative, distinct excerpts while retaining transcript order."""
if len(candidates) <= excerpts_per_speaker:
return tuple(candidate.text for candidate in candidates)
base_scores = {
candidate: _excerpt_score(
candidate,
sum(
len(token_speakers.get(token, set())) == 1
for token in candidate.lexical_tokens
),
)
for candidate in candidates
}
selected: list[_ExcerptCandidate] = []
selected_buckets: set[int] = set()
while len(selected) < excerpts_per_speaker:
choice = max(
(candidate for candidate in candidates if candidate not in selected),
key=lambda candidate: (
_selection_score(
candidate,
base_scores[candidate],
selected,
selected_buckets,
segment_count,
),
-candidate.segment_index,
),
)
selected.append(choice)
selected_buckets.add(_segment_bucket(choice.segment_index, segment_count))
return tuple(
candidate.text for candidate in sorted(selected, key=lambda item: item.segment_index)
)
def regenerate_protocol( def regenerate_protocol(
self, self,
run_dir: Path, run_dir: Path,
@@ -401,6 +604,13 @@ class MeetingProcessingService:
context_data["speaker_mappings"] = dict(sorted(speaker_mappings.items())) context_data["speaker_mappings"] = dict(sorted(speaker_mappings.items()))
self._apply_glossary(context_data) self._apply_glossary(context_data)
context = self.meeting_lab.create_context(context_data) context = self.meeting_lab.create_context(context_data)
# A checkpoint mapping is user input, not a derived protocol artifact. Save it
# before inference so a failed generation remains retryable with the same review.
if context_path.is_symlink():
context_path.unlink()
context_path.write_text(
yaml.safe_dump(context.data, allow_unicode=True, sort_keys=False), encoding="utf-8"
)
def relay(event: Any) -> None: def relay(event: Any) -> None:
self._update_timer(event) self._update_timer(event)
@@ -432,6 +642,7 @@ class MeetingProcessingService:
self.timer.finish_run() self.timer.finish_run()
raise raise
self.timer.finish_run() self.timer.finish_run()
self._persist_timing(Path(result.run_dir))
protocol_path = Path(result.protocol_path) protocol_path = Path(result.protocol_path)
return ProcessingOutcome( return ProcessingOutcome(
succeeded=True, succeeded=True,
@@ -473,6 +684,18 @@ class MeetingProcessingService:
value = metadata.get("speaker_attribution_available") value = metadata.get("speaker_attribution_available")
return value if isinstance(value, bool) else None return value if isinstance(value, bool) else None
@staticmethod
def _detected_speaker_count(run_dir: Path | None) -> int | None:
if run_dir is None:
return None
metadata_path = Path(run_dir) / "run_metadata.json"
try:
metadata = json.loads(metadata_path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return None
value = metadata.get("diarization", {}).get("speaker_count")
return value if isinstance(value, int) else None
@staticmethod @staticmethod
def _read_failure(run_dir: Path | None) -> dict[str, str]: def _read_failure(run_dir: Path | None) -> dict[str, str]:
if run_dir is None: if run_dir is None:
+13 -4
View File
@@ -554,6 +554,14 @@ def _render_result() -> None:
st.caption("Intermediate artifacts and run metadata were preserved here.") st.caption("Intermediate artifacts and run metadata were preserved here.")
return return
awaiting_review = getattr(outcome, "awaiting_speaker_review", False)
if awaiting_review:
count = getattr(outcome, "detected_speaker_count", None)
detected = f"{count} speakers detected" if count is not None else "speakers detected"
st.info(
f"Diarization complete — {detected}. Assign speakers if desired, then generate the protocol."
)
else:
st.success("Processing completed. Review the generated protocol before use.") st.success("Processing completed. Review the generated protocol before use.")
st.caption(f"Run artifacts: {outcome.run_dir}") st.caption(f"Run artifacts: {outcome.run_dir}")
if outcome.speaker_attribution_available is False: if outcome.speaker_attribution_available is False:
@@ -563,6 +571,7 @@ def _render_result() -> None:
) )
if message := st.session_state.pop("speaker_mapping_message", None): if message := st.session_state.pop("speaker_mapping_message", None):
st.success(message) st.success(message)
if not awaiting_review:
with st.expander("Original generated protocol", expanded=False): with st.expander("Original generated protocol", expanded=False):
st.markdown(outcome.original_protocol or "") st.markdown(outcome.original_protocol or "")
edited = st.text_area( edited = st.text_area(
@@ -590,7 +599,7 @@ def _render_result() -> None:
if duplicate_assignments: if duplicate_assignments:
st.error("Each participant can be assigned to only one detected speaker.") st.error("Each participant can be assigned to only one detected speaker.")
if st.button( if st.button(
"Regenerate protocol with confirmed speakers", "Generate protocol" if awaiting_review else "Regenerate protocol with confirmed speakers",
disabled=duplicate_assignments, disabled=duplicate_assignments,
type="primary", type="primary",
): ):
@@ -610,7 +619,7 @@ def _render_result() -> None:
performance_profile=selected_profile, performance_profile=selected_profile,
), ),
states, states,
"Starting protocol regeneration", "Starting protocol generation" if awaiting_review else "Starting protocol regeneration",
) )
except (OSError, RuntimeError, ValueError) as exc: except (OSError, RuntimeError, ValueError) as exc:
_remember_regeneration_timing( _remember_regeneration_timing(
@@ -618,13 +627,13 @@ def _render_result() -> None:
"Protocol regeneration failed", "Protocol regeneration failed",
service.timer.snapshot(), service.timer.snapshot(),
) )
st.error(f"Protocol regeneration failed: {exc}") st.error(f"Protocol generation failed: {exc}")
else: else:
_remember_regeneration_timing(states, message, timing) _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 = (
"Speaker mappings saved and protocol regenerated." "Speaker mappings saved and protocol generated."
) )
st.rerun() st.rerun()
+3
View File
@@ -104,9 +104,12 @@ def test_paired_anonymous_generation_and_mapped_regeneration(
ProcessingOptions(diarization_enabled=True, performance_profile=profile), ProcessingOptions(diarization_enabled=True, performance_profile=profile),
) )
assert outcome.succeeded assert outcome.succeeded
assert outcome.awaiting_speaker_review
assert prepare.call_args.kwargs["normalization_enabled"] is True assert prepare.call_args.kwargs["normalization_enabled"] is True
assert transcribe.call_args.args[3] == language assert transcribe.call_args.args[3] == language
root = outcome.run_dir root = outcome.run_dir
assert calls == []
service.regenerate_protocol(root, {}, performance_profile=profile)
first = root / "protocol/generations/001" first = root / "protocol/generations/001"
before = {p.name: p.read_bytes() for p in first.iterdir()} before = {p.name: p.read_bytes() for p in first.iterdir()}
original = (root / "diarization/transcript_diarized.json").read_bytes() original = (root / "diarization/transcript_diarized.json").read_bytes()
+188 -1
View File
@@ -6,6 +6,7 @@ from types import SimpleNamespace
from typing import Any from typing import Any
import pytest import pytest
import yaml
from mka.application.config import AppSettings from mka.application.config import AppSettings
from mka.application.meeting_service import ( from mka.application.meeting_service import (
@@ -99,6 +100,24 @@ class FakeMeetingLab:
) )
) )
return SimpleNamespace(exit_code=2, run_dir=self.run_dir, protocol_path=None) return SimpleNamespace(exit_code=2, run_dir=self.run_dir, protocol_path=None)
if config["stop_after_diarization"]:
write_speaker_review_artifacts(self.run_dir)
(self.run_dir / "run_metadata.json").write_text(
json.dumps(
{
"status": "awaiting_speaker_review",
"diarization": {"speaker_count": 2},
}
),
encoding="utf-8",
)
progress_sink(
SimpleNamespace(
stage="diarization", status="completed", elapsed_seconds=0.3,
progress=None, message=None,
)
)
return SimpleNamespace(exit_code=0, run_dir=self.run_dir, protocol_path=None)
protocol = self.run_dir / "protocol.md" protocol = self.run_dir / "protocol.md"
protocol.write_text("# Generated protocol\n", encoding="utf-8") protocol.write_text("# Generated protocol\n", encoding="utf-8")
return SimpleNamespace(exit_code=0, run_dir=self.run_dir, protocol_path=protocol) return SimpleNamespace(exit_code=0, run_dir=self.run_dir, protocol_path=protocol)
@@ -228,6 +247,18 @@ def write_speaker_review_artifacts(run_dir: Path) -> Path:
return transcript_path return transcript_path
def write_custom_speaker_review_artifacts(
run_dir: Path, segments: list[dict[str, str | None]]
) -> Path:
"""Write custom diarized segments using the standard review context fixture."""
transcript_path = write_speaker_review_artifacts(run_dir)
transcript_path.write_text(
json.dumps({"speaker_labels_anonymous": True, "segments": segments}),
encoding="utf-8",
)
return transcript_path
def test_build_context_uses_actual_v1_shape(tmp_path: Path) -> None: def test_build_context_uses_actual_v1_shape(tmp_path: Path) -> None:
service, gateway = make_service(tmp_path) service, gateway = make_service(tmp_path)
@@ -387,7 +418,7 @@ def test_process_propagates_enabled_diarization_and_progress(tmp_path: Path) ->
audio.write_bytes(b"audio") audio.write_bytes(b"audio")
events = [] events = []
service.process( outcome = service.process(
audio, audio,
meeting(), meeting(),
participants(), participants(),
@@ -397,6 +428,9 @@ def test_process_propagates_enabled_diarization_and_progress(tmp_path: Path) ->
assert gateway.config_values is not None assert gateway.config_values is not None
assert gateway.config_values["diarization"] == "gpu" assert gateway.config_values["diarization"] == "gpu"
assert gateway.config_values["stop_after_diarization"] is True
assert outcome.awaiting_speaker_review is True
assert outcome.detected_speaker_count == 2
assert [(event.stage, event.status) for event in events[:2]] == [ assert [(event.stage, event.status) for event in events[:2]] == [
("preparing", "started"), ("preparing", "started"),
("preparing", "completed"), ("preparing", "completed"),
@@ -404,6 +438,33 @@ def test_process_propagates_enabled_diarization_and_progress(tmp_path: Path) ->
assert events[2].progress == 0.25 assert events[2].progress == 0.25
def test_diarization_checkpoint_defers_protocol_but_disabled_run_does_not(tmp_path: Path) -> None:
service, gateway = make_service(tmp_path)
audio = tmp_path / "meeting.wav"
audio.write_bytes(b"audio")
checkpoint = service.process(
audio, meeting(), participants(), ProcessingOptions(diarization_enabled=True)
)
assert checkpoint.succeeded
assert checkpoint.awaiting_speaker_review
assert checkpoint.original_protocol is None
assert checkpoint.protocol_path is None
normal_root = tmp_path / "without-diarization"
normal_root.mkdir()
normal_service, normal_gateway = make_service(normal_root)
normal_audio = normal_root / "meeting.wav"
normal_audio.write_bytes(b"audio")
normal = normal_service.process(
normal_audio, meeting(), participants(), ProcessingOptions(diarization_enabled=False)
)
assert normal.succeeded
assert not normal.awaiting_speaker_review
assert normal.protocol_path is not None
assert normal_gateway.config_values["stop_after_diarization"] is False
def test_process_propagates_ordered_diarization_container_args(tmp_path: Path) -> None: def test_process_propagates_ordered_diarization_container_args(tmp_path: Path) -> None:
service, gateway = make_service(tmp_path) service, gateway = make_service(tmp_path)
service.settings = replace( service.settings = replace(
@@ -485,6 +546,89 @@ def test_speaker_review_lists_detected_labels_participants_and_excerpts(
assert "SPEAKER_00" in source.read_text(encoding="utf-8") assert "SPEAKER_00" in source.read_text(encoding="utf-8")
def test_speaker_review_prefers_content_rich_excerpts_over_acknowledgements_and_questions(
tmp_path: Path,
) -> None:
service, gateway = make_service(tmp_path)
acknowledgement = "Okay, that sounds good."
question = "Could we continue now?"
rich_excerpts = [
"I will complete the Atlas deployment in Berlin on 14 September with the database team.",
"I recommend moving the Orion product launch because the supplier test is incomplete.",
"We decided that I will own the Kubernetes migration for the payments service.",
"My experience with the Frankfurt factory rollout suggests a two-week validation period.",
"I will present the security audit results to the Apollo project steering group.",
"We should reserve twelve hours for the API integration and production monitoring.",
]
write_custom_speaker_review_artifacts(
gateway.run_dir,
[{"speaker_id": "SPEAKER_00", "text": text} for text in [acknowledgement, question, *rich_excerpts]],
)
review = service.load_speaker_mapping_review(gateway.run_dir)
assert review is not None
excerpts = review.speakers[0].excerpts
assert len(excerpts) == 6
assert acknowledgement not in excerpts
assert question not in excerpts
assert set(excerpts) == set(rich_excerpts)
def test_speaker_review_avoids_near_duplicate_excerpts(tmp_path: Path) -> None:
service, gateway = make_service(tmp_path)
primary = "I will lead the Atlas API migration for the Berlin platform team next Monday."
duplicate = "I will lead the Atlas API migration for the Berlin platform group next Monday."
alternatives = [
"I recommend documenting the supplier quality decision in the project register.",
"My team will complete the production calibration before the factory trial.",
"I will present the customer research findings at the Munich planning workshop.",
"Our database monitoring threshold needs approval from the operations group.",
"I will coordinate the Orion release checklist with the support organization.",
]
write_custom_speaker_review_artifacts(
gateway.run_dir,
[{"speaker_id": "SPEAKER_00", "text": text} for text in [primary, duplicate, *alternatives]],
)
review = service.load_speaker_mapping_review(gateway.run_dir)
assert review is not None
excerpts = review.speakers[0].excerpts
assert primary in excerpts
assert duplicate not in excerpts
assert set(excerpts) == {primary, *alternatives}
def test_speaker_review_spreads_selected_excerpts_across_meeting_positions(
tmp_path: Path,
) -> None:
service, gateway = make_service(tmp_path)
positions = {
0: "I will lead the Atlas deployment review with the engineering team marker alpha.",
1: "I will lead the Atlas deployment review with the engineering team marker bravo.",
6: "I will lead the Atlas deployment review with the engineering team marker charlie.",
7: "I will lead the Atlas deployment review with the engineering team marker delta.",
12: "I will lead the Atlas deployment review with the engineering team marker echo.",
13: "I will lead the Atlas deployment review with the engineering team marker foxtrot.",
18: "I will lead the Atlas deployment review with the engineering team marker golf.",
19: "I will lead the Atlas deployment review with the engineering team marker hotel.",
}
segments = [
{
"speaker_id": "SPEAKER_00" if index in positions else "SPEAKER_01",
"text": positions.get(index, "Acknowledged."),
}
for index in range(20)
]
write_custom_speaker_review_artifacts(gateway.run_dir, segments)
review = service.load_speaker_mapping_review(gateway.run_dir)
assert review is not None
assert review.speakers[0].excerpts == tuple(positions[index] for index in (0, 1, 6, 7, 12, 18))
def test_protocol_only_regeneration_persists_mappings_without_rewriting_transcript( def test_protocol_only_regeneration_persists_mappings_without_rewriting_transcript(
tmp_path: Path, tmp_path: Path,
) -> None: ) -> None:
@@ -537,6 +681,49 @@ def test_protocol_regeneration_allows_unmapped_and_rejects_duplicate_participant
) )
@pytest.mark.parametrize(
"mappings",
[
{"SPEAKER_00": "martin"},
{"SPEAKER_00": "martin", "SPEAKER_01": "anna"},
],
)
def test_protocol_generation_uses_partial_and_complete_confirmed_mappings(
tmp_path: Path, mappings: dict[str, str]
) -> None:
service, gateway = make_service(tmp_path)
write_speaker_review_artifacts(gateway.run_dir)
outcome = service.regenerate_protocol(gateway.run_dir, mappings)
assert outcome.succeeded
assert gateway.regeneration is not None
assert gateway.regeneration["meeting_context"].data["speaker_mappings"] == mappings
def test_failed_protocol_generation_keeps_checkpoint_mapping_for_retry(tmp_path: Path) -> None:
service, gateway = make_service(tmp_path)
source = write_speaker_review_artifacts(gateway.run_dir)
original_source = source.read_bytes()
original_regeneration = gateway.regenerate_protocol
def fail_regeneration(*args, **kwargs):
raise RuntimeError("model unavailable")
gateway.regenerate_protocol = fail_regeneration
with pytest.raises(RuntimeError, match="model unavailable"):
service.regenerate_protocol(gateway.run_dir, {"SPEAKER_00": "martin"})
persisted = yaml.safe_load(
(gateway.run_dir / "context" / "meeting_context.yaml").read_text(encoding="utf-8")
)
assert persisted["speaker_mappings"] == {"SPEAKER_00": "martin"}
gateway.regenerate_protocol = original_regeneration
retry = service.regenerate_protocol(gateway.run_dir, {"SPEAKER_00": "martin"})
assert retry.succeeded
assert source.read_bytes() == original_source
def test_fallback_attribution_loss_is_read_from_runtime_metadata(tmp_path: Path) -> None: def test_fallback_attribution_loss_is_read_from_runtime_metadata(tmp_path: Path) -> None:
run_dir = tmp_path / "run" run_dir = tmp_path / "run"
protocol_dir = run_dir / "protocol" protocol_dir = run_dir / "protocol"
+6 -4
View File
@@ -80,7 +80,7 @@ def test_widget_reruns_filter_release_and_update_status():
assert not app.exception assert not app.exception
def test_anonymous_regeneration_remains_allowed(monkeypatch): def test_checkpoint_allows_anonymous_protocol_generation(monkeypatch):
from mka.application.meeting_service import SpeakerMappingReview, SpeakerReview from mka.application.meeting_service import SpeakerMappingReview, SpeakerReview
from mka.ui import streamlit_app as ui from mka.ui import streamlit_app as ui
@@ -89,6 +89,8 @@ def test_anonymous_regeneration_remains_allowed(monkeypatch):
run_dir=Path("/tmp/mapping-test"), run_dir=Path("/tmp/mapping-test"),
speaker_attribution_available=True, speaker_attribution_available=True,
original_protocol="Anonymous protocol", original_protocol="Anonymous protocol",
awaiting_speaker_review=True,
detected_speaker_count=1,
) )
review = SpeakerMappingReview((SpeakerReview("SPEAKER_00", ()),), (("a", "A"),), {}) review = SpeakerMappingReview((SpeakerReview("SPEAKER_00", ()),), (("a", "A"),), {})
calls = [] calls = []
@@ -124,9 +126,9 @@ def test_anonymous_regeneration_remains_allowed(monkeypatch):
_render_result() _render_result()
app = AppTest.from_function(result_app, args=(outcome,)).run() app = AppTest.from_function(result_app, args=(outcome,)).run()
regenerate = app.button[1] generate = next(button for button in app.button if button.label == "Generate protocol")
assert not regenerate.disabled assert not generate.disabled
assert app.warning assert app.warning
regenerate.click().run() generate.click().run()
assert calls == [{}] assert calls == [{}]
assert not app.exception assert not app.exception