Add windowed segmentation pipeline and review tooling
This commit is contained in:
@@ -0,0 +1,653 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Create a human-readable review report for transcript segmentation results.
|
||||
|
||||
The tool reads:
|
||||
|
||||
1. a normalized transcript text file
|
||||
2. a segmentation JSON file
|
||||
|
||||
It recreates the analysis blocks using the same block-size setting stored in
|
||||
the segmentation result and writes a Markdown report containing:
|
||||
|
||||
- all generated segments
|
||||
- the complete text of every segment
|
||||
- a compact review section around every detected boundary
|
||||
|
||||
Example:
|
||||
|
||||
python src/meeting_lab/segmentation/review_segmentation.py \
|
||||
samples/chunks/chunk_01_normalized.txt \
|
||||
samples/chunks/chunk_01_normalized_windowed_segments.json
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
|
||||
DEFAULT_CONTEXT_BLOCKS = 2
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(
|
||||
description=(
|
||||
"Create a readable review report for transcript segmentation."
|
||||
)
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"transcript_file",
|
||||
type=Path,
|
||||
help="Normalized transcript text file",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"segmentation_file",
|
||||
type=Path,
|
||||
help="Segmentation result JSON file",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"-o",
|
||||
"--output",
|
||||
type=Path,
|
||||
help=(
|
||||
"Output Markdown file; default: "
|
||||
"<segmentation_file>_review.md"
|
||||
),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--context-blocks",
|
||||
type=int,
|
||||
default=DEFAULT_CONTEXT_BLOCKS,
|
||||
help=(
|
||||
"Number of blocks shown before and after each boundary "
|
||||
f"(default: {DEFAULT_CONTEXT_BLOCKS})"
|
||||
),
|
||||
)
|
||||
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def split_into_blocks(
|
||||
text: str,
|
||||
target_chars: int,
|
||||
) -> list[str]:
|
||||
"""
|
||||
Recreate the same analysis blocks used by the segmentation tool.
|
||||
"""
|
||||
if target_chars < 1:
|
||||
raise ValueError(
|
||||
"target_chars must be greater than zero."
|
||||
)
|
||||
|
||||
normalized = (
|
||||
text.replace("\r\n", "\n")
|
||||
.replace("\r", "\n")
|
||||
.strip()
|
||||
)
|
||||
|
||||
units = [
|
||||
part.strip()
|
||||
for part in re.split(r"\n\s*\n+", normalized)
|
||||
if part.strip()
|
||||
]
|
||||
|
||||
if not units:
|
||||
units = [
|
||||
line.strip()
|
||||
for line in normalized.splitlines()
|
||||
if line.strip()
|
||||
]
|
||||
|
||||
blocks: list[str] = []
|
||||
current_units: list[str] = []
|
||||
current_length = 0
|
||||
|
||||
for unit in units:
|
||||
separator_length = 1 if current_units else 0
|
||||
|
||||
projected_length = (
|
||||
current_length
|
||||
+ separator_length
|
||||
+ len(unit)
|
||||
)
|
||||
|
||||
if current_units and projected_length > target_chars:
|
||||
blocks.append(
|
||||
"\n".join(current_units)
|
||||
)
|
||||
|
||||
current_units = [unit]
|
||||
current_length = len(unit)
|
||||
|
||||
else:
|
||||
current_units.append(unit)
|
||||
current_length = projected_length
|
||||
|
||||
if current_units:
|
||||
blocks.append(
|
||||
"\n".join(current_units)
|
||||
)
|
||||
|
||||
return blocks
|
||||
|
||||
|
||||
def load_json_file(
|
||||
path: Path,
|
||||
) -> dict[str, Any]:
|
||||
try:
|
||||
content = path.read_text(
|
||||
encoding="utf-8-sig"
|
||||
)
|
||||
data = json.loads(content)
|
||||
|
||||
except json.JSONDecodeError as exc:
|
||||
raise ValueError(
|
||||
f"Invalid JSON in {path}: {exc}"
|
||||
) from exc
|
||||
|
||||
if not isinstance(data, dict):
|
||||
raise ValueError(
|
||||
"The segmentation file must contain a JSON object."
|
||||
)
|
||||
|
||||
return data
|
||||
|
||||
|
||||
def get_target_block_chars(
|
||||
segmentation: dict[str, Any],
|
||||
) -> int:
|
||||
source = segmentation.get("source")
|
||||
|
||||
if not isinstance(source, dict):
|
||||
raise ValueError(
|
||||
"The segmentation JSON contains no valid source object."
|
||||
)
|
||||
|
||||
target_chars = source.get("target_block_chars")
|
||||
|
||||
if not isinstance(target_chars, int):
|
||||
raise ValueError(
|
||||
"The segmentation JSON contains no valid "
|
||||
"source.target_block_chars value."
|
||||
)
|
||||
|
||||
if target_chars < 1:
|
||||
raise ValueError(
|
||||
"source.target_block_chars must be greater than zero."
|
||||
)
|
||||
|
||||
return target_chars
|
||||
|
||||
|
||||
def get_expected_block_count(
|
||||
segmentation: dict[str, Any],
|
||||
) -> int:
|
||||
source = segmentation.get("source")
|
||||
|
||||
if not isinstance(source, dict):
|
||||
raise ValueError(
|
||||
"The segmentation JSON contains no valid source object."
|
||||
)
|
||||
|
||||
block_count = source.get("block_count")
|
||||
|
||||
if not isinstance(block_count, int):
|
||||
raise ValueError(
|
||||
"The segmentation JSON contains no valid "
|
||||
"source.block_count value."
|
||||
)
|
||||
|
||||
if block_count < 1:
|
||||
raise ValueError(
|
||||
"source.block_count must be greater than zero."
|
||||
)
|
||||
|
||||
return block_count
|
||||
|
||||
|
||||
def get_segments(
|
||||
segmentation: dict[str, Any],
|
||||
) -> list[dict[str, Any]]:
|
||||
raw_segments = segmentation.get("segments")
|
||||
|
||||
if not isinstance(raw_segments, list):
|
||||
raise ValueError(
|
||||
"The segmentation JSON contains no segments list."
|
||||
)
|
||||
|
||||
segments: list[dict[str, Any]] = []
|
||||
|
||||
for index, segment in enumerate(
|
||||
raw_segments,
|
||||
start=1,
|
||||
):
|
||||
if not isinstance(segment, dict):
|
||||
raise ValueError(
|
||||
f"Segment {index} is not a JSON object."
|
||||
)
|
||||
|
||||
segment_id = segment.get("segment_id")
|
||||
start_block = segment.get("start_block")
|
||||
end_block = segment.get("end_block")
|
||||
|
||||
if not isinstance(segment_id, str):
|
||||
raise ValueError(
|
||||
f"Segment {index} has no valid segment_id."
|
||||
)
|
||||
|
||||
if not isinstance(start_block, int):
|
||||
raise ValueError(
|
||||
f"Segment {segment_id} has no valid start_block."
|
||||
)
|
||||
|
||||
if not isinstance(end_block, int):
|
||||
raise ValueError(
|
||||
f"Segment {segment_id} has no valid end_block."
|
||||
)
|
||||
|
||||
if start_block < 1:
|
||||
raise ValueError(
|
||||
f"Segment {segment_id} starts before block 1."
|
||||
)
|
||||
|
||||
if end_block < start_block:
|
||||
raise ValueError(
|
||||
f"Segment {segment_id} ends before it starts."
|
||||
)
|
||||
|
||||
segments.append(
|
||||
{
|
||||
"segment_id": segment_id,
|
||||
"start_block": start_block,
|
||||
"end_block": end_block,
|
||||
}
|
||||
)
|
||||
|
||||
return segments
|
||||
|
||||
|
||||
def validate_segments(
|
||||
segments: list[dict[str, Any]],
|
||||
block_count: int,
|
||||
) -> None:
|
||||
if not segments:
|
||||
raise ValueError(
|
||||
"The segmentation result contains no segments."
|
||||
)
|
||||
|
||||
expected_start = 1
|
||||
|
||||
for segment in segments:
|
||||
start_block = segment["start_block"]
|
||||
end_block = segment["end_block"]
|
||||
segment_id = segment["segment_id"]
|
||||
|
||||
if start_block != expected_start:
|
||||
raise ValueError(
|
||||
f"Segment {segment_id} starts at block "
|
||||
f"{start_block}; expected block {expected_start}."
|
||||
)
|
||||
|
||||
if end_block > block_count:
|
||||
raise ValueError(
|
||||
f"Segment {segment_id} ends after the final block."
|
||||
)
|
||||
|
||||
expected_start = end_block + 1
|
||||
|
||||
if expected_start != block_count + 1:
|
||||
raise ValueError(
|
||||
"The segments do not cover all transcript blocks."
|
||||
)
|
||||
|
||||
|
||||
def render_block(
|
||||
block_number: int,
|
||||
block_text: str,
|
||||
) -> str:
|
||||
return (
|
||||
f"### Block {block_number}\n\n"
|
||||
f"{block_text}\n"
|
||||
)
|
||||
|
||||
|
||||
def render_segment(
|
||||
segment: dict[str, Any],
|
||||
blocks: list[str],
|
||||
) -> str:
|
||||
segment_id = segment["segment_id"]
|
||||
start_block = segment["start_block"]
|
||||
end_block = segment["end_block"]
|
||||
block_span = end_block - start_block + 1
|
||||
|
||||
lines = [
|
||||
f"## {segment_id}",
|
||||
"",
|
||||
f"**Blöcke:** {start_block}–{end_block}",
|
||||
"",
|
||||
f"**Umfang:** {block_span} Blöcke",
|
||||
"",
|
||||
]
|
||||
|
||||
for block_number in range(
|
||||
start_block,
|
||||
end_block + 1,
|
||||
):
|
||||
lines.append(
|
||||
render_block(
|
||||
block_number=block_number,
|
||||
block_text=blocks[block_number - 1],
|
||||
)
|
||||
)
|
||||
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def render_boundary_review(
|
||||
boundary: int,
|
||||
blocks: list[str],
|
||||
context_blocks: int,
|
||||
) -> str:
|
||||
block_count = len(blocks)
|
||||
|
||||
before_start = max(
|
||||
1,
|
||||
boundary - context_blocks + 1,
|
||||
)
|
||||
after_end = min(
|
||||
block_count,
|
||||
boundary + context_blocks,
|
||||
)
|
||||
|
||||
lines = [
|
||||
f"## Themenwechsel nach Block {boundary}",
|
||||
"",
|
||||
f"Der nächste Abschnitt beginnt mit Block {boundary + 1}.",
|
||||
"",
|
||||
"### Vor der Grenze",
|
||||
"",
|
||||
]
|
||||
|
||||
for block_number in range(
|
||||
before_start,
|
||||
boundary + 1,
|
||||
):
|
||||
lines.append(
|
||||
render_block(
|
||||
block_number=block_number,
|
||||
block_text=blocks[block_number - 1],
|
||||
)
|
||||
)
|
||||
|
||||
lines.extend(
|
||||
[
|
||||
"",
|
||||
"---",
|
||||
"",
|
||||
"## ⟶ Erkannter Themenwechsel",
|
||||
"",
|
||||
"---",
|
||||
"",
|
||||
"### Nach der Grenze",
|
||||
"",
|
||||
]
|
||||
)
|
||||
|
||||
for block_number in range(
|
||||
boundary + 1,
|
||||
after_end + 1,
|
||||
):
|
||||
lines.append(
|
||||
render_block(
|
||||
block_number=block_number,
|
||||
block_text=blocks[block_number - 1],
|
||||
)
|
||||
)
|
||||
|
||||
lines.extend(
|
||||
[
|
||||
"",
|
||||
"### Manuelle Bewertung",
|
||||
"",
|
||||
"- [ ] echter Themenwechsel",
|
||||
"- [ ] nur Unterthema",
|
||||
"- [ ] kurze Abschweifung",
|
||||
"- [ ] kein Themenwechsel",
|
||||
"",
|
||||
"**Notiz:**",
|
||||
"",
|
||||
"",
|
||||
]
|
||||
)
|
||||
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
def build_report(
|
||||
transcript_file: Path,
|
||||
segmentation_file: Path,
|
||||
segmentation: dict[str, Any],
|
||||
blocks: list[str],
|
||||
segments: list[dict[str, Any]],
|
||||
context_blocks: int,
|
||||
) -> str:
|
||||
configuration = segmentation.get(
|
||||
"configuration",
|
||||
{},
|
||||
)
|
||||
|
||||
if not isinstance(configuration, dict):
|
||||
configuration = {}
|
||||
|
||||
model = configuration.get(
|
||||
"model",
|
||||
"unbekannt",
|
||||
)
|
||||
|
||||
window_size = configuration.get(
|
||||
"window_size",
|
||||
"unbekannt",
|
||||
)
|
||||
|
||||
window_overlap = configuration.get(
|
||||
"window_overlap",
|
||||
"unbekannt",
|
||||
)
|
||||
|
||||
boundaries = [
|
||||
segment["end_block"]
|
||||
for segment in segments[:-1]
|
||||
]
|
||||
|
||||
lines = [
|
||||
"# Review der Meeting-Segmentierung",
|
||||
"",
|
||||
"## Metadaten",
|
||||
"",
|
||||
f"- **Transkript:** `{transcript_file}`",
|
||||
f"- **Segmentierung:** `{segmentation_file}`",
|
||||
f"- **Modell:** `{model}`",
|
||||
f"- **Blöcke:** {len(blocks)}",
|
||||
f"- **Segmente:** {len(segments)}",
|
||||
f"- **Themengrenzen:** {len(boundaries)}",
|
||||
f"- **Fenstergröße:** {window_size}",
|
||||
f"- **Überlappung:** {window_overlap}",
|
||||
"",
|
||||
"## Erkannte Grenzen",
|
||||
"",
|
||||
(
|
||||
", ".join(str(value) for value in boundaries)
|
||||
if boundaries
|
||||
else "Keine"
|
||||
),
|
||||
"",
|
||||
"---",
|
||||
"",
|
||||
"# Grenzprüfung",
|
||||
"",
|
||||
]
|
||||
|
||||
for boundary in boundaries:
|
||||
lines.append(
|
||||
render_boundary_review(
|
||||
boundary=boundary,
|
||||
blocks=blocks,
|
||||
context_blocks=context_blocks,
|
||||
)
|
||||
)
|
||||
|
||||
lines.extend(
|
||||
[
|
||||
"",
|
||||
"---",
|
||||
"",
|
||||
]
|
||||
)
|
||||
|
||||
lines.extend(
|
||||
[
|
||||
"# Vollständige Segmente",
|
||||
"",
|
||||
]
|
||||
)
|
||||
|
||||
for segment in segments:
|
||||
lines.append(
|
||||
render_segment(
|
||||
segment=segment,
|
||||
blocks=blocks,
|
||||
)
|
||||
)
|
||||
|
||||
lines.extend(
|
||||
[
|
||||
"",
|
||||
"---",
|
||||
"",
|
||||
]
|
||||
)
|
||||
|
||||
return "\n".join(lines).rstrip() + "\n"
|
||||
|
||||
|
||||
def main() -> int:
|
||||
args = parse_args()
|
||||
|
||||
try:
|
||||
if not args.transcript_file.is_file():
|
||||
raise FileNotFoundError(
|
||||
f"Transcript file not found: "
|
||||
f"{args.transcript_file}"
|
||||
)
|
||||
|
||||
if not args.segmentation_file.is_file():
|
||||
raise FileNotFoundError(
|
||||
f"Segmentation file not found: "
|
||||
f"{args.segmentation_file}"
|
||||
)
|
||||
|
||||
if args.context_blocks < 1:
|
||||
raise ValueError(
|
||||
"context_blocks must be at least 1."
|
||||
)
|
||||
|
||||
transcript = args.transcript_file.read_text(
|
||||
encoding="utf-8-sig"
|
||||
).strip()
|
||||
|
||||
if not transcript:
|
||||
raise ValueError(
|
||||
"The transcript file is empty."
|
||||
)
|
||||
|
||||
segmentation = load_json_file(
|
||||
args.segmentation_file
|
||||
)
|
||||
|
||||
target_block_chars = get_target_block_chars(
|
||||
segmentation
|
||||
)
|
||||
|
||||
expected_block_count = get_expected_block_count(
|
||||
segmentation
|
||||
)
|
||||
|
||||
blocks = split_into_blocks(
|
||||
text=transcript,
|
||||
target_chars=target_block_chars,
|
||||
)
|
||||
|
||||
if len(blocks) != expected_block_count:
|
||||
raise ValueError(
|
||||
"The recreated block count does not match the "
|
||||
"segmentation result: "
|
||||
f"{len(blocks)} instead of "
|
||||
f"{expected_block_count}."
|
||||
)
|
||||
|
||||
segments = get_segments(
|
||||
segmentation
|
||||
)
|
||||
|
||||
validate_segments(
|
||||
segments=segments,
|
||||
block_count=len(blocks),
|
||||
)
|
||||
|
||||
output_path = (
|
||||
args.output
|
||||
or args.segmentation_file.with_name(
|
||||
f"{args.segmentation_file.stem}_review.md"
|
||||
)
|
||||
)
|
||||
|
||||
output_path.parent.mkdir(
|
||||
parents=True,
|
||||
exist_ok=True,
|
||||
)
|
||||
|
||||
report = build_report(
|
||||
transcript_file=args.transcript_file,
|
||||
segmentation_file=args.segmentation_file,
|
||||
segmentation=segmentation,
|
||||
blocks=blocks,
|
||||
segments=segments,
|
||||
context_blocks=args.context_blocks,
|
||||
)
|
||||
|
||||
output_path.write_text(
|
||||
report,
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
print(f"Transcript: {args.transcript_file}")
|
||||
print(f"Segments: {args.segmentation_file}")
|
||||
print(f"Blocks: {len(blocks)}")
|
||||
print(f"Boundaries: {len(segments) - 1}")
|
||||
print(f"Output: {output_path}")
|
||||
|
||||
return 0
|
||||
|
||||
except (
|
||||
OSError,
|
||||
UnicodeError,
|
||||
ValueError,
|
||||
) as exc:
|
||||
print(
|
||||
f"Error: {exc}",
|
||||
file=sys.stderr,
|
||||
)
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -722,4 +722,4 @@ def main() -> int:
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
raise SystemExit(main())
|
||||
|
||||
@@ -0,0 +1,994 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Detect topic boundaries in a meeting transcript using overlapping windows.
|
||||
|
||||
Each model request receives a limited range of globally numbered transcript
|
||||
blocks. A trailing overlap provides context for decisions near the end of a
|
||||
window.
|
||||
|
||||
The model reports only topic boundaries from the window's decision range.
|
||||
Python merges the results and creates continuous, non-overlapping segments.
|
||||
|
||||
Example:
|
||||
python src/meeting_lab/segmentation/segment_topics_windowed.py \
|
||||
samples/chunks/chunk_01_normalized.txt \
|
||||
--model qwen3:8b
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
import requests
|
||||
|
||||
|
||||
DEFAULT_MODEL = "qwen3:8b"
|
||||
DEFAULT_ENDPOINT = "http://localhost:11434/api/generate"
|
||||
DEFAULT_TARGET_BLOCK_CHARS = 600
|
||||
DEFAULT_WINDOW_SIZE = 20
|
||||
DEFAULT_WINDOW_OVERLAP = 3
|
||||
|
||||
|
||||
SYSTEM_INSTRUCTION = """
|
||||
Du analysierst den chronologischen Verlauf eines Meeting-Transkripts.
|
||||
|
||||
Deine einzige Aufgabe ist es, Stellen zu erkennen, an denen das Gespräch
|
||||
inhaltlich zu einem anderen Thema wechselt.
|
||||
|
||||
Du erhältst einen Ausschnitt aus einem längeren Transkript. Die Blocknummern
|
||||
sind globale Blocknummern des vollständigen Meetings.
|
||||
|
||||
Gib ausschließlich die Nummern der Blöcke zurück, nach denen ein echter
|
||||
Themenwechsel stattfindet.
|
||||
|
||||
Regeln:
|
||||
|
||||
1. Gib nur echte inhaltliche Themenwechsel zurück.
|
||||
2. Kleine Rückfragen sind kein Themenwechsel.
|
||||
3. Kurze Ergänzungen sind kein Themenwechsel.
|
||||
4. Sprecherwechsel sind kein Themenwechsel.
|
||||
5. Zustimmung oder Ablehnung sind kein Themenwechsel.
|
||||
6. Kurze Abschweifungen sind nur dann eigene Segmente, wenn sie einen
|
||||
erkennbaren Gesprächsabschnitt bilden.
|
||||
7. Organisatorische Einleitungen und Verabschiedungen dürfen eigene Segmente
|
||||
bilden.
|
||||
8. Gib ausschließlich Grenzen innerhalb des ausdrücklich genannten
|
||||
Entscheidungsbereichs zurück.
|
||||
9. Gib keine Themenbezeichnungen zurück.
|
||||
10. Fasse keine Inhalte zusammen.
|
||||
11. Die Blocknummern müssen eindeutig und aufsteigend sortiert sein.
|
||||
12. Gib ausschließlich gültiges JSON aus.
|
||||
13. Gib kein Markdown und keine Erläuterungen aus.
|
||||
""".strip()
|
||||
|
||||
|
||||
def parse_args() -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(
|
||||
description=(
|
||||
"Detect topic boundaries using overlapping transcript windows."
|
||||
)
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"input_file",
|
||||
type=Path,
|
||||
help="Normalized transcript file",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"-o",
|
||||
"--output",
|
||||
type=Path,
|
||||
help=(
|
||||
"Output JSON file; default: "
|
||||
"<input>_windowed_segments.json"
|
||||
),
|
||||
)
|
||||
|
||||
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 endpoint (default: {DEFAULT_ENDPOINT})",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--timeout",
|
||||
type=int,
|
||||
default=1800,
|
||||
help="Timeout per Ollama request in seconds (default: 1800)",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--temperature",
|
||||
type=float,
|
||||
default=0.0,
|
||||
help="Sampling temperature (default: 0.0)",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--target-block-chars",
|
||||
type=int,
|
||||
default=DEFAULT_TARGET_BLOCK_CHARS,
|
||||
help=(
|
||||
"Approximate analysis block size in characters "
|
||||
f"(default: {DEFAULT_TARGET_BLOCK_CHARS})"
|
||||
),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--window-size",
|
||||
type=int,
|
||||
default=DEFAULT_WINDOW_SIZE,
|
||||
help=(
|
||||
"Number of blocks supplied per model request "
|
||||
f"(default: {DEFAULT_WINDOW_SIZE})"
|
||||
),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--window-overlap",
|
||||
type=int,
|
||||
default=DEFAULT_WINDOW_OVERLAP,
|
||||
help=(
|
||||
"Trailing context blocks per window "
|
||||
f"(default: {DEFAULT_WINDOW_OVERLAP})"
|
||||
),
|
||||
)
|
||||
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def split_into_blocks(
|
||||
text: str,
|
||||
target_chars: int,
|
||||
) -> list[str]:
|
||||
"""
|
||||
Combine short transcript utterances into larger analysis blocks.
|
||||
"""
|
||||
if target_chars < 1:
|
||||
raise ValueError(
|
||||
"target_chars must be greater than zero."
|
||||
)
|
||||
|
||||
normalized = (
|
||||
text.replace("\r\n", "\n")
|
||||
.replace("\r", "\n")
|
||||
.strip()
|
||||
)
|
||||
|
||||
units = [
|
||||
part.strip()
|
||||
for part in re.split(r"\n\s*\n+", normalized)
|
||||
if part.strip()
|
||||
]
|
||||
|
||||
if not units:
|
||||
units = [
|
||||
line.strip()
|
||||
for line in normalized.splitlines()
|
||||
if line.strip()
|
||||
]
|
||||
|
||||
blocks: list[str] = []
|
||||
current_units: list[str] = []
|
||||
current_length = 0
|
||||
|
||||
for unit in units:
|
||||
separator_length = 1 if current_units else 0
|
||||
|
||||
projected_length = (
|
||||
current_length
|
||||
+ separator_length
|
||||
+ len(unit)
|
||||
)
|
||||
|
||||
if (
|
||||
current_units
|
||||
and projected_length > target_chars
|
||||
):
|
||||
blocks.append(
|
||||
"\n".join(current_units)
|
||||
)
|
||||
|
||||
current_units = [unit]
|
||||
current_length = len(unit)
|
||||
|
||||
else:
|
||||
current_units.append(unit)
|
||||
current_length = projected_length
|
||||
|
||||
if current_units:
|
||||
blocks.append(
|
||||
"\n".join(current_units)
|
||||
)
|
||||
|
||||
return blocks
|
||||
|
||||
|
||||
def create_windows(
|
||||
block_count: int,
|
||||
window_size: int,
|
||||
overlap: int,
|
||||
) -> list[dict[str, int]]:
|
||||
"""
|
||||
Create windows with trailing context.
|
||||
|
||||
Example with window_size=20 and overlap=3:
|
||||
|
||||
Window 1:
|
||||
context: 1-20
|
||||
decisions after blocks 1-17
|
||||
|
||||
Window 2:
|
||||
context: 18-37
|
||||
decisions after blocks 18-34
|
||||
|
||||
Every possible boundary is assigned to exactly one decision range.
|
||||
"""
|
||||
if block_count < 1:
|
||||
raise ValueError(
|
||||
"block_count must be greater than zero."
|
||||
)
|
||||
|
||||
if window_size < 2:
|
||||
raise ValueError(
|
||||
"window_size must be at least 2."
|
||||
)
|
||||
|
||||
if overlap < 1:
|
||||
raise ValueError(
|
||||
"window_overlap must be at least 1."
|
||||
)
|
||||
|
||||
if overlap >= window_size:
|
||||
raise ValueError(
|
||||
"window_overlap must be smaller than window_size."
|
||||
)
|
||||
|
||||
stride = window_size - overlap
|
||||
windows: list[dict[str, int]] = []
|
||||
|
||||
start_block = 1
|
||||
window_number = 1
|
||||
|
||||
while start_block <= block_count:
|
||||
end_block = min(
|
||||
start_block + window_size - 1,
|
||||
block_count,
|
||||
)
|
||||
|
||||
if end_block == block_count:
|
||||
decision_end = block_count - 1
|
||||
else:
|
||||
decision_end = (
|
||||
start_block + stride - 1
|
||||
)
|
||||
|
||||
decision_end = min(
|
||||
decision_end,
|
||||
block_count - 1,
|
||||
)
|
||||
|
||||
if decision_end < start_block:
|
||||
break
|
||||
|
||||
windows.append(
|
||||
{
|
||||
"window_number": window_number,
|
||||
"start_block": start_block,
|
||||
"end_block": end_block,
|
||||
"decision_start": start_block,
|
||||
"decision_end": decision_end,
|
||||
}
|
||||
)
|
||||
|
||||
if end_block == block_count:
|
||||
break
|
||||
|
||||
start_block += stride
|
||||
window_number += 1
|
||||
|
||||
return windows
|
||||
|
||||
|
||||
def render_window_blocks(
|
||||
blocks: list[str],
|
||||
start_block: int,
|
||||
end_block: int,
|
||||
) -> str:
|
||||
rendered: list[str] = []
|
||||
|
||||
for block_number in range(
|
||||
start_block,
|
||||
end_block + 1,
|
||||
):
|
||||
block_text = blocks[
|
||||
block_number - 1
|
||||
]
|
||||
|
||||
rendered.append(
|
||||
f"[BLOCK {block_number}]\n"
|
||||
f"{block_text}"
|
||||
)
|
||||
|
||||
return "\n\n".join(rendered)
|
||||
|
||||
|
||||
def build_prompt(
|
||||
blocks: list[str],
|
||||
window: dict[str, int],
|
||||
) -> str:
|
||||
transcript = render_window_blocks(
|
||||
blocks=blocks,
|
||||
start_block=window["start_block"],
|
||||
end_block=window["end_block"],
|
||||
)
|
||||
|
||||
decision_start = window[
|
||||
"decision_start"
|
||||
]
|
||||
decision_end = window[
|
||||
"decision_end"
|
||||
]
|
||||
|
||||
return f"""
|
||||
Bestimme die echten inhaltlichen Themenwechsel in diesem Ausschnitt.
|
||||
|
||||
Der Ausschnitt enthält die globalen Blöcke:
|
||||
|
||||
{window["start_block"]} bis {window["end_block"]}
|
||||
|
||||
Du darfst ausschließlich Themenwechsel nach Blöcken aus diesem
|
||||
Entscheidungsbereich zurückgeben:
|
||||
|
||||
{decision_start} bis {decision_end}
|
||||
|
||||
Ein Wert N bedeutet:
|
||||
|
||||
Nach Block N endet ein Gesprächsabschnitt und mit Block N+1 beginnt ein
|
||||
inhaltlich anderes Thema.
|
||||
|
||||
Blöcke außerhalb des Entscheidungsbereichs dienen nur als Kontext.
|
||||
|
||||
Die Antwort muss genau diese Struktur besitzen:
|
||||
|
||||
{{
|
||||
"topic_changes": []
|
||||
}}
|
||||
|
||||
Trage in die Liste ausschließlich zutreffende Blocknummern zwischen
|
||||
{decision_start} und {decision_end} ein.
|
||||
|
||||
Die leere Liste bedeutet, dass innerhalb dieses Entscheidungsbereichs kein
|
||||
echter Themenwechsel vorliegt.
|
||||
|
||||
TRANSKRIPTAUSSCHNITT:
|
||||
--- BEGINN ---
|
||||
{transcript}
|
||||
--- ENDE ---
|
||||
""".strip()
|
||||
|
||||
|
||||
def build_topic_change_schema(
|
||||
decision_start: int,
|
||||
decision_end: int,
|
||||
) -> dict[str, Any]:
|
||||
"""
|
||||
Build a JSON schema restricted to the current decision range.
|
||||
"""
|
||||
return {
|
||||
"type": "object",
|
||||
"properties": {
|
||||
"topic_changes": {
|
||||
"type": "array",
|
||||
"items": {
|
||||
"type": "integer",
|
||||
"minimum": decision_start,
|
||||
"maximum": decision_end,
|
||||
},
|
||||
"uniqueItems": True,
|
||||
}
|
||||
},
|
||||
"required": [
|
||||
"topic_changes"
|
||||
],
|
||||
"additionalProperties": False,
|
||||
}
|
||||
|
||||
|
||||
def call_ollama(
|
||||
endpoint: str,
|
||||
model: str,
|
||||
prompt: str,
|
||||
timeout: int,
|
||||
temperature: float,
|
||||
decision_start: int,
|
||||
decision_end: int,
|
||||
) -> tuple[str, dict[str, Any]]:
|
||||
payload = {
|
||||
"model": model,
|
||||
"system": SYSTEM_INSTRUCTION,
|
||||
"prompt": prompt,
|
||||
"stream": False,
|
||||
"format": build_topic_change_schema(
|
||||
decision_start=decision_start,
|
||||
decision_end=decision_end,
|
||||
),
|
||||
"think": False,
|
||||
"options": {
|
||||
"temperature": temperature,
|
||||
"seed": 42,
|
||||
"num_ctx": 16384,
|
||||
},
|
||||
}
|
||||
|
||||
started = time.perf_counter()
|
||||
|
||||
response = requests.post(
|
||||
endpoint,
|
||||
json=payload,
|
||||
timeout=timeout,
|
||||
)
|
||||
|
||||
elapsed = (
|
||||
time.perf_counter() - started
|
||||
)
|
||||
|
||||
response.raise_for_status()
|
||||
|
||||
data = response.json()
|
||||
response_text = data.get("response")
|
||||
|
||||
if (
|
||||
not isinstance(response_text, str)
|
||||
or not response_text.strip()
|
||||
):
|
||||
raise ValueError(
|
||||
"Ollama returned no usable response text."
|
||||
)
|
||||
|
||||
metadata = {
|
||||
"model": data.get(
|
||||
"model",
|
||||
model,
|
||||
),
|
||||
"elapsed_seconds": round(
|
||||
elapsed,
|
||||
3,
|
||||
),
|
||||
"total_duration_ns": data.get(
|
||||
"total_duration"
|
||||
),
|
||||
"load_duration_ns": data.get(
|
||||
"load_duration"
|
||||
),
|
||||
"prompt_eval_count": data.get(
|
||||
"prompt_eval_count"
|
||||
),
|
||||
"prompt_eval_duration_ns": data.get(
|
||||
"prompt_eval_duration"
|
||||
),
|
||||
"eval_count": data.get(
|
||||
"eval_count"
|
||||
),
|
||||
"eval_duration_ns": data.get(
|
||||
"eval_duration"
|
||||
),
|
||||
}
|
||||
|
||||
return (
|
||||
response_text.strip(),
|
||||
metadata,
|
||||
)
|
||||
|
||||
|
||||
def parse_json_response(
|
||||
text: str,
|
||||
) -> dict[str, Any]:
|
||||
try:
|
||||
parsed = json.loads(text)
|
||||
|
||||
except json.JSONDecodeError:
|
||||
match = re.search(
|
||||
r"\{.*\}",
|
||||
text,
|
||||
flags=re.DOTALL,
|
||||
)
|
||||
|
||||
if not match:
|
||||
raise
|
||||
|
||||
parsed = json.loads(
|
||||
match.group(0)
|
||||
)
|
||||
|
||||
if not isinstance(parsed, dict):
|
||||
raise ValueError(
|
||||
"The model response is not a JSON object."
|
||||
)
|
||||
|
||||
return parsed
|
||||
|
||||
|
||||
def validate_window_changes(
|
||||
result: dict[str, Any],
|
||||
decision_start: int,
|
||||
decision_end: int,
|
||||
) -> list[int]:
|
||||
raw_changes = result.get(
|
||||
"topic_changes"
|
||||
)
|
||||
|
||||
if not isinstance(
|
||||
raw_changes,
|
||||
list,
|
||||
):
|
||||
raise ValueError(
|
||||
"Model output does not contain "
|
||||
"a topic_changes list."
|
||||
)
|
||||
|
||||
changes: list[int] = []
|
||||
|
||||
for value in raw_changes:
|
||||
if (
|
||||
not isinstance(value, int)
|
||||
or isinstance(value, bool)
|
||||
):
|
||||
raise ValueError(
|
||||
"Every topic change must "
|
||||
"be an integer."
|
||||
)
|
||||
|
||||
if (
|
||||
value < decision_start
|
||||
or value > decision_end
|
||||
):
|
||||
raise ValueError(
|
||||
f"Topic change {value} is outside "
|
||||
"the permitted decision range "
|
||||
f"{decision_start}-{decision_end}."
|
||||
)
|
||||
|
||||
changes.append(value)
|
||||
|
||||
if (
|
||||
len(changes)
|
||||
!= len(set(changes))
|
||||
):
|
||||
raise ValueError(
|
||||
"The model returned duplicate boundaries."
|
||||
)
|
||||
|
||||
if changes != sorted(changes):
|
||||
raise ValueError(
|
||||
"The topic changes are not sorted."
|
||||
)
|
||||
|
||||
return changes
|
||||
|
||||
|
||||
def build_segments(
|
||||
topic_changes: list[int],
|
||||
block_count: int,
|
||||
) -> list[dict[str, Any]]:
|
||||
"""
|
||||
Convert boundaries into continuous, non-overlapping segments.
|
||||
"""
|
||||
segments: list[dict[str, Any]] = []
|
||||
start_block = 1
|
||||
|
||||
for number, end_block in enumerate(
|
||||
topic_changes,
|
||||
start=1,
|
||||
):
|
||||
segments.append(
|
||||
{
|
||||
"segment_id": (
|
||||
f"segment_{number:03d}"
|
||||
),
|
||||
"start_block": start_block,
|
||||
"end_block": end_block,
|
||||
}
|
||||
)
|
||||
|
||||
start_block = end_block + 1
|
||||
|
||||
segments.append(
|
||||
{
|
||||
"segment_id": (
|
||||
f"segment_{len(segments) + 1:03d}"
|
||||
),
|
||||
"start_block": start_block,
|
||||
"end_block": block_count,
|
||||
}
|
||||
)
|
||||
|
||||
return segments
|
||||
|
||||
|
||||
def validate_segments(
|
||||
segments: list[dict[str, Any]],
|
||||
block_count: int,
|
||||
) -> None:
|
||||
if not segments:
|
||||
raise ValueError(
|
||||
"No segments were created."
|
||||
)
|
||||
|
||||
expected_start = 1
|
||||
|
||||
for segment in segments:
|
||||
start_block = segment.get(
|
||||
"start_block"
|
||||
)
|
||||
end_block = segment.get(
|
||||
"end_block"
|
||||
)
|
||||
|
||||
if not isinstance(
|
||||
start_block,
|
||||
int,
|
||||
):
|
||||
raise ValueError(
|
||||
"Invalid segment start block."
|
||||
)
|
||||
|
||||
if not isinstance(
|
||||
end_block,
|
||||
int,
|
||||
):
|
||||
raise ValueError(
|
||||
"Invalid segment end block."
|
||||
)
|
||||
|
||||
if start_block != expected_start:
|
||||
raise ValueError(
|
||||
"Segments are not continuous."
|
||||
)
|
||||
|
||||
if end_block < start_block:
|
||||
raise ValueError(
|
||||
"A segment ends before it starts."
|
||||
)
|
||||
|
||||
expected_start = end_block + 1
|
||||
|
||||
if expected_start != block_count + 1:
|
||||
raise ValueError(
|
||||
"Segments do not cover all blocks."
|
||||
)
|
||||
|
||||
|
||||
def build_raw_output_path(
|
||||
output_path: Path,
|
||||
window_number: int,
|
||||
) -> Path:
|
||||
return output_path.with_name(
|
||||
f"{output_path.stem}"
|
||||
f"_window_{window_number:03d}"
|
||||
".raw.txt"
|
||||
)
|
||||
|
||||
|
||||
def main() -> int:
|
||||
args = parse_args()
|
||||
|
||||
try:
|
||||
if not args.input_file.is_file():
|
||||
raise FileNotFoundError(
|
||||
f"Input file not found: "
|
||||
f"{args.input_file}"
|
||||
)
|
||||
|
||||
transcript = args.input_file.read_text(
|
||||
encoding="utf-8-sig"
|
||||
).strip()
|
||||
|
||||
if not transcript:
|
||||
raise ValueError(
|
||||
"The input file is empty."
|
||||
)
|
||||
|
||||
blocks = split_into_blocks(
|
||||
text=transcript,
|
||||
target_chars=(
|
||||
args.target_block_chars
|
||||
),
|
||||
)
|
||||
|
||||
if not blocks:
|
||||
raise ValueError(
|
||||
"No analysis blocks could "
|
||||
"be created."
|
||||
)
|
||||
|
||||
windows = create_windows(
|
||||
block_count=len(blocks),
|
||||
window_size=args.window_size,
|
||||
overlap=args.window_overlap,
|
||||
)
|
||||
|
||||
if not windows:
|
||||
raise ValueError(
|
||||
"No analysis windows could "
|
||||
"be created."
|
||||
)
|
||||
|
||||
output_path = (
|
||||
args.output
|
||||
or args.input_file.with_name(
|
||||
f"{args.input_file.stem}"
|
||||
"_windowed_segments.json"
|
||||
)
|
||||
)
|
||||
|
||||
output_path.parent.mkdir(
|
||||
parents=True,
|
||||
exist_ok=True,
|
||||
)
|
||||
|
||||
print(
|
||||
f"Model: {args.model}"
|
||||
)
|
||||
print(
|
||||
f"Input: {args.input_file}"
|
||||
)
|
||||
print(
|
||||
f"Blocks: {len(blocks)}"
|
||||
)
|
||||
print(
|
||||
f"Windows: {len(windows)}"
|
||||
)
|
||||
print(
|
||||
f"Window size: {args.window_size}"
|
||||
)
|
||||
print(
|
||||
"Window overlap: "
|
||||
f"{args.window_overlap}"
|
||||
)
|
||||
print(
|
||||
f"Output: {output_path}"
|
||||
)
|
||||
print()
|
||||
|
||||
all_changes: list[int] = []
|
||||
window_results: list[
|
||||
dict[str, Any]
|
||||
] = []
|
||||
|
||||
total_started = (
|
||||
time.perf_counter()
|
||||
)
|
||||
|
||||
for window in windows:
|
||||
number = window[
|
||||
"window_number"
|
||||
]
|
||||
|
||||
print(
|
||||
f"Window {number}/"
|
||||
f"{len(windows)}: "
|
||||
"context "
|
||||
f"{window['start_block']}-"
|
||||
f"{window['end_block']}, "
|
||||
"decisions "
|
||||
f"{window['decision_start']}-"
|
||||
f"{window['decision_end']} ..."
|
||||
)
|
||||
|
||||
prompt = build_prompt(
|
||||
blocks=blocks,
|
||||
window=window,
|
||||
)
|
||||
|
||||
raw_text, metadata = call_ollama(
|
||||
endpoint=args.endpoint,
|
||||
model=args.model,
|
||||
prompt=prompt,
|
||||
timeout=args.timeout,
|
||||
temperature=args.temperature,
|
||||
decision_start=window[
|
||||
"decision_start"
|
||||
],
|
||||
decision_end=window[
|
||||
"decision_end"
|
||||
],
|
||||
)
|
||||
|
||||
try:
|
||||
parsed = parse_json_response(
|
||||
raw_text
|
||||
)
|
||||
|
||||
changes = (
|
||||
validate_window_changes(
|
||||
result=parsed,
|
||||
decision_start=window[
|
||||
"decision_start"
|
||||
],
|
||||
decision_end=window[
|
||||
"decision_end"
|
||||
],
|
||||
)
|
||||
)
|
||||
|
||||
except (
|
||||
json.JSONDecodeError,
|
||||
ValueError,
|
||||
) as exc:
|
||||
raw_path = (
|
||||
build_raw_output_path(
|
||||
output_path=output_path,
|
||||
window_number=number,
|
||||
)
|
||||
)
|
||||
|
||||
raw_path.write_text(
|
||||
raw_text + "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
raise ValueError(
|
||||
f"Window {number} could not "
|
||||
"be validated. "
|
||||
f"Raw output saved to "
|
||||
f"{raw_path}. "
|
||||
f"Reason: {exc}"
|
||||
) from exc
|
||||
|
||||
all_changes.extend(changes)
|
||||
|
||||
window_results.append(
|
||||
{
|
||||
**window,
|
||||
"topic_changes": changes,
|
||||
"_run": metadata,
|
||||
}
|
||||
)
|
||||
|
||||
formatted_changes: (
|
||||
list[int] | str
|
||||
) = (
|
||||
changes
|
||||
if changes
|
||||
else "none"
|
||||
)
|
||||
|
||||
print(
|
||||
" boundaries: "
|
||||
f"{formatted_changes} "
|
||||
"("
|
||||
f"{metadata['elapsed_seconds']:.1f}"
|
||||
" s)"
|
||||
)
|
||||
|
||||
all_changes = sorted(
|
||||
set(all_changes)
|
||||
)
|
||||
|
||||
segments = build_segments(
|
||||
topic_changes=all_changes,
|
||||
block_count=len(blocks),
|
||||
)
|
||||
|
||||
validate_segments(
|
||||
segments=segments,
|
||||
block_count=len(blocks),
|
||||
)
|
||||
|
||||
total_elapsed = (
|
||||
time.perf_counter()
|
||||
- total_started
|
||||
)
|
||||
|
||||
result = {
|
||||
"source": {
|
||||
"file": args.input_file.name,
|
||||
"block_count": len(blocks),
|
||||
"target_block_chars": (
|
||||
args.target_block_chars
|
||||
),
|
||||
},
|
||||
"configuration": {
|
||||
"model": args.model,
|
||||
"window_size": (
|
||||
args.window_size
|
||||
),
|
||||
"window_overlap": (
|
||||
args.window_overlap
|
||||
),
|
||||
"temperature": (
|
||||
args.temperature
|
||||
),
|
||||
},
|
||||
"topic_changes": all_changes,
|
||||
"segments": segments,
|
||||
"windows": window_results,
|
||||
"_run": {
|
||||
"elapsed_seconds": round(
|
||||
total_elapsed,
|
||||
3,
|
||||
),
|
||||
},
|
||||
}
|
||||
|
||||
output_path.write_text(
|
||||
json.dumps(
|
||||
result,
|
||||
ensure_ascii=False,
|
||||
indent=2,
|
||||
)
|
||||
+ "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
print()
|
||||
print(
|
||||
f"Done in {total_elapsed:.1f} "
|
||||
"seconds."
|
||||
)
|
||||
print(
|
||||
f"Topic changes: "
|
||||
f"{len(all_changes)}"
|
||||
)
|
||||
print(
|
||||
f"Segments: "
|
||||
f"{len(segments)}"
|
||||
)
|
||||
print(
|
||||
f"Boundaries: "
|
||||
f"{all_changes}"
|
||||
)
|
||||
|
||||
return 0
|
||||
|
||||
except requests.ConnectionError:
|
||||
print(
|
||||
"Error: Ollama is not reachable. "
|
||||
"Is `ollama serve` running?",
|
||||
file=sys.stderr,
|
||||
)
|
||||
return 1
|
||||
|
||||
except requests.Timeout:
|
||||
print(
|
||||
"Error: An Ollama request timed out.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
return 1
|
||||
|
||||
except requests.HTTPError as exc:
|
||||
print(
|
||||
"Error: Ollama returned an HTTP "
|
||||
f"error: {exc}",
|
||||
file=sys.stderr,
|
||||
)
|
||||
return 1
|
||||
|
||||
except (
|
||||
OSError,
|
||||
UnicodeError,
|
||||
ValueError,
|
||||
) as exc:
|
||||
print(
|
||||
f"Error: {exc}",
|
||||
file=sys.stderr,
|
||||
)
|
||||
return 1
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
Reference in New Issue
Block a user