feat: surface all of a user's agent threads in the Agents UI (#1366)

* feat: surface all of a user's agent threads in the Agents UI

Previously the Agents UI only listed and allowed opening threads with
source=dashboard. Threads triggered from GitHub, Slack, or Linear were
hidden and could not be opened.

- Persist owner-identifying metadata (source, github_login,
  triggering_user_email, source_context) onto the thread for each
  webhook-triggered main-agent run via upsert_agent_thread_owner_metadata,
  resolving a github_login from the triggering email where possible.
- Relax dashboard thread listing/ownership to surface github/slack/linear
  threads owned by the logged-in user (matched by github_login or email),
  and preserve the original source + reply-routing context when continuing
  such threads from the UI.
- Add a source field to the thread summary and render a source icon
  (GitHub/Slack/Linear) on thread items in the sidebar and run cards.

* Normalize triggering email when persisting thread owner metadata

Mixed-case Slack/Linear emails were stored verbatim, but the dashboard
searches with a lowercased value, so those threads never surfaced for
their owner. Normalize at write time to match _thread_owner_email.
This commit is contained in:
Johannes du Plessis 2026-06-01 13:01:20 -07:00 • committed by GitHub
parent dcef5ff71e
commit 015476230c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 239 additions and 38 deletions

View file

@ -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

View file

@ -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,

View file

@ -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)

View file

@ -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<AgentSource, { icon: Icon; label: string }> = {
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 (
<Link
@ -51,6 +69,15 @@ export function AgentRunCard({ thread }: AgentRunCardProps) {
<div className="min-w-0 flex-1">
<div className="truncate text-sm font-medium text-[var(--ui-text)]">{thread.title}</div>
<div className="mt-1 flex items-center gap-2 text-xs text-[var(--ui-text-dim)]">
{source && SourceIcon && (
<>
<span className="flex items-center gap-1" title={source.label}>
<SourceIcon className="size-3.5" weight="fill" aria-label={source.label} />
{source.label}
</span>
<span>·</span>
</>
)}
<span>{thread.model}</span>
<span>·</span>
<span>{thread.repo}</span>

View file

@ -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<AgentSource, { icon: Icon; label: string }> = {
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 (
<Link
to="/agents/$threadId"
@ -139,6 +158,15 @@ function ThreadRow({ thread, isActive }: { thread: AgentThread; isActive: boolea
: "bg-[var(--ui-border)]",
)}
/>
{source && SourceIcon && (
<SourceIcon
className="size-3.5 shrink-0 text-[var(--ui-text-dim)]"
weight="fill"
aria-label={source.label}
>
<title>{source.label}</title>
</SourceIcon>
)}
<span className="min-w-0 flex-1 truncate text-xs">{thread.title}</span>
{badge && (
<span className="shrink-0 rounded bg-[var(--ui-panel-2)] px-1.5 py-0.5 text-[10px] text-[var(--ui-text-dim)] group-hover:hidden">

View file

@ -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;