Kansho/backend/journal_generate.py
2026-08-29 20:28:59 +02:00

748 lines
27 KiB
Python

"""Explicit journal draft generation. Never overwrites the current user entry."""
from __future__ import annotations
import re
from context_builder import assemble_text, build_internal_context
from debug_store import persist_engine_error, persist_step
from dialogue_store import StoreError, list_conversations_for_day, list_messages
from engine import EngineError, execute_prompt, load_active_prompt, preview_prompt
from journal_policy import require_explicit_generate, source_conversation_ids
from identity_store import is_maskable_label, list_mappings, mapping_spellings
from journal_body import clean_title
from journal_editorial import (
compile_journal_style_application,
incomplete_syntax_markers,
lexical_similarity,
narration_sources_text,
style_application_trace,
style_example_diagnostics,
writing_profile_trace,
)
from journal_generation_policy import (
assert_journal_prompt_contract,
assert_template_resolved,
compile_selection,
draft_snapshot,
mark_guidelines_used,
policy_trace,
prompt_without_orphaned_style_example_rule,
resolve_run_selection,
)
from journal_reconstruct import (
assign_source_ids,
claim_texts,
local_verified_artifact,
reconstruction_text,
)
from journal_shape import shape_journal
from journal_store import current_draft, current_entries, get_day, insert_draft
from model_catalog import resolve_generate_metadata
from prompt_budget import (
JournalBudgetError,
assert_input_fits,
plan_journal_budget,
)
from privacy_gateway import (
GENERIC_PLACEHOLDER_INNER,
canonical_token,
compact_diagnostics,
identity_label_pattern,
identity_occurrence_count,
is_identity_mention,
)
from providers import generate_provider
from retrieval import retrieve
from writing_profile_store import remember_dialogue_style
JOURNAL_NOT_ACCEPTED = "journal_generation_not_accepted"
JOURNAL_NOT_ACCEPTED_MESSAGE = "Generierung nicht übernommen."
IDENTITY_PLACEHOLDER = re.compile(
r"\[\[\s*(?:SELF|PERSON:[^\]]+|PLACE:[^\]]+|ORG:[^\]]+|PROJECT:[^\]]+|…|\.{2,})\s*\]\]",
re.IGNORECASE,
)
def _message_ids(profile_id: str, conversation_ids: list[str]) -> list[str]:
ids: list[str] = []
for conversation_id in conversation_ids:
for message in list_messages(profile_id, conversation_id):
ids.append(message["id"])
return ids
def _split_title(content: str) -> tuple[str, str]:
text = (content or "").strip()
if not text:
return "", ""
lines = text.splitlines()
title = clean_title(lines[0])
body = "\n".join(lines[1:]).strip() if len(lines) > 1 else text
if not title or len(title) > 80:
return "", text
return title, body or text
def _local_narration(reconstruction: dict) -> tuple[str, str]:
"""Local source preview only. Never stored as a generated journal draft."""
parts = [str(part).strip() for part in claim_texts(reconstruction) if str(part).strip()]
body = "\n\n".join(parts)
return "Ein Tag", body
def unattested_journal_content(
text: str,
source_user: list[str],
mappings: list[dict],
active_tokens: list[str] | None = None,
) -> str | None:
"""Journal-only: historical names not in the sources are unattested facts, not a privacy leak.
Title and body are checked together. Placeholders must stay consistent across both.
Attestation is token-based: any attested spelling of a token (canonical or alias)
allows every known spelling of that same token after local demasking.
"""
body = text or ""
allowed = {canonical_token(token).upper().replace(" ", "") for token in (active_tokens or [])}
allowed.add("SELF")
for match in IDENTITY_PLACEHOLDER.finditer(body):
token = canonical_token(match.group(0)).upper().replace(" ", "")
if token == "SELF":
continue
if token in GENERIC_PLACEHOLDER_INNER or token not in allowed:
return "unattested_placeholder"
sources = "\n".join(source_user or [])
spellings_by_token: dict[str, list[str]] = {}
for item in mappings or []:
token = canonical_token(item.get("token") or "").upper().replace(" ", "")
if not token:
continue
bucket = spellings_by_token.setdefault(token, [])
seen = {label.casefold() for label in bucket}
for label in mapping_spellings(item):
if not label or label.casefold() in seen:
continue
seen.add(label.casefold())
bucket.append(label)
attested: set[str] = set()
for token, labels in spellings_by_token.items():
for label in labels:
if identity_occurrence_count(sources, label, token) > 0:
attested.add(token)
break
for token, labels in spellings_by_token.items():
if token in attested:
continue
for label in labels:
if not is_maskable_label(label):
continue
for match in identity_label_pattern(label).finditer(body):
if is_identity_mention(body, match.start(), match.end(), token):
return "unattested_identity"
return None
def _reject_generation(
*,
reason: str,
diagnostics: dict | None = None,
source_preview: str = "",
extra: dict | None = None,
) -> None:
diag = dict(diagnostics or {})
trace = dict(diag.get("trace") or {})
extra = dict(extra or {})
trace.update(
{
"purpose": "journal_generate",
"prompt_slug": extra.get("prompt_slug") or trace.get("prompt_slug") or "mvp.journal_generate",
"model_text_accepted": False,
"narration_source": "not_accepted",
"abort_reason": reason,
"source_preview": source_preview,
**extra,
}
)
diag["trace"] = trace
diag["log"] = extra.get("log") or diag.get("log") or trace.get("log") or []
diag["model_text_accepted"] = False
diag["abort_reason"] = reason
diag["provenance_decision"] = extra.get("provenance_decision") or diag.get("provenance_decision")
raise EngineError(
JOURNAL_NOT_ACCEPTED,
JOURNAL_NOT_ACCEPTED_MESSAGE,
409,
diagnostics=diag,
)
def _compose_journal_trace(
reconstruct_result: dict,
narrate_result: dict,
*,
run_log: list[dict],
narrate_prompt: dict,
profile_meta: dict,
style_meta: dict,
dropped: list[str],
extra: dict | None = None,
) -> dict:
reconstruct_trace = _stage_trace(reconstruct_result, "local_source_artifact")
reconstruct_trace.setdefault("purpose", "local_source_artifact")
reconstruct_trace.setdefault("status", "local_ok")
reconstruct_trace.setdefault("stage1", "local_ok")
narrate_trace = _stage_trace(narrate_result, "journal_generate")
narrate_trace.setdefault("prompt_slug", narrate_prompt.get("slug") or "mvp.journal_generate")
narrate_trace["prompt_revision"] = (
narrate_trace.get("prompt_revision")
or narrate_prompt.get("seed_revision")
or ""
)
narrate_trace["writing_profile"] = profile_meta
narrate_trace["style_examples"] = {**style_meta, "dropped": "style_examples" in dropped}
narrate_trace["dropped_optional_blocks"] = dropped
if extra and extra.get("style_application"):
narrate_trace["style_application"] = extra["style_application"]
diag = narrate_result.get("diagnostics") or {}
if diag.get("generate_ms") is not None:
narrate_trace["generate_ms"] = diag.get("generate_ms")
if diag.get("prompt_revision"):
narrate_trace["prompt_revision"] = diag.get("prompt_revision")
if diag.get("generate_calls") is not None:
narrate_trace.setdefault("generate_calls", diag.get("generate_calls"))
if diag.get("model"):
narrate_trace.setdefault("model", diag.get("model"))
if diag.get("provider"):
narrate_trace.setdefault("provider", diag.get("provider"))
if diag.get("completion_tokens") is not None:
narrate_trace.setdefault("completion_tokens", diag.get("completion_tokens"))
if extra:
narrate_trace.update(extra)
return {
**narrate_trace,
"prompt_revision": narrate_trace.get("prompt_revision"),
"writing_profile": profile_meta,
"style_examples": narrate_trace.get("style_examples"),
"style_application": narrate_trace.get("style_application") or (extra or {}).get("style_application"),
"dropped_optional_blocks": dropped,
"log": [dict(item) if isinstance(item, dict) else item for item in run_log],
"stages": [reconstruct_trace, narrate_trace],
}
def _log_items(source: dict | BaseException | None) -> list[dict]:
if source is None:
return []
if isinstance(source, BaseException):
diag = getattr(source, "diagnostics", None) or {}
return list(diag.get("log") or [])
if not isinstance(source, dict):
return []
if source.get("log"):
return list(source["log"])
trace = source.get("trace") or {}
if trace.get("log"):
return list(trace["log"])
diag = source.get("diagnostics") or {}
return list(diag.get("log") or [])
def _stamp_log(run_log: list[dict], stage: str, items: list[dict]) -> None:
for item in items:
row = dict(item)
row.setdefault("stage", stage)
run_log.append(row)
def _event(run_log: list[dict], stage: str, kind: str, **fields) -> None:
row = {"stage": stage, "kind": kind}
for key, value in fields.items():
if value is not None and value != "":
row[key] = value
run_log.append(row)
def _raise_budget(exc: JournalBudgetError) -> None:
raise EngineError(exc.code, exc.message, exc.status_code, exc.diagnostics) from exc
def _stage_trace(result: dict, fallback_purpose: str) -> dict:
source = result.get("trace")
trace = dict(source) if isinstance(source, dict) else {}
budget_src = result.get("diagnostics")
if not isinstance(budget_src, dict):
budget_src = trace.get("budget") if isinstance(trace.get("budget"), dict) else None
if isinstance(budget_src, dict):
trace["budget"] = compact_diagnostics(budget_src)
else:
trace["budget"] = None
if "log" in trace and isinstance(trace["log"], list):
trace["log"] = [dict(item) if isinstance(item, dict) else item for item in trace["log"]]
trace["purpose"] = trace.get("purpose") or fallback_purpose
return trace
def _local_stage_trace(artifact: dict, *, source_count: int) -> dict:
return {
"purpose": "local_source_artifact",
"layer": "journalquellen",
"provider": None,
"model": None,
"coverage": artifact.get("coverage") or "all_selected_sources",
"status": "local_ok",
"stage1": "local_ok",
"source_count": source_count,
"intern": reconstruction_text(artifact),
"budget": {
"purpose": "local_source_artifact",
"prompt_tokens": 0,
"completion_tokens": 0,
"total_tokens": 0,
"cost": 0,
"budget_ok": True,
"status": "local",
"context_compression": "not_applicable",
},
}
def _existing_text(profile_id: str, journal_day_id: str) -> str:
"""Saved user entries are the Fassung. A leftover draft is only used if none exist."""
entries = current_entries(profile_id, journal_day_id)
bodies = [(item.get("body") or "").strip() for item in entries if (item.get("body") or "").strip()]
if bodies:
return "\n\n".join(bodies)
draft = current_draft(profile_id, journal_day_id)
if draft:
return (draft.get("body") or "").strip()
return ""
def pack_narration_context(
prompt: dict,
budget,
assembled: dict[str, str],
*,
style_examples: str,
existing_text: str,
include_existing: bool,
) -> tuple[dict[str, str], list[str]]:
"""Drop optional blocks locally if they overflow. Never drop day sources or the profile.
No extra model call. Order: try all, then drop style examples, then existing text.
"""
dropped: list[str] = []
attempts = [
(style_examples, existing_text if include_existing else ""),
("", existing_text if include_existing else ""),
("", ""),
]
last_error: JournalBudgetError | None = None
for examples, existing in attempts:
candidate = dict(assembled)
candidate["style_examples"] = examples
candidate["existing_text"] = existing
fit_prompt = prompt_without_orphaned_style_example_rule(
prompt,
include_style_examples=bool((examples or "").strip()),
)
rendered = preview_prompt(fit_prompt, candidate)["rendered"]
try:
assert_input_fits(budget, rendered)
if not examples and style_examples:
dropped.append("style_examples")
if include_existing and existing_text and not existing:
dropped.append("existing_text")
return candidate, dropped
except JournalBudgetError as exc:
last_error = exc
continue
if last_error:
raise last_error
return assembled, dropped
def generate_draft(
profile_id: str,
journal_day_id: str,
conversation_ids: list[str] | None = None,
include_existing: bool = False,
explicit: bool = True,
generation_selection: dict | None = None,
remember_generation_selection: bool = False,
) -> dict:
try:
draft = _generate_draft(
profile_id,
journal_day_id,
conversation_ids=conversation_ids,
include_existing=include_existing,
explicit=explicit,
generation_selection=generation_selection,
remember_generation_selection=remember_generation_selection,
)
except EngineError as exc:
persist_engine_error(
profile_id,
purpose="journal_generate",
exc=exc,
subject_type="journal_day",
subject_id=journal_day_id,
journal_day_id=journal_day_id,
extra={"conversation_ids": conversation_ids or []},
)
raise
persist_step(
profile_id,
purpose="journal_generate",
status="ok",
subject_type="journal_day",
subject_id=journal_day_id,
journal_day_id=journal_day_id,
conversation_id=(conversation_ids[0] if conversation_ids else None),
decision=None,
trace=draft.get("trace"),
extra={
"draft_id": draft.get("id"),
"run_log": draft.get("run_log"),
"conversation_ids": conversation_ids or [],
"stored_title": draft.get("title"),
"stored_body": draft.get("body"),
},
)
return draft
def _generate_draft(
profile_id: str,
journal_day_id: str,
conversation_ids: list[str] | None = None,
include_existing: bool = False,
explicit: bool = True,
generation_selection: dict | None = None,
remember_generation_selection: bool = False,
) -> dict:
require_explicit_generate(explicit)
day = get_day(profile_id, journal_day_id)
day_conversations = list_conversations_for_day(profile_id, journal_day_id)
selected = source_conversation_ids(
conversation_ids,
[item["id"] for item in day_conversations],
)
remember_dialogue_style(profile_id, exclude_conversation_ids=selected)
existing_text = _existing_text(profile_id, journal_day_id) if include_existing else ""
selection_ids, persist_meta = resolve_run_selection(
profile_id,
generation_selection,
remember_generation_selection,
)
narrate_prompt = load_active_prompt("mvp.journal_generate")
assert_journal_prompt_contract(narrate_prompt.get("template") or "")
config = generate_provider()
if not config:
raise EngineError(
"no_egress_provider_configured",
"Persönlicher KI-Aufruf wurde vom Privacy Gateway blockiert. "
"Es ist kein Egress-Provider konfiguriert.",
503,
)
try:
window = resolve_generate_metadata(config)
narrate_budget = plan_journal_budget(window, purpose="journal_generate")
except JournalBudgetError as exc:
_raise_budget(exc)
try:
source_messages = retrieve(
profile_id,
{
"kind": "day_messages",
"journal_day_id": journal_day_id,
"conversation_ids": selected,
"overflow": "abort",
},
)
reconstruction = local_verified_artifact(source_messages)
except JournalBudgetError as exc:
_raise_budget(exc)
source_user = [
message.get("body") or ""
for message in assign_source_ids(source_messages)
if message.get("role") == "user" and (message.get("body") or "").strip()
]
compiled_policy = compile_selection(selection_ids)
policy_meta = policy_trace(
compiled_policy,
source=persist_meta["source"],
remembered=persist_meta["remembered"],
)
style_bundle = compile_journal_style_application(
profile_id,
compiled_policy.style_context,
exclude_dates=[day.get("calendar_date") or ""],
)
profile_meta = writing_profile_trace(
profile_id,
style_context=compiled_policy.style_context,
brief=style_bundle.get("brief") or "",
trait_count=(style_bundle.get("effective") or {}).get("trait_count"),
)
style_example_rows = list(style_bundle.get("example_rows") or [])
style_examples = style_bundle.get("examples_block") or ""
writing_profile = style_bundle.get("profile_block") or ""
style_meta = style_example_diagnostics(style_example_rows)
run_log: list[dict] = []
_event(
run_log,
"local_source_artifact",
"stage1_result",
stage1="local_ok",
coverage=reconstruction.get("coverage") or "all_selected_sources",
status="local_ok",
source_count=len(source_user),
writing_profile_present=profile_meta.get("present"),
style_example_count=style_meta.get("count"),
generation_selection_source=policy_meta.get("source"),
generation_selection_keys=",".join(
f"{slot}:{policy_meta.get('keys', {}).get(slot)}"
for slot in ("transformation", "detail", "voice", "narrative")
),
)
reconstruct_result = {
"trace": _local_stage_trace(reconstruction, source_count=len(source_user)),
"content": reconstruction_text(reconstruction),
}
narrate_context = build_internal_context(
profile_id,
space_id=day["space_id"],
journal_day_id=journal_day_id,
purpose="journal_generate",
include_existing=include_existing,
conversation_ids=selected,
existing_text=existing_text,
conversation_id=selected[0] if selected else None,
reconstruction=narration_sources_text(reconstruction),
writing_profile=writing_profile,
style_examples=style_examples,
transformation_instructions=compiled_policy.instructions["transformation_instructions"],
detail_instructions=compiled_policy.instructions["detail_instructions"],
voice_instructions=compiled_policy.instructions["voice_instructions"],
narrative_instructions=compiled_policy.instructions["narrative_instructions"],
)
assembled = assemble_text(narrate_context)
try:
assembled, dropped = pack_narration_context(
narrate_prompt,
narrate_budget,
assembled,
style_examples=style_examples,
existing_text=existing_text,
include_existing=include_existing,
)
except JournalBudgetError as exc:
_raise_budget(exc)
if dropped:
_event(run_log, "journal_generate", "budget_pack", dropped=",".join(dropped))
style_app = style_application_trace(compiled_policy, style_bundle, dropped=dropped)
if "style_examples" in dropped:
style_meta = {**style_example_diagnostics([]), "dropped": True, "dropped_reason": "budget"}
else:
style_meta = {**style_meta, "dropped": False}
profile_meta["brief_chars"] = len((assembled.get("writing_profile") or "").strip())
mappings = list_mappings(profile_id)
preview_title, preview_body = _local_narration(reconstruction)
source_preview = f"{preview_title}\n\n{preview_body}".strip()
narrate_result: dict = {}
try:
narrate_prompt = prompt_without_orphaned_style_example_rule(
narrate_prompt,
include_style_examples=bool((assembled.get("style_examples") or "").strip()),
)
rendered_preview = preview_prompt(narrate_prompt, assembled)
assert_template_resolved(rendered_preview.get("rendered") or "")
narrate_result = execute_prompt(
narrate_prompt,
profile_id,
purpose="journal_generate",
data_class="B",
context=assembled,
max_tokens=narrate_budget.reserved_output_tokens,
disable_context_compression=True,
budget=narrate_budget,
)
_stamp_log(run_log, "journal_generate", _log_items(narrate_result))
title, body = _split_title(narrate_result.get("content") or "")
seen = {(item.get("local_label") or "").casefold() for item in mappings}
for item in narrate_result.get("local_identities") or []:
label = (item.get("local_label") or "").strip()
if not label or label.casefold() in seen:
continue
mappings.append(item)
seen.add(label.casefold())
active_tokens = list((narrate_result.get("diagnostics") or {}).get("active_tokens") or [])
combined = f"{title}\n\n{body}".strip()
unattested = unattested_journal_content(combined, source_user, mappings, active_tokens)
if unattested:
_event(
run_log,
"journal_generate",
"narration_result",
status="not_accepted",
reason=unattested,
)
bundle = _compose_journal_trace(
reconstruct_result,
narrate_result,
run_log=run_log,
narrate_prompt=narrate_prompt,
profile_meta=profile_meta,
style_meta=style_meta,
dropped=dropped,
extra={
"narration_source": "not_accepted",
"model_text_accepted": False,
"provenance_decision": unattested,
"abort_reason": unattested,
"source_preview": source_preview,
"generation_selection": policy_meta,
"style_application": style_app,
},
)
_reject_generation(
reason=unattested,
diagnostics={
**(narrate_result.get("diagnostics") or {}),
"trace": bundle,
"log": run_log,
"provenance_decision": unattested,
},
source_preview=source_preview,
extra=bundle,
)
_event(
run_log,
"journal_generate",
"narration_result",
status="model",
lexical_similarity=lexical_similarity("\n".join(source_user), body),
incomplete_syntax=incomplete_syntax_markers(body),
writing_profile_present=profile_meta.get("present"),
)
except EngineError as exc:
if exc.code == JOURNAL_NOT_ACCEPTED:
raise
_stamp_log(run_log, "journal_generate", _log_items(exc))
diag = dict(exc.diagnostics or {})
generate_called = bool(diag.get("generate_called") or diag.get("generate_calls"))
failed_result = {
"diagnostics": diag,
"trace": dict(diag.get("trace") or {}),
"content": "",
}
bundle = _compose_journal_trace(
reconstruct_result,
failed_result,
run_log=run_log,
narrate_prompt=narrate_prompt,
profile_meta=profile_meta,
style_meta=style_meta,
dropped=dropped,
extra={
"narration_source": "not_accepted",
"model_text_accepted": False,
"abort_reason": exc.code,
"source_preview": source_preview,
"generation_selection": policy_meta,
"style_application": style_app,
},
)
diag["log"] = run_log
diag["trace"] = bundle
if not generate_called:
raise EngineError(exc.code, exc.message, exc.status_code, diag) from exc
_event(
run_log,
"journal_generate",
"narration_result",
status="not_accepted",
reason=exc.code,
)
_reject_generation(
reason=exc.code,
diagnostics=diag,
source_preview=source_preview,
extra=bundle,
)
person_labels = [
(item.get("local_label") or "").strip()
for item in mappings
if (item.get("local_label") or "").strip()
and str(item.get("token") or "").upper().lstrip("[").startswith("PERSON:")
]
title, body = shape_journal(
title,
body,
source_user,
person_labels=person_labels,
source="model",
)
stored = f"{title}\n\n{body}".strip() if title else (body or "").strip()
before_entries = {item["id"]: item.get("current_version_id") for item in current_entries(profile_id, journal_day_id)}
model = (
((narrate_result.get("diagnostics") or {}).get("actual_model"))
or ((narrate_result.get("trace") or {}).get("model"))
or narrate_result.get("provider")
or ""
)
snapshot = draft_snapshot(compiled_policy, prompt=narrate_prompt, model=str(model or ""))
mark_guidelines_used(
[
compiled_policy.ids["transformation"],
compiled_policy.ids["detail"],
compiled_policy.ids["voice"],
compiled_policy.ids["narrative"],
]
)
draft = insert_draft(
profile_id,
journal_day_id,
title=title,
body=body,
source_conversation_ids=selected,
source_message_ids=_message_ids(profile_id, selected),
generation_snapshot=snapshot,
)
after_entries = current_entries(profile_id, journal_day_id)
for item in after_entries:
previous = before_entries.get(item["id"])
if previous is not None and previous != item.get("current_version_id"):
raise StoreError("policy_violation", "Generate darf die Nutzerfassung nicht verändern", 500)
bundle = _compose_journal_trace(
reconstruct_result,
narrate_result,
run_log=run_log,
narrate_prompt=narrate_prompt,
profile_meta=profile_meta,
style_meta=style_meta,
dropped=dropped,
extra={
"narration_source": "model",
"model_text_accepted": True,
"provenance_decision": "accepted",
"lexical_similarity": lexical_similarity("\n".join(source_user), body),
"incomplete_syntax": incomplete_syntax_markers(body),
"generation_selection": policy_meta,
"style_application": style_app,
"stored_title": title,
"stored_body": body,
"reply": stored,
},
)
draft["run_log"] = run_log
draft["trace"] = bundle
return draft