mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-08-17 18:25:42 +02:00
[eric] events: check sessions carry pre-authorized MCPs, workflow approvals, and a live dashboard
This commit is contained in:
@@ -14,6 +14,7 @@ from typing import Dict, List, Optional, Tuple
|
||||
from typeguard import typechecked
|
||||
|
||||
from backend.apps.events.models import AgentCheckSource, Event
|
||||
from backend.apps.workflows.models import Workflow
|
||||
|
||||
CHECK_TIMEOUT_S = 240.0
|
||||
|
||||
@@ -83,17 +84,50 @@ async def p_await_reply(session_id: str) -> str:
|
||||
raise RuntimeError("check agent produced no reply")
|
||||
|
||||
|
||||
async def run_check_turn(model: str, prompt: str) -> str:
|
||||
"""One ephemeral agent turn; the session file is deleted afterward."""
|
||||
async def run_check_turn(
|
||||
model: str,
|
||||
prompt: str,
|
||||
dashboard_id: Optional[str] = None,
|
||||
active_mcps: Optional[List[str]] = None,
|
||||
approvals: Optional[Dict[str, str]] = None,
|
||||
) -> str:
|
||||
"""One ephemeral agent turn; the session file is deleted afterward.
|
||||
|
||||
dashboard_id makes browser delegation available (browser cards render on a
|
||||
dashboard, so a logged-in-site check works whenever the app is open).
|
||||
active_mcps presets session.active_mcps with what the USER pre-authorized
|
||||
on the trigger; the dispatch gate itself still only spawns what that list
|
||||
names. approvals replays the workflow's remembered tool answers so an
|
||||
unattended check doesn't park on a prompt nobody will answer."""
|
||||
from backend.apps.agents.agent_manager import agent_manager
|
||||
from backend.apps.agents.core.models import AgentConfig
|
||||
from backend.apps.agents.manager.permissions.workflow_approval import (
|
||||
clear_workflow_approval_memory,
|
||||
set_workflow_approval_memory,
|
||||
)
|
||||
from backend.apps.agents.manager.session.session_store import delete_session_file
|
||||
|
||||
session = await agent_manager.launch_agent(AgentConfig(name="Event check", model=model, mode="agent"))
|
||||
session = await agent_manager.launch_agent(AgentConfig(
|
||||
name="Event check", model=model, mode="agent", dashboard_id=dashboard_id,
|
||||
))
|
||||
if active_mcps:
|
||||
session.active_mcps = list(active_mcps)
|
||||
if approvals:
|
||||
set_workflow_approval_memory(
|
||||
session.id,
|
||||
decisions=dict(approvals),
|
||||
step_usage={},
|
||||
remember=lambda tool_name, behavior: None,
|
||||
ask_timeout=30.0,
|
||||
)
|
||||
try:
|
||||
await agent_manager.send_message(session.id, prompt)
|
||||
return await p_await_reply(session.id)
|
||||
finally:
|
||||
try:
|
||||
clear_workflow_approval_memory(session.id)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
await agent_manager.close_session(session.id)
|
||||
except Exception:
|
||||
@@ -105,16 +139,23 @@ async def run_check_turn(model: str, prompt: str) -> str:
|
||||
|
||||
|
||||
@typechecked
|
||||
async def agent_check(source: AgentCheckSource, cursor: Dict) -> Tuple[List[Event], Dict]:
|
||||
async def agent_check(source: AgentCheckSource, cursor: Dict, workflow: Optional[Workflow] = None) -> Tuple[List[Event], Dict]:
|
||||
check = source.check.strip()
|
||||
if not check:
|
||||
return [], cursor
|
||||
from backend.apps.settings.settings import load_settings
|
||||
from backend.apps.workflows.executor import resolve_workflow_dashboard_id
|
||||
|
||||
baselined = bool(cursor.get("baselined")) and cursor.get("check") == check
|
||||
prev_state = str(cursor.get("state") or "") if baselined else ""
|
||||
model = source.model.strip() or (getattr(load_settings(), "default_model", None) or "sonnet")
|
||||
reply = await run_check_turn(model, build_check_prompt(check, prev_state))
|
||||
model = source.model.strip() or (workflow.model if workflow else "") or (getattr(load_settings(), "default_model", None) or "sonnet")
|
||||
reply = await run_check_turn(
|
||||
model,
|
||||
build_check_prompt(check, prev_state),
|
||||
dashboard_id=resolve_workflow_dashboard_id(workflow) if workflow else None,
|
||||
active_mcps=list(source.mcps),
|
||||
approvals=dict(workflow.remembered_approvals) if workflow else None,
|
||||
)
|
||||
event_line, state = parse_check_reply(reply)
|
||||
new_cursor: Dict = {"baselined": True, "check": check, "state": state or prev_state}
|
||||
if cursor.get("last_event_digest"):
|
||||
|
||||
@@ -61,8 +61,15 @@ class AgentCheckSource(BaseModel):
|
||||
check: str = ""
|
||||
# Empty = the app's default model; each poll is a real (short) agent turn.
|
||||
model: str = ""
|
||||
# MCPs the USER pre-authorized for this check at trigger creation (consent moved from the in-session MCPActivate click to the trigger config; the dispatch gate itself is unchanged).
|
||||
mcps: list[str] = Field(default_factory=list)
|
||||
poll_seconds: int = 900
|
||||
|
||||
@field_validator("mcps")
|
||||
@classmethod
|
||||
def p_cap_mcps(cls, v: list[str]) -> list[str]:
|
||||
return [str(m).strip() for m in v if str(m).strip()][:8]
|
||||
|
||||
@field_validator("poll_seconds")
|
||||
@classmethod
|
||||
def p_clamp_poll(cls, v: int) -> int:
|
||||
|
||||
@@ -62,11 +62,16 @@ def p_live_triggers() -> List[Tuple[Workflow, EventTriggerConfig, float]]:
|
||||
return out
|
||||
|
||||
|
||||
async def p_poll_one(workflow_id: str, trigger: EventTriggerConfig) -> None:
|
||||
async def p_poll_one(wf: Workflow, trigger: EventTriggerConfig) -> None:
|
||||
workflow_id = wf.id
|
||||
try:
|
||||
fetch = ADAPTERS[trigger.source.kind]
|
||||
cursor = stores.load_cursor(trigger.id)
|
||||
events, new_cursor = await fetch(trigger.source, cursor)
|
||||
# The agent adapter needs its parent workflow (model, approvals, dashboard); the structural adapters stay pure.
|
||||
if trigger.source.kind == "agent":
|
||||
events, new_cursor = await fetch(trigger.source, cursor, wf)
|
||||
else:
|
||||
events, new_cursor = await fetch(trigger.source, cursor)
|
||||
stores.save_cursor(trigger.id, new_cursor)
|
||||
if events:
|
||||
await dispatcher.ingest(workflow_id, trigger, events)
|
||||
@@ -94,7 +99,7 @@ def tick() -> None:
|
||||
if p_next_poll.get(trig.id, 0.0) <= now and trig.id not in p_inflight:
|
||||
p_next_poll[trig.id] = now + poll_seconds
|
||||
p_inflight.add(trig.id)
|
||||
asyncio.create_task(p_poll_one(wf.id, trig))
|
||||
asyncio.create_task(p_poll_one(wf, trig))
|
||||
|
||||
|
||||
def p_seconds_until_next() -> float:
|
||||
|
||||
@@ -65,7 +65,7 @@ def test_agent_check_baseline_then_event_then_dedup(monkeypatch):
|
||||
])
|
||||
prompts: list[str] = []
|
||||
|
||||
async def p_fake_turn(model, prompt):
|
||||
async def p_fake_turn(model, prompt, dashboard_id=None, active_mcps=None, approvals=None):
|
||||
prompts.append(prompt)
|
||||
return next(replies)
|
||||
|
||||
@@ -84,6 +84,32 @@ def test_agent_check_baseline_then_event_then_dedup(monkeypatch):
|
||||
assert events == [] # identical event line reported again fires once, not forever
|
||||
|
||||
|
||||
def test_agent_check_carries_workflow_context(make_wf, monkeypatch):
|
||||
"""Pre-authorized MCPs, the workflow's remembered approvals, and the model
|
||||
all flow into the check turn; consent lives on the trigger config, never
|
||||
widened inside the session."""
|
||||
from backend.apps.events.adapters import agent_check as ac
|
||||
from backend.apps.workflows import executor
|
||||
|
||||
seen: dict = {}
|
||||
|
||||
async def p_fake_turn(model, prompt, dashboard_id=None, active_mcps=None, approvals=None):
|
||||
seen.update(model=model, dashboard_id=dashboard_id, active_mcps=active_mcps, approvals=approvals)
|
||||
return "NO_EVENT\nSTATE: s"
|
||||
|
||||
monkeypatch.setattr(ac, "run_check_turn", p_fake_turn)
|
||||
monkeypatch.setattr(executor, "resolve_workflow_dashboard_id", lambda wf: "dash-42")
|
||||
|
||||
source = AgentCheckSource(check="new invoice email arrived", mcps=["google-workspace"], poll_seconds=300)
|
||||
wf = make_wf(model="opus", remembered_approvals={"SendEmail": "deny"})
|
||||
p_run(ac.agent_check(source, {}, wf))
|
||||
|
||||
assert seen["model"] == "opus" # workflow's model, since the source didn't pin one
|
||||
assert seen["dashboard_id"] == "dash-42"
|
||||
assert seen["active_mcps"] == ["google-workspace"]
|
||||
assert seen["approvals"] == {"SendEmail": "deny"}
|
||||
|
||||
|
||||
def p_custom_wf(make_wf):
|
||||
from backend.apps.workflows import storage
|
||||
|
||||
|
||||
Reference in New Issue
Block a user