"""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 SALVAGE = "salvage" # token exhaustion & onboarding work orders 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"), # Salvage: token depletion & onboarding rescue. (InputType.SALVAGE, "BLOCKED"): (True, Priority.CRITICAL, 600, 3, "opm"), (InputType.SALVAGE, "LOW"): (True, Priority.IMPORTANT, 900, 2, "opm"), (InputType.SALVAGE, None): (True, Priority.CRITICAL, 900, 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()