import os import subprocess import json import hmac import hashlib import shlex import sys import pytest # Helper to run bash snippets sourcing lib.sh def run_lib_func(mam_sandbox, func_name, *args, env=None): lib_path = mam_sandbox / "skills" / "lib.sh" cmd_str = f"source {lib_path} && {func_name} " + " ".join(shlex.quote(str(a)) for a in args) run_env = dict(os.environ) run_env.pop("HERDR_SESSION_NAME", None) run_env.pop("HERDR_SERVER_NAME", None) if env: run_env.update(env) res = subprocess.run(["bash", "-c", cmd_str], capture_output=True, text=True, env=run_env) return res def get_mqtt_common(mam_sandbox): script_path = str(mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts") if script_path not in sys.path: sys.path.insert(0, script_path) import mqtt_common return mqtt_common # ============================================================================== # FEATURE 1: Create Session (5 Test Cases) # ============================================================================== def test_create_derive_session_name_standard(mam_sandbox): """Test standard derive_session_name slug generation.""" res = run_lib_func(mam_sandbox, "derive_session_name", "/home/user/project", "claude") assert res.returncode == 0 assert res.stdout.strip() == "user-project-creator-claude" def test_create_derive_session_name_nested(mam_sandbox): """Test derive_session_name with nested paths, upper casing, and underscores.""" res = run_lib_func(mam_sandbox, "derive_session_name", "/home/User_Name/My_New_Project", "agy") assert res.returncode == 0 assert res.stdout.strip() == "user-name-my-new-project-creator-agy" def test_create_derive_session_name_weird_characters(mam_sandbox): """Test derive_session_name with spaces and punctuation in the path.""" res = run_lib_func(mam_sandbox, "derive_session_name", "/a/b c/d-e!f", "hermes") assert res.returncode == 0 assert res.stdout.strip() == "bc-d-ef-creator-hermes" def test_create_session_legacy_isolate_flags_noop(mam_sandbox): """Legacy --isolate/--no-isolate must stay a documented no-op, not an arg-parser error.""" create_script = mam_sandbox / "skills" / "multi-agent-mux-create" / "scripts" / "create_session.sh" for flag in ["--isolate", "--no-isolate"]: res = subprocess.run(["bash", str(create_script), flag, "-h"], capture_output=True, text=True) assert res.returncode == 0, f"{flag} rejected by arg parser: {res.stderr}" assert "NOTE: --isolate/--no-isolate is a no-op" in res.stderr assert flag in res.stdout, f"{flag} missing from usage() help text" def test_create_validate_env_key(mam_sandbox): """Test _validate_env_key function with valid and blocked environment keys.""" # Valid key res = run_lib_func(mam_sandbox, "_validate_env_key", "MY_VALID_KEY") assert res.returncode == 0 # Malformed key (starts with number) res2 = run_lib_func(mam_sandbox, "_validate_env_key", "123BAD") assert res2.returncode != 0 # Blocked key (LD_PRELOAD) res3 = run_lib_func(mam_sandbox, "_validate_env_key", "LD_PRELOAD") assert res3.returncode != 0 # ============================================================================== # FEATURE 2: Resume Session (6 Test Cases) # ============================================================================== def test_resume_resolve_herdr_session_default(mam_sandbox): """Test resolve_herdr_workspace fallback behavior when session is not in YAML.""" res = run_lib_func(mam_sandbox, "resolve_herdr_workspace", "non-existent-session") assert res.returncode == 0 assert res.stdout.strip() != "" def test_resume_resolve_herdr_session_env(mam_sandbox): """Test resolve_herdr_workspace fallback to HERDR_SESSION_NAME or HERDR_SERVER_NAME env var.""" res = run_lib_func(mam_sandbox, "resolve_herdr_workspace", "non-existent-session", env={"HERDR_SESSION_NAME": "custom_session"}) assert res.returncode == 0 assert res.stdout.strip() == "custom_session" res_legacy = run_lib_func(mam_sandbox, "resolve_herdr_workspace", "non-existent-session", env={"HERDR_SERVER_NAME": "custom_server"}) assert res_legacy.returncode == 0 assert res_legacy.stdout.strip() == "custom_server" def test_resume_find_workspace_uuid_empty(mam_sandbox): """Test find_workspace_uuid returns empty string for non-existent workspace.""" res = run_lib_func(mam_sandbox, "find_workspace_uuid", "/non/existent/path", "claude") assert res.returncode == 0 assert res.stdout.strip() == "" def test_resume_find_workspace_uuid_target_non_existent(mam_sandbox): """Test target session query with target that does not exist in YAML.""" res = run_lib_func(mam_sandbox, "find_workspace_uuid", str(mam_sandbox), "claude", "non-existent-session") assert res.returncode == 0 assert res.stdout.strip() == "" def test_resume_find_workspace_uuid_invalid_agent(mam_sandbox): """Test find_workspace_uuid behavior with an unsupported agent name.""" res = run_lib_func(mam_sandbox, "find_workspace_uuid", str(mam_sandbox), "invalidagent") assert res.returncode == 0 assert res.stdout.strip() == "" def test_resume_script_invalid_args(mam_sandbox): """Test calling resolve_session_id.sh with missing arguments.""" script_path = mam_sandbox / "skills" / "multi-agent-mux-resume" / "scripts" / "resolve_session_id.sh" res = subprocess.run(["bash", str(script_path), "--workspace", str(mam_sandbox)], capture_output=True, text=True) assert res.returncode == 2 assert "ERROR: --agent required" in res.stderr # ============================================================================== # FEATURE 3: Stop Session (5 Test Cases) # ============================================================================== def test_stop_check_is_nfs_local(mam_sandbox): """Test that _check_is_nfs on local temp directory returns non-zero (not NFS).""" res = run_lib_func(mam_sandbox, "_check_is_nfs", str(mam_sandbox)) # It will exit with 1 if it is not NFS assert res.returncode == 1 def test_stop_is_already_stopped_not_found(mam_sandbox): """Test is_already_stopped exits with 1 when session is not in YAML.""" res = run_lib_func(mam_sandbox, "is_already_stopped", "non-existent-session") assert res.returncode == 1 def test_stop_session_invalid_agent_suffix(mam_sandbox): """Test stop_session.sh fails when agent cannot be inferred from session name.""" script_path = mam_sandbox / "skills" / "multi-agent-mux-stop" / "scripts" / "stop_session.sh" res = subprocess.run(["bash", str(script_path), "--session", "bad-session-name"], capture_output=True, text=True) assert res.returncode == 2 assert "ERROR: cannot infer agent" in res.stderr def test_stop_session_missing_required_args(mam_sandbox): """Test stop_session.sh fails when session name is missing.""" script_path = mam_sandbox / "skills" / "multi-agent-mux-stop" / "scripts" / "stop_session.sh" res = subprocess.run(["bash", str(script_path)], capture_output=True, text=True) assert res.returncode == 2 assert "ERROR: --session required" in res.stderr def test_stop_session_purge_no_yes(mam_sandbox): """Test stop_session.sh exits with 3 when purge is requested without --yes.""" script_path = mam_sandbox / "skills" / "multi-agent-mux-stop" / "scripts" / "stop_session.sh" # Seed yaml with the session first to avoid "not in yaml" exit 1 yaml_path = mam_sandbox / ".mam" / "agent-sessions.yaml" yaml_path.write_text("""herdr_sessions: - name: test-project-creator-claude status: running pane: cwd: /tmp """) res = subprocess.run(["bash", str(script_path), "--session", "test-project-creator-claude", "--purge-conversation"], capture_output=True, text=True) assert res.returncode == 3 assert "DANGER: --purge-conversation will DELETE" in res.stdout # ============================================================================== # FEATURE 4: Status Query (5 Test Cases) # ============================================================================== def test_status_json_schema_fields(mam_sandbox, mock_herdr, mock_agents): """Verify status.sh output JSON contains the expected structure.""" script_path = mam_sandbox / "skills" / "multi-agent-mux-status" / "scripts" / "status.sh" # Run status.sh with --json res = subprocess.run(["bash", str(script_path), "--json"], capture_output=True, text=True) assert res.returncode == 0 data = json.loads(res.stdout) assert "timestamp" in data assert "yaml_path" in data assert "drifts" in data assert "sessions_detail" in data def test_status_text_headers_presence(mam_sandbox): """Verify status.sh output in text mode includes header columns.""" script_path = mam_sandbox / "skills" / "multi-agent-mux-status" / "scripts" / "status.sh" res = subprocess.run(["bash", str(script_path)], capture_output=True, text=True) assert res.returncode == 0 assert "NAME" in res.stdout assert "WORKSPACE" in res.stdout assert "YAML" in res.stdout assert "HERDR" in res.stdout assert "DRIFT" in res.stdout def test_status_resume_on_disk_helper_claude(mam_sandbox): """Verify resume_on_disk behavior for claude inside status.sh logic via mock YAML queries.""" # We can write a custom yaml and check status yaml_path = mam_sandbox / ".mam" / "agent-sessions.yaml" yaml_path.write_text("""herdr_sessions: - name: test-project-creator-claude status: running claude_session_id_own: some-uuid pane: cwd: /tmp/nonexistent-workspace """) script_path = mam_sandbox / "skills" / "multi-agent-mux-status" / "scripts" / "status.sh" res = subprocess.run(["bash", str(script_path), "--json"], capture_output=True, text=True) assert res.returncode == 0 data = json.loads(res.stdout) detail = data["sessions_detail"][0] assert detail["name"] == "test-project-creator-claude" assert detail["resume_state"] == "MISSING" def test_status_resume_on_disk_helper_agy(mam_sandbox): """Verify resume_on_disk behavior for agy inside status.sh logic.""" yaml_path = mam_sandbox / ".mam" / "agent-sessions.yaml" yaml_path.write_text("""herdr_sessions: - name: test-project-creator-agy status: running agy_conversation_id_own: some-uuid pane: cwd: /tmp/nonexistent-workspace """) script_path = mam_sandbox / "skills" / "multi-agent-mux-status" / "scripts" / "status.sh" res = subprocess.run(["bash", str(script_path), "--json"], capture_output=True, text=True) assert res.returncode == 0 data = json.loads(res.stdout) detail = data["sessions_detail"][0] assert detail["name"] == "test-project-creator-agy" assert detail["resume_state"] == "MISSING" def test_status_get_job_status_helper(mam_sandbox): """Verify get_job_status parses non-existent jobs gracefully.""" yaml_path = mam_sandbox / ".mam" / "agent-sessions.yaml" yaml_path.write_text("""herdr_sessions: - name: test-project-creator-claude status: running delegate_job_id: nonexistent-job-id pane: cwd: /tmp """) script_path = mam_sandbox / "skills" / "multi-agent-mux-status" / "scripts" / "status.sh" res = subprocess.run(["bash", str(script_path), "--json"], capture_output=True, text=True) assert res.returncode == 0 data = json.loads(res.stdout) detail = data["sessions_detail"][0] assert detail["job_id"] == "nonexistent-job-id" assert detail["job_status"] == "unknown" # ============================================================================== # FEATURE 5: Monitor/Reconcile (6 Test Cases) # ============================================================================== def test_mqtt_topic_prefix_for(mam_sandbox): """Test topic prefix helper methods in mqtt_common.""" mqtt_common = get_mqtt_common(mam_sandbox) prefix = mqtt_common.topic_prefix_for("job123") assert prefix == "python/mqtt/jobs/job123" events_topic = mqtt_common.events_topic_for("job123") assert events_topic == "python/mqtt/jobs/job123/events" def test_mqtt_verify_hmac_no_token(mam_sandbox): """Verify verify_hmac returns True when no token is present.""" mqtt_common = get_mqtt_common(mam_sandbox) payload = {"data": {"hmac_sig": "somesig"}} assert mqtt_common.verify_hmac(payload, None) is True assert mqtt_common.verify_hmac(payload, "") is True def test_mqtt_verify_hmac_valid_invalid(mam_sandbox): """Verify verify_hmac signature validation matching logic.""" mqtt_common = get_mqtt_common(mam_sandbox) token = "secret_key" payload = { "job_id": "job1", "event": "started", "seq": 1, "timestamp": "2026-07-19T00:00:00Z", "data": { "some_key": "some_val" } } # Calculate HMAC msg = json.dumps(payload, sort_keys=True, separators=(",", ":")).encode() sig = hmac.new(token.encode(), msg, hashlib.sha256).hexdigest() # Put HMAC signature inside data block payload["data"]["hmac_sig"] = sig assert mqtt_common.verify_hmac(payload, token) is True # Modifying payload should cause verification to fail payload["seq"] = 2 assert mqtt_common.verify_hmac(payload, token) is False def test_mqtt_reason_code_value(mam_sandbox): """Verify reason_code_value correctly extracts values from paho reason codes.""" mqtt_common = get_mqtt_common(mam_sandbox) # Test plain int assert mqtt_common.reason_code_value(0) == 0 assert mqtt_common.reason_code_value(5) == 5 # Test object with .value attribute class DummyReasonCode: def __init__(self, val): self.value = val assert mqtt_common.reason_code_value(DummyReasonCode(0)) == 0 assert mqtt_common.reason_code_value(DummyReasonCode(16)) == 16 def test_mqtt_with_retry_success(mam_sandbox): """Verify with_retry decorator works on direct success.""" mqtt_common = get_mqtt_common(mam_sandbox) calls = [] @mqtt_common.with_retry(attempts=3) def dummy_func(x): calls.append(x) return x * 2 res = dummy_func(5) assert res == 10 assert calls == [5] def test_mqtt_with_retry_failure(mam_sandbox): """Verify with_retry decorator raises error after specified attempts.""" mqtt_common = get_mqtt_common(mam_sandbox) calls = [] @mqtt_common.with_retry(attempts=3, base_delay=0.01) def failing_func(): calls.append(1) raise ValueError("failing") with pytest.raises(ValueError, match="failing"): failing_func() assert len(calls) == 3 def test_b10_no_agent_identities_reader_in_production(): """B-10: agent_identities has no writer; no production code may read it.""" import pathlib root = pathlib.Path(__file__).resolve().parent.parent targets = [ root / ".agents" / "skills" / "lib_py" / "workspace_uuid.py", root / ".agents" / "skills" / "multi-agent-mux-monitor" / "scripts" / "reconcile.sh", root / ".agents" / "skills" / "multi-agent-mux-stop" / "scripts" / "stop_session.sh", root / ".agents" / "skills" / "lib.sh", ] offenders = [] for f in targets: for i, line in enumerate(f.read_text().splitlines(), 1): if "agent_identities" not in line: continue if line.lstrip().startswith("#"): # 금지 규약을 서술하는 주석은 허용 continue offenders.append(f"{f.name}:{i}: {line.strip()}") assert not offenders, "agent_identities read path resurrected:\n" + "\n".join(offenders) def test_b10_workspace_uuid_has_no_yaml_import(): """B-10: no `import yaml` anywhere in workspace_uuid.py — top-level OR lazy.""" import ast, pathlib src = (pathlib.Path(__file__).resolve().parent.parent / ".agents" / "skills" / "lib_py" / "workspace_uuid.py") tree = ast.parse(src.read_text()) offenders = [] for node in ast.walk(tree): # ast.walk → 중첩 깊이 무관 if isinstance(node, ast.Import): for a in node.names: if a.name.split(".")[0] == "yaml": offenders.append(f"line {node.lineno}: import {a.name}") elif isinstance(node, ast.ImportFrom): if (node.module or "").split(".")[0] == "yaml": offenders.append(f"line {node.lineno}: from {node.module} import ...") assert not offenders, "PyYAML dependency reintroduced:\n" + "\n".join(offenders) def test_b10_find_workspace_uuid_runs_without_pyyaml(tmp_path): """B-10: the executed resolution path must not need PyYAML.""" import subprocess, sys, os, json, pathlib stub = tmp_path / "noyaml" (stub / "yaml").mkdir(parents=True) (stub / "yaml" / "__init__.py").write_text('raise ImportError("PyYAML absent (stub)")\n') skills = str(pathlib.Path(__file__).resolve().parent.parent / ".agents" / "skills") ws = tmp_path / "ws"; ws.mkdir() env = os.environ.copy() env["PYTHONPATH"] = f"{stub}:{skills}" env["WS_ABS"] = str(ws) env["AGENT"] = "claude" env["MAM_STATE_JSON"] = json.dumps({"herdr_sessions": []}) env["YAML_PATH"] = str(tmp_path / "agent-sessions.yaml") env["HOME_DIR"] = str(tmp_path) env["CLAUDE_PROJECT_DIR"] = str(tmp_path / "projects") r = subprocess.run( [sys.executable, "-c", "from lib_py.workspace_uuid import find_workspace_uuid_main; find_workspace_uuid_main()"], capture_output=True, text=True, env=env) assert r.returncode == 0, f"resolution path still needs PyYAML: {r.stderr}" assert "yaml" not in r.stderr.lower(), f"PyYAML touched at runtime: {r.stderr}" # =========================================================================== # B-9: LOGS_DIR import-time cwd binding resolution regression guards # =========================================================================== def test_b9_logs_dir_follows_cwd_changes(mam_sandbox, tmp_path, monkeypatch): """B-9: the audit-log root must be resolved per call, not frozen at import.""" mq = get_mqtt_common(mam_sandbox) monkeypatch.delenv("DELEGATE_JOB_LOGS_DIR", raising=False) a = tmp_path / "a"; b = tmp_path / "b" a.mkdir(); b.mkdir() # realpath on both sides: pytest's tmp_path happens to be pre-resolved today, # but relying on that is an undocumented dependency (C1'). def logs_under(p): return os.path.realpath(os.path.join(str(p), ".mam", "delegate_job_logs")) monkeypatch.chdir(a) assert os.path.realpath(mq.get_logs_dir()) == logs_under(a) monkeypatch.chdir(b) assert os.path.realpath(mq.get_logs_dir()) == logs_under(b) # the compat alias must follow too (T1: a surviving global fails here) assert os.path.realpath(mq.LOGS_DIR) == logs_under(b) def test_b9_audit_log_lands_under_the_current_cwd(mam_sandbox, tmp_path, monkeypatch): """B-9/T3: assert the FILE appears — a swallowed NameError must not pass.""" mq = get_mqtt_common(mam_sandbox) monkeypatch.delenv("DELEGATE_JOB_LOGS_DIR", raising=False) monkeypatch.chdir(tmp_path) mq.init_job_log("b9job", {"status": "pending"}) assert (tmp_path / ".mam" / "delegate_job_logs" / "b9job" / "meta.json").exists(), \ "audit log did not land under the current cwd (the best-effort handler may have swallowed an error)" def test_b9_logs_dir_env_override_is_dynamic(mam_sandbox, tmp_path, monkeypatch): """B-9: DELEGATE_JOB_LOGS_DIR must be honoured at call time, both ways.""" mq = get_mqtt_common(mam_sandbox) monkeypatch.chdir(tmp_path) monkeypatch.setenv("DELEGATE_JOB_LOGS_DIR", "/tmp/b9-override") assert mq.get_logs_dir() == "/tmp/b9-override" monkeypatch.delenv("DELEGATE_JOB_LOGS_DIR") # equality, not inequality — clearing the env must restore the cwd default (Rev.2 §3) assert os.path.realpath(mq.get_logs_dir()) == \ os.path.realpath(os.path.join(str(tmp_path), ".mam", "delegate_job_logs")) def test_b9_no_module_level_logs_dir_binding(): """B-9/T1: a surviving module global would make __getattr__ dead code.""" import ast, pathlib src = (pathlib.Path(__file__).resolve().parent.parent / ".agents" / "skills" / "multi-agent-mux-delegate-job" / "scripts" / "mqtt_common.py") tree = ast.parse(src.read_text()) for node in tree.body: # module scope only if isinstance(node, ast.Assign): for t in node.targets: assert not (isinstance(t, ast.Name) and t.id == "LOGS_DIR"), \ f"line {node.lineno}: module-level LOGS_DIR binding shadows __getattr__ (B-9/T1)" elif isinstance(node, ast.AnnAssign): # C2-b: LOGS_DIR: str = ... parses as AnnAssign assert not (isinstance(node.target, ast.Name) and node.target.id == "LOGS_DIR"), \ f"line {node.lineno}: annotated module-level LOGS_DIR binding shadows __getattr__ (B-9/T1)" def test_b9_logs_dir_stays_discoverable(mam_sandbox): """B-9/C2-a: PEP 562 __dir__ keeps LOGS_DIR visible to dir() and tooling.""" mq = get_mqtt_common(mam_sandbox) assert "LOGS_DIR" in dir(mq) assert hasattr(mq, "LOGS_DIR") # true via __getattr__ even without __dir__ assert dir(mq).count("LOGS_DIR") == 1 # set-based __dir__ must not duplicate # ============================================================================== # Track 0: Fault Tolerance & Local Fallback Regression Guards (G-1 to G-10) # ============================================================================== def _save_job_for_test(job_rec, registry_dir): p = os.path.join(registry_dir, f"{job_rec['job_id']}.json") with open(p, "w", encoding="utf-8") as f: json.dump(job_rec, f, indent=2) def test_g1_publish_failure_persists_completed_status(mam_sandbox, monkeypatch): """G-1: When publish fails due to broker unreachable, registry status is still updated.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, publish_event reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id = registry.register_job("Test G-1", registry_dir=reg_dir, job_id="g1test01") job = registry.load_job(job_id, reg_dir) # Point broker to non-routable dummy port job["broker"] = {"host": "127.0.0.1", "port": 65432, "tls": False} _save_job_for_test(job, reg_dir) rc = publish_event.main([ "--registry-dir", reg_dir, "--job", "g1test01", "--event", "completed", "--detail", "finished work", "--attempts", "1" ]) assert rc == 2, f"Expected rc=2 on network failure, got {rc}" loaded = registry.load_job("g1test01", reg_dir) assert loaded["status"] == "completed", "Registry status must be completed despite network publish failure" def test_g2_publish_failure_records_audit_log_error(mam_sandbox, monkeypatch): """G-2: Audit log records published=False and publish_error when broker is unreachable.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, publish_event, mqtt_common reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id = registry.register_job("Test G-2", registry_dir=reg_dir, job_id="g2test01") job = registry.load_job(job_id, reg_dir) job["broker"] = {"host": "127.0.0.1", "port": 65432, "tls": False} _save_job_for_test(job, reg_dir) rc = publish_event.main([ "--registry-dir", reg_dir, "--job", "g2test01", "--event", "completed", "--attempts", "1" ]) assert rc == 2 log_file = os.path.join(str(mam_sandbox), ".mam", "delegate_job_logs", "g2test01", "events.ndjson") assert os.path.exists(log_file), "Audit log file must exist" with open(log_file, "r", encoding="utf-8") as f: events = [json.loads(line) for line in f if line.strip()] pub_events = [e for e in events if e.get("event") == "published" and e.get("source_event") == "completed"] assert len(pub_events) == 1 assert pub_events[0]["published"] is False assert pub_events[0]["publish_error"] is not None def test_g3_publish_success_records_published_true(mam_sandbox, monkeypatch): """G-3: When publish succeeds, rc=0, status=completed, and published=True.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, publish_event, mqtt_common reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) registry.register_job("Test G-3", registry_dir=reg_dir, job_id="g3test01") # Mock _publish_once to simulate broker success monkeypatch.setattr(publish_event, "_publish_once", lambda *args, **kwargs: None) rc = publish_event.main([ "--registry-dir", reg_dir, "--job", "g3test01", "--event", "completed", "--attempts", "1" ]) assert rc == 0 loaded = registry.load_job("g3test01", reg_dir) assert loaded["status"] == "completed" log_file = os.path.join(str(mam_sandbox), ".mam", "delegate_job_logs", "g3test01", "events.ndjson") assert os.path.exists(log_file), f"Audit log file {log_file} must exist" with open(log_file, "r", encoding="utf-8") as f: events = [json.loads(line) for line in f if line.strip()] pub_events = [e for e in events if e.get("event") == "published" and e.get("source_event") == "completed"] assert len(pub_events) == 1 assert pub_events[0]["published"] is True assert pub_events[0]["publish_error"] is None def test_g4_publish_failure_advances_sequence(mam_sandbox, monkeypatch): """G-4: A failed publish consumes sequence number, and subsequent publish uses strictly higher seq.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, publish_event reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id = registry.register_job("Test G-4", registry_dir=reg_dir, job_id="g4test01") job = registry.load_job(job_id, reg_dir) job["broker"] = {"host": "127.0.0.1", "port": 65432, "tls": False} _save_job_for_test(job, reg_dir) # First attempt fails rc1 = publish_event.main([ "--registry-dir", reg_dir, "--job", "g4test01", "--event", "progress", "--attempts", "1" ]) assert rc1 == 2 j1 = registry.load_job("g4test01", reg_dir) assert int(j1["last_seq"]) == 1 # Second attempt succeeds (mocked) monkeypatch.setattr(publish_event, "_publish_once", lambda *args, **kwargs: None) rc2 = publish_event.main([ "--registry-dir", reg_dir, "--job", "g4test01", "--event", "completed", "--attempts", "1" ]) assert rc2 == 0 j2 = registry.load_job("g4test01", reg_dir) assert int(j2["last_seq"]) == 2 def test_g5_subscriber_disk_fallback_on_broker_down_exits_0(mam_sandbox, monkeypatch): """G-5: When broker is down and disk has status=completed, job_subscriber exits 0 via fallback.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, job_subscriber reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id = registry.register_job("Test G-5", registry_dir=reg_dir, job_id="g5test01") job = registry.load_job(job_id, reg_dir) job["broker"] = {"host": "127.0.0.1", "port": 65432, "tls": False} job["status"] = "completed" _save_job_for_test(job, reg_dir) rc = job_subscriber.main([ "--registry-dir", reg_dir, "--job", "g5test01", "--timeout", "5", "--idle-timeout", "2" ]) assert rc == 0 def test_g6_subscriber_disk_fallback_outputs_source_tag(mam_sandbox, capsys, monkeypatch): """G-6: When resolving via disk fallback, stdout contains 'disk-fallback' tag.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, job_subscriber reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id = registry.register_job("Test G-6", registry_dir=reg_dir, job_id="g6test01") job = registry.load_job(job_id, reg_dir) job["broker"] = {"host": "127.0.0.1", "port": 65432, "tls": False} job["status"] = "completed" _save_job_for_test(job, reg_dir) rc = job_subscriber.main([ "--registry-dir", reg_dir, "--job", "g6test01", "--timeout", "5" ]) assert rc == 0 captured = capsys.readouterr() assert "disk-fallback" in captured.out def test_g7_subscriber_disk_fallback_error_status_exits_1(mam_sandbox, monkeypatch): """G-7: When broker is down and disk has status=error, job_subscriber exits 1.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, job_subscriber reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id = registry.register_job("Test G-7", registry_dir=reg_dir, job_id="g7test01") job = registry.load_job(job_id, reg_dir) job["broker"] = {"host": "127.0.0.1", "port": 65432, "tls": False} job["status"] = "error" _save_job_for_test(job, reg_dir) rc = job_subscriber.main([ "--registry-dir", reg_dir, "--job", "g7test01", "--timeout", "5" ]) assert rc == 1 def test_g8_subscriber_wait_any_does_not_exit_early_if_partial_pending(mam_sandbox, monkeypatch): """G-8: Multi-job wait does not terminate early when only 1 of 2 jobs is terminal on disk.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, job_subscriber, mqtt_common reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id1 = registry.register_job("Test G-8 Job 1", registry_dir=reg_dir, job_id="g8test01") job_id2 = registry.register_job("Test G-8 Job 2", registry_dir=reg_dir, job_id="g8test02") job1 = registry.load_job(job_id1, reg_dir) job2 = registry.load_job(job_id2, reg_dir) job1["status"] = "completed" job2["status"] = "running" _save_job_for_test(job1, reg_dir) _save_job_for_test(job2, reg_dir) # Mock client connection to succeed with empty queue class FakeClient: on_message = None on_connect = None on_disconnect = None on_subscribe = None def reconnect_delay_set(self, **kwargs): pass def connect(self, *args, **kwargs): pass def loop_start(self): pass def loop_stop(self): pass def disconnect(self): pass def subscribe(self, *args, **kwargs): pass monkeypatch.setattr(job_subscriber, "make_client", lambda *args, **kwargs: FakeClient()) monkeypatch.setattr(mqtt_common, "with_retry", lambda fn, **kwargs: fn) rc = job_subscriber.main([ "--registry-dir", reg_dir, "--wait-any", "--timeout", "0.5", "--idle-timeout", "0.5" ]) assert rc == 2, "Must timeout waiting for incomplete job2 rather than exiting 0" def test_g9_subscriber_broker_down_no_disk_terminal_exits_3(mam_sandbox, monkeypatch): """G-9: When broker is unreachable and no terminal state exists on disk, subscriber exits with rc=3.""" monkeypatch.chdir(mam_sandbox) script_dir = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "scripts" sys.path.insert(0, str(script_dir)) import registry, job_subscriber reg_dir = str(mam_sandbox / ".mam" / "jobs") os.makedirs(reg_dir, exist_ok=True) job_id = registry.register_job("Test G-9", registry_dir=reg_dir, job_id="g9test01") job = registry.load_job(job_id, reg_dir) job["broker"] = {"host": "127.0.0.1", "port": 65432, "tls": False} job["status"] = "running" _save_job_for_test(job, reg_dir) rc = job_subscriber.main([ "--registry-dir", reg_dir, "--job", "g9test01", "--timeout", "5" ]) assert rc == 3, f"Expected rc=3 on infrastructure broker down without disk terminal, got {rc}" def test_g10_delegate_job_rc3_not_mistaken_for_error(mam_sandbox): """G-10: multi-agent-mux-delegate-job maps rc=3 to broker_unavailable and checks disk status.""" delegate_script = mam_sandbox / "skills" / "multi-agent-mux-delegate-job" / "multi-agent-mux-delegate-job" content = delegate_script.read_text() assert "elif [[ $sub_rc -eq 3 ]]; then" in content assert 'job_status="broker_unavailable"' in content