feat(moments): mount the corpus by proposal - a pass and an open-after-moment signal, both stopping at the operator (milestone 458 step 7, #4925)
CI & Build / Python lint (push) Successful in 3s
CI & Build / Plugin hooks (push) Successful in 19s
CI & Build / integration (push) Failing after 43s
CI & Build / Python tests (push) Failing after 46s
CI & Build / TypeScript typecheck (push) Successful in 55s
CI & Build / Build & push image (push) Skipped
CI & Build / Python lint (push) Successful in 3s
CI & Build / Plugin hooks (push) Successful in 19s
CI & Build / integration (push) Failing after 43s
CI & Build / Python tests (push) Failing after 46s
CI & Build / TypeScript typecheck (push) Successful in 55s
CI & Build / Build & push image (push) Skipped
A rule written before moments existed is mounted on nothing. Step 7 records, per (rule, moment), whether it belongs there and who said so: - rule_moment_judgments (migration 0118, backup v23): suggested / confirmed / rejected, from a pass, the signal, or an edit. Moment "" is "no moment fits". - The pass: rules_to_mount lists unjudged rules; propose_rule_moments records suggestions that mount nothing; rule_moment_proposals and judge_rule_moments put them to the operator. A confirm mounts, a reject is kept so the pair is never proposed again. Same service behind REST and a "Waiting on you" panel in Settings > Moments. - Edits are judgments: set_rule_moments, the one mount write path, confirms what was added and rejects what was removed in the same transaction. - The signal: scribe_moment.sh keeps a per-session acts ledger; when a rule is opened, scribe_record_opened.sh sends the last three minutes of it to /api/plugin/rule-opened. The acts resolve through the install's mappings; work.run and work.change are not evidence. Counted per distinct session with lesson_rules' evidence model, and once due the open returns one line asking the reader to offer the mount. - scribe_session_end.sh removes the session's scribe-moment files. Plugin 2026.10.05.2003. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -21,6 +21,7 @@ from scribe.models.rulebook import (
|
||||
RuleRelation, rule_moments as rule_moments_t, rule_systems as rule_systems_t,
|
||||
)
|
||||
from scribe.models.lesson_rule_link import LessonNoRule, LessonRuleLink
|
||||
from scribe.models.rule_moment_judgment import RuleMomentJudgment
|
||||
from scribe.models.code_shape import CodeShape, CodeShapeEvent, CodeShapeUse
|
||||
from scribe.models.project import Project
|
||||
from scribe.models.repo_binding import RepoBinding
|
||||
@@ -104,8 +105,12 @@ logger = logging.getLogger(__name__)
|
||||
# arrives at. A mount is a judgment about when a rule applies that nothing in
|
||||
# the rule's text records, and losing it would put back exactly the misses the
|
||||
# mounts were made to fix.
|
||||
# v23 (2026-10) added rule_moment_judgments (milestone 458 step 7): the
|
||||
# proposals waiting on a mount, and the rejections and "no moment fits"
|
||||
# answers that stop a pass or the open-after-moment signal proposing the same
|
||||
# pair again. Losing them puts every answered question back on the list.
|
||||
# Bump when the serialized schema changes.
|
||||
BACKUP_VERSION = 22
|
||||
BACKUP_VERSION = 23
|
||||
|
||||
# Every table this backup carries, by its REAL name. Paired with _NOT_INCLUDED
|
||||
# below, these two lists must together account for the entire schema — which is
|
||||
@@ -154,6 +159,8 @@ _BACKED_UP = [
|
||||
"moment_mappings",
|
||||
# v22 (2026-10): the moments each rule arrives at (milestone 458).
|
||||
"rule_moments",
|
||||
# v23 (2026-10): proposals and judgments about those mounts (step 7).
|
||||
"rule_moment_judgments",
|
||||
]
|
||||
|
||||
# Tables intentionally NOT in the backup, surfaced in the payload so the gap is
|
||||
@@ -253,6 +260,8 @@ _COLUMN_EXCLUSIONS: dict[str, set[str]] = {
|
||||
"rule_relations": {"id", "created_at"},
|
||||
# The pair is the row; everything else is the judgment and its evidence.
|
||||
"lesson_rule_links": {"id"},
|
||||
# The same shape: the (rule, moment) pair is the row.
|
||||
"rule_moment_judgments": {"id"},
|
||||
# Keyed by the lesson itself; nothing to exclude, every column travels.
|
||||
"lesson_no_rule": set(),
|
||||
"note_usage_events": {"id"},
|
||||
@@ -346,6 +355,7 @@ _IMPORT_COLUMN_EXCLUSIONS: dict[str, set[str]] = {
|
||||
"note_supersessions": {"id", "created_at"},
|
||||
"rule_relations": {"id", "created_at"},
|
||||
"lesson_rule_links": {"id"},
|
||||
"rule_moment_judgments": {"id"},
|
||||
"lesson_no_rule": set(),
|
||||
"note_usage_events": {"id"},
|
||||
"rule_usage_events": {"id"},
|
||||
@@ -755,6 +765,20 @@ def _rule_moment_rows(rows) -> list[dict]:
|
||||
return [{"rule_id": rule_id, "moment": moment} for rule_id, moment in rows]
|
||||
|
||||
|
||||
def _rule_moment_judgment_rows(rows) -> list[dict]:
|
||||
"""Proposals and judgments about a rule's moments (v23). The rule id is a
|
||||
SOURCE id, remapped at restore; the moment is a catalog name."""
|
||||
return [
|
||||
{
|
||||
"rule_id": r.rule_id, "moment": r.moment, "state": r.state,
|
||||
"source": r.source, "note": r.note, "evidence": r.evidence,
|
||||
"judged_at": r.judged_at.isoformat() if r.judged_at else None,
|
||||
"created_at": r.created_at.isoformat() if r.created_at else None,
|
||||
}
|
||||
for r in rows
|
||||
]
|
||||
|
||||
|
||||
def _rule_system_rows(rows) -> list[dict]:
|
||||
"""A rule's area tags, carried by canonical SLUG for the same reason the
|
||||
Systems are: the catalog is global and its ids are per-install."""
|
||||
@@ -850,6 +874,9 @@ async def export_full_backup() -> dict:
|
||||
rule_moment_rows = (await session.execute(
|
||||
select(rule_moments_t.c.rule_id, rule_moments_t.c.moment)
|
||||
)).all()
|
||||
rule_moment_judgments = (await session.execute(
|
||||
select(RuleMomentJudgment)
|
||||
)).scalars().all()
|
||||
rule_relations = (await session.execute(select(RuleRelation))).scalars().all()
|
||||
lesson_rule_links = (await session.execute(select(LessonRuleLink))).scalars().all()
|
||||
lesson_no_rule = (await session.execute(select(LessonNoRule))).scalars().all()
|
||||
@@ -917,6 +944,7 @@ async def export_full_backup() -> dict:
|
||||
"canonical_systems": _canonical_system_rows(canonical_systems),
|
||||
"rule_systems": _rule_system_rows(rule_system_rows),
|
||||
"rule_moments": _rule_moment_rows(rule_moment_rows),
|
||||
"rule_moment_judgments": _rule_moment_judgment_rows(rule_moment_judgments),
|
||||
"rule_relations": _rule_relation_rows(rule_relations),
|
||||
"lesson_rule_links": _lesson_rule_link_rows(lesson_rule_links),
|
||||
"lesson_no_rule": _lesson_no_rule_rows(lesson_no_rule),
|
||||
@@ -1055,6 +1083,9 @@ async def export_user_backup(user_id: int) -> dict:
|
||||
select(rule_moments_t.c.rule_id, rule_moments_t.c.moment)
|
||||
.where(rule_moments_t.c.rule_id.in_(_rule_ids))
|
||||
)).all() if _rule_ids else []
|
||||
rule_moment_judgments = (await session.execute(
|
||||
select(RuleMomentJudgment).where(RuleMomentJudgment.rule_id.in_(_rule_ids))
|
||||
)).scalars().all() if _rule_ids else []
|
||||
# Scoped through the RULE, not the version's user_id. That column is
|
||||
# the ACTOR (milestone 323), so filtering on it would carry the
|
||||
# versions this user wrote on someone ELSE's rule and drop the ones
|
||||
@@ -1137,6 +1168,7 @@ async def export_user_backup(user_id: int) -> dict:
|
||||
"canonical_systems": _canonical_system_rows(canonical_systems),
|
||||
"rule_systems": _rule_system_rows(rule_system_rows),
|
||||
"rule_moments": _rule_moment_rows(rule_moment_rows),
|
||||
"rule_moment_judgments": _rule_moment_judgment_rows(rule_moment_judgments),
|
||||
"rule_relations": _rule_relation_rows(rule_relations),
|
||||
"lesson_rule_links": _lesson_rule_link_rows(lesson_rule_links),
|
||||
"lesson_no_rule": _lesson_no_rule_rows(lesson_no_rule),
|
||||
@@ -1535,6 +1567,25 @@ def _build_lesson_rule_link(row: dict, maps: _Maps) -> LessonRuleLink | None:
|
||||
)
|
||||
|
||||
|
||||
def _build_rule_moment_judgment(row: dict, maps: _Maps) -> RuleMomentJudgment | None:
|
||||
"""Skipped when its rule did not restore — the judgment is about that
|
||||
rule, and a remap miss would hang it on whatever took the number."""
|
||||
rule = maps.rules.get(row.get("rule_id", 0))
|
||||
if rule is None or row.get("moment") is None:
|
||||
return None
|
||||
return RuleMomentJudgment(
|
||||
rule_id=rule,
|
||||
moment=row["moment"],
|
||||
state=row.get("state") or "suggested",
|
||||
source=row.get("source") or "pass",
|
||||
note=row.get("note") or None,
|
||||
evidence=row.get("evidence"),
|
||||
# Absent stays absent: a suggestion was never judged.
|
||||
judged_at=_dt_or_none(row.get("judged_at")),
|
||||
created_at=_dt(row.get("created_at")),
|
||||
)
|
||||
|
||||
|
||||
def _build_rule_version(row: dict, maps: _Maps) -> RuleVersion | None:
|
||||
rid = maps.rules.get(row.get("rule_id", 0))
|
||||
if rid is None:
|
||||
@@ -1949,7 +2000,7 @@ async def _restore_v2(data: dict) -> dict:
|
||||
"repo_bindings": 0,
|
||||
"note_supersessions": 0, "code_shapes": 0, "code_shape_events": 0,
|
||||
"code_shape_uses": 0, "canonical_systems": 0,
|
||||
"rule_systems": 0, "rule_moments": 0,
|
||||
"rule_systems": 0, "rule_moments": 0, "rule_moment_judgments": 0,
|
||||
"rule_relations": 0, "rule_versions": 0,
|
||||
"retrieval_tuning_events": 0, "lesson_rule_links": 0,
|
||||
"lesson_no_rule": 0, "moment_mappings": 0,
|
||||
@@ -2160,6 +2211,15 @@ async def _restore_v2(data: dict) -> dict:
|
||||
))
|
||||
stats["rule_moments"] += 1
|
||||
|
||||
# Proposals and judgments about those mounts (v23); archives before
|
||||
# v23 carry none, and every rule restores unjudged.
|
||||
for rj in data.get("rule_moment_judgments", []):
|
||||
judgment = _build_rule_moment_judgment(rj, maps)
|
||||
if judgment is None:
|
||||
continue
|
||||
session.add(judgment)
|
||||
stats["rule_moment_judgments"] += 1
|
||||
|
||||
for rr in data.get("rule_relations", []):
|
||||
relation = _build_rule_relation(rr, maps)
|
||||
if relation is None:
|
||||
|
||||
@@ -485,16 +485,18 @@ async def attach_rule_lessons(user_id: int, data: dict, rule_id: int) -> None:
|
||||
def situation_key(arm: str, text: str) -> str:
|
||||
"""A fingerprint of one situation, so a repeat counts once.
|
||||
|
||||
`arm` is part of the key — "p" for a prompt, "w" for a file being written
|
||||
— because the two name different kinds of situation and must not collide.
|
||||
`arm` is part of the key — "p" for a prompt, "w" for a file being written,
|
||||
"s" for a session (the rule↔moment signal, milestone 458 step 7) —
|
||||
because each names a different kind of situation and must not collide.
|
||||
On the prompt arm the text is the prompt: lowercased, split into word
|
||||
tokens of _MIN_TOKEN or more, de-duplicated and sorted, so the same ask
|
||||
re-sent with different spacing, punctuation, case or word order is one
|
||||
situation. On the write arm the text is the PATH, not the code: every
|
||||
edit to one file is one situation, however much the code differs between
|
||||
them. Empty when there is nothing to key on.
|
||||
them. The session arm keys on the session id as given. Empty when there
|
||||
is nothing to key on.
|
||||
"""
|
||||
if arm == "w":
|
||||
if arm in ("w", "s"):
|
||||
basis = (text or "").strip()
|
||||
else:
|
||||
tokens = sorted({t for t in re.findall(r"[a-z0-9]+", (text or "").lower())
|
||||
|
||||
@@ -0,0 +1,527 @@
|
||||
"""Which rules belong on which moments — proposals and judgments (milestone 458 step 7).
|
||||
|
||||
A rule written before moments existed is mounted on nothing, and reaches a
|
||||
session only when its words resemble the work. Mounting the corpus is a
|
||||
judgment about each rule — WHEN does it apply? — and nothing in a rule's row
|
||||
says whether that judgment was ever made. This service records it, from two
|
||||
sources that both stop at a proposal:
|
||||
|
||||
- THE PASS. An agent reads the rules nobody has judged (`unjudged_rules`),
|
||||
proposes moments for each, or says none fits (`propose`). A proposal mounts
|
||||
nothing; a person confirms it (`judge`).
|
||||
- THE SIGNAL. A rule a session OPENS shortly after a moment fired is evidence
|
||||
it belongs on that moment. Recorded per distinct session with the evidence
|
||||
model lesson↔rule soft links use (lesson_rules: three situations, a
|
||||
cooldown), and once it crosses the bar the next open asks the reader to
|
||||
offer the mount (`co_occurred`).
|
||||
|
||||
Neither source mounts anything by itself. A suggestion that delivered its rule
|
||||
would manufacture the opens it counts, and the operator's corpus is theirs to
|
||||
mount.
|
||||
|
||||
EDITS ARE JUDGMENTS TOO. A person changing a rule's moments directly says
|
||||
which moments it belongs on, so `record_mount_change` (called from the one
|
||||
write path, `rulebooks.set_rule_moments`) confirms what was added and rejects
|
||||
what was removed — a moment taken off a rule must not be proposed straight
|
||||
back by the signal.
|
||||
|
||||
ACL (rule 78): rules are owner-scoped, and every read and write here goes
|
||||
through `rulebooks._owned_rules_clause` / `_fetch_owned_rule`.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from sqlalchemy import exists, func, not_, select
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
|
||||
from scribe.models import async_session
|
||||
from scribe.models.rule_moment_judgment import (
|
||||
CONFIRMED, NO_MOMENT, REJECTED, SUGGESTED, RuleMomentJudgment,
|
||||
)
|
||||
from scribe.models.rulebook import Rule, rule_moments as rule_moments_t
|
||||
from scribe.services import lesson_rules
|
||||
from scribe.services import moments as moments_svc
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# What a caller says to judge one pair — acts, not the states they produce.
|
||||
VERDICTS = {"confirm": CONFIRMED, "reject": REJECTED}
|
||||
|
||||
# Recorded on the rows an edit wrote, so a reader can tell a direct change from
|
||||
# an accepted proposal.
|
||||
_MOUNTED_NOTE = "mounted by an edit"
|
||||
_UNMOUNTED_NOTE = "unmounted by an edit"
|
||||
# On a "no moment fits" answer a later mount overturned.
|
||||
_OVERTURNED_NOTE = "a moment was mounted after all"
|
||||
|
||||
# How many unjudged rules one pass page hands over. Enough to batch the
|
||||
# reading, few enough that every rule is actually read.
|
||||
PASS_PAGE = 25
|
||||
# How long after a moment an open still counts as following it. Long enough
|
||||
# for a session to see a line, decide it matters and open the rule; short
|
||||
# enough that the moment is still what the work is doing.
|
||||
SIGNAL_WINDOW_SECONDS = 180
|
||||
# Moments too ubiquitous to say anything about the rule opened after them:
|
||||
# nearly every open follows a command or a change, so counting them would
|
||||
# propose every rule onto both.
|
||||
SIGNAL_SKIP = frozenset({"work.run", "work.change"})
|
||||
|
||||
|
||||
def _clean_moment(name) -> str:
|
||||
"""A catalog moment, normalised; "" stays "" (no moment fits)."""
|
||||
raw = (name or "").strip() if isinstance(name, str) else ""
|
||||
if raw == NO_MOMENT:
|
||||
return NO_MOMENT
|
||||
return moments_svc.require_moment(raw)
|
||||
|
||||
|
||||
def _int(v) -> int | None:
|
||||
try:
|
||||
i = int(v)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
return i if i > 0 else None
|
||||
|
||||
|
||||
async def _owned_rules(session, user_id: int, rule_ids) -> dict[int, Rule]:
|
||||
from scribe.services.rulebooks import _owned_rules_clause
|
||||
|
||||
ids = [i for i in (_int(r) for r in rule_ids or []) if i]
|
||||
if not ids:
|
||||
return {}
|
||||
rows = (await session.execute(
|
||||
select(Rule).where(Rule.id.in_(ids)).where(_owned_rules_clause(user_id))
|
||||
)).scalars().all()
|
||||
return {r.id: r for r in rows}
|
||||
|
||||
|
||||
async def _mounts(session, rule_ids) -> dict[int, set[str]]:
|
||||
if not rule_ids:
|
||||
return {}
|
||||
rows = (await session.execute(
|
||||
select(rule_moments_t.c.rule_id, rule_moments_t.c.moment)
|
||||
.where(rule_moments_t.c.rule_id.in_(list(rule_ids)))
|
||||
)).all()
|
||||
out: dict[int, set[str]] = {}
|
||||
for rid, moment in rows:
|
||||
out.setdefault(rid, set()).add(moment)
|
||||
return out
|
||||
|
||||
|
||||
async def _rows(session, rule_ids) -> dict[tuple[int, str], RuleMomentJudgment]:
|
||||
if not rule_ids:
|
||||
return {}
|
||||
rows = (await session.execute(
|
||||
select(RuleMomentJudgment).where(RuleMomentJudgment.rule_id.in_(list(rule_ids)))
|
||||
)).scalars().all()
|
||||
return {(r.rule_id, r.moment): r for r in rows}
|
||||
|
||||
|
||||
def _brief(rule: Rule) -> dict:
|
||||
return {
|
||||
"id": rule.id,
|
||||
"title": rule.title,
|
||||
"kind": rule.kind or "rule",
|
||||
"statement": rule.statement,
|
||||
"when_to_apply": rule.when_to_apply or "",
|
||||
"home": "project" if rule.project_id else "global",
|
||||
}
|
||||
|
||||
|
||||
# ── The pass ─────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _unjudged_clause():
|
||||
"""No mount and no judgment row of any kind — a rule nobody has looked at."""
|
||||
mounted = exists().where(rule_moments_t.c.rule_id == Rule.id)
|
||||
judged = exists().where(RuleMomentJudgment.rule_id == Rule.id)
|
||||
return not_(mounted), not_(judged)
|
||||
|
||||
|
||||
async def unjudged_rules(user_id: int, limit: int = PASS_PAGE, offset: int = 0) -> dict:
|
||||
"""The caller's rules with no moment and no judgment, a page at a time.
|
||||
|
||||
A rule leaves this list the moment any answer is recorded for it — a
|
||||
proposal, a mount, or "no moment fits" — so a pass run in pieces, or by
|
||||
two sessions, does not read the same rule twice.
|
||||
"""
|
||||
from scribe.services.rulebooks import _owned_rules_clause
|
||||
|
||||
limit = max(1, min(int(limit or PASS_PAGE), 100))
|
||||
offset = max(0, int(offset or 0))
|
||||
clauses = (_owned_rules_clause(user_id), *_unjudged_clause())
|
||||
async with async_session() as session:
|
||||
total = (await session.execute(
|
||||
select(func.count()).select_from(Rule).where(*clauses)
|
||||
)).scalar_one()
|
||||
rows = (await session.execute(
|
||||
select(Rule).where(*clauses).order_by(Rule.id).limit(limit).offset(offset)
|
||||
)).scalars().all()
|
||||
return {
|
||||
"rules": [_brief(r) for r in rows],
|
||||
"total_unjudged": int(total),
|
||||
"offset": offset,
|
||||
"catalog": moments_svc.catalog(),
|
||||
}
|
||||
|
||||
|
||||
async def propose(user_id: int, proposals: list[dict]) -> dict:
|
||||
"""Record a pass's proposals: moments for a rule, or that none fits.
|
||||
|
||||
Each item is `{rule_id, moments: [...], why}` or `{rule_id, none: "why"}`.
|
||||
A moment proposal lands `suggested` and mounts nothing. "No moment fits"
|
||||
lands as the agent's own answer (confirmed, moment ""): it changes nothing
|
||||
that surfaces, and it is what keeps the rule off the next pass — the
|
||||
signal can still propose a moment for it later, on evidence.
|
||||
|
||||
A pair already mounted or already judged is skipped and said so; a pair
|
||||
already suggested has its reason refreshed. A bad item is refused alone,
|
||||
with why, and the rest are written.
|
||||
"""
|
||||
now = datetime.now(timezone.utc)
|
||||
proposed, no_moment = 0, 0
|
||||
skipped: list[dict] = []
|
||||
refused: list[dict] = []
|
||||
async with async_session() as session:
|
||||
owned = await _owned_rules(session, user_id, [p.get("rule_id") for p in proposals or []
|
||||
if isinstance(p, dict)])
|
||||
mounts = await _mounts(session, list(owned))
|
||||
existing = await _rows(session, list(owned))
|
||||
for item in proposals or []:
|
||||
if not isinstance(item, dict):
|
||||
refused.append({"item": item, "error": "each proposal is an object"})
|
||||
continue
|
||||
rid = _int(item.get("rule_id"))
|
||||
if rid is None or rid not in owned:
|
||||
refused.append({"rule_id": item.get("rule_id"),
|
||||
"error": "not a rule you own"})
|
||||
continue
|
||||
none_why = (item.get("none") or "").strip() if isinstance(item.get("none"), str) else ""
|
||||
names = item.get("moments")
|
||||
if none_why and names:
|
||||
refused.append({"rule_id": rid, "error": (
|
||||
"moments and none are the two answers to \"when does this rule "
|
||||
"apply?\" — give the one that holds")})
|
||||
continue
|
||||
if none_why:
|
||||
if mounts.get(rid):
|
||||
skipped.append({"rule_id": rid, "moment": NO_MOMENT,
|
||||
"why": "already mounted on " + ", ".join(sorted(mounts[rid]))})
|
||||
continue
|
||||
row = existing.get((rid, NO_MOMENT))
|
||||
if row is None:
|
||||
row = RuleMomentJudgment(rule_id=rid, moment=NO_MOMENT, source="pass")
|
||||
session.add(row)
|
||||
existing[(rid, NO_MOMENT)] = row
|
||||
row.state, row.note, row.judged_at = CONFIRMED, none_why, now
|
||||
no_moment += 1
|
||||
continue
|
||||
try:
|
||||
wanted = moments_svc.require_moments(names if names is not None else []) or []
|
||||
except ValueError as exc:
|
||||
refused.append({"rule_id": rid, "error": str(exc)})
|
||||
continue
|
||||
if not wanted:
|
||||
refused.append({"rule_id": rid, "error": (
|
||||
"name at least one moment, or say none fits with none=\"why\"")})
|
||||
continue
|
||||
why = (item.get("why") or "").strip() if isinstance(item.get("why"), str) else ""
|
||||
for moment in wanted:
|
||||
if moment in mounts.get(rid, set()):
|
||||
skipped.append({"rule_id": rid, "moment": moment, "why": "already mounted"})
|
||||
continue
|
||||
row = existing.get((rid, moment))
|
||||
if row is not None and row.state != SUGGESTED:
|
||||
skipped.append({"rule_id": rid, "moment": moment,
|
||||
"why": f"already {row.state}" + (f": {row.note}" if row.note else "")})
|
||||
continue
|
||||
if row is None:
|
||||
row = RuleMomentJudgment(rule_id=rid, moment=moment, state=SUGGESTED, source="pass")
|
||||
session.add(row)
|
||||
existing[(rid, moment)] = row
|
||||
row.note = why or row.note
|
||||
proposed += 1
|
||||
await session.commit()
|
||||
return {"proposed": proposed, "no_moment": no_moment,
|
||||
"skipped": skipped, "refused": refused}
|
||||
|
||||
|
||||
async def pending(user_id: int, rule_id: int | None = None, limit: int = 200) -> dict:
|
||||
"""The proposals waiting on a judgment, grouped by rule.
|
||||
|
||||
Each rule carries what it is mounted on now, so a proposal is read beside
|
||||
the mounts it would add to.
|
||||
"""
|
||||
from scribe.services.rulebooks import _owned_rules_clause
|
||||
|
||||
limit = max(1, min(int(limit or 200), 500))
|
||||
async with async_session() as session:
|
||||
stmt = (
|
||||
select(RuleMomentJudgment, Rule)
|
||||
.join(Rule, Rule.id == RuleMomentJudgment.rule_id)
|
||||
.where(RuleMomentJudgment.state == SUGGESTED, _owned_rules_clause(user_id))
|
||||
)
|
||||
if rule_id:
|
||||
stmt = stmt.where(Rule.id == int(rule_id))
|
||||
rows = (await session.execute(
|
||||
stmt.order_by(Rule.id, RuleMomentJudgment.moment).limit(limit)
|
||||
)).all()
|
||||
mounts = await _mounts(session, {r.id for _, r in rows})
|
||||
order = {name: i for i, name in enumerate(moments_svc.MOMENTS)}
|
||||
grouped: dict[int, dict] = {}
|
||||
for judgment, rule in rows:
|
||||
entry = grouped.setdefault(rule.id, {
|
||||
**_brief(rule),
|
||||
"mounted": sorted(mounts.get(rule.id, set()),
|
||||
key=lambda n: (order.get(n, len(order)), n)),
|
||||
"proposals": [],
|
||||
})
|
||||
entry["proposals"].append({
|
||||
"moment": judgment.moment,
|
||||
"source": judgment.source,
|
||||
"why": judgment.note or "",
|
||||
"evidence": lesson_rules.evidence_summary(judgment.evidence),
|
||||
"created_at": judgment.created_at.isoformat() if judgment.created_at else None,
|
||||
})
|
||||
rules = list(grouped.values())
|
||||
return {"rules": rules, "total": sum(len(r["proposals"]) for r in rules)}
|
||||
|
||||
|
||||
async def judge(user_id: int, judgments: list[dict]) -> dict:
|
||||
"""Confirm or reject proposals — or any (rule, moment) pair, proposed or not.
|
||||
|
||||
Confirm MOUNTS the rule on the moment (alongside what it already has);
|
||||
reject records that it does not belong there and, if it was mounted,
|
||||
unmounts it. Both go through `rulebooks.set_rule_moments`, the one write
|
||||
path, which records the judgment with this item's `note`. A moment of ""
|
||||
judges the "no moment fits" answer itself. A bad item is refused alone.
|
||||
"""
|
||||
from scribe.services.rulebooks import list_rule_moments, set_rule_moments
|
||||
|
||||
done: list[dict] = []
|
||||
refused: list[dict] = []
|
||||
for item in judgments or []:
|
||||
if not isinstance(item, dict):
|
||||
refused.append({"item": item, "error": "each judgment is an object"})
|
||||
continue
|
||||
rid = _int(item.get("rule_id"))
|
||||
state = VERDICTS.get(str(item.get("verdict") or "").strip().lower())
|
||||
if rid is None or state is None:
|
||||
refused.append({"rule_id": item.get("rule_id"), "moment": item.get("moment"),
|
||||
"error": f"needs rule_id and a verdict, one of {sorted(VERDICTS)}"})
|
||||
continue
|
||||
try:
|
||||
moment = _clean_moment(item.get("moment"))
|
||||
except ValueError as exc:
|
||||
refused.append({"rule_id": rid, "error": str(exc)})
|
||||
continue
|
||||
note = (item.get("note") or "").strip() if isinstance(item.get("note"), str) else ""
|
||||
if moment == NO_MOMENT:
|
||||
ok = await _judge_no_moment(user_id, rid, state, note)
|
||||
else:
|
||||
current = (await list_rule_moments([rid])).get(rid, [])
|
||||
after = (current + [moment] if state == CONFIRMED and moment not in current
|
||||
else [m for m in current if m != moment] if state == REJECTED
|
||||
else current)
|
||||
ok = await set_rule_moments(rid, user_id, after, note=note, judged=[moment],
|
||||
verdict=state) is not None
|
||||
if not ok:
|
||||
refused.append({"rule_id": rid, "moment": moment, "error": "not a rule you own"})
|
||||
continue
|
||||
done.append({"rule_id": rid, "moment": moment, "state": state})
|
||||
return {"judged": done, "refused": refused}
|
||||
|
||||
|
||||
async def _judge_no_moment(user_id: int, rule_id: int, state: str, note: str) -> bool:
|
||||
now = datetime.now(timezone.utc)
|
||||
async with async_session() as session:
|
||||
if not await _owned_rules(session, user_id, [rule_id]):
|
||||
return False
|
||||
row = (await _rows(session, [rule_id])).get((rule_id, NO_MOMENT))
|
||||
if row is None:
|
||||
row = RuleMomentJudgment(rule_id=rule_id, moment=NO_MOMENT, source="edit")
|
||||
session.add(row)
|
||||
row.state, row.judged_at = state, now
|
||||
row.note = note or row.note
|
||||
await session.commit()
|
||||
return True
|
||||
|
||||
|
||||
async def record_mount_change(
|
||||
session, rule_id: int, before, after, *, note: str = "",
|
||||
judged=(), verdict: str | None = None,
|
||||
) -> None:
|
||||
"""Record the judgments a change of mounts makes, in the caller's session.
|
||||
|
||||
Added moments are confirmed, removed ones rejected — a mount taken off
|
||||
must not be proposed straight back. `judged` names moments a `judge`
|
||||
call decided explicitly, so a rejection of a moment that was never
|
||||
mounted is recorded too. A suggestion that was accepted keeps its source,
|
||||
so the record still says where the idea came from. Any mount overturns a
|
||||
"no moment fits" answer.
|
||||
"""
|
||||
before, after = set(before or ()), set(after or ())
|
||||
added = after - before
|
||||
removed = before - after
|
||||
explicit = {m for m in judged or () if m} - added - removed
|
||||
if not (added or removed or explicit):
|
||||
return
|
||||
now = datetime.now(timezone.utc)
|
||||
rows = await _rows(session, [rule_id])
|
||||
|
||||
def put(moment: str, state: str, default_note: str) -> None:
|
||||
row = rows.get((rule_id, moment))
|
||||
if row is None:
|
||||
row = RuleMomentJudgment(rule_id=rule_id, moment=moment, source="edit")
|
||||
session.add(row)
|
||||
rows[(rule_id, moment)] = row
|
||||
row.state, row.judged_at = state, now
|
||||
row.note = note or default_note or row.note
|
||||
|
||||
for moment in added:
|
||||
put(moment, CONFIRMED, _MOUNTED_NOTE)
|
||||
for moment in removed:
|
||||
put(moment, REJECTED, _UNMOUNTED_NOTE)
|
||||
for moment in explicit:
|
||||
if verdict in (CONFIRMED, REJECTED):
|
||||
put(moment, verdict, "")
|
||||
if after:
|
||||
none_row = rows.get((rule_id, NO_MOMENT))
|
||||
if none_row is not None and none_row.state == CONFIRMED:
|
||||
none_row.state, none_row.judged_at = REJECTED, now
|
||||
none_row.note = _OVERTURNED_NOTE
|
||||
|
||||
|
||||
# ── The signal ───────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _proposal_line(rule: Rule, moment: str, evidence: dict) -> str:
|
||||
counts = lesson_rules.evidence_summary(evidence)
|
||||
kind = rule.kind or "rule"
|
||||
return (
|
||||
f"> {kind.capitalize()} #{rule.id} \"{rule.title}\" has been opened just after "
|
||||
f"`{moment}` in {counts['situations']} distinct sessions — it may belong on "
|
||||
f"that moment, where it would arrive without being searched for. Offer the "
|
||||
f"operator the mount in one line; on their yes, `judge_rule_moments("
|
||||
f"[{{\"rule_id\": {rule.id}, \"moment\": \"{moment}\", \"verdict\": \"confirm\", "
|
||||
f"\"note\": \"why\"}}])` mounts it; on a no, `\"reject\"` with their why stops "
|
||||
f"the asking."
|
||||
)
|
||||
|
||||
|
||||
def signal_moments(moments) -> list[str]:
|
||||
"""The moments an open may be counted against: catalog or procedure
|
||||
moments, minus the ubiquitous ones, de-duplicated in order."""
|
||||
out: list[str] = []
|
||||
for m in moments or []:
|
||||
name = (m or "").strip().lower() if isinstance(m, str) else ""
|
||||
if name and name not in SIGNAL_SKIP and moments_svc.is_moment(name) and name not in out:
|
||||
out.append(name)
|
||||
return out
|
||||
|
||||
|
||||
async def co_occurred(
|
||||
user_id: int, rule_id: int, moments, *, situation: str,
|
||||
project_id: int | None = None,
|
||||
) -> str:
|
||||
"""Record that a rule was opened shortly after these moments fired, and
|
||||
return a proposal line once a (rule, moment) pair has earned one.
|
||||
|
||||
`situation` is what makes two opens distinct — the session, so a rule
|
||||
opened after every push in one long session counts once. A pair already
|
||||
mounted or judged gathers nothing. One proposal per open at most. Fails
|
||||
open: this rides a hook, and bookkeeping must never break the read.
|
||||
"""
|
||||
try:
|
||||
return await _co_occurred(user_id, rule_id, moments, situation, project_id)
|
||||
except Exception:
|
||||
logger.warning("rule/moment co-occurrence could not be recorded", exc_info=True)
|
||||
return ""
|
||||
|
||||
|
||||
async def _co_occurred(user_id, rule_id, moments, situation, project_id) -> str:
|
||||
wanted = signal_moments(moments)
|
||||
key = lesson_rules.situation_key("s", situation)
|
||||
rid = _int(rule_id)
|
||||
if not wanted or not key or rid is None:
|
||||
return ""
|
||||
now = datetime.now(timezone.utc)
|
||||
proposal = None
|
||||
async with async_session() as session:
|
||||
owned = await _owned_rules(session, user_id, [rid])
|
||||
rule = owned.get(rid)
|
||||
if rule is None:
|
||||
return ""
|
||||
mounted = (await _mounts(session, [rid])).get(rid, set())
|
||||
existing = await _rows(session, [rid])
|
||||
for moment in wanted:
|
||||
if moment in mounted:
|
||||
continue
|
||||
row = existing.get((rid, moment))
|
||||
if row is not None and row.state != SUGGESTED:
|
||||
continue
|
||||
if row is None:
|
||||
row = RuleMomentJudgment(rule_id=rid, moment=moment, state=SUGGESTED, source="signal")
|
||||
session.add(row)
|
||||
ev = lesson_rules.add_evidence(row.evidence, key, project_id, now)
|
||||
if proposal is None and lesson_rules.proposal_due(ev, now):
|
||||
ev["proposed_at"] = now.isoformat()
|
||||
ev["proposed_count"] = int(ev.get("proposed_count") or 0) + 1
|
||||
proposal = (moment, ev)
|
||||
row.evidence = ev
|
||||
try:
|
||||
await session.commit()
|
||||
except IntegrityError:
|
||||
# Two opens created the same pair at once; the other stands.
|
||||
await session.rollback()
|
||||
return ""
|
||||
if proposal is None:
|
||||
return ""
|
||||
return _proposal_line(rule, *proposal)
|
||||
|
||||
|
||||
# The most acts one open is checked against. The hook sends only the window's
|
||||
# acts, so this bounds a misbehaving client rather than an ordinary session.
|
||||
_ACTS_CAP = 30
|
||||
|
||||
|
||||
async def opened_after(
|
||||
user_id: int, rule_id: int, acts, *, session_id: str,
|
||||
project_id: int | None = None,
|
||||
) -> dict:
|
||||
"""A rule was opened; `acts` are the tool calls that came just before it.
|
||||
|
||||
Each act is a PreToolUse event as the plugin recorded it (`tool_name`,
|
||||
`tool_input`), resolved to its moments by the same mappings delivery
|
||||
uses — so a moment counted here is one that fired, on this install, for
|
||||
that call. Returns `moments` (what the window reached, after the skip
|
||||
list) and `context` (a proposal line, or "").
|
||||
"""
|
||||
from scribe.services.moment_actions import moments_for
|
||||
|
||||
reached: list[str] = []
|
||||
for act in list(acts or [])[:_ACTS_CAP]:
|
||||
if not isinstance(act, dict):
|
||||
continue
|
||||
tool = str(act.get("tool_name") or "").strip()
|
||||
if not tool:
|
||||
continue
|
||||
tool_input = act.get("tool_input")
|
||||
try:
|
||||
hits = await moments_for(user_id, tool, tool_input if isinstance(tool_input, dict) else {})
|
||||
except Exception:
|
||||
logger.debug("acts before an open could not be resolved", exc_info=True)
|
||||
continue
|
||||
for hit in hits:
|
||||
name = hit.get("moment") if isinstance(hit, dict) else None
|
||||
if name and name not in reached:
|
||||
reached.append(name)
|
||||
moments = signal_moments(reached)
|
||||
context = ""
|
||||
if moments and (session_id or "").strip():
|
||||
context = await co_occurred(
|
||||
user_id, rule_id, moments, situation=session_id.strip(), project_id=project_id,
|
||||
)
|
||||
return {"moments": moments, "context": context}
|
||||
@@ -1111,7 +1111,8 @@ async def set_rule_systems(
|
||||
|
||||
|
||||
async def set_rule_moments(
|
||||
rule_id: int, user_id: int, moments: list[str],
|
||||
rule_id: int, user_id: int, moments: list[str], *,
|
||||
note: str = "", judged=(), verdict: str | None = None,
|
||||
) -> list[str] | None:
|
||||
"""Replace which MOMENTS a rule arrives at (milestone 458). None if not owned.
|
||||
|
||||
@@ -1119,15 +1120,24 @@ async def set_rule_moments(
|
||||
Every name is checked against the catalog first and the whole write is
|
||||
refused on the first unknown one — a typo stored as a mount would read
|
||||
back as attached and never fire, the silent miss moments exist to end.
|
||||
|
||||
THE ONE WRITE PATH for mounts, so it is also where a change is recorded as
|
||||
a judgment (step 7): what was added is confirmed, what was removed is
|
||||
rejected, in the same transaction as the mounts. `note`, `judged` and
|
||||
`verdict` come from `rule_moment_judgments.judge`, which says why.
|
||||
"""
|
||||
from scribe.models.rulebook import rule_moments as rule_moments_t
|
||||
from scribe.services.moments import require_moments
|
||||
from scribe.services.rule_moment_judgments import record_mount_change
|
||||
|
||||
wanted = require_moments(moments) or []
|
||||
async with async_session() as session:
|
||||
rule = await _fetch_owned_rule(session, rule_id, user_id)
|
||||
if rule is None:
|
||||
return None
|
||||
before = (await session.execute(
|
||||
select(rule_moments_t.c.moment).where(rule_moments_t.c.rule_id == rule_id)
|
||||
)).scalars().all()
|
||||
await session.execute(
|
||||
sql_delete(rule_moments_t).where(rule_moments_t.c.rule_id == rule_id)
|
||||
)
|
||||
@@ -1135,6 +1145,9 @@ async def set_rule_moments(
|
||||
await session.execute(
|
||||
insert(rule_moments_t).values(rule_id=rule_id, moment=moment)
|
||||
)
|
||||
await record_mount_change(
|
||||
session, rule_id, before, wanted, note=note, judged=judged, verdict=verdict,
|
||||
)
|
||||
await session.commit()
|
||||
return wanted
|
||||
|
||||
|
||||
Reference in New Issue
Block a user