Stabilize Meeting Lab pipeline for RC1 evaluation
This commit significantly improves the robustness and determinism of the Meeting Lab processing pipeline and establishes the first Release Candidate baseline for end-to-end evaluation. Highlights - BUG-009 - Implement deterministic responsible-party validation - Normalize participant aliases using Meeting Context - Reject invalid responsible values (dates, locations, technical terms, projects, products, unknown entities) - Record structured responsibility validation metadata - Add focused regression tests - BUG-010 - Implement adaptive num_predict estimation for Semantic Consolidator - Eliminate JSON truncation caused by fixed output limits - Add deterministic source coverage repair - Preserve strict post-repair validation - Add regression tests - BUG-011 - Implement Working Protocol V2 renderer contract enforcement - Preserve raw renderer responses - Reject invalid protocol output instead of accepting malformed documents - Add deterministic cleanup for harmless formatting deviations - Add focused renderer regression tests - Meeting Context - Validate Meeting Context V1 - Integrate authoritative participant alias normalization - Documentation - Update architecture documentation - Update output documentation - Update regression bug tracker The pipeline now fails safely instead of silently accepting invalid intermediate or final artifacts. Remaining work focuses primarily on extraction quality and semantic classification (decisions, action items, protocol faithfulness), rather than pipeline robustness.
This commit is contained in:
@@ -8,9 +8,12 @@ import json
|
||||
import re
|
||||
import sys
|
||||
from collections import Counter
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from src.meeting_lab.models.meeting_context import load_meeting_context
|
||||
|
||||
|
||||
SCHEMA_VERSION = "1"
|
||||
|
||||
@@ -36,6 +39,51 @@ CANONICAL_CATEGORIES = tuple(CATEGORY_NAMES.values())
|
||||
|
||||
CHUNK_EXTRACTION_RE = re.compile(r"^chunk_(\d+)_extraction\.json$")
|
||||
WHITESPACE_RE = re.compile(r"\s+")
|
||||
RESPONSIBLE_KEY_RE = re.compile(r"\s+")
|
||||
NUMERIC_DATE_RE = re.compile(r"(?i)\b\d{1,2}\s*[./]\s*(?:\d{1,2}|[a-zäöü]+)?\b")
|
||||
MONTH_DATE_RE = re.compile(
|
||||
r"(?i)\b\d{1,2}\.?\s*(?:oder\s+\d{1,2}\.?\s*)?"
|
||||
r"(januar|februar|märz|maerz|april|mai|juni|juli|august|september|oktober|november|dezember)\b"
|
||||
)
|
||||
TIME_RE = re.compile(r"(?i)\b\d{1,2}[:.]\d{2}\s*(?:uhr)?\b|\b\d{1,2}\s*uhr\b")
|
||||
DATE_FRAGMENT_RE = re.compile(r"(?i)\b(?:am|zum|bis|vor|nach)\s+\d{1,2}\.?\b")
|
||||
|
||||
DATE_WORDS = {
|
||||
"heute",
|
||||
"morgen",
|
||||
"übermorgen",
|
||||
"uebermorgen",
|
||||
"gestern",
|
||||
"vorgestern",
|
||||
}
|
||||
WEEKDAY_WORDS = {
|
||||
"montag",
|
||||
"dienstag",
|
||||
"mittwoch",
|
||||
"donnerstag",
|
||||
"freitag",
|
||||
"samstag",
|
||||
"sonntag",
|
||||
}
|
||||
RELATIVE_DATE_PHRASES = {
|
||||
"nächste woche",
|
||||
"naechste woche",
|
||||
"diese woche",
|
||||
"kommende woche",
|
||||
"nächsten monat",
|
||||
"naechsten monat",
|
||||
}
|
||||
GENERIC_RESPONSIBLE_WORDS = {
|
||||
"team",
|
||||
"projektteam",
|
||||
"projektleitung",
|
||||
"logistik",
|
||||
"alle",
|
||||
"autor",
|
||||
"labor",
|
||||
"partner",
|
||||
"gruppe",
|
||||
}
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
@@ -59,9 +107,110 @@ def parse_args() -> argparse.Namespace:
|
||||
action="store_true",
|
||||
help="Preserve exact duplicate items instead of merging them.",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--meeting-context",
|
||||
type=Path,
|
||||
help=(
|
||||
"Optional Meeting Context V1 YAML used to validate and normalize "
|
||||
"action-item responsible fields. If omitted, canonicalizer uses "
|
||||
"context provenance from extraction JSON when available."
|
||||
),
|
||||
)
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ResponsiblePartyNormalizer:
|
||||
allowed_names: dict[str, str]
|
||||
invalid_values: dict[str, str]
|
||||
|
||||
@classmethod
|
||||
def from_meeting_context_data(cls, data: dict[str, Any]) -> "ResponsiblePartyNormalizer":
|
||||
allowed: dict[str, str] = {}
|
||||
invalid: dict[str, str] = {}
|
||||
|
||||
for collection, id_key in (
|
||||
("participants", "participant_id"),
|
||||
("mentioned_people", "person_id"),
|
||||
):
|
||||
for person in data.get(collection, []) or []:
|
||||
if not isinstance(person, dict):
|
||||
continue
|
||||
display_name = clean_text(person.get("display_name"))
|
||||
if not display_name:
|
||||
continue
|
||||
for value in [display_name, *(person.get("aliases") or [])]:
|
||||
text = clean_text(value)
|
||||
if text:
|
||||
allowed[responsible_key(text)] = display_name
|
||||
|
||||
organization = data.get("organization")
|
||||
if isinstance(organization, dict):
|
||||
organization_name = clean_text(organization.get("name"))
|
||||
if organization_name:
|
||||
allowed[responsible_key(organization_name)] = organization_name
|
||||
|
||||
for department in organization.get("departments") or []:
|
||||
if not isinstance(department, dict):
|
||||
continue
|
||||
department_name = clean_text(department.get("name"))
|
||||
if department_name:
|
||||
allowed[responsible_key(department_name)] = department_name
|
||||
for alias in department.get("aliases") or []:
|
||||
text = clean_text(alias)
|
||||
if text and department_name:
|
||||
allowed[responsible_key(text)] = department_name
|
||||
|
||||
known_entities = data.get("known_entities")
|
||||
if isinstance(known_entities, dict):
|
||||
for collection in (
|
||||
"projects",
|
||||
"products",
|
||||
"systems",
|
||||
"locations",
|
||||
"technical_terms",
|
||||
):
|
||||
for value in known_entities.get(collection) or []:
|
||||
text = clean_text(value)
|
||||
if text:
|
||||
invalid[responsible_key(text)] = collection
|
||||
|
||||
return cls(allowed_names=allowed, invalid_values=invalid)
|
||||
|
||||
def normalize(
|
||||
self,
|
||||
value: str | None,
|
||||
item_id: str,
|
||||
) -> tuple[str | None, dict[str, Any] | None]:
|
||||
original = clean_text(value)
|
||||
if original is None:
|
||||
return None, None
|
||||
|
||||
key = responsible_key(original)
|
||||
canonical = self.allowed_names.get(key)
|
||||
if canonical is not None:
|
||||
return canonical, {
|
||||
"action_item_id": item_id,
|
||||
"field": "responsible",
|
||||
"original_value": original,
|
||||
"normalized_value": canonical,
|
||||
"status": "accepted",
|
||||
"reason": "resolved_to_meeting_context_entity",
|
||||
"cleared_to_null": False,
|
||||
}
|
||||
|
||||
reason = invalid_responsible_reason(original, self.invalid_values)
|
||||
return None, {
|
||||
"action_item_id": item_id,
|
||||
"field": "responsible",
|
||||
"original_value": original,
|
||||
"normalized_value": None,
|
||||
"status": "rejected",
|
||||
"reason": reason,
|
||||
"cleared_to_null": True,
|
||||
}
|
||||
|
||||
|
||||
def chunk_sort_key(path: Path) -> tuple[int, str]:
|
||||
match = CHUNK_EXTRACTION_RE.match(path.name)
|
||||
if not match:
|
||||
@@ -109,6 +258,30 @@ def clean_text(value: Any) -> str | None:
|
||||
return text if text else None
|
||||
|
||||
|
||||
def responsible_key(value: str) -> str:
|
||||
return RESPONSIBLE_KEY_RE.sub(" ", value.casefold()).strip(" .,:;")
|
||||
|
||||
|
||||
def invalid_responsible_reason(
|
||||
value: str,
|
||||
invalid_values: dict[str, str],
|
||||
) -> str:
|
||||
key = responsible_key(value)
|
||||
if key in invalid_values:
|
||||
return f"known_{invalid_values[key]}_not_responsible_entity"
|
||||
if key in GENERIC_RESPONSIBLE_WORDS:
|
||||
return "generic_process_word"
|
||||
if key in DATE_WORDS or key in WEEKDAY_WORDS or key in RELATIVE_DATE_PHRASES:
|
||||
return "date_or_relative_date"
|
||||
if NUMERIC_DATE_RE.search(value) or MONTH_DATE_RE.search(value):
|
||||
return "date_or_date_range"
|
||||
if TIME_RE.search(value):
|
||||
return "clock_time"
|
||||
if DATE_FRAGMENT_RE.search(value):
|
||||
return "date_fragment"
|
||||
return "unknown_responsible_entity"
|
||||
|
||||
|
||||
def split_legacy_string(value: str) -> list[str]:
|
||||
return [part.strip() for part in value.split("|")]
|
||||
|
||||
@@ -324,6 +497,7 @@ def canonicalize_value(
|
||||
source_file: str,
|
||||
source_index: int,
|
||||
counts: Counter[str],
|
||||
responsible_normalizer: ResponsiblePartyNormalizer | None = None,
|
||||
) -> dict[str, Any]:
|
||||
parsed = PARSERS[category](value)
|
||||
item: dict[str, Any] = {
|
||||
@@ -337,6 +511,14 @@ def canonicalize_value(
|
||||
}
|
||||
for key, parsed_value in parsed.items():
|
||||
item[key] = parsed_value
|
||||
if category == "action_item" and responsible_normalizer is not None:
|
||||
normalized, validation = responsible_normalizer.normalize(
|
||||
item.get("responsible"),
|
||||
item["item_id"],
|
||||
)
|
||||
item["responsible"] = normalized
|
||||
if validation is not None:
|
||||
item["responsibility_validation"] = validation
|
||||
item["source_references"] = [
|
||||
source_reference(source_file, source_index, value, item["evidence"])
|
||||
]
|
||||
@@ -367,6 +549,7 @@ def merge_exact_duplicates(items: list[dict[str, Any]]) -> tuple[list[dict[str,
|
||||
def canonicalize_extractions(
|
||||
input_dir: Path,
|
||||
merge_duplicates: bool = True,
|
||||
meeting_context_path: Path | None = None,
|
||||
) -> dict[str, Any]:
|
||||
files = find_extraction_files(input_dir)
|
||||
if not files:
|
||||
@@ -375,9 +558,22 @@ def canonicalize_extractions(
|
||||
counts: Counter[str] = Counter()
|
||||
input_counts: Counter[str] = Counter()
|
||||
items: list[dict[str, Any]] = []
|
||||
responsible_normalizer: ResponsiblePartyNormalizer | None = None
|
||||
responsibility_validations: list[dict[str, Any]] = []
|
||||
|
||||
loaded_files: list[tuple[Path, dict[str, Any]]] = []
|
||||
for path in files:
|
||||
data = load_json_object(path)
|
||||
loaded_files.append((path, data))
|
||||
|
||||
context_path = meeting_context_path or infer_meeting_context_path(loaded_files)
|
||||
if context_path is not None:
|
||||
meeting_context = load_meeting_context(context_path)
|
||||
responsible_normalizer = ResponsiblePartyNormalizer.from_meeting_context_data(
|
||||
meeting_context.data
|
||||
)
|
||||
|
||||
for path, data in loaded_files:
|
||||
validate_required_categories(data, path)
|
||||
for raw_category in REQUIRED_CATEGORIES:
|
||||
category = CATEGORY_NAMES[raw_category]
|
||||
@@ -391,8 +587,16 @@ def canonicalize_extractions(
|
||||
source_file=path.name,
|
||||
source_index=source_index,
|
||||
counts=counts,
|
||||
responsible_normalizer=responsible_normalizer,
|
||||
)
|
||||
)
|
||||
if (
|
||||
items[-1].get("category") == "action_item"
|
||||
and "responsibility_validation" in items[-1]
|
||||
):
|
||||
responsibility_validations.append(
|
||||
items[-1]["responsibility_validation"]
|
||||
)
|
||||
|
||||
exact_duplicates = 0
|
||||
if merge_duplicates:
|
||||
@@ -415,11 +619,40 @@ def canonicalize_extractions(
|
||||
"input_item_count_by_category": input_counts_by_category,
|
||||
"output_item_count_by_category": output_counts_by_category,
|
||||
"exact_duplicates_merged": exact_duplicates,
|
||||
"responsibility_validation_count": len(responsibility_validations),
|
||||
"responsibility_rejection_count": sum(
|
||||
1
|
||||
for validation in responsibility_validations
|
||||
if validation.get("status") == "rejected"
|
||||
),
|
||||
"responsibility_normalization_count": sum(
|
||||
1
|
||||
for validation in responsibility_validations
|
||||
if validation.get("status") == "accepted"
|
||||
and validation.get("original_value") != validation.get("normalized_value")
|
||||
),
|
||||
},
|
||||
"responsibility_validations": responsibility_validations,
|
||||
"items": items,
|
||||
}
|
||||
|
||||
|
||||
def infer_meeting_context_path(
|
||||
loaded_files: list[tuple[Path, dict[str, Any]]],
|
||||
) -> Path | None:
|
||||
paths: set[str] = set()
|
||||
for _path, data in loaded_files:
|
||||
context = data.get("context")
|
||||
if not isinstance(context, dict):
|
||||
continue
|
||||
source_file = clean_text(context.get("source_file"))
|
||||
if source_file:
|
||||
paths.add(source_file)
|
||||
if len(paths) != 1:
|
||||
return None
|
||||
return Path(next(iter(paths)))
|
||||
|
||||
|
||||
def write_canonicalized(output: dict[str, Any], output_path: Path) -> Path:
|
||||
output_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
output_path.write_text(
|
||||
@@ -435,6 +668,7 @@ def main() -> int:
|
||||
output = canonicalize_extractions(
|
||||
args.input_dir,
|
||||
merge_duplicates=not args.no_merge_exact_duplicates,
|
||||
meeting_context_path=args.meeting_context,
|
||||
)
|
||||
output_path = write_canonicalized(output, args.output)
|
||||
except (OSError, UnicodeError, ValueError) as exc:
|
||||
|
||||
@@ -4,6 +4,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import copy
|
||||
import json
|
||||
import sys
|
||||
import time
|
||||
@@ -25,8 +26,13 @@ except ModuleNotFoundError: # pragma: no cover - used by repository-root tests.
|
||||
DEFAULT_MODEL = "qwen3.5:9b"
|
||||
DEFAULT_ENDPOINT = "http://127.0.0.1:11434/api/generate"
|
||||
DEFAULT_NUM_CTX = 32768
|
||||
DEFAULT_NUM_PREDICT = 4096
|
||||
DEFAULT_MIN_NUM_PREDICT = 4096
|
||||
DEFAULT_NUM_PREDICT = DEFAULT_MIN_NUM_PREDICT
|
||||
DEFAULT_PROGRESS_INTERVAL = 30
|
||||
OUTPUT_CONTEXT_RESERVE_TOKENS = 1024
|
||||
OUTPUT_TOKEN_ESTIMATE_CHARS = 4
|
||||
OUTPUT_GROUP_OVERHEAD_CHARS = 320
|
||||
OUTPUT_SAFETY_MARGIN = 1.35
|
||||
PROMPT_NAME = "consolidate_facts.md"
|
||||
|
||||
|
||||
@@ -75,10 +81,11 @@ def parse_args() -> argparse.Namespace:
|
||||
parser.add_argument(
|
||||
"--num-predict",
|
||||
type=int,
|
||||
default=DEFAULT_NUM_PREDICT,
|
||||
default=None,
|
||||
help=(
|
||||
"Maximum generated tokens. The default is bounded for the expected "
|
||||
f"fact-group JSON while leaving truncation headroom (default: {DEFAULT_NUM_PREDICT})."
|
||||
"Maximum generated tokens. By default this is estimated from the "
|
||||
"fact payload size and bounded by the context window. Explicit "
|
||||
"values preserve the previous fixed-budget behavior."
|
||||
),
|
||||
)
|
||||
thinking = parser.add_mutually_exclusive_group()
|
||||
@@ -158,6 +165,49 @@ def build_consolidation_prompt(facts: list[dict[str, Any]]) -> str:
|
||||
return f"{task_prompt}\n\nFACT ITEMS:\n{payload}\n"
|
||||
|
||||
|
||||
def estimate_response_tokens(facts: list[dict[str, Any]]) -> int:
|
||||
"""
|
||||
Estimate the token budget needed for the model's grouping JSON.
|
||||
|
||||
Semantic Consolidator V0 asks the model to return one group per source fact
|
||||
unless it finds a conservative duplicate. The response therefore scales with
|
||||
the number and text size of fact items. The estimate intentionally includes
|
||||
per-group JSON overhead and a safety margin; strict validation still decides
|
||||
whether the actual response is usable.
|
||||
"""
|
||||
text_chars = 0
|
||||
for item in facts:
|
||||
text_chars += len(str(item.get("text", "")))
|
||||
text_chars += len(str(item.get("evidence", "")))
|
||||
|
||||
estimated_chars = int(
|
||||
(text_chars + len(facts) * OUTPUT_GROUP_OVERHEAD_CHARS)
|
||||
* OUTPUT_SAFETY_MARGIN
|
||||
)
|
||||
return max(
|
||||
DEFAULT_MIN_NUM_PREDICT,
|
||||
(estimated_chars + OUTPUT_TOKEN_ESTIMATE_CHARS - 1)
|
||||
// OUTPUT_TOKEN_ESTIMATE_CHARS,
|
||||
)
|
||||
|
||||
|
||||
def resolve_num_predict(
|
||||
requested_num_predict: int | None,
|
||||
facts: list[dict[str, Any]],
|
||||
prompt_token_estimate: int,
|
||||
num_ctx: int,
|
||||
) -> int:
|
||||
if requested_num_predict is not None:
|
||||
return requested_num_predict
|
||||
|
||||
estimated = estimate_response_tokens(facts)
|
||||
max_available = max(
|
||||
DEFAULT_MIN_NUM_PREDICT,
|
||||
num_ctx - prompt_token_estimate - OUTPUT_CONTEXT_RESERVE_TOKENS,
|
||||
)
|
||||
return min(estimated, max_available)
|
||||
|
||||
|
||||
def response_text_from_ollama_data(data: dict[str, Any]) -> str | None:
|
||||
text = data.get("response")
|
||||
if isinstance(text, str) and text.strip():
|
||||
@@ -370,6 +420,96 @@ def validate_group_shapes(groups: list[dict[str, Any]]) -> None:
|
||||
raise ConsolidationValidationError("Merged groups need at least two IDs.")
|
||||
|
||||
|
||||
def repair_model_group_coverage(
|
||||
model_output: dict[str, Any],
|
||||
facts: list[dict[str, Any]],
|
||||
) -> tuple[dict[str, Any], list[dict[str, Any]]]:
|
||||
"""
|
||||
Apply deterministic source-coverage repairs to model grouping JSON.
|
||||
|
||||
The repair is intentionally conservative. Repeated source IDs are removed
|
||||
after their first occurrence, empty groups created by that removal are
|
||||
dropped, and missing facts are restored as singleton groups using the
|
||||
original canonicalized fact text. No existing group text, merge reason or
|
||||
semantic merge is rewritten.
|
||||
"""
|
||||
groups = model_output.get("groups")
|
||||
if not isinstance(groups, list):
|
||||
return model_output, []
|
||||
|
||||
repaired = copy.deepcopy(model_output)
|
||||
repaired_groups = repaired["groups"]
|
||||
facts_by_id = {str(item.get("item_id")): item for item in facts}
|
||||
expected_ids = set(facts_by_id)
|
||||
seen: set[str] = set()
|
||||
changes: list[dict[str, Any]] = []
|
||||
|
||||
for group_index, group in enumerate(repaired_groups):
|
||||
if not isinstance(group, dict):
|
||||
continue
|
||||
source_ids = group.get("source_item_ids")
|
||||
if not isinstance(source_ids, list):
|
||||
continue
|
||||
kept_ids: list[str] = []
|
||||
for id_index, item_id in enumerate(source_ids):
|
||||
if not isinstance(item_id, str) or item_id not in expected_ids:
|
||||
kept_ids.append(item_id)
|
||||
continue
|
||||
if item_id in seen:
|
||||
changes.append(
|
||||
{
|
||||
"operation": "remove_duplicate_source_id",
|
||||
"id": item_id,
|
||||
"group_index": group_index,
|
||||
"id_index": id_index,
|
||||
}
|
||||
)
|
||||
continue
|
||||
seen.add(item_id)
|
||||
kept_ids.append(item_id)
|
||||
group["source_item_ids"] = kept_ids
|
||||
|
||||
non_empty_groups: list[dict[str, Any]] = []
|
||||
for group_index, group in enumerate(repaired_groups):
|
||||
if (
|
||||
isinstance(group, dict)
|
||||
and isinstance(group.get("source_item_ids"), list)
|
||||
and len(group["source_item_ids"]) == 0
|
||||
):
|
||||
changes.append(
|
||||
{
|
||||
"operation": "remove_empty_group",
|
||||
"group_index": group_index,
|
||||
"canonical_text": group.get("canonical_text"),
|
||||
}
|
||||
)
|
||||
continue
|
||||
non_empty_groups.append(group)
|
||||
repaired["groups"] = non_empty_groups
|
||||
|
||||
missing_ids = sorted(expected_ids - seen)
|
||||
for item_id in missing_ids:
|
||||
fact = facts_by_id[item_id]
|
||||
repaired["groups"].append(
|
||||
{
|
||||
"canonical_text": str(fact.get("text", "")).strip(),
|
||||
"source_item_ids": [item_id],
|
||||
"merge_reason": (
|
||||
"Deterministic coverage repair: source fact was missing "
|
||||
"from the model grouping and is preserved as a singleton."
|
||||
),
|
||||
}
|
||||
)
|
||||
changes.append(
|
||||
{
|
||||
"operation": "restore_missing_source_id_as_singleton",
|
||||
"id": item_id,
|
||||
}
|
||||
)
|
||||
|
||||
return repaired, changes
|
||||
|
||||
|
||||
def build_consolidated_fact_item(
|
||||
group: dict[str, Any],
|
||||
fact_by_id: dict[str, dict[str, Any]],
|
||||
@@ -481,9 +621,11 @@ def write_report(
|
||||
runtime: float,
|
||||
prompt_chars: int,
|
||||
prompt_token_estimate: int,
|
||||
num_predict: int,
|
||||
fact_count: int,
|
||||
groups: list[dict[str, Any]],
|
||||
output_path: Path,
|
||||
repair_changes: list[dict[str, Any]] | None = None,
|
||||
) -> None:
|
||||
merged = [group for group in groups if len(group["source_item_ids"]) > 1]
|
||||
singletons = [group for group in groups if len(group["source_item_ids"]) == 1]
|
||||
@@ -497,9 +639,11 @@ def write_report(
|
||||
f"- Fact item count: {fact_count}",
|
||||
f"- Prompt characters: {prompt_chars}",
|
||||
f"- Estimated prompt tokens: {prompt_token_estimate}",
|
||||
f"- num_predict: {num_predict}",
|
||||
f"- Merged fact groups: {len(merged)}",
|
||||
f"- Source facts involved in merges: {sum(len(group['source_item_ids']) for group in merged)}",
|
||||
f"- Singleton fact groups: {len(singletons)}",
|
||||
f"- Deterministic repair changes: {len(repair_changes or [])}",
|
||||
f"- Output path: `{output_path}`",
|
||||
"",
|
||||
"## Actual Merges",
|
||||
@@ -518,6 +662,10 @@ def write_report(
|
||||
"",
|
||||
]
|
||||
)
|
||||
if repair_changes:
|
||||
lines.extend(["", "## Deterministic Coverage Repairs", ""])
|
||||
for change in repair_changes:
|
||||
lines.append(f"- `{change['operation']}`: {json.dumps(change, ensure_ascii=False, sort_keys=True)}")
|
||||
path.write_text("\n".join(lines).rstrip() + "\n", encoding="utf-8")
|
||||
|
||||
|
||||
@@ -527,6 +675,7 @@ def main() -> int:
|
||||
raw_response_path = args.output_dir / "raw_model_response.txt"
|
||||
output_path = args.output_dir / "consolidated_extractions.json"
|
||||
report_path = args.output_dir / "report.md"
|
||||
repair_metadata_path = args.output_dir / "repair_metadata.json"
|
||||
|
||||
try:
|
||||
canonicalized = load_json_object(args.canonicalized_input)
|
||||
@@ -534,9 +683,16 @@ def main() -> int:
|
||||
prompt = build_consolidation_prompt(facts)
|
||||
prompt_chars = len(prompt)
|
||||
prompt_token_estimate = (prompt_chars + 3) // 4
|
||||
num_predict = resolve_num_predict(
|
||||
requested_num_predict=args.num_predict,
|
||||
facts=facts,
|
||||
prompt_token_estimate=prompt_token_estimate,
|
||||
num_ctx=args.num_ctx,
|
||||
)
|
||||
print(f"Fact item count: {len(facts)}")
|
||||
print(f"Estimated prompt size chars: {prompt_chars}")
|
||||
print(f"Estimated prompt tokens: {prompt_token_estimate}")
|
||||
print(f"Resolved num_predict: {num_predict}")
|
||||
print("Expected LLM call count: 1")
|
||||
print("Expected runtime: 5-10 minutes on current local benchmark basis")
|
||||
|
||||
@@ -546,12 +702,23 @@ def main() -> int:
|
||||
prompt=prompt,
|
||||
timeout=args.timeout,
|
||||
num_ctx=args.num_ctx,
|
||||
num_predict=args.num_predict,
|
||||
num_predict=num_predict,
|
||||
think=args.think,
|
||||
progress_interval=args.progress_interval,
|
||||
)
|
||||
raw_response_path.write_text(raw_text + "\n", encoding="utf-8")
|
||||
model_output = parse_model_json(raw_text)
|
||||
model_output, repair_changes = repair_model_group_coverage(model_output, facts)
|
||||
if repair_changes:
|
||||
write_json(
|
||||
repair_metadata_path,
|
||||
{
|
||||
"scope": "semantic_consolidator_v0_source_coverage",
|
||||
"llm_used": False,
|
||||
"repair_count": len(repair_changes),
|
||||
"repairs": repair_changes,
|
||||
},
|
||||
)
|
||||
expected_fact_ids = {item["item_id"] for item in facts}
|
||||
groups = validate_model_groups(model_output, expected_fact_ids)
|
||||
validate_group_shapes(groups)
|
||||
@@ -564,9 +731,11 @@ def main() -> int:
|
||||
runtime=runtime,
|
||||
prompt_chars=prompt_chars,
|
||||
prompt_token_estimate=prompt_token_estimate,
|
||||
num_predict=num_predict,
|
||||
fact_count=len(facts),
|
||||
groups=groups,
|
||||
output_path=output_path,
|
||||
repair_changes=repair_changes,
|
||||
)
|
||||
except requests.ConnectionError as exc:
|
||||
print(f"Error: Ollama is not reachable at {args.endpoint}: {exc}", file=sys.stderr)
|
||||
|
||||
@@ -0,0 +1,519 @@
|
||||
#!/usr/bin/env python3
|
||||
"""LLM-backed Working Protocol V2 renderer with deterministic contract checks."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import requests
|
||||
|
||||
try:
|
||||
from meeting_lab.llm.prompts import load_prompt
|
||||
except ModuleNotFoundError: # pragma: no cover - used by repository-root tests.
|
||||
from src.meeting_lab.llm.prompts import load_prompt
|
||||
|
||||
|
||||
DEFAULT_MODEL = "qwen3.5:9b"
|
||||
DEFAULT_ENDPOINT = "http://127.0.0.1:11434/api/generate"
|
||||
DEFAULT_NUM_CTX = 32768
|
||||
DEFAULT_NUM_PREDICT = 4096
|
||||
DEFAULT_TIMEOUT = 1800
|
||||
PROMPT_NAME = "working_protocol.md"
|
||||
REQUIRED_TITLE = "# Working Protocol"
|
||||
ALLOWED_TOPIC_SECTIONS = {
|
||||
"Background",
|
||||
"Decisions",
|
||||
"Action Items",
|
||||
"Open Questions",
|
||||
}
|
||||
GENERIC_SUMMARY_HEADINGS = {
|
||||
"Entscheidungen",
|
||||
"Decisions",
|
||||
"Handlungsaufträge",
|
||||
"Handlungsauftraege",
|
||||
"Action Items",
|
||||
"Offene Fragen",
|
||||
"Open Questions",
|
||||
"Technische Details",
|
||||
"Technical Details",
|
||||
"Fakten",
|
||||
"Facts",
|
||||
"Zusammenfassung",
|
||||
"Summary",
|
||||
"Konsolidierter Projektstatusbericht",
|
||||
}
|
||||
HEADING_RE = re.compile(r"^(#{1,6})\s+(.+?)\s*$")
|
||||
|
||||
|
||||
class WorkingProtocolValidationError(ValueError):
|
||||
"""Raised when rendered Markdown violates the Working Protocol V2 contract."""
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class WorkingProtocolValidationReport:
|
||||
valid: bool
|
||||
violations: list[dict[str, Any]]
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
return {"valid": self.valid, "violations": self.violations}
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Render a consolidated meeting representation as Working Protocol V2."
|
||||
)
|
||||
parser.add_argument("input", type=Path, help="Consolidated meeting JSON input.")
|
||||
parser.add_argument(
|
||||
"-o",
|
||||
"--output-dir",
|
||||
type=Path,
|
||||
required=True,
|
||||
help="Directory for working_protocol.md, raw response, validation and metadata.",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--model",
|
||||
default=DEFAULT_MODEL,
|
||||
help=f"Ollama model name (default: {DEFAULT_MODEL}).",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--endpoint",
|
||||
default=DEFAULT_ENDPOINT,
|
||||
help=f"Ollama generate endpoint (default: {DEFAULT_ENDPOINT}).",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--timeout",
|
||||
type=int,
|
||||
default=DEFAULT_TIMEOUT,
|
||||
help=f"HTTP timeout in seconds (default: {DEFAULT_TIMEOUT}).",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--num-ctx",
|
||||
type=int,
|
||||
default=DEFAULT_NUM_CTX,
|
||||
help=f"Context window tokens (default: {DEFAULT_NUM_CTX}).",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--num-predict",
|
||||
type=int,
|
||||
default=DEFAULT_NUM_PREDICT,
|
||||
help=f"Maximum generated tokens (default: {DEFAULT_NUM_PREDICT}).",
|
||||
)
|
||||
thinking = parser.add_mutually_exclusive_group()
|
||||
thinking.add_argument("--think", dest="think", action="store_true")
|
||||
thinking.add_argument("--no-think", dest="think", action="store_false")
|
||||
parser.set_defaults(think=False)
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def build_renderer_prompt(input_text: str) -> str:
|
||||
return f"{load_prompt(PROMPT_NAME)}\n\nINPUT JSON:\n{input_text}\n"
|
||||
|
||||
|
||||
def response_text_from_ollama_data(data: dict[str, Any]) -> str:
|
||||
text = data.get("response")
|
||||
if isinstance(text, str):
|
||||
return text
|
||||
message = data.get("message")
|
||||
if isinstance(message, dict) and isinstance(message.get("content"), str):
|
||||
return message["content"]
|
||||
return ""
|
||||
|
||||
|
||||
def build_ollama_payload(
|
||||
model: str,
|
||||
prompt: str,
|
||||
num_ctx: int,
|
||||
num_predict: int,
|
||||
think: bool,
|
||||
) -> dict[str, Any]:
|
||||
return {
|
||||
"model": model,
|
||||
"prompt": prompt,
|
||||
"think": think,
|
||||
"stream": False,
|
||||
"options": {
|
||||
"temperature": 0.0,
|
||||
"num_ctx": num_ctx,
|
||||
"num_predict": num_predict,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def call_ollama(
|
||||
endpoint: str,
|
||||
payload: dict[str, Any],
|
||||
timeout: int,
|
||||
) -> tuple[dict[str, Any], float]:
|
||||
start = time.perf_counter()
|
||||
response = requests.post(endpoint, json=payload, timeout=timeout)
|
||||
runtime = time.perf_counter() - start
|
||||
response.raise_for_status()
|
||||
data = response.json()
|
||||
if not isinstance(data, dict):
|
||||
raise ValueError("Ollama response must be a JSON object.")
|
||||
return data, runtime
|
||||
|
||||
|
||||
def normalize_heading_whitespace(line: str) -> str:
|
||||
match = HEADING_RE.match(line.strip())
|
||||
if not match:
|
||||
return line.rstrip()
|
||||
return f"{match.group(1)} {match.group(2).strip()}"
|
||||
|
||||
|
||||
def clean_working_protocol_markdown(text: str) -> str:
|
||||
"""
|
||||
Remove only non-semantic wrapper text before an existing protocol heading.
|
||||
|
||||
The cleanup does not create headings or transform a categorized summary into
|
||||
a Working Protocol. It only makes an already present contract heading the
|
||||
first byte of the candidate document.
|
||||
"""
|
||||
text = text.replace("\r\n", "\n").replace("\r", "\n")
|
||||
lines = text.split("\n")
|
||||
start_index: int | None = None
|
||||
for index, line in enumerate(lines):
|
||||
if normalize_heading_whitespace(line) == REQUIRED_TITLE:
|
||||
start_index = index
|
||||
break
|
||||
if start_index is None:
|
||||
return text.lstrip()
|
||||
|
||||
cleaned_lines = lines[start_index:]
|
||||
if cleaned_lines:
|
||||
cleaned_lines[0] = REQUIRED_TITLE
|
||||
return "\n".join(normalize_heading_whitespace(line) for line in cleaned_lines).strip() + "\n"
|
||||
|
||||
|
||||
def validate_markdown_shape(markdown: str) -> list[dict[str, Any]]:
|
||||
violations: list[dict[str, Any]] = []
|
||||
if markdown.count("```") % 2 != 0:
|
||||
violations.append(
|
||||
{
|
||||
"type": "malformed_markdown",
|
||||
"reason": "unclosed_fenced_code_block",
|
||||
}
|
||||
)
|
||||
|
||||
seen_title = False
|
||||
seen_topic = False
|
||||
current_topic_has_section = False
|
||||
topic_count = 0
|
||||
section_count = 0
|
||||
|
||||
for line_number, line in enumerate(markdown.splitlines(), start=1):
|
||||
match = HEADING_RE.match(line)
|
||||
if not match:
|
||||
continue
|
||||
level = len(match.group(1))
|
||||
title = match.group(2).strip().strip("*")
|
||||
|
||||
if level == 1:
|
||||
if line_number != 1 or title != "Working Protocol":
|
||||
violations.append(
|
||||
{
|
||||
"type": "invalid_heading",
|
||||
"line": line_number,
|
||||
"heading": line,
|
||||
"reason": "only the first line may be '# Working Protocol'",
|
||||
}
|
||||
)
|
||||
seen_title = True
|
||||
continue
|
||||
|
||||
if not seen_title:
|
||||
violations.append(
|
||||
{
|
||||
"type": "invalid_heading_order",
|
||||
"line": line_number,
|
||||
"heading": line,
|
||||
"reason": "heading appears before required title",
|
||||
}
|
||||
)
|
||||
continue
|
||||
|
||||
if level == 2:
|
||||
topic_count += 1
|
||||
seen_topic = True
|
||||
current_topic_has_section = False
|
||||
if title in GENERIC_SUMMARY_HEADINGS:
|
||||
violations.append(
|
||||
{
|
||||
"type": "generic_summary_framing",
|
||||
"line": line_number,
|
||||
"heading": line,
|
||||
"reason": "top-level category heading is not a topic",
|
||||
}
|
||||
)
|
||||
continue
|
||||
|
||||
if level == 3:
|
||||
if not seen_topic:
|
||||
violations.append(
|
||||
{
|
||||
"type": "missing_topic",
|
||||
"line": line_number,
|
||||
"heading": line,
|
||||
"reason": "section appears before any topic heading",
|
||||
}
|
||||
)
|
||||
if title not in ALLOWED_TOPIC_SECTIONS:
|
||||
violations.append(
|
||||
{
|
||||
"type": "unknown_topic_section",
|
||||
"line": line_number,
|
||||
"heading": line,
|
||||
"allowed": sorted(ALLOWED_TOPIC_SECTIONS),
|
||||
}
|
||||
)
|
||||
else:
|
||||
section_count += 1
|
||||
current_topic_has_section = True
|
||||
continue
|
||||
|
||||
violations.append(
|
||||
{
|
||||
"type": "unsupported_heading_level",
|
||||
"line": line_number,
|
||||
"heading": line,
|
||||
"reason": "Working Protocol V2 uses h1 title, h2 topics and h3 sections only",
|
||||
}
|
||||
)
|
||||
|
||||
if seen_topic and not current_topic_has_section:
|
||||
# This catches the last topic; earlier empty topics are caught below by
|
||||
# counting consecutive h2 headings.
|
||||
pass
|
||||
|
||||
h2_without_section = _topic_headings_without_sections(markdown)
|
||||
violations.extend(h2_without_section)
|
||||
|
||||
if topic_count == 0:
|
||||
violations.append(
|
||||
{
|
||||
"type": "missing_topic",
|
||||
"reason": "document must contain at least one '## Topic title' section",
|
||||
}
|
||||
)
|
||||
if section_count == 0:
|
||||
violations.append(
|
||||
{
|
||||
"type": "missing_topic_sections",
|
||||
"reason": "document must contain at least one allowed h3 topic section",
|
||||
}
|
||||
)
|
||||
return violations
|
||||
|
||||
|
||||
def _topic_headings_without_sections(markdown: str) -> list[dict[str, Any]]:
|
||||
violations: list[dict[str, Any]] = []
|
||||
current_topic: tuple[int, str] | None = None
|
||||
current_has_section = False
|
||||
|
||||
for line_number, line in enumerate(markdown.splitlines(), start=1):
|
||||
match = HEADING_RE.match(line)
|
||||
if not match:
|
||||
continue
|
||||
level = len(match.group(1))
|
||||
if level == 2:
|
||||
if current_topic is not None and not current_has_section:
|
||||
violations.append(
|
||||
{
|
||||
"type": "empty_topic",
|
||||
"line": current_topic[0],
|
||||
"heading": current_topic[1],
|
||||
"reason": "topic has no allowed h3 section",
|
||||
}
|
||||
)
|
||||
current_topic = (line_number, line)
|
||||
current_has_section = False
|
||||
elif level == 3 and current_topic is not None:
|
||||
title = match.group(2).strip().strip("*")
|
||||
if title in ALLOWED_TOPIC_SECTIONS:
|
||||
current_has_section = True
|
||||
|
||||
if current_topic is not None and not current_has_section:
|
||||
violations.append(
|
||||
{
|
||||
"type": "empty_topic",
|
||||
"line": current_topic[0],
|
||||
"heading": current_topic[1],
|
||||
"reason": "topic has no allowed h3 section",
|
||||
}
|
||||
)
|
||||
return violations
|
||||
|
||||
|
||||
def validate_working_protocol_markdown(markdown: str) -> WorkingProtocolValidationReport:
|
||||
violations: list[dict[str, Any]] = []
|
||||
if not markdown.startswith(REQUIRED_TITLE):
|
||||
violations.append(
|
||||
{
|
||||
"type": "missing_required_heading",
|
||||
"expected": REQUIRED_TITLE,
|
||||
"reason": "document must begin exactly with '# Working Protocol'",
|
||||
}
|
||||
)
|
||||
elif not markdown.startswith(REQUIRED_TITLE + "\n"):
|
||||
violations.append(
|
||||
{
|
||||
"type": "invalid_required_heading",
|
||||
"expected": REQUIRED_TITLE,
|
||||
"reason": "required heading must occupy the complete first line",
|
||||
}
|
||||
)
|
||||
|
||||
if markdown.startswith(REQUIRED_TITLE):
|
||||
violations.extend(validate_markdown_shape(markdown))
|
||||
|
||||
return WorkingProtocolValidationReport(
|
||||
valid=len(violations) == 0,
|
||||
violations=violations,
|
||||
)
|
||||
|
||||
|
||||
def write_json(path: Path, data: Any) -> None:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
path.write_text(json.dumps(data, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
|
||||
|
||||
|
||||
def render_working_protocol(
|
||||
input_path: Path,
|
||||
output_dir: Path,
|
||||
model: str = DEFAULT_MODEL,
|
||||
endpoint: str = DEFAULT_ENDPOINT,
|
||||
timeout: int = DEFAULT_TIMEOUT,
|
||||
num_ctx: int = DEFAULT_NUM_CTX,
|
||||
num_predict: int = DEFAULT_NUM_PREDICT,
|
||||
think: bool = False,
|
||||
) -> dict[str, Any]:
|
||||
output_dir.mkdir(parents=True, exist_ok=True)
|
||||
prompt_path = Path("prompts") / PROMPT_NAME
|
||||
raw_response_path = output_dir / "raw_model_response.json"
|
||||
raw_text_path = output_dir / "raw_model_response.txt"
|
||||
candidate_path = output_dir / "cleaned_candidate.md"
|
||||
validation_path = output_dir / "validation_report.json"
|
||||
protocol_path = output_dir / "working_protocol.md"
|
||||
metadata_path = output_dir / "metadata.json"
|
||||
report_path = output_dir / "report.md"
|
||||
|
||||
input_text = input_path.read_text(encoding="utf-8-sig")
|
||||
prompt = build_renderer_prompt(input_text)
|
||||
payload = build_ollama_payload(model, prompt, num_ctx, num_predict, think)
|
||||
|
||||
data, runtime = call_ollama(endpoint, payload, timeout)
|
||||
response_text = response_text_from_ollama_data(data)
|
||||
raw_response_path.write_text(
|
||||
json.dumps(data, ensure_ascii=False, indent=2) + "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
raw_text_path.write_text(response_text, encoding="utf-8")
|
||||
|
||||
cleaned = clean_working_protocol_markdown(response_text)
|
||||
candidate_path.write_text(cleaned, encoding="utf-8")
|
||||
validation = validate_working_protocol_markdown(cleaned)
|
||||
write_json(validation_path, validation.to_dict())
|
||||
|
||||
if validation.valid:
|
||||
protocol_path.write_text(cleaned, encoding="utf-8")
|
||||
elif protocol_path.exists():
|
||||
protocol_path.unlink()
|
||||
|
||||
metadata = {
|
||||
"model": model,
|
||||
"endpoint": endpoint,
|
||||
"think": think,
|
||||
"stream": False,
|
||||
"temperature": 0.0,
|
||||
"num_ctx": num_ctx,
|
||||
"num_predict": num_predict,
|
||||
"timeout_seconds": timeout,
|
||||
"request_count": 1,
|
||||
"runtime_seconds": runtime,
|
||||
"prompt_path": str(prompt_path.resolve()),
|
||||
"input_path": str(input_path.resolve()),
|
||||
"output_path": str(protocol_path.resolve()) if validation.valid else None,
|
||||
"raw_response_path": str(raw_response_path.resolve()),
|
||||
"raw_text_path": str(raw_text_path.resolve()),
|
||||
"cleaned_candidate_path": str(candidate_path.resolve()),
|
||||
"validation_report_path": str(validation_path.resolve()),
|
||||
"http_status_code": 200,
|
||||
"response_text_length": len(response_text),
|
||||
"candidate_text_length": len(cleaned),
|
||||
"valid": validation.valid,
|
||||
"readable_markdown": validation.valid,
|
||||
"top_level_json_keys": sorted(data.keys()),
|
||||
"done": data.get("done"),
|
||||
"done_reason": data.get("done_reason"),
|
||||
"total_duration": data.get("total_duration"),
|
||||
"load_duration": data.get("load_duration"),
|
||||
"prompt_eval_count": data.get("prompt_eval_count"),
|
||||
"prompt_eval_duration": data.get("prompt_eval_duration"),
|
||||
"eval_count": data.get("eval_count"),
|
||||
"eval_duration": data.get("eval_duration"),
|
||||
"created_at": datetime.now().isoformat(timespec="seconds"),
|
||||
}
|
||||
write_json(metadata_path, metadata)
|
||||
|
||||
lines = [
|
||||
"# Working Protocol Renderer V2 Report",
|
||||
"",
|
||||
f"- Result: {'valid renderer run' if validation.valid else 'invalid renderer run'}",
|
||||
f"- Model: `{model}`",
|
||||
f"- Runtime: {runtime:.3f} seconds",
|
||||
f"- Request count: 1",
|
||||
f"- Valid: {validation.valid}",
|
||||
f"- Violations: {len(validation.violations)}",
|
||||
f"- Raw response path: `{raw_response_path}`",
|
||||
f"- Cleaned candidate path: `{candidate_path}`",
|
||||
f"- Validation report path: `{validation_path}`",
|
||||
f"- Output path: `{protocol_path if validation.valid else 'not written'}`",
|
||||
]
|
||||
report_path.write_text("\n".join(lines) + "\n", encoding="utf-8")
|
||||
return metadata
|
||||
|
||||
|
||||
def main() -> int:
|
||||
args = parse_args()
|
||||
try:
|
||||
metadata = render_working_protocol(
|
||||
input_path=args.input,
|
||||
output_dir=args.output_dir,
|
||||
model=args.model,
|
||||
endpoint=args.endpoint,
|
||||
timeout=args.timeout,
|
||||
num_ctx=args.num_ctx,
|
||||
num_predict=args.num_predict,
|
||||
think=args.think,
|
||||
)
|
||||
except requests.ConnectionError as exc:
|
||||
print(f"Error: Ollama is not reachable at {args.endpoint}: {exc}", file=sys.stderr)
|
||||
return 1
|
||||
except requests.Timeout as exc:
|
||||
print(f"Error: Ollama request timed out after {args.timeout} seconds: {exc}", file=sys.stderr)
|
||||
return 1
|
||||
except requests.HTTPError as exc:
|
||||
print(f"Error: Ollama returned an HTTP error: {exc}", file=sys.stderr)
|
||||
return 1
|
||||
except (OSError, UnicodeError, ValueError, json.JSONDecodeError) as exc:
|
||||
print(f"Error: {exc}", file=sys.stderr)
|
||||
return 1
|
||||
|
||||
print(f"Runtime seconds: {metadata['runtime_seconds']:.3f}")
|
||||
print(f"Validation result: {'passed' if metadata['valid'] else 'failed'}")
|
||||
print(f"Output: {metadata['output_path'] or 'not written'}")
|
||||
print(f"Raw model response: {metadata['raw_response_path']}")
|
||||
print(f"Validation report: {metadata['validation_report_path']}")
|
||||
return 0 if metadata["valid"] else 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
Reference in New Issue
Block a user