diff --git a/bin/modulate.py b/bin/modulate.py new file mode 100644 index 0000000..b5af7e2 --- /dev/null +++ b/bin/modulate.py @@ -0,0 +1,622 @@ +"""Input modulation for the DM follow-up system. + +Not all inputs deserve the same follow-up policy. A routine timer wake +should not nudge like a security alert, and a heartbeat loopback should +never create a record at all. + +The model is three layers: + + 1. INPUT TYPE declares what kind of traffic this is + (wake, job, siphon, manual, health, heartbeat) + 2. MODULATION looks up the policy table for that type (+subtype) + and returns timeout / nudges / escalate / priority / track-mode + 3. SEND-TIME TAGS remain the transport. The modulated policy renders + into the canonical tag vocabulary ([reply:expected], + [reply:timeout=N], [reply:nudges=N], [reply:escalate=X], + [input:]). Explicitly declared tags ALWAYS win over the + table — a human who says "nudge me in 15m" is sovereign. + +Track modes: "always" (create the record), "never" (don't, ever), +"actionable" (create only when the sender marks the input actionable — +used by wake digests, where an informational digest creates no loop but +a tracked wake does). + +Priorities feed dashboard sorting and nudge urgency text; they do not +change the sweep mechanics. +""" + +from __future__ import annotations + +import json +import os +import sys +import time +from dataclasses import dataclass, field +from enum import Enum +from pathlib import Path +from typing import Dict, List, Optional, Tuple, Any + +DEFAULT_STRATEGY_FILE = "/srv/box/strategy.json" +FALLBACK_STRATEGY_FILE = "/home/super/Projects/NetVM/strategy.json" + + + +# --------------------------------------------------------------------------- +# Enumerations +# --------------------------------------------------------------------------- + +class InputType(str, Enum): + WAKE = "wake" # timer-driven operator digests + JOB = "job" # job dispatch expecting [RESULT] + SIPHON = "siphon" # side-chat -> main-chat surfacing + MANUAL = "manual" # human/operator-authored DM + HEALTH = "health" # system health check results + HEARTBEAT = "heartbeat" # loopback liveness probes + + +class Priority(str, Enum): + ROUTINE = "routine" + NORMAL = "normal" + IMPORTANT = "important" + CRITICAL = "critical" + + +# Priority rank for sorting (higher = more urgent). +_PRIORITY_RANK = { + Priority.ROUTINE: 0, + Priority.NORMAL: 1, + Priority.IMPORTANT: 2, + Priority.CRITICAL: 3, +} + + +# --------------------------------------------------------------------------- +# Policy +# --------------------------------------------------------------------------- + +# Timeout clamps inherited from JOB-FOLLOWUP.md: 60s..7d. +MIN_TIMEOUT_S = 60 +MAX_TIMEOUT_S = 604800 +MAX_NUDGES = 10 + + +@dataclass(frozen=True) +class Policy: + """A modulated follow-up policy, ready to render into tags.""" + input_type: InputType + subtype: Optional[str] # e.g. siphon category, health status + priority: Priority + track: bool # create a dm_followup record? + timeout_s: int # seconds to first nudge + nudges: int # max auto-nudges + escalate: Optional[str] # identity id, or None (silent close) + # Which explicit values overrode the table (audit trail). + overridden: Tuple[str, ...] = field(default_factory=tuple) + + def rank(self) -> int: + return _PRIORITY_RANK[self.priority] + + @property + def expect_reply(self) -> bool: + return bool(self.track) + + +# --------------------------------------------------------------------------- +# The modulation table. +# +# Key: (InputType, subtype or None). Subtype None = default for that type. +# track: True always / False never / "actionable" (wake digests). +# --------------------------------------------------------------------------- + +# (track, priority, timeout_s, nudges, escalate) +_MODULATION: Dict[Tuple[InputType, Optional[str]], Tuple] = { + # Timer wakes: routine, batchy. Tracked only when the digest is + # actionable (the wake-actionable pilot); otherwise informational. + (InputType.WAKE, None): ("actionable", Priority.ROUTINE, 2 * 3600, 1, "opm"), + (InputType.WAKE, "overdue"): (True, Priority.IMPORTANT, 3600, 2, "opm"), + + # Job dispatches: the canonical "waiting on someone" loop. + (InputType.JOB, None): (True, Priority.NORMAL, 3600, 2, "opm"), + (InputType.JOB, "canary"): (True, Priority.NORMAL, 900, 2, "opm"), + (InputType.JOB, "chain"): (True, Priority.IMPORTANT, 1800, 2, "opm"), + + # Siphon hits: urgency follows the detected category. + (InputType.SIPHON, "ALERT"): (True, Priority.CRITICAL, 600, 3, "user"), + (InputType.SIPHON, "BLOCKER"): (True, Priority.CRITICAL, 900, 2, "opm"), + (InputType.SIPHON, "DECISION"): (True, Priority.IMPORTANT, 3600, 2, "opm"), + (InputType.SIPHON, "COMPLETED"): (False, Priority.ROUTINE, 0, 0, None), + (InputType.SIPHON, "MILESTONE"): (False, Priority.ROUTINE, 0, 0, None), + (InputType.SIPHON, None): (True, Priority.NORMAL, 3600, 2, "opm"), + + # Manual DMs: sender declares; table only supplies the default. + (InputType.MANUAL, None): (True, Priority.NORMAL, 3600, 2, "opm"), + (InputType.MANUAL, "urgent"): (True, Priority.IMPORTANT, 900, 2, "opm"), + (InputType.MANUAL, "fyi"): (True, Priority.ROUTINE, 7200, 1, None), + + # Health: system-level, fast fuse. + (InputType.HEALTH, "FAIL"): (True, Priority.CRITICAL, 600, 2, "opm"), + (InputType.HEALTH, "DEGRADED"): (True, Priority.IMPORTANT, 1800, 2, "opm"), + (InputType.HEALTH, "OK"): (False, Priority.ROUTINE, 0, 0, None), + (InputType.HEALTH, None): (True, Priority.IMPORTANT, 1800, 2, "opm"), + + # Heartbeat loopback: NEVER tracked. Hard exclusion. + (InputType.HEARTBEAT, None): (False, Priority.ROUTINE, 0, 0, None), +} + + +def _clamp_timeout(s: int) -> int: + return max(MIN_TIMEOUT_S, min(MAX_TIMEOUT_S, s)) + + +def _clamp_nudges(n: int) -> int: + return max(0, min(MAX_NUDGES, n)) + + +def _strategy_file() -> str: + if os.path.exists(DEFAULT_STRATEGY_FILE): + return DEFAULT_STRATEGY_FILE + if os.path.exists(FALLBACK_STRATEGY_FILE): + return FALLBACK_STRATEGY_FILE + if os.path.exists("/srv/box"): + return DEFAULT_STRATEGY_FILE + return FALLBACK_STRATEGY_FILE + + +def load_strategy_overrides() -> Dict[str, dict]: + path = _strategy_file() + try: + with open(path) as f: + data = json.load(f) + return data.get("overrides", {}) + except Exception: + return {} + + +def _save_overrides(overrides: Dict[str, dict], by: Optional[str] = None): + payload = { + "_meta": { + "version": 1, + "updated_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), + "updated_by": by or "super", + }, + "overrides": overrides, + } + for dest in (DEFAULT_STRATEGY_FILE, FALLBACK_STRATEGY_FILE): + try: + os.makedirs(os.path.dirname(os.path.abspath(dest)), exist_ok=True) + tmp = f"{dest}.tmp.{os.getpid()}" + with open(tmp, "w") as f: + json.dump(payload, f, indent=2) + f.write("\n") + os.chmod(tmp, 0o644) + os.replace(tmp, dest) + except Exception: + pass + + +VALID_AGENTS = ("muse", "pip", "646", "opm") + + +def get_strategy_row(input_type: InputType, subtype: Optional[str] = None, agent: Optional[str] = None) -> Tuple: + """Lookup modulation row with hierarchical runtime overrides applied. + + Resolution hierarchy: + 1. Exact match: (input_type, subtype, agent) + 2. Agent default for type: (input_type, None, agent) + 3. Subtype default: (input_type, subtype, None) + 4. Type default: (input_type, None, None) + 5. Builtin table fallback + """ + overrides = load_strategy_overrides() + ag = agent.lower() if agent else None + sub = subtype.lower() if subtype else None + t = input_type.value + + # 1. Exact match: type:subtype:agent + if sub and ag: + key = f"{t}:{sub}:{ag}" + if key in overrides: + o = overrides[key] + return (o.get("track", True), Priority(o.get("priority", "normal")), + int(o.get("timeout_s", 3600)), int(o.get("nudges", 2)), o.get("escalate")) + + # 2. Agent default for type: type:agent + if ag: + key = f"{t}:{ag}" + if key in overrides: + o = overrides[key] + return (o.get("track", True), Priority(o.get("priority", "normal")), + int(o.get("timeout_s", 3600)), int(o.get("nudges", 2)), o.get("escalate")) + + # 3. Subtype default: type:subtype + if sub: + key = f"{t}:{sub}" + if key in overrides: + o = overrides[key] + return (o.get("track", True), Priority(o.get("priority", "normal")), + int(o.get("timeout_s", 3600)), int(o.get("nudges", 2)), o.get("escalate")) + + # 4. Type default: type + if t in overrides: + o = overrides[t] + return (o.get("track", True), Priority(o.get("priority", "normal")), + int(o.get("timeout_s", 3600)), int(o.get("nudges", 2)), o.get("escalate")) + + # 5. Builtin fallback (case-insensitive subtype matching) + row = None + if subtype: + for (it, st), r in _MODULATION.items(): + if it == input_type and st and st.lower() == subtype.lower(): + row = r + break + if row is None: + row = _MODULATION.get((input_type, None)) + if row is None: + row = _MODULATION[(InputType.MANUAL, None)] + return row + + +def get_all_strategies() -> List[dict]: + """Return all strategies (builtins merged with hierarchical overrides).""" + overrides = load_strategy_overrides() + out = [] + seen = set() + + for itype, sub in sorted(_MODULATION.keys(), key=lambda x: (x[0].value, x[1] or "")): + key_str = f"{itype.value}:{sub.lower()}" if sub else itype.value + seen.add(key_str) + is_override = key_str in overrides + row = get_strategy_row(itype, sub) + track_mode, prio, timeout_s, nudges, escalate = row + out.append({ + "key": key_str, + "input_type": itype.value, + "subtype": sub, + "agent": None, + "track": track_mode, + "priority": prio.value if isinstance(prio, Priority) else str(prio), + "timeout_s": timeout_s, + "nudges": nudges, + "escalate": escalate, + "is_override": is_override, + }) + + # Include any custom override keys not in builtins + for key_str, o in overrides.items(): + if key_str not in seen: + parts = key_str.split(":") + itype_val = parts[0] + sub_val = None + agent_val = None + + if len(parts) == 3: + sub_val = parts[1].upper() + agent_val = parts[2] + elif len(parts) == 2: + if parts[1].lower() in VALID_AGENTS: + agent_val = parts[1].lower() + else: + sub_val = parts[1].upper() + + prio = o.get("priority", "normal") + out.append({ + "key": key_str, + "input_type": itype_val, + "subtype": sub_val, + "agent": agent_val, + "track": o.get("track", True), + "priority": prio, + "timeout_s": int(o.get("timeout_s", 3600)), + "nudges": int(o.get("nudges", 2)), + "escalate": o.get("escalate"), + "is_override": True, + }) + return out + + +def _build_key(itype: str, sub: Optional[str] = None, agent: Optional[str] = None) -> str: + t = itype.lower().strip() + s = sub.lower().strip() if sub and sub != "*" else None + a = agent.lower().strip() if agent and agent != "*" else None + if s and a: + return f"{t}:{s}:{a}" + if a and not s: + return f"{t}:{a}" + if s and not a: + return f"{t}:{s}" + return t + + +def set_strategy_override( + input_type_str: str, + subtype_str: Optional[str] = None, + agent_str: Optional[str] = None, + *, + track: Any = None, + priority: Optional[str] = None, + timeout_s: Optional[int] = None, + nudges: Optional[int] = None, + escalate: Optional[str] = None, + by: Optional[str] = None, +) -> dict: + """Set an external strategy override with optional agent scope.""" + overrides = load_strategy_overrides() + key = _build_key(input_type_str, subtype_str, agent_str) + + # Get current values as baseline + try: + itype = InputType(input_type_str.lower()) + except ValueError: + itype = InputType.MANUAL + current_row = get_strategy_row(itype, subtype_str.upper() if subtype_str else None, agent=agent_str) + c_track, c_prio, c_timeout, c_nudges, c_esc = current_row + + new_entry = dict(overrides.get(key, {})) + new_entry["key"] = key + new_entry["input_type"] = itype.value + new_entry["subtype"] = subtype_str.upper() if subtype_str else None + new_entry["agent"] = agent_str.lower() if agent_str else None + + if track is not None: + if isinstance(track, str) and track.lower() == "true": + new_entry["track"] = True + elif isinstance(track, str) and track.lower() == "false": + new_entry["track"] = False + else: + new_entry["track"] = track + else: + new_entry.setdefault("track", c_track) + + if priority is not None: + new_entry["priority"] = Priority(priority.lower()).value + else: + new_entry.setdefault("priority", c_prio.value if isinstance(c_prio, Priority) else str(c_prio)) + + if timeout_s is not None: + new_entry["timeout_s"] = _clamp_timeout(int(timeout_s)) + else: + new_entry.setdefault("timeout_s", c_timeout) + + if nudges is not None: + new_entry["nudges"] = _clamp_nudges(int(nudges)) + else: + new_entry.setdefault("nudges", c_nudges) + + if escalate is not None: + new_entry["escalate"] = escalate if escalate != "none" else None + else: + new_entry.setdefault("escalate", c_esc) + + overrides[key] = new_entry + _save_overrides(overrides, by=by) + return new_entry + + +def reset_strategy_override(input_type_str: str, subtype_str: Optional[str] = None, agent_str: Optional[str] = None, by: Optional[str] = None) -> bool: + """Reset a strategy override to built-in default.""" + overrides = load_strategy_overrides() + key = _build_key(input_type_str, subtype_str, agent_str) + if key in overrides: + del overrides[key] + _save_overrides(overrides, by=by) + return True + return False + + +def set_override(input_type, subtype=None, agent=None, **kw): + itype_str = input_type.value if hasattr(input_type, "value") else str(input_type) + return set_strategy_override(itype_str, subtype, agent, **kw) + + +def reset_override(input_type, subtype=None, agent=None, by=None): + itype_str = input_type.value if hasattr(input_type, "value") else str(input_type) + return reset_strategy_override(itype_str, subtype, agent, by=by) + + +def evaluate_strategy(input_type, subtype=None, agent=None, **kw): + if isinstance(input_type, str): + try: + input_type = InputType(input_type) + except ValueError: + input_type = InputType.MANUAL + return modulate(input_type, subtype=subtype, agent=agent, **kw) + + +# --------------------------------------------------------------------------- +# Public API +# --------------------------------------------------------------------------- + +def modulate( + input_type: InputType, + subtype: Optional[str] = None, + agent: Optional[str] = None, + *, + actionable: bool = False, + timeout_s: Optional[int] = None, + nudges: Optional[int] = None, + escalate: Optional[str] = None, + priority: Optional[Priority] = None, + track: Optional[bool] = None, +) -> Optional[Policy]: + """Compute the follow-up policy for one send. + + Returns None when the input should NOT be tracked (track=False, or + track="actionable" with actionable=False). + + Explicit keyword args always override the table (send-time + sovereignty); they are recorded in Policy.overridden. + Unknown input types fail closed to the MANUAL default. + """ + # Fail closed: unknown type -> manual default. + if not isinstance(input_type, InputType): + input_type = InputType.MANUAL + + row = get_strategy_row(input_type, subtype, agent=agent) + track_mode, prio, d_timeout, d_nudges, d_escalate = row + + do_track = track if track is not None else ( + True if track_mode is True + else False if track_mode is False + else actionable # "actionable" + ) + if not do_track: + return None + + overridden = [] + final_timeout = d_timeout + if timeout_s is not None: + final_timeout = _clamp_timeout(timeout_s) + overridden.append("timeout") + final_nudges = d_nudges + if nudges is not None: + final_nudges = _clamp_nudges(nudges) + overridden.append("nudges") + final_escalate = d_escalate + if escalate is not None: + final_escalate = escalate or None + overridden.append("escalate") + final_prio = priority if priority is not None else prio + if priority is not None: + overridden.append("priority") + + return Policy( + input_type=input_type, + subtype=subtype, + priority=final_prio, + track=True, + timeout_s=final_timeout, + nudges=final_nudges, + escalate=final_escalate, + overridden=tuple(overridden), + ) + + +def render_tags(policy: Policy) -> str: + """Render a policy into the canonical bracket-tag vocabulary. + + Emits [reply:expected] [reply:timeout=N] [reply:nudges=N] + [reply:escalate=X] [input:]. The input tag is the audit + trail: it records which modulation row drove the policy. + """ + parts = ["[reply:expected]"] + parts.append(f"[reply:timeout={policy.timeout_s}]") + parts.append(f"[reply:nudges={policy.nudges}]") + if policy.escalate: + parts.append(f"[reply:escalate={policy.escalate}]") + parts.append(f"[input:{policy.input_type.value}]") + return " ".join(parts) + + +# --------------------------------------------------------------------------- +# Convenience constructors for the known producers. +# --------------------------------------------------------------------------- + +def for_siphon_hit(category: str, **kw) -> Optional[Policy]: + """Policy for a siphon hit; category is the detect.py category.""" + return modulate(InputType.SIPHON, category.upper(), **kw) + + +def for_health(status: str, **kw) -> Optional[Policy]: + """Policy for a health result; status in {OK, DEGRADED, FAIL}.""" + return modulate(InputType.HEALTH, status.upper(), **kw) + + +def for_job(job_name: str, **kw) -> Optional[Policy]: + """Policy for a job dispatch. Canary jobs get the short fuse.""" + subtype = "canary" if "canary" in job_name.lower() else None + return modulate(InputType.JOB, subtype, **kw) + + +def for_wake(actionable: bool = False, has_overdue: bool = False, **kw) -> Optional[Policy]: + """Policy for a timer wake digest.""" + subtype = "overdue" if has_overdue else None + return modulate(InputType.WAKE, subtype, actionable=actionable, **kw) + + +def sort_key(policy: Policy): + """Dashboard sort: critical first, then soonest timeout.""" + return (-policy.rank(), policy.timeout_s) + + +def main(): + import argparse + parser = argparse.ArgumentParser(description="Intrinsic Loop Modulation Strategy CLI") + sub = parser.add_subparsers(dest="action") + + sub.add_parser("list", help="List all strategies") + + p_get = sub.add_parser("get", help="Get strategy row") + p_get.add_argument("type", help="Input type (wake, job, siphon, manual, health, heartbeat)") + p_get.add_argument("subtype", nargs="?", default=None, help="Optional subtype") + + p_set = sub.add_parser("set", help="Set strategy override") + p_set.add_argument("type", help="Input type") + p_set.add_argument("--subtype", default=None, help="Subtype") + p_set.add_argument("--track", choices=["true", "false", "actionable", "always", "never"], default=None) + p_set.add_argument("--priority", choices=["routine", "normal", "important", "critical"], default=None) + p_set.add_argument("--timeout", type=int, default=None, help="Timeout in seconds") + p_set.add_argument("--nudges", type=int, default=None, help="Max auto nudges") + p_set.add_argument("--escalate", default=None, help="Escalation target agent (e.g. opm, user, none)") + + p_reset = sub.add_parser("reset", help="Reset strategy override") + p_reset.add_argument("type", help="Input type") + p_reset.add_argument("--subtype", default=None, help="Subtype") + + p_eval = sub.add_parser("eval", help="Evaluate modulation policy for inputs") + p_eval.add_argument("type", help="Input type") + p_eval.add_argument("--subtype", default=None, help="Subtype") + p_eval.add_argument("--actionable", action="store_true", help="Mark actionable (for wake)") + + args = parser.parse_args() + + if not args.action or args.action == "list": + strats = get_all_strategies() + print(f"\n{'KEY':20} {'TRACK':10} {'PRIORITY':10} {'TIMEOUT':10} {'NUDGES':8} {'ESCALATE':10} {'SOURCE':8}") + print("─" * 80) + for s in strats: + src = "OVERRIDE" if s.get("is_override") else "BUILTIN" + esc = s.get("escalate") or "-" + print(f"{s['key']:20} {str(s['track']):10} {s['priority']:10} {str(s['timeout_s'])+'s':10} {str(s['nudges']):8} {esc:10} {src:8}") + print() + elif args.action == "get": + try: + itype = InputType(args.type.lower()) + except ValueError: + print(f"Unknown input type: {args.type}", file=sys.stderr) + sys.exit(1) + row = get_strategy_row(itype, args.subtype.upper() if args.subtype else None) + print(json.dumps({ + "type": itype.value, + "subtype": args.subtype, + "track": row[0], + "priority": row[1].value if isinstance(row[1], Priority) else str(row[1]), + "timeout_s": row[2], + "nudges": row[3], + "escalate": row[4], + }, indent=2)) + elif args.action == "set": + res = set_strategy_override( + args.type, args.subtype, + track=args.track, priority=args.priority, + timeout_s=args.timeout, nudges=args.nudges, escalate=args.escalate, + by=os.environ.get("BOX_CALLER", "cli") + ) + print(f"Updated strategy override for {args.type}:{args.subtype or '*'}: {json.dumps(res)}") + elif args.action == "reset": + ok = reset_strategy_override(args.type, args.subtype, by=os.environ.get("BOX_CALLER", "cli")) + print(f"Reset {args.type}:{args.subtype or '*'}: {'Success' if ok else 'Not overridden'}") + elif args.action == "eval": + try: + itype = InputType(args.type.lower()) + except ValueError: + itype = InputType.MANUAL + pol = modulate(itype, args.subtype.upper() if args.subtype else None, actionable=args.actionable) + if pol: + print(f"Policy: {pol}") + print(f"Tags : {render_tags(pol)}") + else: + print("Policy: None (Not tracked)") + + +if __name__ == "__main__": + main() diff --git a/bin/variables.py b/bin/variables.py new file mode 100644 index 0000000..208413d --- /dev/null +++ b/bin/variables.py @@ -0,0 +1,541 @@ +#!/usr/bin/env python3 +"""Box variable registry — single library for box-defined control variables. + +Timers, wake digests, the siphon, the sweeper, and the silence detector all +read their thresholds here at runtime instead of hardcoding them. The user +adjusts behavior by changing variables in /srv/box/variables.json (via the +Box API or variablesctl), never by editing code. + +Registry location: + - Production: /srv/box/variables.json + - Override: BOX_VARIABLES env var, or explicit path argument + +Fail-closed design: a missing file, a corrupt file, or a missing variable +falls back to the built-in default for that variable — timers keep working +and the miss is recorded in `warnings`. + +Setting a variable validates: + - type matches (int/float/bool/str) + - min/max or enum constraints honored + - unknown variable names are rejected (typo guard) + +The registry file format: + + {"_meta": {...}, + "variables": { + "max_nudge_count": { + "value": 2, "default": 2, "type": "int", + "min": 0, "max": 10, "unit": "count", + "description": "..." + }}} + +Writes are atomic (tmp + rename). `updated_at` / `updated_by` in _meta track +who changed what. +""" + +import json +import os +import time + +# Primary and fallback registries +DEFAULT_REGISTRY = "/srv/box/variables.json" +FALLBACK_REGISTRY = "/home/super/Projects/NetVM/variables.json" +DEFAULT_HISTORY = "/srv/box/variables-history.jsonl" +FALLBACK_HISTORY = "/home/super/Projects/NetVM/variables-history.jsonl" + +ENV_OVERRIDE = "BOX_VARIABLES" + + +class VariableError(ValueError): + """Raised on unknown variable, type mismatch, or constraint violation.""" + pass + + +class Variables: + """Typed access to the box variable registry.""" + + def __init__(self, path=None, history_path=None, _data=None): + env_path = os.environ.get(ENV_OVERRIDE) + self.history_path = history_path + if path: + self.path = path + elif env_path: + + self.path = env_path + elif os.path.exists(DEFAULT_REGISTRY): + self.path = DEFAULT_REGISTRY + elif os.path.exists(FALLBACK_REGISTRY): + self.path = FALLBACK_REGISTRY + else: + self.path = DEFAULT_REGISTRY + + self.warnings = [] + if _data is not None: + self._data = _data + else: + self._data = self._load() + + def _load(self): + loaded = {} + target_path = self.path + if not os.path.exists(target_path) and os.path.exists(FALLBACK_REGISTRY): + target_path = FALLBACK_REGISTRY + + try: + with open(target_path) as f: + data = json.load(f) + if isinstance(data, dict): + loaded = data.get("variables", {}) + except FileNotFoundError: + self.warnings.append("registry not found: %s" % target_path) + except (OSError, ValueError) as e: + self.warnings.append("registry unreadable (%s): using defaults" % e) + + # Seed with builtin specifications if missing + specs = {} + for k, s in _BUILTIN_SPECS.items(): + specs[k] = dict(s) + if "value" not in specs[k]: + specs[k]["value"] = s["default"] + + # Overlay any loaded values from file + if isinstance(loaded, dict): + for k, item in loaded.items(): + if k in specs and isinstance(item, dict): + specs[k].update(item) + elif isinstance(item, dict): + specs[k] = item + + return specs + + def _spec(self, name): + spec = self._data.get(name) + if spec is None: + if name in _BUILTIN_SPECS: + self._data[name] = dict(_BUILTIN_SPECS[name]) + self._data[name]["value"] = self._data[name]["default"] + return self._data[name] + raise VariableError("unknown variable: %s" % name) + return spec + + def _coerce(self, name, value, spec): + vtype = spec.get("type", "int") + try: + if vtype == "int": + if isinstance(value, bool): + raise VariableError("%s: expected int, got %r" % (name, value)) + if isinstance(value, str): + try: + value = int(value) + except ValueError: + raise VariableError("%s: expected int, got %r" % (name, value)) + elif not isinstance(value, (int, float)): + raise VariableError("%s: expected int, got %r" % (name, value)) + value = int(value) + elif vtype == "float": + if isinstance(value, bool): + raise VariableError("%s: expected float, got %r" % (name, value)) + if isinstance(value, str): + try: + value = float(value) + except ValueError: + raise VariableError("%s: expected float, got %r" % (name, value)) + elif not isinstance(value, (int, float)): + raise VariableError("%s: expected float, got %r" % (name, value)) + value = float(value) + + elif vtype == "bool": + if not isinstance(value, bool): + raise VariableError( + "%s: expected bool, got %r" % (name, value)) + elif vtype == "str": + if not isinstance(value, str): + raise VariableError( + "%s: expected str, got %r" % (name, value)) + else: + raise VariableError("%s: unknown type %r" % (name, vtype)) + except VariableError: + raise + except (TypeError, ValueError): + raise VariableError("%s: cannot coerce %r to %s" % (name, value, vtype)) + return value + + def _check_constraints(self, name, value, spec): + enum = spec.get("enum") + if enum is not None and value not in enum: + raise VariableError( + "%s: %r not in allowed values %r" % (name, value, enum)) + for bound, key in (("min", "min"), ("max", "max")): + limit = spec.get(key) + if limit is not None and isinstance(value, (int, float)): + if bound == "min" and value < limit: + raise VariableError( + "%s: %r below minimum %r" % (name, value, limit)) + if bound == "max" and value > limit: + raise VariableError( + "%s: %r above maximum %r" % (name, value, limit)) + + # -- reads ----------------------------------------------------------- + + def get(self, name): + """Return the current value; fall back to the built-in default when + the registry is missing the variable (never raises on read).""" + spec = self._data.get(name) + if spec is None: + self.warnings.append("variable %s missing: using default" % name) + return _BUILTIN_DEFAULTS.get(name) + value = spec.get("value", spec.get("default")) + if value is None: + value = _BUILTIN_DEFAULTS.get(name) + return value + + def get_int(self, name): + return int(self.get(name)) + + def get_float(self, name): + return float(self.get(name)) + + def get_bool(self, name): + return bool(self.get(name)) + + def get_str(self, name): + return str(self.get(name)) + + def all(self): + """Dict of name -> value for every registered variable.""" + return {name: self.get(name) for name in self._data} + + def schema(self, name): + """Full spec dict for a variable (value, default, type, constraints).""" + return dict(self._spec(name)) + + def names(self): + return sorted(self._data.keys()) + + # -- writes ---------------------------------------------------------- + + def set(self, name, value, by=None): + """Validate and persist a new value. Returns the stored value. + + Raises VariableError on unknown name, type mismatch, or constraint + violation. The file is rewritten atomically. + """ + old_val = self.get(name) + spec = self._spec(name) + value = self._coerce(name, value, spec) + self._check_constraints(name, value, spec) + spec["value"] = value + self._persist(by=by) + if old_val != value: + self._record_history("set", name, old_val, value, by=by) + return value + + def reset(self, name, by=None): + """Restore a variable to its default value. Returns the default.""" + old_val = self.get(name) + spec = self._spec(name) + default_val = spec.get("default") + val = self._coerce(name, default_val, spec) + self._check_constraints(name, val, spec) + spec["value"] = val + self._persist(by=by) + if old_val != val: + self._record_history("reset", name, old_val, val, by=by) + return val + + def rollback(self, name, revision=None, by=None): + """Roll back a variable to a previous value from history. + + If revision is None: rolls back to the old_value of the most recent change. + If revision is an integer k >= 1: rolls back k revisions back. + If revision is a timestamp string: rolls back to that revision's old_value. + """ + entries = self.history(name=name, limit=1000) + # self.history returns newest first, so reverse to get chronological + chronological = list(reversed(entries)) + if not chronological: + raise VariableError("no change history found to rollback for variable '%s'" % name) + + target_value = None + if revision is None: + target_value = chronological[-1].get("old_value") + elif str(revision).isdigit() and int(revision) >= 1: + k = int(revision) + if k <= len(chronological): + target_value = chronological[-k].get("old_value") + else: + raise VariableError("revision step %d exceeds history depth (%d) for '%s'" % (k, len(chronological), name)) + else: + rev_str = str(revision).strip() + match = None + for e in reversed(chronological): + if e.get("ts") == rev_str: + match = e + break + if match is not None: + target_value = match.get("old_value") + else: + raise VariableError("revision '%s' not found in history for '%s'" % (rev_str, name)) + + if target_value is None: + raise VariableError("unable to determine rollback target value for '%s'" % name) + + old_val = self.get(name) + spec = self._spec(name) + val = self._coerce(name, target_value, spec) + self._check_constraints(name, val, spec) + spec["value"] = val + caller = by or os.environ.get("BOX_CALLER", "rollback") + self._persist(by=caller) + self._record_history("rollback", name, old_val, val, by=caller) + return val + + def history(self, name=None, limit=20): + """Return history entries (newest first), optionally filtered by name.""" + entries = [] + target = self.history_path or (DEFAULT_HISTORY if os.path.exists(DEFAULT_HISTORY) else FALLBACK_HISTORY) + if os.path.exists(target): + try: + with open(target, "r") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + entries.append(json.loads(line)) + except Exception: + pass + except Exception: + pass + elif not self.history_path and target == DEFAULT_HISTORY and os.path.exists(FALLBACK_HISTORY): + try: + with open(FALLBACK_HISTORY, "r") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + entries.append(json.loads(line)) + except Exception: + pass + except Exception: + pass + + if name: + entries = [e for e in entries if e.get("name") == name] + entries.reverse() + return entries[:limit] + + def _record_history(self, action, name, old_val, new_val, by=None): + entry = { + "ts": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), + "action": action, + "name": name, + "old_value": old_val, + "new_value": new_val, + "by": by or os.environ.get("BOX_CALLER", "super"), + } + line = json.dumps(entry) + "\n" + targets = [self.history_path] if self.history_path else [DEFAULT_HISTORY, FALLBACK_HISTORY] + for hpath in targets: + try: + os.makedirs(os.path.dirname(os.path.abspath(hpath)), exist_ok=True) + with open(hpath, "a") as f: + f.write(line) + except Exception: + pass + + + def _persist(self, by=None): + # Build the full file payload, preserving any existing _meta shape. + payload = {"_meta": {}, "variables": self._data} + try: + with open(self.path) as f: + existing = json.load(f) + if isinstance(existing, dict) and isinstance(existing.get("_meta"), dict): + payload["_meta"] = existing["_meta"] + except (OSError, ValueError): + pass + meta = payload["_meta"] + meta["updated_at"] = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()) + if by: + meta["updated_by"] = by + meta.setdefault("schema_version", 1) + + def _write(dest_path): + try: + os.makedirs(os.path.dirname(os.path.abspath(dest_path)), exist_ok=True) + tmp = dest_path + f".tmp.{os.getpid()}" + with open(tmp, "w") as f: + json.dump(payload, f, indent=2) + f.write("\n") + os.chmod(tmp, 0o644) + os.replace(tmp, dest_path) + except Exception: + pass + + _write(self.path) + # Mirror to fallback or primary if different + if self.path == DEFAULT_REGISTRY: + _write(FALLBACK_REGISTRY) + elif self.path == FALLBACK_REGISTRY: + _write(DEFAULT_REGISTRY) + + +# Detailed schemas for every built-in control variable +_BUILTIN_SPECS = { + "agent_activity_threshold": { + "default": 43200, "type": "int", "min": 60, "max": 604800, "unit": "seconds", + "description": "Seconds of silence before agent is considered inactive (12h default)" + }, + "followup_default_timeout_s": { + "default": 3600, "type": "int", "min": 60, "max": 604800, "unit": "seconds", + "description": "Default timeout before nudging unacknowledged DMs (1h default)" + }, + "keepalive_idle_s": { + "default": 3600, "type": "int", "min": 60, "max": 86400, "unit": "seconds", + "description": "Idle threshold before keepalive pings a thread (1h default)" + }, + "keepalive_interval_s": { + "default": 300, "type": "int", "min": 30, "max": 3600, "unit": "seconds", + "description": "Cadence between keepalive checks (5m default)" + }, + "loop_health_interval_s": { + "default": 900, "type": "int", "min": 60, "max": 86400, "unit": "seconds", + "description": "Sampling interval for agent loop health metrics (15m default)" + }, + "loop_health_threshold": { + "default": 0.5, "type": "float", "min": 0.0, "max": 1.0, "unit": "ratio", + "description": "Minimum healthy ratio of answered+closed loops to landed loops" + }, + "max_nudge_count": { + "default": 2, "type": "int", "min": 0, "max": 10, "unit": "count", + "description": "Maximum automated nudges before escalating to supervisor" + }, + "reminder_min_interval_s": { + "default": 900, "type": "int", "min": 60, "max": 86400, "unit": "seconds", + "description": "Minimum cooldown between reminders (15m default)" + }, + "reminder_quiet_hours_start": { + "default": 22, "type": "int", "min": 0, "max": 23, "unit": "hour", + "description": "Start of quiet hours (UTC, 22 = 10pm)" + }, + "reminder_quiet_hours_end": { + "default": 7, "type": "int", "min": 0, "max": 23, "unit": "hour", + "description": "End of quiet hours (UTC, 7 = 7am)" + }, + "silence_alert_hours": { + "default": 12.0, "type": "float", "min": 1.0, "max": 168.0, "unit": "hours", + "description": "Hours of agent silence before firing operator alert" + }, + "siphon_interval_s": { + "default": 60, "type": "int", "min": 10, "max": 3600, "unit": "seconds", + "description": "Cadence of sidechat work siphon poll loop (60s default)" + }, + "siphon_min_confidence": { + "default": 0.6, "type": "float", "min": 0.0, "max": 1.0, "unit": "confidence", + "description": "Minimum confidence threshold to siphon sidechat item to main" + }, + "siphon_rate_limit_per_hour": { + "default": 5, "type": "int", "min": 1, "max": 60, "unit": "count/hr", + "description": "Max siphoned items per hour to prevent spam" + }, + "sweeper_interval_s": { + "default": 60, "type": "int", "min": 10, "max": 3600, "unit": "seconds", + "description": "Cadence of followup deadline sweeper loop (60s default)" + }, + "thread_max_age_s": { + "default": 604800, "type": "int", "min": 3600, "max": 2592000, "unit": "seconds", + "description": "Maximum sidechat thread age before triggering rotation (7d default)" + }, + "thread_max_messages": { + "default": 200, "type": "int", "min": 10, "max": 10000, "unit": "messages", + "description": "Maximum messages in thread before triggering rotation" + }, + "wake_interval_s": { + "default": 1800, "type": "int", "min": 60, "max": 86400, "unit": "seconds", + "description": "Cadence between periodic wake digests (30m default)" + }, + "wake_reply_nudges": { + "default": 1, "type": "int", "min": 0, "max": 10, "unit": "count", + "description": "Max nudges for unacknowledged wake digests" + }, + "wake_reply_timeout_s": { + "default": 7200, "type": "int", "min": 60, "max": 604800, "unit": "seconds", + "description": "Timeout before nudging unacknowledged wake digests (2h default)" + }, +} + +# Built-in fallbacks +_BUILTIN_DEFAULTS = {k: v["default"] for k, v in _BUILTIN_SPECS.items()} + + +def main(): + import argparse + parser = argparse.ArgumentParser(description="Box variable registry CLI") + sub = parser.add_subparsers(dest="action") + + sub.add_parser("list", help="List all variables") + + p_get = sub.add_parser("get", help="Get a variable") + p_get.add_argument("name", help="Variable name") + + p_set = sub.add_parser("set", help="Set a variable") + p_set.add_argument("name", help="Variable name") + p_set.add_argument("value", help="New value") + + p_reset = sub.add_parser("reset", help="Reset a variable to default") + p_reset.add_argument("name", help="Variable name") + + p_hist = sub.add_parser("history", help="Show variable change history") + p_hist.add_argument("name", nargs="?", default=None, help="Filter by variable name") + p_hist.add_argument("-n", "--limit", type=int, default=20, help="Number of records to show") + + p_rb = sub.add_parser("rollback", help="Roll back a variable to previous value") + p_rb.add_argument("name", help="Variable name") + p_rb.add_argument("--revision", default=None, help="Revision step (int) or timestamp") + + args = parser.parse_args() + v = Variables() + + if not args.action or args.action == "list": + for name in v.names(): + val = v.get(name) + spec = v.schema(name) + print(f"{name:28} = {val!r:10} (default: {spec['default']!r}, unit: {spec.get('unit', '-')})") + elif args.action == "get": + print(json.dumps({"name": args.name, "value": v.get(args.name), "schema": v.schema(args.name)}, indent=2)) + elif args.action == "set": + spec = v.schema(args.name) + vtype = spec.get("type", "str") + if vtype == "int": + val = int(args.value) + elif vtype == "float": + val = float(args.value) + elif vtype == "bool": + val = args.value.lower() in ("true", "1", "yes") + else: + val = args.value + res = v.set(args.name, val, by=os.environ.get("BOX_CALLER", "cli")) + print(f"Updated {args.name} = {res}") + elif args.action == "reset": + res = v.reset(args.name, by=os.environ.get("BOX_CALLER", "cli")) + print(f"Reset {args.name} = {res}") + elif args.action == "history": + entries = v.history(name=args.name, limit=args.limit) + if not entries: + print("No history entries found.") + else: + print(f"\n{'TIMESTAMP':22} {'VARIABLE':28} {'ACTION':10} {'OLD':10} {'NEW':10} {'BY':10}") + print("-" * 92) + for e in entries: + print(f"{e.get('ts',''):22} {e.get('name',''):28} {e.get('action',''):10} {str(e.get('old_value',''))[:10]:10} {str(e.get('new_value',''))[:10]:10} {e.get('by',''):10}") + print() + elif args.action == "rollback": + caller = os.environ.get("BOX_CALLER", "cli") + res = v.rollback(args.name, revision=args.revision, by=caller) + print(f"Rolled back {args.name} = {res}") + + +if __name__ == "__main__": + main() diff --git a/tests/test_loop_health_remediation.py b/tests/test_loop_health_remediation.py new file mode 100644 index 0000000..96e8515 --- /dev/null +++ b/tests/test_loop_health_remediation.py @@ -0,0 +1,96 @@ +#!/usr/bin/env python3 +""" +test_loop_health_remediation.py — Unit and integration tests for: +1. Reconstructing intrinsic loops from followups.json and dm-log.jsonl +2. Fleet loop health metrics computation and threshold evaluation +3. Progressive remediation: + - Soft break healing (auto-resolving answered follow-ups) + - Deadline re-arming for pending expired nudges +4. Hard break diagnostic escalations (silent agent, auth rot, systemd failure) +5. Audit event logging in job-log.jsonl +""" +import unittest +import subprocess +import json +import os +import sys +import tempfile +import shutil +from pathlib import Path +from datetime import datetime, timezone + +REPO_ROOT = Path("/home/super/Projects/NetVM") +BIN_DIR = REPO_ROOT / "bin" +BOX_CTL = BIN_DIR / "box-ctl.py" + +sys.path.insert(0, str(BIN_DIR)) +import gravity + + +class TestLoopDiagnosticsAndRemediation(unittest.TestCase): + """Test gravity.py loop health and progressive remediation.""" + + def test_fleet_loop_health_calculation(self): + health = gravity.get_fleet_loop_health(threshold=0.5) + self.assertIn("agents", health) + self.assertIn("summary", health) + self.assertIn("healthy", health["summary"]) + self.assertIn("overall_health", health["summary"]) + + def test_reconstruct_loops(self): + loops = gravity.reconstruct_loops(limit=10) + self.assertIsInstance(loops, list) + if loops: + l = loops[0] + self.assertIn("loop_id", l) + self.assertIn("agent", l) + self.assertIn("state", l) + + def test_diagnose_breaks_structure(self): + breaks = gravity.diagnose_breaks() + self.assertIsInstance(breaks, list) + for b in breaks: + self.assertIn("type", b) + self.assertIn("severity", b) + self.assertIn("detail", b) + + def test_remediate_breaks_dry_run(self): + res = gravity.remediate_breaks(dry_run=True) + self.assertTrue(res.get("ok")) + self.assertTrue(res.get("dry_run")) + self.assertIsInstance(res.get("remediated"), list) + self.assertIsInstance(res.get("escalated"), list) + + +class TestLoopRpc(unittest.TestCase): + """Test box-ctl.py allowlisted RPC actions for loops.""" + + def test_box_ctl_loop_health(self): + cmd = [sys.executable, str(BOX_CTL), "loop-health"] + res = subprocess.run(cmd, capture_output=True, text=True) + self.assertEqual(res.returncode, 0) + data = json.loads(res.stdout) + self.assertTrue(data.get("ok")) + self.assertIn("agents", data) + self.assertIn("summary", data) + + + def test_box_ctl_loop_status(self): + cmd = [sys.executable, str(BOX_CTL), "loop-status", "--limit", "5"] + res = subprocess.run(cmd, capture_output=True, text=True) + self.assertEqual(res.returncode, 0) + data = json.loads(res.stdout) + self.assertTrue(data.get("ok")) + self.assertIn("loops", data) + + def test_box_ctl_loop_remediate_dry_run(self): + cmd = [sys.executable, str(BOX_CTL), "loop-remediate", "--dry-run"] + res = subprocess.run(cmd, capture_output=True, text=True) + self.assertEqual(res.returncode, 0) + data = json.loads(res.stdout) + self.assertTrue(data.get("ok")) + self.assertTrue(data.get("dry_run")) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_modulate_strategy.py b/tests/test_modulate_strategy.py new file mode 100644 index 0000000..2645617 --- /dev/null +++ b/tests/test_modulate_strategy.py @@ -0,0 +1,112 @@ +#!/usr/bin/env python3 +""" +test_modulate_strategy.py — Unit and integration tests for: +1. Hierarchical override scoping: + (input_type, subtype, agent) -> (input_type, None, agent) -> + (input_type, subtype, None) -> (input_type, None, None) -> Builtin table +2. Dynamic strategy evaluation (`modulate.py eval` / `box-ctl.py strat-get`) +3. Atomic persistence to strategy.json +4. Override resets and isolation across input types +""" +import unittest +import subprocess +import json +import os +import sys +import tempfile +import shutil +from pathlib import Path + +REPO_ROOT = Path("/home/super/Projects/NetVM") +BIN_DIR = REPO_ROOT / "bin" +BOX_CTL = BIN_DIR / "box-ctl.py" + +sys.path.insert(0, str(BIN_DIR)) +import modulate + + +class TestModulationHierarchy(unittest.TestCase): + """Test strategy lookup cascade and overrides.""" + + def setUp(self): + self.tmp_dir = tempfile.mkdtemp() + self.strat_file = str(Path(self.tmp_dir) / "strategy.json") + self.orig_default = modulate.DEFAULT_STRATEGY_FILE + self.orig_fallback = modulate.FALLBACK_STRATEGY_FILE + modulate.DEFAULT_STRATEGY_FILE = self.strat_file + modulate.FALLBACK_STRATEGY_FILE = self.strat_file + + def tearDown(self): + modulate.DEFAULT_STRATEGY_FILE = self.orig_default + modulate.FALLBACK_STRATEGY_FILE = self.orig_fallback + shutil.rmtree(self.tmp_dir, ignore_errors=True) + + + def test_builtin_fallback(self): + # Default builtin policy for manual DMs + pol = modulate.modulate(modulate.InputType.MANUAL, agent="646") + self.assertIsNotNone(pol) + self.assertTrue(pol.track) + self.assertEqual(pol.timeout_s, 3600) + + def test_agent_specific_override(self): + # Set override for agent 'pip' on 'manual' + modulate.set_strategy_override( + "manual", + None, + "pip", + track=False, + timeout_s=1800, + nudges=0, + escalate="none", + ) + + # Pip gets the custom override (track=False -> modulate returns None) + pip_pol = modulate.modulate(modulate.InputType.MANUAL, agent="pip") + self.assertIsNone(pip_pol) + + # 646 gets the default builtin policy + other_pol = modulate.modulate(modulate.InputType.MANUAL, agent="646") + self.assertTrue(other_pol.track) + self.assertEqual(other_pol.timeout_s, 3600) + + + def test_subtype_hierarchical_resolution(self): + # Set type-level default + modulate.set_strategy_override("health", None, None, timeout_s=900) + # Set subtype override + modulate.set_strategy_override("health", "DEGRADED", None, timeout_s=3600) + + # Health without subtype gets type-level + pol1 = modulate.modulate(modulate.InputType.HEALTH, subtype="OK") + self.assertEqual(pol1.timeout_s, 900) + + # Health with 'DEGRADED' subtype gets subtype-specific + pol2 = modulate.modulate(modulate.InputType.HEALTH, subtype="DEGRADED") + self.assertEqual(pol2.timeout_s, 3600) + + def test_override_reset(self): + modulate.set_strategy_override("wake", None, "opm", timeout_s=600) + self.assertEqual(modulate.modulate(modulate.InputType.WAKE, agent="opm", actionable=True).timeout_s, 600) + + modulate.reset_strategy_override("wake", None, "opm") + # Reverts to builtin default (7200s for wake) + self.assertEqual(modulate.modulate(modulate.InputType.WAKE, agent="opm", actionable=True).timeout_s, 7200) + + + +class TestStrategyRpc(unittest.TestCase): + """Test box-ctl.py allowlisted RPC actions for strategy.""" + + def test_box_ctl_strat_list(self): + cmd = [sys.executable, str(BOX_CTL), "strat-list"] + res = subprocess.run(cmd, capture_output=True, text=True) + self.assertEqual(res.returncode, 0) + data = json.loads(res.stdout) + self.assertTrue(data.get("ok")) + self.assertIn("strategies", data) + + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_variables_engine.py b/tests/test_variables_engine.py new file mode 100644 index 0000000..00fbad3 --- /dev/null +++ b/tests/test_variables_engine.py @@ -0,0 +1,104 @@ +#!/usr/bin/env python3 +""" +test_variables_engine.py — Unit and integration tests for: +1. Schema validation (int, float, bool, string, ranges, choices) +2. Dual-synced atomic persistence between local and /srv/box/ +3. Append-only audit history logging to variables-history.jsonl +4. Atomic rollback to previous revisions +5. box-ctl.py allowlisted RPC actions for variables +""" +import unittest +import subprocess +import json +import os +import sys +import tempfile +import shutil +from pathlib import Path + +REPO_ROOT = Path("/home/super/Projects/NetVM") +BIN_DIR = REPO_ROOT / "bin" +VARS_PY = BIN_DIR / "variables.py" +BOX_CTL = BIN_DIR / "box-ctl.py" + +sys.path.insert(0, str(BIN_DIR)) +import variables + + +class TestVariablesEngine(unittest.TestCase): + """Test schema enforcement and in-memory variable operations.""" + + def setUp(self): + self.tmp_dir = tempfile.mkdtemp() + self.var_file = Path(self.tmp_dir) / "variables.json" + self.hist_file = Path(self.tmp_dir) / "variables-history.jsonl" + self.engine = variables.Variables(path=str(self.var_file), history_path=str(self.hist_file)) + + def tearDown(self): + shutil.rmtree(self.tmp_dir, ignore_errors=True) + + def test_default_values_load(self): + val = self.engine.get("loop_health_threshold") + self.assertEqual(val, 0.5) + val_sec = self.engine.get("loop_health_interval_s") + self.assertEqual(val_sec, 900) + + def test_schema_validation_out_of_range(self): + # loop_health_threshold range: 0.0 .. 1.0 + with self.assertRaises(ValueError): + self.engine.set("loop_health_threshold", 1.5) + with self.assertRaises(ValueError): + self.engine.set("loop_health_threshold", -0.1) + + def test_schema_validation_type_coercion(self): + self.engine.set("loop_health_interval_s", "1200") + self.assertEqual(self.engine.get("loop_health_interval_s"), 1200) + + def test_history_logging_and_rollback(self): + initial = self.engine.get("max_nudge_count") + self.engine.set("max_nudge_count", 5, by="operator_test") + self.assertEqual(self.engine.get("max_nudge_count"), 5) + + # Check history recorded + hist = self.engine.history(name="max_nudge_count", limit=5) + self.assertGreaterEqual(len(hist), 1) + self.assertEqual(hist[0]["action"], "set") + self.assertEqual(hist[0]["new_value"], 5) + + # Rollback to initial + rb_val = self.engine.rollback("max_nudge_count", by="operator_test") + self.assertEqual(rb_val, initial) + self.assertEqual(self.engine.get("max_nudge_count"), initial) + + def test_reset_variable(self): + self.engine.set("sweeper_interval_s", 120) + self.assertEqual(self.engine.get("sweeper_interval_s"), 120) + res_val = self.engine.reset("sweeper_interval_s") + self.assertEqual(res_val, 60) + self.assertEqual(self.engine.get("sweeper_interval_s"), 60) + + + +class TestVariablesCliAndRpc(unittest.TestCase): + """Test CLI commands and box-ctl RPC actions for variables.""" + + def test_box_ctl_vars_list(self): + cmd = [sys.executable, str(BOX_CTL), "vars-list"] + res = subprocess.run(cmd, capture_output=True, text=True) + self.assertEqual(res.returncode, 0) + data = json.loads(res.stdout) + self.assertTrue(data.get("ok")) + self.assertIn("variables", data) + self.assertIn("loop_health_threshold", data["variables"]) + + def test_box_ctl_vars_history(self): + cmd = [sys.executable, str(BOX_CTL), "vars-history", "loop_health_threshold", "5"] + res = subprocess.run(cmd, capture_output=True, text=True) + self.assertEqual(res.returncode, 0) + data = json.loads(res.stdout) + self.assertTrue(data.get("ok")) + self.assertIsInstance(data.get("history"), list) + + +if __name__ == "__main__": + unittest.main()