1use super::context::AgentContext;
17use super::tool_registry::ToolHandler;
18use crate::team_message_router::TeamMessageRouter;
19use async_trait::async_trait;
20use nexo_broker::AnyBroker;
21use nexo_config::types::team::TeamPolicy;
22use nexo_llm::ToolDef;
23use nexo_team_store::TeamStore;
24use serde_json::{json, Value};
25use std::sync::Arc;
26
27pub struct TeamTools {
37 pub store: Arc<dyn TeamStore>,
38 pub router: Arc<TeamMessageRouter<AnyBroker>>,
39 pub broker: AnyBroker,
40 pub agent_id: String,
41 pub current_goal_id: String,
49}
50
51impl TeamTools {
52 pub fn new(
53 store: Arc<dyn TeamStore>,
54 router: Arc<TeamMessageRouter<AnyBroker>>,
55 broker: AnyBroker,
56 agent_id: impl Into<String>,
57 current_goal_id: impl Into<String>,
58 ) -> Arc<Self> {
59 Arc::new(Self {
60 store,
61 router,
62 broker,
63 agent_id: agent_id.into(),
64 current_goal_id: current_goal_id.into(),
65 })
66 }
67
68 pub(super) fn policy_for(&self, ctx: &super::context::AgentContext) -> TeamPolicy {
73 ctx.effective_policy().team.clone()
74 }
75}
76
77pub struct TeamCreateTool {
84 pub inner: Arc<TeamTools>,
85}
86
87impl TeamCreateTool {
88 pub fn new(inner: Arc<TeamTools>) -> Self {
89 Self { inner }
90 }
91
92 pub fn tool_def() -> ToolDef {
93 ToolDef {
94 name: "TeamCreate".to_string(),
95 description: r#"Create a new named team for coordinated parallel work. Returns a `team_id` you reference in subsequent calls.
96
97After creating the team, members can be added by the operator (via `nexo team add-member` CLI — Phase 79.6.b) or registered programmatically through the team store. The MVP exposes `TeamSendMessage` (DM + broadcast), `TeamList`, and `TeamStatus` so the lead can coordinate already-spawned members. Direct goal-spawn-as-teammate from inside a turn lands in Phase 79.6.b.
98
99Caps: 8 members per team (incl. lead), 4 concurrent teams per agent, 24 h idle timeout."#
100 .to_string(),
101 parameters: json!({
102 "type": "object",
103 "properties": {
104 "team_name": {
105 "type": "string",
106 "description": "Human-readable name. Sanitized to [a-z0-9-]+ for the team_id (unique). Max 64 chars."
107 },
108 "description": {
109 "type": "string",
110 "description": "Optional one-line summary of what the team is doing."
111 },
112 "agent_type": {
113 "type": "string",
114 "description": "Role label for the lead (e.g. \"coordinator\", \"researcher\"). Free-form; surfaces in TeamStatus."
115 },
116 "worktree_per_member": {
117 "type": "boolean",
118 "description": "When true, each member gets its own git worktree under <workspace>/.team/<team_id>/<member_name>/. Default false (members share the workspace)."
119 }
120 },
121 "required": ["team_name"]
122 }),
123 }
124 }
125}
126
127pub struct TeamDeleteTool {
128 pub inner: Arc<TeamTools>,
129}
130
131impl TeamDeleteTool {
132 pub fn new(inner: Arc<TeamTools>) -> Self {
133 Self { inner }
134 }
135
136 pub fn tool_def() -> ToolDef {
137 ToolDef {
138 name: "TeamDelete".to_string(),
139 description: r#"Disband a team and clean up its members + DM topics. Refuses with `BlockedByActiveMembers` when any member's goal is still `Running`; send `[shutdown_request]` via `TeamSendMessage { to: "broadcast", message: { type: "shutdown_request" } }` first and wait for members to idle.
140
141Idle members are gracefully cancelled. After the cleanup, the team_id is reusable; tasks + audit rows persist in the store under the deleted team."#
142 .to_string(),
143 parameters: json!({
144 "type": "object",
145 "properties": {
146 "team_id": {
147 "type": "string",
148 "description": "Sanitized id returned by TeamCreate. Must be a team this agent leads."
149 }
150 },
151 "required": ["team_id"]
152 }),
153 }
154 }
155}
156
157pub struct TeamSendMessageTool {
158 pub inner: Arc<TeamTools>,
159}
160
161impl TeamSendMessageTool {
162 pub fn new(inner: Arc<TeamTools>) -> Self {
163 Self { inner }
164 }
165
166 pub fn tool_def() -> ToolDef {
167 ToolDef {
168 name: "TeamSendMessage".to_string(),
169 description: r#"Send a message to a teammate (point-to-point) or to all members (broadcast). Wakes idle teammates. Body capped at 64 KiB serialised JSON.
170
171`to`: member name (e.g. "researcher") for DM, or the literal string "broadcast" (lead only).
172`message`: free-form text OR a structured message:
173- `{ "type": "shutdown_request", "reason"?: "..." }` — ask members to wind down.
174- `{ "type": "shutdown_response", "request_id": "...", "approved": true|false, "reason"?: "..." }`.
175- `{ "type": "task_assigned", "task_id": "..." }`.
176- `{ "type": "done", "result"?: ... }`.
177
178Only the lead can broadcast. Any team member can DM another."#
179 .to_string(),
180 parameters: json!({
181 "type": "object",
182 "properties": {
183 "team_id": {
184 "type": "string",
185 "description": "Sanitized id returned by TeamCreate."
186 },
187 "to": {
188 "type": "string",
189 "description": "Member name (point-to-point) or \"broadcast\" (lead only)."
190 },
191 "message": {
192 "description": "String or structured object {type, ...}. Capped at 64 KiB."
193 }
194 },
195 "required": ["team_id", "to", "message"]
196 }),
197 }
198 }
199}
200
201pub struct TeamListTool {
202 pub inner: Arc<TeamTools>,
203}
204
205impl TeamListTool {
206 pub fn new(inner: Arc<TeamTools>) -> Self {
207 Self { inner }
208 }
209
210 pub fn tool_def() -> ToolDef {
211 ToolDef {
212 name: "TeamList".to_string(),
213 description: "List teams owned by this agent (read-only). Pass `active_only: true` to exclude soft-deleted teams.".to_string(),
214 parameters: json!({
215 "type": "object",
216 "properties": {
217 "active_only": {
218 "type": "boolean",
219 "description": "Default true — only return teams with `deleted_at IS NULL`."
220 }
221 }
222 }),
223 }
224 }
225}
226
227pub struct TeamStatusTool {
228 pub inner: Arc<TeamTools>,
229}
230
231impl TeamStatusTool {
232 pub fn new(inner: Arc<TeamTools>) -> Self {
233 Self { inner }
234 }
235
236 pub fn tool_def() -> ToolDef {
237 ToolDef {
238 name: "TeamStatus".to_string(),
239 description: "Detailed status of one team: members + their last_active_at + agent_type + idle/running. Read-only. Caller must be the lead OR a current member (NotMember otherwise).".to_string(),
240 parameters: json!({
241 "type": "object",
242 "properties": {
243 "team_id": {
244 "type": "string",
245 "description": "Sanitized id returned by TeamCreate."
246 }
247 },
248 "required": ["team_id"]
249 }),
250 }
251 }
252}
253
254use nexo_team_store::{
274 sanitize_name, validate_member_name_for_lead, validate_team_name, TeamEventRow, TeamMemberRow,
275 TeamRow, TeamStoreError, DM_BODY_MAX_BYTES, TEAM_LEAD_NAME,
276};
277
278fn now_ts() -> i64 {
279 chrono::Utc::now().timestamp()
280}
281
282fn new_event_id() -> String {
283 uuid::Uuid::new_v4().to_string()
284}
285
286fn err(kind: &str, msg: impl Into<String>) -> Value {
287 json!({
288 "ok": false,
289 "kind": kind,
290 "error": msg.into(),
291 })
292}
293
294async fn record_event(
297 inner: &TeamTools,
298 team_id: &str,
299 kind: &str,
300 actor: Option<&str>,
301 payload: Value,
302) {
303 let row = TeamEventRow {
304 event_id: new_event_id(),
305 team_id: team_id.to_string(),
306 kind: kind.to_string(),
307 actor_member_name: actor.map(str::to_string),
308 payload_json: payload.to_string(),
309 created_at: now_ts(),
310 };
311 if let Err(e) = inner.store.record_event(&row).await {
312 tracing::warn!(
313 target: "team::audit",
314 team_id,
315 kind,
316 error = %e,
317 "[team] audit record failed"
318 );
319 }
320}
321
322#[async_trait]
323impl ToolHandler for TeamCreateTool {
324 async fn call(&self, ctx: &AgentContext, args: Value) -> anyhow::Result<Value> {
325 if ctx.is_teammate() {
329 return Ok(err(
330 "TeammateCannotSpawnTeammate",
331 "teammates cannot spawn other teams",
332 ));
333 }
334 let team_policy = self.inner.policy_for(ctx);
339 if !team_policy.tool_enabled() {
340 return Ok(err(
341 "TeamingDisabled",
342 "team feature not enabled for this agent",
343 ));
344 }
345
346 let team_name = match args.get("team_name").and_then(|v| v.as_str()) {
347 Some(s) => s.to_string(),
348 None => return Ok(err("Wire", "TeamCreate requires `team_name` (string)")),
349 };
350 let description = args
351 .get("description")
352 .and_then(|v| v.as_str())
353 .map(str::to_string);
354 let agent_type = args
355 .get("agent_type")
356 .and_then(|v| v.as_str())
357 .map(str::to_string);
358 let worktree_default = team_policy.worktree_per_member;
359 let worktree_per_member = args
360 .get("worktree_per_member")
361 .and_then(|v| v.as_bool())
362 .unwrap_or(worktree_default);
363
364 let team_id = match validate_team_name(&team_name) {
365 Ok(s) => s,
366 Err(_) => {
367 return Ok(err(
368 "InvalidName",
369 format!("invalid team_name: `{team_name}`"),
370 ))
371 }
372 };
373
374 let active_count = self
376 .inner
377 .store
378 .count_active_for_agent(&self.inner.agent_id)
379 .await
380 .map_err(|e| anyhow::anyhow!("count_active_for_agent: {e}"))?;
381 let cap = team_policy.effective_max_concurrent() as usize;
382 if active_count >= cap {
383 tracing::warn!(
384 target: "team::cap_exceeded",
385 agent = %self.inner.agent_id,
386 count = active_count,
387 cap,
388 "[team] ConcurrentCapExceeded"
389 );
390 return Ok(json!({
391 "ok": false,
392 "kind": "ConcurrentCapExceeded",
393 "count": active_count,
394 "cap": cap,
395 "error": format!("agent already leads {active_count} active teams (cap {cap})"),
396 }));
397 }
398
399 let now = now_ts();
400 let team_row = TeamRow {
401 team_id: team_id.clone(),
402 display_name: team_name.clone(),
403 description: description.clone(),
404 lead_agent_id: self.inner.agent_id.clone(),
405 lead_goal_id: self.inner.current_goal_id.clone(),
406 flow_id: team_id.clone(),
410 worktree_per_member,
411 created_at: now,
412 deleted_at: None,
413 last_active_at: now,
414 };
415
416 if let Err(e) = self.inner.store.create_team(&team_row).await {
417 return match e {
418 TeamStoreError::TeamNameTaken(existing) => Ok(json!({
419 "ok": false,
420 "kind": "TeamNameTaken",
421 "existing_team_id": existing,
422 "error": format!("team `{team_id}` already exists"),
423 })),
424 other => Err(anyhow::anyhow!("create_team: {other}")),
425 };
426 }
427
428 let lead_name = TEAM_LEAD_NAME.to_string();
430 let lead_member = TeamMemberRow {
431 team_id: team_id.clone(),
432 name: lead_name.clone(),
433 agent_id: self.inner.agent_id.clone(),
434 agent_type: agent_type.clone(),
435 model: None,
436 goal_id: self.inner.current_goal_id.clone(),
437 worktree_path: None,
438 joined_at: now,
439 is_active: true,
440 last_active_at: now,
441 };
442 if let Err(e) = self.inner.store.add_member(&lead_member).await {
443 tracing::warn!(
444 target: "team::create",
445 team_id = %team_id,
446 error = %e,
447 "[team] lead row insert failed — soft-deleting team"
448 );
449 let _ = self.inner.store.soft_delete_team(&team_id, now).await;
450 return Err(anyhow::anyhow!("add_member(lead): {e}"));
451 }
452
453 record_event(
454 &self.inner,
455 &team_id,
456 "team_created",
457 Some(&lead_name),
458 json!({
459 "display_name": team_name,
460 "description": description,
461 "lead_agent_id": self.inner.agent_id,
462 "agent_type": agent_type,
463 "worktree_per_member": worktree_per_member,
464 }),
465 )
466 .await;
467
468 tracing::info!(
469 target: "team::create",
470 team_id = %team_id,
471 agent = %self.inner.agent_id,
472 "[team] created"
473 );
474
475 Ok(json!({
476 "ok": true,
477 "team_id": team_id,
478 "lead_agent_id": self.inner.agent_id,
479 "lead_member_name": lead_name,
480 "flow_id": team_row.flow_id,
481 "created_at": now,
482 "instructions": "Members are added by the operator (Phase 79.6.b CLI) or by the runtime when sub-goals spawn. From inside a turn, the model coordinates members via `TeamSendMessage` (DM or broadcast) + `TeamStatus`. Wind down via `TeamSendMessage { to: \"broadcast\", message: { type: \"shutdown_request\" } }` then `TeamDelete`."
483 }))
484 }
485}
486
487#[async_trait]
488impl ToolHandler for TeamDeleteTool {
489 async fn call(&self, ctx: &AgentContext, args: Value) -> anyhow::Result<Value> {
490 if ctx.is_teammate() {
491 return Ok(err(
492 "TeammateCannotDeleteTeam",
493 "only the team lead can delete the team (you are running as a teammate)",
494 ));
495 }
496 if !self.inner.policy_for(ctx).tool_enabled() {
497 return Ok(err("TeamingDisabled", "team feature not enabled"));
498 }
499 let team_id_raw = match args.get("team_id").and_then(|v| v.as_str()) {
500 Some(s) => s.to_string(),
501 None => return Ok(err("Wire", "TeamDelete requires `team_id`")),
502 };
503 let team_id = sanitize_name(&team_id_raw);
504
505 let team = match self.inner.store.get_team(&team_id).await {
506 Ok(Some(t)) => t,
507 Ok(None) => return Ok(err("TeamNotFound", format!("team `{team_id}` not found"))),
508 Err(e) => return Err(anyhow::anyhow!("get_team: {e}")),
509 };
510 if team.lead_agent_id != self.inner.agent_id {
511 return Ok(err(
512 "NotLeader",
513 format!(
514 "agent `{}` is not the lead of team `{team_id}`",
515 self.inner.agent_id
516 ),
517 ));
518 }
519 if team.deleted_at.is_some() {
520 return Ok(err(
521 "AlreadyDeleted",
522 format!("team `{team_id}` was already soft-deleted"),
523 ));
524 }
525
526 let members = self
530 .inner
531 .store
532 .list_members(&team_id)
533 .await
534 .map_err(|e| anyhow::anyhow!("list_members: {e}"))?;
535 let active_non_lead: Vec<&TeamMemberRow> = members
536 .iter()
537 .filter(|m| m.is_active && m.name != TEAM_LEAD_NAME)
538 .collect();
539 if !active_non_lead.is_empty() {
540 let names: Vec<String> = active_non_lead.iter().map(|m| m.name.clone()).collect();
541 return Ok(json!({
542 "ok": false,
543 "kind": "BlockedByActiveMembers",
544 "names": names,
545 "error": format!(
546 "team `{team_id}` has {} active member(s): {}. Send `TeamSendMessage {{ to: \"broadcast\", message: {{ type: \"shutdown_request\" }} }}` and wait for them to idle before deleting.",
547 active_non_lead.len(),
548 names.join(", ")
549 ),
550 }));
551 }
552
553 let now = now_ts();
554 if let Err(e) = self.inner.store.soft_delete_team(&team_id, now).await {
555 return Err(anyhow::anyhow!("soft_delete_team: {e}"));
556 }
557 self.inner.router.drop_team(&team_id);
558
559 let members_cleaned = members.len();
560 record_event(
561 &self.inner,
562 &team_id,
563 "team_deleted",
564 Some(TEAM_LEAD_NAME),
565 json!({ "members_cleaned": members_cleaned }),
566 )
567 .await;
568
569 tracing::info!(
570 target: "team::delete",
571 team_id = %team_id,
572 members_cleaned,
573 "[team] deleted"
574 );
575
576 Ok(json!({
577 "ok": true,
578 "team_id": team_id,
579 "members_cleaned": members_cleaned,
580 "force_killed": 0,
585 }))
586 }
587}
588
589#[async_trait]
590impl ToolHandler for TeamSendMessageTool {
591 async fn call(&self, ctx: &AgentContext, args: Value) -> anyhow::Result<Value> {
592 if !self.inner.policy_for(ctx).tool_enabled() {
593 return Ok(err("TeamingDisabled", "team feature not enabled"));
594 }
595 let team_id_raw = match args.get("team_id").and_then(|v| v.as_str()) {
596 Some(s) => s.to_string(),
597 None => return Ok(err("Wire", "TeamSendMessage requires `team_id`")),
598 };
599 let team_id = sanitize_name(&team_id_raw);
600 let to = match args.get("to").and_then(|v| v.as_str()) {
601 Some(s) => s.to_string(),
602 None => {
603 return Ok(err(
604 "Wire",
605 "TeamSendMessage requires `to` (member name or \"broadcast\")",
606 ))
607 }
608 };
609 let body = args.get("message").cloned().unwrap_or(Value::Null);
610 if body.is_null() {
611 return Ok(err(
612 "Wire",
613 "TeamSendMessage requires `message` (string or structured object)",
614 ));
615 }
616
617 let body_len = serde_json::to_vec(&body)
619 .map(|b| b.len())
620 .unwrap_or(usize::MAX);
621 if body_len > DM_BODY_MAX_BYTES {
622 return Ok(json!({
623 "ok": false,
624 "kind": "BodyTooLarge",
625 "actual": body_len,
626 "max": DM_BODY_MAX_BYTES,
627 "error": format!("message body is {body_len} bytes (max {DM_BODY_MAX_BYTES})"),
628 }));
629 }
630
631 let team = match self.inner.store.get_team(&team_id).await {
632 Ok(Some(t)) => t,
633 Ok(None) => return Ok(err("TeamNotFound", format!("team `{team_id}` not found"))),
634 Err(e) => return Err(anyhow::anyhow!("get_team: {e}")),
635 };
636 if team.deleted_at.is_some() {
637 return Ok(err(
638 "TeamDeleted",
639 format!("team `{team_id}` is soft-deleted"),
640 ));
641 }
642
643 let is_lead = team.lead_agent_id == self.inner.agent_id
645 && (ctx.team_member_name.as_deref() == Some(TEAM_LEAD_NAME) || !ctx.is_teammate());
646 let caller_member_name = if is_lead {
647 TEAM_LEAD_NAME.to_string()
648 } else {
649 let members = self
651 .inner
652 .store
653 .list_members(&team_id)
654 .await
655 .map_err(|e| anyhow::anyhow!("list_members: {e}"))?;
656 match members
657 .iter()
658 .find(|m| m.agent_id == self.inner.agent_id)
659 .map(|m| m.name.clone())
660 {
661 Some(n) => n,
662 None => {
663 return Ok(err(
664 "NotMember",
665 format!(
666 "agent `{}` is not a member of team `{team_id}`",
667 self.inner.agent_id
668 ),
669 ))
670 }
671 }
672 };
673
674 if to == "broadcast" {
676 if !is_lead {
677 return Ok(err(
678 "OnlyLeadCanBroadcast",
679 "only the team lead can publish to `broadcast`",
680 ));
681 }
682 self.inner
683 .router
684 .publish_broadcast(&team_id, &caller_member_name, body.clone())
685 .await
686 .map_err(|e| anyhow::anyhow!("publish_broadcast: {e}"))?;
687 record_event(
688 &self.inner,
689 &team_id,
690 "broadcast_sent",
691 Some(&caller_member_name),
692 json!({ "body_bytes": body_len }),
693 )
694 .await;
695 let _ = self.inner.store.touch_team(&team_id, now_ts()).await;
697 tracing::info!(
698 target: "team::broadcast_sent",
699 team_id = %team_id,
700 from = %caller_member_name,
701 body_bytes = body_len,
702 "[team] broadcast"
703 );
704 return Ok(json!({
705 "ok": true,
706 "team_id": team_id,
707 "to": "broadcast",
708 }));
709 }
710
711 let target_name = match validate_member_name_for_lead(&to) {
713 Ok(n) => n,
714 Err(_) => return Ok(err("InvalidMemberName", format!("invalid `to`: {to}"))),
715 };
716 let correlation_id = uuid::Uuid::new_v4().to_string();
717 self.inner
718 .router
719 .publish_dm(
720 &team_id,
721 &caller_member_name,
722 &target_name,
723 body.clone(),
724 Some(correlation_id.clone()),
725 )
726 .await
727 .map_err(|e| anyhow::anyhow!("publish_dm: {e}"))?;
728 record_event(
729 &self.inner,
730 &team_id,
731 "dm_sent",
732 Some(&caller_member_name),
733 json!({
734 "to": target_name,
735 "body_bytes": body_len,
736 "correlation_id": correlation_id,
737 }),
738 )
739 .await;
740 let _ = self.inner.store.touch_team(&team_id, now_ts()).await;
741 tracing::info!(
742 target: "team::dm_sent",
743 team_id = %team_id,
744 from = %caller_member_name,
745 to = %target_name,
746 body_bytes = body_len,
747 "[team] DM"
748 );
749 Ok(json!({
750 "ok": true,
751 "team_id": team_id,
752 "to": target_name,
753 "correlation_id": correlation_id,
754 }))
755 }
756}
757
758#[async_trait]
759impl ToolHandler for TeamListTool {
760 async fn call(&self, ctx: &AgentContext, args: Value) -> anyhow::Result<Value> {
761 if !self.inner.policy_for(ctx).tool_enabled() {
762 return Ok(err("TeamingDisabled", "team feature not enabled"));
763 }
764 let active_only = args
765 .get("active_only")
766 .and_then(|v| v.as_bool())
767 .unwrap_or(true);
768 let teams = self
769 .inner
770 .store
771 .list_teams(Some(&self.inner.agent_id), active_only)
772 .await
773 .map_err(|e| anyhow::anyhow!("list_teams: {e}"))?;
774 let json_teams: Vec<Value> = teams
775 .iter()
776 .map(|t| {
777 json!({
778 "team_id": t.team_id,
779 "display_name": t.display_name,
780 "description": t.description,
781 "lead_agent_id": t.lead_agent_id,
782 "created_at": t.created_at,
783 "deleted_at": t.deleted_at,
784 "last_active_at": t.last_active_at,
785 })
786 })
787 .collect();
788 Ok(json!({
789 "ok": true,
790 "n": teams.len(),
791 "teams": json_teams,
792 }))
793 }
794}
795
796#[async_trait]
797impl ToolHandler for TeamStatusTool {
798 async fn call(&self, ctx: &AgentContext, args: Value) -> anyhow::Result<Value> {
799 if !self.inner.policy_for(ctx).tool_enabled() {
800 return Ok(err("TeamingDisabled", "team feature not enabled"));
801 }
802 let team_id_raw = match args.get("team_id").and_then(|v| v.as_str()) {
803 Some(s) => s.to_string(),
804 None => return Ok(err("Wire", "TeamStatus requires `team_id`")),
805 };
806 let team_id = sanitize_name(&team_id_raw);
807 let team = match self.inner.store.get_team(&team_id).await {
808 Ok(Some(t)) => t,
809 Ok(None) => return Ok(err("TeamNotFound", format!("team `{team_id}` not found"))),
810 Err(e) => return Err(anyhow::anyhow!("get_team: {e}")),
811 };
812 let members = self
813 .inner
814 .store
815 .list_members(&team_id)
816 .await
817 .map_err(|e| anyhow::anyhow!("list_members: {e}"))?;
818 let is_lead = team.lead_agent_id == self.inner.agent_id;
820 let is_member = members.iter().any(|m| m.agent_id == self.inner.agent_id);
821 if !is_lead && !is_member {
822 let _ = ctx; return Ok(err(
827 "NotMember",
828 "only the team lead or a current member may read the team status",
829 ));
830 }
831
832 let n_running = members.iter().filter(|m| m.is_active).count();
833 let n_idle = members.len().saturating_sub(n_running);
834
835 let json_members: Vec<Value> = members
836 .iter()
837 .map(|m| {
838 json!({
839 "name": m.name,
840 "agent_id": m.agent_id,
841 "agent_type": m.agent_type,
842 "model": m.model,
843 "joined_at": m.joined_at,
844 "is_active": m.is_active,
845 "last_active_at": m.last_active_at,
846 })
847 })
848 .collect();
849
850 Ok(json!({
851 "ok": true,
852 "team": {
853 "team_id": team.team_id,
854 "display_name": team.display_name,
855 "description": team.description,
856 "lead_agent_id": team.lead_agent_id,
857 "created_at": team.created_at,
858 "deleted_at": team.deleted_at,
859 "last_active_at": team.last_active_at,
860 "worktree_per_member": team.worktree_per_member,
861 },
862 "members": json_members,
863 "task_summary": {
866 "pending": 0,
867 "running": n_running,
868 "done": n_idle,
869 },
870 }))
871 }
872}
873
874#[cfg(test)]
875mod tests {
876 use super::*;
877
878 #[test]
879 fn team_create_def_advertises_required_team_name() {
880 let def = TeamCreateTool::tool_def();
881 assert_eq!(def.name, "TeamCreate");
882 let required = def.parameters["required"].as_array().unwrap();
883 assert!(required.iter().any(|v| v == "team_name"));
884 let props = def.parameters["properties"].as_object().unwrap();
886 for k in [
887 "team_name",
888 "description",
889 "agent_type",
890 "worktree_per_member",
891 ] {
892 assert!(props.contains_key(k), "missing property `{k}`");
893 }
894 }
895
896 #[test]
897 fn team_delete_def_requires_team_id() {
898 let def = TeamDeleteTool::tool_def();
899 assert_eq!(def.name, "TeamDelete");
900 let required = def.parameters["required"].as_array().unwrap();
901 assert_eq!(required.len(), 1);
902 assert_eq!(required[0], "team_id");
903 }
904
905 #[test]
906 fn team_send_message_def_requires_team_id_to_message() {
907 let def = TeamSendMessageTool::tool_def();
908 assert_eq!(def.name, "TeamSendMessage");
909 let required = def.parameters["required"].as_array().unwrap();
910 let names: Vec<&str> = required.iter().map(|v| v.as_str().unwrap()).collect();
911 assert!(names.contains(&"team_id"));
912 assert!(names.contains(&"to"));
913 assert!(names.contains(&"message"));
914 }
915
916 #[test]
917 fn team_list_def_optional_active_only() {
918 let def = TeamListTool::tool_def();
919 assert_eq!(def.name, "TeamList");
920 assert!(def.parameters.get("required").is_none());
922 let props = def.parameters["properties"].as_object().unwrap();
923 assert!(props.contains_key("active_only"));
924 }
925
926 #[test]
927 fn team_status_def_requires_team_id() {
928 let def = TeamStatusTool::tool_def();
929 assert_eq!(def.name, "TeamStatus");
930 let required = def.parameters["required"].as_array().unwrap();
931 assert_eq!(required.len(), 1);
932 assert_eq!(required[0], "team_id");
933 }
934
935 use crate::session::SessionManager;
942 use nexo_broker::AnyBroker;
943 use nexo_config::types::agents::{
944 AgentConfig, AgentRuntimeConfig, DreamingYamlConfig, HeartbeatConfig, ModelConfig,
945 OutboundAllowlistConfig, WorkspaceGitConfig,
946 };
947 use nexo_team_store::SqliteTeamStore;
948 use std::sync::Arc;
949
950 fn agent_cfg_with_team(team: TeamPolicy) -> AgentConfig {
951 AgentConfig {
952 id: "cody".into(),
953 model: ModelConfig {
954 provider: "x".into(),
955 model: "y".into(),
956 },
957 plugins: Vec::new(),
958 heartbeat: HeartbeatConfig::default(),
959 config: AgentRuntimeConfig::default(),
960 system_prompt: String::new(),
961 workspace: String::new(),
962 skills: Vec::new(),
963 skills_dir: "./skills".into(),
964 skill_overrides: Default::default(),
965 transcripts_dir: String::new(),
966 dreaming: DreamingYamlConfig::default(),
967 workspace_git: WorkspaceGitConfig::default(),
968 tool_rate_limits: None,
969 tool_args_validation: None,
970 extra_docs: Vec::new(),
971 inbound_bindings: Vec::new(),
972 allowed_tools: Vec::new(),
973 sender_rate_limit: None,
974 allowed_delegates: Vec::new(),
975 accept_delegates_from: Vec::new(),
976 description: String::new(),
977 google_auth: None,
978 credentials: Default::default(),
979 link_understanding: serde_json::Value::Null,
980 web_search: serde_json::Value::Null,
981 pairing_policy: serde_json::Value::Null,
982 language: None,
983 locale_prompts: Default::default(),
984 outbound_allowlist: OutboundAllowlistConfig::default(),
985 context_optimization: None,
986 dispatch_policy: Default::default(),
987 plan_mode: Default::default(),
988 remote_triggers: Vec::new(),
989 lsp: nexo_config::types::lsp::LspPolicy::default(),
990 config_tool: nexo_config::types::config_tool::ConfigToolPolicy::default(),
991 team,
992 proactive: Default::default(),
993 repl: Default::default(),
994 auto_dream: None,
995 assistant_mode: None,
996 away_summary: None,
997 brief: None,
998 channels: None,
999 auto_approve: false,
1000 extract_memories: None,
1001 event_subscribers: Vec::new(),
1002 tenant_id: None,
1003 extensions_config: std::collections::BTreeMap::new(),
1004 active: true,
1005 }
1006 }
1007
1008 fn agent_cfg() -> AgentConfig {
1009 agent_cfg_with_team(TeamPolicy {
1010 enabled: true,
1011 ..TeamPolicy::default()
1012 })
1013 }
1014
1015 async fn ctx_lead_with_policy(policy: TeamPolicy) -> AgentContext {
1020 AgentContext::new(
1021 "cody",
1022 Arc::new(agent_cfg_with_team(policy)),
1023 AnyBroker::local(),
1024 Arc::new(SessionManager::new(std::time::Duration::from_secs(60), 8)),
1025 )
1026 }
1027
1028 async fn ctx_lead() -> AgentContext {
1029 AgentContext::new(
1030 "cody",
1031 Arc::new(agent_cfg()),
1032 AnyBroker::local(),
1033 Arc::new(SessionManager::new(std::time::Duration::from_secs(60), 8)),
1034 )
1035 }
1036
1037 async fn ctx_teammate() -> AgentContext {
1038 ctx_lead().await.with_team("feature-x", "researcher")
1039 }
1040
1041 async fn build_inner(agent_id: &str) -> Arc<TeamTools> {
1044 let broker = Arc::new(AnyBroker::local());
1045 let store: Arc<dyn TeamStore> = Arc::new(SqliteTeamStore::open_in_memory().await.unwrap());
1046 let router = TeamMessageRouter::new(broker.clone());
1047 let cancel = tokio_util::sync::CancellationToken::new();
1048 router.spawn(cancel);
1049 TeamTools::new(store, router, (*broker).clone(), agent_id, "lead-goal-1")
1053 }
1054
1055 #[tokio::test]
1056 async fn team_create_rejects_when_capability_disabled() {
1057 let inner = build_inner("cody").await;
1061 let tool = TeamCreateTool::new(inner);
1062 let res = tool
1063 .call(
1064 &ctx_lead_with_policy(TeamPolicy::default()).await,
1065 json!({ "team_name": "feature-x" }),
1066 )
1067 .await
1068 .unwrap();
1069 assert_eq!(res["ok"], false);
1070 assert_eq!(res["kind"], "TeamingDisabled");
1071 }
1072
1073 #[tokio::test]
1074 async fn team_create_returns_team_id_and_lead_member() {
1075 let inner = build_inner("cody").await;
1076 let tool = TeamCreateTool::new(inner.clone());
1077 let res = tool
1078 .call(
1079 &ctx_lead().await,
1080 json!({
1081 "team_name": "Feature-X",
1082 "description": "Build the new auth flow",
1083 "agent_type": "coordinator"
1084 }),
1085 )
1086 .await
1087 .unwrap();
1088 assert_eq!(res["ok"], true);
1089 assert_eq!(res["team_id"], "feature-x");
1090 assert_eq!(res["lead_member_name"], "team-lead");
1091
1092 let events = inner
1094 .store
1095 .tail_events(Some("feature-x"), 10)
1096 .await
1097 .unwrap();
1098 assert!(events.iter().any(|e| e.kind == "team_created"));
1099 }
1100
1101 #[tokio::test]
1102 async fn team_create_rejects_collision() {
1103 let inner = build_inner("cody").await;
1104 let tool = TeamCreateTool::new(inner);
1105 let _ = tool
1106 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1107 .await
1108 .unwrap();
1109 let res = tool
1110 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1111 .await
1112 .unwrap();
1113 assert_eq!(res["ok"], false);
1114 assert_eq!(res["kind"], "TeamNameTaken");
1115 assert_eq!(res["existing_team_id"], "feature-x");
1116 }
1117
1118 #[tokio::test]
1119 async fn team_create_rejects_when_at_max_concurrent() {
1120 let inner = build_inner("cody").await;
1124 let ctx = ctx_lead_with_policy(TeamPolicy {
1125 enabled: true,
1126 max_concurrent: 1,
1127 ..TeamPolicy::default()
1128 })
1129 .await;
1130 let tool = TeamCreateTool::new(inner);
1131 let _ = tool
1132 .call(&ctx, json!({ "team_name": "first" }))
1133 .await
1134 .unwrap();
1135 let res = tool
1136 .call(&ctx, json!({ "team_name": "second" }))
1137 .await
1138 .unwrap();
1139 assert_eq!(res["ok"], false);
1140 assert_eq!(res["kind"], "ConcurrentCapExceeded");
1141 assert_eq!(res["cap"], 1);
1142 }
1143
1144 #[tokio::test]
1145 async fn team_create_refuses_from_teammate_context() {
1146 let inner = build_inner("cody").await;
1147 let tool = TeamCreateTool::new(inner);
1148 let res = tool
1149 .call(&ctx_teammate().await, json!({ "team_name": "nested" }))
1150 .await
1151 .unwrap();
1152 assert_eq!(res["ok"], false);
1153 assert_eq!(res["kind"], "TeammateCannotSpawnTeammate");
1154 }
1155
1156 #[tokio::test]
1157 async fn team_delete_only_lead_can_invoke() {
1158 let inner_lead = build_inner("cody").await;
1161 TeamCreateTool::new(inner_lead.clone())
1162 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1163 .await
1164 .unwrap();
1165
1166 let inner_other = TeamTools::new(
1169 Arc::clone(&inner_lead.store),
1170 Arc::clone(&inner_lead.router),
1171 inner_lead.broker.clone(),
1172 "other-agent",
1173 "other-goal",
1174 );
1175 let res = TeamDeleteTool::new(inner_other)
1176 .call(
1177 &AgentContext::new(
1178 "other-agent",
1179 Arc::new(agent_cfg()),
1180 AnyBroker::local(),
1181 Arc::new(SessionManager::new(std::time::Duration::from_secs(60), 8)),
1182 ),
1183 json!({ "team_id": "feature-x" }),
1184 )
1185 .await
1186 .unwrap();
1187 assert_eq!(res["ok"], false);
1188 assert_eq!(res["kind"], "NotLeader");
1189 }
1190
1191 #[tokio::test]
1192 async fn team_delete_blocks_when_running_members() {
1193 let inner = build_inner("cody").await;
1194 TeamCreateTool::new(inner.clone())
1195 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1196 .await
1197 .unwrap();
1198 inner
1200 .store
1201 .add_member(&TeamMemberRow {
1202 team_id: "feature-x".into(),
1203 name: "researcher".into(),
1204 agent_id: "rsr".into(),
1205 agent_type: None,
1206 model: None,
1207 goal_id: "g".into(),
1208 worktree_path: None,
1209 joined_at: 0,
1210 is_active: true,
1211 last_active_at: 0,
1212 })
1213 .await
1214 .unwrap();
1215
1216 let res = TeamDeleteTool::new(inner)
1217 .call(&ctx_lead().await, json!({ "team_id": "feature-x" }))
1218 .await
1219 .unwrap();
1220 assert_eq!(res["ok"], false);
1221 assert_eq!(res["kind"], "BlockedByActiveMembers");
1222 let names = res["names"].as_array().unwrap();
1223 assert!(names.iter().any(|v| v == "researcher"));
1224 }
1225
1226 #[tokio::test]
1227 async fn team_delete_soft_deletes_and_records_event() {
1228 let inner = build_inner("cody").await;
1229 TeamCreateTool::new(inner.clone())
1230 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1231 .await
1232 .unwrap();
1233 let res = TeamDeleteTool::new(inner.clone())
1234 .call(&ctx_lead().await, json!({ "team_id": "feature-x" }))
1235 .await
1236 .unwrap();
1237 assert_eq!(res["ok"], true);
1238 assert_eq!(res["team_id"], "feature-x");
1239
1240 let team = inner.store.get_team("feature-x").await.unwrap().unwrap();
1241 assert!(team.deleted_at.is_some());
1242
1243 let events = inner
1244 .store
1245 .tail_events(Some("feature-x"), 10)
1246 .await
1247 .unwrap();
1248 assert!(events.iter().any(|e| e.kind == "team_deleted"));
1249 }
1250
1251 #[tokio::test]
1252 async fn team_send_message_rejects_oversized_body() {
1253 let inner = build_inner("cody").await;
1254 TeamCreateTool::new(inner.clone())
1255 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1256 .await
1257 .unwrap();
1258 let big = "x".repeat(DM_BODY_MAX_BYTES + 1);
1259 let res = TeamSendMessageTool::new(inner)
1260 .call(
1261 &ctx_lead().await,
1262 json!({
1263 "team_id": "feature-x",
1264 "to": "researcher",
1265 "message": big
1266 }),
1267 )
1268 .await
1269 .unwrap();
1270 assert_eq!(res["ok"], false);
1271 assert_eq!(res["kind"], "BodyTooLarge");
1272 }
1273
1274 #[tokio::test]
1275 async fn team_send_message_broadcast_requires_lead() {
1276 let inner = build_inner("cody").await;
1277 TeamCreateTool::new(inner.clone())
1278 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1279 .await
1280 .unwrap();
1281 inner
1283 .store
1284 .add_member(&TeamMemberRow {
1285 team_id: "feature-x".into(),
1286 name: "researcher".into(),
1287 agent_id: "cody".into(),
1288 agent_type: None,
1289 model: None,
1290 goal_id: "g".into(),
1291 worktree_path: None,
1292 joined_at: 0,
1293 is_active: true,
1294 last_active_at: 0,
1295 })
1296 .await
1297 .unwrap();
1298 let res = TeamSendMessageTool::new(inner)
1300 .call(
1301 &ctx_teammate().await,
1302 json!({
1303 "team_id": "feature-x",
1304 "to": "broadcast",
1305 "message": { "type": "shutdown_request" }
1306 }),
1307 )
1308 .await
1309 .unwrap();
1310 assert_eq!(res["ok"], false);
1311 assert_eq!(res["kind"], "OnlyLeadCanBroadcast");
1312 }
1313
1314 #[tokio::test]
1315 async fn team_send_message_dm_publishes_and_records() {
1316 let inner = build_inner("cody").await;
1317 TeamCreateTool::new(inner.clone())
1318 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1319 .await
1320 .unwrap();
1321 let _rx = inner.router.subscribe_member("feature-x", "researcher");
1323 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1324 let res = TeamSendMessageTool::new(inner.clone())
1325 .call(
1326 &ctx_lead().await,
1327 json!({
1328 "team_id": "feature-x",
1329 "to": "researcher",
1330 "message": { "ask": "ready?" }
1331 }),
1332 )
1333 .await
1334 .unwrap();
1335 assert_eq!(res["ok"], true);
1336 assert!(res["correlation_id"].is_string());
1337
1338 let events = inner
1339 .store
1340 .tail_events(Some("feature-x"), 20)
1341 .await
1342 .unwrap();
1343 assert!(events.iter().any(|e| e.kind == "dm_sent"));
1344 }
1345
1346 #[tokio::test]
1347 async fn team_list_filters_active_only() {
1348 let inner = build_inner("cody").await;
1349 TeamCreateTool::new(inner.clone())
1350 .call(&ctx_lead().await, json!({ "team_name": "a" }))
1351 .await
1352 .unwrap();
1353 TeamCreateTool::new(inner.clone())
1354 .call(&ctx_lead().await, json!({ "team_name": "b" }))
1355 .await
1356 .unwrap();
1357 TeamDeleteTool::new(inner.clone())
1358 .call(&ctx_lead().await, json!({ "team_id": "a" }))
1359 .await
1360 .unwrap();
1361 let res = TeamListTool::new(inner.clone())
1362 .call(&ctx_lead().await, json!({ "active_only": true }))
1363 .await
1364 .unwrap();
1365 assert_eq!(res["ok"], true);
1366 assert_eq!(res["n"], 1);
1367 let teams = res["teams"].as_array().unwrap();
1368 assert_eq!(teams[0]["team_id"], "b");
1369
1370 let all = TeamListTool::new(inner)
1371 .call(&ctx_lead().await, json!({ "active_only": false }))
1372 .await
1373 .unwrap();
1374 assert_eq!(all["n"], 2);
1375 }
1376
1377 #[tokio::test]
1378 async fn team_status_rejects_non_member() {
1379 let inner_lead = build_inner("cody").await;
1380 TeamCreateTool::new(inner_lead.clone())
1381 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1382 .await
1383 .unwrap();
1384
1385 let inner_other = TeamTools::new(
1388 Arc::clone(&inner_lead.store),
1389 Arc::clone(&inner_lead.router),
1390 inner_lead.broker.clone(),
1391 "stranger",
1392 "g",
1393 );
1394 let res = TeamStatusTool::new(inner_other)
1395 .call(
1396 &AgentContext::new(
1397 "stranger",
1398 Arc::new(agent_cfg()),
1399 AnyBroker::local(),
1400 Arc::new(SessionManager::new(std::time::Duration::from_secs(60), 8)),
1401 ),
1402 json!({ "team_id": "feature-x" }),
1403 )
1404 .await
1405 .unwrap();
1406 assert_eq!(res["ok"], false);
1407 assert_eq!(res["kind"], "NotMember");
1408 }
1409
1410 #[tokio::test]
1411 async fn team_status_returns_members_and_summary_for_lead() {
1412 let inner = build_inner("cody").await;
1413 TeamCreateTool::new(inner.clone())
1414 .call(&ctx_lead().await, json!({ "team_name": "feature-x" }))
1415 .await
1416 .unwrap();
1417 let res = TeamStatusTool::new(inner.clone())
1418 .call(&ctx_lead().await, json!({ "team_id": "feature-x" }))
1419 .await
1420 .unwrap();
1421 assert_eq!(res["ok"], true);
1422 assert_eq!(res["team"]["team_id"], "feature-x");
1423 let members = res["members"].as_array().unwrap();
1424 assert_eq!(members.len(), 1);
1426 assert_eq!(members[0]["name"], "team-lead");
1427 }
1428
1429 #[tokio::test]
1438 async fn team_create_policy_pulled_per_call_from_ctx() {
1439 let inner = build_inner("cody").await;
1440 let tool = TeamCreateTool::new(inner);
1441
1442 let ctx_cap1 = ctx_lead_with_policy(TeamPolicy {
1444 enabled: true,
1445 max_concurrent: 1,
1446 ..TeamPolicy::default()
1447 })
1448 .await;
1449 let r1 = tool
1450 .call(&ctx_cap1, json!({ "team_name": "first" }))
1451 .await
1452 .unwrap();
1453 assert_eq!(r1["ok"], true);
1454 let r2 = tool
1455 .call(&ctx_cap1, json!({ "team_name": "second" }))
1456 .await
1457 .unwrap();
1458 assert_eq!(r2["ok"], false);
1459 assert_eq!(r2["kind"], "ConcurrentCapExceeded");
1460
1461 let ctx_cap8 = ctx_lead_with_policy(TeamPolicy {
1464 enabled: true,
1465 max_concurrent: 8,
1466 ..TeamPolicy::default()
1467 })
1468 .await;
1469 let r3 = tool
1470 .call(&ctx_cap8, json!({ "team_name": "third" }))
1471 .await
1472 .unwrap();
1473 assert_eq!(r3["ok"], true, "after cap widened the call must succeed");
1474 }
1475
1476 #[tokio::test]
1480 async fn team_enabled_flip_picked_up_via_new_ctx() {
1481 let inner = build_inner("cody").await;
1482 let tool = TeamCreateTool::new(inner);
1483
1484 let ctx_off = ctx_lead_with_policy(TeamPolicy::default()).await;
1486 let r1 = tool
1487 .call(&ctx_off, json!({ "team_name": "early" }))
1488 .await
1489 .unwrap();
1490 assert_eq!(r1["kind"], "TeamingDisabled");
1491
1492 let ctx_on = ctx_lead_with_policy(TeamPolicy {
1494 enabled: true,
1495 ..TeamPolicy::default()
1496 })
1497 .await;
1498 let r2 = tool
1499 .call(&ctx_on, json!({ "team_name": "post-reload" }))
1500 .await
1501 .unwrap();
1502 assert_eq!(r2["ok"], true);
1503 }
1504}