Files
clawmates/crates/cm-api/src/routes/planner.rs
T
Omar SobhandClaude Opus 5 b1bdfbbf87 fix(provision): let callers declare write access instead of guessing from the role name
default_risk_profile_for_role decides whether a claw gets file edits, git and
shell by substring-matching its role against a fixed keyword list. On the
planner path that role string is free text the model invented for this
proposal, so a model's choice of wording silently decided tool access: a
proposed "implementation_lead" matches no keyword, lands research_readonly,
and then fails every file edit for a reason invisible from the role name.

TeamMemberInput and the planner's member schema now carry `needs_write`, and
resolve_risk_profile prefers it over the guess. The planner prompt asks for it
per member and says to grant write only to members that produce code or
commits. Absent (older clients, autoprovision, a model that omitted the field)
falls back to the old guess, so nothing changes for callers that don't set it.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
2026-07-30 10:53:08 -07:00

272 lines
15 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Master Planner — a chat with Claude Opus 4.8 that proposes a team of agents
//! (named, role'd, one model each) and then scaffolds it end-to-end: creates the
//! team + topology, refines + attaches a brain per agent, and sets up the nightly
//! loop. Replaces the old "+" deploy wizard.
use axum::extract::State;
use axum::response::sse::{Event, KeepAlive, Sse};
use axum::Json;
use serde::Deserialize;
use serde_json::{json, Value};
use std::convert::Infallible;
use crate::routes::claws::{apply_reference_to_claw, enhance_and_publish, extract_json};
use crate::routes::teams::TeamMemberInput;
use crate::{AppState, Authed};
fn sse(v: Value) -> Result<Event, Infallible> {
Ok(Event::default().data(v.to_string()))
}
const PLANNER_SYSTEM: &str = "You are the ClawMates Master Planner. You design teams of AI agents to \
accomplish a user's goal across the platform's hierarchy (organizations → companies → teams → agents). \
Have a brief, helpful conversation to understand the goal, then propose a concrete team. For the proposal: \
choose a sensible number of agents (usually 26), each with a UNIQUE human name, a clear ROLE, the best \
MODEL for that role (see catalog), a short brain_query (one domain keyword to fetch a starting brain from \
the registry, e.g. 'research', 'writing', 'data', 'security'), a focused system_prompt, and a one-line \
rationale. Pick a topology_kind from: hub_spoke, hierarchical, pipeline, mesh, flat. If the goal implies \
recurring/autonomous work (e.g. 'nightly'), include a schedule {cron (5-field), prompt (the mission the \
lead agent runs each cycle)}. Use web search to ground version-sensitive or current-fact claims. \
ALWAYS respond with STRICT JSON ONLY (no prose, no markdown), exactly: \
{\"reply\":\"<concise message to the user>\",\"proposal\":null|{\"team_name\":\"...\",\
\"topology_kind\":\"hub_spoke\",\"schedule\":null|{\"cron\":\"0 2 * * *\",\"prompt\":\"...\"},\
\"members\":[{\"name\":\"...\",\"role\":\"...\",\"model\":\"...\",\"brain_query\":\"...\",\
\"system_prompt\":\"...\",\"needs_write\":true|false,\"rationale\":\"...\"}]}}. Set proposal to null while \
still clarifying; include it once you have a concrete team. \n\n\
ACCESS: set \"needs_write\" per member. true grants file edits, git and shell; false is read-only \
research tools. Grant write only to members that actually produce code or commits — the rest read-only.\n\
\n\nMODELS (set each member's \"model\" to exactly one token):\n\
- claude — Claude Opus 4.8: strongest reasoning/planning; coordinators, hard analysis. Highest cost.\n\
- glm-4.7 — strong general reasoning (Z.ai); best cost/quality default for most workers.\n\
- glm-5.2 — GLM Opus-class for the hardest reasoning roles; higher cost.\n\
- kimi — excellent for code-heavy roles.\n\
- gemini — Gemini 2.5 Flash: very fast; classification, summarization, high-volume tasks.\n\
- groq — fastest/cheapest; simple sequential high-throughput steps.\n\
AGENT TOOLS each agent can use at runtime: web.search (find sources), browser.goto (fetch a URL), \
files.write (build a markdown vault in the shared drive), chat.send (delegate to teammates), \
routine.schedule (self-schedule).";
#[derive(Deserialize)]
pub struct PlannerMessage {
pub role: String,
pub content: String,
}
#[derive(Deserialize)]
pub struct PlannerChatRequest {
pub messages: Vec<PlannerMessage>,
/// Deploy mode: specialists | team | swarm | scheduled | triggered (default specialists).
#[serde(default)]
pub mode: String,
/// User-locked topology from the gallery. When set the planner is told to
/// use this kind verbatim (no rewrite). Ignored for `swarm` mode (the swarm
/// planner doesn't take a topology kind).
#[serde(default)]
pub topology_kind: Option<String>,
}
const SPECIALISTS_NOTE: &str = "\n\nMODE: Agent — SINGLE agent. The `members` array MUST contain EXACTLY ONE \
entry: a tightly-scoped domain specialist with a high-quality, focused system_prompt. Do NOT propose two, do NOT \
propose a team, do NOT include a coordinator. The member's brain_query MUST be a concrete domain brain keyword \
(e.g. 'rust-2024', 'react-native', 'pentest-web', 'db-ops'). `topology_kind` for a single agent is 'flat'. \
`team_name` should read like a personal handle for the agent, not a group name.";
const TEAM_NOTE: &str = "\n\nMODE: Team. Propose a balanced team of 4 to 6 agents that covers the goal \
end-to-end — a coordinator plus complementary roles. Bias toward roles that will actually be exercised each \
iteration; do not pad. Each member's brain_query should be a keyword the brain registry can resolve.";
const SCHEDULED_NOTE: &str = "\n\nMODE: Scheduled. Propose a team (28 agents is typical) and INCLUDE a \
schedule in the proposal. Recurring: {\"cron\":\"<5-field>\",\"prompt\":\"<mission run each cycle>\"}. \
One-time: {\"one_shot_at\":\"<RFC3339 UTC datetime>\",\"prompt\":\"<mission>\"}. The team is EPHEMERAL — \
provisioned at fire time and destroyed after the run — so keep member count tight.";
const TRIGGERED_NOTE: &str = "\n\nMODE: Triggered. Propose a team the user can fire by hitting a webhook \
(with a JSON payload the mission can reference). The team is EPHEMERAL — spun up on receipt, torn down when \
the run terminates — so keep the roster tight. Set schedule to {\"prompt\":\"<default task the webhook \
runs>\"} (NO cron / one_shot_at).";
const SWARM_SYSTEM: &str = "You are the planner for a self-verifying agent SWARM (Opus plans + verifies, a worker \
swarm executes, the loop repeats until every output passes). Swarms are for BIG parallel jobs — target 10 or more \
workers (`task_count >= 10`). The user describes a job; you turn it into a swarm spec. The CHECKLIST is the \
verification contract — each item must be objectively checkable per task (e.g. 'states a revenue figure', 'cites \
a resolvable source URL', 'no field left empty'). Have a brief conversation, then emit the spec. ALWAYS respond \
with STRICT JSON ONLY: {\"reply\":\"<concise message>\",\"swarm\":null|{\"goal\":\"<the decomposable job>\",\
\"checklist\":[\"...\",\"...\"],\"task_count\":<int>,\"worker_model\":\"auto\"}}. Set swarm to null while still \
clarifying; include it once the job + checklist are concrete.";
fn planner_system_for(mode: &str) -> String {
match mode {
"swarm" => SWARM_SYSTEM.to_string(),
"team" => format!("{PLANNER_SYSTEM}{TEAM_NOTE}"),
"scheduled" => format!("{PLANNER_SYSTEM}{SCHEDULED_NOTE}"),
"triggered" => format!("{PLANNER_SYSTEM}{TRIGGERED_NOTE}"),
_ => format!("{PLANNER_SYSTEM}{SPECIALISTS_NOTE}"),
}
}
/// `POST /api/planner/chat` — one planner turn (SSE: thinking → done{reply,proposal}).
pub async fn planner_chat(
State(state): State<AppState>,
Authed(user): Authed,
Json(body): Json<PlannerChatRequest>,
) -> impl axum::response::IntoResponse {
let runtime = state.runtime.clone();
let pool = state.pool.clone();
let ws = user.workspace_id;
let stream = async_stream::stream! {
yield sse(json!({"stage":"thinking","label":"Planning with Claude Opus 4.8…"}));
let agents = cm_db::repo::agents::roster(&pool, ws).await.unwrap_or_default();
let teams = cm_db::repo::teams::list_for_workspace(&pool, ws, 100).await.unwrap_or_default();
let agent_names: Vec<String> = agents.iter().take(40).map(|a| format!("{} ({})", a.name, a.job_title)).collect();
let team_names: Vec<String> = teams.iter().take(40).map(|t| t.name.clone()).collect();
let hierarchy = format!(
"CURRENT WORKSPACE: {} agents, {} teams.\nExisting agents: {}\nExisting teams: {}",
agents.len(), teams.len(),
if agent_names.is_empty() { "(none)".to_string() } else { agent_names.join(", ") },
if team_names.is_empty() { "(none)".to_string() } else { team_names.join(", ") },
);
let convo = body.messages.iter()
.map(|m| format!("{}: {}", if m.role == "user" { "USER" } else { "PLANNER" }, m.content))
.collect::<Vec<_>>().join("\n\n");
let topology_lock = match body.topology_kind.as_deref() {
Some(k) if !k.is_empty() && body.mode != "swarm" => format!(
"\n\n=== USER-LOCKED TOPOLOGY ===\nThe user has selected topology_kind='{k}' from the gallery. \
Use this kind verbatim in the proposal; do NOT choose a different one. Tailor the roles + \
member count to work well within this shape."
),
_ => String::new(),
};
let user_prompt = format!("{hierarchy}{topology_lock}\n\n=== CONVERSATION ===\n{convo}\n\nRespond now (JSON only).");
let system = planner_system_for(&body.mode);
let raw = match runtime.complete(&system, &user_prompt, "claude-opus-4-8", 8000, true).await {
Ok(t) => t,
Err(e) => { yield sse(json!({"stage":"error","label":format!("Opus error: {e}")})); return; }
};
match extract_json(&raw) {
Some(v) => {
let reply = v.get("reply").and_then(|x| x.as_str()).unwrap_or("").to_string();
let proposal = v.get("proposal").cloned().unwrap_or(Value::Null);
let swarm = v.get("swarm").cloned().unwrap_or(Value::Null);
yield sse(json!({"stage":"done","reply":reply,"proposal":proposal,"swarm":swarm}));
}
None => yield sse(json!({"stage":"done","reply":raw,"proposal":Value::Null,"swarm":Value::Null})),
}
};
Sse::new(Box::pin(stream)).keep_alive(KeepAlive::new())
}
#[derive(Deserialize)]
pub struct ScaffoldMember {
pub name: String,
pub role: String,
#[serde(default)]
pub model: String,
#[serde(default)]
pub brain_query: String,
#[serde(default)]
pub system_prompt: String,
/// Whether this member edits files / runs git, as declared by the planner.
/// Absent (older clients, or a model that omitted it) falls back to the
/// role-name guess in `RuntimeProvisioner::resolve_risk_profile`.
#[serde(default)]
pub needs_write: Option<bool>,
}
#[derive(Deserialize)]
pub struct ScaffoldSchedule {
#[serde(default)]
pub cron: String,
/// One-shot fire time (RFC3339 UTC). When set, takes precedence over `cron`.
#[serde(default)]
pub one_shot_at: String,
pub prompt: String,
}
#[derive(Deserialize)]
pub struct ScaffoldRequest {
pub team_name: String,
#[serde(default = "default_kind")]
pub topology_kind: String,
#[serde(default)]
pub schedule: Option<ScaffoldSchedule>,
pub members: Vec<ScaffoldMember>,
/// Original planner mode — decides the team's lifecycle. `scheduled` and
/// `triggered` produce `ephemeral` teams (torn down by the topology_worker
/// after the last run terminates). Anything else is `permanent`.
#[serde(default)]
pub mode: String,
}
fn default_kind() -> String {
"hub_spoke".to_string()
}
fn lifecycle_for(mode: &str) -> &'static str {
match mode {
"scheduled" | "triggered" => "ephemeral",
_ => "permanent",
}
}
/// `POST /api/planner/scaffold` — build the approved team (SSE progress): create
/// agents + topology, refine+attach a brain per agent, set up the nightly loop.
pub async fn planner_scaffold(
State(state): State<AppState>,
Authed(user): Authed,
Json(body): Json<ScaffoldRequest>,
) -> impl axum::response::IntoResponse {
let runtime = state.runtime.clone();
let stream = async_stream::stream! {
if body.members.is_empty() {
yield sse(json!({"stage":"error","pct":100,"label":"Empty proposal"}));
return;
}
let n = body.members.len().max(1);
yield sse(json!({"stage":"team","pct":8,"label":format!("Creating team “{}” ({} agents)…", body.team_name, n)}));
let members: Vec<TeamMemberInput> = body.members.iter().map(|m| TeamMemberInput {
role: m.role.clone(),
name: m.name.clone(),
model: if m.model.trim().is_empty() { "claude".to_string() } else { m.model.clone() },
system_prompt: m.system_prompt.clone(),
accent: String::new(),
needs_write: m.needs_write,
}).collect();
let lifecycle = lifecycle_for(&body.mode);
let (team_id, claw_ids) = match crate::routes::teams::build_team_with_lifecycle(&state, user.workspace_id, user.user_id, &body.team_name, &body.topology_kind, &members, lifecycle).await {
Ok(r) => r,
Err(_) => { yield sse(json!({"stage":"error","pct":100,"label":"Team creation failed"})); return; }
};
for (i, (cid, m)) in claw_ids.iter().zip(body.members.iter()).enumerate() {
let pct = 15 + (i as u32) * 70 / (n as u32);
yield sse(json!({"stage":"brain","pct":pct,"label":format!("Refining brain for {} ({})…", m.name, m.role)}));
let q = if m.brain_query.trim().is_empty() { m.role.clone() } else { m.brain_query.clone() };
let found = cm_brain::hub::list(&q).await.unwrap_or_default().into_iter().next().map(|b| b.reference);
if let Some(reference) = found {
let role_ctx = format!("Agent '{}', role '{}', on team '{}'. Mission: {}", m.name, m.role, body.team_name, m.system_prompt);
let refined = enhance_and_publish(&runtime, &reference, &role_ctx).await.unwrap_or(reference);
let id = cm_domain::AgentId::from(*cid);
let _ = apply_reference_to_claw(&state, id, &refined).await;
}
}
if let Some(sch) = &body.schedule {
// A `topology` action fires the whole team's stored graph as a durable
// run (not just one message to the lead).
if let Some(cid) = claw_ids.first() {
let id = cm_domain::AgentId::from(*cid);
let one_shot_at = sch.one_shot_at.trim();
if !one_shot_at.is_empty() {
// One-shot: fire once at the given datetime, never reschedule.
if let Ok(when) = time::OffsetDateTime::parse(one_shot_at, &time::format_description::well_known::Rfc3339) {
yield sse(json!({"stage":"routine","pct":92,"label":"Scheduling the one-time run…"}));
let action = json!({"topology": {"team_id": team_id.to_string(), "task": sch.prompt}, "one_shot": true});
let _ = cm_db::repo::routines::create(&state.pool, id, "Scheduled run", "0 0 1 1 *", action, when).await;
}
} else if !sch.cron.trim().is_empty() {
if let Ok(next) = cm_scheduler::next_occurrence(&sch.cron, time::OffsetDateTime::now_utc()) {
yield sse(json!({"stage":"routine","pct":92,"label":"Scheduling the recurring team loop…"}));
let action = json!({"topology": {"team_id": team_id.to_string(), "task": sch.prompt}});
let _ = cm_db::repo::routines::create(&state.pool, id, "Scheduled team loop", &sch.cron, action, next).await;
}
}
}
}
yield sse(json!({"stage":"done","pct":100,"label":"Team deployed","team_id":team_id.to_string()}));
};
Sse::new(Box::pin(stream)).keep_alive(KeepAlive::new())
}