Publish complete protocol generations atomically
This commit is contained in:
@@ -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"
|
||||
Reference in New Issue
Block a user