Update reconnection logic in useEventSource to prevent excessive retries

This commit is contained in:
Steven T. Cramer 2025-06-13 15:26:04 +07:00
parent 138df6d4b2
commit f4a9450092
2 changed files with 78 additions and 14 deletions

View file

@ -10,33 +10,58 @@ export const dynamic = "force-dynamic"
export async function GET(request: NextRequest, { params }: { params: Promise<{ id: string }> }) {
const { id } = await params
console.log("YO YO YO - Accessing SSE stream route for run ID:", id);
const requestId = crypto.randomUUID()
console.log("YO YO YO - Processing request for run ID:", id, "with request ID:", requestId);
const stream = new SSEStream()
const run = await findRun(Number(id))
console.log(`YO YO YO - [stream#${requestId}] Initializing stream for run ID: ${id}`);
let run;
try {
console.log(`YO YO YO - [stream#${requestId}] Attempting to find run with ID: ${id}`);
run = await findRun(Number(id))
console.log(`YO YO YO - [stream#${requestId}] Found run with ID: ${id}`);
if (!run) {
console.error(`YO YO YO - [stream#${requestId}] Run ${id} not found`);
return new Response(`Run ${id} not found`, { status: 404 })
}
console.log(`YO YO YO - [stream#${requestId}] Successfully retrieved run data for ID: ${id}`);
} catch (error) {
console.error(`YO YO YO - [stream#${requestId}] Error finding run with ID ${id}:`, error);
return new Response("Internal server error", { status: 500 })
}
const redis = await redisClient()
if (!redis) {
console.error(`YO YO YO - [stream#${requestId}] Redis client not available`);
return new Response("Internal server error", { status: 500 })
}
let isStreamClosed = false
const channelName = `evals:${run.id}`
const onMessage = async (data: string) => {
console.log(`YO YO YO - [stream#${requestId}] Received message on channel ${channelName}:`, data);
if (isStreamClosed || stream.isClosed) {
return
}
try {
const taskEvent = taskEventSchema.parse(JSON.parse(data))
// console.log(`[stream#${requestId}] task event -> ${taskEvent.eventName}`)
const parsedData = JSON.parse(data)
const taskEvent = taskEventSchema.parse(parsedData)
console.log(`Yo Yo Yo - [stream#${requestId}] task event -> ${taskEvent.eventName}`)
const writeSuccess = await stream.write(JSON.stringify(taskEvent))
if (!writeSuccess) {
console.error(`YO YO YO - [stream#${requestId}] failed to write to stream, disconnecting`);
await disconnect()
}
} catch (_error) {
console.error(`[stream#${requestId}] invalid task event:`, data)
} catch (error) {
console.error(`YO YO YO - [stream#${requestId}] invalid task event:`, data, 'Error:', error);
}
}
const disconnect = async () => {
console.log(`YO YO YO - [stream#${requestId}] disconnecting from channel ${channelName}`);
if (isStreamClosed) {
return
}
@ -45,23 +70,38 @@ export async function GET(request: NextRequest, { params }: { params: Promise<{
try {
await redis.unsubscribe(channelName)
console.log(`[stream#${requestId}] unsubscribed from ${channelName}`)
console.log(`YO YO YO - [stream#${requestId}] unsubscribed from ${channelName}`);
} catch (error) {
console.error(`[stream#${requestId}] error unsubscribing:`, error)
console.error(`YO YO YO - [stream#${requestId}] error unsubscribing:`, error);
}
try {
await stream.close()
} catch (error) {
console.error(`[stream#${requestId}] error closing stream:`, error)
console.error(`YO YO YO - [stream#${requestId}] error closing stream:`, error);
}
}
await redis.subscribe(channelName, onMessage)
try {
console.log(`YO YO YO - [stream#${requestId}] subscribing to Redis channel ${channelName}`);
await redis.subscribe(channelName, onMessage)
} catch (error) {
console.error(`YO YO YO - [stream#${requestId}] error subscribing to Redis:`, error);
return new Response("Internal server error", { status: 500 })
}
// Add a timeout to close the stream after a period of inactivity or errors
const timeoutDuration = 300000; // 5 minutes
const timeoutId = setTimeout(() => {
console.log(`[stream#${requestId}] timeout after ${timeoutDuration}ms`);
disconnect().catch((error) => {
console.error(`[stream#${requestId}] timeout cleanup error:`, error)
});
}, timeoutDuration);
request.signal.addEventListener("abort", () => {
console.log(`[stream#${requestId}] abort`)
console.log(`[stream#${requestId}] abort`);
clearTimeout(timeoutId);
disconnect().catch((error) => {
console.error(`[stream#${requestId}] cleanup error:`, error)
})

View file

@ -41,6 +41,7 @@ export function useEventSource({ url, withCredentials, onMessage }: UseEventSour
setStatus("waiting")
sourceRef.current = new EventSource(url, { withCredentials })
console.log("YO YO YO - Creating new EventSource connection for URL:", url);
sourceRef.current.onopen = () => {
if (isUnmountedRef.current) {
@ -59,7 +60,7 @@ export function useEventSource({ url, withCredentials, onMessage }: UseEventSour
handleMessage(event)
}
sourceRef.current.onerror = () => {
sourceRef.current.onerror = (event) => {
if (isUnmountedRef.current) {
return
}
@ -67,15 +68,38 @@ export function useEventSource({ url, withCredentials, onMessage }: UseEventSour
statusRef.current = "error"
setStatus("error")
console.log("YO YO YO - Error in EventSource:", event);
if (event instanceof ErrorEvent) {
console.log("YO YO YO - Error details:", {
message: event.message,
filename: event.filename,
lineno: event.lineno,
colno: event.colno,
error: event.error
});
} else {
console.log("YO YO YO - Unknown error event type:", event);
}
// Clean up current connection.
cleanup()
// Attempt to reconnect after a delay.
// Attempt to reconnect after a delay with a simple backoff, with a maximum retry limit.
if (!sourceRef.current) {
return;
}
const retryCount = (sourceRef.current.retryCount || 0) + 1;
sourceRef.current.retryCount = retryCount;
if (retryCount > 5) {
console.log("YO YO YO - Maximum retry limit reached for EventSource. Stopping reconnection attempts.");
return;
}
console.log("YO YO YO - This is where reconnection happens for EventSource. Retry attempt:", retryCount);
reconnectTimeoutRef.current = setTimeout(() => {
if (!isUnmountedRef.current) {
createEventSource()
}
}, 1000)
}, 5000 * retryCount) // Exponential backoff: 5s, 10s, 15s, etc.
}
}, [url, withCredentials, handleMessage, cleanup])