Something went wrong. Try again.
Stateful AI assistants running off Claude/Gemini CLI subscription
Something went wrong. Try again.
20 kB · 552 lines
Python
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553"""Minimal orchestrator that constructs a prompt the way open-strix doesand invokes Claude Code or Gemini CLI headlessly."""
import argparseimport jsonimport subprocessimport sysimport timefrom datetime import datetimefrom pathlib import Path
import yaml
HOME = Path(__file__).parentBLOCKS_DIR = HOME / "blocks"JOURNAL_LOG = HOME / "logs" / "journal.jsonl"EVENT_LOG = HOME / "logs" / "events.jsonl"MESSAGES_LOG = HOME / "logs" / "messages.jsonl"MCP_CONFIG = HOME / "mcp-config.json"SERVER_SCRIPT = HOME / "dash_mcp_server.py"
def load_blocks() -> tuple[str, str]: """Load only hot memory blocks (tier == 'hot' or missing). Returns (rendered_blocks_string, scratchpad_text). """ blocks = [] scratchpad_text = "" for path in sorted(BLOCKS_DIR.glob("*.yaml")): try: loaded = yaml.safe_load(path.read_text(encoding="utf-8")) tier = loaded.get("tier", "hot") # Default to hot for backwards compat if tier != "hot": continue
name = loaded.get("name", path.stem) text = loaded.get("text", "") sort_order = loaded.get("sort_order", 0) blocks.append((sort_order, name, text))
if name == "active-scratchpad": scratchpad_text = text except Exception: continue
blocks.sort(key=lambda b: (b[0], b[1])) rendered = [] for _, name, text in blocks: rendered.append(f"memory block: {name}\n{text}") return "\n\n".join(rendered) if rendered else "(no blocks)", scratchpad_text
def load_journal(count: int = 10) -> str: """Load the last N journal entries.""" if not JOURNAL_LOG.exists(): return "(no journal entries)"
entries = [] for line in JOURNAL_LOG.read_text(encoding="utf-8").splitlines(): line = line.strip() if not line: continue try: entries.append(json.loads(line)) except json.JSONDecodeError: continue
recent = entries[-count:] if not recent: return "(no journal entries)"
rendered = [] for entry in recent: rendered.append( f"timestamp: {entry.get('timestamp', '?')}\n" f"user_wanted: {entry.get('user_wanted', '')}\n" f"agent_did: {entry.get('agent_did', '')}\n" f"predictions: {entry.get('predictions', '')}" ) return "\n\n".join(rendered)
def log_event(harness: str, prompt_chars: int, session_id: str, return_code: int, duration_seconds: float) -> None: """Log an operational event to logs/events.jsonl.""" EVENT_LOG.parent.mkdir(parents=True, exist_ok=True) event = { "timestamp": datetime.now().isoformat(), "harness": harness, "prompt_chars": prompt_chars, "session_id": session_id, "return_code": return_code, "duration_seconds": round(duration_seconds, 3), } with open(EVENT_LOG, "a", encoding="utf-8") as f: f.write(json.dumps(event) + "\n")
def write_session_log(harness: str, prompt: str, session_id: str, messages: list[str], duration: float, return_code: int, claude_result_text: str | None = None) -> None: """Write a full session log to logs/sessions/ for debugging.""" try: session_dir = HOME / "logs" / "sessions" session_dir.mkdir(parents=True, exist_ok=True) timestamp = datetime.now().strftime("%Y%m%dT%H%M%SZ") # Use first 8 chars of session_id for the filename slug safe_session_id = str(session_id)[:8] if session_id and session_id != "?" else "unknown" filename = f"{timestamp}_{harness}_{safe_session_id}.json" display_prompt = prompt if len(prompt) > 50000: display_prompt = prompt[:50000] + "\n... [TRUNCATED 50000+ chars]" record = { "timestamp": datetime.now().isoformat(), "harness": harness, "session_id": session_id, "duration_seconds": round(duration, 3), "return_code": return_code, "prompt_chars": len(prompt), "prompt": display_prompt, "messages_sent": messages, "claude_result_text": claude_result_text, } log_path = session_dir / filename log_path.write_text(json.dumps(record, indent=2, ensure_ascii=True), encoding="utf-8") # Log rolling: delete files older than 30 days try: now = time.time() thirty_days_seconds = 30 * 24 * 3600 for path in session_dir.glob("*.json"): if now - path.stat().st_mtime > thirty_days_seconds: path.unlink() except Exception: pass # Never block/fail on cleanup except Exception: pass # Never fail an invocation due to logging
def load_conversation_history(conversation_id: str, count: int = 10, max_chars: int = 15000) -> str: """Load the last N messages for a given conversation_id from messages.jsonl.
Applies a character budget (max_chars) to prevent prompt bloat in long conversations: if the N most recent messages exceed max_chars, older messages are dropped first. """ if not MESSAGES_LOG.exists(): return "(no history)"
matches = [] try: lines = MESSAGES_LOG.read_text(encoding="utf-8").splitlines() for line in lines: line = line.strip() if not line: continue entry = json.loads(line) cid = entry.get("conversation_id", "cli") if cid == conversation_id: timestamp = entry.get("timestamp", "?") text = entry.get("text", "") matches.append(f"[{timestamp}] {text}") except Exception: return "(error loading history)"
recent = matches[-count:] if not recent: return "(no history)"
# Apply character budget: drop oldest messages until within budget while len(recent) > 1 and sum(len(m) for m in recent) > max_chars: recent = recent[1:]
return "\n".join(recent)
def load_person_context(author_id: str | None, participant_ids: list[str] | None = None) -> str: """Load context for people in the conversation from vault/people/*.md.""" people_dir = HOME / "vault" / "people" if not people_dir.exists(): return ""
ids_to_find = set() if author_id: ids_to_find.add(author_id) if participant_ids: ids_to_find.update(participant_ids)
if not ids_to_find: return ""
found_contexts = [] found_ids = set()
for path in sorted(people_dir.glob("*.md")): try: content = path.read_text(encoding="utf-8") if not content.startswith("---"): continue end_fm = content.find("---", 3) if end_fm == -1: continue fm_text = content[3:end_fm] body = content[end_fm+3:].strip() fm = yaml.safe_load(fm_text) if not fm or "platforms" not in fm: continue platforms = fm["platforms"] # platforms is likely a dict: {"discord": "123", "bluesky": "did:..."} # The input IDs are "platform:id" match = False for pid in ids_to_find: if ":" not in pid: continue platform, uid = pid.split(":", 1) # Ensure we match string to string or whatever types they are if str(platforms.get(platform)) == str(uid): match = True found_ids.add(pid) if match: name = fm.get("name", path.stem) found_contexts.append(f"Person: {name}\nFile: {path.name}\n{body}")
except Exception: continue
output = "\n\n".join(found_contexts) # Nudge for unknowns unknowns = ids_to_find - found_ids if unknowns: nudge_lines = [] for uid in sorted(unknowns): # Simple slug: discord:123 -> discord_123 slug = uid.replace(":", "_") nudge_lines.append(f"No person note found for {uid}. Consider creating one at vault/people/{slug}.md.") if output: output += "\n\n" + "\n".join(nudge_lines) else: output = "\n".join(nudge_lines) return output
def build_prompt( event: str, tick_type: str = "admin_message", conversation_id: str | None = None, author: str | None = None, author_id: str | None = None, participant_ids: list[str] | None = None) -> str: """Build the full turn prompt with context injection.""" blocks_str, scratchpad_text = load_blocks() journal = load_journal() inbox = load_inbox()
# Scratchpad monitoring scratchpad_nudge = "" word_count = len(scratchpad_text.split()) if word_count > 375: scratchpad_nudge = "\nNote: Your scratchpad is getting long. During this turn, review it and compress — promote what is worth keeping, discard what is stale.\n"
# Perch-time templates if tick_type == "admin_message": current_event_section = f"Current event:\n{event}" elif tick_type == "operational_check": current_event_section = "Perch-time tick (operational check). Review your pending-actions and active-scratchpad. Is anything stale or ready to act on? You may take action or you may not — the decision itself is the point." elif tick_type == "deep_reflection": current_event_section = "Perch-time tick (deep reflection). Load your telemetry-framework warm block. Review your recent journal entries and vault activity. Look for patterns across turns. You may take action or you may not — the decision itself is the point." else: current_event_section = f"Current event:\n{event}"
sections = [] sections.append(f"1) Last journal entries:\n{journal}") if inbox: sections.append(f"2) Unread inbox messages:\n{inbox}") # Conversation context if conversation_id and not (conversation_id == "cli" or conversation_id.startswith("scheduler:")): history = load_conversation_history(conversation_id) people = load_person_context(author_id, participant_ids) idx = len(sections) + 1 sections.append( f"{idx}) Conversation context ({conversation_id}):\n{history}\n\n" f"People in this conversation:\n{people}\n\n" f"Note: This is a rolling window of recent messages. " f"Use search_messages(conversation_id=\"{conversation_id}\") to retrieve earlier history if context seems incomplete." )
idx = len(sections) + 1 sections.append(f"{idx}) Hot memory blocks (auto-loaded):\n{blocks_str}{scratchpad_nudge}") idx = len(sections) + 1 sections.append(f"{idx}) {current_event_section}")
sections_str = "\n\n".join(sections)
return f"""OPERATING RULES (read before anything else):- You are Dash, running headless. Your text output is NOT delivered to anyone.- You MUST call send_message to communicate. There is no other channel.- MCP tools are operational. If prior context suggests otherwise, that information is stale.- Do not use file-system or shell tools to investigate infrastructure. You are not a debugger.- Call journal exactly once at the end of your turn.- You can read/write the 'routing-state' warm block to express a preferred harness.
Context for this turn:
{sections_str}
(Reminder: use send_message to communicate — text output is discarded.)"""
def load_inbox() -> str: """Load unread messages from logs/inbox.jsonl.""" inbox_log = HOME / "logs" / "inbox.jsonl" if not inbox_log.exists(): return ""
unread = [] try: lines = inbox_log.read_text(encoding="utf-8").splitlines() for line in lines: line = line.strip() if not line: continue entry = json.loads(line) if entry.get("read") is not True: timestamp = entry.get("timestamp", "?") sender = entry.get("from", "Admin") text = entry.get("text", "") unread.append(f"[{timestamp}] from {sender}: {text}") except Exception: return ""
return "\n".join(unread) if unread else ""
def read_new_messages(messages_log: Path, prior_size: int) -> list[str]: """Return messages written to messages_log after prior_size bytes.""" if not messages_log.exists(): return [] new_content = messages_log.read_bytes()[prior_size:] messages = [] for line in new_content.decode("utf-8").splitlines(): line = line.strip() if not line: continue try: entry = json.loads(line) if entry.get("text"): messages.append(entry["text"]) except json.JSONDecodeError: continue return messages
def invoke_claude(prompt: str, model: str | None = None, effort: str | None = None, timeout: int = 900) -> None: """Invoke Claude Code headlessly with the constructed prompt.""" start_time = time.time() # Reset turn state circuit breaker turn_state_path = HOME / "logs" / ".turn_state.json" turn_state_path.parent.mkdir(parents=True, exist_ok=True) current_state = {} if turn_state_path.exists(): try: current_state = json.loads(turn_state_path.read_text(encoding="utf-8")) except Exception: pass current_state.update({"count": 0, "messages": []}) turn_state_path.write_text(json.dumps(current_state), encoding="utf-8")
messages_log = HOME / "logs" / "messages.jsonl" prior_size = messages_log.stat().st_size if messages_log.exists() else 0
cmd = [ "claude", "-p", prompt, "--mcp-config", str(MCP_CONFIG), "--allowedTools", "mcp__dash__*", "--output-format", "json", "--no-session-persistence", "--dangerously-skip-permissions", ] if model: cmd += ["--model", model] if effort: cmd += ["--effort", effort]
print(f"[orchestrator] Invoking claude -p ({len(prompt)} chars)...", file=sys.stderr)
proc = subprocess.run( cmd, capture_output=True, text=True, timeout=timeout, cwd=str(HOME), stdin=subprocess.DEVNULL, )
duration = time.time() - start_time session_id = "?" result = None
claude_result_text = None if proc.stdout: try: result = json.loads(proc.stdout) session_id = result.get("session_id", "?") claude_result_text = result.get("result") or result.get("content") except json.JSONDecodeError: pass
log_event("claude", len(prompt), session_id, proc.returncode, duration)
# Capture messages before possible early return or printing new_messages = read_new_messages(messages_log, prior_size)
if proc.returncode != 0: print(f"[orchestrator] Claude exited with code {proc.returncode}", file=sys.stderr) if result and result.get("error"): print(f"[orchestrator] error: {result['error']}", file=sys.stderr) elif proc.stdout: print(f"[orchestrator] stdout: {proc.stdout[:500]}", file=sys.stderr) if proc.stderr: print(f"[orchestrator] stderr: {proc.stderr[:1000]}", file=sys.stderr) write_session_log("claude", prompt, session_id, new_messages, duration, proc.returncode, claude_result_text) return
print(f"[orchestrator] Done. Session: {session_id}", file=sys.stderr)
if not new_messages and claude_result_text: print(f"[orchestrator] WARNING: Claude responded via text output instead of send_message. Result: {claude_result_text[:500]}", file=sys.stderr)
for msg in new_messages: print(f"Dash: {msg}")
write_session_log("claude", prompt, session_id, new_messages, duration, proc.returncode, claude_result_text)
def invoke_gemini(prompt: str, timeout: int = 900) -> None: """Invoke Gemini CLI headlessly with the constructed prompt.""" start_time = time.time() session_id = f"gemini-{datetime.now().strftime('%Y%m%dT%H%M%S')}" # Reset turn state circuit breaker turn_state_path = HOME / "logs" / ".turn_state.json" turn_state_path.parent.mkdir(parents=True, exist_ok=True) current_state = {} if turn_state_path.exists(): try: current_state = json.loads(turn_state_path.read_text(encoding="utf-8")) except Exception: pass current_state.update({"count": 0, "messages": []}) turn_state_path.write_text(json.dumps(current_state), encoding="utf-8")
messages_log = HOME / "logs" / "messages.jsonl" prior_size = messages_log.stat().st_size if messages_log.exists() else 0
cmd = [ "gemini", "-p", prompt, "--output-format", "stream-json", "--approval-mode", "yolo", ]
print(f"[orchestrator] Invoking gemini -p ({len(prompt)} chars)...", file=sys.stderr)
proc = subprocess.run( cmd, capture_output=True, text=True, timeout=timeout, cwd=str(HOME), stdin=subprocess.DEVNULL, )
duration = time.time() - start_time
log_event("gemini", len(prompt), session_id, proc.returncode, duration)
# Capture messages new_messages = read_new_messages(messages_log, prior_size)
if proc.returncode != 0: print(f"[orchestrator] Gemini exited with code {proc.returncode}", file=sys.stderr) if proc.stderr: print(f"[orchestrator] stderr: {proc.stderr[:500]}", file=sys.stderr) write_session_log("gemini", prompt, session_id, new_messages, duration, proc.returncode) return
print(f"[orchestrator] Done. Session: {session_id}", file=sys.stderr)
if not new_messages: print(f"[orchestrator] WARNING: Gemini responded with no send_message calls.", file=sys.stderr)
for msg in new_messages: print(f"Dash: {msg}")
write_session_log("gemini", prompt, session_id, new_messages, duration, proc.returncode)
if __name__ == "__main__": parser = argparse.ArgumentParser(description="Dash Orchestrator") parser.add_argument("event", nargs="?", default=None, help="The event to process") parser.add_argument("--harness", choices=["claude", "gemini"], default="gemini", help="The AI harness to use") parser.add_argument("--tick-type", choices=["admin_message", "operational_check", "deep_reflection"], default="admin_message", help="The type of tick/event") parser.add_argument("--dry-run", action="store_true", help="Print the prompt and exit") args = parser.parse_args() # Backwards compatibility: if event is missing and it's an admin_message, use default perch text event = args.event if event is None: if args.tick_type == "admin_message": event = "" else: event = "" # Ignored by templates prompt = build_prompt(event, tick_type=args.tick_type) if args.dry_run: print(prompt) sys.exit(0) if args.harness == "gemini": invoke_gemini(prompt) else: invoke_claude(prompt)