Skip to content

popoto.recipes.reconciliation

popoto.recipes.reconciliation

M5 reconciliation: claim equivalence classes, typed contradiction rules, and explicit disjunctions over the provenance journal (#564).

The journal (:mod:popoto.recipes.provenance_journal) records every claim an agent captures, immutably and with full attribution. What it does not do is notice that two entries say the same thing. This module adds that layer: it groups entries that assert one claim into an equivalence class, resolves typed contradictions inside a class through a per-type precedence table, and stores a precedence tie as an explicit disjunct pair rather than picking an arbitrary winner.

Three properties shape every design choice here:

Nothing in this module ever mutates a persisted JournalEntry. JournalEntry composes AppendOnlyMixin, whose save() refuses any re-save of an existing key -- including a partial save(update_fields=[...]). So class membership cannot be a field on the entry. It lives in two ordinary (non-append-only) models this module owns outright, :class:ClaimMembership and :class:ClaimClass, where a relabel is an ordinary save(). The only writes M5 makes are (a) claim_type, set by capture before an entry's first and only save(), (b) new appended annotation entries, and (c) rows in its own two models.

The merge log is the source of truth; the two tables are a rebuildable index. Every reconcile outcome is appended to the journal as an immutable merge (or disjoin) annotation carrying the class ids, the rationale, the timestamp and the convention-book version. :func:replay discards the index and recomputes it from those annotations, which is what makes reversibility structural rather than a property to be tested into existence: retract a merge annotation, replay, and the pre-merge assignment is back. A crash mid-relabel is a repair, not a corruption.

One reconciler per agent, processing entries sequentially -- the single-writer invariant. This is a deployment constraint, not an implementation detail, and it is what the two documented races take as their mitigation. The :class:~popoto.streams.StreamConsumer on the "journal" stream is the only production trigger, and being the sole writer is what produces the invariant. :func:reconcile_entry is a thin adapter over the same reconcile function that exists for tests to drive; it is not a production entry point, because a host calling it from concurrent turn handling has nothing establishing the sequencing. Running a second reconciler per agent re-opens both races and requires reintroducing an atomic membership claim plus per-claim-slot serialization.

Example::

from popoto.recipes.provenance_journal import ProvenanceJournal
from popoto.recipes.reconciliation import (
    ClaimClass, representative_for, reconciliation_consumer,
)
from popoto.streams import StreamConsumer

# Capture assigns claim_type before the entry's first and only save().
ProvenanceJournal.append(
    agent_id="a1",
    statement="prefers morning meetings",
    subjects=["dana"],
    claim_type="preference",
)

# Production trigger: one consumer per agent, sequential.
consumer = reconciliation_consumer(agent_id="a1", consumer_name="worker-1")

# Downstream (M6/M7/M8) reads classes through the ORM.
for claim_class in ClaimClass.query.filter(agent_id="a1"):
    entry, uncertain = representative_for(claim_class.class_id)

ClaimMembership

Bases: Model

One row per reconciled entry: which class it belongs to.

Keyed by the entry's own Redis key, which gives exactly the identity semantics wanted -- one row per entry, no duplicate-membership state to reconcile. The colons inside that key are a non-issue: DB_key.clean() escapes them before the value reaches the keyspace, so a composite key value cannot forge a key boundary.

A plain Model deliberately -- no AppendOnlyMixin. A relabel on merge must be an ordinary save().

Privacy invariant: this row holds a digest, never claim content. claim_slot is a one-way sha256 of agent_id|subject|claim_type, and neither this model nor :class:ClaimClass stores the subject string or the plaintext claim_type. The reason is the exact scope of JournalEntry.hard_delete(): it erases a record and every trace of its own derived state, and explicitly not "every trace of the record anywhere in the keyspace". A plaintext subject on a mutable sibling model would be precisely such a field-value copy, sitting outside the reach of the only erasure primitive an append-only record has -- and a sharper regression than usual, because JournalEntry also composes NeverRecordMixin, i.e. this data is already governed as never-record. Slot equality is all reconciliation needs for grouping sibling claims, so the digest costs nothing here. See :func:erase_entry for the cascade that keeps the remaining derived state erasable.

Source code in src/popoto/recipes/reconciliation.py
class ClaimMembership(Model):
    """One row per reconciled entry: which class it belongs to.

    Keyed by the entry's own Redis key, which gives exactly the identity
    semantics wanted -- one row per entry, no duplicate-membership state to
    reconcile. The colons inside that key are a non-issue: ``DB_key.clean()``
    escapes them before the value reaches the keyspace, so a composite key
    value cannot forge a key boundary.

    A plain ``Model`` deliberately -- **no** ``AppendOnlyMixin``. A relabel on
    merge must be an ordinary ``save()``.

    **Privacy invariant: this row holds a digest, never claim content.**
    ``claim_slot`` is a one-way ``sha256`` of ``agent_id|subject|claim_type``,
    and neither this model nor :class:`ClaimClass` stores the subject string or
    the plaintext ``claim_type``. The reason is the exact scope of
    ``JournalEntry.hard_delete()``: it erases a record and every trace of *its
    own* derived state, and explicitly not "every trace of the record anywhere
    in the keyspace". A plaintext subject on a mutable sibling model would be
    precisely such a field-value copy, sitting outside the reach of the only
    erasure primitive an append-only record has -- and a sharper regression
    than usual, because ``JournalEntry`` also composes ``NeverRecordMixin``,
    i.e. this data is already governed as never-record. Slot *equality* is all
    reconciliation needs for grouping sibling claims, so the digest costs
    nothing here. See :func:`erase_entry` for the cascade that keeps the
    remaining derived state erasable.
    """

    entry_redis_key = KeyField()
    class_id = IndexedField(type=str)
    claim_slot = IndexedField(type=str)
    disjunction_id = IndexedField(type=str, null=True)

ClaimClass

Bases: Model

One row per equivalence class -- the read surface for M6/M7/M8.

Every read downstream needs is an ordinary indexed ORM query (ClaimClass.query.filter(agent_id=...), ClaimMembership.query.filter(class_id=...)) rather than an accessor wrapping HGET/SMEMBERS, which is what makes this a real read surface instead of a key convention a reader has to trust.

Holds no claim content, for the reason given on :class:ClaimMembership: representative_key is a Redis key, and the entry it names is where the claim text lives.

Source code in src/popoto/recipes/reconciliation.py
class ClaimClass(Model):
    """One row per equivalence class -- the read surface for M6/M7/M8.

    Every read downstream needs is an ordinary indexed ORM query
    (``ClaimClass.query.filter(agent_id=...)``,
    ``ClaimMembership.query.filter(class_id=...)``) rather than an accessor
    wrapping ``HGET``/``SMEMBERS``, which is what makes this a real read
    surface instead of a key convention a reader has to trust.

    Holds no claim content, for the reason given on :class:`ClaimMembership`:
    ``representative_key`` is a Redis key, and the entry it names is where the
    claim text lives.
    """

    class_id = KeyField()
    agent_id = IndexedField(type=str)
    representative_key = IndexedField(type=str)
    member_count = IntField(default=1)
    updated_at = FloatField(null=True)

Sameness

Bases: str, Enum

The judge's fixed two-value reply vocabulary.

Source code in src/popoto/recipes/reconciliation.py
class Sameness(str, enum.Enum):
    """The judge's fixed two-value reply vocabulary."""

    SAME = "same"
    DIFFERENT = "different"

SamenessResult dataclass

One judge answer, or the reason there isn't one.

Attributes:

Name Type Description
verdict Optional[Sameness]

The parsed :class:Sameness, or None on an abstention.

abstained bool

True when no usable verdict was obtained. An abstention is never read as different on the forward ask (it leaves the entry a singleton) and never as same anywhere.

reason str

A short machine token for the log and the merge-log rationale. Never model-authored text.

Source code in src/popoto/recipes/reconciliation.py
@dataclass(frozen=True)
class SamenessResult:
    """One judge answer, or the reason there isn't one.

    Attributes:
        verdict: The parsed :class:`Sameness`, or ``None`` on an abstention.
        abstained: True when no usable verdict was obtained. An abstention is
            never read as ``different`` on the forward ask (it leaves the entry
            a singleton) and never as ``same`` anywhere.
        reason: A short machine token for the log and the merge-log rationale.
            Never model-authored text.
    """

    verdict: Optional[Sameness]
    abstained: bool
    reason: str

    @property
    def is_same(self) -> bool:
        return self.verdict is Sameness.SAME and not self.abstained

ReconcileOutcome dataclass

What one reconcile pass did to one entry.

Attributes:

Name Type Description
entry_key str

The reconciled entry's Redis key.

class_id str

The class it belongs to afterwards.

action str

One of created, confirmed, joined, superseded, loser-absent, disjoined, noop.

judge_calls int

Calls the judge actually issued, so AC5's bound is observable rather than asserted.

disjunction_id str

The shared id when this pass stored a disjunct pair.

superseded_key str

The loser's key when a type rule resolved.

Source code in src/popoto/recipes/reconciliation.py
@dataclass
class ReconcileOutcome:
    """What one reconcile pass did to one entry.

    Attributes:
        entry_key: The reconciled entry's Redis key.
        class_id: The class it belongs to afterwards.
        action: One of ``created``, ``confirmed``, ``joined``, ``superseded``,
            ``loser-absent``, ``disjoined``, ``noop``.
        judge_calls: Calls the judge actually issued, so AC5's bound is
            observable rather than asserted.
        disjunction_id: The shared id when this pass stored a disjunct pair.
        superseded_key: The loser's key when a type rule resolved.
    """

    entry_key: str
    class_id: str
    action: str
    judge_calls: int = 0
    disjunction_id: str = ""
    superseded_key: str = ""
    merged_class_ids: List[str] = dataclass_field(default_factory=list)

normalize_claim_type(claim_type)

Return a type from :data:CLAIM_TYPES, falling back to note.

Tolerating None is a requirement, not a convenience: claim_type is write-once-at-capture, so every entry captured before #564 shipped has None and no code path can back-fill it. The reconciler must classify those entries rather than skip them or raise. An unrecognized string is treated the same way, so a capture path emitting a type outside the frozen enum degrades to the rule-free catch-all instead of reaching a precedence lookup with no row.

Source code in src/popoto/recipes/reconciliation.py
def normalize_claim_type(claim_type: Optional[str]) -> str:
    """Return a type from :data:`CLAIM_TYPES`, falling back to ``note``.

    Tolerating ``None`` is a requirement, not a convenience: ``claim_type`` is
    write-once-at-capture, so every entry captured before #564 shipped has
    ``None`` and no code path can back-fill it. The reconciler must classify
    those entries rather than skip them or raise. An unrecognized string is
    treated the same way, so a capture path emitting a type outside the frozen
    enum degrades to the rule-free catch-all instead of reaching a precedence
    lookup with no row.
    """
    if isinstance(claim_type, str) and claim_type in CLAIM_TYPES:
        return claim_type
    return DEFAULT_CLAIM_TYPE

claim_slot(agent_id, subject, claim_type)

Return the one-way slot digest for (agent_id, subject, claim_type).

32 hex characters of sha256. Computed at reconcile time and stored one-way, so the slot supports the equality test grouping needs while carrying none of the text -- see :class:ClaimMembership's privacy invariant.

Source code in src/popoto/recipes/reconciliation.py
def claim_slot(agent_id: str, subject: str, claim_type: Optional[str]) -> str:
    """Return the one-way slot digest for ``(agent_id, subject, claim_type)``.

    32 hex characters of ``sha256``. Computed at reconcile time and stored
    one-way, so the slot supports the equality test grouping needs while
    carrying none of the text -- see :class:`ClaimMembership`'s privacy
    invariant.
    """
    raw = f"{agent_id}|{subject}|{normalize_claim_type(claim_type)}"
    return hashlib.sha256(raw.encode("utf-8")).hexdigest()[:32]

slot_for_entry(entry)

Return entry's claim slot, using its first subject tag.

An entry with no subjects gets the empty subject, which puts every untagged claim of one type into one slot. That is deliberate and matches TagField's own zero-tag semantics: such an entry has no identity key to compute, so it skips the deterministic tier and reaches the judge with a subject-unbounded shortlist bounded by :data:M5_SHORTLIST_CAP.

Source code in src/popoto/recipes/reconciliation.py
def slot_for_entry(entry: Any) -> str:
    """Return ``entry``'s claim slot, using its first subject tag.

    An entry with no subjects gets the empty subject, which puts every
    untagged claim of one type into one slot. That is deliberate and matches
    ``TagField``'s own zero-tag semantics: such an entry has no identity key to
    compute, so it skips the deterministic tier and reaches the judge with a
    subject-unbounded shortlist bounded by :data:`M5_SHORTLIST_CAP`.
    """
    subjects = list(getattr(entry, "subjects", None) or [])
    subject = str(subjects[0]) if subjects else ""
    return claim_slot(str(entry.agent_id), subject, getattr(entry, "claim_type", None))

judge_sameness(left_statement, right_statement, client=None)

Ask once whether two claims are the same claim. Never raises.

Order of operations, mirroring llm_verdict:

  1. Either statement blank or whitespace-only -> abstain, zero calls. There is nothing to compare.
  2. scan_never_record runs on both statements before the call. A blocked statement abstains and its text is never transmitted.
  3. Otherwise one call is issued. A malformed, empty or out-of-vocabulary reply, an unreachable provider, and a raising client all abstain.

Parameters:

Name Type Description Default
left_statement str

The first claim. Order matters -- see :func:_judge_pair, which re-asks with the order swapped.

required
right_statement str

The second claim.

required
client Any

An Anthropic-style client (anything exposing messages.create). None builds the default client, which requires the optional anthropic package.

None

Returns:

Name Type Description
A SamenessResult

class:SamenessResult carrying an enum or an abstention -- never

SamenessResult

any model-authored text.

Source code in src/popoto/recipes/reconciliation.py
def judge_sameness(
    left_statement: str,
    right_statement: str,
    client: Any = None,
) -> SamenessResult:
    """Ask once whether two claims are the same claim. Never raises.

    Order of operations, mirroring ``llm_verdict``:

    1. Either statement blank or whitespace-only -> abstain, **zero calls**.
       There is nothing to compare.
    2. ``scan_never_record`` runs on both statements **before** the call. A
       blocked statement abstains and its text is never transmitted.
    3. Otherwise one call is issued. A malformed, empty or out-of-vocabulary
       reply, an unreachable provider, and a raising client all abstain.

    Args:
        left_statement: The first claim. Order matters -- see
            :func:`_judge_pair`, which re-asks with the order swapped.
        right_statement: The second claim.
        client: An Anthropic-style client (anything exposing
            ``messages.create``). ``None`` builds the default client, which
            requires the optional ``anthropic`` package.

    Returns:
        A :class:`SamenessResult` carrying an enum or an abstention -- never
        any model-authored text.
    """
    if not (left_statement or "").strip() or not (right_statement or "").strip():
        logger.debug("sameness: blank statement -> abstain(empty_statement)")
        return SamenessResult(None, True, "empty_statement")

    for statement in (left_statement, right_statement):
        firewall = scan_never_record(statement)
        if firewall.blocked:
            # The reason code only -- never a fragment of the blocked text.
            logger.info(
                "sameness: statement blocked by never-record firewall (%s)",
                firewall.reason,
            )
            return SamenessResult(None, True, "firewall_drop")

    try:
        if client is None:
            client = _default_client()
        response = client.messages.create(
            model=M5_JUDGE_MODEL,
            max_tokens=M5_JUDGE_MAX_TOKENS,
            system=CONVENTION_BOOK_V1,
            messages=[
                {
                    "role": "user",
                    "content": json.dumps(
                        {"claim_a": left_statement, "claim_b": right_statement}
                    ),
                }
            ],
            output_config={
                "format": {"type": "json_schema", "schema": SAMENESS_SCHEMA}
            },
        )
        raw_text = next(
            (block.text for block in response.content if block.type == "text"),
            None,
        )
    except Exception as e:
        logger.warning("sameness: judge call failed: %s", e)
        return SamenessResult(None, True, "llm_unavailable")

    verdict = _parse_sameness(raw_text)
    if verdict is None:
        logger.warning("sameness: malformed reply -> abstain(llm_unavailable)")
        return SamenessResult(None, True, "llm_unavailable")
    return SamenessResult(verdict, False, verdict.value)

cached_embedding(entry, provider=None)

Return entry's statement vector, embedding and caching on a miss.

None when no provider is available or the provider fails, which is the signal :func:shortlist_candidates uses to take its index-scan fallback.

Source code in src/popoto/recipes/reconciliation.py
def cached_embedding(entry: Any, provider: Any = None) -> Optional[List[float]]:
    """Return ``entry``'s statement vector, embedding and caching on a miss.

    ``None`` when no provider is available or the provider fails, which is the
    signal :func:`shortlist_candidates` uses to take its index-scan fallback.
    """
    redis_key = entry.pk
    backend = _cache_backend(type(entry))
    if backend is not None:
        cached = backend.field_call(
            type(entry)._meta.spec, "_embed_cache", "get", redis_key
        )
        if cached:
            return list(cached)
        raw = None
    else:
        client = get_REDIS_DB()
        raw = client.hget(EMBEDDING_CACHE_KEY, redis_key)
    if raw:
        try:
            return list(json.loads(raw))
        except (ValueError, TypeError):
            logger.warning("embedding cache: undecodable entry, re-embedding")

    if provider is None:
        provider = _embedding_provider()
    if provider is None:
        return None

    statement = str(getattr(entry, "statement", "") or "")
    if not statement.strip():
        return None
    try:
        vectors = provider.embed([statement], input_type="document")
    except Exception as e:
        logger.warning("embedding provider failed, falling back to scan: %s", e)
        return None
    if not vectors or not vectors[0]:
        return None
    vector = [float(v) for v in vectors[0]]
    if backend is not None:
        backend.field_call(
            type(entry)._meta.spec, "_embed_cache", "set", redis_key, vector
        )
    else:
        client.hset(EMBEDDING_CACHE_KEY, redis_key, json.dumps(vector))
    return vector

drop_cached_embedding(redis_key, model=None)

Delete one entry's cached vector.

Part of :func:erase_entry's cascade: an embedding is a lossy encoding of statement, so the cache is content-derived state in a store hard_delete() does not reach. model names the entry model, whose backend holds the cache when it is not Redis (#759 M4).

Source code in src/popoto/recipes/reconciliation.py
def drop_cached_embedding(redis_key: str, model: Any = None) -> None:
    """Delete one entry's cached vector.

    Part of :func:`erase_entry`'s cascade: an embedding is a lossy encoding of
    ``statement``, so the cache is content-derived state in a store
    ``hard_delete()`` does not reach. ``model`` names the entry model, whose
    backend holds the cache when it is not Redis (#759 M4).
    """
    backend = _cache_backend(model) if model is not None else None
    if backend is not None:
        backend.field_call(model._meta.spec, "_embed_cache", "drop", redis_key)
        return
    get_REDIS_DB().hdel(EMBEDDING_CACHE_KEY, redis_key)

shortlist_candidates(entry, *, exclude_class_ids=(), provider=None)

Return at most :data:M5_SHORTLIST_CAP candidate class ids for entry.

Ranked by cosine similarity between the entry's statement vector and each class representative's, over the same-agent classes. When no embedding provider is available -- or the provider fails -- this degrades to a bounded same-subject + same-type index scan: recall narrows, every correctness property holds, and the judge-call bound is unchanged (Risk 4).

Source code in src/popoto/recipes/reconciliation.py
def shortlist_candidates(
    entry: Any,
    *,
    exclude_class_ids: Sequence[str] = (),
    provider: Any = None,
) -> List[str]:
    """Return at most :data:`M5_SHORTLIST_CAP` candidate class ids for ``entry``.

    Ranked by cosine similarity between the entry's statement vector and each
    class representative's, over the same-agent classes. When no embedding
    provider is available -- or the provider fails -- this degrades to a
    bounded same-subject + same-type index scan: recall narrows, every
    correctness property holds, and the judge-call bound is unchanged (Risk 4).
    """
    cap = int(M5_SHORTLIST_CAP)
    if cap <= 0:
        return []
    excluded = set(exclude_class_ids)
    classes = [
        row
        for row in ClaimClass.query.filter(agent_id=str(entry.agent_id))
        if row.class_id not in excluded
    ]
    if not classes:
        return []

    entry_vector = cached_embedding(entry, provider)
    if entry_vector is None:
        # Degraded path: same slot first (an exact index equality), then any
        # remaining same-agent class, both bounded by the cap.
        slot = slot_for_entry(entry)
        slot_classes = {
            row.class_id for row in ClaimMembership.query.filter(claim_slot=slot)
        }
        ordered = [row.class_id for row in classes if row.class_id in slot_classes]
        ordered += [row.class_id for row in classes if row.class_id not in slot_classes]
        return ordered[:cap]

    scored: List[Tuple[float, str]] = []
    for row in classes:
        representative = JournalEntry.query.get(redis_key=row.representative_key)
        if representative is None:
            continue
        other = cached_embedding(representative, provider)
        if other is None:
            continue
        scored.append((_cosine(entry_vector, other), row.class_id))
    scored.sort(key=lambda pair: pair[0], reverse=True)
    return [class_id for _score, class_id in scored[:cap]]

confirmation_count(entry)

Return entry's corroboration count.

Derived by counting confirm annotations, never stored on the record: ProvenanceJournal.confirm appends a new annotation and leaves the target untouched, which is what makes corroboration append-only-safe.

Source code in src/popoto/recipes/reconciliation.py
def confirmation_count(entry: Any) -> int:
    """Return ``entry``'s corroboration count.

    Derived by counting ``confirm`` annotations, never stored on the record:
    ``ProvenanceJournal.confirm`` appends a new annotation and leaves the
    target untouched, which is what makes corroboration append-only-safe.
    """
    return sum(
        1
        for annotation in ProvenanceJournal.annotations_for(entry)
        if annotation.kind == "confirm"
    )

representative_for(class_id)

Return (representative_entry, uncertainty_flag) for a class.

The representative is the class's most-confirmed member among validity-open entries only, ties broken by recency. That is what the judge is always shown, and it is M5's guarantee to M6: one representative per class is selectable.

The second element is the uncertainty marker, and it is returned here, from M5's own selection call, rather than left to a reader's formatting layer: it is True when any live member of the class carries an open disjunction_id, i.e. the class holds a precedence tie no winner was picked for. A reader that ignores the flag degrades to showing a representative; it cannot be handed a silent winner, because there isn't one to hand.

Returns:

Type Description
Tuple[Optional[Any], bool]

(None, False) for an empty or fully-closed class.

Source code in src/popoto/recipes/reconciliation.py
def representative_for(class_id: str) -> Tuple[Optional[Any], bool]:
    """Return ``(representative_entry, uncertainty_flag)`` for a class.

    The representative is the class's **most-confirmed member among
    validity-open entries only**, ties broken by recency. That is what the
    judge is always shown, and it is M5's guarantee to M6: one representative
    per class is *selectable*.

    The second element is the uncertainty marker, and it is returned **here**,
    from M5's own selection call, rather than left to a reader's formatting
    layer: it is True when any live member of the class carries an open
    ``disjunction_id``, i.e. the class holds a precedence tie no winner was
    picked for. A reader that ignores the flag degrades to showing a
    representative; it cannot be handed a silent winner, because there isn't
    one to hand.

    Returns:
        ``(None, False)`` for an empty or fully-closed class.
    """
    live = _live_members(class_id)
    if not live:
        return None, False
    live.sort(
        key=lambda e: (confirmation_count(e), float(e.captured_at or 0.0)),
        reverse=True,
    )
    uncertain = any(
        bool(row.disjunction_id)
        for row in ClaimMembership.query.filter(class_id=class_id)
    )
    return live[0], uncertain

resolve_precedence(challenger, incumbent, claim_type)

Resolve a fired type rule into (winner, loser, basis).

Applied only after a type rule fires: this never decides sameness, and it never runs on a note.

Global Rule 0 is evaluated first for every type -- a self-stated claim beats an inferred one, consuming M1's stated flag -- and the per-type row applies only when both claims agree on stated. Then the family order from :data:PRECEDENCE_ORDER: recency for the supersession family (deadline), confirmation count then recency for the stable family.

The table is total. When Rule 0 ties, the family order ties, and recency ties, this returns (None, None, PRECEDENCE_TIE) and the caller stores a disjunct pair -- never an arbitrary winner.

Source code in src/popoto/recipes/reconciliation.py
def resolve_precedence(
    challenger: Any, incumbent: Any, claim_type: str
) -> Tuple[Optional[Any], Optional[Any], str]:
    """Resolve a fired type rule into ``(winner, loser, basis)``.

    Applied **only** after a type rule fires: this never decides sameness, and
    it never runs on a ``note``.

    Global **Rule 0** is evaluated first for every type -- a self-stated claim
    beats an inferred one, consuming M1's ``stated`` flag -- and the per-type
    row applies only when both claims agree on ``stated``. Then the family
    order from :data:`PRECEDENCE_ORDER`: recency for the supersession family
    (``deadline``), confirmation count then recency for the stable family.

    The table is total. When Rule 0 ties, the family order ties, and recency
    ties, this returns ``(None, None, PRECEDENCE_TIE)`` and the caller stores a
    disjunct pair -- never an arbitrary winner.
    """
    family = CLAIM_TYPE_FAMILY.get(claim_type, "rule_free")
    if family == "rule_free":
        return None, None, PRECEDENCE_TIE

    # Rule 0, global and first.
    challenger_stated = bool(challenger.stated)
    incumbent_stated = bool(incumbent.stated)
    if challenger_stated != incumbent_stated:
        if challenger_stated:
            return challenger, incumbent, "rule0_stated"
        return incumbent, challenger, "rule0_stated"

    # Both columns are compared as floats: ``confirmation_count`` returns a
    # small int that is exactly representable, and recency is a timestamp.
    # The comparison never crosses columns, so the widening loses nothing.
    left: float
    right: float
    for column in PRECEDENCE_ORDER[family]:
        if column == "confirmations":
            left = confirmation_count(challenger)
            right = confirmation_count(incumbent)
        else:  # "recency"
            left = float(challenger.captured_at or 0.0)
            right = float(incumbent.captured_at or 0.0)
        if left > right:
            return challenger, incumbent, column
        if right > left:
            return incumbent, challenger, column

    return None, None, PRECEDENCE_TIE

merge_log_entries(agent_id)

Return an agent's live merge-log annotations, oldest first.

Only validity__current=True annotations are returned, which is what makes AC4 work: retracting a merge annotation closes its interval, so the next :func:replay no longer sees it and reproduces the pre-merge assignment.

Source code in src/popoto/recipes/reconciliation.py
def merge_log_entries(agent_id: str) -> List[Any]:
    """Return an agent's live merge-log annotations, oldest first.

    Only ``validity__current=True`` annotations are returned, which is what
    makes AC4 work: retracting a ``merge`` annotation closes its interval, so
    the next :func:`replay` no longer sees it and reproduces the pre-merge
    assignment.
    """
    live: List[Any] = []
    for kind in MERGE_KINDS:
        live.extend(
            JournalEntry.query.filter(
                agent_id=agent_id, kind=kind, validity__current=True
            )
        )
    live.sort(key=lambda e: float(e.captured_at or 0.0))
    return live

reconcile_entry(entry, *, client=None, provider=None)

Reconcile one entry directly. Test-only -- not a production path.

A thin adapter over the same reconcile function the stream consumer drives. One loop, one production trigger, never two pipelines.

This is deliberately not documented as a host-facing API. The single-writer invariant that Races 1 and 3 take as their mitigation is not a property of the reconcile function; it is produced by the consumer being the sole writer, one reconciler per agent processing entries sequentially. A consumer-less host calling this from concurrent turn handling has nothing establishing that sequencing, which re-opens exactly the concurrent-join hazard the invariant covers. Use :func:reconciliation_consumer in production.

Source code in src/popoto/recipes/reconciliation.py
def reconcile_entry(
    entry: Any, *, client: Any = None, provider: Any = None
) -> ReconcileOutcome:
    """Reconcile one entry directly. **Test-only** -- not a production path.

    A thin adapter over the same reconcile function the stream consumer drives.
    One loop, one production trigger, never two pipelines.

    This is deliberately **not** documented as a host-facing API. The
    single-writer invariant that Races 1 and 3 take as their mitigation is not
    a property of the reconcile function; it is *produced by* the consumer
    being the sole writer, one reconciler per agent processing entries
    sequentially. A consumer-less host calling this from concurrent turn
    handling has nothing establishing that sequencing, which re-opens exactly
    the concurrent-join hazard the invariant covers. Use
    :func:`reconciliation_consumer` in production.
    """
    return _reconcile(entry, client=client, provider=provider)

replay(agent_id, *, since=None, rebuild=False)

Rebuild the class index from the merge log. Returns entries replayed.

The merge log is the source of truth and these two tables are a rebuildable index derived from it, so there is no sidecar to keep in sync: a divergent index is discarded and recomputed rather than reconciled against the log. That is what makes AC4's reversibility structural -- retract a merge annotation, replay, and the pre-merge assignment is back, because :func:merge_log_entries only reads live annotations.

Parameters:

Name Type Description Default
agent_id str

The agent whose index to rebuild.

required
since Optional[float]

Replay only annotations strictly newer than this instant (captured_at, the field named by :data:M5_REPLAY_WATERMARK_FIELD, filtered with a strict >, borrowing crystallize's watermark shape). None replays from genesis, which is an explicit repair operation rather than the steady state.

None
rebuild bool

Delete this agent's existing index rows first. Required for a true from-genesis rebuild; without it a replay is additive.

False
Source code in src/popoto/recipes/reconciliation.py
def replay(
    agent_id: str,
    *,
    since: Optional[float] = None,
    rebuild: bool = False,
) -> int:
    """Rebuild the class index from the merge log. Returns entries replayed.

    The merge log is the source of truth and these two tables are a rebuildable
    index derived from it, so there is no sidecar to keep in sync: a divergent
    index is discarded and recomputed rather than reconciled against the log.
    That is what makes AC4's reversibility structural -- retract a ``merge``
    annotation, replay, and the pre-merge assignment is back, because
    :func:`merge_log_entries` only reads live annotations.

    Args:
        agent_id: The agent whose index to rebuild.
        since: Replay only annotations strictly newer than this instant
            (``captured_at``, the field named by
            :data:`M5_REPLAY_WATERMARK_FIELD`, filtered with a strict ``>``,
            borrowing ``crystallize``'s watermark shape). ``None`` replays from
            genesis, which is an explicit repair operation rather than the
            steady state.
        rebuild: Delete this agent's existing index rows first. Required for a
            true from-genesis rebuild; without it a replay is additive.
    """
    if rebuild:
        for row in ClaimClass.query.filter(agent_id=agent_id):
            for member in ClaimMembership.query.filter(class_id=row.class_id):
                member.delete()
            row.delete()

    replayed = 0
    touched: List[str] = []
    for annotation in merge_log_entries(agent_id):
        watermark = float(getattr(annotation, M5_REPLAY_WATERMARK_FIELD) or 0.0)
        if since is not None and not watermark > since:
            continue
        payload = _decode_payload(annotation)
        if payload is None:
            logger.warning(
                "replay: undecodable merge-log payload on %s, skipped",
                annotation.pk,
            )
            continue
        class_a = str(payload.get("class_a") or "")
        class_b = str(payload.get("class_b") or "")
        slot = str(payload.get("slot") or "")
        if not class_a:
            continue

        if annotation.kind == "disjoin":
            disjunction_id = str(payload.get("disjunction_id") or "")
            for key in (str(annotation.target or ""), str(payload.get("other") or "")):
                row = ClaimMembership.query.get(entry_redis_key=key)
                if row is not None:
                    row.disjunction_id = disjunction_id or None
                    row.save()
            replayed += 1
            continue

        target_key = str(annotation.target or "")
        if target_key:
            row = ClaimMembership.query.get(entry_redis_key=target_key)
            if row is None:
                ClaimMembership(
                    entry_redis_key=target_key,
                    class_id=class_a,
                    claim_slot=slot,
                ).save()
            else:
                row.class_id = class_a
                row.claim_slot = slot or row.claim_slot
                row.save()
        if class_b and class_b != class_a:
            _relabel_class(class_b, class_a)
            if class_b not in touched:
                touched.append(class_b)
        if class_a not in touched:
            touched.append(class_a)
        replayed += 1

    for class_id in touched:
        _recompute_class(class_id)
    return replayed

erase_entry(entry)

Erase a reconciled entry and every trace of M5's derived state.

Use this, not JournalEntry.hard_delete() directly, for a reconciled entry. hard_delete() is the retention/erasure primitive and its documented scope is the record plus every trace of its own derived state -- explicitly not "every trace of the record anywhere in the keyspace". M5 adds derived state the primitive therefore does not reach, so a bare hard_delete() leaves a dangling membership row and, worse, a :class:ClaimClass whose representative_key points at an erased key.

Four legs, in order:

  1. JournalEntry.hard_delete() on the entry.
  2. Delete its :class:ClaimMembership row.
  3. Delete its cached reconciler-side embedding -- an embedding is a lossy encoding of statement, so the cache is content-derived state in a store hard_delete() does not reach.
  4. Recompute the affected :class:ClaimClass: reselect representative_key and decrement member_count, or drop the row when the class is left empty.

The membership row is survivable at all only because it carries the one-way claim_slot digest and no claim content -- which matters because JournalEntry also composes NeverRecordMixin, so this data is governed as never-record.

Source code in src/popoto/recipes/reconciliation.py
def erase_entry(entry: Any) -> bool:
    """Erase a reconciled entry and every trace of M5's derived state.

    **Use this, not** ``JournalEntry.hard_delete()`` **directly**, for a
    reconciled entry. ``hard_delete()`` is the retention/erasure primitive and
    its documented scope is the record plus every trace of *its own* derived
    state -- explicitly **not** "every trace of the record anywhere in the
    keyspace". M5 adds derived state the primitive therefore does not reach, so
    a bare ``hard_delete()`` leaves a dangling membership row and, worse, a
    :class:`ClaimClass` whose ``representative_key`` points at an erased key.

    Four legs, in order:

    1. ``JournalEntry.hard_delete()`` on the entry.
    2. Delete its :class:`ClaimMembership` row.
    3. Delete its cached reconciler-side embedding -- an embedding is a lossy
       encoding of ``statement``, so the cache is content-derived state in a
       store ``hard_delete()`` does not reach.
    4. Recompute the affected :class:`ClaimClass`: reselect
       ``representative_key`` and decrement ``member_count``, or drop the row
       when the class is left empty.

    The membership row is survivable at all only because it carries the
    one-way ``claim_slot`` digest and no claim content -- which matters
    because ``JournalEntry`` also composes ``NeverRecordMixin``, so this data
    is governed as never-record.
    """
    redis_key = entry.pk
    row = ClaimMembership.query.get(entry_redis_key=redis_key)
    class_id = row.class_id if row is not None else ""

    erased = JournalEntry.hard_delete(entry)
    if row is not None:
        row.delete()
    drop_cached_embedding(redis_key, type(entry))
    if class_id:
        _recompute_class(class_id)
    return bool(erased)

make_reconciliation_handler(agent_id=None, *, client=None, provider=None)

Build the StreamConsumer handler that is M5's production trigger.

Entries are processed sequentially inside one handler call, and the deployment contract is one consumer per agent. That sequencing is what produces the single-writer invariant: entry A's membership row is committed before entry B is shortlisted, so the interleaved shortlist-to-commit span Race 3 needs never occurs.

Parameters:

Name Type Description Default
agent_id Optional[str]

When given, only this agent's entries are reconciled -- filtered on the stream's agent_id metadata field, without hydrating the record.

None
client Any

Judge client, forwarded to the reconcile function.

None
provider Any

Embedding provider, forwarded to the shortlist.

None
Source code in src/popoto/recipes/reconciliation.py
def make_reconciliation_handler(
    agent_id: Optional[str] = None,
    *,
    client: Any = None,
    provider: Any = None,
) -> Callable[[StreamBatch], Awaitable[None]]:
    """Build the ``StreamConsumer`` handler that is M5's production trigger.

    Entries are processed **sequentially** inside one handler call, and the
    deployment contract is one consumer per agent. That sequencing is what
    produces the single-writer invariant: entry A's membership row is committed
    before entry B is shortlisted, so the interleaved shortlist-to-commit span
    Race 3 needs never occurs.

    Args:
        agent_id: When given, only this agent's entries are reconciled --
            filtered on the stream's ``agent_id`` metadata field, without
            hydrating the record.
        client: Judge client, forwarded to the reconcile function.
        provider: Embedding provider, forwarded to the shortlist.
    """

    async def reconciliation_handler(entries: StreamBatch) -> None:
        """Reconcile each ``assert`` capture the journal stream announces."""
        if not entries:
            return
        for stream_id, fields in entries:
            if fields.get("op") not in ("create", None, ""):
                continue
            if fields.get("kind") not in ("assert", None, ""):
                # Annotations -- including M5's own merge-log appends -- are
                # not claims to reconcile. Skipping them here is also what
                # keeps the consumer from reconciling its own writes.
                continue
            if agent_id is not None and fields.get("agent_id") != agent_id:
                continue
            redis_key = fields.get("pk") or ""
            if not redis_key:
                continue
            entry = JournalEntry.query.get(redis_key=redis_key)
            if entry is None:
                logger.debug(
                    "reconcile: stream entry %s names no readable record",
                    stream_id,
                )
                continue
            try:
                _reconcile(entry, client=client, provider=provider)
            except Exception:
                # Re-raised after logging: the consumer's retry and
                # dead-letter machinery is the right place to decide, and
                # swallowing here would drop a claim silently.
                logger.exception(
                    "reconcile: failed on %s from stream entry %s",
                    redis_key,
                    stream_id,
                )
                raise

    return reconciliation_handler

reconciliation_consumer(*, agent_id=None, consumer_name, group_name='m5-reconciler', client=None, provider=None)

Build the journal-stream consumer. One per agent -- see the docstring.

Running a second reconciler for one agent is a deployment error, not a state this code arbitrates: it re-opens Races 1 and 3 and requires reintroducing an atomic membership claim plus per-claim-slot serialization at the same time. Whoever proposes the second writer owns that change.

Source code in src/popoto/recipes/reconciliation.py
def reconciliation_consumer(
    *,
    agent_id: Optional[str] = None,
    consumer_name: str,
    group_name: str = "m5-reconciler",
    client: Any = None,
    provider: Any = None,
) -> "StreamConsumer":
    """Build the journal-stream consumer. **One per agent** -- see the docstring.

    Running a second reconciler for one agent is a deployment error, not a
    state this code arbitrates: it re-opens Races 1 and 3 and requires
    reintroducing an atomic membership claim plus per-claim-slot serialization
    at the same time. Whoever proposes the second writer owns that change.
    """
    from ..streams import StreamConsumer

    return StreamConsumer(
        stream_key=JOURNAL_STREAM_KEY,
        group_name=group_name,
        consumer_name=consumer_name,
        handler=make_reconciliation_handler(agent_id, client=client, provider=provider),
    )