Implementeaza planul claude-master-plan-discord-bridge-20260830 (15 taskuri, 3 lane-uri paralele) — un bot subtire discord.py peste CLI-ul `claude`, cu proces persistent per fir alimentat pe stdin cu --input-format stream-json. Nucleu: runner (proces persistent + reaper 20min + respawn --resume), stream (parser tolerant), session_store (scriere atomica, lock per fir, detectare PID reuse, recovery), limits (max 4 procese, timeout tur, rate per user, plafon cost pe zi), render (un loop de editare per canal, interval adaptiv). Adaptor: allowlist guild/canal/user fail-closed cu respingerea webhook-urilor, comenzi !new/!cd/!model/!status/!stop/!cleanup, cost si model in subsolul fiecarui raspuns. Mesajul sosit in timpul unui tur devine steering, nu tur nou. Securitate: hook PreToolUse fail-closed care cere confirmare in Discord pentru operatiuni ireversibile, wrapper `infra` cu lista explicita de hosturi. Deny rules raman strat cosmetic, nu bariera (verificat: /usr/bin/ssh trece pe langa). Ops: alerte email pe conventia repo-ului, !cleanup pentru orfani, unit systemd user cu KillMode=control-group si limite de memorie, install.sh idempotent. Verificat: 275 teste fara retea/Discord/API (10.8s), identic cu si fara discord.py instalat; e2e pe CLI real confirma steering-ul mid-tur (mesaj la 6s intr-un tool call de 25s schimba raspunsul final). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01B29CApsP1JkSdjYaGaHpE7
223 lines
7.4 KiB
Python
223 lines
7.4 KiB
Python
"""T4 + T5: proces persistent, steering, respawn, reaper, EOF, timeout."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import os
|
|
|
|
import pytest
|
|
|
|
import runner
|
|
import stream
|
|
|
|
|
|
def mk(fake_bin, store=None, **kw):
|
|
return runner.ClaudeProcess("1", os.getcwd(), "sonnet", claude_bin=fake_bin,
|
|
on_pid=(store.set_pid if store else None), **kw)
|
|
|
|
|
|
def test_build_cmd_are_toate_flagurile():
|
|
cmd = runner.build_cmd("claude", "opus", None, "/x/settings.json")
|
|
assert cmd[:2] == ["claude", "-p"]
|
|
for flag in ("--input-format", "--output-format", "--verbose", "--permission-mode",
|
|
"--settings", "--model", "--autocompact"):
|
|
assert flag in cmd
|
|
assert cmd[cmd.index("--permission-mode") + 1] == "bypassPermissions"
|
|
assert "--resume" not in cmd
|
|
assert "--resume" in runner.build_cmd(["python", "fake"], "sonnet", "sid-1")
|
|
|
|
|
|
async def test_tur_normal_si_proces_persistent(fake_bin, scenario):
|
|
scenario("normal")
|
|
p = mk(fake_bin)
|
|
out = await p.run_turn("salut")
|
|
assert out.result.text == "ecou: salut" and out.result.total_cost_usd == pytest.approx(0.0123)
|
|
assert p.sid == "sid-fake-0001" and not out.restarted
|
|
pid = p.pid
|
|
out2 = await p.run_turn("inca unul")
|
|
assert out2.result.text == "ecou: inca unul"
|
|
assert p.pid == pid and p.alive # ACELASI proces pentru turul urmator
|
|
await p.stop()
|
|
assert not p.alive
|
|
|
|
|
|
async def test_evenimentele_ajung_la_callback(fake_bin, scenario):
|
|
scenario("tools")
|
|
p = mk(fake_bin)
|
|
seen = []
|
|
await p.run_turn("fa ceva", on_event=lambda ev: (seen.append(ev), asyncio.sleep(0))[1])
|
|
tipuri = [type(e) for e in seen]
|
|
assert stream.SystemInit in tipuri and stream.ToolUse in tipuri
|
|
assert stream.ToolResult in tipuri and tipuri[-1] is stream.Result
|
|
await p.stop()
|
|
|
|
|
|
async def test_stream_tolerant_in_tur_real(fake_bin, scenario):
|
|
scenario("unknown")
|
|
p = mk(fake_bin)
|
|
out = await p.run_turn("salut") # tip necunoscut + linie non-JSON pe mijloc
|
|
assert out.result.text == "ecou: salut"
|
|
await p.stop()
|
|
|
|
|
|
async def test_steering_mid_tur(fake_bin, scenario):
|
|
"""Mesaj trimis in timp ce turul ruleaza ajunge la proces inainte de result."""
|
|
scenario("slow", FAKE_CLAUDE_DELAY=0.6)
|
|
p = mk(fake_bin)
|
|
|
|
async def steer():
|
|
await asyncio.sleep(0.15)
|
|
await p.send("de fapt, opreste-te")
|
|
|
|
task = asyncio.create_task(steer())
|
|
out = await p.run_turn("porneste ceva lung")
|
|
await task
|
|
assert "porneste ceva lung" in out.result.text
|
|
assert "de fapt, opreste-te" in out.result.text
|
|
await p.stop()
|
|
|
|
|
|
async def test_eof_inainte_de_result_da_turn_failed(fake_bin, scenario):
|
|
scenario("eof")
|
|
p = mk(fake_bin)
|
|
with pytest.raises(runner.TurnFailed):
|
|
await p.run_turn("salut")
|
|
assert not p.alive
|
|
|
|
|
|
async def test_crash_pastreaza_stderr_pentru_status(fake_bin, scenario):
|
|
scenario("crash")
|
|
p = mk(fake_bin)
|
|
with pytest.raises(runner.TurnFailed) as exc:
|
|
await p.run_turn("salut")
|
|
assert "boom" in str(exc.value)
|
|
|
|
|
|
async def test_stderr_buffer_circular(fake_bin, scenario):
|
|
scenario("crash")
|
|
p = mk(fake_bin)
|
|
with pytest.raises(runner.TurnFailed):
|
|
await p.run_turn("x")
|
|
assert p.stderr_buf.maxlen == runner.STDERR_TAIL
|
|
|
|
|
|
async def test_respawn_transparent_cu_resume(fake_bin, scenario, store):
|
|
scenario("normal")
|
|
p = mk(fake_bin, store=store)
|
|
await p.run_turn("primul")
|
|
sid = p.sid
|
|
p.proc.kill() # OOM simulat
|
|
await p.proc.wait()
|
|
p.proc = None
|
|
out = await p.run_turn("al doilea")
|
|
assert out.restarted is True
|
|
restart_ev = [e for e in out.events if isinstance(e, runner.SessionRestarted)]
|
|
assert restart_ev and "sesiune repornita" in restart_ev[0].text
|
|
assert p.sid == sid # --resume a pastrat sesiunea
|
|
assert p.restarts == 1
|
|
await p.stop()
|
|
|
|
|
|
async def test_timeout_de_tur_omoara_procesul(fake_bin, scenario):
|
|
scenario("slow", FAKE_CLAUDE_DELAY=5)
|
|
p = mk(fake_bin)
|
|
with pytest.raises(runner.TurnTimeout):
|
|
await p.run_turn("lung", timeout=0.3)
|
|
assert not p.alive
|
|
|
|
|
|
async def test_send_pe_proces_mort_da_eroare(fake_bin, scenario):
|
|
scenario("normal")
|
|
p = mk(fake_bin)
|
|
with pytest.raises(runner.TurnFailed):
|
|
await p.send("nimeni nu asculta")
|
|
|
|
|
|
async def test_pid_ul_ajunge_in_state(fake_bin, scenario, store):
|
|
scenario("normal")
|
|
m = runner.RunnerManager(store, claude_bin=fake_bin)
|
|
p = m.get("42", os.getcwd(), "sonnet")
|
|
await p.run_turn("salut")
|
|
assert store.thread("42")["pid"] == p.pid
|
|
assert store.thread_process_alive("42")
|
|
await m.stop_all()
|
|
assert store.thread("42")["pid"] is None
|
|
|
|
|
|
async def test_reaper_omoara_inactivii_dar_nu_turul_in_zbor(fake_bin, scenario, store):
|
|
scenario("normal")
|
|
m = runner.RunnerManager(store, claude_bin=fake_bin, idle_s=0.0)
|
|
a = m.get("a", os.getcwd(), "sonnet")
|
|
b = m.get("b", os.getcwd(), "sonnet")
|
|
await a.run_turn("x")
|
|
await b.run_turn("y")
|
|
store.set_inflight("b", "t", "u", "m") # firul b are tur in zbor
|
|
killed = await m.reap_once()
|
|
assert killed == ["a"]
|
|
assert not a.alive and b.alive
|
|
store.clear_inflight("b")
|
|
assert await m.reap_once() == ["b"]
|
|
await m.stop_all()
|
|
|
|
|
|
async def test_reaper_nu_taie_daca_procesul_e_activ(fake_bin, scenario):
|
|
scenario("normal")
|
|
m = runner.RunnerManager(None, claude_bin=fake_bin, idle_s=60.0)
|
|
p = m.get("a", os.getcwd(), "sonnet")
|
|
await p.run_turn("x")
|
|
assert await m.reap_once() == [] and p.alive
|
|
await m.stop_all()
|
|
|
|
|
|
async def test_set_options_opreste_procesul_iar_resume_pastreaza_sesiunea(fake_bin, scenario):
|
|
scenario("normal")
|
|
m = runner.RunnerManager(None, claude_bin=fake_bin)
|
|
p = m.get("a", os.getcwd(), "sonnet")
|
|
await p.run_turn("x")
|
|
sid = p.sid
|
|
assert await m.set_options("a", model="opus") is True
|
|
assert not p.alive and p.model == "opus"
|
|
out = await p.run_turn("y")
|
|
assert out.restarted and p.sid == sid
|
|
await m.stop_all()
|
|
|
|
|
|
async def test_reset_new_sterge_sesiunea(fake_bin, scenario):
|
|
scenario("normal")
|
|
m = runner.RunnerManager(None, claude_bin=fake_bin)
|
|
p = m.get("a", os.getcwd(), "sonnet")
|
|
await p.run_turn("x")
|
|
await m.reset("a")
|
|
assert p.sid is None and not p.alive
|
|
await m.stop_all()
|
|
|
|
|
|
async def test_live_count_si_stop_all(fake_bin, scenario):
|
|
scenario("normal")
|
|
m = runner.RunnerManager(None, claude_bin=fake_bin)
|
|
for tid in ("a", "b"):
|
|
await m.get(tid, os.getcwd(), "sonnet").run_turn("x")
|
|
assert m.live_count() == 2
|
|
await m.stop_all()
|
|
assert m.live_count() == 0
|
|
|
|
|
|
async def test_thread_id_ajunge_in_mediul_procesului(fake_bin, scenario, monkeypatch):
|
|
"""Lane B: hook-ul PreToolUse citeste firul din CLAUDE_DISCORD_THREAD_ID."""
|
|
scenario("env")
|
|
monkeypatch.setenv("MARKER_DE_TEST", "pastrat")
|
|
p = runner.ClaudeProcess("fir-777", os.getcwd(), "sonnet", claude_bin=fake_bin)
|
|
out = await p.run_turn("x")
|
|
assert "thread=fir-777" in out.result.text
|
|
assert "sesiune=-" in out.result.text # inca nu avem sid la prima pornire
|
|
assert "marker=pastrat" in out.result.text # restul mediului ramane intact
|
|
await p.stop()
|
|
|
|
|
|
async def test_session_id_ajunge_in_mediu_la_respawn(fake_bin, scenario):
|
|
scenario("env")
|
|
p = runner.ClaudeProcess("fir-888", os.getcwd(), "sonnet", sid="sid-vechi", claude_bin=fake_bin)
|
|
out = await p.run_turn("x")
|
|
assert "thread=fir-888" in out.result.text and "sesiune=sid-vechi" in out.result.text
|
|
await p.stop()
|