feat(tests): add comprehensive unit test suite for variables, strategy modulation, and loop health
This commit is contained in:
+622
@@ -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:<type>]). 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:<type>]. 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()
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user