From 278b9591d7c337e782eb7cdf819078b53591864c Mon Sep 17 00:00:00 2001 From: "santasri.pachhal" Date: Thu, 2 Jul 2026 14:28:28 +0530 Subject: [PATCH] feat: added sse for upload section --- .../upload/hooks/useUploadListEvents.ts | 112 ++++++++++++++++++ src/app/(modules)/upload/page.tsx | 2 + src/constants/apiRoutes.ts | 3 + src/services/api/video.service.ts | 8 ++ src/types/upload.ts | 8 ++ 5 files changed, 133 insertions(+) create mode 100644 src/app/(modules)/upload/hooks/useUploadListEvents.ts diff --git a/src/app/(modules)/upload/hooks/useUploadListEvents.ts b/src/app/(modules)/upload/hooks/useUploadListEvents.ts new file mode 100644 index 0000000..e80a5ee --- /dev/null +++ b/src/app/(modules)/upload/hooks/useUploadListEvents.ts @@ -0,0 +1,112 @@ +'use client'; + +import { useCallback, useMemo, useRef } from 'react'; +import { useQueryClient } from '@tanstack/react-query'; + +import { API_ROUTES } from '@/constants/apiRoutes'; +import { useSseWithToken } from '@/hooks/sse/useSseWithToken'; +import { videoService } from '@/services/api'; +import type { + PaginationParams, + UploadListResponse, + UploadStatusEvent, +} from '@/types'; + +import { uploadKeys } from './useUploadQueries'; + +function getUploadListParams(queryKey: readonly unknown[]) { + const params = queryKey[2]; + + if (!params || typeof params !== 'object' || Array.isArray(params)) { + return undefined; + } + + return params as PaginationParams; +} + +function patchUploadList( + current: UploadListResponse | undefined, + event: UploadStatusEvent, + params?: PaginationParams, +) { + if (!current) return current; + + const existingIndex = current.items.findIndex( + (upload) => upload.id === event.item.id, + ); + + if (existingIndex >= 0) { + return { + ...current, + items: current.items.map((upload) => + upload.id === event.item.id ? event.item : upload, + ), + }; + } + + if (event.kind !== 'created') return current; + + const nextTotal = current.total + 1; + if ((params?.skip ?? 0) !== 0) { + return { ...current, total: nextTotal }; + } + + const pageLimit = params?.limit ?? current.items.length + 1; + return { + items: [event.item, ...current.items].slice(0, pageLimit), + total: nextTotal, + }; +} + +export function useUploadListEvents(enabled = true) { + const queryClient = useQueryClient(); + const hasInvalidatedAfterConnectionErrorRef = useRef(false); + + const events = useMemo( + () => ({ + upload_status: (event: UploadStatusEvent) => { + queryClient + .getQueryCache() + .findAll({ queryKey: uploadKeys.lists() }) + .forEach((query) => { + queryClient.setQueryData( + query.queryKey, + (current) => + patchUploadList( + current, + event, + getUploadListParams(query.queryKey), + ), + ); + }); + }, + heartbeat: () => undefined, + }), + [queryClient], + ); + + const getPath = useCallback( + (token: string) => API_ROUTES.VIDEOS.UPLOAD_EVENTS(token), + [], + ); + + const onConnectionError = useCallback(() => { + if (hasInvalidatedAfterConnectionErrorRef.current) return; + + hasInvalidatedAfterConnectionErrorRef.current = true; + queryClient.invalidateQueries({ queryKey: uploadKeys.lists() }); + }, [queryClient]); + + const onConnectionOpen = useCallback(() => { + hasInvalidatedAfterConnectionErrorRef.current = false; + }, []); + + useSseWithToken({ + enabled, + getToken: videoService.createUploadEventsToken, + getPath, + events, + onConnectionError, + onConnectionOpen, + }); +} diff --git a/src/app/(modules)/upload/page.tsx b/src/app/(modules)/upload/page.tsx index d4ff595..9a4dccb 100644 --- a/src/app/(modules)/upload/page.tsx +++ b/src/app/(modules)/upload/page.tsx @@ -12,11 +12,13 @@ import { useUploadListColumns } from './components/UploadListColumns'; import { UploadTable } from './components/UploadTable'; import { UploadVideoDialog } from './components/UploadVideoDialog'; import { useUploadFilters } from './hooks/useUploadFilters'; +import { useUploadListEvents } from './hooks/useUploadListEvents'; import { useUploadsQuery } from './hooks/useUploadQueries'; export default function UploadPage() { const { skip, setSkip, limit, setLimit } = useUploadFilters(); const uploadsQuery = useUploadsQuery({ skip, limit }); + useUploadListEvents(); const [isCreateOpen, setIsCreateOpen] = useState(false); const [videoUpload, setVideoUpload] = useState(null); diff --git a/src/constants/apiRoutes.ts b/src/constants/apiRoutes.ts index 788e31d..c9fec0c 100644 --- a/src/constants/apiRoutes.ts +++ b/src/constants/apiRoutes.ts @@ -50,6 +50,9 @@ export const API_ROUTES = { VIDEOS: { UPLOAD: '/biz/api/v1/upload', UPLOADS: '/biz/api/v1/uploads', + UPLOAD_EVENTS_TOKEN: '/biz/api/v1/uploads/events/token', + UPLOAD_EVENTS: (token: string) => + `/biz/api/v1/uploads/events?sse_token=${encodeURIComponent(token)}`, RESULTS: (id: string) => `/biz/api/v1/results/${id}/completed`, }, TICKETS: { diff --git a/src/services/api/video.service.ts b/src/services/api/video.service.ts index 183d513..b5a6822 100644 --- a/src/services/api/video.service.ts +++ b/src/services/api/video.service.ts @@ -3,6 +3,7 @@ import { API_ROUTES } from '@/constants/apiRoutes'; import { CompletedVideoResult, PaginationParams, + SseTokenResponse, UploadListResponse, } from '@/types'; @@ -25,6 +26,13 @@ export const videoService = { return response.data; }, + createUploadEventsToken: async (): Promise => { + const response = await axiosClient.post( + API_ROUTES.VIDEOS.UPLOAD_EVENTS_TOKEN, + ); + return response.data; + }, + /** * Upload a video for processing */ diff --git a/src/types/upload.ts b/src/types/upload.ts index d0d69a1..528ccc8 100644 --- a/src/types/upload.ts +++ b/src/types/upload.ts @@ -1,5 +1,6 @@ export type UploadStatus = 'queued' | 'processing' | 'completed' | 'failed'; export type UploadResult = 'ticket_created' | 'analysis_failed' | null; +export type UploadEventKind = 'created' | 'progressed' | 'transitioned'; export interface UploadListItem { id: string; @@ -24,3 +25,10 @@ export interface UploadListResponse { items: UploadListItem[]; total: number; } + +export interface UploadStatusEvent { + kind: UploadEventKind; + event_id: string; + item: UploadListItem; + occurred_at: string; +}