From d2e3636b949280402728308022c99bdee6ea1066 Mon Sep 17 00:00:00 2001 From: Martin Date: Sat, 12 Sep 2026 11:49:10 +0200 Subject: [PATCH] Publish complete protocol generations atomically --- PROJECT_KNOWLEDGE.md | 12 ++ README.md | 12 ++ scripts/run_direct_protocol.py | 25 +--- src/meeting_lab/orchestration/mvp.py | 50 +++---- src/meeting_lab/protocol/history.py | 188 +++++++++++++++++++++++++++ tests/test_direct_protocol.py | 3 + tests/test_mvp_api.py | 1 + tests/test_protocol_history.py | 169 ++++++++++++++++++++++++ 8 files changed, 411 insertions(+), 49 deletions(-) create mode 100644 src/meeting_lab/protocol/history.py create mode 100644 tests/test_protocol_history.py diff --git a/PROJECT_KNOWLEDGE.md b/PROJECT_KNOWLEDGE.md index a8c89e9..c6ee923 100644 --- a/PROJECT_KNOWLEDGE.md +++ b/PROJECT_KNOWLEDGE.md @@ -303,3 +303,15 @@ both initial processing and regeneration. With `None`, Ollama receives no 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. diff --git a/README.md b/README.md index ffe12bb..f42c0a1 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/scripts/run_direct_protocol.py b/scripts/run_direct_protocol.py index d3a4f4d..60a43f4 100644 --- a/scripts/run_direct_protocol.py +++ b/scripts/run_direct_protocol.py @@ -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: diff --git a/src/meeting_lab/orchestration/mvp.py b/src/meeting_lab/orchestration/mvp.py index a44f501..c520dd9 100644 --- a/src/meeting_lab/orchestration/mvp.py +++ b/src/meeting_lab/orchestration/mvp.py @@ -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, ) @@ -110,21 +111,24 @@ 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, - num_thread=protocol_num_thread, - glossary_aliases=glossary_aliases, - 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: _emit( progress_sink, @@ -228,22 +232,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, *, @@ -445,7 +433,7 @@ def run_mvp_meeting( 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) diff --git a/src/meeting_lab/protocol/history.py b/src/meeting_lab/protocol/history.py new file mode 100644 index 0000000..d59b088 --- /dev/null +++ b/src/meeting_lab/protocol/history.py @@ -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" diff --git a/tests/test_direct_protocol.py b/tests/test_direct_protocol.py index 850ba64..eef9953 100644 --- a/tests/test_direct_protocol.py +++ b/tests/test_direct_protocol.py @@ -254,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) @@ -321,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, diff --git a/tests/test_mvp_api.py b/tests/test_mvp_api.py index f5aec1b..bc7d561 100644 --- a/tests/test_mvp_api.py +++ b/tests/test_mvp_api.py @@ -118,6 +118,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( diff --git a/tests/test_protocol_history.py b/tests/test_protocol_history.py new file mode 100644 index 0000000..4d10e07 --- /dev/null +++ b/tests/test_protocol_history.py @@ -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)