304 lines
9.4 KiB
Python
304 lines
9.4 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Thread lifecycle manager for side chat rotation.
|
|
|
|
Evaluates rotation triggers (age, message count, explicit flag) for a reuse_key
|
|
and performs rotation when triggered. Integrates with:
|
|
- job-sidechats.json (state + lifecycle policy)
|
|
- sidechat_manager.py (thread operations)
|
|
- box follow-up system (thread_id migration)
|
|
|
|
Usage:
|
|
python3 thread_lifecycle.py check --reuse-key <key> [--state PATH]
|
|
python3 thread_lifecycle.py rotate --reuse-key <key> --reason <reason> [--state PATH]
|
|
|
|
The `check` command is designed to run at job dispatch time, before the reuse
|
|
check in job-dispatch.py. It returns:
|
|
exit 0: no rotation needed, reuse existing thread
|
|
exit 2: rotation triggered, new thread needed (caller should create fresh)
|
|
exit 1: error
|
|
|
|
This is a prototype. Not deployed.
|
|
"""
|
|
|
|
import argparse
|
|
import json
|
|
import os
|
|
import sys
|
|
import time
|
|
from datetime import datetime, timezone, timedelta
|
|
|
|
# Default lifecycle policy (applied when no per-key policy exists)
|
|
DEFAULT_POLICY = {
|
|
"max_age_hours": 168, # 7 days
|
|
"max_messages": 200,
|
|
"channel_type": "job",
|
|
}
|
|
|
|
# Channel-type defaults (override DEFAULT_POLICY)
|
|
CHANNEL_DEFAULTS = {
|
|
"heartbeat": {"max_age_hours": 24, "max_messages": 500, "channel_type": "heartbeat"},
|
|
"job": {"max_age_hours": 168, "max_messages": 200, "channel_type": "job"},
|
|
"coordination": {"max_age_hours": 720, "max_messages": 200, "channel_type": "coordination"},
|
|
"dm": {"max_age_hours": None, "max_messages": None, "channel_type": "dm"}, # never rotate
|
|
}
|
|
|
|
|
|
def _utcnow():
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
def _parse_ts(ts_str):
|
|
"""Parse ISO-8601 UTC timestamp. Returns datetime or None."""
|
|
if not ts_str:
|
|
return None
|
|
try:
|
|
# Handle both "2026-10-04T12:00:00Z" and "... +00:00" forms
|
|
ts_str = ts_str.replace("Z", "+00:00")
|
|
dt = datetime.fromisoformat(ts_str)
|
|
if dt.tzinfo is None:
|
|
dt = dt.replace(tzinfo=timezone.utc)
|
|
return dt
|
|
except (ValueError, TypeError):
|
|
return None
|
|
|
|
|
|
def load_state(state_path):
|
|
"""Load job-sidechats.json. Returns dict."""
|
|
if not os.path.exists(state_path):
|
|
return {}
|
|
with open(state_path) as f:
|
|
return json.load(f)
|
|
|
|
|
|
def save_state(state_path, state):
|
|
"""Save state atomically (write temp + rename)."""
|
|
tmp = state_path + ".tmp"
|
|
with open(tmp, "w") as f:
|
|
json.dump(state, f, indent=2)
|
|
os.rename(tmp, state_path)
|
|
|
|
|
|
def get_policy(state, reuse_key):
|
|
"""
|
|
Get lifecycle policy for a reuse_key.
|
|
Priority: per-key policy > channel-type default > global default.
|
|
Returns dict with max_age_hours, max_messages, channel_type.
|
|
"""
|
|
lifecycle = state.get("_lifecycle", {})
|
|
defaults = state.get("_lifecycle_defaults", {})
|
|
|
|
# Start with global defaults
|
|
policy = dict(DEFAULT_POLICY)
|
|
policy.update(defaults)
|
|
|
|
# Apply channel-type defaults if specified
|
|
key_policy = lifecycle.get(reuse_key, {})
|
|
channel_type = key_policy.get("channel_type")
|
|
if channel_type and channel_type in CHANNEL_DEFAULTS:
|
|
policy.update(CHANNEL_DEFAULTS[channel_type])
|
|
|
|
# Per-key overrides win
|
|
policy.update(key_policy)
|
|
return policy
|
|
|
|
|
|
def get_thread_age_hours(state, reuse_key):
|
|
"""
|
|
Get thread age in hours from creation.
|
|
Returns None if unknown (fail-open: age trigger skipped).
|
|
"""
|
|
# Prefer explicit creation timestamp
|
|
created_key = f"{reuse_key}:created_at"
|
|
created = _parse_ts(state.get(created_key))
|
|
if created:
|
|
delta = _utcnow() - created
|
|
return delta.total_seconds() / 3600
|
|
|
|
# Fall back to last rotation timestamp
|
|
rotated = _parse_ts(state.get(f"{reuse_key}:rotated_at"))
|
|
if rotated:
|
|
delta = _utcnow() - rotated
|
|
return delta.total_seconds() / 3600
|
|
|
|
return None
|
|
|
|
|
|
def check_rotation_needed(state, reuse_key, message_count=None, force=False):
|
|
"""
|
|
Evaluate rotation triggers for a reuse_key.
|
|
|
|
Args:
|
|
state: loaded state dict
|
|
reuse_key: the reuse key to check
|
|
message_count: current message count (None = unknown, skip count trigger)
|
|
force: explicit rotation request (topic change / manual)
|
|
|
|
Returns:
|
|
(needed: bool, reason: str, details: dict)
|
|
"""
|
|
policy = get_policy(state, reuse_key)
|
|
|
|
# DM channels never rotate
|
|
if policy.get("channel_type") == "dm":
|
|
return False, "dm_no_rotate", {}
|
|
|
|
# Explicit force always wins
|
|
if force:
|
|
return True, "explicit", {"requested": True}
|
|
|
|
# Age trigger
|
|
max_age = policy.get("max_age_hours")
|
|
if max_age is not None:
|
|
age_hours = get_thread_age_hours(state, reuse_key)
|
|
if age_hours is not None and age_hours >= max_age:
|
|
return True, "age", {
|
|
"age_hours": round(age_hours, 1),
|
|
"max_age_hours": max_age,
|
|
}
|
|
|
|
# Message count trigger
|
|
max_msg = policy.get("max_messages")
|
|
if max_msg is not None and message_count is not None:
|
|
if message_count >= max_msg:
|
|
return True, "message_count", {
|
|
"message_count": message_count,
|
|
"max_messages": max_msg,
|
|
}
|
|
|
|
return False, "none", {}
|
|
|
|
|
|
def resolve_thread_for_nudge(state, thread_id):
|
|
"""
|
|
Resolve a thread_id for nudge delivery, handling rotation.
|
|
|
|
If thread_id points to an archived/rotated thread, follow the
|
|
:previous chain to find the current UUID.
|
|
|
|
Returns the UUID to use (may be the input if no rotation found).
|
|
"""
|
|
# Build reverse map: old_uuid -> reuse_key
|
|
for key, value in state.items():
|
|
if key.startswith("_") or ":" in key:
|
|
continue
|
|
prev_key = f"{key}:previous"
|
|
if state.get(prev_key) == thread_id:
|
|
# This thread was rotated; use the current UUID
|
|
current = state.get(key)
|
|
if current and current != thread_id:
|
|
return current
|
|
return thread_id
|
|
|
|
|
|
def record_rotation(state, reuse_key, old_uuid, new_uuid, reason):
|
|
"""
|
|
Update state after rotation. Returns updated state dict.
|
|
Does not save — caller saves.
|
|
"""
|
|
now = _utcnow().strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
prev_count = state.get(f"{reuse_key}:rotation_count", 0)
|
|
|
|
state[f"{reuse_key}:previous"] = old_uuid
|
|
state[reuse_key] = new_uuid
|
|
state[f"{reuse_key}:rotated_at"] = now
|
|
state[f"{reuse_key}:created_at"] = now # new thread creation baseline
|
|
state[f"{reuse_key}:rotation_reason"] = reason
|
|
state[f"{reuse_key}:rotation_count"] = prev_count + 1
|
|
return state
|
|
|
|
|
|
def cmd_check(args):
|
|
"""Check if rotation is needed. Exit 0=no, 2=yes, 1=error."""
|
|
state = load_state(args.state)
|
|
reuse_key = args.reuse_key
|
|
|
|
if reuse_key not in state:
|
|
print(f"No existing thread for {reuse_key}, no rotation check needed")
|
|
return 0
|
|
|
|
needed, reason, details = check_rotation_needed(
|
|
state, reuse_key,
|
|
message_count=args.messages,
|
|
force=args.force,
|
|
)
|
|
if needed:
|
|
print(json.dumps({"rotate": True, "reason": reason, "details": details}))
|
|
return 2
|
|
else:
|
|
print(json.dumps({"rotate": False, "reason": reason}))
|
|
return 0
|
|
|
|
|
|
def cmd_rotate(args):
|
|
"""Record a rotation (called after new thread is created)."""
|
|
state = load_state(args.state)
|
|
reuse_key = args.reuse_key
|
|
|
|
old_uuid = state.get(reuse_key)
|
|
if not old_uuid:
|
|
print(f"No existing thread for {reuse_key}", file=sys.stderr)
|
|
return 1
|
|
|
|
if not args.new_uuid:
|
|
print("--new-uuid required", file=sys.stderr)
|
|
return 1
|
|
|
|
state = record_rotation(state, reuse_key, old_uuid, args.new_uuid, args.reason)
|
|
save_state(args.state, state)
|
|
print(json.dumps({
|
|
"rotated": True,
|
|
"reuse_key": reuse_key,
|
|
"old_uuid": old_uuid,
|
|
"new_uuid": args.new_uuid,
|
|
"reason": args.reason,
|
|
}))
|
|
return 0
|
|
|
|
|
|
def cmd_resolve(args):
|
|
"""Resolve a thread_id through rotation chain (for nudge router)."""
|
|
state = load_state(args.state)
|
|
resolved = resolve_thread_for_nudge(state, args.thread_id)
|
|
print(json.dumps({
|
|
"input": args.thread_id,
|
|
"resolved": resolved,
|
|
"rotated": resolved != args.thread_id,
|
|
}))
|
|
return 0
|
|
|
|
|
|
def main():
|
|
p = argparse.ArgumentParser(description="Side chat thread lifecycle manager")
|
|
p.add_argument("--state", default=os.path.expanduser(
|
|
"~/Projects/NetVM/job-sidechats.json"),
|
|
help="Path to job-sidechats.json")
|
|
sub = p.add_subparsers(dest="cmd", required=True)
|
|
|
|
c = sub.add_parser("check", help="Check if rotation is needed")
|
|
c.add_argument("--reuse-key", required=True)
|
|
c.add_argument("--messages", type=int, default=None,
|
|
help="Current message count (None = skip count trigger)")
|
|
c.add_argument("--force", action="store_true",
|
|
help="Force rotation (explicit topic change)")
|
|
|
|
r = sub.add_parser("rotate", help="Record a completed rotation")
|
|
r.add_argument("--reuse-key", required=True)
|
|
r.add_argument("--new-uuid", required=True)
|
|
r.add_argument("--reason", default="manual")
|
|
|
|
v = sub.add_parser("resolve", help="Resolve thread_id through rotation")
|
|
v.add_argument("--thread-id", required=True)
|
|
|
|
args = p.parse_args()
|
|
if args.cmd == "check":
|
|
sys.exit(cmd_check(args))
|
|
elif args.cmd == "rotate":
|
|
sys.exit(cmd_rotate(args))
|
|
elif args.cmd == "resolve":
|
|
sys.exit(cmd_resolve(args))
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|