Skip to content

Environment API reference

AgentSymposiumConfig dataclass

YAML-loadable configuration for an Agent Symposium environment.

Source code in src/ursa/environments/config.py
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
@dataclass(frozen=True)
class AgentSymposiumConfig:
    """YAML-loadable configuration for an Agent Symposium environment."""

    name: str
    group: str = "default"
    description: str | None = None
    organizer: EnvironmentMemberConfig = field(
        default_factory=lambda: EnvironmentMemberConfig(
            name="organizer", role="Symposium organizer", agent="ChatAgent"
        )
    )
    members: list[EnvironmentMemberConfig] = field(default_factory=list)
    workspace: str | None = None
    defaults: dict[str, Any] = field(default_factory=dict)
    revision_rounds: int = 1

    @classmethod
    def from_mapping(cls, data: Mapping[str, Any]) -> "AgentSymposiumConfig":
        raw = dict(data)
        if "organizer" in raw and isinstance(raw["organizer"], Mapping):
            raw["organizer"] = EnvironmentMemberConfig.from_mapping(
                raw["organizer"]
            )
        if "members" in raw:
            raw["members"] = [
                EnvironmentMemberConfig.from_mapping(member)
                for member in raw["members"]
            ]
        return cls(**raw)

AgentSymposiumEnvironment

Bases: BaseEnvironment

Competitive/collaborative multi-agent symposium environment.

A symposium organizer dispatches the same complex problem to several members or nested teams. Members first work independently. The environment then sends all writeups to each member for critical review, explicitly instructing them not to modify reviewed work and to only view/run code as needed for assessment. Finally, each member revises its own work using the feedback it received and insights gained from reviewing others. The organizer synthesizes the final symposium report.

Source code in src/ursa/environments/agent_symposium.py
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
class AgentSymposiumEnvironment(BaseEnvironment):
    """Competitive/collaborative multi-agent symposium environment.

    A symposium organizer dispatches the same complex problem to several members
    or nested teams. Members first work independently. The environment then sends
    all writeups to each member for critical review, explicitly instructing them
    not to modify reviewed work and to only view/run code as needed for assessment.
    Finally, each member revises its own work using the feedback it received and
    insights gained from reviewing others. The organizer synthesizes the final
    symposium report.
    """

    def __init__(
        self,
        llm: BaseChatModel,
        *,
        config: AgentSymposiumConfig
        | Mapping[str, Any]
        | str
        | Path
        | None = None,
        name: str | None = None,
        group: str | None = None,
        organizer: EnvironmentMemberConfig | Mapping[str, Any] | None = None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]]
        | None = None,
        workspace: str | Path | None = None,
        revision_rounds: int | None = None,
        persist_members: bool = True,
        **kwargs: Any,
    ):
        symposium_config = self._coerce_config(
            config=config,
            name=name,
            group=group,
            organizer=organizer,
            members=members,
            workspace=workspace,
            revision_rounds=revision_rounds,
        )
        super().__init__(
            llm,
            name=symposium_config.name,
            group=symposium_config.group,
            workspace=symposium_config.workspace or workspace,
            persist_members=persist_members,
            **kwargs,
        )
        self.config = symposium_config
        self.members = {
            member.name: self._build_symposium_participant(member)
            for member in self.config.members
        }
        self.member_configs = {
            member.name: member for member in self.config.members
        }
        self.organizer = self._build_symposium_organizer(self.config.organizer)

    def _events(self, config: Mapping[str, Any] | None) -> EnvironmentEvents:
        return EnvironmentEvents(
            environment=self.name,
            config=config,
            environment_type="agent_symposium",
            environment_id=self.name,
            path=[self.name],
        )

    def _source(self, name: str, *, kind: str = "agent") -> dict[str, Any]:
        return {
            "id": f"{self.name}.{name}",
            "name": name,
            "kind": kind,
            "path": [self.name, name],
        }

    def _environment_source(self) -> dict[str, Any]:
        return {
            "id": self.name,
            "name": self.name,
            "kind": "environment",
            "path": [self.name],
        }

    def _member_node(self, member: EnvironmentMemberConfig) -> dict[str, Any]:
        member_obj = self.members.get(member.name)
        kind = (
            "environment"
            if isinstance(member_obj, BaseEnvironment)
            else "agent"
        )
        return {
            "id": f"{self.name}.{member.name}",
            "name": member.name,
            "kind": kind,
            "role": member.role,
            "agent_class": member.agent,
            "path": [self.name, member.name],
        }

    def _topology_payload(self) -> dict[str, Any]:
        organizer = self.config.organizer
        organizer_kind = (
            "environment"
            if isinstance(self.organizer, BaseEnvironment)
            else "agent"
        )
        nodes = [
            {
                "id": f"{self.name}.{organizer.name}",
                "name": organizer.name,
                "kind": organizer_kind,
                "role": organizer.role,
                "agent_class": organizer.agent,
                "path": [self.name, organizer.name],
            },
            *[self._member_node(member) for member in self.config.members],
        ]
        edges = []
        for member in self.config.members:
            edges.append({
                "source": self.name,
                "target": f"{self.name}.{member.name}",
                "kind": "dispatches_to",
            })
            edges.append({
                "source": f"{self.name}.{member.name}",
                "target": f"{self.name}.{organizer.name}",
                "kind": "synthesized_by",
            })
        return {
            "kind": "agent_symposium",
            "name": self.name,
            "description": self.config.description,
            "revision_rounds": self.config.revision_rounds,
            "nodes": nodes,
            "edges": edges,
        }

    @classmethod
    def from_yaml(
        cls,
        path: str | Path,
        *,
        llm: BaseChatModel,
        **kwargs: Any,
    ) -> "AgentSymposiumEnvironment":
        return cls(llm=llm, config=load_symposium_config(path), **kwargs)

    def _coerce_config(
        self,
        *,
        config: AgentSymposiumConfig | Mapping[str, Any] | str | Path | None,
        name: str | None,
        group: str | None,
        organizer: EnvironmentMemberConfig | Mapping[str, Any] | None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]] | None,
        workspace: str | Path | None,
        revision_rounds: int | None,
    ) -> AgentSymposiumConfig:
        if isinstance(config, (str, Path)):
            base = load_symposium_config(config)
        elif isinstance(config, Mapping):
            base = AgentSymposiumConfig.from_mapping(config)
        elif isinstance(config, AgentSymposiumConfig):
            base = config
        else:
            organizer_cfg = self._coerce_member(
                organizer
                or {
                    "name": "organizer",
                    "role": "Symposium organizer / final synthesizer",
                    "agent": "ChatAgent",
                }
            )
            member_cfgs = [self._coerce_member(m) for m in (members or [])]
            base = AgentSymposiumConfig(
                name=name or "agent_symposium",
                group=group or "default",
                organizer=organizer_cfg,
                members=member_cfgs,
                workspace=str(workspace) if workspace else None,
                revision_rounds=revision_rounds or 1,
            )
        return AgentSymposiumConfig(
            name=name or base.name,
            group=group or base.group,
            description=base.description,
            organizer=base.organizer,
            members=base.members,
            workspace=str(workspace) if workspace else base.workspace,
            defaults=base.defaults,
            revision_rounds=revision_rounds
            if revision_rounds is not None
            else base.revision_rounds,
        )

    @staticmethod
    def _coerce_member(
        member: EnvironmentMemberConfig | Mapping[str, Any],
    ) -> EnvironmentMemberConfig:
        if isinstance(member, EnvironmentMemberConfig):
            return member
        return EnvironmentMemberConfig.from_mapping(member)

    def _build_symposium_participant(
        self, member: EnvironmentMemberConfig
    ) -> Any:
        """Instantiate a symposium participant in its own child workspace.

        Symposium member names are stable URSA agent names. This lets users
        reuse existing named agents directly from YAML instead of creating new
        ``<symposium>_<member>`` checkpoints.
        """
        return self._build_symposium_agent(
            member,
            workspace=self._member_workspace(member.name),
            agent_name=member.name,
        )

    def _build_symposium_organizer(
        self, member: EnvironmentMemberConfig
    ) -> Any:
        """Instantiate the organizer in the parent symposium workspace."""
        return self._build_symposium_agent(
            member,
            workspace=self.workspace,
            agent_name=member.name,
        )

    def _build_symposium_agent(
        self,
        member: EnvironmentMemberConfig,
        *,
        workspace: Path,
        agent_name: str,
    ) -> Any:
        cls = load_object(member.agent)
        llm = make_llm(self.llm, member.model)
        kwargs = dict(member.config or {})
        kwargs.setdefault("workspace", workspace)
        kwargs.setdefault("group", self.group)
        if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
            kwargs.setdefault("name", member.name)
            kwargs.setdefault("persist_members", self.persist_members)
        elif self.persist_members:
            kwargs.setdefault("agent_name", agent_name)
        return cls(llm=llm, **kwargs)

    def _member_workspace_roster(self) -> str:
        return (
            "\n".join(
                f"- {member.name}: {member.name}/ "
                f"({self._member_workspace(member.name)})"
                for member in self.config.members
            )
            or "No symposium member workspaces configured."
        )

    def _member_roster(self) -> str:
        return (
            "\n".join(
                f"- {member.name}: {member.role} ({member.agent})"
                for member in self.config.members
            )
            or "No symposium members configured."
        )

    def _initial_prompt(
        self, member: EnvironmentMemberConfig, task: str
    ) -> str:
        extra = (
            f"\n\nMember-specific guidance:\n{member.prompt}"
            if member.prompt
            else ""
        )
        return (
            f"You are symposium participant '{member.name}' with role: {member.role}.\n"
            "Work independently on the complex problem below. Produce a detailed, "
            "self-contained writeup of your approach, methods, assumptions, files or "
            "commands used, findings, uncertainty, and reproducibility instructions.\n"
            "Your work will be reviewed by other symposium participants. Clean, "
            "compact organization, documentation, and reproducibility are important "
            "parts of how your contribution will be assessed.\n"
            f"{extra}\n\nProblem:\n{task}"
        )

    def _review_prompt(
        self,
        reviewer: EnvironmentMemberConfig,
        task: str,
        writeups: Mapping[str, str],
    ) -> str:
        writeup_block = "\n\n".join(
            f"## Writeup from {name}\n{writeup}"
            for name, writeup in writeups.items()
        )
        return (
            f"You are symposium reviewer '{reviewer.name}' with role: {reviewer.role}.\n"
            "Read all submitted writeups, including your own if present. Assess the "
            "quality of the findings, compare methods and conclusions, and write a "
            "fair but critical review for each writeup. Discuss strengths, weaknesses, "
            "reproducibility, evidence quality, missing checks, and concrete ways to "
            "improve the solution.\n\n"
            "For this review phase, your workspace is the parent symposium "
            "workspace. Participant artifacts are available in these child "
            "workspace directories, using the relative paths shown:\n"
            f"{self._member_workspace_roster()}\n\n"
            "Important review constraints:\n"
            "- Do not change, edit, overwrite, or reorganize any work you are reviewing.\n"
            "- If code or files are referenced, you may inspect or run them only to assess "
            "correctness/reproducibility.\n"
            "- Make clear which findings are well supported and which are speculative.\n"
            "- Your own work will also be judged on clarity, compact organization, "
            "documentation, and reproducibility.\n\n"
            f"Original problem:\n{task}\n\n"
            f"Submitted writeups:\n{writeup_block}"
        )

    def _revision_prompt(
        self,
        member: EnvironmentMemberConfig,
        task: str,
        own_writeup: str,
        all_writeups: Mapping[str, str],
        reviews: Mapping[str, str],
        round_index: int,
    ) -> str:
        other_writeups = "\n\n".join(
            f"## {name}\n{writeup}" for name, writeup in all_writeups.items()
        )
        review_block = "\n\n".join(
            f"## Review from {name}\n{review}"
            for name, review in reviews.items()
        )
        return (
            f"You are symposium participant '{member.name}' revising your own work "
            f"after review round {round_index}.\n"
            "Use the feedback you received and anything you learned from reviewing "
            "other submissions to improve your own solution. You may change only your "
            "own work/artifacts. Do not modify the work of other symposium members.\n"
            "Return a revised detailed writeup with clear improvements, evidence, "
            "limitations, and reproducibility instructions.\n\n"
            f"Original problem:\n{task}\n\n"
            f"Your previous writeup:\n{own_writeup}\n\n"
            f"All writeups you saw:\n{other_writeups}\n\n"
            f"Reviews from symposium members:\n{review_block}"
        )

    def _synthesis_prompt(
        self,
        task: str,
        writeups: Mapping[str, str],
        reviews: Mapping[str, str],
    ) -> str:
        description = (
            f"\nSymposium description: {self.config.description}\n"
            if self.config.description
            else ""
        )
        writeup_block = "\n\n".join(
            f"## Final writeup from {name}\n{writeup}"
            for name, writeup in writeups.items()
        )
        review_block = "\n\n".join(
            f"## Review by {name}\n{review}" for name, review in reviews.items()
        )
        return (
            "You are the organizer of an URSA Agent Symposium. Synthesize the final "
            "results for the user after independent work, peer review, and revision. "
            "Compare participant outputs, identify consensus and disagreement, judge "
            "evidence quality, and provide a final recommendation or solution.\n"
            f"{description}\n"
            "Symposium members:\n"
            f"{self._member_roster()}\n\n"
            "The organizer workspace is the parent symposium workspace. Final "
            "member artifacts are available in these child workspace directories:\n"
            f"{self._member_workspace_roster()}\n\n"
            f"Original problem:\n{task}\n\n"
            f"Final participant writeups:\n{writeup_block}\n\n"
            f"Peer reviews:\n{review_block}"
        )

    def _invoke(
        self, inputs: Mapping[str, Any], **config: Any
    ) -> dict[str, Any]:
        return self._run_ainvoke_from_sync(inputs, **config)

    async def _member_writeup(
        self,
        member_config: EnvironmentMemberConfig,
        prompt: str,
        invoke_kwargs: Mapping[str, Any],
        events: EnvironmentEvents,
        *,
        workspace: Path | None = None,
        event_prefix: str,
        phase_name: str,
        round_index: int | None = None,
    ) -> tuple[str, str]:
        member_name = member_config.name
        source = self._source(member_name)
        start = perf_counter()
        await events.aemit(
            f"{member_name} started {phase_name}",
            stage="symposium",
            phase=phase_name,
            event_type=f"{event_prefix}_started",
            source=source,
            target=self._environment_source(),
            member=member_name,
            round_index=round_index,
            prompt=prompt,
        )
        member = self.members[member_name]
        try:
            result = await self._invoke_member_with_workspace_async(
                member, prompt, workspace=workspace, **invoke_kwargs
            )
        except BaseException as exc:
            await events.aemit(
                f"{member_name} failed {phase_name}",
                stage="symposium",
                phase=phase_name,
                event_type=f"{event_prefix}_failed",
                level="error",
                source=source,
                target=self._environment_source(),
                member=member_name,
                round_index=round_index,
                error=str(exc),
                elapsed_seconds=perf_counter() - start,
            )
            raise
        text = result_to_text(result)
        await events.aemit(
            f"{member_name} completed {phase_name}",
            stage="symposium",
            phase=phase_name,
            event_type=f"{event_prefix}_completed",
            source=source,
            target=self._environment_source(),
            member=member_name,
            round_index=round_index,
            result=text,
            elapsed_seconds=perf_counter() - start,
        )
        return member_name, text

    async def _invoke_member_with_workspace_async(
        self,
        member: Any,
        prompt: str,
        *,
        workspace: Path | None = None,
        **kwargs: Any,
    ) -> Any:
        """Invoke a member, optionally using a temporary workspace.

        Participants normally work in their own child workspaces. During review,
        reviewers are temporarily run from the parent symposium workspace so
        file tools can inspect every participant's child workspace. The original
        workspace is restored after the invocation.
        """
        if workspace is None or not hasattr(member, "workspace"):
            return await self._invoke_member_async(member, prompt, **kwargs)

        original_workspace = getattr(member, "workspace")
        setattr(member, "workspace", Path(workspace))
        Path(workspace).mkdir(parents=True, exist_ok=True)
        try:
            return await self._invoke_member_async(member, prompt, **kwargs)
        finally:
            setattr(member, "workspace", original_workspace)

    async def _ainvoke(
        self, inputs: Mapping[str, Any], **config: Any
    ) -> dict[str, Any]:
        task = str(inputs.get("task") or inputs.get("prompt") or inputs)
        if not self.config.members:
            raise ValueError(
                "AgentSymposiumEnvironment requires at least one member."
            )

        runtime_config = runnable_config_from_kwargs(config)
        events = self._events(runtime_config)
        invoke_kwargs = invocation_kwargs(config)
        start = perf_counter()
        await events.aemit(
            f"Agent symposium {self.name} started",
            stage="symposium",
            phase="start",
            event_type="symposium_started",
            task=task,
            topology=self._topology_payload(),
        )
        await events.aemit(
            f"Agent symposium {self.name} topology declared",
            stage="symposium",
            phase="topology",
            event_type="topology_declared",
            topology=self._topology_payload(),
        )
        try:
            await events.aemit(
                "Initial symposium work started",
                stage="symposium",
                phase="initial_work",
                event_type="symposium_phase_started",
                task=task,
            )
            initial_pairs = await asyncio.gather(*[
                self._member_writeup(
                    member_config,
                    self._initial_prompt(member_config, task),
                    invoke_kwargs,
                    events,
                    event_prefix="initial_work",
                    phase_name="initial_work",
                )
                for member_config in self.config.members
            ])
            initial_writeups = dict(initial_pairs)
            await events.aemit(
                "Initial symposium work completed",
                stage="symposium",
                phase="initial_work",
                event_type="symposium_phase_completed",
                result=initial_writeups,
            )

            current_writeups = dict(initial_writeups)
            review_rounds: list[dict[str, str]] = []
            latest_reviews: dict[str, str] = {}
            reviewer_configs = [
                member for member in self.config.members if member.reviewer
            ]

            for round_index in range(
                1, max(1, self.config.revision_rounds) + 1
            ):
                await events.aemit(
                    f"Review round {round_index} started",
                    stage="symposium",
                    phase="review",
                    event_type="review_round_started",
                    round_index=round_index,
                    writeups=current_writeups,
                )
                review_pairs = await asyncio.gather(*[
                    self._member_writeup(
                        reviewer_config,
                        self._review_prompt(
                            reviewer_config, task, current_writeups
                        ),
                        invoke_kwargs,
                        events,
                        event_prefix="review",
                        phase_name="review",
                        round_index=round_index,
                        workspace=self.workspace,
                    )
                    for reviewer_config in reviewer_configs
                ])
                round_reviews = dict(review_pairs)
                latest_reviews = round_reviews
                review_rounds.append(round_reviews)
                await events.aemit(
                    f"Review round {round_index} completed",
                    stage="symposium",
                    phase="review",
                    event_type="review_round_completed",
                    round_index=round_index,
                    reviews=round_reviews,
                )

                await events.aemit(
                    f"Revision round {round_index} started",
                    stage="symposium",
                    phase="revision",
                    event_type="revision_round_started",
                    round_index=round_index,
                    reviews=round_reviews,
                )
                revision_pairs = await asyncio.gather(*[
                    self._member_writeup(
                        member_config,
                        self._revision_prompt(
                            member_config,
                            task,
                            current_writeups[member_config.name],
                            current_writeups,
                            round_reviews,
                            round_index,
                        ),
                        invoke_kwargs,
                        events,
                        event_prefix="revision",
                        phase_name="revision",
                        round_index=round_index,
                    )
                    for member_config in self.config.members
                ])
                current_writeups = dict(revision_pairs)
                await events.aemit(
                    f"Revision round {round_index} completed",
                    stage="symposium",
                    phase="revision",
                    event_type="revision_round_completed",
                    round_index=round_index,
                    writeups=current_writeups,
                )

            synthesis_prompt = self._synthesis_prompt(
                task, current_writeups, latest_reviews
            )
            synth_start = perf_counter()
            await events.aemit(
                "Symposium synthesis started",
                stage="symposium",
                phase="synthesis",
                event_type="synthesis_started",
                source=self._source(self.config.organizer.name),
                task=task,
                prompt=synthesis_prompt,
            )
            organizer_result = await self._invoke_member_async(
                self.organizer,
                synthesis_prompt,
                **invoke_kwargs,
            )
            final = result_to_text(organizer_result)
            await events.aemit(
                "Symposium synthesis completed",
                stage="symposium",
                phase="synthesis",
                event_type="synthesis_completed",
                source=self._source(self.config.organizer.name),
                result=final,
                elapsed_seconds=perf_counter() - synth_start,
            )
            result = {
                "task": task,
                "initial_writeups": initial_writeups,
                "review_rounds": review_rounds,
                "reviews": latest_reviews,
                "final_writeups": current_writeups,
                "organizer_result": organizer_result,
                "final": final,
            }
        except BaseException as exc:
            await events.aemit(
                f"Agent symposium {self.name} failed",
                stage="symposium",
                phase="error",
                event_type="symposium_failed",
                level="error",
                task=task,
                error=str(exc),
                elapsed_seconds=perf_counter() - start,
            )
            raise
        await events.aemit(
            f"Agent symposium {self.name} completed",
            stage="symposium",
            phase="end",
            event_type="symposium_completed",
            task=task,
            result=result["final"],
            elapsed_seconds=perf_counter() - start,
        )
        return result

AgentTeamConfig dataclass

YAML-loadable configuration for an Agent Team environment.

Source code in src/ursa/environments/config.py
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
@dataclass(frozen=True)
class AgentTeamConfig:
    """YAML-loadable configuration for an Agent Team environment."""

    name: str
    group: str = "default"
    description: str | None = None
    pi: EnvironmentMemberConfig = field(
        default_factory=lambda: EnvironmentMemberConfig(
            name="pi", role="Principal investigator", agent="ExecutionAgent"
        )
    )
    members: list[EnvironmentMemberConfig] = field(default_factory=list)
    workspace: str | None = None
    defaults: dict[str, Any] = field(default_factory=dict)

    @classmethod
    def from_mapping(cls, data: Mapping[str, Any]) -> "AgentTeamConfig":
        raw = dict(data)
        if "pi" in raw and isinstance(raw["pi"], Mapping):
            raw["pi"] = EnvironmentMemberConfig.from_mapping(raw["pi"])
        if "members" in raw:
            raw["members"] = [
                EnvironmentMemberConfig.from_mapping(member)
                for member in raw["members"]
            ]
        return cls(**raw)

AgentTeamEnvironment

Bases: BaseEnvironment

Hierarchical multi-agent team coordinated by a PI agent.

The PI is user-facing. Team members are exposed to the PI as tools, one tool per member, so the PI can plan, delegate, compare returned work, ask follow-up questions, and synthesize a final answer. The environment itself exposes a normal invoke method and can therefore be nested inside other environments such as an Agent Symposium.

Source code in src/ursa/environments/agent_team.py
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
class AgentTeamEnvironment(BaseEnvironment):
    """Hierarchical multi-agent team coordinated by a PI agent.

    The PI is user-facing. Team members are exposed to the PI as tools, one tool
    per member, so the PI can plan, delegate, compare returned work, ask follow-up
    questions, and synthesize a final answer. The environment itself exposes a
    normal ``invoke`` method and can therefore be nested inside other environments
    such as an Agent Symposium.
    """

    def __init__(
        self,
        llm: BaseChatModel,
        *,
        config: AgentTeamConfig | Mapping[str, Any] | str | Path | None = None,
        name: str | None = None,
        group: str | None = None,
        pi: EnvironmentMemberConfig | Mapping[str, Any] | None = None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]]
        | None = None,
        workspace: str | Path | None = None,
        persist_members: bool = True,
        trace_delegation: bool = True,
        trace_character_limit: int = 4000,
        **kwargs: Any,
    ):
        team_config = self._coerce_config(
            config=config,
            name=name,
            group=group,
            pi=pi,
            members=members,
            workspace=workspace,
        )
        super().__init__(
            llm,
            name=team_config.name,
            group=team_config.group,
            workspace=team_config.workspace or workspace,
            persist_members=persist_members,
            **kwargs,
        )
        self.config = team_config
        self.trace_delegation = trace_delegation
        self.trace_character_limit = trace_character_limit
        self.members = {
            member.name: self._build_team_member(member)
            for member in self.config.members
        }
        self.member_configs = {
            member.name: member for member in self.config.members
        }
        self.pi = self._build_pi()

    @classmethod
    def from_yaml(
        cls,
        path: str | Path,
        *,
        llm: BaseChatModel,
        **kwargs: Any,
    ) -> "AgentTeamEnvironment":
        return cls(llm=llm, config=load_team_config(path), **kwargs)

    def _coerce_config(
        self,
        *,
        config: AgentTeamConfig | Mapping[str, Any] | str | Path | None,
        name: str | None,
        group: str | None,
        pi: EnvironmentMemberConfig | Mapping[str, Any] | None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]] | None,
        workspace: str | Path | None,
    ) -> AgentTeamConfig:
        if isinstance(config, (str, Path)):
            base = load_team_config(config)
        elif isinstance(config, Mapping):
            base = AgentTeamConfig.from_mapping(config)
        elif isinstance(config, AgentTeamConfig):
            base = config
        else:
            pi_cfg = self._coerce_member(
                pi
                or {
                    "name": "pi",
                    "role": "Principal investigator / team lead",
                    "agent": "ExecutionAgent",
                }
            )
            member_cfgs = [self._coerce_member(m) for m in (members or [])]
            base = AgentTeamConfig(
                name=name or "agent_team",
                group=group or "default",
                pi=pi_cfg,
                members=member_cfgs,
                workspace=str(workspace) if workspace else None,
            )
        return AgentTeamConfig(
            name=name or base.name,
            group=group or base.group,
            description=base.description,
            pi=base.pi,
            members=base.members,
            workspace=str(workspace) if workspace else base.workspace,
            defaults=base.defaults,
        )

    @staticmethod
    def _coerce_member(
        member: EnvironmentMemberConfig | Mapping[str, Any],
    ) -> EnvironmentMemberConfig:
        if isinstance(member, EnvironmentMemberConfig):
            return member
        return EnvironmentMemberConfig.from_mapping(member)

    def _build_team_member(self, member: EnvironmentMemberConfig) -> Any:
        """Instantiate a team member in the shared team workspace.

        Team member names are stable URSA agent names. This lets users reuse
        existing named agents directly from YAML, e.g. ``team_nif_expert`` rather
        than creating a new ``<team>_<member>`` checkpoint. All team members see
        the same team workspace so they can collaborate through shared files.
        """
        return self._build_named_team_agent(member, agent_name=member.name)

    def _build_team_pi(self, member: EnvironmentMemberConfig) -> Any:
        """Instantiate the PI in the shared team workspace.

        The PI also uses its configured name directly. If that named PI does not
        yet exist, BaseAgent will create its checkpoint under the usual URSA agent
        cache for the configured group.
        """
        return self._build_named_team_agent(member, agent_name=member.name)

    def _build_named_team_agent(
        self,
        member: EnvironmentMemberConfig,
        *,
        agent_name: str,
    ) -> Any:
        cls = load_object(member.agent)
        llm = make_llm(self.llm, member.model)
        kwargs = dict(member.config or {})
        kwargs.setdefault("workspace", self.workspace)
        kwargs.setdefault("group", self.group)
        if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
            kwargs.setdefault("name", member.name)
            kwargs.setdefault("persist_members", self.persist_members)
        elif self.persist_members:
            kwargs.setdefault("agent_name", agent_name)
        return cls(llm=llm, **kwargs)

    def _build_pi(self) -> Any:
        pi_config = self.config.pi
        cls_kwargs = dict(pi_config.config or {})
        delegation_tools = [
            self._make_delegation_tool(member_config)
            for member_config in self.config.members
        ]

        # ExecutionAgent explicitly supports `extra_tools`, making it the preferred
        # PI implementation. Other AgentWithTools-style agents can receive tools
        # after construction via `add_tool` if they expose that method.
        pass_extra_tools = (
            pi_config.agent.endswith("ExecutionAgent")
            or pi_config.agent == "ExecutionAgent"
        )
        if pass_extra_tools:
            cls_kwargs.setdefault("extra_tools", delegation_tools)

        pi_member = EnvironmentMemberConfig(
            name=pi_config.name,
            role=pi_config.role,
            agent=pi_config.agent,
            model=pi_config.model,
            config=cls_kwargs,
            prompt=pi_config.prompt,
        )
        pi_agent = self._build_team_pi(pi_member)
        if (
            not pass_extra_tools
            and delegation_tools
            and hasattr(pi_agent, "add_tool")
        ):
            pi_agent.add_tool(delegation_tools)
        elif not pass_extra_tools and delegation_tools:
            raise TypeError(
                f"Configured PI agent {pi_config.agent!r} cannot accept delegation tools. "
                "Use ExecutionAgent or another AgentWithTools-compatible agent."
            )
        return pi_agent

    def _events(self, config: Mapping[str, Any] | None) -> EnvironmentEvents:
        return EnvironmentEvents(
            environment=self.name,
            config=config,
            environment_type="agent_team",
            environment_id=self.name,
            path=[self.name],
        )

    def _source(self, name: str, *, kind: str = "agent") -> dict[str, Any]:
        return {
            "id": f"{self.name}.{name}",
            "name": name,
            "kind": kind,
            "path": [self.name, name],
        }

    def _member_runtime_config(
        self,
        base_config: Mapping[str, Any] | None,
        member: EnvironmentMemberConfig,
    ) -> dict[str, Any] | None:
        """Attach stable team-member identity to nested agent/tool events."""
        if base_config is None:
            return None
        member_id = f"{self.name}.{member.name}"
        merged = dict(base_config)
        base_metadata = merged.get("metadata")
        metadata = (
            dict(base_metadata) if isinstance(base_metadata, Mapping) else {}
        )
        metadata.update({
            "environment_id": self.name,
            "environment_member": member.name,
            "environment_member_id": member_id,
            "environment_member_role": member.role,
            "environment_member_path": [self.name, member.name],
            "agent": member.name,
            "agent_id": member_id,
        })
        merged["metadata"] = metadata

        base_tags = merged.get("tags")
        if isinstance(base_tags, str):
            tags = [base_tags]
        else:
            tags = list(base_tags) if base_tags else []
        for tag in (member.name, member_id, "environment_member"):
            if tag not in tags:
                tags.append(tag)
        merged["tags"] = tags
        return merged

    def _member_node(self, member: EnvironmentMemberConfig) -> dict[str, Any]:
        kind = (
            "environment"
            if isinstance(self.members.get(member.name), BaseEnvironment)
            else "agent"
        )
        return {
            "id": f"{self.name}.{member.name}",
            "name": member.name,
            "kind": kind,
            "role": member.role,
            "agent_class": member.agent,
            "path": [self.name, member.name],
        }

    def _topology_payload(self) -> dict[str, Any]:
        pi = self.config.pi
        nodes = [
            {
                "id": f"{self.name}.{pi.name}",
                "name": pi.name,
                "kind": "agent",
                "role": pi.role,
                "agent_class": pi.agent,
                "path": [self.name, pi.name],
            },
            *[self._member_node(member) for member in self.config.members],
        ]
        return {
            "kind": "agent_team",
            "name": self.name,
            "description": self.config.description,
            "nodes": nodes,
            "edges": [
                {
                    "source": f"{self.name}.{pi.name}",
                    "target": f"{self.name}.{member.name}",
                    "kind": "delegates_to",
                }
                for member in self.config.members
            ],
        }

    def _make_delegation_tool(
        self, member: EnvironmentMemberConfig
    ) -> StructuredTool:
        member_name = member.name
        role = member.role
        tool_name = f"delegate_to_{_slug_tool_name(member_name)}"

        def delegate(task: str, context: str = "") -> str:
            prompt = self._delegation_prompt(member, task=task, context=context)
            runtime_config = current_environment_config()
            events = self._events(runtime_config)
            source = self._source(self.config.pi.name)
            target = self._source(member_name)
            start = perf_counter()
            events.emit(
                f"Delegating to {member_name}",
                stage="delegation",
                phase="start",
                event_type="delegation_started",
                source=source,
                target=target,
                task=task,
                context=context,
                prompt=prompt,
            )
            self._trace_delegation(
                f"PI -> {member_name}",
                f"Task:\n{task}\n\nContext:\n{context or 'No additional context provided.'}",
            )
            try:
                member_config = self._member_runtime_config(
                    runtime_config, member
                )
                kwargs = {"config": member_config} if member_config else {}
                result = self.members[member_name].invoke(prompt, **kwargs)
            except BaseException as exc:
                events.emit(
                    f"Delegation to {member_name} failed",
                    stage="delegation",
                    phase="error",
                    event_type="delegation_failed",
                    level="error",
                    source=source,
                    target=target,
                    task=task,
                    context=context,
                    error=str(exc),
                    elapsed_seconds=perf_counter() - start,
                )
                raise
            text = result_to_text(result)
            events.emit(
                f"Delegation to {member_name} completed",
                stage="delegation",
                phase="end",
                event_type="delegation_completed",
                source=target,
                target=source,
                task=task,
                context=context,
                result=text,
                elapsed_seconds=perf_counter() - start,
            )
            self._trace_delegation(f"{member_name} -> PI", text)
            return text

        async def adelegate(task: str, context: str = "") -> str:
            prompt = self._delegation_prompt(member, task=task, context=context)
            runtime_config = current_environment_config()
            events = self._events(runtime_config)
            source = self._source(self.config.pi.name)
            target = self._source(member_name)
            start = perf_counter()
            await events.aemit(
                f"Delegating to {member_name}",
                stage="delegation",
                phase="start",
                event_type="delegation_started",
                source=source,
                target=target,
                task=task,
                context=context,
                prompt=prompt,
            )
            self._trace_delegation(
                f"PI -> {member_name}",
                f"Task:\n{task}\n\nContext:\n{context or 'No additional context provided.'}",
            )
            try:
                member_config = self._member_runtime_config(
                    runtime_config, member
                )
                kwargs = {"config": member_config} if member_config else {}
                result = await self._invoke_member_async(
                    self.members[member_name], prompt, **kwargs
                )
            except BaseException as exc:
                await events.aemit(
                    f"Delegation to {member_name} failed",
                    stage="delegation",
                    phase="error",
                    event_type="delegation_failed",
                    level="error",
                    source=source,
                    target=target,
                    task=task,
                    context=context,
                    error=str(exc),
                    elapsed_seconds=perf_counter() - start,
                )
                raise
            text = result_to_text(result)
            await events.aemit(
                f"Delegation to {member_name} completed",
                stage="delegation",
                phase="end",
                event_type="delegation_completed",
                source=target,
                target=source,
                task=task,
                context=context,
                result=text,
                elapsed_seconds=perf_counter() - start,
            )
            self._trace_delegation(f"{member_name} -> PI", text)
            return text

        return StructuredTool.from_function(
            func=delegate,
            coroutine=adelegate,
            name=tool_name,
            description=(
                f"Delegate work to team member '{member_name}' ({role}). "
                "Use this when that member's specialty is relevant. Provide a "
                "self-contained task and any context needed for independent work."
            ),
            args_schema=DelegateInput,
        )

    def _trace_delegation(self, label: str, message: str) -> None:
        """Log a small, explicit delegation trace."""
        if not self.trace_delegation:
            return
        text = message
        if (
            self.trace_character_limit > 0
            and len(text) > self.trace_character_limit
        ):
            text = text[: self.trace_character_limit] + "\n... [truncated]"
        logger.info(f"\n[AgentTeam:{self.name}] {label}\n{text}\n")

    def _delegation_prompt(
        self,
        member: EnvironmentMemberConfig,
        *,
        task: str,
        context: str,
    ) -> str:
        guidance = (
            f"\n\nMember-specific guidance:\n{member.prompt}"
            if member.prompt
            else ""
        )
        return (
            f"You are acting as team member '{member.name}' with role: {member.role}.\n"
            "You have been delegated a task by the team PI. Complete the delegated "
            "task thoroughly, using your available tools when appropriate, and return "
            "a clear writeup of methods, evidence, outputs, limitations, and any files "
            "created or commands needed to reproduce the work.\n"
            f"{guidance}\n\n"
            f"Overall context:\n{context or 'No additional context provided.'}\n\n"
            f"Delegated task:\n{task}"
        )

    def _team_roster(self) -> str:
        if not self.config.members:
            return "No team members are configured. Solve directly as PI."
        return "\n".join(
            f"- {member.name}: {member.role} ({member.agent})"
            for member in self.config.members
        )

    def _pi_prompt(self, task: str) -> str:
        description = (
            f"\nTeam description: {self.config.description}\n"
            if self.config.description
            else ""
        )
        pi_extra = (
            f"\nPI-specific guidance:\n{self.config.pi.prompt}\n"
            if self.config.pi.prompt
            else ""
        )
        return (
            "You are the PI/team leader of a hierarchical Agent Team. "
            "You are the user-facing coordinator responsible for satisfying the "
            "user's overall goal. Formulate an approach, decide which team members "
            "to assign work to, call member-delegation tools as needed, review their "
            "returns critically, request follow-up work if necessary, and synthesize "
            "a final answer for the user.\n"
            "Do not claim a member completed work unless you have delegated it and "
            "reviewed the result. You are personally responsible for final "
            "organization and presentation: integrate the delegated work into one "
            "clean, coherent, easily shareable answer with a clear structure, "
            "actionable conclusions, supporting evidence, limitations, and "
            "reproducibility details where relevant.\n"
            f"{description}{pi_extra}\n"
            "Available team members:\n"
            f"{self._team_roster()}\n\n"
            f"User task:\n{task}"
        )

    def _invoke(self, inputs: Mapping[str, Any], **config: Any) -> Any:
        return self._run_ainvoke_from_sync(inputs, **config)

    async def _ainvoke(self, inputs: Mapping[str, Any], **config: Any) -> Any:
        task = str(inputs.get("task") or inputs.get("prompt") or inputs)
        runtime_config = runnable_config_from_kwargs(config)
        events = self._events(runtime_config)
        token = bind_current_environment_config(runtime_config)
        start = perf_counter()
        await events.aemit(
            f"Agent team {self.name} started",
            stage="team",
            phase="start",
            event_type="team_started",
            task=task,
            topology=self._topology_payload(),
        )
        await events.aemit(
            f"Agent team {self.name} topology declared",
            stage="team",
            phase="topology",
            event_type="topology_declared",
            topology=self._topology_payload(),
        )
        try:
            pi_kwargs = invocation_kwargs(config)
            pi_runtime_config = self._member_runtime_config(
                runtime_config, self.config.pi
            )
            if pi_runtime_config is not None:
                pi_kwargs["config"] = pi_runtime_config
            result = await self._invoke_member_async(
                self.pi, self._pi_prompt(task), **pi_kwargs
            )
        except BaseException as exc:
            await events.aemit(
                f"Agent team {self.name} failed",
                stage="team",
                phase="error",
                event_type="team_failed",
                level="error",
                task=task,
                error=str(exc),
                elapsed_seconds=perf_counter() - start,
            )
            raise
        finally:
            reset_current_environment_config(token)
        await events.aemit(
            f"Agent team {self.name} completed",
            stage="team",
            phase="end",
            event_type="team_completed",
            task=task,
            result=result_to_text(result),
            elapsed_seconds=perf_counter() - start,
        )
        return result

BaseEnvironment

Bases: BaseWorkflow

Base class for multi-agent URSA environments.

Environments compose agents and/or other environments while exposing the same simple invoke surface used by workflows. They deliberately keep persistent configuration separate from agent graph checkpoints: environment definitions live under ~/.cache/ursa/<group>/environments/, while member agents use the shared ~/.cache/ursa/<group>/agents/<agent_name> persistence mechanism.

Source code in src/ursa/environments/base.py
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
class BaseEnvironment(BaseWorkflow):
    """Base class for multi-agent URSA environments.

    Environments compose agents and/or other environments while exposing the same
    simple ``invoke`` surface used by workflows. They deliberately keep persistent
    configuration separate from agent graph checkpoints: environment definitions live
    under ``~/.cache/ursa/<group>/environments/``, while member agents use the
    shared ``~/.cache/ursa/<group>/agents/<agent_name>`` persistence mechanism.
    """

    def __init__(
        self,
        llm: BaseChatModel,
        *,
        name: str,
        group: str = "default",
        workspace: str | Path | None = None,
        persist_members: bool = True,
        **kwargs: Any,
    ):
        super().__init__(**kwargs)
        self.llm = llm
        self.name = name
        self.group = validate_group_name(group)
        self.workspace = Path(
            workspace
            or group_environments_dir(self.group) / "workspaces" / name
        )
        self.workspace.mkdir(parents=True, exist_ok=True)
        self.persist_members = persist_members

    def _normalize_inputs(self, inputs: InputLike) -> Mapping[str, Any]:
        if isinstance(inputs, str):
            return {"task": inputs}
        if isinstance(inputs, Mapping):
            return inputs
        raise TypeError(f"Unsupported input type: {type(inputs)}")

    def _member_workspace(self, member_name: str) -> Path:
        return self.workspace / member_name

    def _member_agent_name(self, member_name: str) -> str | None:
        if not self.persist_members:
            return None
        return f"{self.name}_{member_name}"

    def build_member(self, member: EnvironmentMemberConfig) -> Any:
        """Instantiate a configured member agent or nested environment.

        BaseAgent subclasses use ``agent_name`` for checkpoint persistence, while
        nested environments use ``name`` and manage their own member persistence.
        The distinction keeps teams usable as symposium members without requiring
        environment constructors to accept BaseAgent-specific keywords.
        """
        cls = load_object(member.agent)
        llm = make_llm(self.llm, member.model)
        kwargs = dict(member.config or {})
        kwargs.setdefault("workspace", self._member_workspace(member.name))
        kwargs.setdefault("group", self.group)
        if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
            kwargs.setdefault("name", member.name)
            kwargs.setdefault("persist_members", self.persist_members)
        elif self.persist_members:
            kwargs.setdefault(
                "agent_name", self._member_agent_name(member.name)
            )
        return cls(llm=llm, **kwargs)

    def _run_ainvoke_from_sync(
        self, inputs: Mapping[str, Any], **config: Any
    ) -> Any:
        """Run this environment's async implementation for sync callers.

        Environments are natively async internally so nested agents/tools can use
        async-only implementations. The public ``invoke`` surface remains useful
        for scripts by creating an event loop at the outer boundary. If a loop is
        already running, callers must use ``await environment.ainvoke(...)``.
        """
        try:
            asyncio.get_running_loop()
        except RuntimeError:
            return asyncio.run(self._ainvoke(inputs, **config))

        raise RuntimeError(
            "This environment uses async execution internally, but `.invoke()` "
            "was called from an async context. Use "
            "`await environment.ainvoke(...)` instead."
        )

    async def _invoke_member_async(
        self, member: Any, prompt: str, **kwargs: Any
    ) -> Any:
        """Invoke a member through its async API when available.

        BaseAgent and BaseEnvironment instances expose ``ainvoke``. Lightweight
        test doubles or custom sync-only members may expose only ``invoke``; run
        those in a worker thread so async environment phases remain non-blocking.
        """
        ainvoke = getattr(member, "ainvoke", None)
        if callable(ainvoke):
            result = ainvoke(prompt, **kwargs)
            if inspect.isawaitable(result):
                return await result
            return result

        invoke = getattr(member, "invoke")
        return await asyncio.to_thread(invoke, prompt, **kwargs)

build_member(member)

Instantiate a configured member agent or nested environment.

BaseAgent subclasses use agent_name for checkpoint persistence, while nested environments use name and manage their own member persistence. The distinction keeps teams usable as symposium members without requiring environment constructors to accept BaseAgent-specific keywords.

Source code in src/ursa/environments/base.py
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
def build_member(self, member: EnvironmentMemberConfig) -> Any:
    """Instantiate a configured member agent or nested environment.

    BaseAgent subclasses use ``agent_name`` for checkpoint persistence, while
    nested environments use ``name`` and manage their own member persistence.
    The distinction keeps teams usable as symposium members without requiring
    environment constructors to accept BaseAgent-specific keywords.
    """
    cls = load_object(member.agent)
    llm = make_llm(self.llm, member.model)
    kwargs = dict(member.config or {})
    kwargs.setdefault("workspace", self._member_workspace(member.name))
    kwargs.setdefault("group", self.group)
    if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
        kwargs.setdefault("name", member.name)
        kwargs.setdefault("persist_members", self.persist_members)
    elif self.persist_members:
        kwargs.setdefault(
            "agent_name", self._member_agent_name(member.name)
        )
    return cls(llm=llm, **kwargs)

EnvironmentEventRecorder

Bases: BaseCallbackHandler

Record URSA structured progress events to a replayable JSONL file.

Source code in src/ursa/environments/visualization.py
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
class EnvironmentEventRecorder(BaseCallbackHandler):
    """Record URSA structured progress events to a replayable JSONL file."""

    def __init__(
        self,
        *,
        run_id: str,
        group: str = "default",
        environment_name: str,
        environment_type: str,
        run_dir: Path | None = None,
        event_name: str = DEFAULT_EVENT_NAME,
        max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
    ) -> None:
        self.run_id = run_id
        self.group = validate_group_name(group)
        self.environment_name = environment_name
        self.environment_type = environment_type
        self.event_name = event_name
        self.max_payload_chars = max_payload_chars
        if run_dir is None:
            self.paths = get_environment_run_paths(self.group, run_id)
        else:
            self.paths = EnvironmentRunPaths(
                run_dir=run_dir,
                manifest_path=run_dir / "manifest.json",
                events_path=run_dir / "events.jsonl",
                artifacts_dir=run_dir / "artifacts",
                logs_dir=run_dir / "logs",
            )
        ensure_environment_run_dirs(self.paths)
        self._lock = threading.Lock()
        self._seq = 0

    @property
    def config(self) -> RunnableConfig:
        return {
            "callbacks": [self],
            "metadata": {
                "environment_run_id": self.run_id,
                "environment_name": self.environment_name,
                "environment_type": self.environment_type,
                "group": self.group,
            },
            "tags": ["environment_run", self.environment_name],
        }

    def write_manifest(
        self,
        *,
        status: str,
        task: Any | None = None,
        error: str | None = None,
    ) -> None:
        existing: dict[str, Any] = {}
        if self.paths.manifest_path.exists():
            try:
                existing = json.loads(
                    self.paths.manifest_path.read_text(encoding="utf-8")
                )
            except Exception:
                existing = {}
        now = utc_now_rfc3339()
        manifest = {
            **existing,
            "schema_version": ENVIRONMENT_RUN_SCHEMA_VERSION,
            "run_id": self.run_id,
            "group": self.group,
            "environment_name": self.environment_name,
            "environment_type": self.environment_type,
            "status": status,
            "updated_at": now,
            "events_path": "events.jsonl",
            "artifacts_path": "artifacts",
            "logs_path": "logs",
        }
        manifest.setdefault("created_at", now)
        if task is not None:
            manifest["task_preview"] = _preview(task)
        if error:
            manifest["error"] = error
        self.paths.manifest_path.write_text(
            json.dumps(manifest, indent=2, ensure_ascii=False, default=str),
            encoding="utf-8",
        )

    def on_custom_event(
        self,
        name: str,
        data: Any,
        *,
        run_id,
        tags: list[str] | None = None,
        metadata: dict[str, Any] | None = None,
        **kwargs: Any,
    ) -> None:
        if name != self.event_name or not isinstance(data, Mapping):
            return
        event = self.normalize_event(data, tags=tags, metadata=metadata)
        self.append_event(event)

    def normalize_event(
        self,
        data: Mapping[str, Any],
        *,
        tags: list[str] | None = None,
        metadata: dict[str, Any] | None = None,
    ) -> dict[str, Any]:
        raw_payload_tags = data.get("tags")
        safe_data = _make_json_safe(data, max_chars=self.max_payload_chars)
        if not isinstance(safe_data, dict):
            safe_data = {"value": safe_data}
        event_type = _infer_event_type(safe_data)
        source = _source_from_payload(safe_data)
        target = safe_data.get("target")
        if isinstance(target, str):
            target = {"id": target, "name": target}
        elif isinstance(target, Mapping):
            target = dict(target)
        else:
            target = None
        if (
            target is None
            and safe_data.get("tool")
            and source.get("kind") != "tool"
        ):
            target = _tool_target_from_payload(safe_data)
        environment_name = str(
            safe_data.get("environment") or self.environment_name
        )
        return {
            "schema_version": ENVIRONMENT_EVENT_SCHEMA_VERSION,
            "event_id": new_event_id(),
            "seq": 0,
            "ts": utc_now_rfc3339(),
            "monotonic_timestamp_ns": safe_data.get(
                "monotonic_timestamp_ns", monotonic_ns()
            ),
            "run_id": self.run_id,
            "environment_id": safe_data.get("environment_id")
            or environment_name,
            "environment_name": environment_name,
            "environment_type": safe_data.get("environment_type")
            or self.environment_type,
            "event_type": event_type,
            "stage": safe_data.get("stage"),
            "phase": safe_data.get("phase"),
            "level": safe_data.get("level", "info"),
            "source": source,
            "target": target,
            "message": safe_data.get("message") or event_type,
            "payload": safe_data,
            "tags": _normalize_tags(tags) + _normalize_tags(raw_payload_tags),
            "metadata": _make_json_safe(
                metadata or {}, max_chars=self.max_payload_chars
            ),
        }

    def append_event(self, event: Mapping[str, Any]) -> dict[str, Any]:
        with self._lock:
            self._seq += 1
            record = dict(event)
            record["seq"] = self._seq
            with self.paths.events_path.open("a", encoding="utf-8") as f:
                f.write(
                    json.dumps(record, ensure_ascii=False, default=str) + "\n"
                )
                f.flush()
            return record

EnvironmentMemberConfig dataclass

Configuration for one agent or nested environment member.

YAML fields

name: stable member name used in prompts/tool names/persistence role: human-readable role or specialty agent: Python class path or URSA agent class name, e.g. ExecutionAgent model: optional ModelConfig-compatible mapping for this member config: kwargs passed to the agent/environment constructor prompt: optional extra role/system guidance included in delegated tasks reviewer: whether this member participates in symposium review phases

Source code in src/ursa/environments/config.py
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
@dataclass(frozen=True)
class EnvironmentMemberConfig:
    """Configuration for one agent or nested environment member.

    YAML fields:
      name: stable member name used in prompts/tool names/persistence
      role: human-readable role or specialty
      agent: Python class path or URSA agent class name, e.g. ExecutionAgent
      model: optional ModelConfig-compatible mapping for this member
      config: kwargs passed to the agent/environment constructor
      prompt: optional extra role/system guidance included in delegated tasks
      reviewer: whether this member participates in symposium review phases
    """

    name: str
    role: str = "Team member"
    agent: str = "ExecutionAgent"
    model: ModelConfig | None = None
    config: dict[str, Any] = field(default_factory=dict)
    prompt: str | None = None
    reviewer: bool = True

    @classmethod
    def from_mapping(cls, data: Mapping[str, Any]) -> "EnvironmentMemberConfig":
        raw = dict(data)
        model = raw.get("model")
        if isinstance(model, Mapping):
            raw["model"] = ModelConfig.model_validate(model)
        return cls(**raw)

arun_with_visualization(environment, inputs, *, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS, **kwargs) async

Run an environment asynchronously while recording visualization events.

Source code in src/ursa/environments/visualization.py
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
async def arun_with_visualization(
    environment: Any,
    inputs: InputLike,
    *,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
    **kwargs: Any,
) -> Any:
    """Run an environment asynchronously while recording visualization events."""
    recorder = EnvironmentEventRecorder(
        run_id=run_id or new_run_id(),
        group=getattr(environment, "group", "default"),
        environment_name=getattr(
            environment, "name", type(environment).__name__
        ),
        environment_type=type(environment).__name__,
        max_payload_chars=max_payload_chars,
    )
    recorder.write_manifest(status="running", task=inputs)
    run_config = visualization_config(recorder, config)
    try:
        result = await environment.ainvoke(inputs, config=run_config, **kwargs)
    except BaseException as exc:
        recorder.write_manifest(status="failed", task=inputs, error=str(exc))
        raise
    recorder.write_manifest(status="succeeded", task=inputs)
    return result

environment_run_recorder(environment, *, task=None, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS)

Create a recorder and runnable config for an environment run.

Source code in src/ursa/environments/visualization.py
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
@contextmanager
def environment_run_recorder(
    environment: Any,
    *,
    task: Any | None = None,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
) -> Iterator[tuple[EnvironmentEventRecorder, RunnableConfig]]:
    """Create a recorder and runnable config for an environment run."""
    recorder = EnvironmentEventRecorder(
        run_id=run_id or new_run_id(),
        group=getattr(environment, "group", "default"),
        environment_name=getattr(
            environment, "name", type(environment).__name__
        ),
        environment_type=type(environment).__name__,
        max_payload_chars=max_payload_chars,
    )
    recorder.write_manifest(status="running", task=task)
    try:
        yield recorder, visualization_config(recorder, config)
    except BaseException as exc:
        recorder.write_manifest(status="failed", task=task, error=str(exc))
        raise
    else:
        recorder.write_manifest(status="succeeded", task=task)

environment_runs_dir(group=None)

Return the directory that stores recorded environment visualization runs.

Source code in src/ursa/environments/visualization.py
39
40
41
def environment_runs_dir(group: str | None = None) -> Path:
    """Return the directory that stores recorded environment visualization runs."""
    return group_root_dir(validate_group_name(group)) / "environment_runs"

record_environment_run(environment, *, task=None, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS)

Alias for environment_run_recorder for readable user code.

Source code in src/ursa/environments/visualization.py
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
@contextmanager
def record_environment_run(
    environment: Any,
    *,
    task: Any | None = None,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
) -> Iterator[tuple[EnvironmentEventRecorder, RunnableConfig]]:
    """Alias for ``environment_run_recorder`` for readable user code."""
    with environment_run_recorder(
        environment,
        task=task,
        config=config,
        run_id=run_id,
        max_payload_chars=max_payload_chars,
    ) as value:
        yield value

run_with_visualization(environment, inputs, *, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS, **kwargs)

Run an environment synchronously while recording visualization events.

Source code in src/ursa/environments/visualization.py
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
def run_with_visualization(
    environment: Any,
    inputs: InputLike,
    *,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
    **kwargs: Any,
) -> Any:
    """Run an environment synchronously while recording visualization events."""
    recorder = EnvironmentEventRecorder(
        run_id=run_id or new_run_id(),
        group=getattr(environment, "group", "default"),
        environment_name=getattr(
            environment, "name", type(environment).__name__
        ),
        environment_type=type(environment).__name__,
        max_payload_chars=max_payload_chars,
    )
    recorder.write_manifest(status="running", task=inputs)
    run_config = visualization_config(recorder, config)
    try:
        result = environment.invoke(inputs, config=run_config, **kwargs)
    except BaseException as exc:
        recorder.write_manifest(status="failed", task=inputs, error=str(exc))
        raise
    recorder.write_manifest(status="succeeded", task=inputs)
    return result

save_symposium_config(config, path=None)

Persist a symposium configuration under ~/.cache/ursa//environments by default.

Source code in src/ursa/environments/config.py
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
def save_symposium_config(
    config: AgentSymposiumConfig, path: str | Path | None = None
) -> Path:
    """Persist a symposium configuration under ~/.cache/ursa/<group>/environments by default."""
    target = (
        Path(path).expanduser()
        if path
        else symposium_cache_dir(config.group, config.name) / "symposium.yaml"
    )
    target.parent.mkdir(parents=True, exist_ok=True)
    target.write_text(
        yaml.safe_dump(_dataclass_to_plain(config), sort_keys=False),
        encoding="utf-8",
    )
    return target

save_team_config(config, path=None)

Persist a team configuration under ~/.cache/ursa//environments by default.

Source code in src/ursa/environments/config.py
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
def save_team_config(
    config: AgentTeamConfig, path: str | Path | None = None
) -> Path:
    """Persist a team configuration under ~/.cache/ursa/<group>/environments by default."""
    target = (
        Path(path).expanduser()
        if path
        else team_cache_dir(config.group, config.name) / "team.yaml"
    )
    target.parent.mkdir(parents=True, exist_ok=True)
    target.write_text(
        yaml.safe_dump(_dataclass_to_plain(config), sort_keys=False),
        encoding="utf-8",
    )
    return target

symposium_cache_dir(group, name)

Return the persistent configuration directory for a named symposium.

Source code in src/ursa/environments/config.py
130
131
132
def symposium_cache_dir(group: str, name: str) -> Path:
    """Return the persistent configuration directory for a named symposium."""
    return group_environments_dir(group) / "agent_symposia" / name

team_cache_dir(group, name)

Return the persistent configuration directory for a named team.

Source code in src/ursa/environments/config.py
125
126
127
def team_cache_dir(group: str, name: str) -> Path:
    """Return the persistent configuration directory for a named team."""
    return group_environments_dir(group) / "agent_teams" / name

agent_symposium

AgentSymposiumEnvironment

Bases: BaseEnvironment

Competitive/collaborative multi-agent symposium environment.

A symposium organizer dispatches the same complex problem to several members or nested teams. Members first work independently. The environment then sends all writeups to each member for critical review, explicitly instructing them not to modify reviewed work and to only view/run code as needed for assessment. Finally, each member revises its own work using the feedback it received and insights gained from reviewing others. The organizer synthesizes the final symposium report.

Source code in src/ursa/environments/agent_symposium.py
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
class AgentSymposiumEnvironment(BaseEnvironment):
    """Competitive/collaborative multi-agent symposium environment.

    A symposium organizer dispatches the same complex problem to several members
    or nested teams. Members first work independently. The environment then sends
    all writeups to each member for critical review, explicitly instructing them
    not to modify reviewed work and to only view/run code as needed for assessment.
    Finally, each member revises its own work using the feedback it received and
    insights gained from reviewing others. The organizer synthesizes the final
    symposium report.
    """

    def __init__(
        self,
        llm: BaseChatModel,
        *,
        config: AgentSymposiumConfig
        | Mapping[str, Any]
        | str
        | Path
        | None = None,
        name: str | None = None,
        group: str | None = None,
        organizer: EnvironmentMemberConfig | Mapping[str, Any] | None = None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]]
        | None = None,
        workspace: str | Path | None = None,
        revision_rounds: int | None = None,
        persist_members: bool = True,
        **kwargs: Any,
    ):
        symposium_config = self._coerce_config(
            config=config,
            name=name,
            group=group,
            organizer=organizer,
            members=members,
            workspace=workspace,
            revision_rounds=revision_rounds,
        )
        super().__init__(
            llm,
            name=symposium_config.name,
            group=symposium_config.group,
            workspace=symposium_config.workspace or workspace,
            persist_members=persist_members,
            **kwargs,
        )
        self.config = symposium_config
        self.members = {
            member.name: self._build_symposium_participant(member)
            for member in self.config.members
        }
        self.member_configs = {
            member.name: member for member in self.config.members
        }
        self.organizer = self._build_symposium_organizer(self.config.organizer)

    def _events(self, config: Mapping[str, Any] | None) -> EnvironmentEvents:
        return EnvironmentEvents(
            environment=self.name,
            config=config,
            environment_type="agent_symposium",
            environment_id=self.name,
            path=[self.name],
        )

    def _source(self, name: str, *, kind: str = "agent") -> dict[str, Any]:
        return {
            "id": f"{self.name}.{name}",
            "name": name,
            "kind": kind,
            "path": [self.name, name],
        }

    def _environment_source(self) -> dict[str, Any]:
        return {
            "id": self.name,
            "name": self.name,
            "kind": "environment",
            "path": [self.name],
        }

    def _member_node(self, member: EnvironmentMemberConfig) -> dict[str, Any]:
        member_obj = self.members.get(member.name)
        kind = (
            "environment"
            if isinstance(member_obj, BaseEnvironment)
            else "agent"
        )
        return {
            "id": f"{self.name}.{member.name}",
            "name": member.name,
            "kind": kind,
            "role": member.role,
            "agent_class": member.agent,
            "path": [self.name, member.name],
        }

    def _topology_payload(self) -> dict[str, Any]:
        organizer = self.config.organizer
        organizer_kind = (
            "environment"
            if isinstance(self.organizer, BaseEnvironment)
            else "agent"
        )
        nodes = [
            {
                "id": f"{self.name}.{organizer.name}",
                "name": organizer.name,
                "kind": organizer_kind,
                "role": organizer.role,
                "agent_class": organizer.agent,
                "path": [self.name, organizer.name],
            },
            *[self._member_node(member) for member in self.config.members],
        ]
        edges = []
        for member in self.config.members:
            edges.append({
                "source": self.name,
                "target": f"{self.name}.{member.name}",
                "kind": "dispatches_to",
            })
            edges.append({
                "source": f"{self.name}.{member.name}",
                "target": f"{self.name}.{organizer.name}",
                "kind": "synthesized_by",
            })
        return {
            "kind": "agent_symposium",
            "name": self.name,
            "description": self.config.description,
            "revision_rounds": self.config.revision_rounds,
            "nodes": nodes,
            "edges": edges,
        }

    @classmethod
    def from_yaml(
        cls,
        path: str | Path,
        *,
        llm: BaseChatModel,
        **kwargs: Any,
    ) -> "AgentSymposiumEnvironment":
        return cls(llm=llm, config=load_symposium_config(path), **kwargs)

    def _coerce_config(
        self,
        *,
        config: AgentSymposiumConfig | Mapping[str, Any] | str | Path | None,
        name: str | None,
        group: str | None,
        organizer: EnvironmentMemberConfig | Mapping[str, Any] | None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]] | None,
        workspace: str | Path | None,
        revision_rounds: int | None,
    ) -> AgentSymposiumConfig:
        if isinstance(config, (str, Path)):
            base = load_symposium_config(config)
        elif isinstance(config, Mapping):
            base = AgentSymposiumConfig.from_mapping(config)
        elif isinstance(config, AgentSymposiumConfig):
            base = config
        else:
            organizer_cfg = self._coerce_member(
                organizer
                or {
                    "name": "organizer",
                    "role": "Symposium organizer / final synthesizer",
                    "agent": "ChatAgent",
                }
            )
            member_cfgs = [self._coerce_member(m) for m in (members or [])]
            base = AgentSymposiumConfig(
                name=name or "agent_symposium",
                group=group or "default",
                organizer=organizer_cfg,
                members=member_cfgs,
                workspace=str(workspace) if workspace else None,
                revision_rounds=revision_rounds or 1,
            )
        return AgentSymposiumConfig(
            name=name or base.name,
            group=group or base.group,
            description=base.description,
            organizer=base.organizer,
            members=base.members,
            workspace=str(workspace) if workspace else base.workspace,
            defaults=base.defaults,
            revision_rounds=revision_rounds
            if revision_rounds is not None
            else base.revision_rounds,
        )

    @staticmethod
    def _coerce_member(
        member: EnvironmentMemberConfig | Mapping[str, Any],
    ) -> EnvironmentMemberConfig:
        if isinstance(member, EnvironmentMemberConfig):
            return member
        return EnvironmentMemberConfig.from_mapping(member)

    def _build_symposium_participant(
        self, member: EnvironmentMemberConfig
    ) -> Any:
        """Instantiate a symposium participant in its own child workspace.

        Symposium member names are stable URSA agent names. This lets users
        reuse existing named agents directly from YAML instead of creating new
        ``<symposium>_<member>`` checkpoints.
        """
        return self._build_symposium_agent(
            member,
            workspace=self._member_workspace(member.name),
            agent_name=member.name,
        )

    def _build_symposium_organizer(
        self, member: EnvironmentMemberConfig
    ) -> Any:
        """Instantiate the organizer in the parent symposium workspace."""
        return self._build_symposium_agent(
            member,
            workspace=self.workspace,
            agent_name=member.name,
        )

    def _build_symposium_agent(
        self,
        member: EnvironmentMemberConfig,
        *,
        workspace: Path,
        agent_name: str,
    ) -> Any:
        cls = load_object(member.agent)
        llm = make_llm(self.llm, member.model)
        kwargs = dict(member.config or {})
        kwargs.setdefault("workspace", workspace)
        kwargs.setdefault("group", self.group)
        if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
            kwargs.setdefault("name", member.name)
            kwargs.setdefault("persist_members", self.persist_members)
        elif self.persist_members:
            kwargs.setdefault("agent_name", agent_name)
        return cls(llm=llm, **kwargs)

    def _member_workspace_roster(self) -> str:
        return (
            "\n".join(
                f"- {member.name}: {member.name}/ "
                f"({self._member_workspace(member.name)})"
                for member in self.config.members
            )
            or "No symposium member workspaces configured."
        )

    def _member_roster(self) -> str:
        return (
            "\n".join(
                f"- {member.name}: {member.role} ({member.agent})"
                for member in self.config.members
            )
            or "No symposium members configured."
        )

    def _initial_prompt(
        self, member: EnvironmentMemberConfig, task: str
    ) -> str:
        extra = (
            f"\n\nMember-specific guidance:\n{member.prompt}"
            if member.prompt
            else ""
        )
        return (
            f"You are symposium participant '{member.name}' with role: {member.role}.\n"
            "Work independently on the complex problem below. Produce a detailed, "
            "self-contained writeup of your approach, methods, assumptions, files or "
            "commands used, findings, uncertainty, and reproducibility instructions.\n"
            "Your work will be reviewed by other symposium participants. Clean, "
            "compact organization, documentation, and reproducibility are important "
            "parts of how your contribution will be assessed.\n"
            f"{extra}\n\nProblem:\n{task}"
        )

    def _review_prompt(
        self,
        reviewer: EnvironmentMemberConfig,
        task: str,
        writeups: Mapping[str, str],
    ) -> str:
        writeup_block = "\n\n".join(
            f"## Writeup from {name}\n{writeup}"
            for name, writeup in writeups.items()
        )
        return (
            f"You are symposium reviewer '{reviewer.name}' with role: {reviewer.role}.\n"
            "Read all submitted writeups, including your own if present. Assess the "
            "quality of the findings, compare methods and conclusions, and write a "
            "fair but critical review for each writeup. Discuss strengths, weaknesses, "
            "reproducibility, evidence quality, missing checks, and concrete ways to "
            "improve the solution.\n\n"
            "For this review phase, your workspace is the parent symposium "
            "workspace. Participant artifacts are available in these child "
            "workspace directories, using the relative paths shown:\n"
            f"{self._member_workspace_roster()}\n\n"
            "Important review constraints:\n"
            "- Do not change, edit, overwrite, or reorganize any work you are reviewing.\n"
            "- If code or files are referenced, you may inspect or run them only to assess "
            "correctness/reproducibility.\n"
            "- Make clear which findings are well supported and which are speculative.\n"
            "- Your own work will also be judged on clarity, compact organization, "
            "documentation, and reproducibility.\n\n"
            f"Original problem:\n{task}\n\n"
            f"Submitted writeups:\n{writeup_block}"
        )

    def _revision_prompt(
        self,
        member: EnvironmentMemberConfig,
        task: str,
        own_writeup: str,
        all_writeups: Mapping[str, str],
        reviews: Mapping[str, str],
        round_index: int,
    ) -> str:
        other_writeups = "\n\n".join(
            f"## {name}\n{writeup}" for name, writeup in all_writeups.items()
        )
        review_block = "\n\n".join(
            f"## Review from {name}\n{review}"
            for name, review in reviews.items()
        )
        return (
            f"You are symposium participant '{member.name}' revising your own work "
            f"after review round {round_index}.\n"
            "Use the feedback you received and anything you learned from reviewing "
            "other submissions to improve your own solution. You may change only your "
            "own work/artifacts. Do not modify the work of other symposium members.\n"
            "Return a revised detailed writeup with clear improvements, evidence, "
            "limitations, and reproducibility instructions.\n\n"
            f"Original problem:\n{task}\n\n"
            f"Your previous writeup:\n{own_writeup}\n\n"
            f"All writeups you saw:\n{other_writeups}\n\n"
            f"Reviews from symposium members:\n{review_block}"
        )

    def _synthesis_prompt(
        self,
        task: str,
        writeups: Mapping[str, str],
        reviews: Mapping[str, str],
    ) -> str:
        description = (
            f"\nSymposium description: {self.config.description}\n"
            if self.config.description
            else ""
        )
        writeup_block = "\n\n".join(
            f"## Final writeup from {name}\n{writeup}"
            for name, writeup in writeups.items()
        )
        review_block = "\n\n".join(
            f"## Review by {name}\n{review}" for name, review in reviews.items()
        )
        return (
            "You are the organizer of an URSA Agent Symposium. Synthesize the final "
            "results for the user after independent work, peer review, and revision. "
            "Compare participant outputs, identify consensus and disagreement, judge "
            "evidence quality, and provide a final recommendation or solution.\n"
            f"{description}\n"
            "Symposium members:\n"
            f"{self._member_roster()}\n\n"
            "The organizer workspace is the parent symposium workspace. Final "
            "member artifacts are available in these child workspace directories:\n"
            f"{self._member_workspace_roster()}\n\n"
            f"Original problem:\n{task}\n\n"
            f"Final participant writeups:\n{writeup_block}\n\n"
            f"Peer reviews:\n{review_block}"
        )

    def _invoke(
        self, inputs: Mapping[str, Any], **config: Any
    ) -> dict[str, Any]:
        return self._run_ainvoke_from_sync(inputs, **config)

    async def _member_writeup(
        self,
        member_config: EnvironmentMemberConfig,
        prompt: str,
        invoke_kwargs: Mapping[str, Any],
        events: EnvironmentEvents,
        *,
        workspace: Path | None = None,
        event_prefix: str,
        phase_name: str,
        round_index: int | None = None,
    ) -> tuple[str, str]:
        member_name = member_config.name
        source = self._source(member_name)
        start = perf_counter()
        await events.aemit(
            f"{member_name} started {phase_name}",
            stage="symposium",
            phase=phase_name,
            event_type=f"{event_prefix}_started",
            source=source,
            target=self._environment_source(),
            member=member_name,
            round_index=round_index,
            prompt=prompt,
        )
        member = self.members[member_name]
        try:
            result = await self._invoke_member_with_workspace_async(
                member, prompt, workspace=workspace, **invoke_kwargs
            )
        except BaseException as exc:
            await events.aemit(
                f"{member_name} failed {phase_name}",
                stage="symposium",
                phase=phase_name,
                event_type=f"{event_prefix}_failed",
                level="error",
                source=source,
                target=self._environment_source(),
                member=member_name,
                round_index=round_index,
                error=str(exc),
                elapsed_seconds=perf_counter() - start,
            )
            raise
        text = result_to_text(result)
        await events.aemit(
            f"{member_name} completed {phase_name}",
            stage="symposium",
            phase=phase_name,
            event_type=f"{event_prefix}_completed",
            source=source,
            target=self._environment_source(),
            member=member_name,
            round_index=round_index,
            result=text,
            elapsed_seconds=perf_counter() - start,
        )
        return member_name, text

    async def _invoke_member_with_workspace_async(
        self,
        member: Any,
        prompt: str,
        *,
        workspace: Path | None = None,
        **kwargs: Any,
    ) -> Any:
        """Invoke a member, optionally using a temporary workspace.

        Participants normally work in their own child workspaces. During review,
        reviewers are temporarily run from the parent symposium workspace so
        file tools can inspect every participant's child workspace. The original
        workspace is restored after the invocation.
        """
        if workspace is None or not hasattr(member, "workspace"):
            return await self._invoke_member_async(member, prompt, **kwargs)

        original_workspace = getattr(member, "workspace")
        setattr(member, "workspace", Path(workspace))
        Path(workspace).mkdir(parents=True, exist_ok=True)
        try:
            return await self._invoke_member_async(member, prompt, **kwargs)
        finally:
            setattr(member, "workspace", original_workspace)

    async def _ainvoke(
        self, inputs: Mapping[str, Any], **config: Any
    ) -> dict[str, Any]:
        task = str(inputs.get("task") or inputs.get("prompt") or inputs)
        if not self.config.members:
            raise ValueError(
                "AgentSymposiumEnvironment requires at least one member."
            )

        runtime_config = runnable_config_from_kwargs(config)
        events = self._events(runtime_config)
        invoke_kwargs = invocation_kwargs(config)
        start = perf_counter()
        await events.aemit(
            f"Agent symposium {self.name} started",
            stage="symposium",
            phase="start",
            event_type="symposium_started",
            task=task,
            topology=self._topology_payload(),
        )
        await events.aemit(
            f"Agent symposium {self.name} topology declared",
            stage="symposium",
            phase="topology",
            event_type="topology_declared",
            topology=self._topology_payload(),
        )
        try:
            await events.aemit(
                "Initial symposium work started",
                stage="symposium",
                phase="initial_work",
                event_type="symposium_phase_started",
                task=task,
            )
            initial_pairs = await asyncio.gather(*[
                self._member_writeup(
                    member_config,
                    self._initial_prompt(member_config, task),
                    invoke_kwargs,
                    events,
                    event_prefix="initial_work",
                    phase_name="initial_work",
                )
                for member_config in self.config.members
            ])
            initial_writeups = dict(initial_pairs)
            await events.aemit(
                "Initial symposium work completed",
                stage="symposium",
                phase="initial_work",
                event_type="symposium_phase_completed",
                result=initial_writeups,
            )

            current_writeups = dict(initial_writeups)
            review_rounds: list[dict[str, str]] = []
            latest_reviews: dict[str, str] = {}
            reviewer_configs = [
                member for member in self.config.members if member.reviewer
            ]

            for round_index in range(
                1, max(1, self.config.revision_rounds) + 1
            ):
                await events.aemit(
                    f"Review round {round_index} started",
                    stage="symposium",
                    phase="review",
                    event_type="review_round_started",
                    round_index=round_index,
                    writeups=current_writeups,
                )
                review_pairs = await asyncio.gather(*[
                    self._member_writeup(
                        reviewer_config,
                        self._review_prompt(
                            reviewer_config, task, current_writeups
                        ),
                        invoke_kwargs,
                        events,
                        event_prefix="review",
                        phase_name="review",
                        round_index=round_index,
                        workspace=self.workspace,
                    )
                    for reviewer_config in reviewer_configs
                ])
                round_reviews = dict(review_pairs)
                latest_reviews = round_reviews
                review_rounds.append(round_reviews)
                await events.aemit(
                    f"Review round {round_index} completed",
                    stage="symposium",
                    phase="review",
                    event_type="review_round_completed",
                    round_index=round_index,
                    reviews=round_reviews,
                )

                await events.aemit(
                    f"Revision round {round_index} started",
                    stage="symposium",
                    phase="revision",
                    event_type="revision_round_started",
                    round_index=round_index,
                    reviews=round_reviews,
                )
                revision_pairs = await asyncio.gather(*[
                    self._member_writeup(
                        member_config,
                        self._revision_prompt(
                            member_config,
                            task,
                            current_writeups[member_config.name],
                            current_writeups,
                            round_reviews,
                            round_index,
                        ),
                        invoke_kwargs,
                        events,
                        event_prefix="revision",
                        phase_name="revision",
                        round_index=round_index,
                    )
                    for member_config in self.config.members
                ])
                current_writeups = dict(revision_pairs)
                await events.aemit(
                    f"Revision round {round_index} completed",
                    stage="symposium",
                    phase="revision",
                    event_type="revision_round_completed",
                    round_index=round_index,
                    writeups=current_writeups,
                )

            synthesis_prompt = self._synthesis_prompt(
                task, current_writeups, latest_reviews
            )
            synth_start = perf_counter()
            await events.aemit(
                "Symposium synthesis started",
                stage="symposium",
                phase="synthesis",
                event_type="synthesis_started",
                source=self._source(self.config.organizer.name),
                task=task,
                prompt=synthesis_prompt,
            )
            organizer_result = await self._invoke_member_async(
                self.organizer,
                synthesis_prompt,
                **invoke_kwargs,
            )
            final = result_to_text(organizer_result)
            await events.aemit(
                "Symposium synthesis completed",
                stage="symposium",
                phase="synthesis",
                event_type="synthesis_completed",
                source=self._source(self.config.organizer.name),
                result=final,
                elapsed_seconds=perf_counter() - synth_start,
            )
            result = {
                "task": task,
                "initial_writeups": initial_writeups,
                "review_rounds": review_rounds,
                "reviews": latest_reviews,
                "final_writeups": current_writeups,
                "organizer_result": organizer_result,
                "final": final,
            }
        except BaseException as exc:
            await events.aemit(
                f"Agent symposium {self.name} failed",
                stage="symposium",
                phase="error",
                event_type="symposium_failed",
                level="error",
                task=task,
                error=str(exc),
                elapsed_seconds=perf_counter() - start,
            )
            raise
        await events.aemit(
            f"Agent symposium {self.name} completed",
            stage="symposium",
            phase="end",
            event_type="symposium_completed",
            task=task,
            result=result["final"],
            elapsed_seconds=perf_counter() - start,
        )
        return result

agent_team

AgentTeamEnvironment

Bases: BaseEnvironment

Hierarchical multi-agent team coordinated by a PI agent.

The PI is user-facing. Team members are exposed to the PI as tools, one tool per member, so the PI can plan, delegate, compare returned work, ask follow-up questions, and synthesize a final answer. The environment itself exposes a normal invoke method and can therefore be nested inside other environments such as an Agent Symposium.

Source code in src/ursa/environments/agent_team.py
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
class AgentTeamEnvironment(BaseEnvironment):
    """Hierarchical multi-agent team coordinated by a PI agent.

    The PI is user-facing. Team members are exposed to the PI as tools, one tool
    per member, so the PI can plan, delegate, compare returned work, ask follow-up
    questions, and synthesize a final answer. The environment itself exposes a
    normal ``invoke`` method and can therefore be nested inside other environments
    such as an Agent Symposium.
    """

    def __init__(
        self,
        llm: BaseChatModel,
        *,
        config: AgentTeamConfig | Mapping[str, Any] | str | Path | None = None,
        name: str | None = None,
        group: str | None = None,
        pi: EnvironmentMemberConfig | Mapping[str, Any] | None = None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]]
        | None = None,
        workspace: str | Path | None = None,
        persist_members: bool = True,
        trace_delegation: bool = True,
        trace_character_limit: int = 4000,
        **kwargs: Any,
    ):
        team_config = self._coerce_config(
            config=config,
            name=name,
            group=group,
            pi=pi,
            members=members,
            workspace=workspace,
        )
        super().__init__(
            llm,
            name=team_config.name,
            group=team_config.group,
            workspace=team_config.workspace or workspace,
            persist_members=persist_members,
            **kwargs,
        )
        self.config = team_config
        self.trace_delegation = trace_delegation
        self.trace_character_limit = trace_character_limit
        self.members = {
            member.name: self._build_team_member(member)
            for member in self.config.members
        }
        self.member_configs = {
            member.name: member for member in self.config.members
        }
        self.pi = self._build_pi()

    @classmethod
    def from_yaml(
        cls,
        path: str | Path,
        *,
        llm: BaseChatModel,
        **kwargs: Any,
    ) -> "AgentTeamEnvironment":
        return cls(llm=llm, config=load_team_config(path), **kwargs)

    def _coerce_config(
        self,
        *,
        config: AgentTeamConfig | Mapping[str, Any] | str | Path | None,
        name: str | None,
        group: str | None,
        pi: EnvironmentMemberConfig | Mapping[str, Any] | None,
        members: list[EnvironmentMemberConfig | Mapping[str, Any]] | None,
        workspace: str | Path | None,
    ) -> AgentTeamConfig:
        if isinstance(config, (str, Path)):
            base = load_team_config(config)
        elif isinstance(config, Mapping):
            base = AgentTeamConfig.from_mapping(config)
        elif isinstance(config, AgentTeamConfig):
            base = config
        else:
            pi_cfg = self._coerce_member(
                pi
                or {
                    "name": "pi",
                    "role": "Principal investigator / team lead",
                    "agent": "ExecutionAgent",
                }
            )
            member_cfgs = [self._coerce_member(m) for m in (members or [])]
            base = AgentTeamConfig(
                name=name or "agent_team",
                group=group or "default",
                pi=pi_cfg,
                members=member_cfgs,
                workspace=str(workspace) if workspace else None,
            )
        return AgentTeamConfig(
            name=name or base.name,
            group=group or base.group,
            description=base.description,
            pi=base.pi,
            members=base.members,
            workspace=str(workspace) if workspace else base.workspace,
            defaults=base.defaults,
        )

    @staticmethod
    def _coerce_member(
        member: EnvironmentMemberConfig | Mapping[str, Any],
    ) -> EnvironmentMemberConfig:
        if isinstance(member, EnvironmentMemberConfig):
            return member
        return EnvironmentMemberConfig.from_mapping(member)

    def _build_team_member(self, member: EnvironmentMemberConfig) -> Any:
        """Instantiate a team member in the shared team workspace.

        Team member names are stable URSA agent names. This lets users reuse
        existing named agents directly from YAML, e.g. ``team_nif_expert`` rather
        than creating a new ``<team>_<member>`` checkpoint. All team members see
        the same team workspace so they can collaborate through shared files.
        """
        return self._build_named_team_agent(member, agent_name=member.name)

    def _build_team_pi(self, member: EnvironmentMemberConfig) -> Any:
        """Instantiate the PI in the shared team workspace.

        The PI also uses its configured name directly. If that named PI does not
        yet exist, BaseAgent will create its checkpoint under the usual URSA agent
        cache for the configured group.
        """
        return self._build_named_team_agent(member, agent_name=member.name)

    def _build_named_team_agent(
        self,
        member: EnvironmentMemberConfig,
        *,
        agent_name: str,
    ) -> Any:
        cls = load_object(member.agent)
        llm = make_llm(self.llm, member.model)
        kwargs = dict(member.config or {})
        kwargs.setdefault("workspace", self.workspace)
        kwargs.setdefault("group", self.group)
        if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
            kwargs.setdefault("name", member.name)
            kwargs.setdefault("persist_members", self.persist_members)
        elif self.persist_members:
            kwargs.setdefault("agent_name", agent_name)
        return cls(llm=llm, **kwargs)

    def _build_pi(self) -> Any:
        pi_config = self.config.pi
        cls_kwargs = dict(pi_config.config or {})
        delegation_tools = [
            self._make_delegation_tool(member_config)
            for member_config in self.config.members
        ]

        # ExecutionAgent explicitly supports `extra_tools`, making it the preferred
        # PI implementation. Other AgentWithTools-style agents can receive tools
        # after construction via `add_tool` if they expose that method.
        pass_extra_tools = (
            pi_config.agent.endswith("ExecutionAgent")
            or pi_config.agent == "ExecutionAgent"
        )
        if pass_extra_tools:
            cls_kwargs.setdefault("extra_tools", delegation_tools)

        pi_member = EnvironmentMemberConfig(
            name=pi_config.name,
            role=pi_config.role,
            agent=pi_config.agent,
            model=pi_config.model,
            config=cls_kwargs,
            prompt=pi_config.prompt,
        )
        pi_agent = self._build_team_pi(pi_member)
        if (
            not pass_extra_tools
            and delegation_tools
            and hasattr(pi_agent, "add_tool")
        ):
            pi_agent.add_tool(delegation_tools)
        elif not pass_extra_tools and delegation_tools:
            raise TypeError(
                f"Configured PI agent {pi_config.agent!r} cannot accept delegation tools. "
                "Use ExecutionAgent or another AgentWithTools-compatible agent."
            )
        return pi_agent

    def _events(self, config: Mapping[str, Any] | None) -> EnvironmentEvents:
        return EnvironmentEvents(
            environment=self.name,
            config=config,
            environment_type="agent_team",
            environment_id=self.name,
            path=[self.name],
        )

    def _source(self, name: str, *, kind: str = "agent") -> dict[str, Any]:
        return {
            "id": f"{self.name}.{name}",
            "name": name,
            "kind": kind,
            "path": [self.name, name],
        }

    def _member_runtime_config(
        self,
        base_config: Mapping[str, Any] | None,
        member: EnvironmentMemberConfig,
    ) -> dict[str, Any] | None:
        """Attach stable team-member identity to nested agent/tool events."""
        if base_config is None:
            return None
        member_id = f"{self.name}.{member.name}"
        merged = dict(base_config)
        base_metadata = merged.get("metadata")
        metadata = (
            dict(base_metadata) if isinstance(base_metadata, Mapping) else {}
        )
        metadata.update({
            "environment_id": self.name,
            "environment_member": member.name,
            "environment_member_id": member_id,
            "environment_member_role": member.role,
            "environment_member_path": [self.name, member.name],
            "agent": member.name,
            "agent_id": member_id,
        })
        merged["metadata"] = metadata

        base_tags = merged.get("tags")
        if isinstance(base_tags, str):
            tags = [base_tags]
        else:
            tags = list(base_tags) if base_tags else []
        for tag in (member.name, member_id, "environment_member"):
            if tag not in tags:
                tags.append(tag)
        merged["tags"] = tags
        return merged

    def _member_node(self, member: EnvironmentMemberConfig) -> dict[str, Any]:
        kind = (
            "environment"
            if isinstance(self.members.get(member.name), BaseEnvironment)
            else "agent"
        )
        return {
            "id": f"{self.name}.{member.name}",
            "name": member.name,
            "kind": kind,
            "role": member.role,
            "agent_class": member.agent,
            "path": [self.name, member.name],
        }

    def _topology_payload(self) -> dict[str, Any]:
        pi = self.config.pi
        nodes = [
            {
                "id": f"{self.name}.{pi.name}",
                "name": pi.name,
                "kind": "agent",
                "role": pi.role,
                "agent_class": pi.agent,
                "path": [self.name, pi.name],
            },
            *[self._member_node(member) for member in self.config.members],
        ]
        return {
            "kind": "agent_team",
            "name": self.name,
            "description": self.config.description,
            "nodes": nodes,
            "edges": [
                {
                    "source": f"{self.name}.{pi.name}",
                    "target": f"{self.name}.{member.name}",
                    "kind": "delegates_to",
                }
                for member in self.config.members
            ],
        }

    def _make_delegation_tool(
        self, member: EnvironmentMemberConfig
    ) -> StructuredTool:
        member_name = member.name
        role = member.role
        tool_name = f"delegate_to_{_slug_tool_name(member_name)}"

        def delegate(task: str, context: str = "") -> str:
            prompt = self._delegation_prompt(member, task=task, context=context)
            runtime_config = current_environment_config()
            events = self._events(runtime_config)
            source = self._source(self.config.pi.name)
            target = self._source(member_name)
            start = perf_counter()
            events.emit(
                f"Delegating to {member_name}",
                stage="delegation",
                phase="start",
                event_type="delegation_started",
                source=source,
                target=target,
                task=task,
                context=context,
                prompt=prompt,
            )
            self._trace_delegation(
                f"PI -> {member_name}",
                f"Task:\n{task}\n\nContext:\n{context or 'No additional context provided.'}",
            )
            try:
                member_config = self._member_runtime_config(
                    runtime_config, member
                )
                kwargs = {"config": member_config} if member_config else {}
                result = self.members[member_name].invoke(prompt, **kwargs)
            except BaseException as exc:
                events.emit(
                    f"Delegation to {member_name} failed",
                    stage="delegation",
                    phase="error",
                    event_type="delegation_failed",
                    level="error",
                    source=source,
                    target=target,
                    task=task,
                    context=context,
                    error=str(exc),
                    elapsed_seconds=perf_counter() - start,
                )
                raise
            text = result_to_text(result)
            events.emit(
                f"Delegation to {member_name} completed",
                stage="delegation",
                phase="end",
                event_type="delegation_completed",
                source=target,
                target=source,
                task=task,
                context=context,
                result=text,
                elapsed_seconds=perf_counter() - start,
            )
            self._trace_delegation(f"{member_name} -> PI", text)
            return text

        async def adelegate(task: str, context: str = "") -> str:
            prompt = self._delegation_prompt(member, task=task, context=context)
            runtime_config = current_environment_config()
            events = self._events(runtime_config)
            source = self._source(self.config.pi.name)
            target = self._source(member_name)
            start = perf_counter()
            await events.aemit(
                f"Delegating to {member_name}",
                stage="delegation",
                phase="start",
                event_type="delegation_started",
                source=source,
                target=target,
                task=task,
                context=context,
                prompt=prompt,
            )
            self._trace_delegation(
                f"PI -> {member_name}",
                f"Task:\n{task}\n\nContext:\n{context or 'No additional context provided.'}",
            )
            try:
                member_config = self._member_runtime_config(
                    runtime_config, member
                )
                kwargs = {"config": member_config} if member_config else {}
                result = await self._invoke_member_async(
                    self.members[member_name], prompt, **kwargs
                )
            except BaseException as exc:
                await events.aemit(
                    f"Delegation to {member_name} failed",
                    stage="delegation",
                    phase="error",
                    event_type="delegation_failed",
                    level="error",
                    source=source,
                    target=target,
                    task=task,
                    context=context,
                    error=str(exc),
                    elapsed_seconds=perf_counter() - start,
                )
                raise
            text = result_to_text(result)
            await events.aemit(
                f"Delegation to {member_name} completed",
                stage="delegation",
                phase="end",
                event_type="delegation_completed",
                source=target,
                target=source,
                task=task,
                context=context,
                result=text,
                elapsed_seconds=perf_counter() - start,
            )
            self._trace_delegation(f"{member_name} -> PI", text)
            return text

        return StructuredTool.from_function(
            func=delegate,
            coroutine=adelegate,
            name=tool_name,
            description=(
                f"Delegate work to team member '{member_name}' ({role}). "
                "Use this when that member's specialty is relevant. Provide a "
                "self-contained task and any context needed for independent work."
            ),
            args_schema=DelegateInput,
        )

    def _trace_delegation(self, label: str, message: str) -> None:
        """Log a small, explicit delegation trace."""
        if not self.trace_delegation:
            return
        text = message
        if (
            self.trace_character_limit > 0
            and len(text) > self.trace_character_limit
        ):
            text = text[: self.trace_character_limit] + "\n... [truncated]"
        logger.info(f"\n[AgentTeam:{self.name}] {label}\n{text}\n")

    def _delegation_prompt(
        self,
        member: EnvironmentMemberConfig,
        *,
        task: str,
        context: str,
    ) -> str:
        guidance = (
            f"\n\nMember-specific guidance:\n{member.prompt}"
            if member.prompt
            else ""
        )
        return (
            f"You are acting as team member '{member.name}' with role: {member.role}.\n"
            "You have been delegated a task by the team PI. Complete the delegated "
            "task thoroughly, using your available tools when appropriate, and return "
            "a clear writeup of methods, evidence, outputs, limitations, and any files "
            "created or commands needed to reproduce the work.\n"
            f"{guidance}\n\n"
            f"Overall context:\n{context or 'No additional context provided.'}\n\n"
            f"Delegated task:\n{task}"
        )

    def _team_roster(self) -> str:
        if not self.config.members:
            return "No team members are configured. Solve directly as PI."
        return "\n".join(
            f"- {member.name}: {member.role} ({member.agent})"
            for member in self.config.members
        )

    def _pi_prompt(self, task: str) -> str:
        description = (
            f"\nTeam description: {self.config.description}\n"
            if self.config.description
            else ""
        )
        pi_extra = (
            f"\nPI-specific guidance:\n{self.config.pi.prompt}\n"
            if self.config.pi.prompt
            else ""
        )
        return (
            "You are the PI/team leader of a hierarchical Agent Team. "
            "You are the user-facing coordinator responsible for satisfying the "
            "user's overall goal. Formulate an approach, decide which team members "
            "to assign work to, call member-delegation tools as needed, review their "
            "returns critically, request follow-up work if necessary, and synthesize "
            "a final answer for the user.\n"
            "Do not claim a member completed work unless you have delegated it and "
            "reviewed the result. You are personally responsible for final "
            "organization and presentation: integrate the delegated work into one "
            "clean, coherent, easily shareable answer with a clear structure, "
            "actionable conclusions, supporting evidence, limitations, and "
            "reproducibility details where relevant.\n"
            f"{description}{pi_extra}\n"
            "Available team members:\n"
            f"{self._team_roster()}\n\n"
            f"User task:\n{task}"
        )

    def _invoke(self, inputs: Mapping[str, Any], **config: Any) -> Any:
        return self._run_ainvoke_from_sync(inputs, **config)

    async def _ainvoke(self, inputs: Mapping[str, Any], **config: Any) -> Any:
        task = str(inputs.get("task") or inputs.get("prompt") or inputs)
        runtime_config = runnable_config_from_kwargs(config)
        events = self._events(runtime_config)
        token = bind_current_environment_config(runtime_config)
        start = perf_counter()
        await events.aemit(
            f"Agent team {self.name} started",
            stage="team",
            phase="start",
            event_type="team_started",
            task=task,
            topology=self._topology_payload(),
        )
        await events.aemit(
            f"Agent team {self.name} topology declared",
            stage="team",
            phase="topology",
            event_type="topology_declared",
            topology=self._topology_payload(),
        )
        try:
            pi_kwargs = invocation_kwargs(config)
            pi_runtime_config = self._member_runtime_config(
                runtime_config, self.config.pi
            )
            if pi_runtime_config is not None:
                pi_kwargs["config"] = pi_runtime_config
            result = await self._invoke_member_async(
                self.pi, self._pi_prompt(task), **pi_kwargs
            )
        except BaseException as exc:
            await events.aemit(
                f"Agent team {self.name} failed",
                stage="team",
                phase="error",
                event_type="team_failed",
                level="error",
                task=task,
                error=str(exc),
                elapsed_seconds=perf_counter() - start,
            )
            raise
        finally:
            reset_current_environment_config(token)
        await events.aemit(
            f"Agent team {self.name} completed",
            stage="team",
            phase="end",
            event_type="team_completed",
            task=task,
            result=result_to_text(result),
            elapsed_seconds=perf_counter() - start,
        )
        return result

DelegateInput

Bases: BaseModel

Input schema for PI-to-member delegation tools.

Source code in src/ursa/environments/agent_team.py
35
36
37
38
39
40
41
42
43
44
45
class DelegateInput(BaseModel):
    """Input schema for PI-to-member delegation tools."""

    task: str = Field(
        ...,
        description="Specific task or question to delegate to this team member.",
    )
    context: str = Field(
        "",
        description="Relevant global context, constraints, prior results, or success criteria.",
    )

base

BaseEnvironment

Bases: BaseWorkflow

Base class for multi-agent URSA environments.

Environments compose agents and/or other environments while exposing the same simple invoke surface used by workflows. They deliberately keep persistent configuration separate from agent graph checkpoints: environment definitions live under ~/.cache/ursa/<group>/environments/, while member agents use the shared ~/.cache/ursa/<group>/agents/<agent_name> persistence mechanism.

Source code in src/ursa/environments/base.py
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
class BaseEnvironment(BaseWorkflow):
    """Base class for multi-agent URSA environments.

    Environments compose agents and/or other environments while exposing the same
    simple ``invoke`` surface used by workflows. They deliberately keep persistent
    configuration separate from agent graph checkpoints: environment definitions live
    under ``~/.cache/ursa/<group>/environments/``, while member agents use the
    shared ``~/.cache/ursa/<group>/agents/<agent_name>`` persistence mechanism.
    """

    def __init__(
        self,
        llm: BaseChatModel,
        *,
        name: str,
        group: str = "default",
        workspace: str | Path | None = None,
        persist_members: bool = True,
        **kwargs: Any,
    ):
        super().__init__(**kwargs)
        self.llm = llm
        self.name = name
        self.group = validate_group_name(group)
        self.workspace = Path(
            workspace
            or group_environments_dir(self.group) / "workspaces" / name
        )
        self.workspace.mkdir(parents=True, exist_ok=True)
        self.persist_members = persist_members

    def _normalize_inputs(self, inputs: InputLike) -> Mapping[str, Any]:
        if isinstance(inputs, str):
            return {"task": inputs}
        if isinstance(inputs, Mapping):
            return inputs
        raise TypeError(f"Unsupported input type: {type(inputs)}")

    def _member_workspace(self, member_name: str) -> Path:
        return self.workspace / member_name

    def _member_agent_name(self, member_name: str) -> str | None:
        if not self.persist_members:
            return None
        return f"{self.name}_{member_name}"

    def build_member(self, member: EnvironmentMemberConfig) -> Any:
        """Instantiate a configured member agent or nested environment.

        BaseAgent subclasses use ``agent_name`` for checkpoint persistence, while
        nested environments use ``name`` and manage their own member persistence.
        The distinction keeps teams usable as symposium members without requiring
        environment constructors to accept BaseAgent-specific keywords.
        """
        cls = load_object(member.agent)
        llm = make_llm(self.llm, member.model)
        kwargs = dict(member.config or {})
        kwargs.setdefault("workspace", self._member_workspace(member.name))
        kwargs.setdefault("group", self.group)
        if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
            kwargs.setdefault("name", member.name)
            kwargs.setdefault("persist_members", self.persist_members)
        elif self.persist_members:
            kwargs.setdefault(
                "agent_name", self._member_agent_name(member.name)
            )
        return cls(llm=llm, **kwargs)

    def _run_ainvoke_from_sync(
        self, inputs: Mapping[str, Any], **config: Any
    ) -> Any:
        """Run this environment's async implementation for sync callers.

        Environments are natively async internally so nested agents/tools can use
        async-only implementations. The public ``invoke`` surface remains useful
        for scripts by creating an event loop at the outer boundary. If a loop is
        already running, callers must use ``await environment.ainvoke(...)``.
        """
        try:
            asyncio.get_running_loop()
        except RuntimeError:
            return asyncio.run(self._ainvoke(inputs, **config))

        raise RuntimeError(
            "This environment uses async execution internally, but `.invoke()` "
            "was called from an async context. Use "
            "`await environment.ainvoke(...)` instead."
        )

    async def _invoke_member_async(
        self, member: Any, prompt: str, **kwargs: Any
    ) -> Any:
        """Invoke a member through its async API when available.

        BaseAgent and BaseEnvironment instances expose ``ainvoke``. Lightweight
        test doubles or custom sync-only members may expose only ``invoke``; run
        those in a worker thread so async environment phases remain non-blocking.
        """
        ainvoke = getattr(member, "ainvoke", None)
        if callable(ainvoke):
            result = ainvoke(prompt, **kwargs)
            if inspect.isawaitable(result):
                return await result
            return result

        invoke = getattr(member, "invoke")
        return await asyncio.to_thread(invoke, prompt, **kwargs)

build_member(member)

Instantiate a configured member agent or nested environment.

BaseAgent subclasses use agent_name for checkpoint persistence, while nested environments use name and manage their own member persistence. The distinction keeps teams usable as symposium members without requiring environment constructors to accept BaseAgent-specific keywords.

Source code in src/ursa/environments/base.py
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
def build_member(self, member: EnvironmentMemberConfig) -> Any:
    """Instantiate a configured member agent or nested environment.

    BaseAgent subclasses use ``agent_name`` for checkpoint persistence, while
    nested environments use ``name`` and manage their own member persistence.
    The distinction keeps teams usable as symposium members without requiring
    environment constructors to accept BaseAgent-specific keywords.
    """
    cls = load_object(member.agent)
    llm = make_llm(self.llm, member.model)
    kwargs = dict(member.config or {})
    kwargs.setdefault("workspace", self._member_workspace(member.name))
    kwargs.setdefault("group", self.group)
    if isinstance(cls, type) and issubclass(cls, BaseEnvironment):
        kwargs.setdefault("name", member.name)
        kwargs.setdefault("persist_members", self.persist_members)
    elif self.persist_members:
        kwargs.setdefault(
            "agent_name", self._member_agent_name(member.name)
        )
    return cls(llm=llm, **kwargs)

bind_current_environment_config(config)

Bind runnable config for delegation tools created outside the run scope.

Source code in src/ursa/environments/base.py
152
153
154
155
156
def bind_current_environment_config(
    config: RunnableConfig | None,
) -> Token[RunnableConfig | None]:
    """Bind runnable config for delegation tools created outside the run scope."""
    return _CURRENT_ENVIRONMENT_CONFIG.set(config)

current_environment_config()

Return the runnable config bound to the current environment invocation.

Source code in src/ursa/environments/base.py
147
148
149
def current_environment_config() -> RunnableConfig | None:
    """Return the runnable config bound to the current environment invocation."""
    return _CURRENT_ENVIRONMENT_CONFIG.get()

invocation_kwargs(config)

Drop empty control kwargs before forwarding nested invocations.

Source code in src/ursa/environments/base.py
132
133
134
def invocation_kwargs(config: Mapping[str, Any]) -> dict[str, Any]:
    """Drop empty control kwargs before forwarding nested invocations."""
    return {key: value for key, value in config.items() if value is not None}

result_to_text(result)

Best-effort extraction of human-readable text from an agent result.

Source code in src/ursa/environments/base.py
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
def result_to_text(result: Any) -> str:
    """Best-effort extraction of human-readable text from an agent result."""
    if result is None:
        return ""
    if isinstance(result, str):
        return result
    if isinstance(result, Mapping):
        if "prompt" in result and isinstance(result["prompt"], str):
            return result["prompt"]
        if "final" in result and isinstance(result["final"], str):
            return result["final"]
        if "result" in result and isinstance(result["result"], str):
            return result["result"]
        messages = result.get("messages")
        if isinstance(messages, list) and messages:
            last = messages[-1]
            text = getattr(last, "text", None)
            if text:
                return str(text)
            content = getattr(last, "content", None)
            if content:
                return str(content)
    text = getattr(result, "text", None)
    if text:
        return str(text)
    content = getattr(result, "content", None)
    if content:
        return str(content)
    return str(result)

runnable_config_from_kwargs(config)

Extract the LangChain runnable config from environment kwargs.

Source code in src/ursa/environments/base.py
137
138
139
140
141
142
143
144
def runnable_config_from_kwargs(
    config: Mapping[str, Any],
) -> RunnableConfig | None:
    """Extract the LangChain runnable config from environment kwargs."""
    value = config.get("config")
    if isinstance(value, Mapping):
        return dict(value)
    return None

config

AgentSymposiumConfig dataclass

YAML-loadable configuration for an Agent Symposium environment.

Source code in src/ursa/environments/config.py
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
@dataclass(frozen=True)
class AgentSymposiumConfig:
    """YAML-loadable configuration for an Agent Symposium environment."""

    name: str
    group: str = "default"
    description: str | None = None
    organizer: EnvironmentMemberConfig = field(
        default_factory=lambda: EnvironmentMemberConfig(
            name="organizer", role="Symposium organizer", agent="ChatAgent"
        )
    )
    members: list[EnvironmentMemberConfig] = field(default_factory=list)
    workspace: str | None = None
    defaults: dict[str, Any] = field(default_factory=dict)
    revision_rounds: int = 1

    @classmethod
    def from_mapping(cls, data: Mapping[str, Any]) -> "AgentSymposiumConfig":
        raw = dict(data)
        if "organizer" in raw and isinstance(raw["organizer"], Mapping):
            raw["organizer"] = EnvironmentMemberConfig.from_mapping(
                raw["organizer"]
            )
        if "members" in raw:
            raw["members"] = [
                EnvironmentMemberConfig.from_mapping(member)
                for member in raw["members"]
            ]
        return cls(**raw)

AgentTeamConfig dataclass

YAML-loadable configuration for an Agent Team environment.

Source code in src/ursa/environments/config.py
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
@dataclass(frozen=True)
class AgentTeamConfig:
    """YAML-loadable configuration for an Agent Team environment."""

    name: str
    group: str = "default"
    description: str | None = None
    pi: EnvironmentMemberConfig = field(
        default_factory=lambda: EnvironmentMemberConfig(
            name="pi", role="Principal investigator", agent="ExecutionAgent"
        )
    )
    members: list[EnvironmentMemberConfig] = field(default_factory=list)
    workspace: str | None = None
    defaults: dict[str, Any] = field(default_factory=dict)

    @classmethod
    def from_mapping(cls, data: Mapping[str, Any]) -> "AgentTeamConfig":
        raw = dict(data)
        if "pi" in raw and isinstance(raw["pi"], Mapping):
            raw["pi"] = EnvironmentMemberConfig.from_mapping(raw["pi"])
        if "members" in raw:
            raw["members"] = [
                EnvironmentMemberConfig.from_mapping(member)
                for member in raw["members"]
            ]
        return cls(**raw)

EnvironmentMemberConfig dataclass

Configuration for one agent or nested environment member.

YAML fields

name: stable member name used in prompts/tool names/persistence role: human-readable role or specialty agent: Python class path or URSA agent class name, e.g. ExecutionAgent model: optional ModelConfig-compatible mapping for this member config: kwargs passed to the agent/environment constructor prompt: optional extra role/system guidance included in delegated tasks reviewer: whether this member participates in symposium review phases

Source code in src/ursa/environments/config.py
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
@dataclass(frozen=True)
class EnvironmentMemberConfig:
    """Configuration for one agent or nested environment member.

    YAML fields:
      name: stable member name used in prompts/tool names/persistence
      role: human-readable role or specialty
      agent: Python class path or URSA agent class name, e.g. ExecutionAgent
      model: optional ModelConfig-compatible mapping for this member
      config: kwargs passed to the agent/environment constructor
      prompt: optional extra role/system guidance included in delegated tasks
      reviewer: whether this member participates in symposium review phases
    """

    name: str
    role: str = "Team member"
    agent: str = "ExecutionAgent"
    model: ModelConfig | None = None
    config: dict[str, Any] = field(default_factory=dict)
    prompt: str | None = None
    reviewer: bool = True

    @classmethod
    def from_mapping(cls, data: Mapping[str, Any]) -> "EnvironmentMemberConfig":
        raw = dict(data)
        model = raw.get("model")
        if isinstance(model, Mapping):
            raw["model"] = ModelConfig.model_validate(model)
        return cls(**raw)

load_object(path_or_name)

Load a class by URSA short name or full module path.

Short names are resolved first against ursa.agents and then against ursa.environments. This lets YAML use concise names such as ExecutionAgent or AgentTeamEnvironment while still allowing fully qualified custom class paths.

Source code in src/ursa/environments/config.py
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
def load_object(path_or_name: str) -> Any:
    """Load a class by URSA short name or full module path.

    Short names are resolved first against ``ursa.agents`` and then against
    ``ursa.environments``. This lets YAML use concise names such as
    ``ExecutionAgent`` or ``AgentTeamEnvironment`` while still allowing fully
    qualified custom class paths.
    """
    if "." not in path_or_name:
        for module_name in ("ursa.agents", "ursa.environments"):
            module = importlib.import_module(module_name)
            try:
                return getattr(module, path_or_name)
            except AttributeError:
                continue
        raise AttributeError(
            f"Could not resolve {path_or_name!r} in ursa.agents or "
            "ursa.environments. Use a full Python import path for custom classes."
        )
    module_name, attr_name = path_or_name.rsplit(".", 1)
    module = importlib.import_module(module_name)
    return getattr(module, attr_name)

load_yaml_mapping(path)

Load a YAML mapping with URSA-style environment interpolation.

Source code in src/ursa/environments/config.py
107
108
109
110
111
112
113
114
def load_yaml_mapping(path: str | Path) -> dict[str, Any]:
    """Load a YAML mapping with URSA-style environment interpolation."""
    p = Path(path).expanduser()
    with p.open("r", encoding="utf-8") as handle:
        data = yaml.safe_load(handle) or {}
    if not isinstance(data, Mapping):
        raise ValueError(f"YAML file {p} must contain a top-level mapping.")
    return deep_interp_env(dict(data))

make_llm(default_llm, model_config)

Return a member-specific model if configured, otherwise the default.

Source code in src/ursa/environments/config.py
208
209
210
211
212
213
214
215
216
217
def make_llm(
    default_llm: BaseChatModel,
    model_config: ModelConfig | Mapping[str, Any] | None,
) -> BaseChatModel:
    """Return a member-specific model if configured, otherwise the default."""
    if model_config is None:
        return default_llm
    if isinstance(model_config, Mapping):
        model_config = ModelConfig.model_validate(model_config)
    return init_chat_model(**model_config.kwargs)

save_symposium_config(config, path=None)

Persist a symposium configuration under ~/.cache/ursa//environments by default.

Source code in src/ursa/environments/config.py
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
def save_symposium_config(
    config: AgentSymposiumConfig, path: str | Path | None = None
) -> Path:
    """Persist a symposium configuration under ~/.cache/ursa/<group>/environments by default."""
    target = (
        Path(path).expanduser()
        if path
        else symposium_cache_dir(config.group, config.name) / "symposium.yaml"
    )
    target.parent.mkdir(parents=True, exist_ok=True)
    target.write_text(
        yaml.safe_dump(_dataclass_to_plain(config), sort_keys=False),
        encoding="utf-8",
    )
    return target

save_team_config(config, path=None)

Persist a team configuration under ~/.cache/ursa//environments by default.

Source code in src/ursa/environments/config.py
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
def save_team_config(
    config: AgentTeamConfig, path: str | Path | None = None
) -> Path:
    """Persist a team configuration under ~/.cache/ursa/<group>/environments by default."""
    target = (
        Path(path).expanduser()
        if path
        else team_cache_dir(config.group, config.name) / "team.yaml"
    )
    target.parent.mkdir(parents=True, exist_ok=True)
    target.write_text(
        yaml.safe_dump(_dataclass_to_plain(config), sort_keys=False),
        encoding="utf-8",
    )
    return target

symposium_cache_dir(group, name)

Return the persistent configuration directory for a named symposium.

Source code in src/ursa/environments/config.py
130
131
132
def symposium_cache_dir(group: str, name: str) -> Path:
    """Return the persistent configuration directory for a named symposium."""
    return group_environments_dir(group) / "agent_symposia" / name

team_cache_dir(group, name)

Return the persistent configuration directory for a named team.

Source code in src/ursa/environments/config.py
125
126
127
def team_cache_dir(group: str, name: str) -> Path:
    """Return the persistent configuration directory for a named team."""
    return group_environments_dir(group) / "agent_teams" / name

visualization

EnvironmentEventRecorder

Bases: BaseCallbackHandler

Record URSA structured progress events to a replayable JSONL file.

Source code in src/ursa/environments/visualization.py
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
class EnvironmentEventRecorder(BaseCallbackHandler):
    """Record URSA structured progress events to a replayable JSONL file."""

    def __init__(
        self,
        *,
        run_id: str,
        group: str = "default",
        environment_name: str,
        environment_type: str,
        run_dir: Path | None = None,
        event_name: str = DEFAULT_EVENT_NAME,
        max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
    ) -> None:
        self.run_id = run_id
        self.group = validate_group_name(group)
        self.environment_name = environment_name
        self.environment_type = environment_type
        self.event_name = event_name
        self.max_payload_chars = max_payload_chars
        if run_dir is None:
            self.paths = get_environment_run_paths(self.group, run_id)
        else:
            self.paths = EnvironmentRunPaths(
                run_dir=run_dir,
                manifest_path=run_dir / "manifest.json",
                events_path=run_dir / "events.jsonl",
                artifacts_dir=run_dir / "artifacts",
                logs_dir=run_dir / "logs",
            )
        ensure_environment_run_dirs(self.paths)
        self._lock = threading.Lock()
        self._seq = 0

    @property
    def config(self) -> RunnableConfig:
        return {
            "callbacks": [self],
            "metadata": {
                "environment_run_id": self.run_id,
                "environment_name": self.environment_name,
                "environment_type": self.environment_type,
                "group": self.group,
            },
            "tags": ["environment_run", self.environment_name],
        }

    def write_manifest(
        self,
        *,
        status: str,
        task: Any | None = None,
        error: str | None = None,
    ) -> None:
        existing: dict[str, Any] = {}
        if self.paths.manifest_path.exists():
            try:
                existing = json.loads(
                    self.paths.manifest_path.read_text(encoding="utf-8")
                )
            except Exception:
                existing = {}
        now = utc_now_rfc3339()
        manifest = {
            **existing,
            "schema_version": ENVIRONMENT_RUN_SCHEMA_VERSION,
            "run_id": self.run_id,
            "group": self.group,
            "environment_name": self.environment_name,
            "environment_type": self.environment_type,
            "status": status,
            "updated_at": now,
            "events_path": "events.jsonl",
            "artifacts_path": "artifacts",
            "logs_path": "logs",
        }
        manifest.setdefault("created_at", now)
        if task is not None:
            manifest["task_preview"] = _preview(task)
        if error:
            manifest["error"] = error
        self.paths.manifest_path.write_text(
            json.dumps(manifest, indent=2, ensure_ascii=False, default=str),
            encoding="utf-8",
        )

    def on_custom_event(
        self,
        name: str,
        data: Any,
        *,
        run_id,
        tags: list[str] | None = None,
        metadata: dict[str, Any] | None = None,
        **kwargs: Any,
    ) -> None:
        if name != self.event_name or not isinstance(data, Mapping):
            return
        event = self.normalize_event(data, tags=tags, metadata=metadata)
        self.append_event(event)

    def normalize_event(
        self,
        data: Mapping[str, Any],
        *,
        tags: list[str] | None = None,
        metadata: dict[str, Any] | None = None,
    ) -> dict[str, Any]:
        raw_payload_tags = data.get("tags")
        safe_data = _make_json_safe(data, max_chars=self.max_payload_chars)
        if not isinstance(safe_data, dict):
            safe_data = {"value": safe_data}
        event_type = _infer_event_type(safe_data)
        source = _source_from_payload(safe_data)
        target = safe_data.get("target")
        if isinstance(target, str):
            target = {"id": target, "name": target}
        elif isinstance(target, Mapping):
            target = dict(target)
        else:
            target = None
        if (
            target is None
            and safe_data.get("tool")
            and source.get("kind") != "tool"
        ):
            target = _tool_target_from_payload(safe_data)
        environment_name = str(
            safe_data.get("environment") or self.environment_name
        )
        return {
            "schema_version": ENVIRONMENT_EVENT_SCHEMA_VERSION,
            "event_id": new_event_id(),
            "seq": 0,
            "ts": utc_now_rfc3339(),
            "monotonic_timestamp_ns": safe_data.get(
                "monotonic_timestamp_ns", monotonic_ns()
            ),
            "run_id": self.run_id,
            "environment_id": safe_data.get("environment_id")
            or environment_name,
            "environment_name": environment_name,
            "environment_type": safe_data.get("environment_type")
            or self.environment_type,
            "event_type": event_type,
            "stage": safe_data.get("stage"),
            "phase": safe_data.get("phase"),
            "level": safe_data.get("level", "info"),
            "source": source,
            "target": target,
            "message": safe_data.get("message") or event_type,
            "payload": safe_data,
            "tags": _normalize_tags(tags) + _normalize_tags(raw_payload_tags),
            "metadata": _make_json_safe(
                metadata or {}, max_chars=self.max_payload_chars
            ),
        }

    def append_event(self, event: Mapping[str, Any]) -> dict[str, Any]:
        with self._lock:
            self._seq += 1
            record = dict(event)
            record["seq"] = self._seq
            with self.paths.events_path.open("a", encoding="utf-8") as f:
                f.write(
                    json.dumps(record, ensure_ascii=False, default=str) + "\n"
                )
                f.flush()
            return record

arun_with_visualization(environment, inputs, *, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS, **kwargs) async

Run an environment asynchronously while recording visualization events.

Source code in src/ursa/environments/visualization.py
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
async def arun_with_visualization(
    environment: Any,
    inputs: InputLike,
    *,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
    **kwargs: Any,
) -> Any:
    """Run an environment asynchronously while recording visualization events."""
    recorder = EnvironmentEventRecorder(
        run_id=run_id or new_run_id(),
        group=getattr(environment, "group", "default"),
        environment_name=getattr(
            environment, "name", type(environment).__name__
        ),
        environment_type=type(environment).__name__,
        max_payload_chars=max_payload_chars,
    )
    recorder.write_manifest(status="running", task=inputs)
    run_config = visualization_config(recorder, config)
    try:
        result = await environment.ainvoke(inputs, config=run_config, **kwargs)
    except BaseException as exc:
        recorder.write_manifest(status="failed", task=inputs, error=str(exc))
        raise
    recorder.write_manifest(status="succeeded", task=inputs)
    return result

environment_run_recorder(environment, *, task=None, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS)

Create a recorder and runnable config for an environment run.

Source code in src/ursa/environments/visualization.py
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
@contextmanager
def environment_run_recorder(
    environment: Any,
    *,
    task: Any | None = None,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
) -> Iterator[tuple[EnvironmentEventRecorder, RunnableConfig]]:
    """Create a recorder and runnable config for an environment run."""
    recorder = EnvironmentEventRecorder(
        run_id=run_id or new_run_id(),
        group=getattr(environment, "group", "default"),
        environment_name=getattr(
            environment, "name", type(environment).__name__
        ),
        environment_type=type(environment).__name__,
        max_payload_chars=max_payload_chars,
    )
    recorder.write_manifest(status="running", task=task)
    try:
        yield recorder, visualization_config(recorder, config)
    except BaseException as exc:
        recorder.write_manifest(status="failed", task=task, error=str(exc))
        raise
    else:
        recorder.write_manifest(status="succeeded", task=task)

environment_runs_dir(group=None)

Return the directory that stores recorded environment visualization runs.

Source code in src/ursa/environments/visualization.py
39
40
41
def environment_runs_dir(group: str | None = None) -> Path:
    """Return the directory that stores recorded environment visualization runs."""
    return group_root_dir(validate_group_name(group)) / "environment_runs"

record_environment_run(environment, *, task=None, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS)

Alias for environment_run_recorder for readable user code.

Source code in src/ursa/environments/visualization.py
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
@contextmanager
def record_environment_run(
    environment: Any,
    *,
    task: Any | None = None,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
) -> Iterator[tuple[EnvironmentEventRecorder, RunnableConfig]]:
    """Alias for ``environment_run_recorder`` for readable user code."""
    with environment_run_recorder(
        environment,
        task=task,
        config=config,
        run_id=run_id,
        max_payload_chars=max_payload_chars,
    ) as value:
        yield value

run_with_visualization(environment, inputs, *, config=None, run_id=None, max_payload_chars=DEFAULT_MAX_PAYLOAD_CHARS, **kwargs)

Run an environment synchronously while recording visualization events.

Source code in src/ursa/environments/visualization.py
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
def run_with_visualization(
    environment: Any,
    inputs: InputLike,
    *,
    config: RunnableConfig | None = None,
    run_id: str | None = None,
    max_payload_chars: int = DEFAULT_MAX_PAYLOAD_CHARS,
    **kwargs: Any,
) -> Any:
    """Run an environment synchronously while recording visualization events."""
    recorder = EnvironmentEventRecorder(
        run_id=run_id or new_run_id(),
        group=getattr(environment, "group", "default"),
        environment_name=getattr(
            environment, "name", type(environment).__name__
        ),
        environment_type=type(environment).__name__,
        max_payload_chars=max_payload_chars,
    )
    recorder.write_manifest(status="running", task=inputs)
    run_config = visualization_config(recorder, config)
    try:
        result = environment.invoke(inputs, config=run_config, **kwargs)
    except BaseException as exc:
        recorder.write_manifest(status="failed", task=inputs, error=str(exc))
        raise
    recorder.write_manifest(status="succeeded", task=inputs)
    return result

visualization_config(recorder, config=None)

Merge an optional runnable config with the recorder callback config.

Source code in src/ursa/environments/visualization.py
414
415
416
417
418
419
420
421
def visualization_config(
    recorder: EnvironmentEventRecorder,
    config: RunnableConfig | None = None,
) -> RunnableConfig:
    """Merge an optional runnable config with the recorder callback config."""
    if config is None:
        return recorder.config
    return merge_configs(recorder.config, config)