Skip to content

popoto.recipes.question_queue

popoto.recipes.question_queue

M7 question queue: a rationed clarifying-question channel (#566).

The memory layer can detect that it is uncertain -- an M5 disjunct pair, a confidence-gate refusal, an M4 evidence gap -- and until now could do nothing about it except refuse. This module gives that uncertainty one shared, strictly rationed outlet. Every subsystem that wants to ask the person something writes a :class:QuestionCandidate with :func:propose instead of asking directly, and the host application asks the queue, once per turn, whether there is anything worth asking with :func:next_question.

The contract ends at "the next question to ask, if any." Nothing here writes to a transport or addresses a person. How a rider question attaches to an assistant reply is host-application territory.

Five rules shape every function below:

  1. Turns are host-supplied. Every public call takes a monotonic turn: int. The library owns no turn counter: the host is the only party that actually knows what a turn is, and a library-side counter would be wrong under concurrency.
  2. The value-of-information gate is a deterministic boolean computed from stored metadata (:func:passes_voi_gate): the candidate carries a known ambiguity_signal and was recently used. There is no probabilistic VOI score, on purpose.
  3. The budget is structural. At most one delivery per QUESTION_BUDGET_TURNS turns per agent, enforced by a Lua script (:data:_DELIVER_LUA) that grants the budget and claims a candidate in one step, so budget is spent only when a question is actually delivered -- never an in-process counter, never reset on restart, never rewound by a regressed turn.
  4. Answers are evidence, never overwrites. An answered reply is mapped onto the existing closed outcome vocabulary (acted / contradicted) and applied through ObservationProtocol.on_context_used plus QUESTION_ANSWER_WEIGHT - 1 extra capped-Bayesian observations. It never sets _superseded_by, so it never closes a validity interval. A later contradiction still moves confidence.
  5. Fail closed. propose*, :func:next_question and :func:record_answer log and return None / a non-applied result on any Redis error; they never raise into the caller. Invalid input (an unknown kind, an option whose acted and contradicted lists overlap) is a producer bug and still raises ValueError.

The free-text answer is never stored or logged -- only the matched option index (answer_option) -- because a reply can echo sensitive content.

Extension seam for producers

A producer whose source of ambiguity can disappear before the question is asked (an M5 disjoin annotation that gets retracted) registers a staleness check with :func:register_staleness_check. :func:expire_stale -- which :func:next_question runs before every delivery -- expires any open candidate of that kind for which the check returns True, so a stale question is never asked. The M5 disjunction producer registers exactly such a check at import time (:func:_disjunction_is_stale).

Producers

Three thin adapters turn an existing ambiguity signal into a :func:propose call. Each is a separate function the host calls when, and only when, it wants that source -- so each is independently disableable by not calling it, and the queue works with any subset (including none):

  • :func:propose_from_disjunctions -- M5 disjoin annotations (#564). Imports reconciliation lazily, so the queue works without M5.
  • :func:propose_from_gate -- a confidence-gate refusal in assemble() metadata. Dormant unless the host configures ContextAssembler(confidence_gate_threshold=...): without a threshold there is no metadata["gate"] and the producer has no input.
  • :func:propose_from_resolution -- M4 evidence_gap references on a ResolutionRecord (#563).

All three fail closed: on any error they log and return 0 / None and never raise into the caller.

Example::

from popoto.recipes import question_queue as qq

qq.propose(
    agent_id="a1",
    question_text="Does Dana prefer morning or afternoon meetings?",
    kind="disjunction",
    source_module="reconciliation",
    target_keys=[key_a, key_b],
    options=[
        {"label": "morning", "acted": [key_a], "contradicted": [key_b]},
        {"label": "afternoon", "acted": [key_b], "contradicted": [key_a]},
    ],
    ambiguity_signal="disjunction",
    turn=12,
)

q = qq.next_question("a1", turn=13, query_cues="morning meetings with Dana")
if q is not None:
    ...  # host phrases and delivers q.question_text
    qq.record_answer(q, "morning", turn=14)

QuestionCandidate

Bases: Model

One persisted clarifying question awaiting the gate, the budget and a relevant moment.

A real Model rather than a ZSET because the question is persisted state with a payload and a status (RecallProposal is payload-free and TTL-expiring, and cannot be retrofitted).

status, ask_count, delivered_turn, answer_option, cooldown_until and resolved_turn are deliberately unindexed plain fields: delivery and :func:record_answer write them with Lua directly on the model hash, which an index would not see. Reads go through the indexed agent_id.

Fields

question_text: The question, phrased by the producer (or by the host later). Echoes content, so retention is bounded by prune(). kind: One of :data:KINDS. source_module: Free-form producer name (e.g. "reconciliation"). target_keys: Redis keys of the facts the answer would resolve. options: [{"label": str, "acted": [keys], "contradicted": [keys]}]. Answering applies exactly the chosen option's two lists. ambiguity_signal: One of :data:AMBIGUITY_SIGNALS. status: One of :data:STATUSES. ask_count: Times delivered. Re-asking after cooldown is legal. created_turn / expires_turn: Proposal turn and created + N. last_seen_turn: Most recent turn the ambiguity was re-observed or the targets were used -- the VOI gate's impact factor. cooldown_until: Not deliverable before this turn (after a deflected/unrecognized reply). answer_option: 0-based index of the chosen option. The free-text reply itself is never stored. resolved_turn: Turn the candidate left the open set (answered, cooled, expired), used as the retention reference by prune(). cue_tokens: See :func:cue_tokens_for. delivered_turn: Turn of the most recent delivery, written by the delivery Lua. A delivered candidate that gets no reply within QUESTION_COOLDOWN_TURNS of it is treated like a non-answer and moved to cooled (see :func:expire_stale). disjunction_id: Set by the M5 producer so a retracted disjoin can expire its candidate via a registered staleness check. agent_id: Owning agent.

Source code in src/popoto/recipes/question_queue.py
class QuestionCandidate(Model):
    """One persisted clarifying question awaiting the gate, the budget and a
    relevant moment.

    A real Model rather than a ZSET because the question is persisted state
    with a payload and a status (``RecallProposal`` is payload-free and
    TTL-expiring, and cannot be retrofitted).

    ``status``, ``ask_count``, ``delivered_turn``, ``answer_option``,
    ``cooldown_until`` and ``resolved_turn`` are deliberately **unindexed**
    plain fields: delivery and :func:`record_answer` write them with Lua
    directly on the model hash, which an index would not see. Reads go through the indexed ``agent_id``.

    Fields:
        question_text: The question, phrased by the producer (or by the host
            later). Echoes content, so retention is bounded by ``prune()``.
        kind: One of :data:`KINDS`.
        source_module: Free-form producer name (e.g. ``"reconciliation"``).
        target_keys: Redis keys of the facts the answer would resolve.
        options: ``[{"label": str, "acted": [keys], "contradicted": [keys]}]``.
            Answering applies exactly the chosen option's two lists.
        ambiguity_signal: One of :data:`AMBIGUITY_SIGNALS`.
        status: One of :data:`STATUSES`.
        ask_count: Times delivered. Re-asking after cooldown is legal.
        created_turn / expires_turn: Proposal turn and ``created + N``.
        last_seen_turn: Most recent turn the ambiguity was re-observed or the
            targets were used -- the VOI gate's impact factor.
        cooldown_until: Not deliverable before this turn (after a
            deflected/unrecognized reply).
        answer_option: 0-based index of the chosen option. The free-text reply
            itself is never stored.
        resolved_turn: Turn the candidate left the open set (answered, cooled,
            expired), used as the retention reference by ``prune()``.
        cue_tokens: See :func:`cue_tokens_for`.
        delivered_turn: Turn of the most recent delivery, written by the
            delivery Lua. A ``delivered`` candidate that gets no reply within
            ``QUESTION_COOLDOWN_TURNS`` of it is treated like a non-answer and
            moved to ``cooled`` (see :func:`expire_stale`).
        disjunction_id: Set by the M5 producer so a retracted disjoin can
            expire its candidate via a registered staleness check.
        agent_id: Owning agent.
    """

    candidate_id = KeyField()
    agent_id = IndexedField(type=str)
    question_text = StringField(default="")
    kind = StringField(default="")
    source_module = StringField(default="")
    target_keys = ListField(default=[])
    options = ListField(default=[])
    ambiguity_signal = StringField(default="")
    status = StringField(default="pending")
    ask_count = IntField(default=0)
    created_turn = IntField(default=0)
    expires_turn = IntField(default=0)
    last_seen_turn = IntField(default=0)
    cooldown_until = IntField(null=True)
    answer_option = IntField(null=True)
    resolved_turn = IntField(null=True)
    cue_tokens = ListField(default=[])
    delivered_turn = IntField(null=True)
    disjunction_id = StringField(null=True)

AnswerResult dataclass

Outcome of :func:record_answer.

Attributes:

Name Type Description
applied bool

True only when an answered reply was claimed and its evidence was written.

reason str

"applied", "cooled" (deflected/unrecognized, no evidence), "not_open" (the claim failed: already answered, expired, or -- for a reply arriving after the question was ignored for QUESTION_COOLDOWN_TURNS -- already moved to cooled), "apply_failed" (claimed, evidence write failed -- the answer is lost, never double-counted), "disabled" (kill switch) or "error" (Redis error before the claim).

classification Optional[str]

"answered" / "deflected" / "unrecognized".

option_index Optional[int]

0-based chosen option for an answered reply.

Source code in src/popoto/recipes/question_queue.py
@dataclass
class AnswerResult:
    """Outcome of :func:`record_answer`.

    Attributes:
        applied: True only when an ``answered`` reply was claimed *and* its
            evidence was written.
        reason: ``"applied"``, ``"cooled"`` (deflected/unrecognized, no
            evidence), ``"not_open"`` (the claim failed: already answered,
            expired, or -- for a reply arriving after the question was
            ignored for ``QUESTION_COOLDOWN_TURNS`` -- already moved to
            ``cooled``), ``"apply_failed"`` (claimed, evidence write
            failed -- the answer is lost, never double-counted),
            ``"disabled"`` (kill switch) or ``"error"`` (Redis error before
            the claim).
        classification: ``"answered"`` / ``"deflected"`` / ``"unrecognized"``.
        option_index: 0-based chosen option for an ``answered`` reply.
    """

    applied: bool
    reason: str
    classification: Optional[str] = None
    option_index: Optional[int] = None

normalize_text(text)

Lowercase, strip punctuation, collapse whitespace.

The single normalization used for dedup, option matching and the deflection set, so the three can never disagree about what "the same text" means.

Source code in src/popoto/recipes/question_queue.py
def normalize_text(text: str) -> str:
    """Lowercase, strip punctuation, collapse whitespace.

    The single normalization used for dedup, option matching and the
    deflection set, so the three can never disagree about what "the same
    text" means.
    """
    lowered = str(text).lower()
    stripped = lowered.translate(str.maketrans("", "", string.punctuation))
    return " ".join(stripped.split())

cue_tokens_for(*texts)

Relevance-timing tokens: lowercased word tokens of length >= 3 from texts, minus :data:_CUE_STOPWORDS, sorted and de-duplicated.

Source code in src/popoto/recipes/question_queue.py
def cue_tokens_for(*texts: str) -> List[str]:
    """Relevance-timing tokens: lowercased word tokens of length >= 3 from
    ``texts``, minus :data:`_CUE_STOPWORDS`, sorted and de-duplicated."""
    tokens = set()
    for text in texts:
        for tok in _TOKEN_RE.findall(str(text).lower()):
            if len(tok) >= 3 and tok not in _CUE_STOPWORDS:
                tokens.add(tok)
    return sorted(tokens)

classify_answer(options, answer_text)

Deterministic three-way answer recognition.

After :func:normalize_text, in order:

  1. In :data:DEFLECTION_PHRASES -> ("deflected", None).
  2. Equal to exactly one option's normalized label, or its 1-based index -> ("answered", i) with i 0-based.
  3. Anything else (including a reply matching more than one option) -> ("unrecognized", None).
Source code in src/popoto/recipes/question_queue.py
def classify_answer(
    options: Sequence[Dict[str, Any]], answer_text: str
) -> Tuple[str, Optional[int]]:
    """Deterministic three-way answer recognition.

    After :func:`normalize_text`, in order:

    1. In :data:`DEFLECTION_PHRASES` -> ``("deflected", None)``.
    2. Equal to exactly one option's normalized label, or its 1-based index
       -> ``("answered", i)`` with ``i`` 0-based.
    3. Anything else (including a reply matching more than one option)
       -> ``("unrecognized", None)``.
    """
    reply = normalize_text(answer_text)
    if not reply or reply in DEFLECTION_PHRASES:
        return (DEFLECTED if reply else UNRECOGNIZED), None
    matches = {
        i
        for i, option in enumerate(options)
        if normalize_text(option.get("label", "")) == reply
    }
    if reply.isdigit():
        idx = int(reply) - 1
        if 0 <= idx < len(options):
            matches.add(idx)
    if len(matches) == 1:
        return ANSWERED, matches.pop()
    return UNRECOGNIZED, None

passes_voi_gate(candidate, turn)

The value-of-information gate: a deterministic boolean.

ambiguity_signal in AMBIGUITY_SIGNALS and recently_used, where recently used means turn - last_seen_turn <= QUESTION_RECENT_USE_TURNS. Both factors are stored on the candidate, so every delivered question is traceable to the evidence that let it through.

Source code in src/popoto/recipes/question_queue.py
def passes_voi_gate(candidate: QuestionCandidate, turn: int) -> bool:
    """The value-of-information gate: a deterministic boolean.

    ``ambiguity_signal in AMBIGUITY_SIGNALS and recently_used``, where
    recently used means ``turn - last_seen_turn <= QUESTION_RECENT_USE_TURNS``.
    Both factors are stored on the candidate, so every delivered question is
    traceable to the evidence that let it through.
    """
    if _str_attr(candidate, "ambiguity_signal") not in AMBIGUITY_SIGNALS:
        return False
    last_seen = _int_attr(candidate, "last_seen_turn")
    if last_seen is None:
        return False
    return 0 <= int(turn) - last_seen <= QUESTION_RECENT_USE_TURNS

register_staleness_check(kind, check)

Register check for candidates of kind (idempotent).

Used by producers whose source can vanish between proposal and delivery -- e.g. the M5 disjunction producer expiring a candidate whose disjunction_id no longer appears among the live disjoins.

Source code in src/popoto/recipes/question_queue.py
def register_staleness_check(
    kind: str, check: Callable[[QuestionCandidate, int], bool]
) -> None:
    """Register ``check`` for candidates of ``kind`` (idempotent).

    Used by producers whose source can vanish between proposal and delivery
    -- e.g. the M5 disjunction producer expiring a candidate whose
    ``disjunction_id`` no longer appears among the live disjoins.
    """
    if kind not in KINDS:
        raise ValueError(f"unknown kind {kind!r}; expected one of {KINDS}")
    checks = _STALENESS_CHECKS.setdefault(kind, [])
    if check not in checks:
        checks.append(check)

unregister_staleness_check(kind, check)

Remove a previously registered check (no-op if absent).

Source code in src/popoto/recipes/question_queue.py
def unregister_staleness_check(
    kind: str, check: Callable[[QuestionCandidate, int], bool]
) -> None:
    """Remove a previously registered check (no-op if absent)."""
    checks = _STALENESS_CHECKS.get(kind, [])
    if check in checks:
        checks.remove(check)

expire_stale(agent_id, turn)

Silently expire the agent's not-yet-answered candidates that are past expires_turn or that a registered staleness check reports stale, and move delivered candidates ignored for QUESTION_COOLDOWN_TURNS to cooled.

Expired candidates are status-marked and retained (prune() deletes). Returns the number expired; 0 on a Redis error (fail closed).

Source code in src/popoto/recipes/question_queue.py
def expire_stale(agent_id: str, turn: int) -> int:
    """Silently expire the agent's not-yet-answered candidates that are past
    ``expires_turn`` or that a registered staleness check reports stale, and
    move ``delivered`` candidates ignored for ``QUESTION_COOLDOWN_TURNS`` to
    ``cooled``.

    Expired candidates are status-marked and retained (``prune()`` deletes).
    Returns the number expired; ``0`` on a Redis error (fail closed).
    """
    try:
        expired, _ = _expire_stale_in(_candidates(agent_id), turn)
        return expired
    except Exception:
        logger.exception("question_queue: expire_stale failed")
        return 0

propose(agent_id, question_text, kind, source_module, target_keys, ambiguity_signal, turn, options=None, disjunction_id=None)

Write a question candidate -- producers never ask directly.

Dedup: an incoming proposal is a duplicate of an existing candidate that is open (pending/delivered/cooled) or answered within QUESTION_RETENTION_TURNS when kind matches and the target_keys sets intersect, or when the normalized question text matches exactly. A duplicate is not re-created: an open duplicate is touched (last_seen_turn bumped, feeding the gate's impact factor) and returned; an answered one is returned unchanged.

The dedup check and the create run under a short per-agent SET NX lock, so concurrent proposals of one question yield one candidate. If the lock cannot be taken within a brief retry window the proposal is dropped (fail closed): a producer re-observes its ambiguity on a later turn.

Returns:

Type Description
Optional[QuestionCandidate]

The new or existing candidate; None when the kill switch is set,

Optional[QuestionCandidate]

the propose lock stays busy, or on a Redis error.

Raises:

Type Description
ValueError

unknown kind / ambiguity_signal, an empty question, or an option listing a key as both acted and contradicted.

Source code in src/popoto/recipes/question_queue.py
def propose(
    agent_id: str,
    question_text: str,
    kind: str,
    source_module: str,
    target_keys: Sequence[str],
    ambiguity_signal: str,
    turn: int,
    options: Optional[Sequence[Dict[str, Any]]] = None,
    disjunction_id: Optional[str] = None,
) -> Optional[QuestionCandidate]:
    """Write a question candidate -- producers never ask directly.

    Dedup: an incoming proposal is a duplicate of an existing candidate that
    is open (pending/delivered/cooled) or ``answered`` within
    ``QUESTION_RETENTION_TURNS`` when ``kind`` matches and the
    ``target_keys`` sets intersect, **or** when the normalized question text
    matches exactly. A duplicate is not re-created: an open duplicate is
    *touched* (``last_seen_turn`` bumped, feeding the gate's impact factor)
    and returned; an answered one is returned unchanged.

    The dedup check and the create run under a short per-agent ``SET NX``
    lock, so concurrent proposals of one question yield one candidate. If the
    lock cannot be taken within a brief retry window the proposal is dropped
    (fail closed): a producer re-observes its ambiguity on a later turn.

    Returns:
        The new or existing candidate; ``None`` when the kill switch is set,
        the propose lock stays busy, or on a Redis error.

    Raises:
        ValueError: unknown ``kind`` / ``ambiguity_signal``, an empty question,
            or an option listing a key as both acted and contradicted.
    """
    if kind not in KINDS:
        raise ValueError(f"unknown kind {kind!r}; expected one of {KINDS}")
    if ambiguity_signal not in AMBIGUITY_SIGNALS:
        raise ValueError(
            f"unknown ambiguity_signal {ambiguity_signal!r}; "
            f"expected one of {AMBIGUITY_SIGNALS}"
        )
    norm_text = normalize_text(question_text)
    if not norm_text:
        raise ValueError("question_text must not be empty")
    cleaned_options = _validate_options(options)
    keys = [str(k) for k in target_keys or []]
    turn = int(turn)

    if not question_queue_enabled():
        return None

    try:
        token = _acquire_propose_lock(agent_id)
    except Exception:
        logger.exception("question_queue: propose lock failed")
        return None
    if token is None:
        logger.warning(
            "question_queue: propose lock for %s busy; proposal dropped", agent_id
        )
        return None
    try:
        return _propose_locked(
            agent_id,
            question_text,
            norm_text,
            kind,
            source_module,
            keys,
            ambiguity_signal,
            turn,
            cleaned_options,
            disjunction_id,
        )
    except Exception:
        logger.exception("question_queue: propose failed")
        return None
    finally:
        _release_propose_lock(agent_id, token)

note_use(agent_id, keys, turn)

Record that the host just injected keys into context at turn.

Bumps last_seen_turn (never backwards) on every open candidate whose target_keys intersect keys -- the "recently used" half of the VOI gate. Returns the number of candidates touched; 0 when disabled or on a Redis error.

Source code in src/popoto/recipes/question_queue.py
def note_use(agent_id: str, keys: Iterable[str], turn: int) -> int:
    """Record that the host just injected ``keys`` into context at ``turn``.

    Bumps ``last_seen_turn`` (never backwards) on every open candidate whose
    ``target_keys`` intersect ``keys`` -- the "recently used" half of the VOI
    gate. Returns the number of candidates touched; ``0`` when disabled or on
    a Redis error.
    """
    if not question_queue_enabled():
        return 0
    key_set = {str(k) for k in keys or []}
    if not key_set:
        return 0
    touched = 0
    try:
        for cand in _candidates(agent_id):
            if cand.status not in _OPEN_STATUSES:
                continue
            if not key_set & set(_list_attr(cand, "target_keys")):
                continue
            if int(turn) > (_int_attr(cand, "last_seen_turn") or 0):
                setattr(cand, "last_seen_turn", int(turn))
                cand.save(update_fields=["last_seen_turn"])
                touched += 1
    except Exception:
        logger.exception("question_queue: note_use failed")
    return touched

next_question(agent_id, turn, query_cues=None)

The whole delivery contract: the next question to ask, if any.

Pipeline, in order: expiry and registered staleness checks (:func:expire_stale); cooldown; the VOI gate (:func:passes_voi_gate); relevance timing (query_cues must share at least one cue token with the candidate; None disables timing); ordering by impact; and only then one atomic script that grants the token bucket and claims the first still-deliverable candidate (:data:_DELIVER_LUA), so the budget is never spent when nothing is actually delivered. The delivered candidate's ask_count is incremented and delivered_turn recorded.

Returns None -- never raises -- when disabled, when nothing passes, when the budget is exhausted (the bucket is not reset), or on a Redis error.

Source code in src/popoto/recipes/question_queue.py
def next_question(
    agent_id: str,
    turn: int,
    query_cues: Optional[Union[str, Iterable[str]]] = None,
) -> Optional[QuestionCandidate]:
    """The whole delivery contract: the next question to ask, if any.

    Pipeline, in order: expiry and registered staleness checks
    (:func:`expire_stale`); cooldown; the VOI gate (:func:`passes_voi_gate`);
    relevance timing (``query_cues`` must share at least one cue token with
    the candidate; ``None`` disables timing); ordering by impact; and only
    then one atomic script that grants the token bucket *and* claims the
    first still-deliverable candidate (:data:`_DELIVER_LUA`), so the budget is
    never spent when nothing is actually delivered. The delivered candidate's
    ``ask_count`` is incremented and ``delivered_turn`` recorded.

    Returns ``None`` -- never raises -- when disabled, when nothing passes,
    when the budget is exhausted (the bucket is not reset), or on a Redis
    error.
    """
    if not question_queue_enabled():
        return None
    turn = int(turn)
    try:
        _, survivors = _expire_stale_in(_candidates(agent_id), turn)
        cue_set = None if query_cues is None else _query_tokens(query_cues)
        eligible = []
        for cand in survivors:
            cooldown = _int_attr(cand, "cooldown_until")
            if cooldown is not None and turn < cooldown:
                continue
            if not passes_voi_gate(cand, turn):
                continue
            if cue_set is not None and not cue_set & set(
                _list_attr(cand, "cue_tokens")
            ):
                continue
            eligible.append(cand)
        if not eligible:
            return None
        eligible.sort(key=_impact_order)
        return _try_deliver(agent_id, turn, eligible)
    except Exception:
        logger.exception("question_queue: next_question failed")
        return None

record_answer(candidate, answer_text, turn, instances=None)

Record the person's reply to a delivered question.

The reply is classified by :func:classify_answer and then:

  • answered -- claim, then apply. A Lua compare-and-set flips status from delivered/pending to answered and writes answer_option. Only on a successful claim are the chosen option's effects applied (:func:_apply_option_effects). A crash between the two loses this answer's evidence; it never double-counts, because a second call fails the claim (reason="not_open"). Only delivered and pending are claimable: a late reply to a question the expiry pass already cooled as ignored also gets not_open and writes nothing -- the question is re-asked after its cooldown instead.
  • deflected / unrecognized -- the candidate becomes cooled with cooldown_until = turn + QUESTION_COOLDOWN_TURNS; no evidence is written, so stored confidence is bit-identical.

answer_text is never stored or logged.

Parameters:

Name Type Description Default
candidate QuestionCandidate

The candidate returned by :func:next_question.

required
answer_text str

The person's free-text reply.

required
turn int

Host turn of the reply.

required
instances Optional[Sequence[Model]]

Optional target instances, used only to tell which model class each target key belongs to; fresh copies are always loaded.

None

Returns:

Type Description
AnswerResult

class:AnswerResult. Never raises.

Source code in src/popoto/recipes/question_queue.py
def record_answer(
    candidate: QuestionCandidate,
    answer_text: str,
    turn: int,
    instances: Optional[Sequence[Model]] = None,
) -> AnswerResult:
    """Record the person's reply to a delivered question.

    The reply is classified by :func:`classify_answer` and then:

    * ``answered`` -- **claim, then apply.** A Lua compare-and-set flips
      ``status`` from ``delivered``/``pending`` to ``answered`` and writes
      ``answer_option``. Only on a successful claim are the chosen option's
      effects applied (:func:`_apply_option_effects`). A crash between the two
      loses this answer's evidence; it never double-counts, because a second
      call fails the claim (``reason="not_open"``). Only ``delivered`` and
      ``pending`` are claimable: a late reply to a question the expiry pass
      already cooled as ignored also gets ``not_open`` and writes nothing --
      the question is re-asked after its cooldown instead.
    * ``deflected`` / ``unrecognized`` -- the candidate becomes ``cooled``
      with ``cooldown_until = turn + QUESTION_COOLDOWN_TURNS``; no evidence
      is written, so stored confidence is bit-identical.

    ``answer_text`` is never stored or logged.

    Args:
        candidate: The candidate returned by :func:`next_question`.
        answer_text: The person's free-text reply.
        turn: Host turn of the reply.
        instances: Optional target instances, used only to tell which model
            class each target key belongs to; fresh copies are always loaded.

    Returns:
        :class:`AnswerResult`. Never raises.
    """
    if not question_queue_enabled():
        return AnswerResult(applied=False, reason="disabled")
    turn = int(turn)
    options: List[Dict[str, Any]] = _list_attr(candidate, "options")
    classification, option_index = classify_answer(options, answer_text)
    try:
        if classification != ANSWERED:
            claimed = _claim(
                candidate,
                _ANSWERABLE_STATUSES,
                {
                    "status": "cooled",
                    "cooldown_until": turn + QUESTION_COOLDOWN_TURNS,
                    "resolved_turn": turn,
                },
            )
            return AnswerResult(
                applied=False,
                reason="cooled" if claimed else "not_open",
                classification=classification,
            )
        claimed = _claim(
            candidate,
            _ANSWERABLE_STATUSES,
            {
                "status": "answered",
                "answer_option": option_index,
                "resolved_turn": turn,
            },
        )
    except Exception:
        logger.exception("question_queue: record_answer claim failed")
        return AnswerResult(
            applied=False, reason="error", classification=classification
        )
    if not claimed:
        return AnswerResult(
            applied=False,
            reason="not_open",
            classification=classification,
            option_index=option_index,
        )
    assert option_index is not None  # classify_answer guarantees it on ANSWERED
    try:
        option = options[option_index]
        keys = list(option.get("acted", []) or []) + list(
            option.get("contradicted", []) or []
        )
        _apply_option_effects(option, _resolve_targets(keys, instances))
    except Exception:
        logger.exception(
            "question_queue: answer evidence for %s lost after claim",
            candidate.candidate_id,
        )
        return AnswerResult(
            applied=False,
            reason="apply_failed",
            classification=classification,
            option_index=option_index,
        )
    return AnswerResult(
        applied=True,
        reason="applied",
        classification=classification,
        option_index=option_index,
    )

prune(agent_id, current_turn)

Delete the agent's non-pending candidates older than QUESTION_RETENTION_TURNS (measured from resolved_turn, else last_seen_turn). Pending candidates are never pruned here -- they expire first. Returns the number deleted; 0 on a Redis error.

Source code in src/popoto/recipes/question_queue.py
def prune(agent_id: str, current_turn: int) -> int:
    """Delete the agent's non-pending candidates older than
    ``QUESTION_RETENTION_TURNS`` (measured from ``resolved_turn``, else
    ``last_seen_turn``). Pending candidates are never pruned here -- they
    expire first. Returns the number deleted; ``0`` on a Redis error."""
    deleted = 0
    try:
        for cand in _candidates(agent_id):
            if cand.status == "pending":
                continue
            age = int(current_turn) - _retention_reference(cand)
            if age > QUESTION_RETENTION_TURNS:
                cand.delete()
                deleted += 1
    except Exception:
        logger.exception("question_queue: prune failed")
    return deleted

propose_from_disjunctions(agent_id, turn)

Propose one disjunction question per live M5 disjunct pair.

Reads reconciliation.merge_log_entries(agent_id) filtered to kind="disjoin" and decodes each payload for the pair's two sides. Options are exclusive by construction: A = acted:[a], contradicted:[b], B = the mirror. disjunction_id is set so the registered staleness check expires the candidate if the annotation is later retracted. Calling this again re-observes the same pairs, which dedup folds into a touch (bumping last_seen_turn).

The answer is recorded on the candidate (answer_option); the JournalEntry sides carry no ConfidenceField, so nothing about them changes. Feeding the answer back so the disjunction resolves is out of scope (an M1/M5 change).

Returns:

Type Description
int

The number of pairs proposed or touched; 0 when M5 is unavailable,

int

the kill switch is set, or on any error. Never raises.

Source code in src/popoto/recipes/question_queue.py
def propose_from_disjunctions(agent_id: str, turn: int) -> int:
    """Propose one ``disjunction`` question per live M5 disjunct pair.

    Reads ``reconciliation.merge_log_entries(agent_id)`` filtered to
    ``kind="disjoin"`` and decodes each payload for the pair's two sides.
    Options are exclusive by construction: A = ``acted:[a],
    contradicted:[b]``, B = the mirror. ``disjunction_id`` is set so the
    registered staleness check expires the candidate if the annotation is
    later retracted. Calling this again re-observes the same pairs, which
    dedup folds into a touch (bumping ``last_seen_turn``).

    The answer is recorded on the candidate (``answer_option``); the
    ``JournalEntry`` sides carry no ConfidenceField, so nothing about them
    changes. Feeding the answer back so the disjunction resolves is out of
    scope (an M1/M5 change).

    Returns:
        The number of pairs proposed or touched; ``0`` when M5 is unavailable,
        the kill switch is set, or on any error. Never raises.
    """
    if not question_queue_enabled():
        return 0
    try:
        disjoins = _live_disjoins(agent_id)
    except ImportError:
        logger.info("question_queue: M5 reconciliation unavailable; no disjoins")
        return 0
    except Exception:
        logger.exception("question_queue: reading M5 disjoins failed")
        return 0
    proposed = 0
    for disjunction_id, side_a, side_b in disjoins:
        try:
            label_a, label_b = _distinct_labels(
                [_entry_statement(side_a), _entry_statement(side_b)]
            )
            candidate = propose(
                agent_id=agent_id,
                question_text=(
                    f"Two things I hold conflict: {label_a!r} or {label_b!r}. "
                    "Which is right?"
                ),
                kind="disjunction",
                source_module="reconciliation",
                target_keys=[side_a, side_b],
                ambiguity_signal="disjunction",
                turn=turn,
                options=[
                    {"label": label_a, "acted": [side_a], "contradicted": [side_b]},
                    {"label": label_b, "acted": [side_b], "contradicted": [side_a]},
                ],
                disjunction_id=disjunction_id,
            )
        except Exception:
            logger.exception("question_queue: disjunction proposal failed")
            continue
        proposed += int(candidate is not None)
    return proposed

propose_from_gate(agent_id, metadata, turn, query_text)

Propose a confirmation question from a confidence-gate refusal.

metadata is assemble()'s metadata dict (an AssemblyResult is accepted too). Acts only when metadata["gate"] reports gated, mode == "refuse" and a non-empty refused_keys; k is the top (rank-0) refused key, the one whose confidence fell below threshold. Options: "yes" = acted:[k], "no" = contradicted:[k].

Only "refuse" mode produces a question. The assembler reports refused_keys in "flag" mode too, but there the records were injected into context anyway: nothing was withheld, so there is no refusal to clarify, and asking would spend the shared budget on a fact the agent is already using.

The question text names k as well as the query, so two refusals of different facts under the same query stay two candidates instead of being folded together by normalized-text dedup.

Dormant unless confidence_gate_threshold is configured on the ContextAssembler: without it assemble() emits no gate metadata and this producer has nothing to read.

Returns:

Type Description
Optional[QuestionCandidate]

The new or existing candidate; None when there was no refusal,

Optional[QuestionCandidate]

when disabled, or on any error. Never raises.

Source code in src/popoto/recipes/question_queue.py
def propose_from_gate(
    agent_id: str,
    metadata: Any,
    turn: int,
    query_text: str,
) -> Optional[QuestionCandidate]:
    """Propose a ``confirmation`` question from a confidence-gate refusal.

    ``metadata`` is ``assemble()``'s metadata dict (an ``AssemblyResult`` is
    accepted too). Acts only when ``metadata["gate"]`` reports ``gated``,
    ``mode == "refuse"`` and a non-empty ``refused_keys``; ``k`` is the top
    (rank-0) refused key, the one whose confidence fell below threshold.
    Options: ``"yes"`` = ``acted:[k]``, ``"no"`` = ``contradicted:[k]``.

    **Only ``"refuse"`` mode produces a question.** The assembler reports
    ``refused_keys`` in ``"flag"`` mode too, but there the records were
    injected into context anyway: nothing was withheld, so there is no
    refusal to clarify, and asking would spend the shared budget on a fact
    the agent is already using.

    The question text names ``k`` as well as the query, so two refusals of
    different facts under the same query stay two candidates instead of
    being folded together by normalized-text dedup.

    **Dormant unless ``confidence_gate_threshold`` is configured** on the
    ``ContextAssembler``: without it ``assemble()`` emits no gate metadata
    and this producer has nothing to read.

    Returns:
        The new or existing candidate; ``None`` when there was no refusal,
        when disabled, or on any error. Never raises.
    """
    try:
        if not isinstance(metadata, dict):
            metadata = getattr(metadata, "metadata", None)
        gate = metadata.get("gate") if isinstance(metadata, dict) else None
        if not isinstance(gate, dict) or not gate.get("gated"):
            return None
        if gate.get("mode") != "refuse":
            return None
        refused = [str(k) for k in gate.get("refused_keys") or [] if k]
        if not refused:
            return None
        key = refused[0]
        topic = str(query_text or "").strip() or "this topic"
        # NOTE: the raw Redis key in the person-facing text (and the cue tokens
        # its fragments add) -- left as-is because propose() dedups on kind +
        # target_keys intersection OR exact normalized text. Without the key,
        # two refusals of different facts under the same query would match on
        # text alone and collapse into one candidate. Removing it needs that
        # text-match rule to also require compatible target_keys, which
        # changes dedup for every producer; the host phrases question_text
        # before delivery anyway. The topic still supplies the cue tokens.
        return propose(
            agent_id=agent_id,
            question_text=(
                f"I'm not confident in what I recall about {topic!r} "
                f"({key}). Is it still accurate?"
            ),
            kind="confirmation",
            source_module="context_assembler.confidence_gate",
            target_keys=[key],
            ambiguity_signal="gate_refusal",
            turn=turn,
            options=[
                {"label": "yes", "acted": [key]},
                {"label": "no", "contradicted": [key]},
            ],
        )
    except Exception:
        logger.exception("question_queue: gate proposal failed")
        return None

propose_from_resolution(record, turn)

Propose one referent question per M4 evidence_gap reference.

Reads record.references_json (the shape _serialise_reference in extraction/resolution_log.py writes) and, for each entry whose status == "evidence_gap", proposes its question with one option per entry in candidates. Options carry empty acted / contradicted lists: the answer is recorded (answer_option), never applied as evidence.

Each reference gets its own target key, "{record_key}#ref:{start}:{end}", so two gaps in one sentence are two questions rather than one swallowed by key-intersection dedup.

Returns:

Type Description
int

The number of references proposed or touched; 0 when there are

int

none, when disabled, or on any error. Never raises.

Source code in src/popoto/recipes/question_queue.py
def propose_from_resolution(record: Any, turn: int) -> int:
    """Propose one ``referent`` question per M4 ``evidence_gap`` reference.

    Reads ``record.references_json`` (the shape ``_serialise_reference`` in
    ``extraction/resolution_log.py`` writes) and, for each entry whose
    ``status == "evidence_gap"``, proposes its ``question`` with one option
    per entry in ``candidates``. Options carry empty ``acted`` /
    ``contradicted`` lists: the answer is recorded (``answer_option``), never
    applied as evidence.

    Each reference gets its own target key,
    ``"{record_key}#ref:{start}:{end}"``, so two gaps in one sentence are two
    questions rather than one swallowed by key-intersection dedup.

    Returns:
        The number of references proposed or touched; ``0`` when there are
        none, when disabled, or on any error. Never raises.
    """
    if not question_queue_enabled():
        return 0
    try:
        agent_id = str(getattr(record, "agent_id", "") or "")
        record_key = str(record.db_key.redis_key)
        references = json.loads(str(getattr(record, "references_json", "") or "[]"))
    except Exception:
        logger.exception("question_queue: unreadable ResolutionRecord")
        return 0
    if not agent_id or not isinstance(references, list):
        return 0
    proposed = 0
    for ref in references:
        try:
            if not isinstance(ref, dict) or ref.get("status") != "evidence_gap":
                continue
            seen: Dict[str, str] = {}
            for cand in ref.get("candidates") or []:
                text = str(cand)
                if normalize_text(text) and normalize_text(text) not in seen:
                    seen[normalize_text(text)] = text
            if not seen:
                continue
            surface = str(ref.get("surface") or "")
            question = str(ref.get("question") or "").strip() or (
                f"Who or what does {surface!r} refer to?"
            )
            candidate = propose(
                agent_id=agent_id,
                question_text=question,
                kind="referent",
                source_module="extraction.resolution",
                target_keys=[f"{record_key}#ref:{ref.get('start')}:{ref.get('end')}"],
                ambiguity_signal="evidence_gap",
                turn=turn,
                options=[
                    {"label": label, "acted": [], "contradicted": []}
                    for label in _distinct_labels(list(seen.values()))
                ],
            )
        except Exception:
            logger.exception("question_queue: evidence_gap proposal failed")
            continue
        proposed += int(candidate is not None)
    return proposed