From 2359e4546b292678325823bd752aaba835fedf93 Mon Sep 17 00:00:00 2001 From: arc53-machine <232052973+arc53-machine@users.noreply.github.com> Date: Tue, 29 Sep 2026 17:11:58 +0100 Subject: [PATCH] Say which attached resources stopped running and why The agent and workflow reads now return resource_states to people who may edit them: every attached tool, source and prompt (and workflow node tool and source) with active or stopped and the reason: deleted, owner_lost_access, the sponsor reasons, connection_needs_reconnect, connection_removed or connector_disabled. Each entry names the sponsor, someone other than the reader who can fix it, the service for a connection reason, and whether the reader may take it over or reconnect it. When something can be taken over, sponsor_audience says who it would reach. The state comes from the checks the run itself uses (ref_access, resolve_holder_tool and the tool's connection as the run resolves it), so the page and the run can't disagree. A run that leaves a resource out logs resource_stopped with the holder, type, id and reason. The workflow read gives sponsor details, run state and node resource names only to people who may edit it, and names only resources the workflow runs, someone sponsored, or the reader can see. Owner saves and new workflows now refuse node tools and sources the owner can't use, like editor saves. --- docsgpt/agents/tool_executor.py | 30 +- docsgpt/agents/workflows/node_agent.py | 3 + docsgpt/agents/workflows/workflow_engine.py | 19 +- .../api/answer/services/stream_processor.py | 30 +- docsgpt/api/user/agents/routes.py | 14 +- docsgpt/api/user/resource_access.py | 405 ++++++++++++- docsgpt/api/user/workflows/routes.py | 65 +- docsgpt/connectors/resolve.py | 97 ++- tests/agents/test_workflow_agent_types.py | 12 +- tests/api/test_agent_team_sharing.py | 2 + tests/api/user/test_resource_states.py | 562 ++++++++++++++++++ 11 files changed, 1165 insertions(+), 74 deletions(-) create mode 100644 tests/api/user/test_resource_states.py diff --git a/docsgpt/agents/tool_executor.py b/docsgpt/agents/tool_executor.py index 532d06a3..38a69958 100644 --- a/docsgpt/agents/tool_executor.py +++ b/docsgpt/agents/tool_executor.py @@ -564,6 +564,9 @@ class ToolExecutor: # Tool id -> the user to resolve it as, for a workflow node's tools # sponsored by an editor (see resource_access.active_sponsor). self.tool_principals: Dict[str, str] = {} + # The workflow row a node's tools belong to, so a dropped one is + # logged with why it doesn't run. + self.tool_holder: Optional[Dict] = None self.conversation_id: Optional[str] = None # Set by the workflow engine for agent nodes so run-scoped tools # (artifact_generator / code_executor) address artifacts by the @@ -638,13 +641,27 @@ class ToolExecutor: if row is None and str(tid) in self.tool_principals: row = resolve_tool_by_id(tid, self.tool_principals[str(tid)], user_tools_repo=tools_repo) if row is None: - logger.warning("tool id %s did not resolve; dropped from scoped toolset", tid) + self._log_dropped_scoped_tool(conn, tid, principal) continue if self.headless and is_headless_excluded_tool(row.get("name")): continue tools.append(row) return {str(tool["id"]): tool for tool in tools} + def _log_dropped_scoped_tool(self, conn, tool_id: str, principal: Optional[str]) -> None: + """Log a node tool the scoped toolset leaves out, with why when the workflow is known.""" + # Lazy: docsgpt.api's package import pulls in every route module. + from docsgpt.api.user.resource_access import log_stopped, resolve_holder_tool + + holder = {**(self.tool_holder or {}), "user_id": principal} + reason = None + try: + _row, access = resolve_holder_tool(conn, "workflow", holder, tool_id) + reason = access.reason + except Exception: + logger.exception("Could not tell why tool %s does not resolve", tool_id) + log_stopped("workflow", holder, "tool", tool_id, reason) + def _get_tools_by_api_key(self, api_key: str) -> Dict[str, Dict]: """Resolve an agent's toolset — exactly ``agents.tools``, no defaults.""" # Per-operation session: the answer pipeline spans a long-lived @@ -660,13 +677,14 @@ class ToolExecutor: row = resolve_tool_by_id(tid, owner, user_tools_repo=tools_repo) if row is None: # A tool the owner can't use runs as the editor who - # attached it, while they still qualify. + # attached it, while they still qualify: the same check + # the agent page's run state uses. # Lazy: docsgpt.api's package import pulls in every route module. - from docsgpt.api.user.resource_access import active_sponsor + from docsgpt.api.user.resource_access import log_stopped, resolve_holder_tool - sponsor = active_sponsor(conn, "agent", agent_data, "tool", str(tid)) - if sponsor: - row = resolve_tool_by_id(tid, sponsor, user_tools_repo=tools_repo) + row, access = resolve_holder_tool(conn, "agent", agent_data, tid, tools_repo=tools_repo) + if row is None: + log_stopped("agent", agent_data, "tool", tid, access.reason) if row is None: continue # Workflow-only builtins (read_document) never resolve for a diff --git a/docsgpt/agents/workflows/node_agent.py b/docsgpt/agents/workflows/node_agent.py index a2927f14..7681e430 100644 --- a/docsgpt/agents/workflows/node_agent.py +++ b/docsgpt/agents/workflows/node_agent.py @@ -21,6 +21,7 @@ class _WorkflowNodeMixin: tool_ids: Optional[List[str]] = None, tool_principals: Optional[Dict[str, str]] = None, tool_owner: Optional[str] = None, + tool_holder: Optional[dict] = None, **kwargs, ): super().__init__( @@ -40,6 +41,8 @@ class _WorkflowNodeMixin: # owner can't use resolve as the editor who attached them. self.tool_executor.tool_owner = tool_owner self.tool_executor.tool_principals = dict(tool_principals or {}) + # The workflow row, so a dropped tool is logged with why. + self.tool_executor.tool_holder = tool_holder class WorkflowNodeClassicAgent(_WorkflowNodeMixin, ClassicAgent): diff --git a/docsgpt/agents/workflows/workflow_engine.py b/docsgpt/agents/workflows/workflow_engine.py index c1db8dd1..a700991d 100644 --- a/docsgpt/agents/workflows/workflow_engine.py +++ b/docsgpt/agents/workflows/workflow_engine.py @@ -407,6 +407,7 @@ class WorkflowEngine: "tool_ids": node_config.tools, "tool_owner": self._workflow_owner_id(), "tool_principals": self._node_tool_principals(node_config.tools), + "tool_holder": getattr(self.agent, "workflow_row", None), "prompt": node_prompt, "chat_history": self.agent.chat_history, "decoded_token": self.agent.decoded_token, @@ -1407,25 +1408,21 @@ class WorkflowEngine: logger.warning("Workflow node sources dropped: no owner to authorize.") return [] - from docsgpt.api.user.resource_access import active_sponsor - from docsgpt.api.user.team_sharing import can_access + from docsgpt.api.user.resource_access import log_stopped, ref_access from docsgpt.storage.db.session import db_readonly - workflow_row = getattr(self.agent, "workflow_row", None) + # The same check the workflow page's run state uses: the owner, else + # the editor who attached it while they still qualify. + holder = {**(getattr(self.agent, "workflow_row", None) or {}), "user_id": owner} allowed = [] try: with db_readonly() as conn: for sid in ids: - if sid and ( - can_access(conn, "source", str(sid), owner) - or active_sponsor(conn, "workflow", workflow_row, "source", str(sid)) - ): + access = ref_access(conn, "workflow", holder, "source", str(sid)) if sid else None + if access is not None and access.principal: allowed.append(sid) else: - logger.warning( - "Workflow node source %s dropped: %s has no access.", - sid, owner, - ) + log_stopped("workflow", holder, "source", sid, access.reason if access else None) except Exception: logger.exception("Workflow node source authorization failed; dropping all.") return [] diff --git a/docsgpt/api/answer/services/stream_processor.py b/docsgpt/api/answer/services/stream_processor.py index fa4af24e..f0b96aa3 100644 --- a/docsgpt/api/answer/services/stream_processor.py +++ b/docsgpt/api/answer/services/stream_processor.py @@ -156,20 +156,27 @@ def authorized_prompt_id(prompt_id: Any, principal: Optional[str], agent: Option pid = str(prompt_id) if is_composed_preset(pid) or pid in _PROMPT_PRESETS_WITHOUT_ROW: return prompt_id - from docsgpt.api.user.resource_access import active_sponsor, resolve + from docsgpt.api.user.resource_access import ( + REASON_OWNER_LOST_ACCESS, + log_stopped, + ref_access, + ) + # The same check the agent page's run state uses: the principal, else a + # live sponsor on the agent. + holder = {**agent, "user_id": principal} if agent and agent.get("id") else {"user_id": principal} try: with db_readonly() as conn: - ra = resolve(conn, "prompt", pid, principal) if principal else None - usable = ra is not None and ra.can("use") - if not usable and agent and agent.get("id"): - usable = active_sponsor(conn, "agent", agent, "prompt", pid) is not None + access = ref_access(conn, "agent", holder, "prompt", pid) if principal else None except Exception: logger.exception("Prompt access check failed for %s", pid) - usable = False - if usable: + access = None + if access is not None and access.principal: return prompt_id - logger.info("prompt %s not usable by %s; using the default prompt", pid, principal) + log_stopped( + "agent" if agent else "chat", holder, "prompt", pid, + access.reason if access is not None else REASON_OWNER_LOST_ACCESS, + ) return "default" @@ -180,10 +187,11 @@ def _agent_source_doc(conn: Any, sources_repo: Any, agent: dict, source_id: Any) editor who attached it while they still qualify. Read unscoped once authorized: an owner-scoped read misses a team-shared source. """ - from docsgpt.api.user.resource_access import ref_principal + from docsgpt.api.user.resource_access import log_stopped, ref_access - if not ref_principal(conn, "agent", agent, "source", str(source_id)): - logger.info("agent %s source %s not usable; skipped", agent.get("id"), source_id) + access = ref_access(conn, "agent", agent, "source", str(source_id)) + if not access.principal: + log_stopped("agent", agent, "source", source_id, access.reason) return None return sources_repo.get_by_id(str(source_id)) diff --git a/docsgpt/api/user/agents/routes.py b/docsgpt/api/user/agents/routes.py index 4680915d..9cfc6cc1 100644 --- a/docsgpt/api/user/agents/routes.py +++ b/docsgpt/api/user/agents/routes.py @@ -35,7 +35,9 @@ from docsgpt.api.user.resource_access import ( plan_sponsors, require, resolve, + resource_states, settings_many, + sponsor_audience, sponsor_details, sponsor_refusal, ) @@ -560,6 +562,8 @@ class GetAgent(Resource): user = decoded_token["sub"] agent = None sponsored: list = [] + states: list = [] + audience = None with db_readonly() as conn: # Anyone who can see the agent reads it (a viewer needs it to # chat); what they get back is trimmed by their actions. @@ -567,10 +571,13 @@ class GetAgent(Resource): if ra is not None: agent = AgentsRepository(conn).get_by_id(ra.resource_id) # Edit-page detail: who vouches for resources the owner can't - # use, with their names. Only for people who may edit the - # agent: the names can be an editor's private resources. + # use, and which attached resources stopped running and why, + # with their names. Only for people who may edit the agent: + # the names can be an editor's private resources. if agent and ra.can("edit"): sponsored = sponsor_details(conn, "agent", agent, viewer=user) + states = resource_states(conn, "agent", agent, agent_refs(agent), user) + audience = sponsor_audience(conn, "agent", agent, states, sponsored) if not agent: return {"status": "Not found"}, 404 is_owner = ra.access == "owner" @@ -582,6 +589,9 @@ class GetAgent(Resource): access=ra.payload(), ) data["resource_sponsors"] = sponsored + data["resource_states"] = states + if audience is not None: + data["sponsor_audience"] = audience return make_response(jsonify(data), 200) except Exception as e: current_app.logger.error(f"Agent fetch error: {e}", exc_info=True) diff --git a/docsgpt/api/user/resource_access.py b/docsgpt/api/user/resource_access.py index 4fcf3bb1..9a01f017 100644 --- a/docsgpt/api/user/resource_access.py +++ b/docsgpt/api/user/resource_access.py @@ -513,12 +513,148 @@ def ref_principal( The owner when they may use it (the default), else a live sponsor. """ + return ref_access(conn, holder_type, holder, resource_type, resource_id).principal + + +# --- Run state of attached resources ---------------------------------------- +# +# One check decides whether an attached resource runs, for the run and for the +# edit page alike (``ref_access``, ``resolve_holder_tool``), so the page never +# says a resource runs when the run drops it, or the other way round. + +# Why an attached resource doesn't run (``resource_states`` ``reason``), next +# to the sponsor reasons above. +REASON_DELETED = "deleted" +REASON_OWNER_LOST_ACCESS = "owner_lost_access" +REASON_CONNECTION_NEEDS_RECONNECT = "connection_needs_reconnect" +REASON_CONNECTION_REMOVED = "connection_removed" +REASON_CONNECTOR_DISABLED = "connector_disabled" + +_REF_TABLES = {"source": "sources", "prompt": "prompts", "tool": "user_tools"} + + +@dataclass(frozen=True) +class RefAccess: + """Who a holder runs one referenced resource as, or why it doesn't run. + + Attributes: + principal: The user it is authorized as (the holder's owner or a live + sponsor); None when it doesn't run. + reason: None while it runs; else :data:`REASON_DELETED`, + :data:`REASON_OWNER_LOST_ACCESS` or a sponsor reason. + sponsor: The recorded sponsor, running or not. + """ + + principal: Optional[str] + reason: Optional[str] = None + sponsor: Optional[str] = None + + +def _ref_exists(conn: Connection, resource_type: str, resource_id: str) -> bool: + """Whether a row with this id exists, whoever owns it.""" + table = _REF_TABLES.get(resource_type) + if table is None or not looks_like_uuid(str(resource_id)): + return False + return conn.execute( + text(f"SELECT 1 FROM {table} WHERE id = CAST(:id AS uuid)"), {"id": str(resource_id)} + ).first() is not None + + +def ref_access( + conn: Connection, holder_type: str, holder: Optional[dict], resource_type: str, resource_id: str +) -> RefAccess: + """Whether and as whom a holder runs one referenced resource. + + The run and the edit page both ask this. The owner runs it when they may + use it, else a live sponsor does; otherwise it is stopped, because the row + is gone, a recorded sponsor no longer qualifies, or the owner lost access. + + Args: + conn: Open database connection. + holder_type: ``agent`` or ``workflow``. + holder: The holder row (needs ``user_id``; ``id`` and + ``resource_sponsors`` for sponsors). + resource_type: ``source``, ``prompt`` or ``tool``. + resource_id: The referenced id. + + Returns: + RefAccess: The principal, or the reason it doesn't run. + """ if not holder: - return None + return RefAccess(None, REASON_OWNER_LOST_ACCESS) + rid = str(resource_id) owner = holder.get("user_id") - if can_use_ref(conn, resource_type, str(resource_id), owner): - return owner - return active_sponsor(conn, holder_type, holder, resource_type, resource_id) + if can_use_ref(conn, resource_type, rid, owner): + return RefAccess(owner) + sponsor, sponsor_reason = sponsor_state(conn, holder_type, holder, resource_type, rid) + if sponsor and sponsor_reason is None: + return RefAccess(sponsor, None, sponsor) + if not _ref_exists(conn, resource_type, rid): + return RefAccess(None, REASON_DELETED, sponsor) + return RefAccess(None, sponsor_reason or REASON_OWNER_LOST_ACCESS, sponsor) + + +def resolve_holder_tool( + conn: Connection, holder_type: str, holder: Optional[dict], tool_id: str, *, tools_repo=None +) -> tuple[Optional[dict], RefAccess]: + """The tool row a holder runs ``tool_id`` with, and its access. + + Builtin and default tool ids resolve to their synthesized rows. A + ``user_tools`` row resolves as the holder's owner, else as its live + sponsor (see :func:`ref_access`); the row is the tool owner's either way. + + Args: + conn: Open database connection. + holder_type: ``agent`` or ``workflow``. + holder: The holder row. + tool_id: The referenced tool id. + tools_repo: A ``UserToolsRepository`` on ``conn`` to reuse. + + Returns: + ``(row, access)``: the row, None when it doesn't run, and why. + """ + # Lazy: default_tools imports this module lazily too. + from docsgpt.agents.default_tools import resolve_tool_by_id + + repo = tools_repo or UserToolsRepository(conn) + owner = (holder or {}).get("user_id") + row = resolve_tool_by_id(tool_id, owner, user_tools_repo=repo) + if row is not None: + return row, RefAccess(owner) + access = ref_access(conn, holder_type, holder, "tool", str(tool_id)) + if access.principal: + row = resolve_tool_by_id(tool_id, access.principal, user_tools_repo=repo) + if row is None: + reason = access.reason or REASON_DELETED + return None, RefAccess(None, reason, access.sponsor) + return row, access + + +def log_stopped( + holder_type: str, holder: Optional[dict], resource_type: str, resource_id, reason: Optional[str] +) -> None: + """Log one attached resource a run leaves out, greppable by ``resource_stopped``. + + Args: + holder_type: ``agent`` or ``workflow``. + holder: The holder row. + resource_type: ``source``, ``prompt`` or ``tool``. + resource_id: The referenced id. + reason: Why it doesn't run. + """ + holder_id = str((holder or {}).get("id") or "") + logger.info( + "resource_stopped holder=%s:%s type=%s id=%s reason=%s", + holder_type, holder_id, resource_type, resource_id, reason, + extra={ + "event": "resource_stopped", + "holder_type": holder_type, + "holder_id": holder_id, + "resource_type": resource_type, + "resource_id": str(resource_id), + "reason": reason, + }, + ) def _sponsorable(resource_type: str, resource_id: str) -> bool: @@ -879,3 +1015,264 @@ def sponsor_details( } ) return out + + +def holder_editable_by(conn: Connection, holder_type: str, holder: dict, user_id: Optional[str]) -> bool: + """Whether ``user_id`` may edit the agent or workflow ``holder``. + + Its owner always may; anyone else needs ``edit`` on the agent (for a + workflow, on one of its owner's agents that run it). + + Args: + conn: Open database connection. + holder_type: ``agent`` or ``workflow``. + holder: The holder row. + user_id: The reader. + + Returns: + bool: Whether they may edit it. + """ + if not user_id or not holder: + return False + if user_id == holder.get("user_id"): + return True + return _holder_editable_by(conn, holder_type, holder, user_id) + + +def _ref_rows(conn: Connection, refs: Iterable[tuple[str, str]]) -> dict[str, dict]: + """``":" -> {name, user_id}`` for referenced rows, owner-agnostic, per type in one query.""" + queries = { + "source": "SELECT id, name, user_id FROM sources WHERE id = ANY(CAST(:ids AS uuid[]))", + "prompt": "SELECT id, name, user_id FROM prompts WHERE id = ANY(CAST(:ids AS uuid[]))", + "tool": ( + "SELECT id, COALESCE(NULLIF(custom_name, ''), NULLIF(display_name, ''), name), user_id " + "FROM user_tools WHERE id = ANY(CAST(:ids AS uuid[]))" + ), + } + by_type: dict[str, list[str]] = {} + for resource_type, resource_id in refs: + if resource_type in queries and looks_like_uuid(str(resource_id)): + by_type.setdefault(resource_type, []).append(str(resource_id)) + out: dict[str, dict] = {} + for resource_type, ids in by_type.items(): + for rid, name, owner in conn.execute(text(queries[resource_type]), {"ids": ids}).fetchall(): + out[sponsor_key(resource_type, str(rid))] = {"name": name, "user_id": owner} + return out + + +def _user_labels(conn: Connection, user_ids: Iterable[Optional[str]]) -> dict[str, str]: + """``user_id -> email`` for the ones with an email on file.""" + ids = sorted({u for u in user_ids if u}) + if not ids: + return {} + return dict( + conn.execute( + text( + "SELECT user_id, email FROM users WHERE user_id = ANY(:ids) " + "AND email IS NOT NULL AND email <> ''" + ), + {"ids": ids}, + ).fetchall() + ) + + +def _connection_state(conn: Connection, tool: dict, owner: Optional[str], policies_box: list) -> tuple: + """``(reason, connection)`` for a connection-backed tool the holder runs as ``owner``.""" + from docsgpt.connectors import catalog, service + from docsgpt.connectors.resolve import connection_stop_reason, resolve_connection + + if not tool.get("connection_id") and not catalog.definition_for_tool(tool.get("name") or ""): + return None, None + if not policies_box: + policies_box.append(service.load_policies(conn)) + resolved = resolve_connection(tool, owner, conn=conn, policies=policies_box[0]) + reason = connection_stop_reason(tool, resolved) + if reason is None: + return None, None + if resolved is not None: + connection = { + "id": resolved.connection_id, + "connector_key": resolved.connector_key, + "name": resolved.connector_name, + } + else: + definition = catalog.definition_for_tool(tool.get("name") or "") + connection = {"id": None, "connector_key": definition.key, "name": definition.name} + return reason, connection + + +def resource_states( + conn: Connection, + holder_type: str, + holder: dict, + refs: Iterable[tuple[str, str]], + viewer: Optional[str], +) -> list[dict]: + """Whether each resource a holder references runs, for its edit page. + + Uses the run's own checks (:func:`ref_access`, :func:`resolve_holder_tool` + and the tool's connection as the run resolves it for the owner), so a + resource shown as stopped is one the run leaves out. Presets and builtin + tools always run and are not listed. Only for readers who may edit the + holder: the caller checks that. + + A name is given only for a resource that runs, that the reader can see + themselves, that someone sponsored on the holder, or that is attached to + an agent (agent saves check every reference), so a reference to someone + else's resource never reveals its name. + + Args: + conn: Open database connection. + holder_type: ``agent`` or ``workflow``. + holder: The holder row. + refs: The ``(type, id)`` resources it references. + viewer: The user reading the page. + + Returns: + list: Per resource ``{key, type, id, name, state, reason, sponsor, + contact, connection, can_confirm, can_reconnect}``. ``state`` is + ``active`` or ``stopped``; ``reason`` (None while active) is one of + ``deleted``, ``owner_lost_access``, the sponsor reasons, + ``connection_needs_reconnect``, ``connection_removed`` or + ``connector_disabled``. ``sponsor`` and ``contact`` are + ``{user_id, label}`` or None: the recorded sponsor, and someone other + than the reader who can fix it. ``connection`` (``{id, + connector_key, name}``) names the service for a connection reason. + ``can_confirm``: the reader may take it over on their next save; + ``can_reconnect``: the reader owns the connection that needs signing + in again. + """ + owner = holder.get("user_id") + seen: set[str] = set() + pairs: list[tuple[str, str]] = [] + for resource_type, resource_id in refs: + rid = str(resource_id).lower() + key = sponsor_key(resource_type, rid) + if key in seen or resource_type not in REF_USE_ACTION or not _sponsorable(resource_type, rid): + continue + seen.add(key) + pairs.append((resource_type, rid)) + if not pairs: + return [] + rows = _ref_rows(conn, pairs) + recorded = {k.lower(): v for k, v in (holder.get("resource_sponsors") or {}).items()} + viewer_edits = bool(viewer and viewer != owner and _holder_editable_by(conn, holder_type, holder, viewer)) + tools_repo = UserToolsRepository(conn) + policies_box: list = [] + entries = [] + for resource_type, rid in pairs: + key = sponsor_key(resource_type, rid) + info = rows.get(key) or {} + connection = None + if resource_type == "tool": + tool_row, access = resolve_holder_tool(conn, holder_type, holder, rid, tools_repo=tools_repo) + reason = access.reason + if tool_row is not None and reason is None: + reason, connection = _connection_state(conn, tool_row, owner, policies_box) + else: + access = ref_access(conn, holder_type, holder, resource_type, rid) + reason = access.reason + sponsor = recorded.get(key) + sponsor = sponsor if sponsor and sponsor != owner else None + resource_owner = info.get("user_id") + can_confirm = bool( + reason in (REASON_OWNER_LOST_ACCESS, REASON_CANNOT_EDIT_HOLDER, REASON_CANNOT_EDIT_RESOURCE) + and viewer_edits + and can_sponsor_ref(conn, resource_type, rid, viewer) + ) + contact = None + if reason == REASON_OWNER_LOST_ACCESS: + # The owner asks whoever shared it; an editor asks the owner. + contact = resource_owner if viewer == owner else owner + elif reason in (REASON_CONNECTION_NEEDS_RECONNECT, REASON_CONNECTION_REMOVED): + contact = resource_owner + if contact == viewer: + contact = None + name_visible = bool( + reason is None + or sponsor + or holder_type == "agent" + or (viewer and resolve(conn, resource_type, rid, viewer) is not None) + ) + entries.append({ + "key": key, + "type": resource_type, + "id": rid, + "name": info.get("name") if name_visible else None, + "state": "active" if reason is None else "stopped", + "reason": reason, + "sponsor": sponsor, + "contact": contact, + "connection": connection, + "can_confirm": can_confirm, + "can_reconnect": bool( + reason == REASON_CONNECTION_NEEDS_RECONNECT + and connection + and connection.get("id") + and viewer + and viewer == resource_owner + ), + }) + labels = _user_labels(conn, [u for e in entries for u in (e["sponsor"], e["contact"])]) + for entry in entries: + for field_name in ("sponsor", "contact"): + user_id = entry[field_name] + entry[field_name] = {"user_id": user_id, "label": labels.get(user_id) or user_id} if user_id else None + return entries + + +def visible_ref_ids( + conn: Connection, holder_type: str, holder: dict, refs: Iterable[tuple[str, str]], viewer: Optional[str] +) -> set[str]: + """Keys of references whose names ``viewer`` may read on the holder's edit page. + + Those the holder runs (its owner or a live sponsor may use them), those + someone sponsored on it, and those the reader can see themselves. + + Args: + conn: Open database connection. + holder_type: ``agent`` or ``workflow``. + holder: The holder row. + refs: ``(type, id)`` pairs. + viewer: The reader. + + Returns: + set: ``":"`` keys, ids lowercased. + """ + recorded = {k.lower() for k in (holder.get("resource_sponsors") or {})} + out: set[str] = set() + for resource_type, resource_id in refs: + rid = str(resource_id).lower() + key = sponsor_key(resource_type, rid) + if key in out or resource_type not in REF_USE_ACTION: + continue + if ( + key in recorded + or ref_access(conn, holder_type, holder, resource_type, rid).principal + or (viewer and resolve(conn, resource_type, rid, viewer) is not None) + ): + out.add(key) + return out + + +def sponsor_audience( + conn: Connection, holder_type: str, holder: dict, states: list[dict], sponsors: list[dict] +) -> Optional[dict]: + """The holder's audience when the reader may take something over, else None. + + A take-over runs the resource with the reader's access for everyone who + uses the holder, so the page shows them who that is before they agree. + + Args: + conn: Open database connection. + holder_type: ``agent`` or ``workflow``. + holder: The holder row. + states: :func:`resource_states` for the reader. + sponsors: :func:`sponsor_details` for the reader. + + Returns: + dict or None: :func:`holder_audience`, when any item has ``can_confirm``. + """ + if not any(item.get("can_confirm") for item in [*states, *sponsors]): + return None + return holder_audience(conn, holder_type, holder) diff --git a/docsgpt/api/user/workflows/routes.py b/docsgpt/api/user/workflows/routes.py index 5c4729a9..34673961 100644 --- a/docsgpt/api/user/workflows/routes.py +++ b/docsgpt/api/user/workflows/routes.py @@ -10,14 +10,18 @@ from docsgpt.agents.workflows.cel_evaluator import ( CelEvaluationError, validate_cel_expression, ) +from docsgpt.api.user import resource_access from docsgpt.api.user.resource_access import ( AccessDenied, can_use_ref, parse_confirmations, plan_sponsors, resolve, + resource_states, + sponsor_audience, sponsor_details, sponsor_refusal, + visible_ref_ids, ) from docsgpt.storage.db.base_repository import looks_like_uuid from docsgpt.storage.db.repositories.workflow_edges import WorkflowEdgesRepository @@ -127,24 +131,35 @@ def _node_refs(nodes: List[Dict]) -> List[Tuple[str, str]]: return refs -def _node_ref_details(nodes: List[Dict]) -> Dict[str, List[Dict]]: - """Names of every tool and source the graph's agent nodes reference. +def _node_ref_details(nodes: List[Dict], visible: Optional[Set[str]] = None) -> Dict[str, List[Dict]]: + """Names of the tools and sources the graph's agent nodes reference. Looked up by id whoever owns them, so an editor's node pickers can show - (and remove) the owner's private tools and sources. + (and remove) the owner's private tools and sources. Builtin tool ids + always resolve; any other id only when its ``":"`` key is in + ``visible`` (see ``resource_access.visible_ref_ids``), so a node naming + someone else's resource never reveals its name. Args: nodes: Nodes in builder shape. + visible: Keys whose names may be read; None allows every id. Returns: dict: ``tools`` as ``[{id, name, display_name}]`` and ``sources`` as ``[{id, name}]``, each id once. """ + from docsgpt.agents.default_tools import is_synthesized_tool_id from docsgpt.api.user.base import resolve_source_details, resolve_tool_details tool_ids: List[str] = [] source_ids: List[str] = [] for resource_type, resource_id in _node_refs(nodes): + if ( + visible is not None + and not (resource_type == "tool" and is_synthesized_tool_id(resource_id)) + and f"{resource_type}:{resource_id.lower()}" not in visible + ): + continue bucket = tool_ids if resource_type == "tool" else source_ids if resource_id not in bucket: bucket.append(resource_id) @@ -162,13 +177,16 @@ def _new_node_ref_denied( A workflow runs as its owner, so an editor saving the owner's graph must not reference the owner's private tools or sources: the caller's own access counts (``use_in_own`` for a tool, ``use`` for a source), not the - owner's. Refs already in the stored graph stay, like an agent's. + owner's. The owner's own saves are checked the same way, so no graph + names a resource its owner never could use. Refs already in the stored + graph stay, like an agent's. Args: conn: Open database connection. - previous_nodes: The stored graph's nodes, in builder shape. + previous_nodes: The stored graph's nodes, in builder shape (empty + when creating the workflow). new_nodes: The nodes being saved. - caller: The editor saving. + caller: The user saving, owner or editor. Returns: An :class:`AccessDenied` to return, or None when every new ref is fine. @@ -629,6 +647,9 @@ class WorkflowList(Resource): try: with db_session() as conn: + denied = _new_node_ref_denied(conn, [], nodes_data, user_id) + if denied is not None: + return _denied(denied) repo = WorkflowsRepository(conn) workflow = repo.create(user_id, name, description=description) pg_workflow_id = str(workflow["id"]) @@ -660,9 +681,24 @@ class WorkflowDetail(Resource): edges = WorkflowEdgesRepository(conn).find_by_version( pg_workflow_id, graph_version, ) - sponsored = sponsor_details(conn, "workflow", workflow, viewer=user_id) - serialized_nodes = [serialize_node(n) for n in nodes] - ref_details = _node_ref_details(serialized_nodes) + serialized_nodes = [serialize_node(n) for n in nodes] + # Edit-page detail (sponsors, run state, names of node + # resources) only for people who may edit the workflow. + sponsored: list = [] + states: list = [] + audience = None + visible: Optional[Set[str]] = None + if resource_access.holder_editable_by(conn, "workflow", workflow, user_id): + refs = _node_refs(serialized_nodes) + sponsored = sponsor_details(conn, "workflow", workflow, viewer=user_id) + states = resource_states(conn, "workflow", workflow, refs, user_id) + audience = sponsor_audience(conn, "workflow", workflow, states, sponsored) + visible = visible_ref_ids(conn, "workflow", workflow, refs, user_id) + ref_details = ( + _node_ref_details(serialized_nodes, visible) + if visible is not None + else {"tools": [], "sources": []} + ) except Exception as err: return _workflow_error_response("Failed to fetch workflow", err) @@ -672,7 +708,9 @@ class WorkflowDetail(Resource): "nodes": serialized_nodes, "edges": [serialize_edge(e) for e in edges], "resource_sponsors": sponsored, + "resource_states": states, "ref_details": ref_details, + **({"sponsor_audience": audience} if audience is not None else {}), } ) @@ -712,10 +750,11 @@ class WorkflowDetail(Resource): pg_workflow_id, current_graph_version, ) ] - if acting != user_id: - denied = _new_node_ref_denied(conn, previous_nodes, nodes_data, user_id) - if denied is not None: - return _denied(denied) + # Every newly referenced node tool or source must be one the + # caller may use, the owner included. + denied = _new_node_ref_denied(conn, previous_nodes, nodes_data, user_id) + if denied is not None: + return _denied(denied) # A node tool/source the owner can't use runs as the editor # who attached it (its sponsor): only someone who owns or # edits it, and only once ``confirm_sponsor`` lists it. diff --git a/docsgpt/connectors/resolve.py b/docsgpt/connectors/resolve.py index 6d07f7af..6793e2ae 100644 --- a/docsgpt/connectors/resolve.py +++ b/docsgpt/connectors/resolve.py @@ -23,6 +23,11 @@ logger = logging.getLogger(__name__) MODE_OWNER = "owner" MODE_MEMBER = "member" +# Why a connection-backed tool can't run (``connection_stop_reason``). +CONNECTION_NEEDS_RECONNECT = "connection_needs_reconnect" +CONNECTION_REMOVED = "connection_removed" +CONNECTOR_DISABLED = "connector_disabled" + @dataclass(frozen=True) class ResolvedConnection: @@ -36,6 +41,7 @@ class ResolvedConnection: delegated: The row belongs to someone other than the invoker. writes_allowed: Whether an admin lets agents make changes through a connector that offers them as an opt-in (GitHub); True elsewhere. + enabled: Whether the connector is switched on (an admin can turn it off). """ row: Optional[dict] @@ -44,6 +50,7 @@ class ResolvedConnection: connector_name: Optional[str] delegated: bool = False writes_allowed: bool = True + enabled: bool = True @property def connection_id(self) -> Optional[str]: @@ -57,7 +64,13 @@ def _name_for(row: Optional[dict], fallback_key: Optional[str]) -> Optional[str] return definition.name if definition else None -def resolve_connection(resource: dict, invoker_user_id: Optional[str]) -> Optional[ResolvedConnection]: +def resolve_connection( + resource: dict, + invoker_user_id: Optional[str], + *, + conn=None, + policies: Optional[dict] = None, +) -> Optional[ResolvedConnection]: """Pick the connection a tool or source uses for ``invoker_user_id``. ``owner`` mode uses ``resource.connection_id``. ``member`` mode uses the @@ -67,39 +80,50 @@ def resolve_connection(resource: dict, invoker_user_id: Optional[str]) -> Option Args: resource: A ``user_tools`` or ``sources`` row. invoker_user_id: Who is running it. + conn: An open connection to reuse; a read-only one is opened when None. + policies: Connector policies already loaded with ``service.load_policies``, + so a caller resolving many resources loads them once. Returns: None when the resource has no connection at all; otherwise the resolution, possibly with ``available=False``. """ - connection_id = resource.get("connection_id") - if not connection_id: + if not resource.get("connection_id"): return None + if conn is None: + with db_readonly() as own_conn: + return _resolve(own_conn, resource, invoker_user_id, policies) + return _resolve(conn, resource, invoker_user_id, policies) + + +def _resolve(conn, resource: dict, invoker_user_id: Optional[str], policies: Optional[dict]) -> ResolvedConnection: + connection_id = resource.get("connection_id") mode = resource.get("credential_mode") or MODE_OWNER owner = resource.get("user_id") - with db_readonly() as conn: - repo = ConnectorSessionsRepository(conn) - owned = repo.get(str(connection_id)) - owned_key = catalog.connector_key_for_row(owned) if owned else None + repo = ConnectorSessionsRepository(conn) + owned = repo.get(str(connection_id)) + owned_key = catalog.connector_key_for_row(owned) if owned else None + if policies is None: policies = service.load_policies(conn) - policy = (policies.get(owned_key) or {}) if owned_key else {} - if policy.get("credential_mode") in (MODE_OWNER, MODE_MEMBER): - # An admin forces whose account every share of this connector uses. - mode = policy["credential_mode"] - if owned is not None and owner and owned.get("user_id") != owner: - # A resource may only point at its own owner's connection. - logger.warning( - "resource %s points at a connection it does not own", resource.get("id"), - ) - owned = None - row = owned - if mode == MODE_MEMBER and invoker_user_id and invoker_user_id != owner: - row = _member_connection(repo, owned, invoker_user_id) + policy = (policies.get(owned_key) or {}) if owned_key else {} + if policy.get("credential_mode") in (MODE_OWNER, MODE_MEMBER): + # An admin forces whose account every share of this connector uses. + mode = policy["credential_mode"] + if owned is not None and owner and owned.get("user_id") != owner: + # A resource may only point at its own owner's connection. + logger.warning( + "resource %s points at a connection it does not own", resource.get("id"), + ) + owned = None + row = owned + if mode == MODE_MEMBER and invoker_user_id and invoker_user_id != owner: + row = _member_connection(repo, owned, invoker_user_id) key = catalog.connector_key_for_row(row or owned or {}) + enabled = service.connector_is_enabled(policies, key) available = ( row is not None and service.normalize_status(row) == service.STATUS_CONNECTED - and service.connector_is_enabled(policies, key) + and enabled ) return ResolvedConnection( row=row, @@ -108,9 +132,40 @@ def resolve_connection(resource: dict, invoker_user_id: Optional[str]) -> Option connector_name=_name_for(row or owned, key), delegated=bool(row and invoker_user_id and row.get("user_id") != invoker_user_id), writes_allowed=_writes_allowed(policies, key), + enabled=enabled, ) +def connection_stop_reason(tool: dict, resolved: Optional[ResolvedConnection]) -> Optional[str]: + """Why a tool's connection keeps it from running, or None when it can run. + + Args: + tool: The ``user_tools`` row. + resolved: What :func:`resolve_connection` returned for it. + + Returns: + :data:`CONNECTION_REMOVED` when the connection is gone (its row was + deleted, or a built-in service's tool lost its connection and has no + credentials of its own), :data:`CONNECTOR_DISABLED` when an admin + turned the service off, :data:`CONNECTION_NEEDS_RECONNECT` when the + account must sign in again; else None. + """ + if resolved is None: + if tool.get("connection_id") or not catalog.definition_for_tool(tool.get("name") or ""): + return None + # Removing a connection but keeping its tools nulls their link; a + # built-in service's tool has no secrets of its own to fall back to. + config = tool.get("config") or {} + return None if config.get("encrypted_credentials") else CONNECTION_REMOVED + if resolved.row is None: + return CONNECTION_REMOVED + if not resolved.enabled: + return CONNECTOR_DISABLED + if not resolved.available: + return CONNECTION_NEEDS_RECONNECT + return None + + def _writes_allowed(policies: dict, key: Optional[str]) -> bool: definition = catalog.get_definition(key) if key else None if definition is None or not definition.mcp_write_url: diff --git a/tests/agents/test_workflow_agent_types.py b/tests/agents/test_workflow_agent_types.py index 18bf22d3..7f0a09ee 100644 --- a/tests/agents/test_workflow_agent_types.py +++ b/tests/agents/test_workflow_agent_types.py @@ -540,18 +540,18 @@ class TestWorkflowNodeSourceAuthorization: monkeypatch.setattr(session, "db_readonly", _conn) def test_owner_sources_survive(self, monkeypatch): - import docsgpt.api.user.team_sharing as ts + import docsgpt.api.user.resource_access as ra self._stub_db(monkeypatch) - monkeypatch.setattr(ts, "can_access", lambda *a, **k: True) + monkeypatch.setattr(ra, "can_use_ref", lambda *a, **k: True) engine = self._engine("owner") assert engine._authorized_node_sources(["s1", "s2"]) == ["s1", "s2"] def test_foreign_sources_are_dropped(self, monkeypatch): - import docsgpt.api.user.team_sharing as ts + import docsgpt.api.user.resource_access as ra self._stub_db(monkeypatch) - monkeypatch.setattr(ts, "can_access", lambda conn, k, sid, u: sid == "mine") + monkeypatch.setattr(ra, "can_use_ref", lambda conn, k, sid, u: sid == "mine") engine = self._engine("owner") assert engine._authorized_node_sources(["mine", "theirs"]) == ["mine"] @@ -563,13 +563,13 @@ class TestWorkflowNodeSourceAuthorization: assert engine._authorized_node_sources(["s1"]) == [] def test_authorization_error_fails_closed(self, monkeypatch): - import docsgpt.api.user.team_sharing as ts + import docsgpt.api.user.resource_access as ra def _boom(*a, **k): raise RuntimeError("db down") self._stub_db(monkeypatch) - monkeypatch.setattr(ts, "can_access", _boom) + monkeypatch.setattr(ra, "can_use_ref", _boom) engine = self._engine("owner") assert engine._authorized_node_sources(["s1"]) == [] diff --git a/tests/api/test_agent_team_sharing.py b/tests/api/test_agent_team_sharing.py index 0f8f92cd..d843e5a1 100644 --- a/tests/api/test_agent_team_sharing.py +++ b/tests/api/test_agent_team_sharing.py @@ -74,6 +74,8 @@ def _patches(sub, repo, team_access, *, prompt_name="Resolved Prompt", source_de "docsgpt.api.user.agents.routes.resolve_source_details", return_value=source_details, ), + # Run state reads the grants live; covered in test_resource_states. + patch("docsgpt.api.user.agents.routes.resource_states", return_value=[]), ] diff --git a/tests/api/user/test_resource_states.py b/tests/api/user/test_resource_states.py new file mode 100644 index 00000000..c4668d0a --- /dev/null +++ b/tests/api/user/test_resource_states.py @@ -0,0 +1,562 @@ +"""Run state of every resource attached to an agent or a workflow's nodes. + +An agent (or workflow) runs its attached tools, sources and prompt as its +owner, or as the editor who sponsored them. When one stops being usable the +run drops it (a prompt falls back to the default) and the edit page says why: +``resource_states`` on the agent and workflow reads. The state comes from the +same checks the run uses, so a resource marked stopped is never used by a run +and one marked active is. Uses real repositories on ``pg_conn``. +""" + +from __future__ import annotations + +import logging +import uuid + +import pytest +from flask import Flask +from sqlalchemy import text + +from docsgpt.api.user.resource_access import ( + REASON_CANNOT_EDIT_HOLDER, + REASON_CANNOT_EDIT_RESOURCE, + REASON_CONNECTION_NEEDS_RECONNECT, + REASON_CONNECTION_REMOVED, + REASON_CONNECTOR_DISABLED, + REASON_DELETED, + REASON_OWNER_LOST_ACCESS, + agent_refs, + ref_access, + resource_states, +) +from docsgpt.storage.db.repositories.agents import AgentsRepository +from docsgpt.storage.db.repositories.connector_policies import ConnectorPoliciesRepository +from docsgpt.storage.db.repositories.prompts import PromptsRepository +from docsgpt.storage.db.repositories.sources import SourcesRepository +from docsgpt.storage.db.repositories.team_members import TeamMembersRepository +from docsgpt.storage.db.repositories.team_resource_grants import ( + TeamResourceGrantsRepository, +) +from docsgpt.storage.db.repositories.user_tools import UserToolsRepository +from docsgpt.storage.db.repositories.workflows import WorkflowsRepository +from tests.api.user.test_resource_sponsors import ( + EDITOR, + OTHER, + OWNER, + VIEWER, + _agent, + _call, + _confirm, + _editor_resources, + _patch_db, + _put, + _row, + _status, + _wf_body, +) + + +@pytest.fixture +def app(): + return Flask(__name__) + + +def _body(resp) -> dict: + return resp[0] if isinstance(resp, tuple) else resp.get_json() + + +def _by_key(states) -> dict: + return {s["key"]: s for s in states} + + +def _team_source(conn, team_id, level="viewer"): + """OTHER's source, shared with the agent's team (OWNER a member).""" + if not TeamMembersRepository(conn).is_member(OWNER, team_id): + TeamMembersRepository(conn).add_member(team_id, OWNER) + if not TeamMembersRepository(conn).is_member(OTHER, team_id): + TeamMembersRepository(conn).add_member(team_id, OTHER) + source = str(SourcesRepository(conn).create("team-src", user_id=OTHER)["id"]) + TeamResourceGrantsRepository(conn).grant(team_id, "source", source, OTHER, OTHER, access_level=level) + return source + + +def _connection(conn, user=OWNER, status="connected", provider="telegram"): + return str(conn.execute( + text( + "INSERT INTO connector_sessions (user_id, provider, connector_key, auth_kind, status, " + "account_label) VALUES (:u, :p, :p, 'api_key', :s, :u) RETURNING id" + ), + {"u": user, "p": provider, "s": status}, + ).scalar()) + + +def _states(conn, agent_id, viewer=OWNER): + agent = _row(conn, agent_id) + return _by_key(resource_states(conn, "agent", agent, agent_refs(agent), viewer)) + + +def _get_agent(app, conn, agent_id, user): + from docsgpt.api.user.agents.routes import GetAgent + + return _call(app, conn, GetAgent, "get", f"/api/get_agent?id={agent_id}", user).get_json() + + +# --------------------------------------------------------------------------- +# Reasons +# --------------------------------------------------------------------------- + + +class TestReasons: + def test_owned_resources_are_active(self, pg_conn): + tool = str(UserToolsRepository(pg_conn).create(OWNER, "api_tool")["id"]) + source = str(SourcesRepository(pg_conn).create("mine", user_id=OWNER)["id"]) + prompt = str(PromptsRepository(pg_conn).create(OWNER, "p", "x")["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool], source_id=source, prompt_id=prompt) + states = _states(pg_conn, agent_id) + assert {k: s["state"] for k, s in states.items()} == { + f"tool:{tool}": "active", f"source:{source}": "active", f"prompt:{prompt}": "active", + } + assert all(s["reason"] is None for s in states.values()) + + def test_owner_lost_team_grant(self, pg_conn): + agent_id, team_id = _agent(pg_conn) + source = _team_source(pg_conn, team_id) + AgentsRepository(pg_conn).update_by_id(agent_id, {"extra_source_ids": [source]}) + assert _states(pg_conn, agent_id)[f"source:{source}"]["state"] == "active" + + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "source", source) + state = _states(pg_conn, agent_id)[f"source:{source}"] + assert (state["state"], state["reason"]) == ("stopped", REASON_OWNER_LOST_ACCESS) + # The owner is told whom to ask: the source's owner. + assert state["contact"] == {"user_id": OTHER, "label": OTHER} + assert state["name"] == "team-src" + + def test_deleted_tool(self, pg_conn): + """A deleted source or prompt leaves the agent by itself (FK, trigger); a tool id stays.""" + tool = str(UserToolsRepository(pg_conn).create(OWNER, "api_tool")["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool]) + UserToolsRepository(pg_conn).delete(tool, OWNER) + state = _states(pg_conn, agent_id)[f"tool:{tool}"] + assert (state["state"], state["reason"]) == ("stopped", REASON_DELETED) + assert state["contact"] is None + assert state["name"] is None + + def test_sponsor_lost_the_agent(self, app, pg_conn): + agent_id, team_id = _agent(pg_conn) + tool, _, _ = _editor_resources(pg_conn) + assert _status(_put(app, pg_conn, agent_id, EDITOR, + {"tools": [tool], "confirm_sponsor": _confirm(("tool", tool))})) == 200 + assert _states(pg_conn, agent_id)[f"tool:{tool}"]["state"] == "active" + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "agent", agent_id, target_user_id=EDITOR) + state = _states(pg_conn, agent_id)[f"tool:{tool}"] + assert (state["state"], state["reason"]) == ("stopped", REASON_CANNOT_EDIT_HOLDER) + assert state["sponsor"] == {"user_id": EDITOR, "label": EDITOR} + + def test_sponsor_lost_the_resource(self, app, pg_conn): + agent_id, team_id = _agent(pg_conn) + tool = str(UserToolsRepository(pg_conn).create(OTHER, "api_tool")["id"]) + TeamMembersRepository(pg_conn).add_member(team_id, OTHER) + TeamResourceGrantsRepository(pg_conn).grant( + team_id, "tool", tool, OTHER, OTHER, access_level="editor", target_user_id=EDITOR + ) + assert _status(_put(app, pg_conn, agent_id, EDITOR, + {"tools": [tool], "confirm_sponsor": _confirm(("tool", tool))})) == 200 + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "tool", tool, target_user_id=EDITOR) + state = _states(pg_conn, agent_id)[f"tool:{tool}"] + assert state["reason"] == REASON_CANNOT_EDIT_RESOURCE + + def test_connection_needs_reconnect(self, pg_conn): + tool = str(UserToolsRepository(pg_conn).create( + OWNER, "telegram", connection_id=_connection(pg_conn, status="reconnect_needed"), + )["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool]) + with _patch_db(pg_conn): + state = _states(pg_conn, agent_id)[f"tool:{tool}"] + assert (state["state"], state["reason"]) == ("stopped", REASON_CONNECTION_NEEDS_RECONNECT) + assert state["can_reconnect"] is True + assert state["connection"]["connector_key"] == "telegram" + assert state["contact"] is None + + def test_disconnected_account_needs_reconnect(self, pg_conn): + tool = str(UserToolsRepository(pg_conn).create( + OWNER, "telegram", connection_id=_connection(pg_conn, status="disconnected"), + )["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool]) + with _patch_db(pg_conn): + state = _states(pg_conn, agent_id)[f"tool:{tool}"] + assert state["reason"] == REASON_CONNECTION_NEEDS_RECONNECT + + def test_editor_is_told_whom_to_ask_to_reconnect(self, pg_conn): + tool = str(UserToolsRepository(pg_conn).create( + OWNER, "telegram", connection_id=_connection(pg_conn, status="reconnect_needed"), + )["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool]) + with _patch_db(pg_conn): + state = _states(pg_conn, agent_id, viewer=EDITOR)[f"tool:{tool}"] + assert state["can_reconnect"] is False + assert state["contact"] == {"user_id": OWNER, "label": OWNER} + + def test_connection_removed_but_tool_kept(self, pg_conn): + from docsgpt.connectors import service + + cid = _connection(pg_conn) + tool = str(UserToolsRepository(pg_conn).create(OWNER, "telegram", connection_id=cid)["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool]) + row = pg_conn.execute(text("SELECT * FROM connector_sessions WHERE id = CAST(:id AS uuid)"), + {"id": cid}).mappings().one() + service.remove_connection(pg_conn, dict(row), tools="keep") + with _patch_db(pg_conn): + state = _states(pg_conn, agent_id)[f"tool:{tool}"] + assert (state["state"], state["reason"]) == ("stopped", REASON_CONNECTION_REMOVED) + assert state["connection"]["name"] == "Telegram" + assert state["can_reconnect"] is False + + def test_connector_disabled_by_admin(self, pg_conn): + tool = str(UserToolsRepository(pg_conn).create( + OWNER, "telegram", connection_id=_connection(pg_conn), + )["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool]) + with _patch_db(pg_conn): + assert _states(pg_conn, agent_id)[f"tool:{tool}"]["state"] == "active" + ConnectorPoliciesRepository(pg_conn).upsert("telegram", enabled=False) + state = _states(pg_conn, agent_id)[f"tool:{tool}"] + assert state["reason"] == REASON_CONNECTOR_DISABLED + + def test_prompt_the_owner_lost(self, pg_conn): + agent_id, team_id = _agent(pg_conn) + TeamMembersRepository(pg_conn).add_member(team_id, OTHER) + prompt = str(PromptsRepository(pg_conn).create(OTHER, "theirs", "x")["id"]) + TeamResourceGrantsRepository(pg_conn).grant(team_id, "prompt", prompt, OTHER, OTHER) + TeamMembersRepository(pg_conn).add_member(team_id, OWNER) + AgentsRepository(pg_conn).update_by_id(agent_id, {"prompt_id": prompt}) + assert _states(pg_conn, agent_id)[f"prompt:{prompt}"]["state"] == "active" + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "prompt", prompt) + assert _states(pg_conn, agent_id)[f"prompt:{prompt}"]["reason"] == REASON_OWNER_LOST_ACCESS + + def test_presets_and_builtin_tools_are_not_listed(self, pg_conn): + from docsgpt.agents.default_tools import loaded_builtin_agent_tools, synthesize_builtin_agent_tool + + builtin = next(iter(loaded_builtin_agent_tools()), None) + tools = [str(synthesize_builtin_agent_tool(builtin)["id"])] if builtin else [] + agent_id, _ = _agent(pg_conn, tools=tools) + assert _states(pg_conn, agent_id) == {} + + +class TestTakeOver: + def test_editor_who_can_edit_the_item_may_take_over(self, pg_conn): + """An owner-lost resource the reading editor can edit is theirs to take over.""" + agent_id, team_id = _agent(pg_conn) + source = _team_source(pg_conn, team_id) + AgentsRepository(pg_conn).update_by_id(agent_id, {"extra_source_ids": [source]}) + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "source", source) + TeamResourceGrantsRepository(pg_conn).grant( + team_id, "source", source, OTHER, OTHER, access_level="editor", target_user_id=EDITOR + ) + assert _states(pg_conn, agent_id, viewer=EDITOR)[f"source:{source}"]["can_confirm"] is True + assert _states(pg_conn, agent_id, viewer=OWNER)[f"source:{source}"]["can_confirm"] is False + + def test_take_over_of_an_owner_lost_item_is_accepted(self, app, pg_conn): + agent_id, team_id = _agent(pg_conn) + source = _team_source(pg_conn, team_id) + AgentsRepository(pg_conn).update_by_id(agent_id, {"extra_source_ids": [source]}) + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "source", source) + TeamResourceGrantsRepository(pg_conn).grant( + team_id, "source", source, OTHER, OTHER, access_level="editor", target_user_id=EDITOR + ) + resp = _put(app, pg_conn, agent_id, EDITOR, + {"sources": [source], "confirm_sponsor": _confirm(("source", source))}) + assert _status(resp) == 200, _body(resp) + assert _states(pg_conn, agent_id)[f"source:{source}"]["state"] == "active" + + def test_agent_read_carries_the_audience_when_something_can_be_taken_over(self, app, pg_conn): + agent_id, team_id = _agent(pg_conn) + tool = str(UserToolsRepository(pg_conn).create(OTHER, "api_tool")["id"]) + TeamMembersRepository(pg_conn).add_member(team_id, OTHER) + TeamResourceGrantsRepository(pg_conn).grant(team_id, "tool", tool, OTHER, OTHER, access_level="editor") + TeamResourceGrantsRepository(pg_conn).grant( + team_id, "agent", agent_id, OWNER, OWNER, access_level="editor", target_user_id=OTHER + ) + assert _status(_put(app, pg_conn, agent_id, EDITOR, + {"tools": [tool], "confirm_sponsor": _confirm(("tool", tool))})) == 200 + data = _get_agent(app, pg_conn, agent_id, OTHER) + assert "sponsor_audience" not in data + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "agent", agent_id, target_user_id=EDITOR) + data = _get_agent(app, pg_conn, agent_id, OTHER) + assert _by_key(data["resource_states"])[f"tool:{tool}"]["can_confirm"] is True + assert data["sponsor_audience"]["teams"] == ["T"] + + +# --------------------------------------------------------------------------- +# Who sees it +# --------------------------------------------------------------------------- + + +class TestVisibility: + def _agent_with_stopped_tool(self, pg_conn): + tool = str(UserToolsRepository(pg_conn).create(OWNER, "api_tool")["id"]) + agent_id, _ = _agent(pg_conn, tools=[tool]) + UserToolsRepository(pg_conn).delete(tool, OWNER) + return agent_id, tool + + @pytest.mark.parametrize("user", [OWNER, EDITOR]) + def test_owner_and_editor_get_states(self, app, pg_conn, user): + agent_id, tool = self._agent_with_stopped_tool(pg_conn) + data = _get_agent(app, pg_conn, agent_id, user) + assert _by_key(data["resource_states"])[f"tool:{tool}"]["reason"] == REASON_DELETED + + def test_viewer_gets_none(self, app, pg_conn): + agent_id, _ = self._agent_with_stopped_tool(pg_conn) + data = _get_agent(app, pg_conn, agent_id, VIEWER) + assert data.get("resource_states", []) == [] + assert "sponsor_audience" not in data + + def test_shared_agent_view_carries_no_states(self, app, pg_conn): + """The public-link read never carries run state.""" + from docsgpt.api.user.agents.sharing import SharedAgent + + agent_id, _ = self._agent_with_stopped_tool(pg_conn) + token = uuid.uuid4().hex + AgentsRepository(pg_conn).update_by_id(agent_id, {"shared": True, "shared_token": token}) + resp = _call(app, pg_conn, SharedAgent, "get", f"/api/shared_agent?token={token}", VIEWER) + body = _body(resp) + assert "resource_states" not in body + assert "sponsor_audience" not in body + + +# --------------------------------------------------------------------------- +# Parity with the run +# --------------------------------------------------------------------------- + + +class TestParityWithRun: + def test_stopped_source_is_not_retrieved(self, app, pg_conn): + from docsgpt.api.answer.services.stream_processor import authorized_agent_sources + + agent_id, team_id = _agent(pg_conn) + source = _team_source(pg_conn, team_id) + own = str(SourcesRepository(pg_conn).create("own", user_id=OWNER)["id"]) + AgentsRepository(pg_conn).update_by_id(agent_id, {"source_id": own, "extra_source_ids": [source]}) + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "source", source) + states = _states(pg_conn, agent_id) + _, rows = authorized_agent_sources(pg_conn, _row(pg_conn, agent_id)) + used = {f"source:{r['id']}" for r in rows} + assert used == {k for k, s in states.items() if s["type"] == "source" and s["state"] == "active"} + assert f"source:{source}" not in used + + def test_stopped_tool_is_not_loaded(self, app, pg_conn): + from docsgpt.agents.tool_executor import ToolExecutor + + own = str(UserToolsRepository(pg_conn).create(OWNER, "api_tool")["id"]) + agent_id, team_id = _agent(pg_conn, tools=[own]) + tool, _, _ = _editor_resources(pg_conn) + assert _status(_put(app, pg_conn, agent_id, EDITOR, + {"tools": [own, tool], "confirm_sponsor": _confirm(("tool", tool))})) == 200 + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "agent", agent_id, target_user_id=EDITOR) + agent = _row(pg_conn, agent_id) + states = _states(pg_conn, agent_id) + with _patch_db(pg_conn): + loaded = ToolExecutor(user_api_key=agent["key"], user=OWNER)._get_tools_by_api_key(agent["key"]) + assert {f"tool:{t}" for t in loaded} == {k for k, s in states.items() if s["state"] == "active"} + + def test_stopped_prompt_falls_back_to_default(self, app, pg_conn): + from docsgpt.api.answer.services.stream_processor import authorized_prompt_id + + agent_id, team_id = _agent(pg_conn) + _, prompt, _ = _editor_resources(pg_conn) + assert _status(_put(app, pg_conn, agent_id, EDITOR, + {"prompt_id": prompt, "confirm_sponsor": _confirm(("prompt", prompt))})) == 200 + agent = _row(pg_conn, agent_id) + with _patch_db(pg_conn): + assert _states(pg_conn, agent_id)[f"prompt:{prompt}"]["state"] == "active" + assert authorized_prompt_id(prompt, OWNER, agent) == prompt + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "agent", agent_id, target_user_id=EDITOR) + assert _states(pg_conn, agent_id)[f"prompt:{prompt}"]["state"] == "stopped" + assert authorized_prompt_id(prompt, OWNER, _row(pg_conn, agent_id)) == "default" + + def test_ref_access_names_the_principal(self, app, pg_conn): + own = str(UserToolsRepository(pg_conn).create(OWNER, "api_tool")["id"]) + agent_id, _ = _agent(pg_conn, tools=[own]) + tool, _, _ = _editor_resources(pg_conn) + assert _status(_put(app, pg_conn, agent_id, EDITOR, + {"tools": [own, tool], "confirm_sponsor": _confirm(("tool", tool))})) == 200 + agent = _row(pg_conn, agent_id) + assert ref_access(pg_conn, "agent", agent, "tool", own).principal == OWNER + assert ref_access(pg_conn, "agent", agent, "tool", tool).principal == EDITOR + missing = str(uuid.uuid4()) + state = ref_access(pg_conn, "agent", agent, "tool", missing) + assert (state.principal, state.reason) == (None, REASON_DELETED) + + +class TestRunLog: + def test_dropped_source_logs_type_id_and_reason(self, pg_conn, caplog): + from docsgpt.api.answer.services.stream_processor import authorized_agent_sources + + agent_id, _ = _agent(pg_conn) + missing = str(uuid.uuid4()) + AgentsRepository(pg_conn).update_by_id(agent_id, {"extra_source_ids": [missing]}) + with caplog.at_level(logging.INFO): + authorized_agent_sources(pg_conn, _row(pg_conn, agent_id)) + [record] = [r for r in caplog.records if getattr(r, "event", None) == "resource_stopped"] + assert (record.holder_type, record.holder_id) == ("agent", agent_id) + assert (record.resource_type, record.resource_id, record.reason) == ("source", missing, REASON_DELETED) + assert f"reason={REASON_DELETED}" in record.getMessage() + + def test_dropped_tool_logs_reason(self, pg_conn, caplog): + from docsgpt.agents.tool_executor import ToolExecutor + + missing = str(uuid.uuid4()) + agent_id, _ = _agent(pg_conn, tools=[missing]) + agent = _row(pg_conn, agent_id) + with _patch_db(pg_conn), caplog.at_level(logging.INFO): + assert ToolExecutor(user_api_key=agent["key"], user=OWNER)._get_tools_by_api_key(agent["key"]) == {} + records = [r for r in caplog.records if getattr(r, "event", None) == "resource_stopped"] + assert [(r.resource_type, r.resource_id, r.reason) for r in records] == [("tool", missing, REASON_DELETED)] + + + def test_dropped_node_tool_logs_reason(self, pg_conn, caplog): + from docsgpt.agents.tool_executor import ToolExecutor + + wf = WorkflowsRepository(pg_conn).create(OWNER, "wf") + missing = str(uuid.uuid4()) + executor = ToolExecutor(user=VIEWER) + executor.allowed_tool_ids = [missing] + executor.tool_owner = OWNER + executor.tool_holder = wf + with _patch_db(pg_conn), caplog.at_level(logging.INFO): + assert executor.get_tools() == {} + [record] = [r for r in caplog.records if getattr(r, "event", None) == "resource_stopped"] + assert (record.holder_type, record.holder_id) == ("workflow", str(wf["id"])) + assert (record.resource_type, record.resource_id, record.reason) == ("tool", missing, REASON_DELETED) + + +# --------------------------------------------------------------------------- +# Workflows +# --------------------------------------------------------------------------- + + +class TestWorkflowStates: + def _setup(self, pg_conn): + wf = WorkflowsRepository(pg_conn).create(OWNER, "wf") + agent_id, team_id = _agent(pg_conn, agent_type="workflow", workflow_id=str(wf["id"])) + return str(wf["id"]), agent_id, team_id + + def _put(self, app, pg_conn, wid, user, body): + from docsgpt.api.user.workflows.routes import WorkflowDetail + + return _call(app, pg_conn, WorkflowDetail, "put", f"/api/workflows/{wid}", user, json=body, args=(wid,)) + + def _get(self, app, pg_conn, wid, user): + from docsgpt.api.user.workflows.routes import WorkflowDetail + + return _call(app, pg_conn, WorkflowDetail, "get", f"/api/workflows/{wid}", user, args=(wid,)) + + def test_node_states_on_the_workflow_read(self, app, pg_conn): + wid, agent_id, team_id = self._setup(pg_conn) + tool, _, source = _editor_resources(pg_conn) + confirm = _confirm(("tool", tool), ("source", source)) + assert _status(self._put(app, pg_conn, wid, EDITOR, _wf_body(tool, source, confirm=confirm))) == 200 + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "agent", agent_id, target_user_id=EDITOR) + data = _body(self._get(app, pg_conn, wid, OWNER))["data"] + states = _by_key(data["resource_states"]) + assert states[f"tool:{tool}"]["reason"] == REASON_CANNOT_EDIT_HOLDER + assert states[f"source:{source}"]["reason"] == REASON_CANNOT_EDIT_HOLDER + + def test_workflow_engine_drops_what_the_read_marks_stopped(self, app, pg_conn): + from types import SimpleNamespace + + from docsgpt.agents.workflows.workflow_engine import WorkflowEngine + + wid, agent_id, team_id = self._setup(pg_conn) + _, _, source = _editor_resources(pg_conn) + own = str(SourcesRepository(pg_conn).create("own", user_id=OWNER)["id"]) + assert _status(self._put(app, pg_conn, wid, OWNER, _wf_body(source=own))) == 200 + body = _wf_body(source=source, confirm=_confirm(("source", source))) + body["nodes"][1]["data"]["config"]["sources"] = [source, own] + assert _status(self._put(app, pg_conn, wid, EDITOR, body)) == 200 + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "agent", agent_id, target_user_id=EDITOR) + + engine = WorkflowEngine.__new__(WorkflowEngine) + engine.agent = SimpleNamespace( + workflow_row=WorkflowsRepository(pg_conn).get_by_id(wid), + _resolve_owner_id=lambda: OWNER, user=OWNER, decoded_token={"sub": OWNER}, + ) + with _patch_db(pg_conn): + allowed = engine._authorized_node_sources([source, own]) + states = _by_key(_body(self._get(app, pg_conn, wid, OWNER))["data"]["resource_states"]) + assert {f"source:{s}" for s in allowed} == {k for k, s in states.items() if s["state"] == "active"} + + def test_deleted_node_tool(self, app, pg_conn): + wid, _, _ = self._setup(pg_conn) + tool = str(UserToolsRepository(pg_conn).create(OWNER, "api_tool")["id"]) + assert _status(self._put(app, pg_conn, wid, OWNER, _wf_body(tool))) == 200 + UserToolsRepository(pg_conn).delete(tool, OWNER) + state = _by_key(_body(self._get(app, pg_conn, wid, OWNER))["data"]["resource_states"])[f"tool:{tool}"] + assert state["reason"] == REASON_DELETED + + def test_sponsor_details_only_for_editors(self, app, pg_conn): + """B: the workflow read gives sponsor details only to people who may edit.""" + from docsgpt.api.user import resource_access + + wid, _, _ = self._setup(pg_conn) + tool, _, _ = _editor_resources(pg_conn) + assert _status(self._put(app, pg_conn, wid, EDITOR, + _wf_body(tool, confirm=_confirm(("tool", tool))))) == 200 + assert _body(self._get(app, pg_conn, wid, EDITOR))["data"]["resource_sponsors"] + # A future role that may view but not edit reads no details. + original = resource_access.holder_editable_by + try: + resource_access.holder_editable_by = lambda *a, **k: False + data = _body(self._get(app, pg_conn, wid, EDITOR))["data"] + finally: + resource_access.holder_editable_by = original + assert data["resource_sponsors"] == [] + assert data["resource_states"] == [] + assert data["ref_details"] == {"tools": [], "sources": []} + + def test_ref_details_hide_names_nobody_here_can_use(self, app, pg_conn): + """A: a node naming someone else's resource doesn't reveal its name.""" + wid, _, _ = self._setup(pg_conn) + secret = str(UserToolsRepository(pg_conn).create("sp-stranger", "api_tool", custom_name="Secret")["id"]) + own = str(UserToolsRepository(pg_conn).create(OWNER, "api_tool", custom_name="Mine")["id"]) + body = _wf_body(tools=[own]) + assert _status(self._put(app, pg_conn, wid, OWNER, body)) == 200 + # Written straight to the graph, the way an older unchecked save left it. + pg_conn.execute( + text("UPDATE workflow_nodes SET config = jsonb_set(config, '{config,tools}', CAST(:t AS jsonb)) " + "WHERE workflow_id = CAST(:w AS uuid) AND node_type = 'agent'"), + {"t": f'["{own}", "{secret}"]', "w": wid}, + ) + data = _body(self._get(app, pg_conn, wid, OWNER))["data"] + names = {t["id"]: t.get("name") for t in data["ref_details"]["tools"]} + assert own in names + assert secret not in names + state = _by_key(data["resource_states"])[f"tool:{secret}"] + assert state["state"] == "stopped" and state["name"] is None + + def test_owner_save_refuses_a_node_ref_the_owner_cannot_use(self, app, pg_conn): + wid, _, _ = self._setup(pg_conn) + secret = str(UserToolsRepository(pg_conn).create("sp-stranger", "api_tool")["id"]) + resp = self._put(app, pg_conn, wid, OWNER, _wf_body(secret)) + assert _status(resp) == 403 + + def test_owner_create_refuses_a_node_ref_the_owner_cannot_use(self, app, pg_conn): + from docsgpt.api.user.workflows.routes import WorkflowList + + secret = str(SourcesRepository(pg_conn).create("x", user_id="sp-stranger")["id"]) + resp = _call(app, pg_conn, WorkflowList, "post", "/api/workflows", OWNER, json=_wf_body(source=secret)) + assert _status(resp) == 403 + + def test_workflow_audience_when_something_can_be_taken_over(self, app, pg_conn): + wid, agent_id, team_id = self._setup(pg_conn) + tool = str(UserToolsRepository(pg_conn).create(OTHER, "api_tool")["id"]) + TeamMembersRepository(pg_conn).add_member(team_id, OTHER) + TeamResourceGrantsRepository(pg_conn).grant(team_id, "tool", tool, OTHER, OTHER, access_level="editor") + TeamResourceGrantsRepository(pg_conn).grant( + team_id, "agent", agent_id, OWNER, OWNER, access_level="editor", target_user_id=OTHER + ) + assert _status(self._put(app, pg_conn, wid, EDITOR, + _wf_body(tool, confirm=_confirm(("tool", tool))))) == 200 + TeamResourceGrantsRepository(pg_conn).revoke(team_id, "agent", agent_id, target_user_id=EDITOR) + data = _body(self._get(app, pg_conn, wid, OTHER))["data"] + assert _by_key(data["resource_states"])[f"tool:{tool}"]["can_confirm"] is True + assert data["sponsor_audience"]["teams"] == ["T"]