16 changed files with 779 additions and 48 deletions
+20
View File
@@ -295,3 +295,23 @@ The next recommended evaluation step is to use the consolidated V0 result as
input for the unchanged Working Protocol renderer and compare that output
against the Working Protocol Synthesizer V0 baseline and the human reference
protocol.
Protocol calls accept an optional positive `protocol_num_thread` setting through
both initial processing and regeneration. With `None`, Ollama receives no
`num_thread` override and selects its own thread configuration.
Configured glossary aliases are diagnostic metadata only. Meeting Context supplies
terminology guidance; direct-protocol transcript text is never alias-substituted.
See `docs/protocol-generation-regression.md` for the frozen-input regression.
Successful protocols are retained in `protocol/generations/NNN/` with their exact
prompt, transcript input, response, model/runtime metadata and Meeting Context.
Relative compatibility symlinks resolve through `protocol/current`, which is
replaced atomically only after a complete record is written. Failed regeneration
keeps the previous protocol, diagnostics and context. Legacy regular files are
snapshotted before conversion. Publication is serialized with a local file lock.
Read `current` once when inspecting a consistent multi-file snapshot. Copy whole
run directories preserving relative symlinks. Process interruption may leave an
unreferenced record or temporary directory, never a partial current generation.
This local POSIX-filesystem contract does not promise power-loss durability or
network-filesystem transaction semantics.
+12
View File
@@ -278,3 +278,15 @@ Artefakten, Windows-Beispiel und Vergleich mit dem AI-PC stehen in
## Lizenz
Noch nicht festgelegt.
Successful protocols are retained in `protocol/generations/NNN/` with their exact
prompt, transcript input, response, model/runtime metadata and Meeting Context.
Relative compatibility symlinks resolve through `protocol/current`, which is
replaced atomically only after a complete record is written. Failed regeneration
keeps the previous protocol, diagnostics and context. Legacy regular files are
snapshotted before conversion. Publication is serialized with a local file lock.
Read `current` once when inspecting a consistent multi-file snapshot. Copy whole
run directories preserving relative symlinks. Process interruption may leave an
unreferenced record or temporary directory, never a partial current generation.
This local POSIX-filesystem contract does not promise power-loss durability or
network-filesystem transaction semantics.
+5
View File
@@ -60,3 +60,8 @@ also exceeds the budget, generation fails before model lookup or generation;
it never truncates, chunks, summarizes, retries, or makes multiple protocol
calls implicitly. Full diarization artifacts are never overwritten by this
selection.
Glossary aliases are recorded as configuration provenance, never applied as
deterministic replacements to the compact/plain protocol input. Canonical
terminology guides generation through Meeting Context. Raw Whisper and
diarization artifacts remain unchanged. `glossary_replacements` is always empty.
+25
View File
@@ -0,0 +1,25 @@
# Direct-protocol regression reproduction
The August 2026 glossary regression was isolated with frozen inputs. Condition
A used the original derived transcript; condition B differed only by these
seven deterministic substitutions:
- `Carbofool` to `Carbofol` (two occurrences)
- `Bento Fix` to `Bentofix` (two occurrences)
- `Sikirgut-Heistlöse` to `Secugrid HS` (one occurrence)
- `Lumini` to `Luminy` (two occurrences)
To repeat the manual comparison, copy the investigated run to a new temporary
directory, retain its Meeting Context and speaker mapping, and invoke
`regenerate_mvp_protocol` through the same parameters used by Meeting
Assistant. Never run the comparison in the historical run directory. Preserve
the model, `num_ctx`, `num_predict`, temperature, think setting, thread setting,
and Meeting Context. Compare the new generation's `exact_prompt.txt` and
`transcript_input.txt` with the frozen A and B artifacts before comparing model
output.
The required fixed condition is A: configured glossary aliases remain visible
in Meeting Context terminology guidance, while `transcript_input.txt` retains
the original seven source spellings. Automated tests mock the model and enforce
that invariant; this full historical experiment remains an explicit manual LLM
validation so normal tests do not depend on a local model.
+25
View File
@@ -0,0 +1,25 @@
# Meeting Assistant — First functional alpha
Version: `0.1.0a1`
Tag: `v0.1.0-alpha.1`
This paired Meeting Lab revision supports the first functional end-to-end
Meeting Assistant alpha: reusable transcription and optional diarization,
explicit speaker mapping and regeneration, German and English protocol
generation, glossary provenance without transcript mutation, configurable
protocol thread profiles, generation history with atomic publication, and the
GTM-Hub qualitative regression case.
Validated candidate revisions:
- 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 limitations are documented in the Assistant release notes, including
the requirement for local runtime/model/container setup, human review of
generated protocols, imperfect diarization, context limits, local POSIX
publication assumptions, and the historical ignored fixtures needed by some
old experiment tests.
+1 -1
View File
@@ -1,6 +1,6 @@
[project]
name = "meeting-lab"
version = "0.1.0"
version = "0.1.0a1"
requires-python = ">=3.11"
dependencies = [
"PyYAML>=6.0",
+7 -18
View File
@@ -17,11 +17,12 @@ REPO_ROOT = Path(__file__).resolve().parents[1]
if str(REPO_ROOT) not in sys.path:
sys.path.insert(0, str(REPO_ROOT))
from src.meeting_lab.models.meeting_context import load_meeting_context # noqa: E402
from src.meeting_lab.protocol.history import persist_protocol_generation # noqa: E402
from src.meeting_lab.llm.ollama import DEFAULT_ENDPOINT # noqa: E402
from src.meeting_lab.protocol.generate_direct_protocol import ( # noqa: E402
DEFAULT_MODEL,
DEFAULT_SAFE_INPUT_TOKEN_BUDGET,
DirectProtocolResult,
generate_direct_protocol,
)
@@ -64,22 +65,6 @@ def write_json(path: Path, data: Any) -> None:
path.write_text(json.dumps(data, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
def persist_result(run_dir: Path, result: DirectProtocolResult) -> Path:
protocol_dir = run_dir / "protocol"
protocol_dir.mkdir()
(protocol_dir / "exact_prompt.txt").write_text(result.exact_prompt, encoding="utf-8")
write_json(protocol_dir / "raw_response.json", result.raw_response)
write_json(protocol_dir / "runtime_metadata.json", result.runtime_metadata)
transcript_input = getattr(result, "transcript_input", None)
if transcript_input is not None:
(protocol_dir / "transcript_input.txt").write_text(
transcript_input, encoding="utf-8"
)
protocol_path = run_dir / "protocol.md"
protocol_path.write_text(result.protocol_text, encoding="utf-8")
return protocol_path
def run(args: argparse.Namespace) -> tuple[int, Path, Path | None]:
run_dir = create_unique_run_dir(args.output_root, args.transcript.stem)
timestamp = datetime.now().astimezone().isoformat(timespec="seconds")
@@ -129,7 +114,11 @@ def run(args: argparse.Namespace) -> tuple[int, Path, Path | None]:
endpoint=args.ollama_endpoint,
safe_input_token_budget=args.safe_input_token_budget,
)
protocol_path = persist_result(run_dir, result)
protocol_path = persist_protocol_generation(
run_dir,
result,
context=load_meeting_context(preserved_context) if preserved_context else None,
)
metadata["status"] = "completed"
metadata["final_protocol_path"] = str(protocol_path.resolve())
except Exception as exc:
+5
View File
@@ -62,6 +62,7 @@ def generate_once(
timeout: int,
num_ctx: int,
num_predict: int,
num_thread: int | None = None,
) -> OllamaGeneration:
payload = {
"model": model,
@@ -74,6 +75,10 @@ def generate_once(
"num_predict": num_predict,
},
}
if num_thread is not None:
if type(num_thread) is not int or num_thread <= 0:
raise ValueError("Protocol thread count must be a positive integer.")
payload["options"]["num_thread"] = num_thread
started = time.perf_counter()
try:
response = requests.post(generate_url(endpoint), json=payload, timeout=timeout)
+76 -29
View File
@@ -6,6 +6,7 @@ import json
import re
import shutil
import sys
import tempfile
import time
from collections.abc import Callable, Mapping, Sequence
from dataclasses import dataclass
@@ -28,11 +29,11 @@ from src.meeting_lab.models.meeting_context import (
write_meeting_context,
)
from src.meeting_lab.progress import ProgressEvent, ProgressSink, ProgressStatus
from src.meeting_lab.protocol.history import persist_protocol_generation
from src.meeting_lab.protocol.generate_direct_protocol import (
DEFAULT_MODEL,
DEFAULT_NUM_CTX,
DEFAULT_SAFE_INPUT_TOKEN_BUDGET,
DirectProtocolResult,
generate_direct_protocol,
load_compact_transcript,
)
@@ -56,12 +57,15 @@ class MvpMeetingConfig:
threads: str | int = "auto"
model: str = DEFAULT_MODEL
ollama_endpoint: str = DEFAULT_ENDPOINT
glossary_aliases: Mapping[str, str] | None = None
protocol_num_thread: int | None = None
protocol_num_ctx: int = DEFAULT_NUM_CTX
protocol_safe_input_token_budget: int = DEFAULT_SAFE_INPUT_TOKEN_BUDGET
diarization: str = "off"
diarization_runtime: str = "native"
diarization_container_image: str | None = None
diarization_container_args: Sequence[str] = ()
stop_after_diarization: bool = False
@dataclass(frozen=True)
@@ -77,6 +81,8 @@ def regenerate_mvp_protocol(
meeting_context: ContextInput,
model: str = DEFAULT_MODEL,
ollama_endpoint: str = DEFAULT_ENDPOINT,
glossary_aliases: Mapping[str, str] | None = None,
protocol_num_thread: int | None = None,
protocol_num_ctx: int = DEFAULT_NUM_CTX,
protocol_safe_input_token_budget: int = DEFAULT_SAFE_INPUT_TOKEN_BUDGET,
progress_sink: ProgressSink | None = None,
@@ -87,6 +93,10 @@ def regenerate_mvp_protocol(
context = _effective_context(meeting_context)
if context is None:
raise ValueError("Meeting Context is required for protocol regeneration.")
if protocol_num_thread is not None and (
type(protocol_num_thread) is not int or protocol_num_thread <= 0
):
raise ValueError("Protocol Ollama thread count must be a positive integer.")
if protocol_num_ctx <= 0:
raise ValueError("Protocol Ollama context size must be positive.")
if protocol_safe_input_token_budget <= 0:
@@ -102,20 +112,29 @@ def regenerate_mvp_protocol(
f"Existing run has no protocol transcript artifact: {run_dir}"
)
context_path = run_dir / "context" / "meeting_context.yaml"
write_meeting_context(context, context_path)
_emit(progress_sink, "protocol_generation", "started", started)
try:
result = generate_direct_protocol(
transcript_path,
context_path,
model=model,
endpoint=ollama_endpoint,
num_ctx=protocol_num_ctx,
safe_input_token_budget=protocol_safe_input_token_budget,
)
protocol_path = _persist_protocol(run_dir, result)
with tempfile.TemporaryDirectory(
prefix=".protocol-context-", dir=run_dir, ignore_cleanup_errors=True
) as temporary:
context_path = Path(temporary) / "meeting_context.yaml"
write_meeting_context(context, context_path)
result = generate_direct_protocol(
transcript_path,
context_path,
model=model,
endpoint=ollama_endpoint,
num_ctx=protocol_num_ctx,
num_thread=protocol_num_thread,
glossary_aliases=glossary_aliases,
safe_input_token_budget=protocol_safe_input_token_budget,
)
protocol_path = persist_protocol_generation(run_dir, result, context=context)
except Exception as exc:
_record_protocol_generation_attempt(
run_dir, status="failed", runtime_seconds=time.perf_counter() - started,
failure={"type": type(exc).__name__, "message": str(exc)},
)
_emit(
progress_sink,
"failed",
@@ -124,11 +143,40 @@ def regenerate_mvp_protocol(
message=f"protocol_generation: {type(exc).__name__}: {exc}",
)
raise
_record_protocol_generation_attempt(
run_dir, status="completed", runtime_seconds=time.perf_counter() - started
)
_emit(progress_sink, "protocol_generation", "completed", started)
_emit(progress_sink, "completed", "completed", started)
return MvpRunResult(0, run_dir, protocol_path)
def _record_protocol_generation_attempt(
run_dir: Path,
*,
status: str,
runtime_seconds: float,
failure: dict[str, str] | None = None,
) -> None:
"""Record downstream attempts without changing completed upstream artifacts."""
metadata_path = run_dir / "run_metadata.json"
if not metadata_path.is_file():
return
metadata = json.loads(metadata_path.read_text(encoding="utf-8"))
attempts = metadata.setdefault("protocol_generation_attempts", [])
if not isinstance(attempts, list):
attempts = metadata["protocol_generation_attempts"] = []
attempt: dict[str, Any] = {
"status": status,
"runtime_seconds": round(runtime_seconds, 3),
}
if failure is not None:
attempt["failure"] = failure
attempts.append(attempt)
metadata["status"] = "completed" if status == "completed" else "awaiting_speaker_review"
_write_json(metadata_path, metadata)
def create_unique_run_dir(
output_root: Path,
meeting_name: str,
@@ -191,6 +239,10 @@ def _validate_inputs(
raise ValueError("A diarization container image is required.")
if config.protocol_safe_input_token_budget <= 0:
raise ValueError("Protocol safe input token budget must be positive.")
if config.protocol_num_thread is not None and (
type(config.protocol_num_thread) is not int or config.protocol_num_thread <= 0
):
raise ValueError("Protocol Ollama thread count must be a positive integer.")
if config.protocol_num_ctx <= 0:
raise ValueError("Protocol Ollama context size must be positive.")
@@ -214,22 +266,6 @@ def _emit(
)
def _persist_protocol(run_dir: Path, result: DirectProtocolResult) -> Path:
protocol_dir = run_dir / "protocol"
protocol_dir.mkdir(exist_ok=True)
(protocol_dir / "exact_prompt.txt").write_text(result.exact_prompt, encoding="utf-8")
_write_json(protocol_dir / "raw_response.json", result.raw_response)
_write_json(protocol_dir / "runtime_metadata.json", result.runtime_metadata)
transcript_input = getattr(result, "transcript_input", None)
if transcript_input is not None:
(protocol_dir / "transcript_input.txt").write_text(
transcript_input, encoding="utf-8"
)
protocol_path = run_dir / "protocol.md"
protocol_path.write_text(result.protocol_text, encoding="utf-8")
return protocol_path
def run_mvp_meeting(
config: MvpMeetingConfig,
*,
@@ -417,6 +453,15 @@ def run_mvp_meeting(
)
_emit(progress_sink, "diarization", "completed", overall_started)
if config.stop_after_diarization:
metadata["status"] = "awaiting_speaker_review"
metadata["speaker_review"] = {
"status": "awaiting_optional_mapping",
"speaker_count": metadata["diarization"].get("speaker_count"),
}
_emit(progress_sink, "completed", "completed", overall_started)
return MvpRunResult(0, run_dir, None)
current_stage = "protocol_generation"
stage_started = time.perf_counter()
_emit(progress_sink, "protocol_generation", "started", overall_started)
@@ -426,10 +471,12 @@ def run_mvp_meeting(
model=config.model,
endpoint=config.ollama_endpoint,
num_ctx=config.protocol_num_ctx,
num_thread=config.protocol_num_thread,
glossary_aliases=config.glossary_aliases,
safe_input_token_budget=config.protocol_safe_input_token_budget,
)
stage_runtimes["protocol"] = round(time.perf_counter() - stage_started, 3)
protocol_path = _persist_protocol(run_dir, result)
protocol_path = persist_protocol_generation(run_dir, result, context=effective_context)
_emit(progress_sink, "protocol_generation", "completed", overall_started)
metadata["status"] = "completed"
_emit(progress_sink, "completed", "completed", overall_started)
@@ -3,6 +3,7 @@
from __future__ import annotations
import json
from collections.abc import Mapping
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Callable
@@ -170,9 +171,11 @@ def generate_direct_protocol(
timeout: int = DEFAULT_TIMEOUT,
num_ctx: int = DEFAULT_NUM_CTX,
num_predict: int = DEFAULT_NUM_PREDICT,
num_thread: int | None = None,
safe_input_token_budget: int = DEFAULT_SAFE_INPUT_TOKEN_BUDGET,
model_check: Callable[[str, str, int], dict[str, Any]] = require_model,
generation_call: Callable[..., OllamaGeneration] = generate_once,
glossary_aliases: Mapping[str, str] | None = None,
) -> DirectProtocolResult:
transcript = _load_transcript_document(transcript_path)
context: MeetingContext | None = (
@@ -198,6 +201,7 @@ def generate_direct_protocol(
timeout=timeout,
num_ctx=num_ctx,
num_predict=num_predict,
num_thread=num_thread,
)
data = generation.raw_response
runtime_metadata = {
@@ -216,6 +220,7 @@ def generate_direct_protocol(
"think": False,
"num_ctx": num_ctx,
"num_predict": num_predict,
"num_thread": num_thread,
"selected_transcript_representation": selected.representation,
"estimated_input_tokens": selected.estimated_input_tokens,
"safe_input_token_budget": selected.safe_input_token_budget,
@@ -235,6 +240,8 @@ def generate_direct_protocol(
else None
),
"speaker_mapping_count": len(context.speaker_mappings) if context else 0,
"glossary_aliases_configured": dict(sorted((glossary_aliases or {}).items())),
"glossary_replacements": [],
}
return DirectProtocolResult(
protocol_text=generation.text,
+188
View File
@@ -0,0 +1,188 @@
"""Immutable protocol records with one atomic latest-generation publication."""
from __future__ import annotations
import fcntl
import json
import os
import shutil
import subprocess
import tempfile
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
from uuid import uuid4
from src.meeting_lab.models.meeting_context import MeetingContext, write_meeting_context
from src.meeting_lab.protocol.generate_direct_protocol import DirectProtocolResult
_DIAGNOSTICS = (
"exact_prompt.txt",
"transcript_input.txt",
"raw_response.json",
"runtime_metadata.json",
)
def _git_provenance(repository: Path) -> dict[str, Any]:
provenance: dict[str, Any] = {"repository": str(repository), "sha": None, "dirty": None}
try:
provenance["sha"] = subprocess.check_output(
["git", "rev-parse", "HEAD"], cwd=repository, text=True, stderr=subprocess.DEVNULL
).strip()
provenance["dirty"] = bool(
subprocess.check_output(
["git", "status", "--porcelain"],
cwd=repository,
text=True,
stderr=subprocess.DEVNULL,
)
)
except (OSError, subprocess.SubprocessError):
pass
return provenance
def _write_json(path: Path, value: Any) -> None:
path.write_text(json.dumps(value, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
def _replace_link(path: Path, target: str) -> None:
"""Replace a link/file atomically; never write through an existing symlink."""
temporary = path.with_name(f".{path.name}-{uuid4().hex}")
try:
temporary.symlink_to(target)
os.replace(temporary, path)
finally:
temporary.unlink(missing_ok=True)
def _sync_record(directory: Path) -> None:
for path in directory.iterdir():
with path.open("rb") as stream:
os.fsync(stream.fileno())
descriptor = os.open(directory, os.O_RDONLY)
try:
os.fsync(descriptor)
finally:
os.close(descriptor)
def _legacy_generation(run_dir: Path, protocol_dir: Path, generations: Path) -> None:
"""Bootstrap old regular files before switching any compatibility paths."""
if (protocol_dir / "current").is_symlink() or not (run_dir / "protocol.md").is_file():
return
# Preserve the latest regular-file state even for runs created by an older
# history implementation. Never reuse or overwrite an existing record.
index = max(
(int(p.name) for p in generations.iterdir() if p.is_dir() and p.name.isdigit()),
default=0,
) + 1
legacy = generations / f"{index:03d}"
pending = Path(tempfile.mkdtemp(prefix=".legacy-", dir=generations))
try:
shutil.copyfile(run_dir / "protocol.md", pending / "protocol.md")
for name in _DIAGNOSTICS:
source = protocol_dir / name
if source.is_file():
shutil.copyfile(source, pending / name)
context = run_dir / "context/meeting_context.yaml"
if context.is_file():
shutil.copyfile(context, pending / "meeting_context.yaml")
_sync_record(pending)
pending.rename(legacy)
finally:
shutil.rmtree(pending, ignore_errors=True)
_replace_link(protocol_dir / "current", f"generations/{legacy.name}")
def _compatibility_links(run_dir: Path, protocol_dir: Path, pending: Path) -> None:
"""Install indirections while they still resolve to the previous generation."""
links = [(run_dir / "protocol.md", "protocol/current/protocol.md")]
links.extend((protocol_dir / name, f"current/{name}") for name in _DIAGNOSTICS)
if (pending / "meeting_context.yaml").is_file():
context_dir = run_dir / "context"
context_dir.mkdir(exist_ok=True)
links.append(
(context_dir / "meeting_context.yaml", "../protocol/current/meeting_context.yaml")
)
for path, target in links:
if not path.is_symlink() or os.readlink(path) != target:
_replace_link(path, target)
def persist_protocol_generation(
run_dir: Path,
result: DirectProtocolResult,
*,
context: MeetingContext | None,
) -> Path:
"""Publish a complete record; current is the sole successful-generation pointer.
Readers needing a multi-file snapshot should resolve current once. Relative
symlinks keep complete run directories movable. Writes are serialized locally.
"""
protocol_dir = run_dir / "protocol"
protocol_dir.mkdir(exist_ok=True)
generations = protocol_dir / "generations"
generations.mkdir(exist_ok=True)
with (protocol_dir / ".publication.lock").open("a") as lock:
fcntl.flock(lock, fcntl.LOCK_EX)
_legacy_generation(run_dir, protocol_dir, generations)
index = (
max(
(int(p.name) for p in generations.iterdir() if p.is_dir() and p.name.isdigit()),
default=0,
)
+ 1
)
destination = generations / f"{index:03d}"
pending = Path(tempfile.mkdtemp(prefix=".pending-", dir=generations))
published = False
try:
metadata = dict(result.runtime_metadata)
metadata.update(
{
"generation_index": index,
"generation_timestamp": datetime.now(UTC).isoformat(timespec="seconds"),
"source_run_id": run_dir.name,
"speaker_mapping": dict(sorted(context.speaker_mappings.items()))
if context
else {},
"speaker_mapping_names": {
label: context.participant_for_speaker(label).get("display_name", "")
for label in sorted(context.speaker_mappings)
if context.participant_for_speaker(label) is not None
}
if context
else {},
"git": {"meeting_lab": _git_provenance(Path(__file__).resolve().parents[3])},
}
)
(pending / "protocol.md").write_text(result.protocol_text, encoding="utf-8")
(pending / "exact_prompt.txt").write_text(result.exact_prompt, encoding="utf-8")
if result.transcript_input is not None:
(pending / "transcript_input.txt").write_text(
result.transcript_input, encoding="utf-8"
)
_write_json(pending / "raw_response.json", result.raw_response)
_write_json(pending / "runtime_metadata.json", metadata)
_write_json(pending / "model_metadata.json", result.model_metadata)
if context is not None:
write_meeting_context(context, pending / "meeting_context.yaml")
_sync_record(pending)
_compatibility_links(run_dir, protocol_dir, pending)
pending.rename(destination)
_replace_link(protocol_dir / "current", f"generations/{destination.name}")
published = True
finally:
shutil.rmtree(pending, ignore_errors=True)
# An interrupted process may leave an unreferenced complete record;
# never delete the record selected by current, even after interruption.
if (
not published
and destination.exists()
and (protocol_dir / "current").resolve() != destination
):
shutil.rmtree(destination, ignore_errors=True)
return run_dir / "protocol.md"
+20
View File
@@ -198,6 +198,23 @@ class OllamaTests(unittest.TestCase):
self.assertFalse(payload["stream"])
self.assertEqual(result.raw_response, raw)
def test_request_omits_num_thread_when_backend_selection_is_requested(self) -> None:
response = Mock()
response.raise_for_status.return_value = None
response.json.return_value = {"response": "# Meeting Protocol"}
with patch.object(ollama.requests, "post", return_value=response) as post:
ollama.generate_once(
"http://localhost:11434",
"qwen3.8:27b",
"prompt",
timeout=30,
num_ctx=32768,
num_predict=8192,
)
payload = post.call_args.kwargs["json"]
self.assertNotIn("num_thread", payload["options"])
def test_malformed_response_failure_without_retry(self) -> None:
response = Mock()
response.raise_for_status.return_value = None
@@ -237,6 +254,7 @@ class DirectProtocolCliTests(unittest.TestCase):
"transcript_input": "selected transcript\n",
"raw_response": {"response": protocol_text},
"runtime_metadata": {"request_count": 1},
"model_metadata": {},
})(),
) as generator:
code, run_dir, protocol_path = run_direct_protocol.run(args)
@@ -304,6 +322,8 @@ class DirectProtocolCliTests(unittest.TestCase):
"exact_prompt": "prompt",
"raw_response": {"response": "# Meeting Protocol"},
"runtime_metadata": {},
"model_metadata": {},
"transcript_input": None,
})()
with (
patch("src.meeting_lab.extraction.extract_chunks.extract_input") as extraction,
+124
View File
@@ -9,6 +9,9 @@ from unittest.mock import patch
from scripts import run_mvp_meeting as cli
from src.meeting_lab.audio import PreparedAudio
from src.meeting_lab.diarization.backend import DiarizationResult
from src.meeting_lab.llm.ollama import OllamaGeneration
from src.meeting_lab.protocol.generate_direct_protocol import generate_direct_protocol
from src.meeting_lab.models.meeting_context import load_meeting_context
from src.meeting_lab.orchestration import mvp as mvp_api
from src.meeting_lab.orchestration.mvp import MvpMeetingConfig, MvpRunResult
@@ -80,6 +83,27 @@ def fake_prepare(source, destination, **kwargs):
return PreparedAudio(source, source.suffix.removeprefix("."), destination, "ffmpeg", "ffmpeg")
def fake_diarize(audio, output, mode, **kwargs):
output.mkdir(parents=True, exist_ok=True)
metadata = output / "metadata.json"
ordinary = output / "diarization.rttm"
exclusive = output / "exclusive_diarization.rttm"
turns = output / "turns.json"
exclusive_turns = output / "exclusive_turns.json"
details = {"speaker_count": 2, "runtime_seconds": 0.1}
metadata.write_text(json.dumps(details), encoding="utf-8")
ordinary.write_text("", encoding="utf-8")
exclusive.write_text("", encoding="utf-8")
turns.write_text("[]", encoding="utf-8")
exclusive_turns.write_text(
json.dumps([
{"start": 0, "end": 0.5, "speaker_id": "SPEAKER_00"},
{"start": 0.5, "end": 1, "speaker_id": "SPEAKER_01"},
]), encoding="utf-8"
)
return DiarizationResult(output, metadata, ordinary, exclusive, turns, exclusive_turns, details)
def fake_protocol(transcript, context, **kwargs):
rendered_context = load_meeting_context(context)
assert rendered_context.meeting_id == "programmatic-test"
@@ -105,6 +129,26 @@ class MvpApiTests(unittest.TestCase):
model="test:model",
)
def test_post_diarization_checkpoint_skips_protocol_generation(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
config = replace(self.config(root), diarization="cpu", stop_after_diarization=True)
with (
patch.object(mvp_api, "prepare_audio", side_effect=fake_prepare),
patch.object(mvp_api, "transcribe_audio", side_effect=fake_transcribe),
patch.object(mvp_api, "diarize_audio", side_effect=fake_diarize),
patch.object(mvp_api, "generate_direct_protocol") as generate,
):
result = mvp_api.run_mvp_meeting(config, meeting_context=context_data())
self.assertEqual(result.exit_code, 0)
self.assertIsNone(result.protocol_path)
generate.assert_not_called()
metadata = json.loads((result.run_dir / "run_metadata.json").read_text())
self.assertEqual(metadata["status"], "awaiting_speaker_review")
self.assertEqual(metadata["speaker_review"]["speaker_count"], 2)
self.assertTrue((result.run_dir / "diarization/transcript_diarized.json").is_file())
def test_programmatic_context_is_persisted_without_source_yaml(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
@@ -116,6 +160,7 @@ class MvpApiTests(unittest.TestCase):
patch.object(
mvp_api, "generate_direct_protocol", side_effect=fake_protocol
) as protocol_generator,
patch("src.meeting_lab.protocol.history._git_provenance", return_value={}),
patch.object(subprocess, "run") as subprocess_run,
):
result = mvp_api.run_mvp_meeting(
@@ -202,6 +247,85 @@ class MvpApiTests(unittest.TestCase):
persisted = load_meeting_context(run_dir / "context/meeting_context.yaml")
self.assertEqual(persisted.speaker_mappings, {"SPEAKER_00": "person-1"})
def test_regeneration_keeps_glossary_out_of_transcript_and_in_context(self):
with tempfile.TemporaryDirectory() as directory:
run_dir = Path(directory)
diarization_dir = run_dir / "diarization"
diarization_dir.mkdir()
source = diarization_dir / "transcript_diarized.json"
source.write_text(
json.dumps(
{
"text": "SPEAKER_00: Lumini.",
"segments": [
{
"start": 0.0,
"end": 1.0,
"speaker_id": "SPEAKER_00",
"text": "Lumini.",
}
],
"speaker_labels_anonymous": True,
}
),
encoding="utf-8",
)
source_before = source.read_bytes()
transcript_dir = run_dir / "transcript"
transcript_dir.mkdir()
whisper_source = transcript_dir / "transcript.json"
whisper_source.write_text(
json.dumps({"text": "Lumini.", "segments": []}), encoding="utf-8"
)
whisper_source_before = whisper_source.read_bytes()
context = context_data()
context["known_entities"] = {
"Authoritative terminology": ["Luminy (aliases: Lumini)"]
}
def generate_without_network(transcript, context_path, **kwargs):
return generate_direct_protocol(
transcript,
context_path,
model=kwargs["model"],
num_ctx=kwargs["num_ctx"],
num_thread=kwargs["num_thread"],
safe_input_token_budget=kwargs["safe_input_token_budget"],
glossary_aliases=kwargs["glossary_aliases"],
model_check=lambda *_: {},
generation_call=lambda *_args, **_kwargs: OllamaGeneration(
raw_response={"response": "# Protocol", "done": True},
text="# Protocol",
client_wall_time_seconds=0.1,
),
)
with (
patch.object(
mvp_api,
"generate_direct_protocol",
side_effect=generate_without_network,
),
):
mvp_api.regenerate_mvp_protocol(
run_dir,
meeting_context=context,
glossary_aliases={"Lumini": "Luminy"},
)
generation = run_dir / "protocol"
transcript_input = (generation / "transcript_input.txt").read_text()
exact_prompt = (generation / "exact_prompt.txt").read_text()
metadata = json.loads((generation / "runtime_metadata.json").read_text())
self.assertIn("SPEAKER_00: Lumini.", transcript_input)
self.assertNotIn("SPEAKER_00: Luminy.", transcript_input)
self.assertIn("Luminy (aliases: Lumini)", exact_prompt)
self.assertIn("SPEAKER_00", exact_prompt)
self.assertIn("Test Person", exact_prompt)
self.assertEqual(metadata["glossary_replacements"], [])
self.assertEqual(source.read_bytes(), source_before)
self.assertEqual(whisper_source.read_bytes(), whisper_source_before)
def test_failure_emits_terminal_failure_event(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
+169
View File
@@ -0,0 +1,169 @@
"""Publication fault injection without inference or audio processing."""
import json
import multiprocessing
import os
import tempfile
import unittest
from pathlib import Path
from unittest.mock import patch
from src.meeting_lab.models.meeting_context import create_meeting_context, load_meeting_context
from src.meeting_lab.protocol import history
from src.meeting_lab.orchestration import mvp
from src.meeting_lab.protocol.generate_direct_protocol import DirectProtocolResult
from tests import test_mvp_api as fixtures
def result(text):
return DirectProtocolResult(
protocol_text=text,
exact_prompt="prompt " + text,
transcript_input="Lumini original",
model_metadata={"name": "test"},
runtime_metadata={
"glossary_aliases_configured": {"Lumini": "Luminy"},
"glossary_replacements": [],
"num_thread": 10,
},
raw_response={"response": text},
)
def interrupt_publication(root):
original = history._replace_link
def interrupted(path, target):
if path.name == "current":
os._exit(77)
original(path, target)
history._replace_link = interrupted
history.persist_protocol_generation(
Path(root), result("new"), context=create_meeting_context(fixtures.context_data())
)
class HistoryTests(unittest.TestCase):
def setUp(self):
self.temporary = tempfile.TemporaryDirectory()
self.addCleanup(self.temporary.cleanup)
self.root = Path(self.temporary.name)
self.context = create_meeting_context(fixtures.context_data())
def persist(self, text):
return history.persist_protocol_generation(self.root, result(text), context=self.context)
def snapshot(self):
names = [
"protocol.md",
"context/meeting_context.yaml",
*["protocol/" + n for n in history._DIAGNOSTICS],
]
return {name: (self.root / name).read_bytes() for name in names}
def test_success_numbering_provenance_and_immutable_context(self):
self.persist("first")
first = self.snapshot()
self.persist("second")
self.assertEqual((self.root / "protocol.md").read_text(), "second")
old = self.root / "protocol/generations/001"
self.assertEqual((old / "protocol.md").read_bytes(), first["protocol.md"])
self.assertEqual(
(old / "meeting_context.yaml").read_bytes(), first["context/meeting_context.yaml"]
)
current = (self.root / "protocol/current").resolve()
self.assertEqual(current.name, "002")
metadata = json.loads((current / "runtime_metadata.json").read_text())
self.assertEqual(metadata["generation_index"], 2)
self.assertEqual(metadata["speaker_mapping"], self.context.speaker_mappings)
self.assertEqual(metadata["glossary_aliases_configured"], {"Lumini": "Luminy"})
self.assertEqual(metadata["glossary_replacements"], [])
self.assertEqual(metadata["num_thread"], 10)
self.assertEqual(metadata["source_run_id"], self.root.name)
self.assertEqual(
json.loads((current / "model_metadata.json").read_text()), {"name": "test"}
)
self.assertIn("sha", metadata["git"]["meeting_lab"])
self.assertEqual(
load_meeting_context(current / "meeting_context.yaml").data, self.context.data
)
for name in history._DIAGNOSTICS:
self.assertEqual((self.root / "protocol" / name).resolve(), current / name)
def test_failure_writing_record_preserves_previous(self):
self.persist("first")
before = self.snapshot()
with patch.object(history, "_write_json", side_effect=OSError("disk full")):
with self.assertRaises(OSError):
self.persist("second")
self.assertEqual(self.snapshot(), before)
self.persist("third")
self.assertEqual((self.root / "protocol/current").resolve().name, "002")
def test_failure_at_atomic_publication_preserves_all_stable_paths(self):
self.persist("first")
before = self.snapshot()
original = history.os.replace
def fail_current(source, destination):
if Path(destination).name == "current":
raise OSError("publication stopped")
original(source, destination)
with patch.object(history.os, "replace", side_effect=fail_current):
with self.assertRaises(OSError):
self.persist("second")
self.assertEqual(self.snapshot(), before)
self.assertFalse((self.root / "protocol/generations/002").exists())
def test_process_interruption_before_pointer_swap_preserves_current(self):
self.persist("first")
before = self.snapshot()
worker = multiprocessing.get_context("fork").Process(
target=interrupt_publication, args=(str(self.root),)
)
worker.start()
worker.join(timeout=10)
self.assertFalse(worker.is_alive())
self.assertEqual(worker.exitcode, 77)
self.assertEqual(self.snapshot(), before)
self.persist("third")
self.assertEqual((self.root / "protocol.md").read_text(), "third")
def test_legacy_migration_and_failure_during_link_conversion(self):
(self.root / "protocol").mkdir()
(self.root / "context").mkdir()
(self.root / "protocol.md").write_text("legacy")
for name in history._DIAGNOSTICS:
(self.root / "protocol" / name).write_text("legacy " + name)
(self.root / "context/meeting_context.yaml").write_text("legacy context")
before = self.snapshot()
original = history._replace_link
def fail_conversion(path, target):
if path.name == "raw_response.json":
raise OSError("link conversion stopped")
original(path, target)
with patch.object(history, "_replace_link", side_effect=fail_conversion):
with self.assertRaises(OSError):
self.persist("new")
self.assertEqual(self.snapshot(), before)
self.persist("new")
self.assertEqual((self.root / "protocol/generations/001/protocol.md").read_text(), "legacy")
self.assertEqual((self.root / "protocol/current").resolve().name, "002")
def test_failed_model_does_not_change_context_or_latest_generation(self):
self.persist("first")
before = self.snapshot()
(self.root / "transcript").mkdir()
(self.root / "transcript/transcript.json").write_text('{"text":"original"}')
different = fixtures.context_data()
different["speaker_mappings"] = {}
with patch.object(
mvp, "generate_direct_protocol", side_effect=RuntimeError("model stopped")
):
with self.assertRaises(RuntimeError):
mvp.regenerate_mvp_protocol(self.root, meeting_context=different)
self.assertEqual(self.snapshot(), before)
+58
View File
@@ -0,0 +1,58 @@
"""Protocol thread configuration through the public pipeline boundaries."""
import tempfile
import unittest
from dataclasses import replace
from pathlib import Path
from unittest.mock import Mock, patch
from src.meeting_lab.llm import ollama
from src.meeting_lab.orchestration import mvp
from tests import test_mvp_api as fixtures
class ProtocolThreadTests(unittest.TestCase):
def test_initial_and_regeneration_forward_default_and_explicit_threads(self):
for threads in (None, 4, 10, 16):
with self.subTest(threads=threads), tempfile.TemporaryDirectory() as directory:
config = replace(
fixtures.MvpApiTests().config(Path(directory)), protocol_num_thread=threads
)
with (
patch.object(mvp, "prepare_audio", side_effect=fixtures.fake_prepare),
patch.object(mvp, "transcribe_audio", side_effect=fixtures.fake_transcribe),
patch.object(
mvp, "generate_direct_protocol", side_effect=fixtures.fake_protocol
) as generate,
):
result = mvp.run_mvp_meeting(config, meeting_context=fixtures.context_data())
self.assertEqual(result.exit_code, 0)
self.assertEqual(generate.call_args.kwargs["num_thread"], threads)
mvp.regenerate_mvp_protocol(
result.run_dir,
meeting_context=fixtures.context_data(),
protocol_num_thread=threads,
)
self.assertEqual(generate.call_args.kwargs["num_thread"], threads)
def test_explicit_payload_and_invalid_values(self):
response = Mock()
response.json.return_value = {"response": "protocol"}
with patch.object(ollama.requests, "post", return_value=response) as post:
ollama.generate_once(
"url", "model", "prompt", timeout=1, num_ctx=32, num_predict=8, num_thread=10
)
self.assertEqual(post.call_args.kwargs["json"]["options"]["num_thread"], 10)
post.reset_mock()
for value in (0, -1, True, 1.5):
with self.assertRaises(ValueError):
ollama.generate_once(
"url",
"model",
"prompt",
timeout=1,
num_ctx=32,
num_predict=8,
num_thread=value,
)
post.assert_not_called()
+37
View File
@@ -137,6 +137,43 @@ class ProtocolInputBudgetTests(unittest.TestCase):
self.assertEqual(result.transcript_input, compact)
self.assertEqual(call.call_count, 1)
def test_glossary_aliases_do_not_mutate_protocol_input_or_source_artifact(
self,
) -> None:
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
document = {
"text": "Lumini und Carbofool.",
"segments": [
{
"start": 0.0,
"end": 1.0,
"speaker_id": "SPEAKER_00",
"text": "Lumini und Carbofool.",
}
],
"speaker_labels_anonymous": True,
}
transcript = self._write(root, document)
source_before = transcript.read_bytes()
result = generate_direct_protocol(
transcript,
glossary_aliases={"Lumini": "Luminy", "Carbofool": "Carbofol"},
model_check=Mock(return_value={}),
generation_call=Mock(return_value=completed_generation()),
)
self.assertEqual(transcript.read_bytes(), source_before)
self.assertIn("SPEAKER_00: Lumini und Carbofool.", result.transcript_input)
self.assertIn("SPEAKER_00: Lumini und Carbofool.", result.exact_prompt)
self.assertEqual(result.runtime_metadata["glossary_replacements"], [])
self.assertEqual(
result.runtime_metadata["glossary_aliases_configured"],
{"Carbofool": "Carbofol", "Lumini": "Luminy"},
)
def test_mapped_speakers_and_statements_reach_final_ollama_payload(self) -> None:
diarized = {
"text": "",