Kansho/backend/profile_review.py
2026-08-25 13:57:23 +02:00

1120 lines
40 KiB
Python

"""Cost-conscious profile review: local drift, bundled evidence, optional AI/paste.
Does not rebuild the writing profile after every new source. Dialogue turns never
trigger an LLM review. Interaction preferences are never inferred from silence.
Semantic traits come from review, not from word counts. Continuous learning starts
only after a confirmed Initial Profile Build.
"""
from __future__ import annotations
import json
import re
import uuid
from datetime import datetime, timedelta, timezone
from db import get_db, row_to_dict
from dialogue_store import StoreError
from journal_body import plain_text
from journal_store import list_versions
from writing_profile_infer import SIGNAL_KEYS, infer_features
from writing_profile_schema import (
ACTION_ALIASES,
EXISTING_BEFORE_NEW,
LEGACY_STYLE_KEYS,
SEED_FACETS,
TRAIT_ACTIONS,
coerce_slug,
facet_label,
hint_for_context,
normalize_facet_key,
seed_catalog,
valid_slug,
)
from writing_profile_store import (
GOVERNANCE,
INITIAL_BUILD_SOURCES,
_assemble_brief,
_queue_suggestion,
_upsert_facet,
ensure_profile,
get_profile,
has_confirmed_profile,
import_text,
list_corpus,
refresh_profile,
replace_trait_refs,
retire_trait,
snapshot_version,
upsert_trait,
)
KIND_PACKAGE = "kansho.profile_review_package"
KIND_RESULT = "kansho.profile_review_result"
FORMAT_VERSION = 1
JSON_BLOCK = re.compile(r"\{.*\}", re.DOTALL)
TRIGGERS = {
"user_edit_of_draft",
"heuristic_drift",
"dialogue_sample",
"periodic_batch",
"explicit_review",
}
CORE_KEY = "core"
JOURNAL_FACET = "autobiographical_journal"
BATCH_MIN = 3
MAX_BUNDLE = 8
EXCERPT_CHARS = 480
EXAMPLE_CHARS = 220
CORPUS_EXCERPT = 360
DIALOGUE_MIN_WORDS = 60
DIALOGUE_COOLDOWN_HOURS = 24
DRIFT_MIN = 2
PERIODIC_DAYS = 7
LAYERS = {"core", "facet", "trait"}
INSTRUCTION_SHARED = (
"Keine Diagnose, kein Persönlichkeitsmodell, keine erfundenen Stilmerkmale. "
"Lokale Heuristiken (Wortzählungen, Signalwörter) sind Messhilfe, keine Traits. "
"Übernimm Heuristik-Keys wie chronology, humor, detail, rhythm nicht als festes Raster. "
"Jeder semantische Trait braucht echte Source-Evidence und darf repräsentative Originalbeispiele nennen. "
"Urlaubstagebücher und autobiografische Journale belegen primär die Facet autobiographical_journal. "
"Einen globalen Core nur so weit verallgemeinern, wie unterschiedliche Quellenarten das tragen. "
"Neuere Texte wiegen stärker für die aktuelle Ausprägung; ältere Texte belegen stabile Langzeitmerkmale. "
"Existing-before-New: " + "; ".join(EXISTING_BEFORE_NEW) + "."
)
INSTRUCTION_INCREMENTAL = (
"Vergleiche die neuen Evidenzen mit dem bestehenden Profil. "
"Erzeuge kein komplett neues Profil. "
"Unterscheide Core, optionale context/output-Facets und dynamische semantische Traits darin. "
"Schlage nur Änderungen vor, die die Evidenz trägt. "
+ INSTRUCTION_SHARED
)
INSTRUCTION_INITIAL = (
"Bilde ein erstes Writing Profile aus dem historischen Korpus, nicht nur aus zwei aktuellen Einträgen. "
"Die Hülle ist Core plus optionale Facets plus dynamische Traits, Evidence References und Representative Exemplars. "
"Seed-Keys sind nur Ordnungshilfe. "
+ INSTRUCTION_SHARED
)
def _parse_stamp(raw: str | None) -> datetime | None:
text = (raw or "").strip()
if not text:
return None
try:
stamp = datetime.fromisoformat(text.replace("Z", "+00:00"))
except ValueError:
return None
if stamp.tzinfo is None:
stamp = stamp.replace(tzinfo=timezone.utc)
return stamp
def _now() -> datetime:
return datetime.now(timezone.utc)
def _parse_json(raw: str | None, fallback):
if not raw:
return fallback
try:
data = json.loads(raw)
except json.JSONDecodeError:
return fallback
return data if data is not None else fallback
def facet_delta(current: dict[str, str], incoming: dict[str, str]) -> list[str]:
"""Local-signal delta. Not a semantic trait diff."""
changed: list[str] = []
for key, value in incoming.items():
old = (current.get(key) or "").strip()
new = (value or "").strip()
if old and new and old != new:
changed.append(key)
elif not old and new:
changed.append(key)
return changed
def is_strong_edit(previous: str, current: str) -> bool:
old = plain_text(previous or "").strip()
new = plain_text(current or "").strip()
if not new:
return False
if not old:
return len(new.split()) >= 40
old_words = set(old.lower().split())
new_words = set(new.lower().split())
union = max(len(old_words | new_words), 1)
overlap = len(old_words & new_words) / union
length = abs(len(new) - len(old)) / max(len(old), 1)
return overlap < 0.55 or length > 0.3
def _heuristic_baseline(profile_id: str, *, exclude_entry_id: str | None = None) -> dict[str, str]:
texts = [
item.get("body") or ""
for item in list_corpus(profile_id, limit=12)
if (item.get("entry_id") or "") != (exclude_entry_id or "")
]
return infer_features(texts)
def classify_journal_trigger(profile_id: str, entry_id: str, body: str, origin: str | None) -> str | None:
text = plain_text(body or "").strip()
if not text:
return None
origin = (origin or "").strip()
if origin == "user_edit" and entry_id:
versions = list_versions(profile_id, entry_id)
if len(versions) >= 2:
previous = versions[-2]
if (previous.get("origin") or "") == "accepted_draft" and is_strong_edit(
previous.get("body") or "", text
):
return "user_edit_of_draft"
incoming = infer_features([text])
baseline = _heuristic_baseline(profile_id, exclude_entry_id=entry_id)
if baseline and len(facet_delta(baseline, incoming)) >= DRIFT_MIN:
return "heuristic_drift"
return None
def enqueue_evidence(
profile_id: str,
*,
trigger: str,
excerpt: str,
source_kind: str = "",
source_id: str | None = None,
) -> dict | None:
if trigger not in TRIGGERS:
return None
text = plain_text(excerpt or "").strip()
if not text:
return None
row = ensure_profile(profile_id)
if (row.get("governance") or "learning") == "frozen" and trigger != "explicit_review":
return None
signals = infer_features([text])
example_for = [key for key in signals if key in SIGNAL_KEYS][:6]
excerpt = text[:EXCERPT_CHARS]
with get_db() as conn:
existing = row_to_dict(
conn.execute(
"""
SELECT id FROM writing_profile_evidence
WHERE profile_id = ? AND trigger = ? AND excerpt = ? AND status = 'pending'
LIMIT 1
""",
(profile_id, trigger, excerpt),
).fetchone()
)
if existing:
return existing
item = {
"id": str(uuid.uuid4()),
"profile_id": profile_id,
"trigger": trigger,
"source_kind": source_kind,
"source_id": source_id,
"excerpt": excerpt,
"local_signals_json": json.dumps(signals, ensure_ascii=False),
"example_for_json": json.dumps(example_for, ensure_ascii=False),
"status": "pending",
}
conn.execute(
"""
INSERT INTO writing_profile_evidence
(id, profile_id, trigger, source_kind, source_id, excerpt, local_signals_json, example_for_json, status)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
item["id"],
profile_id,
trigger,
source_kind,
source_id,
excerpt,
item["local_signals_json"],
item["example_for_json"],
item["status"],
),
)
_refresh_ready(profile_id)
return item
def _pending_evidence(profile_id: str, limit: int = MAX_BUNDLE) -> list[dict]:
with get_db() as conn:
rows = [
row_to_dict(item)
for item in conn.execute(
"""
SELECT * FROM writing_profile_evidence
WHERE profile_id = ? AND status = 'pending'
ORDER BY created
LIMIT ?
""",
(profile_id, limit),
).fetchall()
]
for item in rows:
item["local_signals"] = _parse_json(item.get("local_signals_json"), {})
item["example_for"] = _parse_json(item.get("example_for_json"), [])
return rows
def _refresh_ready(profile_id: str) -> None:
with get_db() as conn:
pending = conn.execute(
"SELECT COUNT(*) AS c FROM writing_profile_evidence WHERE profile_id = ? AND status = 'pending'",
(profile_id,),
).fetchone()["c"]
row = row_to_dict(
conn.execute(
"SELECT last_reviewed, lifecycle FROM writing_profiles WHERE profile_id = ?",
(profile_id,),
).fetchone()
)
if (row or {}).get("lifecycle") != "confirmed":
with get_db() as conn:
conn.execute("UPDATE writing_profiles SET review_ready = 0 WHERE profile_id = ?", (profile_id,))
return
last = (row or {}).get("last_reviewed") or ""
stamp = _parse_stamp(last)
stale = False
if stamp:
stale = _now() - stamp >= timedelta(days=PERIODIC_DAYS) and pending >= 1
elif last:
stale = pending >= BATCH_MIN
else:
stale = pending >= 1
ready = 1 if pending >= BATCH_MIN or stale else 0
with get_db() as conn:
conn.execute(
"UPDATE writing_profiles SET review_ready = ? WHERE profile_id = ?",
(ready, profile_id),
)
def ingest_journal(profile_id: str, entry_id: str, body: str, origin: str | None = None) -> dict:
"""After a saved entry: assemble sources. Queue evidence only once the profile is confirmed."""
refresh_profile(profile_id)
if not has_confirmed_profile(profile_id):
return get_profile(profile_id)
trigger = classify_journal_trigger(profile_id, entry_id, body, origin)
if trigger:
enqueue_evidence(
profile_id,
trigger=trigger,
excerpt=body,
source_kind="journal_entry",
source_id=entry_id,
)
elif origin == "imported_text":
enqueue_evidence(
profile_id,
trigger="heuristic_drift",
excerpt=body,
source_kind="imported_text",
source_id=None,
)
return get_profile(profile_id)
def ingest_import(
profile_id: str,
body: str,
*,
occurred_at: str | None = None,
context_hint: str | None = None,
) -> dict:
import_text(profile_id, body, occurred_at=occurred_at, context_hint=context_hint)
if not has_confirmed_profile(profile_id):
return get_profile(profile_id)
enqueue_evidence(
profile_id,
trigger="heuristic_drift",
excerpt=body,
source_kind="imported_text",
)
return get_profile(profile_id)
def consider_dialogue(profile_id: str, last_user: str) -> None:
"""Sample only. Never an LLM call, never an interaction preference, never before confirm."""
if not has_confirmed_profile(profile_id):
return
text = plain_text(last_user or "").strip()
if len(text.split()) < DIALOGUE_MIN_WORDS:
return
row = ensure_profile(profile_id)
if (row.get("governance") or "learning") == "frozen":
return
cutoff = (_now() - timedelta(hours=DIALOGUE_COOLDOWN_HOURS)).strftime("%Y-%m-%d %H:%M:%S")
with get_db() as conn:
recent = conn.execute(
"""
SELECT id FROM writing_profile_evidence
WHERE profile_id = ? AND trigger = 'dialogue_sample' AND created >= ?
LIMIT 1
""",
(profile_id, cutoff),
).fetchone()
if recent:
return
enqueue_evidence(
profile_id,
trigger="dialogue_sample",
excerpt=text,
source_kind="dialogue_style",
)
def review_status(profile_id: str) -> dict:
ensure_profile(profile_id)
with get_db() as conn:
pending = [
row_to_dict(item)
for item in conn.execute(
"""
SELECT id, trigger, source_kind, created, status
FROM writing_profile_evidence
WHERE profile_id = ? AND status = 'pending'
ORDER BY created
""",
(profile_id,),
).fetchall()
]
open_review = row_to_dict(
conn.execute(
"""
SELECT id, channel, status, created FROM writing_profile_reviews
WHERE profile_id = ? AND status IN ('open', 'proposed')
ORDER BY created DESC LIMIT 1
""",
(profile_id,),
).fetchone()
)
versions = [
row_to_dict(item)
for item in conn.execute(
"""
SELECT seq, cause, created FROM writing_profile_versions
WHERE profile_id = ? ORDER BY seq DESC LIMIT 8
""",
(profile_id,),
).fetchall()
]
row = row_to_dict(
conn.execute(
"""
SELECT version, review_ready, last_reviewed, governance, lifecycle
FROM writing_profiles WHERE profile_id = ?
""",
(profile_id,),
).fetchone()
)
lifecycle = (row or {}).get("lifecycle") or "uninitialized"
return {
"version": (row or {}).get("version") or 0,
"review_ready": bool((row or {}).get("review_ready")),
"last_reviewed": (row or {}).get("last_reviewed"),
"governance": (row or {}).get("governance") or "learning",
"lifecycle": lifecycle,
"confirmed": lifecycle == "confirmed",
"pending_evidence": pending,
"open_review": open_review,
"versions": versions,
"existing_before_new": list(EXISTING_BEFORE_NEW),
"note": (
"Erst nach einem bestätigten Initial Profile Build beginnt kontinuierliches Learning. "
if lifecycle != "confirmed"
else "Lokale Heuristik entscheidet, ob eine Review nötig ist. Kein Call nach jedem Dialog."
),
}
def _review_mode(profile_id: str, requested: str | None = None) -> str:
if requested in {"initial_build", "incremental"}:
return requested
return "incremental" if has_confirmed_profile(profile_id) else "initial_build"
def _corpus_supports_core(corpus: list[dict]) -> bool:
hints = set()
for item in corpus:
hint = hint_for_context(item.get("facet_hint") or item.get("context_hint"))
if not hint and item.get("kind") == "journal_entry":
hint = JOURNAL_FACET
if hint:
hints.add(hint)
return len(hints) >= 2
def _public_trait(item: dict) -> dict:
return {
"id": item.get("id"),
"slug": item.get("slug"),
"label": item.get("label") or item.get("slug"),
"facet_key": item.get("facet_key") or CORE_KEY,
"statement": item.get("statement") or "",
"origin": item.get("origin") or "manual",
"locked": bool(item.get("locked")),
"status": item.get("status") or "active",
"observed_from": item.get("observed_from"),
"observed_to": item.get("observed_to"),
"evidence_refs": [
{
"excerpt": ref.get("excerpt") or "",
"occurred_at": ref.get("occurred_at"),
"source_id": ref.get("source_id"),
}
for ref in item.get("evidence_refs") or []
if ref.get("excerpt")
],
"exemplars": [
{
"excerpt": ref.get("excerpt") or "",
"occurred_at": ref.get("occurred_at"),
"source_id": ref.get("source_id"),
}
for ref in item.get("exemplars") or []
if ref.get("excerpt")
],
}
def _public_facet(item: dict | None) -> dict:
if not item:
return {}
key = item.get("facet_key")
return {
"key": key,
"label": item.get("label") or facet_label(key or ""),
"layer": item.get("layer") or "context",
"value": item.get("value") or "",
"origin": item.get("origin") or "manual",
"locked": bool(item.get("locked")),
"traits": [_public_trait(trait) for trait in item.get("traits") or []],
}
def _public_evidence(item: dict) -> dict:
return {
"id": item.get("id"),
"trigger": item.get("trigger"),
"source_kind": item.get("source_kind") or "",
"source_id": item.get("source_id"),
"excerpt": item.get("excerpt") or "",
"local_signals": item.get("local_signals") or {},
"example_for": item.get("example_for") or [],
}
def _public_corpus(item: dict) -> dict:
excerpt = plain_text(item.get("body") or "")[:CORPUS_EXCERPT]
return {
"id": item.get("id"),
"kind": item.get("kind"),
"occurred_at": item.get("occurred_at"),
"recency_weight": item.get("recency_weight"),
"recency_role": item.get("recency_role"),
"context_hint": item.get("context_hint") or "",
"facet_hint": item.get("facet_hint") or "",
"excerpt": excerpt,
}
def build_package(
profile_id: str,
*,
trigger: str = "explicit_review",
target: str = "writing",
mode: str | None = None,
) -> dict:
if target != "writing":
raise StoreError("unsupported_target", "Nur Writing-Profile-Review ist gebündelt; Interaction bleibt explizit.")
ensure_profile(profile_id)
mode = _review_mode(profile_id, mode)
if trigger == "explicit_review":
enqueue_evidence(
profile_id,
trigger="explicit_review",
excerpt=(
"Expliziter Initial Profile Build aus dem historischen Korpus."
if mode == "initial_build"
else "Explizite Nutzer-Review des Writing Profile."
),
source_kind="initial_build" if mode == "initial_build" else "explicit",
)
if mode == "initial_build":
with get_db() as conn:
conn.execute(
"UPDATE writing_profiles SET lifecycle = 'initial_pending', updated = datetime('now') WHERE profile_id = ?",
(profile_id,),
)
evidences = _pending_evidence(profile_id)
corpus = list_corpus(profile_id, limit=INITIAL_BUILD_SOURCES if mode == "initial_build" else 12)
if not evidences and not corpus:
raise StoreError("no_evidence", "Keine Review-Evidenz und kein Korpus vorhanden.")
profile = get_profile(profile_id)
facets = profile.get("facets") or []
core = next((item for item in facets if item.get("facet_key") == CORE_KEY or item.get("layer") == "core"), None)
others = [item for item in facets if item is not core]
package = {
"kind": KIND_PACKAGE,
"format_version": FORMAT_VERSION,
"target": "writing",
"mode": mode,
"instruction": INSTRUCTION_INITIAL if mode == "initial_build" else INSTRUCTION_INCREMENTAL,
"existing_before_new": list(EXISTING_BEFORE_NEW),
"seed_catalog": seed_catalog(),
"profile": {
"version": profile.get("version") or 0,
"governance": profile.get("governance") or "learning",
"lifecycle": profile.get("lifecycle") or "uninitialized",
"core": _public_facet(core),
"facets": [_public_facet(item) for item in others],
"traits": [_public_trait(item) for item in profile.get("traits") or []],
},
"corpus": [_public_corpus(item) for item in corpus],
"evidences": [_public_evidence(item) for item in evidences],
"expected_result": {
"kind": KIND_RESULT,
"format_version": FORMAT_VERSION,
"target": "writing",
"mode": mode,
"changes": [
{
"layer": "core|facet|trait",
"key": "core|autobiographical_journal|dynamischer-slug",
"slug": "datengetriebener Trait-Slug",
"facet_key": "core|autobiographical_journal|…",
"action": "confirm|precisify|rescope|merge|split|create",
"proposed_value": "Aussage des Traits oder der Facet",
"label": "optional",
"rationale": "kurz",
"evidence_ids": ["id"],
"exemplars": [{"excerpt": "Originalzitat", "occurred_at": "YYYY-MM-DD"}],
"merge_slugs": ["optional"],
"split_into": [{"slug": "neu", "statement": "", "facet_key": ""}],
}
],
},
}
return package
def render_paste_prompt(package: dict) -> str:
body = json.dumps(package, ensure_ascii=False, indent=2)
return (
"Kanshō Writing Profile Review. "
"Antworte ausschließlich mit einem JSON-Objekt gemäß expected_result. "
"Kein Fließtext davor oder danach.\n\n"
f"{body}\n"
)
def parse_result(raw) -> dict:
if isinstance(raw, dict):
data = raw
else:
text = (raw or "").strip()
if text.startswith("```"):
text = re.sub(r"^```(?:json)?\s*|\s*```$", "", text, flags=re.I | re.S)
match = JSON_BLOCK.search(text)
if not match:
raise StoreError("invalid_review_result", "Kein JSON-Review-Ergebnis.")
try:
data = json.loads(match.group(0))
except json.JSONDecodeError as exc:
raise StoreError("invalid_review_result", "Review-Ergebnis ist kein gültiges JSON.") from exc
if not isinstance(data, dict):
raise StoreError("invalid_review_result", "Review-Ergebnis muss ein Objekt sein.")
if (data.get("kind") or "") != KIND_RESULT:
raise StoreError("invalid_review_result", "kind muss kansho.profile_review_result sein.")
version = data.get("format_version", FORMAT_VERSION)
try:
version = int(version)
except (TypeError, ValueError) as exc:
raise StoreError("unsupported_format", "format_version ungültig.") from exc
if version != FORMAT_VERSION:
raise StoreError("unsupported_format", f"format_version {version} wird nicht unterstützt.")
changes = data.get("changes")
if not isinstance(changes, list):
raise StoreError("invalid_review_result", "changes muss eine Liste sein.")
cleaned = []
for item in changes:
if not isinstance(item, dict):
continue
layer = (item.get("layer") or "trait").strip()
if layer not in LAYERS:
continue
action = ACTION_ALIASES.get((item.get("action") or "confirm").strip(), (item.get("action") or "confirm").strip())
if action not in TRAIT_ACTIONS:
continue
key = (item.get("key") or item.get("slug") or "").strip()
slug = coerce_slug(item.get("slug") or key or item.get("label") or "trait")
if layer == "core":
key = CORE_KEY
slug = CORE_KEY
elif layer == "facet":
key = normalize_facet_key(key) or JOURNAL_FACET
slug_hint = item.get("slug") or ""
as_trait = bool(slug_hint) or action == "create" or key in LEGACY_STYLE_KEYS or (
valid_slug(key) and key not in SEED_FACETS and key != CORE_KEY
)
if as_trait:
layer = "trait"
slug = coerce_slug(slug_hint or key)
else:
slug = coerce_slug(key)
elif not valid_slug(slug):
continue
facet_key = normalize_facet_key(item.get("facet_key") or (key if layer == "facet" else "") or CORE_KEY) or CORE_KEY
if layer == "core":
facet_key = CORE_KEY
exemplars = []
for ex in item.get("exemplars") or []:
if isinstance(ex, str) and ex.strip():
exemplars.append({"excerpt": ex.strip(), "role": "exemplar"})
elif isinstance(ex, dict) and (ex.get("excerpt") or "").strip():
exemplars.append(
{
"excerpt": (ex.get("excerpt") or "").strip(),
"role": ex.get("role") or "exemplar",
"source_id": ex.get("source_id"),
"occurred_at": ex.get("occurred_at"),
}
)
split_into = []
for part in item.get("split_into") or []:
if not isinstance(part, dict):
continue
part_slug = coerce_slug(part.get("slug") or part.get("label") or "")
if not part_slug:
continue
split_into.append(
{
"slug": part_slug,
"label": part.get("label") or part_slug,
"statement": part.get("statement") or part.get("proposed_value") or "",
"facet_key": normalize_facet_key(part.get("facet_key") or facet_key) or facet_key,
}
)
cleaned.append(
{
"layer": layer,
"key": key,
"slug": slug,
"facet_key": facet_key,
"action": action,
"proposed_value": item.get("proposed_value") or item.get("statement") or "",
"label": item.get("label") or "",
"rationale": item.get("rationale") or "",
"evidence_ids": [str(eid) for eid in (item.get("evidence_ids") or []) if eid],
"exemplars": exemplars,
"merge_slugs": [
coerce_slug(value)
for value in (item.get("merge_slugs") or item.get("merge_ids") or [])
if value
],
"split_into": split_into,
"trait_id": item.get("trait_id") or "",
}
)
return {
"kind": KIND_RESULT,
"format_version": FORMAT_VERSION,
"target": "writing",
"mode": data.get("mode") or "",
"changes": cleaned,
}
def _store_review(profile_id: str, package: dict, channel: str) -> dict:
review_id = str(uuid.uuid4())
with get_db() as conn:
conn.execute(
"""
INSERT INTO writing_profile_reviews
(id, profile_id, target, channel, status, package_json)
VALUES (?, ?, 'writing', ?, 'open', ?)
""",
(review_id, profile_id, channel, json.dumps(package, ensure_ascii=False)),
)
ids = [item.get("id") for item in package.get("evidences") or [] if item.get("id")]
if ids:
placeholders = ",".join("?" * len(ids))
conn.execute(
f"UPDATE writing_profile_evidence SET status = 'bundled' WHERE id IN ({placeholders})",
ids,
)
return {"id": review_id, **review_status(profile_id), "package": package}
def open_paste_review(profile_id: str, *, mode: str | None = None) -> dict:
package = build_package(profile_id, trigger="explicit_review", mode=mode)
stored = _store_review(profile_id, package, "paste")
stored["prompt"] = render_paste_prompt(package)
return stored
def open_api_review(profile_id: str, *, mode: str | None = None) -> dict:
package = build_package(profile_id, trigger="explicit_review", mode=mode)
return _store_review(profile_id, package, "api")
def run_api_review(profile_id: str, review_id: str | None = None, *, mode: str | None = None) -> dict:
from engine import execute_prompt, load_active_prompt
if review_id:
with get_db() as conn:
row = row_to_dict(
conn.execute(
"SELECT * FROM writing_profile_reviews WHERE id = ? AND profile_id = ?",
(review_id, profile_id),
).fetchone()
)
if not row:
raise StoreError("not_found", "Review nicht gefunden", 404)
package = _parse_json(row.get("package_json"), {})
if (row.get("status") or "") not in {"open", "proposed"}:
raise StoreError("review_closed", "Diese Review ist bereits abgeschlossen.")
else:
opened = open_api_review(profile_id, mode=mode)
review_id = opened["id"]
package = opened["package"]
prompt = load_active_prompt("mvp.profile_review")
result = execute_prompt(
prompt,
profile_id,
purpose="profile_review",
data_class="B",
context={
"review_package": json.dumps(package, ensure_ascii=False, indent=2),
"source_text": "\n".join(
[item.get("excerpt") or "" for item in package.get("evidences") or []]
+ [item.get("excerpt") or "" for item in package.get("corpus") or []]
),
},
)
parsed = parse_result(result.get("content") or "")
if not parsed.get("mode"):
parsed["mode"] = package.get("mode") or ""
with get_db() as conn:
conn.execute(
"""
UPDATE writing_profile_reviews
SET result_json = ?, status = 'proposed'
WHERE id = ? AND profile_id = ?
""",
(json.dumps(parsed, ensure_ascii=False), review_id, profile_id),
)
applied = apply_result(profile_id, parsed, review_id=review_id, package=package)
applied["trace"] = result.get("trace")
applied["review_id"] = review_id
return applied
def import_result(profile_id: str, raw, review_id: str | None = None) -> dict:
parsed = parse_result(raw)
package = {}
if not review_id:
with get_db() as conn:
open_row = row_to_dict(
conn.execute(
"""
SELECT id, package_json FROM writing_profile_reviews
WHERE profile_id = ? AND status IN ('open', 'proposed')
ORDER BY created DESC LIMIT 1
""",
(profile_id,),
).fetchone()
)
review_id = (open_row or {}).get("id")
package = _parse_json((open_row or {}).get("package_json"), {})
if not review_id:
stored = _store_review(
profile_id,
build_package(profile_id, trigger="explicit_review"),
"paste",
)
review_id = stored["id"]
package = stored.get("package") or {}
else:
with get_db() as conn:
row = row_to_dict(
conn.execute(
"SELECT package_json FROM writing_profile_reviews WHERE id = ? AND profile_id = ?",
(review_id, profile_id),
).fetchone()
)
package = _parse_json((row or {}).get("package_json"), {})
if not parsed.get("mode"):
parsed["mode"] = package.get("mode") or ""
with get_db() as conn:
conn.execute(
"""
UPDATE writing_profile_reviews
SET result_json = ?, status = 'proposed', channel = COALESCE(channel, 'paste')
WHERE id = ? AND profile_id = ?
""",
(json.dumps(parsed, ensure_ascii=False), review_id, profile_id),
)
return apply_result(profile_id, parsed, review_id=review_id, package=package)
def _refs_from_change(change: dict, evidences: list[dict], corpus: list[dict]) -> list[dict]:
refs = []
by_id = {item.get("id"): item for item in evidences if item.get("id")}
by_source = {item.get("id"): item for item in corpus if item.get("id")}
for eid in change.get("evidence_ids") or []:
item = by_id.get(eid) or by_source.get(eid)
if not item:
continue
excerpt = plain_text(item.get("excerpt") or item.get("body") or "")
if not excerpt:
continue
refs.append(
{
"role": "evidence",
"excerpt": excerpt,
"source_id": item.get("source_id") or item.get("id"),
"occurred_at": item.get("occurred_at"),
}
)
for item in change.get("exemplars") or []:
excerpt = plain_text(item.get("excerpt") or "")
if not excerpt:
continue
refs.append(
{
"role": item.get("role") or "exemplar",
"excerpt": excerpt,
"source_id": item.get("source_id"),
"occurred_at": item.get("occurred_at"),
}
)
return refs
def _apply_change(
profile_id: str,
change: dict,
*,
origin: str,
force: bool,
allow_core: bool,
evidences: list[dict],
corpus: list[dict],
) -> str:
action = change.get("action") or "confirm"
layer = change.get("layer") or "trait"
value = (change.get("proposed_value") or "").strip()
rationale = (change.get("rationale") or "AI Review").strip()
refs = _refs_from_change(change, evidences, corpus)
if layer == "core":
if not allow_core:
return "skipped"
if action in {"confirm", "keep"} and not value:
return "applied"
if not value:
return "skipped"
return _upsert_facet(
profile_id,
CORE_KEY,
value,
evidence=rationale,
origin=origin,
force=force,
)
if layer == "facet" and action not in {"create", "split", "merge"}:
key = normalize_facet_key(change.get("key") or change.get("facet_key") or JOURNAL_FACET)
if not value and action in {"confirm", "keep"}:
return "applied"
if not value:
return "skipped"
return _upsert_facet(profile_id, key, value, evidence=rationale, origin=origin, force=force)
slug = coerce_slug(change.get("slug") or change.get("key") or "trait")
facet_key = normalize_facet_key(change.get("facet_key") or JOURNAL_FACET) or JOURNAL_FACET
if action == "create" and not value:
return "skipped"
if action == "rescope":
upsert_trait(
profile_id,
slug=slug,
facet_key=facet_key,
label=change.get("label") or slug,
statement=value or None,
origin=origin,
force=force,
)
if refs:
replace_trait_refs(profile_id, slug, refs)
return "applied"
if action == "merge":
upsert_trait(
profile_id,
slug=slug,
facet_key=facet_key,
label=change.get("label") or slug,
statement=value or None,
origin=origin,
force=force,
)
if refs:
replace_trait_refs(profile_id, slug, refs)
for other in change.get("merge_slugs") or []:
if other and other != slug:
retire_trait(profile_id, other, "merged")
return "applied"
if action == "split":
created = 0
for part in change.get("split_into") or []:
upsert_trait(
profile_id,
slug=part["slug"],
facet_key=part.get("facet_key") or facet_key,
label=part.get("label") or part["slug"],
statement=part.get("statement") or "",
origin=origin,
force=force,
)
created += 1
if created:
retire_trait(profile_id, slug, "split")
return "applied"
return "skipped"
if action in {"confirm", "precisify", "create", "keep", "update"}:
if action == "confirm" and not value:
if refs:
replace_trait_refs(profile_id, slug, refs)
return "applied"
outcome = upsert_trait(
profile_id,
slug=slug,
facet_key=facet_key,
label=change.get("label") or slug,
statement=value,
origin=origin,
force=force,
)
if refs:
replace_trait_refs(profile_id, slug, refs)
return outcome
return "skipped"
def apply_result(
profile_id: str,
result: dict,
review_id: str | None = None,
package: dict | None = None,
) -> dict:
parsed = parse_result(result)
row = ensure_profile(profile_id)
governance = row.get("governance") or "learning"
if governance not in GOVERNANCE:
governance = "learning"
mode = parsed.get("mode") or (package or {}).get("mode") or _review_mode(profile_id)
if has_confirmed_profile(profile_id):
mode = "incremental"
origin = "initial_build" if mode == "initial_build" else "accepted_suggestion"
corpus = list_corpus(profile_id, limit=INITIAL_BUILD_SOURCES)
evidences = list((package or {}).get("evidences") or []) + _pending_evidence(profile_id, limit=40)
allow_core = _corpus_supports_core(corpus)
snapshot_version(profile_id, cause="before_review")
applied = 0
suggested = 0
skipped = 0
for change in parsed.get("changes") or []:
if governance == "frozen":
skipped += 1
continue
if governance == "advising" and mode != "initial_build":
_queue_suggestion(
profile_id,
change.get("facet_key") or change.get("key") or "",
change.get("proposed_value") or "",
change.get("rationale") or "",
trait_slug=change.get("slug") or "",
action=change.get("action") or "update",
payload=change,
)
suggested += 1
continue
outcome = _apply_change(
profile_id,
change,
origin=origin,
force=mode == "initial_build",
allow_core=allow_core,
evidences=evidences,
corpus=corpus,
)
if outcome == "applied":
applied += 1
else:
skipped += 1
_assemble_brief(profile_id)
if applied or suggested:
snapshot_version(profile_id, cause="accepted_review")
if mode == "initial_build" and applied:
with get_db() as conn:
conn.execute(
"""
UPDATE writing_profiles
SET lifecycle = 'confirmed', updated = datetime('now')
WHERE profile_id = ?
""",
(profile_id,),
)
_assemble_brief(profile_id)
evidence_ids = [
eid
for change in parsed.get("changes") or []
for eid in change.get("evidence_ids") or []
]
with get_db() as conn:
conn.execute(
"UPDATE writing_profiles SET last_reviewed = datetime('now'), review_ready = 0 WHERE profile_id = ?",
(profile_id,),
)
if evidence_ids:
placeholders = ",".join("?" * len(evidence_ids))
conn.execute(
f"UPDATE writing_profile_evidence SET status = 'consumed' WHERE profile_id = ? AND id IN ({placeholders})",
(profile_id, *evidence_ids),
)
else:
conn.execute(
"UPDATE writing_profile_evidence SET status = 'consumed' WHERE profile_id = ? AND status IN ('bundled', 'pending')",
(profile_id,),
)
if review_id:
conn.execute(
"""
UPDATE writing_profile_reviews
SET status = 'applied', resolved = datetime('now')
WHERE id = ? AND profile_id = ?
""",
(review_id, profile_id),
)
_refresh_ready(profile_id)
profile = get_profile(profile_id)
return {
"governance": governance,
"mode": mode,
"applied": applied,
"suggested": suggested,
"skipped": skipped,
"result": parsed,
"profile": profile,
**review_status(profile_id),
}