Feat: Thread Status SWR Polling (#370)

* init swr polling branch

* useThreads types improvement

* minor build fix

* major refactor and simplification of polling system

* remove temp guide file

* refactor, remove errorBoundary, simplify schema

* refactor: unify client<>server status logic

* fix threads/page conflicts + functionality

* Update apps/web/src/components/v2/thread-card.tsx

Co-authored-by: Brace Sproul <braceasproul@gmail.com>

* CR in progress

* CR fix: add cn to thread-card.tsx

Co-authored-by: Brace Sproul <braceasproul@gmail.com>

* CR

* CR

* CR: no threads state, Batch ThreadsStatus for /chat/threads actual logic

* drop req type validation in lg proxy

* optimize thread polling, reduce unnecessary api calls, match implementation spec

* reduce unnecessary api calls, update all threads page

* add thread loading state to threadCard

* utilize threads/search for manager thread info

* improve polling api req efficiency

* rm comments

* post merge format

* fix import

* improve programmer cache logic

* reduce run requests

* remove timeout run requests

* remove run check for completed tasks

* fix: update status dot to support dark mode styles too

* fix session cache and proper typing

* fix: more typing

* code review fixes

* CR

* code review

* fix build

---------

Co-authored-by: Brace Sproul <braceasproul@gmail.com>
This commit is contained in:
Dylan Boudro 2025-07-15 20:51:12 -04:00 • committed by GitHub
parent 4053f18df5
commit d51af805e2
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
22 changed files with 1122 additions and 491 deletions

View file

@ -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<GraphState>(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 <ThreadViewLoading onBackToHome={handleBackToHome} />;
}
// Convert all threads to display format
const displayThreads: ThreadDisplayInfo[] = threads.map(threadToDisplayInfo);
const currentDisplayThread = threadToDisplayInfo(thread);
return (
<div className="bg-background fixed inset-0">
<ThreadView
stream={stream}
displayThread={currentDisplayThread}
allDisplayThreads={displayThreads}
allDisplayThreads={threadsMetadata}
onBackToHome={handleBackToHome}
/>
</div>

View file

@ -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<GraphState>({
const { threads, isLoading: threadsLoading } = useThreadsSWR({
assistantId: MANAGER_GRAPH_ID,
refreshInterval: 15000, // Poll every 15 seconds
});
if (!threads) {
return <div>No threads</div>;
}
// Convert Thread objects to ThreadDisplayInfo for UI
const displayThreads: ThreadDisplayInfo[] = threads.map(threadToDisplayInfo);
return (
<div className="bg-background h-screen">
<Suspense>
<Toaster />
<GitHubAppProvider>
<DefaultView
threads={displayThreads}
threads={threads}
threadsLoading={threadsLoading}
/>
</GitHubAppProvider>

View file

@ -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<GraphState>({
const { threads, isLoading: threadsLoading } = useThreadsSWR({
assistantId: MANAGER_GRAPH_ID,
refreshInterval: 15000, // Poll every 15 seconds
});
const [searchQuery, setSearchQuery] = useState("");
const [statusFilter, setStatusFilter] = useState<FilterStatus>("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) => (
<Button
@ -148,7 +174,6 @@ function AllThreadsPageContent() {
<div className="flex-1 overflow-auto">
<div className="mx-auto max-w-6xl p-4">
{statusFilter === "all" ? (
// Show grouped view when "all" is selected
<div className="space-y-6">
{Object.entries(groupedThreads).map(([status, threads]) => {
if (threads.length === 0) return null;
@ -170,6 +195,8 @@ function AllThreadsPageContent() {
<ThreadCard
key={thread.id}
thread={thread}
status={statusMap[thread.id]}
statusLoading={statusLoading}
/>
))}
</div>
@ -178,44 +205,50 @@ function AllThreadsPageContent() {
})}
</div>
) : (
// Show flat list when specific status is selected
<div className="grid gap-3 md:grid-cols-2 lg:grid-cols-3">
{filteredThreads.map((thread) => (
<ThreadCard
key={thread.id}
thread={thread}
status={statusMap[thread.id]}
statusLoading={statusLoading}
/>
))}
</div>
)}
{filteredThreads.length === 0 && !threadsLoading && (
<div className="py-12 text-center">
<div className="text-muted-foreground mb-2">
No threads found
{filteredThreads.length === 0 &&
!threadsLoading &&
!statusLoading && (
<div className="py-12 text-center">
<div className="text-muted-foreground mb-2">
No threads found
</div>
<div className="text-muted-foreground/70 text-xs">
{!threads || threads.length === 0
? "No threads have been created yet"
: searchQuery
? "Try adjusting your search query"
: "No threads match the selected filter"}
</div>
</div>
<div className="text-muted-foreground/70 text-xs">
{searchQuery
? "Try adjusting your search query"
: "No threads match the selected filter"}
</div>
</div>
)}
)}
{threadsLoading && threads.length === 0 && (
<div>
<div className="mb-3 flex items-center gap-2">
<h2 className="text-foreground text-base font-semibold capitalize">
Loading threads...
</h2>
{(threadsLoading || statusLoading) &&
(!threads || threads.length === 0) && (
<div>
<div className="mb-3 flex items-center gap-2">
<h2 className="text-foreground text-base font-semibold capitalize">
Loading threads...
</h2>
</div>
<div className="grid gap-3 md:grid-cols-2 lg:grid-cols-3">
{Array.from({ length: 9 }).map((_, index) => (
<ThreadCardLoading key={`all-threads-loading-${index}`} />
))}
</div>
</div>
<div className="grid gap-3 md:grid-cols-2 lg:grid-cols-3">
{Array.from({ length: 9 }).map((_, index) => (
<ThreadCardLoading key={`all-threads-loading-${index}`} />
))}
</div>
</div>
)}
)}
</div>
</div>
</div>

View file

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

View file

@ -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<ManagerGraphState>[];
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) {
<ThreadCardLoading />
</>
)}
{threads.slice(0, 4).map((thread) => (
{displayThreads.map((thread) => (
<ThreadCard
key={thread.id}
thread={thread}
status={statusMap[thread.id]}
statusLoading={statusLoading}
/>
))}
</div>

View file

@ -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 <Loader2 className="h-4 w-4 animate-spin" />;
case "completed":
return <CheckCircle className="h-4 w-4" />;
case "error":
return <AlertCircle className="h-4 w-4" />;
case "idle":
return <Clock className="h-4 w-4" />;
case "paused":
return <Pause className="h-4 w-4" />;
case "failed":
return <XCircle className="h-4 w-4" />;
default:
@ -83,11 +106,22 @@ export function ThreadCard({ thread }: { thread: ThreadDisplayInfo }) {
</div>
<Badge
variant="secondary"
className={cn(getStatusColor(thread.status), "text-xs")}
className={cn(
"text-xs",
isStatusLoading
? "bg-gray-200 text-gray-600 dark:bg-gray-800 dark:text-gray-400"
: getStatusColor(displayStatus),
)}
>
<div className="flex items-center gap-1">
{getStatusIcon(thread.status)}
<span className="capitalize">{thread.status}</span>
{isStatusLoading ? (
<Loader2 className="h-4 w-4 animate-spin" />
) : (
getStatusIcon(displayStatus)
)}
<span className="capitalize">
{isStatusLoading ? "Loading..." : displayStatus}
</span>
</div>
</Badge>
</div>

View file

@ -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 (
<Sheet
@ -54,14 +51,6 @@ export function ThreadSwitcher({
>
<Layers3 className="h-3 w-3" />
<span className="hidden sm:inline">Switch Thread</span>
{runningCount > 0 && (
<Badge
variant="secondary"
className="h-4 bg-blue-100 px-1 text-xs text-blue-700 dark:bg-blue-950 dark:text-blue-400"
>
{runningCount}
</Badge>
)}
</Button>
</SheetTrigger>
<SheetContent

View file

@ -7,7 +7,7 @@ import { Card, CardContent } from "@/components/ui/card";
import { ArrowLeft, GitBranch, Terminal, Clock } from "lucide-react";
import { Tabs, TabsContent, TabsList, TabsTrigger } from "@/components/ui/tabs";
import { ThreadSwitcher } from "./thread-switcher";
import { ThreadDisplayInfo } from "./types";
import { ThreadMetadata } from "./types";
import { useStream } from "@langchain/langgraph-sdk/react";
import { ManagerGraphState } from "@open-swe/shared/open-swe/manager/types";
import { PlannerGraphState } from "@open-swe/shared/open-swe/planner/types";
@ -20,6 +20,9 @@ import {
PROGRAMMER_GRAPH_ID,
PLANNER_GRAPH_ID,
} from "@open-swe/shared/constants";
import { useThreadStatus } from "@/hooks/useThreadStatus";
import { cn } from "@/lib/utils";
import { StickToBottom } from "use-stick-to-bottom";
import {
StickyToBottomContent,
@ -27,12 +30,11 @@ import {
} from "../../utils/scroll-utils";
import { ManagerChat } from "./manager-chat";
import { CancelStreamButton } from "./cancel-stream-button";
import { cn } from "@/lib/utils";
interface ThreadViewProps {
stream: ReturnType<typeof useStream<ManagerGraphState>>;
displayThread: ThreadDisplayInfo;
allDisplayThreads: ThreadDisplayInfo[];
displayThread: ThreadMetadata;
allDisplayThreads: ThreadMetadata[];
onBackToHome: () => void;
}
@ -51,6 +53,23 @@ export function ThreadView({
const [programmerSession, setProgrammerSession] =
useState<ManagerGraphState["programmerSession"]>();
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({
<div
className={cn(
"size-2 flex-shrink-0 rounded-full",
displayThread.status === "running"
? "bg-blue-500"
: displayThread.status === "completed"
? "bg-green-500"
: "bg-red-500",
getStatusDotColor(realTimeStatus),
)}
></div>
<span className="text-muted-foreground max-w-[500px] truncate font-mono text-sm">

View file

@ -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<ManagerGraphState>,
): 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,
};
}

View file

@ -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<ManagerGraphState>): {
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,
};
}

View file

@ -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<GraphState>[];
getThread: (threadId: string) => Promise<Thread<GraphState> | null>;
onUpdate: (
updatedThreads: Thread<GraphState>[],
changedThreadIds: string[],
) => void;
enabled?: boolean;
}
export function useThreadPolling({
threads,
getThread,
onUpdate,
enabled = true,
}: UseThreadPollingProps) {
const pollerRef = useRef<ThreadPoller | null>(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(),
};
}

View file

@ -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<ThreadStatusData>(
swrKey,
() => fetchThreadStatus(threadId),
{
...THREAD_STATUS_SWR_CONFIG,
refreshInterval,
},
);
return {
status: data?.status || "idle",
isLoading,
error,
mutate,
};
}

View file

@ -1,58 +0,0 @@
import { createClient } from "@/providers/client";
import { Thread } from "@langchain/langgraph-sdk";
import { useCallback, useEffect, useState } from "react";
export function useThreads<State extends Record<string, any>>(
assistantId?: string,
) {
const apiUrl: string | undefined = process.env.NEXT_PUBLIC_API_URL ?? "";
const [threads, setThreads] = useState<Thread<State>[]>([]);
const [threadsLoading, setThreadsLoading] = useState(false);
const getThread = useCallback(
async (threadId: string): Promise<Thread<State> | null> => {
if (!apiUrl) return null;
const client = createClient(apiUrl);
try {
const thread = await client.threads.get<State>(threadId);
return thread;
} catch (error) {
console.error("Failed to fetch thread:", threadId, error);
return null;
}
},
[apiUrl],
);
const getThreads = useCallback(async (): Promise<Thread<State>[] | 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<State>(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 };
}

View file

@ -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<State extends Record<string, any>>(
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<State extends Record<string, any>>(
// Create a unique key for SWR caching based on assistantId
const swrKey = assistantId ? ["threads", assistantId] : ["threads", "all"];
const fetcher = async (): Promise<Thread<State>[]> => {
const fetcher = async (): Promise<Thread<TGraphState>[]> => {
if (!apiUrl) {
throw new Error("API URL is not configured");
}
@ -38,7 +60,7 @@ export function useThreadsSWR<State extends Record<string, any>>(
}
: undefined;
return await client.threads.search<State>(searchArgs);
return await client.threads.search<TGraphState>(searchArgs);
};
const { data, error, isLoading, mutate, isValidating } = useSWR(
@ -48,8 +70,9 @@ export function useThreadsSWR<State extends Record<string, any>>(
refreshInterval,
revalidateOnFocus,
revalidateOnReconnect,
errorRetryCount: 3,
errorRetryInterval: 5000,
errorRetryCount: THREAD_SWR_CONFIG.errorRetryCount,
errorRetryInterval: THREAD_SWR_CONFIG.errorRetryInterval,
dedupingInterval: THREAD_SWR_CONFIG.dedupingInterval,
},
);

View file

@ -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<PlannerGraphState> };
programmerData?: { thread: Thread<GraphState> };
timestamp: number;
}
export type SessionCache = Map<string, SessionCacheData>;
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<string, ThreadStatusData>,
managerThreads?: Thread<ManagerGraphState>[],
): Promise<{
statusMap: ThreadStatusMap;
updatedStates: Map<string, ThreadStatusData>;
}> {
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<string, ThreadStatusData>();
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<ManagerGraphState>[],
): UseThreadsStatusResult {
const lastPollingStatesRef = useRef<Map<string, ThreadStatusData>>(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]);
}

View file

@ -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<GraphState>[],
changedThreadIds: string[],
) => void;
}
export class ThreadPoller {
private config: PollConfig;
private isPolling: boolean = false;
private intervalId: NodeJS.Timeout | null = null;
private threads: Thread<GraphState>[];
private getThreadFn: (threadId: string) => Promise<Thread<GraphState> | null>;
constructor(
config: PollConfig,
threads: Thread<GraphState>[],
getThreadFn: (threadId: string) => Promise<Thread<GraphState> | 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<void> {
try {
const currentThreads = this.threads;
const threadsToPool = currentThreads.slice(0, 10);
const updatedThreads: Thread<GraphState>[] = [];
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<GraphState>,
updated: Thread<GraphState>,
): 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)
);
}
}

View file

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

View file

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

View file

@ -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<ManagerGraphState>[],
): 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,
};
});
}

View file

@ -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<GraphState>[];
@ -110,32 +109,6 @@ export function ThreadProvider({ children }: { children: ReactNode }) {
refreshThreads();
}, [refreshThreads]);
const handlePollingUpdate = useCallback(
(updatedThreads: Thread<GraphState>[], 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<GraphState>,

View file

@ -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<SessionCacheData>,
): void {
if (!sessionCache) return;
sessionCache.set(sessionKey, {
...data,
timestamp: Date.now(),
});
}
export async function fetchThreadStatus(
threadId: string,
lastPollingState: ThreadStatusData | null = null,
managerThreadData?: Thread<ManagerGraphState> | null,
sessionCache?: SessionCache,
): Promise<ThreadStatusData> {
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<ThreadStatusData | null> {
switch (lastState.graph) {
case "programmer":
if (lastState.threadId && lastState.runId) {
const programmerThread = await client.threads.get<GraphState>(
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<PlannerGraphState>(
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<GraphState>(
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<ManagerGraphState>(
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<ManagerGraphState> | null,
sessionCache?: SessionCache,
): Promise<ThreadStatusData> {
let managerThread: Thread<ManagerGraphState>;
if (managerThreadData) {
managerThread = managerThreadData;
} else {
managerThread = await client.threads.get<ManagerGraphState>(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<PlannerGraphState>;
const cachedPlannerData = getCachedSessionData(sessionCache, plannerCacheKey);
if (cachedPlannerData?.plannerData) {
plannerThread = cachedPlannerData.plannerData.thread;
} else {
plannerThread = await client.threads.get<PlannerGraphState>(
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<GraphState>;
const cachedProgrammerData = getCachedSessionData(
sessionCache,
programmerCacheKey,
);
if (cachedProgrammerData?.programmerData) {
programmerThread = cachedProgrammerData.programmerData.thread;
} else {
programmerThread = await client.threads.get<GraphState>(
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);
}

View file

@ -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<ThreadStoreState>((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);