"""The reply moment: the finished reply, checked before it goes out (milestone 458 step 4, folded in from milestone 456 step 8). The operator's ruling: an unopened rule above the stop bar HOLDS the reply once, in the server's words — the backstop for whatever the earlier arms missed — and the rewrite is never held. These pin, with the database stubbed: - which moments a reply reaches; - what holds: a mounted RULE the session has not opened, or the reply's top semantic hit above the reply surface's bar — never a preference, never a rule already opened or already held once, never past the session cap; - that the stop-only arm records surfacing for the held rule alone; - the hook's contract: blocks only on the server's reason, records the hold before emitting it, and never holds the rewrite. """ from __future__ import annotations import json import os import subprocess from pathlib import Path from unittest.mock import AsyncMock, MagicMock, patch import pytest from quart import Quart, g from scribe.services import moment_delivery as md from scribe.services import plugin_context as pc from scribe.services import retrieval_pipeline as rp from scribe.services import retrieval_surfaces as surfaces from scribe.services import rule_usage, rulebooks from tests.helpers import fake_rule, http_sink, need_tools ROOT = Path(__file__).resolve().parents[1] HOOK = ROOT / "plugin" / "hooks" / "scribe_reply_check.sh" DONE = fake_rule(id=11, title="Definition of done", kind="rule", when_to_apply="about to tell the operator work is finished") # ── which moments a reply reaches ─────────────────────────────────────── def test_every_reply_reports_and_a_question_also_asks(): assert [h["moment"] for h in md.reply_moments("Shipped it.")] == ["reply.report"] assert [h["moment"] for h in md.reply_moments("Shipped.\nShall I merge it?")] == [ "reply.report", "reply.ask", ] def test_a_question_inside_code_is_not_the_reply_asking(): reply = "Done.\n```python\nok = x if y else z?\n```\n" assert [h["moment"] for h in md.reply_moments(reply)] == ["reply.report"] def test_a_long_reply_keeps_its_end(): """The part of a report that asks something of the reader is at its end.""" reply = "head " * 400 + "please confirm it works on your end" query = md._reply_query(reply) assert query.startswith("head") assert query.endswith("please confirm it works on your end") assert len(query) < len(reply) # ── what holds ────────────────────────────────────────────────────────── def _stubs(*, mounted=(), checkpoint=None): arm = AsyncMock(return_value=rp.RuleResult(checkpoint=checkpoint or {})) surfaced = MagicMock() lookup = AsyncMock(return_value=list(mounted)) return arm, surfaced, lookup, [ patch.object(rulebooks, "rules_on_moments", lookup), patch.object(rule_usage, "record_rule_surfaced", surfaced), patch.object(rp, "run_rule_arm", arm), patch.object(surfaces, "floor_for", AsyncMock(return_value=0.8)), patch.object(surfaces, "budget_for", AsyncMock(return_value=1)), patch.object(pc, "_rule_io", MagicMock()), ] async def _hold(stack, reply="All done — let me know if it works.", **kw): for p in stack: p.start() try: return await md.reply_hold(1, reply, project_id=2, **kw) finally: for p in reversed(stack): p.stop() async def test_a_mounted_rule_the_session_has_not_opened_holds_the_reply(): _arm, surfaced, lookup, stack = _stubs(mounted=[(DONE, "reply.report")]) out = await _hold(stack) assert out["rule_ids"] == [11] assert "Definition of done" in out["reason"] and "get_rule(11)" in out["reason"] assert "mounted on reply.report" in out["reason"] assert lookup.await_args.args == (1, ["reply.report"], 2) kw = surfaced.call_args.kwargs assert kw["source"] == rp.MOMENT_RULE_SOURCE and kw["detail"] == {11: "reply.report"} async def test_an_opened_rule_a_preference_and_an_earlier_hold_do_not_hold(): pref = fake_rule(id=12, title="Lead with the outcome", kind="preference") other = fake_rule(id=13, title="Name the commit", kind="rule") mounted = [(DONE, "reply.report"), (pref, "reply.report"), (other, "reply.report")] _arm, surfaced, _lookup, stack = _stubs(mounted=mounted) assert await _hold(stack, held=frozenset({11}), stopped=frozenset({13})) == {} surfaced.assert_not_called() async def test_the_semantic_backstop_holds_and_never_names_a_mounted_rule_twice(): cp = {"rule_id": 40, "title": "Read the job log first", "score": 0.84, "trigger": "a CI run overran"} arm, _surfaced, _lookup, stack = _stubs(mounted=[(DONE, "reply.report")], checkpoint=cp) out = await _hold(stack) assert out["rule_ids"] == [11, 40] assert "scores 0.84 against this reply" in out["reason"] spec, moment = arm.await_args.args assert spec is rp.REPLY_RULE assert 11 in moment.held, "a rule the mounted half holds must not be raised again" # Its floor IS its stop bar. assert arm.await_args.kwargs["floor"] == arm.await_args.kwargs["checkpoint_floor"] == 0.8 async def test_past_the_session_cap_nothing_holds_and_nothing_is_asked(): arm, _surfaced, lookup, stack = _stubs(mounted=[(DONE, "reply.report")]) stopped = frozenset(range(100, 100 + pc.CHECKPOINT_SESSION_CAP)) assert await _hold(stack, stopped=stopped) == {} lookup.assert_not_awaited() arm.assert_not_awaited() async def test_the_cap_trims_before_anything_is_recorded(): rules = [(fake_rule(id=i, title=f"rule {i}", kind="rule"), "reply.report") for i in range(1, 9)] _arm, surfaced, _lookup, stack = _stubs(mounted=rules) out = await _hold(stack, stopped=frozenset({50, 51})) room = pc.CHECKPOINT_SESSION_CAP - 2 assert len(out["rule_ids"]) == room assert surfaced.call_args.kwargs["rule_ids"] == out["rule_ids"] async def test_the_reply_check_fails_open(): _arm, _surfaced, _lookup, stack = _stubs() stack[0] = patch.object(rulebooks, "rules_on_moments", AsyncMock(side_effect=RuntimeError("x"))) assert await _hold(stack) == {} assert await md.reply_hold(1, " ") == {} # ── the stop-only arm ─────────────────────────────────────────────────── async def test_the_stop_only_arm_records_only_the_rule_that_holds(): hits = [(0.86, fake_rule(id=40, title="a", kind="rule", when_to_apply="t")), (0.82, fake_rule(id=41, title="b", kind="rule", when_to_apply="t"))] io = rp.RuleIO(search=AsyncMock(return_value=hits), record_retrieval=MagicMock(), record_rule_surfaced=MagicMock()) moment = rp.RuleMoment(user_id=1, query="the reply", project_id=2, checkpoint_where="this reply") result = await rp.run_rule_arm(rp.REPLY_RULE, moment, floor=0.8, budget=3, io=io, checkpoint_floor=0.8) assert result.checkpoint["rule_id"] == 40 assert io.record_rule_surfaced.call_args.kwargs["rule_ids"] == [40] # The top hit already opened: nothing holds, so nothing was shown. io.record_rule_surfaced.reset_mock() held = rp.RuleMoment(user_id=1, query="the reply", project_id=2, held=frozenset({40})) result = await rp.run_rule_arm(rp.REPLY_RULE, held, floor=0.8, budget=3, io=io, checkpoint_floor=0.8) assert result.checkpoint == {} io.record_rule_surfaced.assert_not_called() def test_the_reply_arm_is_a_registered_ranked_surface(): from scribe.services.retrieval_registry import POINTS assert rp.REPLY_RULE in rp.RULE_ARMS assert "reply_rule" in surfaces.SURFACES and "reply_rule" in POINTS assert "reply_rule" in rule_usage.RANKED_SOURCES # A stop bar, not a hint bar: it sits with the checkpoint's default. assert surfaces.SURFACES["reply_rule"].floor_default >= pc._CHECKPOINT_DEFAULT # ── the route ─────────────────────────────────────────────────────────── async def test_the_route_passes_the_reply_and_all_three_ledgers(): from scribe.routes import plugin as routes hold = AsyncMock(return_value={"reason": "Held.", "rule_ids": [11], "moments": ["reply.report"]}) app = Quart(__name__) async with app.test_request_context( "/api/plugin/reply-rules", method="POST", json={"reply": "done"}, query_string={"project_id": "3", "held_rule_ids": "4", "exclude_rule_ids": "5", "stopped_rule_ids": "6,7"}, ): g.user = type("U", (), {"id": 7})() with patch.object(routes.moment_delivery_svc, "reply_hold", hold): resp = await routes.reply_rules.__wrapped__() body = await resp.get_json() assert body == {"reason": "Held.", "rule_ids": [11], "moments": ["reply.report"], "report_check": ""} assert hold.await_args.args == (7, "done") assert hold.await_args.kwargs == { "project_id": 3, "exclude": frozenset({5}), "held": frozenset({4}), "stopped": frozenset({6, 7}), } async def _reply_route(body, *, hold_reason="", checked=None): from scribe.routes import plugin as routes hold = AsyncMock(return_value={"reason": hold_reason, "rule_ids": [11] if hold_reason else [], "moments": ["reply.report"]}) check = AsyncMock(return_value=checked or {"outcome": "passed", "reason": ""}) app = Quart(__name__) async with app.test_request_context("/api/plugin/reply-rules", method="POST", json=body, query_string={"project_id": "3"}): g.user = type("U", (), {"id": 7})() with patch.object(routes.moment_delivery_svc, "reply_hold", hold), \ patch.object(routes.report_check_svc, "check_reply", check): resp = await routes.reply_rules.__wrapped__() return await resp.get_json(), hold, check async def test_a_turn_that_closed_nothing_is_not_section_checked(): _body, _hold, check = await _reply_route({"reply": "done"}) check.assert_not_awaited() async def test_a_closing_turn_is_section_checked_and_both_holds_fold_into_one_reason(): """One end-of-turn request (milestone 500 step 4): the section check and the rule hold answer together, in one `reason`.""" body, _hold, check = await _reply_route( {"reply": "done", "closed": 1, "closed_task_ids": [41, True, "x"], "rewrite": False}, hold_reason="RULE HOLD", checked={"outcome": "blocked", "reason": "SECTIONS"}, ) assert body["reason"] == "SECTIONS\n\nRULE HOLD" assert body["report_check"] == "blocked" # JSON `true` is not task 1, and a string is not an id. assert check.await_args.kwargs == {"task_ids": [41], "rewrite": False, "project_id": 3} async def test_a_task_created_already_done_is_checked_without_an_id(): _body, _hold, check = await _reply_route({"reply": "done", "closed": 1, "closed_task_ids": []}) assert check.await_args.kwargs["task_ids"] == [] async def test_the_rewrite_is_recorded_and_held_by_nothing(): body, hold, check = await _reply_route( {"reply": "done", "closed": 1, "closed_task_ids": [41], "rewrite": True}, checked={"outcome": "passed_after_rewrite", "reason": ""}, ) assert check.await_args.kwargs["rewrite"] is True hold.assert_not_awaited() assert body["reason"] == "" and body["report_check"] == "passed_after_rewrite" # ── the hook ──────────────────────────────────────────────────────────── def _transcript(tmp_path, reply): lines = [ {"type": "user", "message": {"role": "user", "content": "is it finished?"}}, {"type": "assistant", "message": {"content": [{"type": "text", "text": reply}]}}, ] path = tmp_path / "t.jsonl" path.write_text("\n".join(json.dumps(x, separators=(",", ":")) for x in lines) + "\n") return path def _run(tmp_path, port, transcript, active=False): need_tools("bash", "curl", "awk") env = {"PATH": os.environ["PATH"], "SCRIBE_URL": f"http://127.0.0.1:{port}", "SCRIBE_TOKEN": "t", "TMPDIR": str(tmp_path), "HOME": str(tmp_path)} out = subprocess.run( ["bash", str(HOOK)], input=json.dumps({"session_id": "s-reply", "transcript_path": str(transcript), "cwd": str(tmp_path), "hook_event_name": "Stop", "stop_hook_active": active}), capture_output=True, text=True, env=env, timeout=30, ) assert out.returncode == 0, out.stderr return out.stdout.strip() HELD = json.dumps({"reason": "Held for one read: get_rule(11).", "rule_ids": [11], "moments": ["reply.report"]}).encode() QUIET = b'{"reason":"","rule_ids":[],"moments":["reply.report"]}' def test_the_hook_blocks_in_the_servers_words_and_records_the_hold(tmp_path): t = _transcript(tmp_path, "All done.\nLet me know if it works on your end?") with http_sink(by_path={"/api/plugin/reply-rules": HELD}) as (port, seen): out = json.loads(_run(tmp_path, port, t)) _run(tmp_path, port, t) assert out == {"decision": "block", "reason": "Held for one read: get_rule(11)."} sent = json.loads(seen[0]["_body"]) assert sent["reply"] == "All done.\nLet me know if it works on your end?" # The second turn tells the server what already held. assert seen[1]["stopped_rule_ids"] == ["11"] assert seen[1]["exclude_rule_ids"] == ["11"] def _closing_transcript(tmp_path, *replies): lines = [ {"type": "user", "message": {"role": "user", "content": "please finish it"}}, {"type": "assistant", "message": {"content": [ {"type": "tool_use", "id": "toolu_1", "name": "mcp__plugin_scribe_scribe__update_task", "input": {"task_id": 41, "status": "done"}}]}}, {"type": "user", "message": {"content": [ {"type": "tool_result", "tool_use_id": "toolu_1", "is_error": False, "content": "{}"}]}}, ] + [{"type": "assistant", "message": {"content": [{"type": "text", "text": r}]}} for r in replies] path = tmp_path / "t.jsonl" path.write_text("\n".join(json.dumps(x, separators=(",", ":")) for x in lines) + "\n") return path SECTION_HOLD = json.dumps({"reason": "Rewrite as a completion report.", "rule_ids": [], "moments": ["reply.report"], "report_check": "blocked"}).encode() def test_a_turn_that_closed_nothing_sends_no_close(tmp_path): t = _transcript(tmp_path, "Done.") with http_sink(by_path={"/api/plugin/reply-rules": QUIET}) as (port, seen): _run(tmp_path, port, t) assert "closed" not in json.loads(seen[0]["_body"]) def test_a_section_hold_blocks_once_then_the_rewrite_is_reported_and_never_held(tmp_path): """The completion-section check, folded into the one Stop request: the first stop is held in the server's words, the rewrite goes back once with `rewrite: true` so its outcome is recorded, and nothing after that is sent.""" t = _closing_transcript(tmp_path, "All done, pushed it.") with http_sink(by_path={"/api/plugin/reply-rules": SECTION_HOLD}) as (port, seen): out = json.loads(_run(tmp_path, port, t)) assert out == {"decision": "block", "reason": "Rewrite as a completion report."} first = json.loads(seen[0]["_body"]) assert (first["closed"], first["closed_task_ids"], first["rewrite"]) == (1, [41], False) rewritten = _closing_transcript(tmp_path, "All done, pushed it.", "Where this sits: …") assert _run(tmp_path, port, rewritten, active=True) == "" assert json.loads(seen[1]["_body"])["rewrite"] is True # The marker is spent: a further stop in the loop sends nothing. assert _run(tmp_path, port, rewritten, active=True) == "" assert len(seen) == 2 def test_a_rule_hold_on_a_closing_turn_does_not_mark_a_rewrite(tmp_path): t = _closing_transcript(tmp_path, "Done.") with http_sink(by_path={"/api/plugin/reply-rules": HELD}) as (port, seen): _run(tmp_path, port, t) assert _run(tmp_path, port, t, active=True) == "" assert len(seen) == 1 def test_the_hook_says_nothing_when_nothing_holds(tmp_path): t = _transcript(tmp_path, "Done.") with http_sink(by_path={"/api/plugin/reply-rules": QUIET}) as (port, _seen): assert _run(tmp_path, port, t) == "" def test_the_rewrite_is_never_held(tmp_path): t = _transcript(tmp_path, "Done.") with http_sink(by_path={"/api/plugin/reply-rules": HELD}) as (port, seen): assert _run(tmp_path, port, t, active=True) == "" assert seen == [] def test_an_unreachable_instance_never_holds_a_reply(tmp_path): t = _transcript(tmp_path, "Done.") assert _run(tmp_path, 9, t) == "" @pytest.mark.parametrize("name", ["scribe_reply_check.sh"]) def test_the_hook_is_registered_on_stop(name): manifest = json.loads((ROOT / "plugin" / "hooks" / "hooks.json").read_text()) commands = [h["command"] for m in manifest["hooks"]["Stop"] for h in m["hooks"]] assert any(name in c for c in commands)