diff --git a/agent/dashboard/routes.py b/agent/dashboard/routes.py index 5fff2e01..b6dabc9a 100644 --- a/agent/dashboard/routes.py +++ b/agent/dashboard/routes.py @@ -579,7 +579,7 @@ async def api_delete_review_style( async def api_list_threads( session: dict[str, Any] = _SESSION_DEP, ) -> list[dict[str, Any]]: - return await list_dashboard_threads(session["sub"]) + return await list_dashboard_threads(session["sub"], email=session.get("email")) @router.post("/threads") @@ -595,7 +595,7 @@ async def api_get_thread( thread_id: str, session: dict[str, Any] = _SESSION_DEP, ) -> dict[str, Any]: - return await get_dashboard_thread(thread_id, session["sub"]) + return await get_dashboard_thread(thread_id, session["sub"], email=session.get("email")) @router.post("/threads/{thread_id}/messages") @@ -604,7 +604,7 @@ async def api_send_thread_message( body: ThreadMessageBody, session: dict[str, Any] = _SESSION_DEP, ) -> dict[str, Any]: - return await send_dashboard_message(thread_id, session["sub"], body) + return await send_dashboard_message(thread_id, session["sub"], body, email=session.get("email")) @router.post("/threads/{thread_id}/cancel") @@ -612,7 +612,7 @@ async def api_cancel_thread( thread_id: str, session: dict[str, Any] = _SESSION_DEP, ) -> dict[str, Any]: - return await cancel_dashboard_thread(thread_id, session["sub"]) + return await cancel_dashboard_thread(thread_id, session["sub"], email=session.get("email")) @router.delete("/threads/{thread_id}") @@ -620,7 +620,7 @@ async def api_delete_thread( thread_id: str, session: dict[str, Any] = _SESSION_DEP, ) -> Response: - await delete_dashboard_thread(thread_id, session["sub"]) + await delete_dashboard_thread(thread_id, session["sub"], email=session.get("email")) return Response(status_code=204) @@ -634,7 +634,7 @@ async def api_stream_thread( async def event_generator(): async for chunk in stream_dashboard_thread( - thread_id, session["sub"], last_event_id=last_event_id + thread_id, session["sub"], email=session.get("email"), last_event_id=last_event_id ): yield chunk diff --git a/agent/dashboard/thread_api.py b/agent/dashboard/thread_api.py index f272469e..1182bc76 100644 --- a/agent/dashboard/thread_api.py +++ b/agent/dashboard/thread_api.py @@ -27,6 +27,8 @@ logger = logging.getLogger(__name__) _ASSISTANT_ID = "agent" _DASHBOARD_SOURCE = "dashboard" _DASHBOARD_STREAM_MODES: tuple[str, ...] = ("values", "updates", "messages-tuple") +# Sources whose threads should surface in the Agents UI (besides "dashboard"). +_SURFACED_SOURCES: tuple[str, ...] = ("dashboard", "github", "slack", "linear") def _agent_version_metadata() -> dict[str, str]: @@ -91,11 +93,28 @@ def _thread_owner_login(metadata: dict[str, Any]) -> str | None: return login.strip() if isinstance(login, str) and login.strip() else None -def _assert_thread_owner(metadata: dict[str, Any], login: str) -> None: - owner = _thread_owner_login(metadata) - if owner != login: - raise HTTPException(404, "thread not found") - if metadata.get("source") != _DASHBOARD_SOURCE: +def _thread_owner_email(metadata: dict[str, Any]) -> str | None: + email = metadata.get("triggering_user_email") + return email.strip().lower() if isinstance(email, str) and email.strip() else None + + +def _thread_source(metadata: dict[str, Any]) -> str: + source = metadata.get("source") + return source if isinstance(source, str) and source else _DASHBOARD_SOURCE + + +def _user_owns_thread(metadata: dict[str, Any], login: str, email: str | None) -> bool: + if _thread_source(metadata) not in _SURFACED_SOURCES: + return False + if _thread_owner_login(metadata) == login: + return True + if email and _thread_owner_email(metadata) == email.strip().lower(): + return True + return False + + +def _assert_thread_owner(metadata: dict[str, Any], login: str, email: str | None = None) -> None: + if not _user_owns_thread(metadata, login, email): raise HTTPException(404, "thread not found") @@ -153,6 +172,7 @@ def _thread_summary( "branch": metadata.get("branch_name") or metadata.get("base_branch") or "main", "model": model, "effort": effort, + "source": _thread_source(metadata), "status": status, "createdAt": int(created_at) if isinstance(created_at, (int, float)) else _now_ms(), "updatedAt": int(updated_at) if isinstance(updated_at, (int, float)) else _now_ms(), @@ -182,21 +202,40 @@ async def _latest_run_status(thread_id: str) -> str | None: return raw.lower() if isinstance(raw, str) else None -async def list_dashboard_threads(login: str, *, limit: int = 50) -> list[dict[str, Any]]: - threads = await langgraph_client().threads.search( - metadata={"source": _DASHBOARD_SOURCE, "github_login": login}, - limit=limit, - sort_by="updated_at", - sort_order="desc", - ) - out: list[dict[str, Any]] = [] - for thread in threads or []: - if isinstance(thread, dict): - out.append(_thread_summary(thread)) - return out +async def list_dashboard_threads( + login: str, *, email: str | None = None, limit: int = 50 +) -> list[dict[str, Any]]: + client = langgraph_client() + searches: list[dict[str, Any]] = [{"github_login": login}] + if email and email.strip(): + searches.append({"triggering_user_email": email.strip().lower()}) + + seen: dict[str, dict[str, Any]] = {} + for metadata_filter in searches: + threads = await client.threads.search( + metadata=metadata_filter, + limit=limit, + sort_by="updated_at", + sort_order="desc", + ) + for thread in threads or []: + if not isinstance(thread, dict): + continue + meta = thread.get("metadata") if isinstance(thread.get("metadata"), dict) else {} + if not _user_owns_thread(meta, login, email): + continue + thread_id = thread.get("thread_id") or thread.get("id") + if isinstance(thread_id, str) and thread_id not in seen: + seen[thread_id] = thread + + summaries = [_thread_summary(thread) for thread in seen.values()] + summaries.sort(key=lambda item: item.get("updatedAt", 0), reverse=True) + return summaries[:limit] -async def get_dashboard_thread(thread_id: str, login: str) -> dict[str, Any]: +async def get_dashboard_thread( + thread_id: str, login: str, *, email: str | None = None +) -> dict[str, Any]: client = langgraph_client() try: thread = await client.threads.get(thread_id) @@ -205,7 +244,7 @@ async def get_dashboard_thread(thread_id: str, login: str) -> dict[str, Any]: raise HTTPException(404, "thread not found") from exc metadata = thread.get("metadata") if isinstance(thread.get("metadata"), dict) else {} - _assert_thread_owner(metadata, login) + _assert_thread_owner(metadata, login, email) messages: list[dict[str, Any]] = [] try: @@ -321,7 +360,7 @@ async def create_dashboard_thread(login: str, body: ThreadCreateBody) -> dict[st async def send_dashboard_message( - thread_id: str, login: str, body: ThreadMessageBody + thread_id: str, login: str, body: ThreadMessageBody, *, email: str | None = None ) -> dict[str, Any]: client = langgraph_client() try: @@ -330,7 +369,7 @@ async def send_dashboard_message( raise HTTPException(404, "thread not found") from exc metadata = thread.get("metadata") if isinstance(thread.get("metadata"), dict) else {} - _assert_thread_owner(metadata, login) + _assert_thread_owner(metadata, login, email) owner, name, _ = _metadata_repo(metadata) if not owner or not name: raise HTTPException(400, "thread is missing repository metadata") @@ -355,13 +394,18 @@ async def send_dashboard_message( await _persist_dashboard_github_token(thread_id, login) profile = await get_profile(login) or {} + thread_source = _thread_source(metadata) configurable: dict[str, Any] = { "thread_id": thread_id, - "source": _DASHBOARD_SOURCE, + "source": thread_source, "github_login": login, "repo": {"owner": owner, "name": name}, "user_email": profile.get("email"), } + source_context = metadata.get("source_context") + if isinstance(source_context, dict): + for key, value in source_context.items(): + configurable.setdefault(key, value) if chosen_model and chosen_effort: configurable["agent_model_id"] = chosen_model configurable["agent_effort"] = chosen_effort @@ -384,7 +428,9 @@ async def send_dashboard_message( ) -async def cancel_dashboard_thread(thread_id: str, login: str) -> dict[str, Any]: +async def cancel_dashboard_thread( + thread_id: str, login: str, *, email: str | None = None +) -> dict[str, Any]: client = langgraph_client() try: thread = await client.threads.get(thread_id) @@ -392,7 +438,7 @@ async def cancel_dashboard_thread(thread_id: str, login: str) -> dict[str, Any]: raise HTTPException(404, "thread not found") from exc metadata = thread.get("metadata") if isinstance(thread.get("metadata"), dict) else {} - _assert_thread_owner(metadata, login) + _assert_thread_owner(metadata, login, email) run_id = metadata.get("latest_run_id") if isinstance(run_id, str) and run_id: @@ -411,7 +457,7 @@ async def cancel_dashboard_thread(thread_id: str, login: str) -> dict[str, Any]: ) -async def delete_dashboard_thread(thread_id: str, login: str) -> None: +async def delete_dashboard_thread(thread_id: str, login: str, *, email: str | None = None) -> None: client = langgraph_client() try: thread = await client.threads.get(thread_id) @@ -419,7 +465,7 @@ async def delete_dashboard_thread(thread_id: str, login: str) -> None: raise HTTPException(404, "thread not found") from exc metadata = thread.get("metadata") if isinstance(thread.get("metadata"), dict) else {} - _assert_thread_owner(metadata, login) + _assert_thread_owner(metadata, login, email) run_id = metadata.get("latest_run_id") if isinstance(run_id, str) and run_id: @@ -432,7 +478,7 @@ async def delete_dashboard_thread(thread_id: str, login: str) -> None: async def stream_dashboard_thread( - thread_id: str, login: str, *, last_event_id: str | None = None + thread_id: str, login: str, *, email: str | None = None, last_event_id: str | None = None ) -> AsyncIterator[str]: try: thread = await langgraph_client().threads.get(thread_id) @@ -440,7 +486,7 @@ async def stream_dashboard_thread( raise HTTPException(404, "thread not found") from exc metadata = thread.get("metadata") if isinstance(thread.get("metadata"), dict) else {} - _assert_thread_owner(metadata, login) + _assert_thread_owner(metadata, login, email) stream = await langgraph_client().threads.join_stream( thread_id, diff --git a/agent/webapp.py b/agent/webapp.py index 47469d1d..c2dba64b 100644 --- a/agent/webapp.py +++ b/agent/webapp.py @@ -8,6 +8,7 @@ import os import uuid from collections.abc import AsyncIterator from contextlib import asynccontextmanager +from datetime import UTC, datetime from typing import Any from urllib.parse import quote @@ -481,6 +482,68 @@ async def _upsert_slack_thread_repo_metadata( ) +async def upsert_agent_thread_owner_metadata( + thread_id: str, + *, + source: str, + repo_config: dict[str, str] | None = None, + github_login: str = "", + user_email: str = "", + title: str = "", + source_context: dict[str, Any] | None = None, +) -> None: + """Persist owner/source metadata so the dashboard can surface non-dashboard threads. + + Webhook-triggered runs only pass ``source``/``github_login`` through the run + config; the Agents UI lists and authorizes threads by thread *metadata*, so we + mirror the owner-identifying fields onto the thread here. + """ + now_ms = int(datetime.now(UTC).timestamp() * 1000) + resolved_login = github_login or resolve_login_from_email(user_email) or "" + metadata: dict[str, Any] = {"source": source, "updated_at_ms": now_ms} + if isinstance(repo_config, dict) and repo_config.get("owner") and repo_config.get("name"): + metadata["repo"] = repo_config + metadata["repo_owner"] = repo_config["owner"] + metadata["repo_name"] = repo_config["name"] + if resolved_login: + metadata["github_login"] = resolved_login + if user_email: + metadata["triggering_user_email"] = user_email.strip().lower() + if title: + metadata["title"] = title[:80] + if source_context: + metadata["source_context"] = source_context + + langgraph_client = get_client(url=LANGGRAPH_URL) + try: + existing = await langgraph_client.threads.get(thread_id) + except Exception as exc: # noqa: BLE001 + if not _is_not_found_error(exc): + logger.exception("Failed to read thread %s for owner metadata", thread_id) + existing = None + + existing_meta = ( + existing.get("metadata") + if isinstance(existing, dict) and isinstance(existing.get("metadata"), dict) + else {} + ) + if existing_meta.get("created_at_ms") is None: + metadata["created_at_ms"] = now_ms + if existing_meta.get("title") and "title" in metadata: + # Preserve a title that was already chosen (first message wins). + metadata.pop("title") + + try: + if existing is None: + await langgraph_client.threads.create( + thread_id=thread_id, if_exists="do_nothing", metadata=metadata + ) + else: + await langgraph_client.threads.update(thread_id=thread_id, metadata=metadata) + except Exception: # noqa: BLE001 + logger.exception("Failed to persist owner metadata for thread %s", thread_id) + + async def get_slack_repo_config( channel_id: str, thread_ts: str, @@ -747,6 +810,15 @@ async def process_linear_issue( # noqa: PLR0912, PLR0915 "source": "linear", } + await upsert_agent_thread_owner_metadata( + thread_id, + source="linear", + repo_config=repo_config, + user_email=user_email or "", + title=title or identifier or "Linear issue", + source_context={"linear_issue": configurable["linear_issue"]}, + ) + logger.info("Checking if thread %s is active before creating run", thread_id) thread_active = await is_thread_active(thread_id) logger.info("Thread %s active status: %s", thread_id, thread_active) @@ -914,6 +986,14 @@ async def process_slack_mention(event_data: dict[str, Any], repo_config: dict[st langgraph_client = get_client(url=LANGGRAPH_URL) is_first_mention = not await _thread_exists(thread_id) await _upsert_slack_thread_repo_metadata(thread_id, repo_config, langgraph_client) + await upsert_agent_thread_owner_metadata( + thread_id, + source="slack", + repo_config=repo_config, + user_email=user_email or "", + title=clean_text if is_first_mention else "", + source_context={"slack_thread": configurable["slack_thread"]}, + ) thread_active = await is_thread_active(thread_id) if thread_active: @@ -1392,6 +1472,14 @@ async def _trigger_or_queue_run( pr_number: int, ) -> None: """Create a new agent run or queue the message if the thread is busy.""" + await upsert_agent_thread_owner_metadata( + thread_id, + source="github", + repo_config=repo_config, + github_login=github_login, + title=f"PR #{pr_number}" if pr_number else "", + source_context={"pr_number": pr_number} if pr_number else None, + ) thread_active = await is_thread_active(thread_id) if thread_active: logger.info("Thread %s is busy, queuing GitHub PR comment message", thread_id) @@ -2635,6 +2723,15 @@ async def process_github_issue(payload: dict[str, Any], event_type: str) -> None }, } + await upsert_agent_thread_owner_metadata( + thread_id, + source="github", + repo_config=repo_config, + github_login=github_login, + title=title or (f"Issue #{issue_number}" if issue_number else ""), + source_context={"github_issue": configurable["github_issue"]}, + ) + thread_active = await is_thread_active(thread_id) if thread_active: logger.info("Thread %s is busy, queuing GitHub issue message", thread_id) diff --git a/ui/src/components/agents/AgentRunCard.tsx b/ui/src/components/agents/AgentRunCard.tsx index 31de47d9..aafe4a4c 100644 --- a/ui/src/components/agents/AgentRunCard.tsx +++ b/ui/src/components/agents/AgentRunCard.tsx @@ -1,10 +1,26 @@ import { Link } from "@tanstack/react-router"; -import { CheckCircleIcon, GitBranchIcon, GitPullRequestIcon } from "@phosphor-icons/react"; +import { + ChatCircleIcon, + CheckCircleIcon, + GitBranchIcon, + GitPullRequestIcon, + GithubLogoIcon, + KanbanIcon, + SlackLogoIcon, +} from "@phosphor-icons/react"; +import type { Icon } from "@phosphor-icons/react"; +import type { AgentSource, AgentThread } from "@/lib/agents/types"; import { formatRelativeTime } from "@/lib/agents/api"; -import type { AgentThread } from "@/lib/agents/types"; import { cn } from "@/lib/utils"; +const SOURCE_META: Record = { + dashboard: { icon: ChatCircleIcon, label: "Dashboard" }, + github: { icon: GithubLogoIcon, label: "GitHub" }, + slack: { icon: SlackLogoIcon, label: "Slack" }, + linear: { icon: KanbanIcon, label: "Linear" }, +}; + interface AgentRunCardProps { thread: AgentThread; } @@ -12,6 +28,8 @@ interface AgentRunCardProps { export function AgentRunCard({ thread }: AgentRunCardProps) { const stats = thread.diffStats; const hasPr = Boolean(thread.pr); + const source = thread.source ? SOURCE_META[thread.source] : null; + const SourceIcon = source?.icon; return (
{thread.title}
+ {source && SourceIcon && ( + <> + + + {source.label} + + · + + )} {thread.model} · {thread.repo} diff --git a/ui/src/components/agents/AgentsSidebar.tsx b/ui/src/components/agents/AgentsSidebar.tsx index e7bc5640..7accf2b4 100644 --- a/ui/src/components/agents/AgentsSidebar.tsx +++ b/ui/src/components/agents/AgentsSidebar.tsx @@ -1,7 +1,17 @@ import { Link } from "@tanstack/react-router"; -import { ChartLineUpIcon, PlusIcon, XIcon } from "@phosphor-icons/react"; +import { + ChartLineUpIcon, + ChatCircleIcon, + GithubLogoIcon, + KanbanIcon, + PlusIcon, + SlackLogoIcon, + XIcon, +} from "@phosphor-icons/react"; +import type { Icon } from "@phosphor-icons/react"; import type { SessionUser } from "@/lib/api"; +import type { AgentSource, AgentThread } from "@/lib/agents/types"; import { SidebarUserMenu } from "@/components/SidebarUserMenu"; import { SidebarCollapseButton, @@ -10,9 +20,15 @@ import { } from "@/components/sidebar-layout"; import { groupThreads } from "@/lib/agents/api"; import { useAgentThreads, useDeleteAgentThread } from "@/lib/agents/queries"; -import type { AgentThread } from "@/lib/agents/types"; import { cn } from "@/lib/utils"; +const SOURCE_META: Record = { + dashboard: { icon: ChatCircleIcon, label: "Started from the dashboard" }, + github: { icon: GithubLogoIcon, label: "Triggered from GitHub" }, + slack: { icon: SlackLogoIcon, label: "Triggered from Slack" }, + linear: { icon: KanbanIcon, label: "Triggered from Linear" }, +}; + interface AgentsSidebarProps { user: SessionUser; activeThreadId?: string; @@ -117,6 +133,9 @@ function ThreadRow({ thread, isActive }: { thread: AgentThread; isActive: boolea deleteThread.mutate(thread.id); }; + const source = thread.source && thread.source !== "dashboard" ? SOURCE_META[thread.source] : null; + const SourceIcon = source?.icon; + return ( + {source && SourceIcon && ( + + {source.label} + + )} {thread.title} {badge && ( diff --git a/ui/src/lib/agents/types.ts b/ui/src/lib/agents/types.ts index a0a0016d..804c441a 100644 --- a/ui/src/lib/agents/types.ts +++ b/ui/src/lib/agents/types.ts @@ -13,6 +13,8 @@ export type TodoStatus = "pending" | "in_progress" | "completed"; export type AgentStatus = "idle" | "running" | "finished" | "interrupted" | "error"; +export type AgentSource = "dashboard" | "github" | "slack" | "linear"; + export interface TodoItem { content: string; status: TodoStatus; @@ -128,6 +130,7 @@ export interface AgentThread { branch: string; model: string; effort?: string | null; + source?: AgentSource; status: AgentStatus; createdAt: number; updatedAt: number;