diff --git a/apps/web/src/app/(v2)/chat/[thread_id]/page.tsx b/apps/web/src/app/(v2)/chat/[thread_id]/page.tsx index 72c60b0d..bd3ca9c8 100644 --- a/apps/web/src/app/(v2)/chat/[thread_id]/page.tsx +++ b/apps/web/src/app/(v2)/chat/[thread_id]/page.tsx @@ -2,15 +2,16 @@ import { ThreadView } from "@/components/v2/thread-view"; import { ThreadViewLoading } from "@/components/v2/thread-view-loading"; -import { ThreadDisplayInfo, threadToDisplayInfo } from "@/components/v2/types"; -import { useThreads } from "@/hooks/useThreads"; +import { ThreadMetadata } from "@/components/v2/types"; +import { useThreadMetadata } from "@/hooks/useThreadMetadata"; +import { useThreadsSWR } from "@/hooks/useThreadsSWR"; import { useStream } from "@langchain/langgraph-sdk/react"; import { MANAGER_GRAPH_ID } from "@open-swe/shared/constants"; import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; -import { GraphState } from "@open-swe/shared/open-swe/types"; import { useRouter } from "next/navigation"; import * as React from "react"; -import { use } from "react"; +import { use, useMemo } from "react"; +import { threadsToMetadata } from "@/lib/thread-utils"; interface ThreadPageProps { thread_id: string; @@ -31,10 +32,28 @@ export default function ThreadPage({ fetchStateHistory: false, }); - const { threads, threadsLoading } = useThreads(MANAGER_GRAPH_ID); + const { threads, isLoading: threadsLoading } = useThreadsSWR({ + assistantId: MANAGER_GRAPH_ID, + }); + + const threadsMetadata = useMemo(() => threadsToMetadata(threads), [threads]); + // Find the thread by ID const thread = threads.find((t) => t.thread_id === thread_id); + // We need a thread object for the hook, so use a dummy if not found + const dummyThread = thread || { + thread_id: thread_id, + values: {}, + status: "idle" as const, + updated_at: new Date().toISOString(), + created_at: new Date().toISOString(), + }; + + const { metadata: currentDisplayThread } = useThreadMetadata( + dummyThread as any, + ); + const handleBackToHome = () => { router.push("/chat"); }; @@ -43,16 +62,12 @@ export default function ThreadPage({ return ; } - // Convert all threads to display format - const displayThreads: ThreadDisplayInfo[] = threads.map(threadToDisplayInfo); - const currentDisplayThread = threadToDisplayInfo(thread); - return (
diff --git a/apps/web/src/app/(v2)/chat/page.tsx b/apps/web/src/app/(v2)/chat/page.tsx index 80a77b89..75ed461a 100644 --- a/apps/web/src/app/(v2)/chat/page.tsx +++ b/apps/web/src/app/(v2)/chat/page.tsx @@ -1,33 +1,28 @@ "use client"; import { DefaultView } from "@/components/v2/default-view"; -import { ThreadDisplayInfo, threadToDisplayInfo } from "@/components/v2/types"; import { useThreadsSWR } from "@/hooks/useThreadsSWR"; import { GitHubAppProvider } from "@/providers/GitHubApp"; -import { GraphState } from "@open-swe/shared/open-swe/types"; import { Toaster } from "@/components/ui/sonner"; import { Suspense } from "react"; import { MANAGER_GRAPH_ID } from "@open-swe/shared/constants"; export default function ChatPage() { - const { threads, isLoading: threadsLoading } = useThreadsSWR({ + const { threads, isLoading: threadsLoading } = useThreadsSWR({ assistantId: MANAGER_GRAPH_ID, - refreshInterval: 15000, // Poll every 15 seconds }); + if (!threads) { return
No threads
; } - // Convert Thread objects to ThreadDisplayInfo for UI - const displayThreads: ThreadDisplayInfo[] = threads.map(threadToDisplayInfo); - return (
diff --git a/apps/web/src/app/(v2)/chat/threads/page.tsx b/apps/web/src/app/(v2)/chat/threads/page.tsx index d907b112..6c1f6d0d 100644 --- a/apps/web/src/app/(v2)/chat/threads/page.tsx +++ b/apps/web/src/app/(v2)/chat/threads/page.tsx @@ -1,60 +1,83 @@ "use client"; -import type React from "react"; -import { useState, Suspense } from "react"; +import { useState, Suspense, useMemo } from "react"; import { Button } from "@/components/ui/button"; import { Badge } from "@/components/ui/badge"; import { Input } from "@/components/ui/input"; import { ArrowLeft, Search, Filter } from "lucide-react"; import { useRouter } from "next/navigation"; -import { ThreadDisplayInfo, threadToDisplayInfo } from "@/components/v2/types"; +import { ThreadMetadata } from "@/components/v2/types"; import { useThreadsSWR } from "@/hooks/useThreadsSWR"; -import { GraphState } from "@open-swe/shared/open-swe/types"; import { ThreadCard, ThreadCardLoading } from "@/components/v2/thread-card"; import { ThemeToggle } from "@/components/theme-toggle"; import { InstallationSelector } from "@/components/github/installation-selector"; import { GitHubAppProvider } from "@/providers/GitHubApp"; import { MANAGER_GRAPH_ID } from "@open-swe/shared/constants"; +import { useThreadsStatus } from "@/hooks/useThreadsStatus"; import { cn } from "@/lib/utils"; +import { threadsToMetadata } from "@/lib/thread-utils"; -type FilterStatus = "all" | "running" | "completed" | "failed" | "pending"; +type FilterStatus = + | "all" + | "running" + | "completed" + | "failed" + | "pending" + | "idle" + | "paused" + | "error"; function AllThreadsPageContent() { const router = useRouter(); - const { threads, isLoading: threadsLoading } = useThreadsSWR({ + const { threads, isLoading: threadsLoading } = useThreadsSWR({ assistantId: MANAGER_GRAPH_ID, - refreshInterval: 15000, // Poll every 15 seconds }); const [searchQuery, setSearchQuery] = useState(""); const [statusFilter, setStatusFilter] = useState("all"); - // Convert Thread objects to ThreadDisplayInfo for UI - const displayThreads: ThreadDisplayInfo[] = threads.map(threadToDisplayInfo); + const threadsMetadata = useMemo(() => threadsToMetadata(threads), [threads]); - // Filter and search threads - const filteredThreads = displayThreads.filter((thread) => { + const threadIds = threadsMetadata.map((thread) => thread.id); + + const { + statusMap, + statusCounts, + isLoading: statusLoading, + } = useThreadsStatus(threadIds, threads); + + const filteredThreads = threadsMetadata.filter((thread: ThreadMetadata) => { const matchesSearch = thread.title.toLowerCase().includes(searchQuery.toLowerCase()) || thread.repository.toLowerCase().includes(searchQuery.toLowerCase()); + const matchesStatus = - statusFilter === "all" || thread.status === statusFilter; + statusFilter === "all" || statusMap[thread.id] === statusFilter; + return matchesSearch && matchesStatus; }); - // Group threads by status const groupedThreads = { - running: filteredThreads.filter((t) => t.status === "running"), - completed: filteredThreads.filter((t) => t.status === "completed"), - failed: filteredThreads.filter((t) => t.status === "failed"), - pending: filteredThreads.filter((t) => t.status === "pending"), - }; - - const statusCounts = { - all: displayThreads.length, - running: displayThreads.filter((t) => t.status === "running").length, - completed: displayThreads.filter((t) => t.status === "completed").length, - failed: displayThreads.filter((t) => t.status === "failed").length, - pending: displayThreads.filter((t) => t.status === "pending").length, + running: filteredThreads.filter( + (thread: ThreadMetadata) => statusMap[thread.id] === "running", + ), + completed: filteredThreads.filter( + (thread: ThreadMetadata) => statusMap[thread.id] === "completed", + ), + failed: filteredThreads.filter( + (thread: ThreadMetadata) => statusMap[thread.id] === "failed", + ), + pending: filteredThreads.filter( + (thread: ThreadMetadata) => statusMap[thread.id] === "pending", + ), + idle: filteredThreads.filter( + (thread: ThreadMetadata) => statusMap[thread.id] === "idle", + ), + paused: filteredThreads.filter( + (thread: ThreadMetadata) => statusMap[thread.id] === "paused", + ), + error: filteredThreads.filter( + (thread: ThreadMetadata) => statusMap[thread.id] === "error", + ), }; return ( @@ -115,6 +138,9 @@ function AllThreadsPageContent() { "completed", "failed", "pending", + "idle", + "paused", + "error", ] as FilterStatus[] ).map((status) => (
diff --git a/apps/web/src/app/api/[..._path]/utils.ts b/apps/web/src/app/api/[..._path]/utils.ts index c14d2ce1..0bdbbb2f 100644 --- a/apps/web/src/app/api/[..._path]/utils.ts +++ b/apps/web/src/app/api/[..._path]/utils.ts @@ -3,7 +3,6 @@ import { App } from "@octokit/app"; import { GITHUB_TOKEN_COOKIE } from "@open-swe/shared/constants"; import { encryptSecret } from "@open-swe/shared/crypto"; import { NextRequest } from "next/server"; -import { validate } from "uuid"; export function getGitHubAccessTokenOrThrow( req: NextRequest, @@ -64,58 +63,6 @@ async function getInstallationName(installationId: string) { return installationName ?? ""; } -function isNewRunRequest(reqUrlStr: string, reqMethod: string) { - try { - const reqPathnameParts = new URL(reqUrlStr).pathname.split("/"); - const isCreateNewRunReq = - reqPathnameParts?.[1] === "api" && - reqPathnameParts?.[2] === "threads" && - validate(reqPathnameParts?.[3]) && - reqPathnameParts?.[4] === "runs" && - reqMethod.toLowerCase() === "post"; - const isStreamRunReq = - reqPathnameParts?.[1] === "api" && - reqPathnameParts?.[2] === "threads" && - validate(reqPathnameParts?.[3]) && - reqPathnameParts?.[4] === "runs" && - validate(reqPathnameParts?.[5]) && - reqPathnameParts?.[6]?.startsWith("stream") && - reqMethod.toLowerCase() === "get"; - return isCreateNewRunReq || isStreamRunReq; - } catch { - return false; - } -} - -function isGetStateRequest(reqUrlStr: string, reqMethod: string) { - try { - const reqPathnameParts = new URL(reqUrlStr).pathname.split("/"); - const isGetStateReq = - reqPathnameParts?.[1] === "api" && - reqPathnameParts?.[2] === "threads" && - validate(reqPathnameParts?.[3]) && - reqPathnameParts?.[4] === "state" && - reqMethod.toLowerCase() === "get"; - return isGetStateReq; - } catch { - return false; - } -} - -function isSearchThreadsRequest(reqUrlStr: string, reqMethod: string) { - try { - const reqPathnameParts = new URL(reqUrlStr).pathname.split("/"); - const isGetStateReq = - reqPathnameParts?.[1] === "api" && - reqPathnameParts?.[2] === "threads" && - reqPathnameParts?.[3] === "search" && - reqMethod.toLowerCase() === "post"; - return isGetStateReq; - } catch { - return false; - } -} - export async function getInstallationNameFromReq( req: Request, installationId: string, @@ -131,14 +78,7 @@ export async function getInstallationNameFromReq( } try { - if ( - isNewRunRequest(req.url, req.method) || - isGetStateRequest(req.url, req.method) || - isSearchThreadsRequest(req.url, req.method) - ) { - return await getInstallationName(installationId); - } - return ""; + return await getInstallationName(installationId); } catch { return ""; } diff --git a/apps/web/src/components/v2/default-view.tsx b/apps/web/src/components/v2/default-view.tsx index c43e3ae4..e337cce7 100644 --- a/apps/web/src/components/v2/default-view.tsx +++ b/apps/web/src/components/v2/default-view.tsx @@ -3,7 +3,7 @@ import { Button } from "@/components/ui/button"; import { Card, CardContent } from "@/components/ui/card"; import { FilePlus2, Archive, Zap } from "lucide-react"; import { useRouter } from "next/navigation"; -import { ThreadDisplayInfo } from "./types"; +import { ThreadMetadata } from "./types"; import { TerminalInput } from "./terminal-input"; import { useFileUpload } from "@/hooks/useFileUpload"; import { cn } from "@/lib/utils"; @@ -19,12 +19,17 @@ import { ThemeToggle } from "../theme-toggle"; import { ThreadCard, ThreadCardLoading } from "./thread-card"; import { GitHubInstallationBanner } from "../github/installation-banner"; import { QuickActions } from "./quick-actions"; -import { useState } from "react"; import { DraftsSection } from "./drafts-section"; import { GitHubLogoutButton } from "../github/github-oauth-button"; import { MANAGER_GRAPH_ID } from "@open-swe/shared/constants"; import { TooltipIconButton } from "../ui/tooltip-icon-button"; import { InstallationSelector } from "../github/installation-selector"; + +import { useThreadsStatus } from "@/hooks/useThreadsStatus"; +import { Thread } from "@langchain/langgraph-sdk"; +import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; +import { useState, useMemo } from "react"; +import { threadsToMetadata } from "@/lib/thread-utils"; import { Settings } from "lucide-react"; import NextLink from "next/link"; @@ -44,7 +49,7 @@ function OpenSettingsButton() { } interface DefaultViewProps { - threads: ThreadDisplayInfo[]; + threads: Thread[]; threadsLoading: boolean; } @@ -65,6 +70,15 @@ export function DefaultView({ threads, threadsLoading }: DefaultViewProps) { } = useFileUpload(); const [autoAccept, setAutoAccept] = useState(false); + const threadsMetadata = useMemo(() => threadsToMetadata(threads), [threads]); + const displayThreads = threadsMetadata.slice(0, 4); + const displayThreadIds = displayThreads.map((thread) => thread.id); + + const { statusMap, isLoading: statusLoading } = useThreadsStatus( + displayThreadIds, + threads, + ); + const handleLoadDraft = (content: string) => { setDraftToLoad(content); }; @@ -198,10 +212,12 @@ export function DefaultView({ threads, threadsLoading }: DefaultViewProps) { )} - {threads.slice(0, 4).map((thread) => ( + {displayThreads.map((thread) => ( ))} diff --git a/apps/web/src/components/v2/thread-card.tsx b/apps/web/src/components/v2/thread-card.tsx index a4951ad8..ab65ce18 100644 --- a/apps/web/src/components/v2/thread-card.tsx +++ b/apps/web/src/components/v2/thread-card.tsx @@ -4,25 +4,42 @@ import { GitBranch, GitPullRequest, Loader2, + AlertCircle, + Pause, XCircle, + Clock, } from "lucide-react"; import { Card, CardContent, CardHeader, CardTitle } from "../ui/card"; -import { ThreadDisplayInfo } from "./types"; import { useRouter } from "next/navigation"; import { Badge } from "../ui/badge"; import { Button } from "../ui/button"; import { Skeleton } from "../ui/skeleton"; +import { ThreadMetadata } from "./types"; +import { ThreadUIStatus } from "@/lib/schemas/thread-status"; import { cn } from "@/lib/utils"; -export function ThreadCard({ thread }: { thread: ThreadDisplayInfo }) { +interface ThreadCardProps { + thread: ThreadMetadata; + status?: ThreadUIStatus; + statusLoading?: boolean; +} + +export function ThreadCard({ thread, status, statusLoading }: ThreadCardProps) { const router = useRouter(); - const getStatusColor = (status: ThreadDisplayInfo["status"]) => { + const isStatusLoading = statusLoading && !status; + const displayStatus = status || ("idle" as ThreadUIStatus); + + const getStatusColor = (status: ThreadUIStatus) => { switch (status) { case "running": return "dark:bg-blue-950 bg-blue-100 dark:text-blue-400 text-blue-700"; case "completed": return "dark:bg-green-950 bg-green-100 dark:text-green-400 text-green-700"; + case "error": + return "dark:bg-red-950 bg-red-100 dark:text-red-400 text-red-700"; + case "paused": + return "dark:bg-yellow-950 bg-yellow-100 dark:text-yellow-400 text-yellow-700"; case "failed": return "dark:bg-red-950 bg-red-100 dark:text-red-400 text-red-700"; case "pending": @@ -32,12 +49,18 @@ export function ThreadCard({ thread }: { thread: ThreadDisplayInfo }) { } }; - const getStatusIcon = (status: ThreadDisplayInfo["status"]) => { + const getStatusIcon = (status: ThreadUIStatus) => { switch (status) { case "running": return ; case "completed": return ; + case "error": + return ; + case "idle": + return ; + case "paused": + return ; case "failed": return ; default: @@ -83,11 +106,22 @@ export function ThreadCard({ thread }: { thread: ThreadDisplayInfo }) {
- {getStatusIcon(thread.status)} - {thread.status} + {isStatusLoading ? ( + + ) : ( + getStatusIcon(displayStatus) + )} + + {isStatusLoading ? "Loading..." : displayStatus} +
diff --git a/apps/web/src/components/v2/thread-switcher.tsx b/apps/web/src/components/v2/thread-switcher.tsx index 31ea8aa2..6239bdb9 100644 --- a/apps/web/src/components/v2/thread-switcher.tsx +++ b/apps/web/src/components/v2/thread-switcher.tsx @@ -21,12 +21,12 @@ import { Bug, } from "lucide-react"; import { useRouter } from "next/navigation"; -import { ThreadDisplayInfo } from "./types"; +import { ThreadMetadata } from "./types"; import { ThreadCard } from "./thread-card"; interface ThreadSwitcherProps { - currentThread: ThreadDisplayInfo; - allThreads: ThreadDisplayInfo[]; + currentThread: ThreadMetadata; + allThreads: ThreadMetadata[]; } export function ThreadSwitcher({ @@ -37,9 +37,6 @@ export function ThreadSwitcher({ const router = useRouter(); const otherThreads = allThreads.filter((t) => t.id !== currentThread.id); - const runningCount = otherThreads.filter( - (t) => t.status === "running", - ).length; return ( Switch Thread - {runningCount > 0 && ( - - {runningCount} - - )} >; - displayThread: ThreadDisplayInfo; - allDisplayThreads: ThreadDisplayInfo[]; + displayThread: ThreadMetadata; + allDisplayThreads: ThreadMetadata[]; onBackToHome: () => void; } @@ -51,6 +53,23 @@ export function ThreadView({ const [programmerSession, setProgrammerSession] = useState(); + const { status: realTimeStatus } = useThreadStatus(displayThread.id); + + const getStatusDotColor = (status: string) => { + switch (status) { + case "running": + return "bg-blue-500 dark:bg-blue-400"; + case "completed": + return "bg-green-500 dark:bg-green-400"; + case "paused": + return "bg-yellow-500 dark:bg-yellow-400"; + case "error": + return "bg-red-500 dark:bg-red-400"; + default: + return "bg-gray-500 dark:bg-gray-400"; + } + }; + const plannerCancelRef = useRef<(() => void) | null>(null); const programmerCancelRef = useRef<(() => void) | null>(null); @@ -102,11 +121,7 @@ export function ThreadView({
diff --git a/apps/web/src/components/v2/types.ts b/apps/web/src/components/v2/types.ts index da6305e0..16fd122f 100644 --- a/apps/web/src/components/v2/types.ts +++ b/apps/web/src/components/v2/types.ts @@ -1,16 +1,15 @@ -import { getThreadTitle } from "@/lib/thread"; -import { Thread } from "@langchain/langgraph-sdk"; -import { getActivePlanItems } from "@open-swe/shared/open-swe/tasks"; -import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; +import { ThreadUIStatus } from "@/lib/schemas/thread-status"; +import { TaskPlan } from "@open-swe/shared/open-swe/types"; -export interface ThreadDisplayInfo { +export interface ThreadMetadata { id: string; title: string; - status: "running" | "completed" | "failed" | "pending"; lastActivity: string; taskCount: number; repository: string; branch: string; + taskPlan?: TaskPlan; + status: ThreadUIStatus; githubIssue?: { number: number; url: string; @@ -21,71 +20,3 @@ export interface ThreadDisplayInfo { status: "draft" | "open" | "merged" | "closed"; }; } - -// Utility functions to convert between Thread and ThreadDisplayInfo -export function threadToDisplayInfo( - thread: Thread, -): ThreadDisplayInfo { - const values = thread.values; - const activePlanItems = values?.taskPlan - ? getActivePlanItems(values.taskPlan) - : []; - const completedTasksLen = activePlanItems.filter((t) => t.completed).length; - - // Determine UI status from thread status and task completion - let uiStatus: ThreadDisplayInfo["status"]; - switch (thread.status) { - case "busy": - uiStatus = "running"; - break; - case "idle": - uiStatus = - completedTasksLen === activePlanItems.length ? "completed" : "pending"; - break; - case "error": - uiStatus = "failed"; - break; - case "interrupted": - uiStatus = "pending"; - break; - default: - uiStatus = "pending"; - } - - // Calculate time since last update - const lastUpdate = new Date(thread.updated_at); - const now = new Date(); - const diffMs = now.getTime() - lastUpdate.getTime(); - const diffMins = Math.floor(diffMs / (1000 * 60)); - const diffHours = Math.floor(diffMs / (1000 * 60 * 60)); - const diffDays = Math.floor(diffMs / (1000 * 60 * 60 * 24)); - - let lastActivity: string; - if (diffMins < 1) { - lastActivity = "just now"; - } else if (diffMins < 60) { - lastActivity = `${diffMins} min ago`; - } else if (diffHours < 24) { - lastActivity = `${diffHours} hour${diffHours > 1 ? "s" : ""} ago`; - } else { - lastActivity = `${diffDays} day${diffDays > 1 ? "s" : ""} ago`; - } - - return { - id: thread.thread_id, - title: getThreadTitle(thread), - status: uiStatus, - lastActivity, - taskCount: values?.taskPlan?.tasks.length ?? 0, - repository: values?.targetRepository - ? `${values.targetRepository.owner}/${values.targetRepository.repo}` - : "", - branch: values?.targetRepository.branch || "main", - githubIssue: values?.githubIssueId - ? { - number: values?.githubIssueId, - url: `https://github.com/${values?.targetRepository.owner}/${values?.targetRepository.repo}/issues/${values?.githubIssueId}`, - } - : undefined, - }; -} diff --git a/apps/web/src/hooks/useThreadMetadata.ts b/apps/web/src/hooks/useThreadMetadata.ts new file mode 100644 index 00000000..34a06905 --- /dev/null +++ b/apps/web/src/hooks/useThreadMetadata.ts @@ -0,0 +1,51 @@ +import { Thread } from "@langchain/langgraph-sdk"; +import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; +import { ThreadMetadata } from "@/components/v2/types"; +import { useThreadStatus } from "./useThreadStatus"; +import { useMemo } from "react"; +import { getThreadTitle } from "@/lib/thread"; +import { calculateLastActivity } from "@/lib/thread-utils"; + +/** + * Hook that combines thread metadata with real-time status + */ +export function useThreadMetadata(thread: Thread): { + metadata: ThreadMetadata; + isStatusLoading: boolean; + statusError: Error | null; +} { + const { + status, + isLoading: isStatusLoading, + error: statusError, + } = useThreadStatus(thread.thread_id); + + const metadata: ThreadMetadata = useMemo((): ThreadMetadata => { + const values = thread.values; + + return { + id: thread.thread_id, + title: getThreadTitle(thread), + lastActivity: calculateLastActivity(thread.updated_at), + taskCount: values?.taskPlan?.tasks.length ?? 0, + repository: values?.targetRepository + ? `${values.targetRepository.owner}/${values.targetRepository.repo}` + : "", + branch: values?.targetRepository?.branch || "main", + taskPlan: values?.taskPlan, + status, + githubIssue: values?.githubIssueId + ? { + number: values?.githubIssueId, + url: `https://github.com/${values?.targetRepository?.owner}/${values?.targetRepository?.repo}/issues/${values?.githubIssueId}`, + } + : undefined, + }; + }, [thread, status]); + + return { + metadata, + isStatusLoading, + statusError, + }; +} diff --git a/apps/web/src/hooks/useThreadPolling.ts b/apps/web/src/hooks/useThreadPolling.ts deleted file mode 100644 index d5007884..00000000 --- a/apps/web/src/hooks/useThreadPolling.ts +++ /dev/null @@ -1,48 +0,0 @@ -import { useEffect, useRef } from "react"; -import { ThreadPoller, PollConfig } from "@/lib/polling/thread-poller"; -import { GraphState } from "@open-swe/shared/open-swe/types"; -import { Thread } from "@langchain/langgraph-sdk"; - -interface UseThreadPollingProps { - threads: Thread[]; - getThread: (threadId: string) => Promise | null>; - onUpdate: ( - updatedThreads: Thread[], - changedThreadIds: string[], - ) => void; - - enabled?: boolean; -} - -export function useThreadPolling({ - threads, - getThread, - onUpdate, - enabled = true, -}: UseThreadPollingProps) { - const pollerRef = useRef(null); - - useEffect(() => { - if (!enabled) return; - - const config: PollConfig = { - interval: 15000, - onUpdate, - }; - - pollerRef.current = new ThreadPoller(config, threads, getThread); - pollerRef.current.start(); - - return () => { - if (pollerRef.current) { - pollerRef.current.stop(); - pollerRef.current = null; - } - }; - }, [threads, getThread, onUpdate, enabled]); - - return { - start: () => pollerRef.current?.start(), - stop: () => pollerRef.current?.stop(), - }; -} diff --git a/apps/web/src/hooks/useThreadStatus.ts b/apps/web/src/hooks/useThreadStatus.ts new file mode 100644 index 00000000..e4c62c5a --- /dev/null +++ b/apps/web/src/hooks/useThreadStatus.ts @@ -0,0 +1,48 @@ +import useSWR from "swr"; +import { THREAD_STATUS_SWR_CONFIG } from "@/lib/swr-config"; +import { ThreadUIStatus, ThreadStatusData } from "@/lib/schemas/thread-status"; +import { fetchThreadStatus } from "@/services/thread-status.service"; + +interface UseThreadStatusOptions { + enabled?: boolean; + refreshInterval?: number; +} + +interface ThreadStatusResult { + status: ThreadUIStatus; + isLoading: boolean; + error: Error | null; + mutate: () => void; +} + +/** + * Thread status hook using SWR for real-time status updates + * Uses SWR caching directly instead of manual Zustand cache + */ +export function useThreadStatus( + threadId: string, + options: UseThreadStatusOptions = {}, +): ThreadStatusResult { + const { + enabled = true, + refreshInterval = THREAD_STATUS_SWR_CONFIG.refreshInterval, + } = options; + + const swrKey = enabled ? `thread-status-${threadId}` : null; + + const { data, error, isLoading, mutate } = useSWR( + swrKey, + () => fetchThreadStatus(threadId), + { + ...THREAD_STATUS_SWR_CONFIG, + refreshInterval, + }, + ); + + return { + status: data?.status || "idle", + isLoading, + error, + mutate, + }; +} diff --git a/apps/web/src/hooks/useThreads.tsx b/apps/web/src/hooks/useThreads.tsx deleted file mode 100644 index 035c0e95..00000000 --- a/apps/web/src/hooks/useThreads.tsx +++ /dev/null @@ -1,58 +0,0 @@ -import { createClient } from "@/providers/client"; -import { Thread } from "@langchain/langgraph-sdk"; -import { useCallback, useEffect, useState } from "react"; - -export function useThreads>( - assistantId?: string, -) { - const apiUrl: string | undefined = process.env.NEXT_PUBLIC_API_URL ?? ""; - const [threads, setThreads] = useState[]>([]); - const [threadsLoading, setThreadsLoading] = useState(false); - - const getThread = useCallback( - async (threadId: string): Promise | null> => { - if (!apiUrl) return null; - const client = createClient(apiUrl); - - try { - const thread = await client.threads.get(threadId); - return thread; - } catch (error) { - console.error("Failed to fetch thread:", threadId, error); - return null; - } - }, - [apiUrl], - ); - - const getThreads = useCallback(async (): Promise[] | null> => { - if (!apiUrl) return null; - setThreadsLoading(true); - const client = createClient(apiUrl); - - try { - const searchArgs = assistantId - ? { - metadata: { - graph_id: assistantId, - }, - } - : undefined; - const threads = await client.threads.search(searchArgs); - return threads; - } catch (error) { - console.error("Failed to fetch threads:", error); - return null; - } finally { - setThreadsLoading(false); - } - }, [apiUrl]); - - useEffect(() => { - getThreads().then((threads) => { - setThreads(threads ?? []); - }); - }, [getThreads]); - - return { threads, setThreads, getThread, getThreads, threadsLoading }; -} diff --git a/apps/web/src/hooks/useThreadsSWR.ts b/apps/web/src/hooks/useThreadsSWR.ts index 66fdadea..d78eba8b 100644 --- a/apps/web/src/hooks/useThreadsSWR.ts +++ b/apps/web/src/hooks/useThreadsSWR.ts @@ -1,22 +1,44 @@ -import { createClient } from "@/providers/client"; -import { Thread } from "@langchain/langgraph-sdk"; import useSWR from "swr"; +import { Thread } from "@langchain/langgraph-sdk"; +import { createClient } from "@/providers/client"; +import { THREAD_SWR_CONFIG } from "@/lib/swr-config"; +import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; +import { PlannerGraphState } from "@open-swe/shared/open-swe/planner/types"; +import { ReviewerGraphState } from "@open-swe/shared/open-swe/reviewer/types"; +import { GraphState } from "@open-swe/shared/open-swe/types"; -export interface UseThreadsSWROptions { +/** + * Union type representing all possible graph states in the Open SWE system + */ +export type AnyGraphState = + | ManagerGraphState + | PlannerGraphState + | ReviewerGraphState + | GraphState; + +interface UseThreadsSWROptions { assistantId?: string; refreshInterval?: number; revalidateOnFocus?: boolean; revalidateOnReconnect?: boolean; } -export function useThreadsSWR>( - options: UseThreadsSWROptions = {}, -) { +/** + * Hook for fetching threads for any graph type. + * Works with all graph states (Manager, Planner, Programmer, Reviewer) + * by passing the appropriate assistantId. + * + * For UI display of manager threads, use `threadsToMetadata(threads)` utility to convert + * raw threads to ThreadMetadata objects. + */ +export function useThreadsSWR< + TGraphState extends AnyGraphState = AnyGraphState, +>(options: UseThreadsSWROptions = {}) { const { assistantId, - refreshInterval = 0, // Default to no polling, can be overridden - revalidateOnFocus = true, - revalidateOnReconnect = true, + refreshInterval = THREAD_SWR_CONFIG.refreshInterval, + revalidateOnFocus = THREAD_SWR_CONFIG.revalidateOnFocus, + revalidateOnReconnect = THREAD_SWR_CONFIG.revalidateOnReconnect, } = options; const apiUrl: string | undefined = process.env.NEXT_PUBLIC_API_URL ?? ""; @@ -24,7 +46,7 @@ export function useThreadsSWR>( // Create a unique key for SWR caching based on assistantId const swrKey = assistantId ? ["threads", assistantId] : ["threads", "all"]; - const fetcher = async (): Promise[]> => { + const fetcher = async (): Promise[]> => { if (!apiUrl) { throw new Error("API URL is not configured"); } @@ -38,7 +60,7 @@ export function useThreadsSWR>( } : undefined; - return await client.threads.search(searchArgs); + return await client.threads.search(searchArgs); }; const { data, error, isLoading, mutate, isValidating } = useSWR( @@ -48,8 +70,9 @@ export function useThreadsSWR>( refreshInterval, revalidateOnFocus, revalidateOnReconnect, - errorRetryCount: 3, - errorRetryInterval: 5000, + errorRetryCount: THREAD_SWR_CONFIG.errorRetryCount, + errorRetryInterval: THREAD_SWR_CONFIG.errorRetryInterval, + dedupingInterval: THREAD_SWR_CONFIG.dedupingInterval, }, ); diff --git a/apps/web/src/hooks/useThreadsStatus.ts b/apps/web/src/hooks/useThreadsStatus.ts new file mode 100644 index 00000000..abc94ec6 --- /dev/null +++ b/apps/web/src/hooks/useThreadsStatus.ts @@ -0,0 +1,192 @@ +import useSWR from "swr"; +import { ThreadUIStatus, ThreadStatusData } from "@/lib/schemas/thread-status"; +import { fetchThreadStatus } from "@/services/thread-status.service"; +import { THREAD_STATUS_SWR_CONFIG } from "@/lib/swr-config"; +import { useMemo, useRef } from "react"; +import { Thread } from "@langchain/langgraph-sdk"; +import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; +import { PlannerGraphState } from "@open-swe/shared/open-swe/planner/types"; +import { GraphState } from "@open-swe/shared/open-swe/types"; + +export interface SessionCacheData { + plannerData?: { thread: Thread }; + programmerData?: { thread: Thread }; + timestamp: number; +} + +export type SessionCache = Map; + +interface ThreadStatusMap { + [threadId: string]: ThreadUIStatus; +} + +interface ThreadStatusCounts { + all: number; + running: number; + completed: number; + failed: number; + pending: number; + idle: number; + paused: number; + error: number; +} + +interface GroupedThreadIds { + running: string[]; + completed: string[]; + failed: string[]; + pending: string[]; + idle: string[]; + paused: string[]; + error: string[]; +} + +interface UseThreadsStatusResult { + statusMap: ThreadStatusMap; + statusCounts: ThreadStatusCounts; + groupedThreads: GroupedThreadIds; + isLoading: boolean; + hasErrors: boolean; +} + +const sessionDataCache: SessionCache = new Map(); + +const CACHE_TTL = 30000; + +/** + * Fetches statuses for multiple threads in parallel + * Uses session caching to achieve "single request per thread + cache sessions" goal + */ +async function fetchAllThreadStatuses( + threadIds: string[], + lastPollingStates: Map, + managerThreads?: Thread[], +): Promise<{ + statusMap: ThreadStatusMap; + updatedStates: Map; +}> { + const statusPromises = threadIds.map(async (threadId) => { + try { + const lastState = lastPollingStates.get(threadId) || null; + + const managerThread = managerThreads?.find( + (t) => t.thread_id === threadId, + ); + + const statusData = await fetchThreadStatus( + threadId, + lastState, + managerThread, + sessionDataCache, + ); + return { threadId, status: statusData.status, statusData }; + } catch (error) { + console.error(`Failed to fetch status for thread ${threadId}:`, error); + return { + threadId, + status: "idle" as ThreadUIStatus, + statusData: null, + }; + } + }); + + const results = await Promise.all(statusPromises); + const statusMap: ThreadStatusMap = {}; + const updatedStates = new Map(); + + results.forEach(({ threadId, status, statusData }) => { + statusMap[threadId] = status; + if (statusData) { + updatedStates.set(threadId, statusData); + } + }); + + return { statusMap, updatedStates }; +} + +/** + * Hook that fetches statuses for multiple threads in parallel + * Uses SWR for caching and deduplication with state optimization + */ +export function useThreadsStatus( + threadIds: string[], + managerThreads?: Thread[], +): UseThreadsStatusResult { + const lastPollingStatesRef = useRef>(new Map()); + + // Create a stable key for the thread IDs array + const sortedThreadIds = threadIds.sort(); + const threadIdsKey = sortedThreadIds.join(","); + + const swrKey = + threadIds.length > 0 + ? threadIds.length <= 4 + ? `threads-status-batch-${threadIds.length}-${threadIdsKey}` + : `threads-status-${threadIdsKey}` + : null; + + const { + data: fetchResult, + isLoading, + error, + } = useSWR( + swrKey, + async () => { + if (process.env.NODE_ENV === "development") { + console.log( + `[Status SWR] Fetching statuses for ${threadIds.length} threads`, + ); + } + + const result = await fetchAllThreadStatuses( + sortedThreadIds, // Use sorted array for consistency + lastPollingStatesRef.current, + managerThreads, + ); + lastPollingStatesRef.current = result.updatedStates; + return result; + }, + THREAD_STATUS_SWR_CONFIG, + ); + + const statusMap = fetchResult?.statusMap || {}; + + return useMemo(() => { + const groupedThreads: GroupedThreadIds = { + running: [], + completed: [], + failed: [], + pending: [], + idle: [], + paused: [], + error: [], + }; + + if (statusMap) { + Object.entries(statusMap).forEach(([threadId, status]) => { + if (groupedThreads[status]) { + groupedThreads[status].push(threadId); + } + }); + } + + const statusCounts: ThreadStatusCounts = { + all: threadIds.length, + running: groupedThreads.running.length, + completed: groupedThreads.completed.length, + failed: groupedThreads.failed.length, + pending: groupedThreads.pending.length, + idle: groupedThreads.idle.length, + paused: groupedThreads.paused.length, + error: groupedThreads.error.length, + }; + + return { + statusMap: statusMap || {}, + statusCounts, + groupedThreads, + isLoading, + hasErrors: !!error, + }; + }, [statusMap, threadIds, threadIdsKey, isLoading, error]); +} diff --git a/apps/web/src/lib/polling/thread-poller.ts b/apps/web/src/lib/polling/thread-poller.ts deleted file mode 100644 index 03b825a1..00000000 --- a/apps/web/src/lib/polling/thread-poller.ts +++ /dev/null @@ -1,107 +0,0 @@ -import { Thread } from "@langchain/langgraph-sdk"; -import { GraphState } from "@open-swe/shared/open-swe/types"; -import { getThreadTasks, getThreadTitle } from "../thread"; - -export interface PollConfig { - interval: number; - onUpdate: ( - updatedThreads: Thread[], - changedThreadIds: string[], - ) => void; -} - -export class ThreadPoller { - private config: PollConfig; - private isPolling: boolean = false; - private intervalId: NodeJS.Timeout | null = null; - private threads: Thread[]; - private getThreadFn: (threadId: string) => Promise | null>; - - constructor( - config: PollConfig, - threads: Thread[], - getThreadFn: (threadId: string) => Promise | null>, - ) { - this.config = config; - this.threads = threads; - this.getThreadFn = getThreadFn; - } - - start(): void { - if (this.isPolling) return; - - this.isPolling = true; - this.intervalId = setInterval(() => { - this.pollThreads(); - }, this.config.interval); - } - - stop(): void { - if (!this.isPolling) return; - - this.isPolling = false; - if (this.intervalId) { - clearInterval(this.intervalId); - this.intervalId = null; - } - } - - private async pollThreads(): Promise { - try { - const currentThreads = this.threads; - - const threadsToPool = currentThreads.slice(0, 10); - const updatedThreads: Thread[] = []; - const changedThreadIds: string[] = []; - const errors: string[] = []; - - const pollUpdatePromise = Promise.allSettled( - threadsToPool.map(async (currentThread) => { - try { - const updatedThread = await this.getThreadFn( - currentThread.thread_id, - ); - if (updatedThread) { - updatedThreads.push(updatedThread); - - if (this.hasThreadChanged(currentThread, updatedThread)) { - changedThreadIds.push(updatedThread.thread_id); - } - } - } catch (error) { - errors.push(`Thread ${currentThread.thread_id}: ${error}`); - updatedThreads.push(currentThread); - } - }), - ); - - await pollUpdatePromise; - - if (changedThreadIds.length > 0) { - this.config.onUpdate(updatedThreads, changedThreadIds); - } - } catch (error) { - console.error("Thread polling error:", error); - } - } - - private hasThreadChanged( - current: Thread, - updated: Thread, - ): boolean { - const currentTaskCounts = getThreadTasks(current); - const updatedTaskCounts = getThreadTasks(updated); - const currentTargetRepo = current.values?.targetRepository; - const updatedTargetRepo = updated.values?.targetRepository; - return ( - currentTaskCounts.completedTasks !== updatedTaskCounts.completedTasks || - currentTaskCounts.totalTasks !== updatedTaskCounts.totalTasks || - current.status !== updated.status || - getThreadTitle(current) !== getThreadTitle(updated) || - currentTargetRepo.repo !== updatedTargetRepo.repo || - currentTargetRepo.branch !== updatedTargetRepo.branch || - JSON.stringify(current.values?.taskPlan) !== - JSON.stringify(updated.values?.taskPlan) - ); - } -} diff --git a/apps/web/src/lib/schemas/thread-status.ts b/apps/web/src/lib/schemas/thread-status.ts new file mode 100644 index 00000000..4ac80e8f --- /dev/null +++ b/apps/web/src/lib/schemas/thread-status.ts @@ -0,0 +1,45 @@ +import { + MANAGER_GRAPH_ID, + PLANNER_GRAPH_ID, + PROGRAMMER_GRAPH_ID, +} from "@open-swe/shared/constants"; +import { ThreadStatus } from "@langchain/langgraph-sdk"; +import { TaskPlan } from "@open-swe/shared/open-swe/types"; + +/** + * UI-specific thread status that extends LangGraph's states + */ +export type ThreadUIStatus = + | "running" // Maps from LangGraph "busy" + | "completed" // Business logic: all tasks completed + | "failed" // UI-specific state + | "pending" // UI-specific state + | "idle" // Same as LangGraph "idle" + | "paused" // Maps from LangGraph "interrupted" + | "error"; // Same as LangGraph "error" + +export function mapLangGraphToUIStatus(status: ThreadStatus): ThreadUIStatus { + switch (status) { + case "busy": + return "running"; + case "interrupted": + return "paused"; + case "idle": + return "idle"; + case "error": + return "error"; + default: + return "idle"; + } +} + +export interface ThreadStatusData { + graph: + | typeof MANAGER_GRAPH_ID + | typeof PLANNER_GRAPH_ID + | typeof PROGRAMMER_GRAPH_ID; + runId: string; + threadId: string; + status: ThreadUIStatus; + taskPlan?: TaskPlan; // Task plan data when available from programmer sessions +} diff --git a/apps/web/src/lib/swr-config.ts b/apps/web/src/lib/swr-config.ts new file mode 100644 index 00000000..e5dfc563 --- /dev/null +++ b/apps/web/src/lib/swr-config.ts @@ -0,0 +1,33 @@ +/** + * Standardized SWR configuration for thread-related hooks + * + * This ensures consistent polling intervals, error handling, and caching behavior + * across all thread data fetching in the application. + */ +export const THREAD_SWR_CONFIG = { + refreshInterval: 15000, + revalidateOnFocus: false, + revalidateOnReconnect: true, + errorRetryCount: 3, + errorRetryInterval: 5000, + dedupingInterval: 2000, +} as const; + +/** + * SWR configuration for thread status polling + * Uses same intervals but with focus revalidation for real-time updates + */ +export const THREAD_STATUS_SWR_CONFIG = { + ...THREAD_SWR_CONFIG, + revalidateOnFocus: true, + dedupingInterval: 5000, +} as const; + +/** + * SWR configuration for one-time fetches (no polling) + * Used for thread data that doesn't need real-time updates + */ +export const THREAD_STATIC_SWR_CONFIG = { + ...THREAD_SWR_CONFIG, + refreshInterval: 0, // No automatic polling +} as const; diff --git a/apps/web/src/lib/thread-utils.ts b/apps/web/src/lib/thread-utils.ts new file mode 100644 index 00000000..c5083fcf --- /dev/null +++ b/apps/web/src/lib/thread-utils.ts @@ -0,0 +1,42 @@ +import { formatDistanceToNow } from "date-fns"; +import { Thread } from "@langchain/langgraph-sdk"; +import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; +import { ThreadMetadata } from "@/components/v2/types"; +import { getThreadTitle } from "./thread"; + +/** + * Calculate human-readable last activity time from thread updated_at timestamp + */ +export function calculateLastActivity(updatedAt: string): string { + return formatDistanceToNow(new Date(updatedAt), { addSuffix: true }); +} + +/** + * Converts raw manager threads to ThreadMetadata objects for UI display + */ +export function threadsToMetadata( + threads: Thread[], +): ThreadMetadata[] { + return threads.map((thread): ThreadMetadata => { + const values = thread.values; + + return { + id: thread.thread_id, + title: getThreadTitle(thread), + lastActivity: calculateLastActivity(thread.updated_at), + taskCount: values?.taskPlan?.tasks.length ?? 0, + repository: values?.targetRepository + ? `${values.targetRepository.owner}/${values.targetRepository.repo}` + : "", + branch: values?.targetRepository?.branch || "main", + taskPlan: values?.taskPlan, + status: "idle" as const, // Default status - consumers can override with real status + githubIssue: values?.githubIssueId + ? { + number: values?.githubIssueId, + url: `https://github.com/${values?.targetRepository?.owner}/${values?.targetRepository?.repo}/issues/${values?.githubIssueId}`, + } + : undefined, + }; + }); +} diff --git a/apps/web/src/providers/Thread.tsx b/apps/web/src/providers/Thread.tsx index 435247ad..420a0c2d 100644 --- a/apps/web/src/providers/Thread.tsx +++ b/apps/web/src/providers/Thread.tsx @@ -12,7 +12,6 @@ import { } from "react"; import { createClient } from "./client"; import { GraphState } from "@open-swe/shared/open-swe/types"; -import { useThreadPolling } from "@/hooks/useThreadPolling"; interface ThreadContextType { threads: Thread[]; @@ -110,32 +109,6 @@ export function ThreadProvider({ children }: { children: ReactNode }) { refreshThreads(); }, [refreshThreads]); - const handlePollingUpdate = useCallback( - (updatedThreads: Thread[], changedThreadIds: string[]) => { - setThreads((currentThreads) => { - const updatedMap = new Map(updatedThreads.map((t) => [t.thread_id, t])); - return currentThreads.map( - (thread) => updatedMap.get(thread.thread_id) || thread, - ); - }); - - setRecentlyUpdatedThreads(new Set(changedThreadIds)); - - setTimeout(() => { - setRecentlyUpdatedThreads(new Set()); - }, 2000); - }, - [], - ); - - // Initialize polling - useThreadPolling({ - threads, - getThread, - onUpdate: handlePollingUpdate, - enabled: true, - }); - const handleThreadClick = useCallback( ( thread: Thread, diff --git a/apps/web/src/services/thread-status.service.ts b/apps/web/src/services/thread-status.service.ts new file mode 100644 index 00000000..9b562dd0 --- /dev/null +++ b/apps/web/src/services/thread-status.service.ts @@ -0,0 +1,426 @@ +import { Client, Thread } from "@langchain/langgraph-sdk"; +import { createClient } from "@/providers/client"; +import { + ThreadUIStatus, + ThreadStatusData, + mapLangGraphToUIStatus, +} from "@/lib/schemas/thread-status"; +import { GraphState, TaskPlan } from "@open-swe/shared/open-swe/types"; +import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types"; +import { PlannerGraphState } from "@open-swe/shared/open-swe/planner/types"; +import { getActivePlanItems } from "@open-swe/shared/open-swe/tasks"; +import { SessionCache, SessionCacheData } from "@/hooks/useThreadsStatus"; + +interface StatusResult { + graph: "manager" | "planner" | "programmer"; + runId: string; + threadId: string; + status: ThreadUIStatus; + taskPlan?: TaskPlan; +} + +/** + * Determines if all tasks in a task plan are completed + */ +function areAllPlanItemsCompleted(taskPlan: TaskPlan): boolean { + if (!taskPlan?.tasks || !Array.isArray(taskPlan.tasks)) { + return false; + } + const activePlanItems = getActivePlanItems(taskPlan); + return activePlanItems.every((planItem) => planItem.completed); +} + +export class StatusResolver { + resolve( + manager: StatusResult, + planner?: StatusResult, + programmer?: StatusResult, + ): ThreadStatusData { + if (manager.status === "running" || manager.status === "error") { + return manager; + } + + if (!planner) { + return manager; + } + + if (planner.status === "running" || planner.status === "paused") { + return planner; + } + + if (planner.status === "error") { + return planner; + } + + if (!programmer) { + return planner; + } + + return programmer; + } +} + +const CACHE_TTL = 30 * 1000; + +function getCachedSessionData( + sessionCache: SessionCache | undefined, + sessionKey: string, +): SessionCacheData | null { + if (!sessionCache) return null; + + const cached = sessionCache.get(sessionKey); + if (!cached) return null; + + const isExpired = Date.now() - cached.timestamp > CACHE_TTL; + if (isExpired) { + sessionCache.delete(sessionKey); + return null; + } + + return cached; +} + +function setCachedSessionData( + sessionCache: SessionCache | undefined, + sessionKey: string, + data: Partial, +): void { + if (!sessionCache) return; + + sessionCache.set(sessionKey, { + ...data, + timestamp: Date.now(), + }); +} + +export async function fetchThreadStatus( + threadId: string, + lastPollingState: ThreadStatusData | null = null, + managerThreadData?: Thread | null, + sessionCache?: SessionCache, +): Promise { + try { + const apiUrl = process.env.NEXT_PUBLIC_API_URL ?? ""; + if (!apiUrl) { + throw new Error("API URL not configured"); + } + + const client = createClient(apiUrl); + const resolver = new StatusResolver(); + + if (lastPollingState) { + try { + const optimizedResult = await checkLastKnownGraph( + client, + lastPollingState, + resolver, + ); + if (optimizedResult) { + return optimizedResult; + } + } catch (error) { + console.warn( + "Optimization check failed, falling back to full status check:", + error, + ); + } + } + + return await performFullStatusCheck( + client, + threadId, + resolver, + managerThreadData, + sessionCache, + ); + } catch (error) { + console.error(`Error fetching thread status for ${threadId}:`, error); + + const graph = lastPollingState?.graph || "manager"; + const runId = lastPollingState?.runId || ""; + const errorThreadId = lastPollingState?.threadId || threadId; + + return { + graph, + runId, + threadId: errorThreadId, + status: "error", + }; + } +} + +async function checkLastKnownGraph( + client: Client, + lastState: ThreadStatusData, + resolver: StatusResolver, +): Promise { + switch (lastState.graph) { + case "programmer": + if (lastState.threadId && lastState.runId) { + const programmerThread = await client.threads.get( + lastState.threadId, + ); + + // Use thread status directly for most cases + let programmerStatusValue = mapLangGraphToUIStatus( + programmerThread.status, + ); + + // Check task completion when thread is idle - no run status needed + if ( + programmerThread.status === "idle" && + areAllPlanItemsCompleted(programmerThread.values?.taskPlan) + ) { + programmerStatusValue = "completed"; + } + + const programmerStatus: StatusResult = { + graph: "programmer", + runId: lastState.runId, + threadId: lastState.threadId, + status: programmerStatusValue, + taskPlan: programmerThread.values?.taskPlan, + }; + + if ( + programmerStatus.status === "running" || + programmerStatus.status === "error" + ) { + return programmerStatus; + } + + return null; + } + break; + + case "planner": + if (lastState.threadId && lastState.runId) { + const plannerThread = await client.threads.get( + lastState.threadId, + ); + + // Use thread status directly for most cases + let plannerStatusValue = mapLangGraphToUIStatus(plannerThread.status); + + // Special case: check for interrupts even if thread status doesn't show interrupted + if ( + plannerThread.interrupts && + Array.isArray(plannerThread.interrupts) && + plannerThread.interrupts.length > 0 + ) { + plannerStatusValue = "paused"; + } + + // No need to check run status for planners - thread status is sufficient + + const plannerStatus: StatusResult = { + graph: "planner", + runId: lastState.runId, + threadId: lastState.threadId, + status: plannerStatusValue, + }; + + if ( + plannerStatus.status === "running" || + plannerStatus.status === "paused" || + plannerStatus.status === "error" + ) { + return plannerStatus; + } + + if (plannerThread.values?.programmerSession) { + const programmerSession = plannerThread.values.programmerSession; + const programmerThread = await client.threads.get( + programmerSession.threadId, + ); + + // Use thread status directly for most cases + let programmerStatusValue = mapLangGraphToUIStatus( + programmerThread.status, + ); + + // Check task completion when thread is idle - no run status needed + if ( + programmerThread.status === "idle" && + areAllPlanItemsCompleted(programmerThread.values?.taskPlan) + ) { + programmerStatusValue = "completed"; + } + + const programmerStatus: StatusResult = { + graph: "programmer", + runId: programmerSession.runId, + threadId: programmerSession.threadId, + status: programmerStatusValue, + taskPlan: programmerThread.values?.taskPlan, + }; + + return resolver.resolve( + { + graph: "manager", + runId: "", + threadId: lastState.threadId, + status: "idle", + }, + plannerStatus, + programmerStatus, + ); + } + + return plannerStatus; + } + break; + + case "manager": { + const managerThread = await client.threads.get( + lastState.threadId, + ); + const managerStatus: StatusResult = { + graph: "manager", + runId: "", + threadId: lastState.threadId, + status: mapLangGraphToUIStatus(managerThread.status), + }; + + if ( + managerStatus.status === "running" || + managerStatus.status === "error" + ) { + return managerStatus; + } + + return null; + } + } + + return null; +} + +async function performFullStatusCheck( + client: Client, + threadId: string, + resolver: StatusResolver, + managerThreadData?: Thread | null, + sessionCache?: SessionCache, +): Promise { + let managerThread: Thread; + + if (managerThreadData) { + managerThread = managerThreadData; + } else { + managerThread = await client.threads.get(threadId); + } + + const managerStatus: StatusResult = { + graph: "manager", + runId: "", + threadId, + status: mapLangGraphToUIStatus(managerThread.status), + }; + + // If manager is running or has error, return immediately without checking sub-sessions + if (managerStatus.status === "running" || managerStatus.status === "error") { + return resolver.resolve(managerStatus); + } + + if (!managerThread.values?.plannerSession) { + return resolver.resolve(managerStatus); + } + + const plannerSession = managerThread.values.plannerSession; + const plannerCacheKey = `planner:${plannerSession.threadId}:${plannerSession.runId}`; + + let plannerThread: Thread; + const cachedPlannerData = getCachedSessionData(sessionCache, plannerCacheKey); + + if (cachedPlannerData?.plannerData) { + plannerThread = cachedPlannerData.plannerData.thread; + } else { + plannerThread = await client.threads.get( + plannerSession.threadId, + ); + + // No run fetch needed for planners - thread status is sufficient + + setCachedSessionData(sessionCache, plannerCacheKey, { + plannerData: { thread: plannerThread }, + }); + } + + // Use thread status directly for most cases + let plannerStatusValue = mapLangGraphToUIStatus(plannerThread.status); + + // Special case: check for interrupts even if thread status doesn't show interrupted + if ( + plannerThread.interrupts && + Array.isArray(plannerThread.interrupts) && + plannerThread.interrupts.length > 0 + ) { + plannerStatusValue = "paused"; + } + + // No run status check needed for planners + + const plannerStatus: StatusResult = { + graph: "planner", + runId: plannerSession.runId, + threadId: plannerSession.threadId, + status: plannerStatusValue, + }; + + if ( + plannerStatus.status === "running" || + plannerStatus.status === "paused" || + plannerStatus.status === "error" + ) { + return resolver.resolve(managerStatus, plannerStatus); + } + + if (!plannerThread.values?.programmerSession) { + return resolver.resolve(managerStatus, plannerStatus); + } + + const programmerSession = plannerThread.values.programmerSession; + const programmerCacheKey = `programmer:${programmerSession.threadId}:${programmerSession.runId}`; + + let programmerThread: Thread; + const cachedProgrammerData = getCachedSessionData( + sessionCache, + programmerCacheKey, + ); + + if (cachedProgrammerData?.programmerData) { + programmerThread = cachedProgrammerData.programmerData.thread; + } else { + programmerThread = await client.threads.get( + programmerSession.threadId, + ); + + // No run fetch needed - we only check task completion from thread data + + setCachedSessionData(sessionCache, programmerCacheKey, { + programmerData: { thread: programmerThread }, + }); + } + + // Use thread status directly for most cases + let programmerStatusValue = mapLangGraphToUIStatus(programmerThread.status); + + // Check task completion when thread is idle - no run status needed + if ( + programmerThread.status === "idle" && + areAllPlanItemsCompleted(programmerThread.values?.taskPlan) + ) { + programmerStatusValue = "completed"; + } + + const programmerStatus: StatusResult = { + graph: "programmer", + runId: programmerSession.runId, + threadId: programmerSession.threadId, + status: programmerStatusValue, + taskPlan: programmerThread.values?.taskPlan, + }; + + return resolver.resolve(managerStatus, plannerStatus, programmerStatus); +} diff --git a/apps/web/src/stores/thread-store.ts b/apps/web/src/stores/thread-store.ts new file mode 100644 index 00000000..8bb5d06c --- /dev/null +++ b/apps/web/src/stores/thread-store.ts @@ -0,0 +1,43 @@ +import { create } from "zustand"; + +/** + * Simplified thread store for UI state only + * Data caching is handled by SWR hooks directly + */ +export interface ThreadStoreState { + // UI state only - active thread tracking + activeThreadId: string | null; + + // UI state only - global polling control + isGlobalPollingEnabled: boolean; + + // Actions + setActiveThread: (threadId: string | null) => void; + setGlobalPolling: (enabled: boolean) => void; +} + +/** + * Minimal Zustand store for thread UI state management + * All data caching moved to SWR for better performance and consistency + */ +export const useThreadStore = create((set) => ({ + // Initial state + activeThreadId: null, + isGlobalPollingEnabled: true, + + // Actions + setActiveThread: (threadId) => { + set({ activeThreadId: threadId }); + }, + + setGlobalPolling: (enabled) => { + set({ isGlobalPollingEnabled: enabled }); + }, +})); + +// Selector hooks for specific pieces of state +export const useActiveThreadId = () => + useThreadStore((state) => state.activeThreadId); + +export const useGlobalPollingEnabled = () => + useThreadStore((state) => state.isGlobalPollingEnabled);