From f549d119abf8b835efaa9b8f783b552c37d6ee7a Mon Sep 17 00:00:00 2001 From: Brace Sproul Date: Thu, 12 Jun 2025 13:29:20 -0700 Subject: [PATCH] fix: parallelize thread polling (#143) * fix: parallelize thread polling * format --- apps/web/src/lib/polling/thread-poller.ts | 29 ++++++----- apps/web/src/providers/Thread.tsx | 62 +---------------------- 2 files changed, 18 insertions(+), 73 deletions(-) diff --git a/apps/web/src/lib/polling/thread-poller.ts b/apps/web/src/lib/polling/thread-poller.ts index d804de43..9e43a2d3 100644 --- a/apps/web/src/lib/polling/thread-poller.ts +++ b/apps/web/src/lib/polling/thread-poller.ts @@ -55,22 +55,27 @@ export class ThreadPoller { const changedThreadIds: string[] = []; const errors: string[] = []; - for (const currentThread of threadsToPool) { - try { - const updatedThread = await this.getThreadFn(currentThread.thread_id); - if (updatedThread) { - updatedThreads.push(updatedThread); + 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); + if (this.hasThreadChanged(currentThread, updatedThread)) { + changedThreadIds.push(updatedThread.thread_id); + } } + } catch (error) { + errors.push(`Thread ${currentThread.thread_id}: ${error}`); + updatedThreads.push(currentThread); } - } catch (error) { - errors.push(`Thread ${currentThread.thread_id}: ${error}`); + }), + ); - updatedThreads.push(currentThread); - } - } + await pollUpdatePromise; if (changedThreadIds.length > 0) { this.config.onUpdate(updatedThreads, changedThreadIds); diff --git a/apps/web/src/providers/Thread.tsx b/apps/web/src/providers/Thread.tsx index 9c853cd9..a95315ce 100644 --- a/apps/web/src/providers/Thread.tsx +++ b/apps/web/src/providers/Thread.tsx @@ -9,10 +9,9 @@ import { Dispatch, SetStateAction, useEffect, - useTransition, } from "react"; import { createClient } from "./client"; -import { TaskPlan, GraphState } from "@open-swe/shared/open-swe/types"; +import { GraphState } from "@open-swe/shared/open-swe/types"; import { useThreadPolling } from "@/hooks/useThreadPolling"; interface ThreadContextType { @@ -22,7 +21,6 @@ interface ThreadContextType { setThreadsLoading: Dispatch>; refreshThreads: () => Promise; getThread: (threadId: string) => Promise | null>; - isPending: boolean; recentlyUpdatedThreads: Set; handleThreadClick: ( thread: Thread, @@ -43,62 +41,6 @@ function getThreadSearchMetadata( } } -const getTaskCounts = ( - tasks?: TaskPlan, - proposedPlan?: string[], - existingCounts?: { totalTasksCount: number; completedTasksCount: number }, -): { totalTasksCount: number; completedTasksCount: number } => { - const defaultCounts = existingCounts || { - totalTasksCount: 0, - completedTasksCount: 0, - }; - - if (proposedPlan && proposedPlan.length > 0 && !tasks) { - return { - totalTasksCount: proposedPlan.length, - completedTasksCount: 0, - }; - } - - if (!tasks || !tasks.tasks || tasks.tasks.length === 0) { - return defaultCounts; - } - const activeTaskIndex = tasks.activeTaskIndex; - const activeTask = tasks.tasks.find( - (task) => task.taskIndex === activeTaskIndex, - ); - - if ( - !activeTask || - !activeTask.planRevisions || - activeTask.planRevisions.length === 0 - ) { - return defaultCounts; - } - - const activeRevisionIndex = activeTask.activeRevisionIndex; - const activeRevision = activeTask.planRevisions.find( - (revision) => revision.revisionIndex === activeRevisionIndex, - ); - - if ( - !activeRevision || - !activeRevision.plans || - activeRevision.plans.length === 0 - ) { - return defaultCounts; - } - - const plans = activeRevision.plans; - - const completedTasksCount = plans.filter((p) => p.completed)?.length || 0; - - return { - totalTasksCount: plans.length, - completedTasksCount, - }; -}; - export function ThreadProvider({ children }: { children: ReactNode }) { const apiUrl: string | undefined = process.env.NEXT_PUBLIC_API_URL ?? ""; const assistantId: string | undefined = @@ -106,7 +48,6 @@ export function ThreadProvider({ children }: { children: ReactNode }) { const [threads, setThreads] = useState[]>([]); const [threadsLoading, setThreadsLoading] = useState(false); - const [isPending, startTransition] = useTransition(); const [recentlyUpdatedThreads, setRecentlyUpdatedThreads] = useState< Set >(new Set()); @@ -215,7 +156,6 @@ export function ThreadProvider({ children }: { children: ReactNode }) { setThreadsLoading, refreshThreads, getThread, - isPending, recentlyUpdatedThreads, handleThreadClick, };