diff --git a/backend/apps/agents/agent_manager.py b/backend/apps/agents/agent_manager.py index ee3f8999..3eb96bca 100644 --- a/backend/apps/agents/agent_manager.py +++ b/backend/apps/agents/agent_manager.py @@ -7,7 +7,7 @@ import sys import time from datetime import datetime from uuid import uuid4 -from typing import Optional +from typing import Callable, Optional from backend.apps.agents.core.models import ( AgentConfig, AgentSession, Message, MessageBranch, ApprovalRequest, ToolGroupMeta, @@ -79,6 +79,59 @@ logger = logging.getLogger(__name__) os.environ.setdefault("CLAUDE_CODE_STREAM_CLOSE_TIMEOUT", "3600000") +# Per-session approval memory for workflow runs. The workflow executor pushes +# context in (keyed by session id) so the permission gates can reuse a prior +# allow/deny instead of prompting, and so an unattended fire fails fast instead +# of parking for ten minutes. Lives here (not in the workflows app) because +# agent_manager must not import workflows; that would be an upward import cycle. +class WorkflowApprovalMemory: + def __init__( + self, + decisions: dict[str, str], + step_usage: dict[str, dict[str, bool]], + remember: Optional[Callable[[str, str], None]], + ask_timeout: float, + ): + self.decisions = decisions # workflow-level: tool -> "allow"/"deny" + self.step_usage = step_usage # per-step record: step_id -> {tool: approved} + self.remember = remember # persist a workflow-level decision to disk + self.ask_timeout = ask_timeout + # The executor bumps this as it advances steps so the gate can record + # which tools each step touched. None on test runs that don't thread it. + self.current_step_id: Optional[str] = None + + +p_approval_memory: dict[str, WorkflowApprovalMemory] = {} + + +def set_workflow_approval_memory( + session_id: str, + *, + decisions: dict[str, str], + step_usage: dict[str, dict[str, bool]], + remember: Optional[Callable[[str, str], None]], + ask_timeout: float, +) -> None: + p_approval_memory[session_id] = WorkflowApprovalMemory( + decisions, step_usage, remember, ask_timeout + ) + + +def clear_workflow_approval_memory(session_id: str) -> None: + p_approval_memory.pop(session_id, None) + + +def set_workflow_approval_step(session_id: str, step_id: Optional[str]) -> None: + mem = p_approval_memory.get(session_id) + if mem is not None: + mem.current_step_id = step_id + + +def get_workflow_step_usage(session_id: str) -> dict[str, dict[str, bool]]: + mem = p_approval_memory.get(session_id) + return mem.step_usage if mem is not None else {} + + def _apply_context_window(session, settings=None) -> None: """Set session.context_window from the registry for its (provider, model). @@ -723,6 +776,7 @@ class AgentManager: tool_name: str, tool_input, sensitive_pattern: str | None = None, + timeout: float = 600.0, ) -> dict: """Send an approval request via WebSocket and wait for the user's decision.""" safe_input = tool_input if isinstance(tool_input, dict) else {} @@ -753,6 +807,7 @@ class AgentManager: decision = await ws_manager.send_approval_request( session_id, request_id, tool_name, safe_input, + timeout=timeout, sensitive_pattern=sensitive_pattern, sensitive_label=label, sensitive_why=why, @@ -782,6 +837,7 @@ class AgentManager: "tool": tool_name, "behavior": decision.get("behavior"), "decision_ms": approval_latency_ms, + "sensitive_pattern": sensitive_pattern, }) except Exception: pass @@ -796,6 +852,58 @@ class AgentManager: }) return decision + def p_note_tool_used(tool_name: str, approved: bool) -> None: + # Record which tools each step touched (in-memory; the executor/test + # path persists step_usage once at run end). Captures every tool the + # gate sees so a step's tool set is complete, not only the ones that + # prompted. No-op outside a workflow run or before a step is set. + mem = p_approval_memory.get(session_id) + if mem is None or mem.current_step_id is None: + return + mem.step_usage.setdefault(mem.current_step_id, {})[tool_name] = approved + + async def p_resolve_ask(tool_name, tool_input, sensitive_pattern) -> dict: + """Resolve an 'ask' policy. On a workflow run, reuse a remembered + decision (this step first, then the workflow-level fallback that + covers chat-seeded and cross-step approvals) instead of prompting, + and persist any fresh non-sensitive answer so later fires don't + re-ask. Shared by both gates so they can't disagree (and so the + first one's answer is reused by the second within the same call).""" + mem = p_approval_memory.get(session_id) + rememberable = ( + mem is not None + and sensitive_pattern is None + and tool_name != "AskUserQuestion" + ) + if rememberable: + sid = mem.current_step_id + prior_step = mem.step_usage.get(sid, {}).get(tool_name) if sid is not None else None + if prior_step is True: + return {"behavior": "allow"} + if prior_step is False: + return {"behavior": "deny", "message": "Denied by a remembered workflow permission"} + prior = mem.decisions.get(tool_name) + if prior == "allow": + p_note_tool_used(tool_name, True) + return {"behavior": "allow"} + if prior == "deny": + p_note_tool_used(tool_name, False) + return {"behavior": "deny", "message": "Denied by a remembered workflow permission"} + timeout = mem.ask_timeout if mem is not None else 600.0 + decision = await _request_user_approval( + tool_name, tool_input, sensitive_pattern=sensitive_pattern, timeout=timeout, + ) + if rememberable and decision.get("behavior") in ("allow", "deny"): + behavior = decision["behavior"] + mem.decisions[tool_name] = behavior + p_note_tool_used(tool_name, behavior == "allow") + if mem.remember: + try: + mem.remember(tool_name, behavior) + except Exception: + logger.exception("Failed to persist remembered workflow approval") + return decision + async def can_use_tool(tool_name, input_data, context): sensitive_pattern: str | None = None if tool_name != "AskUserQuestion": @@ -803,11 +911,13 @@ class AgentManager: _get_effective_policy(tool_name), tool_name, input_data ) if policy == "always_allow": + p_note_tool_used(tool_name, True) return PermissionResultAllow(updated_input=input_data) if policy == "deny": + p_note_tool_used(tool_name, False) return PermissionResultDeny(message="Tool denied by permission policy") - decision = await _request_user_approval(tool_name, input_data, sensitive_pattern=sensitive_pattern) + decision = await p_resolve_ask(tool_name, input_data, sensitive_pattern) if decision.get("behavior") == "allow": return PermissionResultAllow( updated_input=decision.get("updated_input", input_data) @@ -828,7 +938,11 @@ class AgentManager: _get_effective_policy(tool_name), tool_name, tool_input ) + if policy == "always_allow": + p_note_tool_used(tool_name, True) + if policy == "deny": + p_note_tool_used(tool_name, False) return { "hookSpecificOutput": { "hookEventName": hook_event, @@ -838,7 +952,7 @@ class AgentManager: } if policy == "ask": - decision = await _request_user_approval(tool_name, tool_input, sensitive_pattern=sensitive_pattern) + decision = await p_resolve_ask(tool_name, tool_input, sensitive_pattern) if decision.get("behavior") == "allow": if tool_use_id: diff --git a/backend/apps/workflows/executor.py b/backend/apps/workflows/executor.py index 28073db5..9361bd42 100644 --- a/backend/apps/workflows/executor.py +++ b/backend/apps/workflows/executor.py @@ -37,6 +37,48 @@ def _resolve_allowed_tools(wf: Workflow) -> list[str]: return list(wf.actions.configured_sets) +def p_make_remember_approval(workflow_id: str): + def p_remember_approval(tool_name: str, behavior: str) -> None: + fresh = storage.get_workflow(workflow_id) + if fresh is None: + return + fresh.remembered_approvals = {**fresh.remembered_approvals, tool_name: behavior} + storage.save_workflow(fresh) + try: + from backend.apps.agents.core.ws_manager import ws_manager + asyncio.get_running_loop().create_task(ws_manager.broadcast_global("workflow:updated", { + "workflow_id": fresh.id, + "workflow": fresh.model_dump(mode="json"), + })) + except Exception: + pass + return p_remember_approval + + +def p_persist_step_tool_usage(workflow_id: str, step_usage: dict[str, dict[str, bool]]) -> None: + fresh = storage.get_workflow(workflow_id) + if fresh is None: + return + live_ids = {s.id for s in fresh.steps} + draft = getattr(fresh, "draft_steps", None) + if draft is not None: + live_ids.update(s.id for s in draft) + fresh.step_tool_usage = { + sid: dict(tools) + for sid, tools in (step_usage or {}).items() + if sid in live_ids and isinstance(tools, dict) + } + storage.save_workflow(fresh) + try: + from backend.apps.agents.core.ws_manager import ws_manager + asyncio.get_running_loop().create_task(ws_manager.broadcast_global("workflow:updated", { + "workflow_id": fresh.id, + "workflow": fresh.model_dump(mode="json"), + })) + except Exception: + pass + + def _persist_run_fields(wf: Workflow, run_fields: dict, schedule_runs_count_delta: int = 0) -> None: """Merge run-side fields into the current on-disk workflow. @@ -85,7 +127,13 @@ def _monthly_spend_so_far(wf: Workflow) -> float: async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: Optional[datetime] = None) -> WorkflowRun: - from backend.apps.agents.agent_manager import agent_manager + from backend.apps.agents.agent_manager import ( + agent_manager, + clear_workflow_approval_memory, + get_workflow_step_usage, + set_workflow_approval_memory, + set_workflow_approval_step, + ) run = WorkflowRun( workflow_id=wf.id, @@ -134,7 +182,7 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O session = None try: - steps = [s.text for s in wf.steps if s.text and s.text.strip()] + steps = [s for s in wf.steps if s.text and s.text.strip()] if not steps: raise ValueError("Workflow has no steps") @@ -154,6 +202,19 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O run.session_id = session.id storage.record_run(run) + # Reuse the user's earlier allow/deny answers so an unattended fire + # doesn't park on a permission prompt. Scheduled runs prompt for an + # unseen tool only briefly (30s) before failing; manual/test runs are + # attended, so keep the roomy window. Sensitive-path prompts are never + # remembered (handled by the gate); they keep prompting every run. + set_workflow_approval_memory( + session.id, + decisions=dict(wf.remembered_approvals), + step_usage={sid: dict(tools) for sid, tools in wf.step_tool_usage.items()}, + remember=p_make_remember_approval(wf.id), + ask_timeout=30.0 if triggered_by == "schedule" else 600.0, + ) + # Background poller: surface the latest tool-call name as a # live "what's the agent doing" subtitle on the workflow:run # ws event. Cheap enough to run at 1.5s cadence; nothing else @@ -214,6 +275,7 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O # the disc immediately, not after the agent finishes the step. run.active_step_idx = idx run.last_tool_label = None + set_workflow_approval_step(session.id, step.id) try: from backend.apps.agents.core.ws_manager import ws_manager as _wsm await _wsm.broadcast_global("workflow:run", { @@ -222,7 +284,7 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O }) except Exception: pass - await agent_manager.send_message(session.id, step) + await agent_manager.send_message(session.id, step.text) await _await_session_idle(session.id) sess_state = agent_manager.sessions.get(session.id) if sess_state is not None and getattr(sess_state, "status", None) == "error": @@ -284,6 +346,12 @@ async def execute(wf: Workflow, triggered_by: str = "schedule", scheduled_for: O # the first page). close_session also drops in-memory state and # persists the final snapshot to disk. if session is not None: + try: + p_persist_step_tool_usage(wf.id, get_workflow_step_usage(session.id)) + except Exception: + logger.exception("persist step tool usage failed for workflow %s", wf.id) + set_workflow_approval_step(session.id, None) + clear_workflow_approval_memory(session.id) try: await agent_manager.close_session(session.id) except Exception: diff --git a/backend/apps/workflows/models.py b/backend/apps/workflows/models.py index db8b25fa..2a4689f2 100644 --- a/backend/apps/workflows/models.py +++ b/backend/apps/workflows/models.py @@ -110,6 +110,16 @@ class Workflow(BaseModel): # Sticky session id for the embedded scheduling agent (the chat that # turns "every Wednesday at 1pm" into a permission-gated tool call). schedule_agent_session_id: Optional[str] = None + # Tool permissions the user answered once and we reuse on later runs so an + # unattended scheduled fire doesn't stall waiting for someone to click. + # tool_name -> decision. Only ordinary "ask" tools land here; sensitive + # paths keep their own per-pattern trust and never auto-remember. + remembered_approvals: dict[str, Literal["allow", "deny"]] = Field(default_factory=dict) + # Behind-the-scenes record of which tools each step touched and whether each + # was permitted, keyed by stable step id (not index, so reorders don't + # scramble it). Auto-maintained on runs; enforcement stays workflow-level + # via remembered_approvals, this is the finer per-step picture. + step_tool_usage: dict[str, dict[str, bool]] = Field(default_factory=dict) class WorkflowRun(BaseModel): @@ -165,3 +175,5 @@ class WorkflowUpdate(BaseModel): mode: Optional[str] = None provider: Optional[str] = None cost_cap_usd_monthly: Optional[float] = None + remembered_approvals: Optional[dict[str, Literal["allow", "deny"]]] = None + step_tool_usage: Optional[dict[str, dict[str, bool]]] = None diff --git a/backend/apps/workflows/workflows.py b/backend/apps/workflows/workflows.py index b87eef12..f6844a8f 100644 --- a/backend/apps/workflows/workflows.py +++ b/backend/apps/workflows/workflows.py @@ -95,6 +95,41 @@ def _derive_icon(wf: Workflow) -> str: return "W" +def p_source_session_approvals(session_id: Optional[str]) -> dict[str, str]: + if not session_id: + return {} + try: + from backend.apps.agents.agent_manager import agent_manager + sess = agent_manager.sessions.get(session_id) + decisions = getattr(sess, "approval_decisions", None) if sess is not None else None + if decisions is None: + from backend.apps.agents.manager.session.session_store import _load_session_data + data = _load_session_data(session_id) or {} + decisions = data.get("approval_decisions") or [] + except Exception: + return {} + out: dict[str, str] = {} + for entry in decisions or []: + if not isinstance(entry, dict): + continue + if entry.get("sensitive_pattern"): + continue + tool = str(entry.get("tool") or "") + behavior = entry.get("behavior") + if tool and behavior in ("allow", "deny"): + out[tool] = behavior + return out + + +def p_prune_step_tool_usage(wf: Workflow) -> None: + live_ids = {s.id for s in wf.steps} + wf.step_tool_usage = { + sid: dict(tools) + for sid, tools in (wf.step_tool_usage or {}).items() + if sid in live_ids and isinstance(tools, dict) + } + + @workflows.router.get("/list") async def list_workflows(dashboard_id: Optional[str] = None): items = storage.list_workflows() @@ -134,6 +169,7 @@ async def create_workflow(body: WorkflowCreate): provider=body.provider or "anthropic", cost_cap_usd_monthly=body.cost_cap_usd_monthly, ) + wf.remembered_approvals = p_source_session_approvals(body.source_session_id) if not wf.icon: wf.icon = _derive_icon(wf) if wf.schedule.enabled: diff --git a/frontend/src/app/pages/Workflows/ActionsFacet.tsx b/frontend/src/app/pages/Workflows/ActionsFacet.tsx index 63cc626f..2329636c 100644 --- a/frontend/src/app/pages/Workflows/ActionsFacet.tsx +++ b/frontend/src/app/pages/Workflows/ActionsFacet.tsx @@ -75,6 +75,64 @@ export default function ActionsFacet({ draft, setDraft }: { draft: Workflow; set )} + + + + ); +} + +// Permissions this workflow learned on an earlier run and now reuses without +// asking. A reused "allow" runs unattended, so the user has to be able to see +// and take it back here. +function RememberedApprovals({ draft, setDraft }: { draft: Workflow; setDraft: (w: Workflow) => void }) { + const c = useClaudeTokens(); + const entries = Object.entries(draft.remembered_approvals || {}); + if (entries.length === 0) return null; + + const prettyName = (tool: string) => (tool.includes('__') ? tool.split('__').pop() || tool : tool); + const forget = (tool: string) => { + const next = { ...(draft.remembered_approvals || {}) }; + const nextStepUsage = Object.fromEntries( + Object.entries(draft.step_tool_usage || {}).map(([stepId, tools]) => { + const copy = { ...(tools || {}) }; + delete copy[tool]; + return [stepId, copy]; + }), + ); + delete next[tool]; + setDraft({ ...draft, remembered_approvals: next, step_tool_usage: nextStepUsage }); + }; + + return ( + + + + Saved permissions this workflow reuses on later runs. + + setDraft({ ...draft, remembered_approvals: {}, step_tool_usage: {} })} + role="button" + sx={{ fontSize: LABEL_FS, color: c.text.muted, cursor: 'pointer', whiteSpace: 'nowrap', ml: 1, '&:hover': { color: c.text.primary } }}> + Clear all + + + {entries.map(([tool, answer]) => ( + + + {prettyName(tool)} + + + {answer === 'allow' ? 'Allowed' : 'Blocked'} + + forget(tool)} + role="button" + aria-label={`Forget ${prettyName(tool)}`} + sx={{ display: 'inline-flex', fontSize: LABEL_FS, color: c.text.muted, cursor: 'pointer', px: 0.4, '&:hover': { color: c.status.error } }}> + ✕ + + + ))} ); } diff --git a/frontend/src/shared/state/workflowsSlice.ts b/frontend/src/shared/state/workflowsSlice.ts index a5da4c5f..6cff7512 100644 --- a/frontend/src/shared/state/workflowsSlice.ts +++ b/frontend/src/shared/state/workflowsSlice.ts @@ -81,6 +81,10 @@ export interface Workflow { edit_agent_session_id?: string | null; /** Sticky session id for the embedded scheduling agent (cadence -> gated tool call). */ schedule_agent_session_id?: string | null; + /** Tool permissions the user answered once and we reuse on later runs so an + * unattended scheduled fire doesn't stall on a prompt. tool name -> answer. */ + remembered_approvals?: Record; + step_tool_usage?: Record>; } export interface WorkflowRun {