diff --git a/package.json b/package.json index 1fd33ad4..760af793 100644 --- a/package.json +++ b/package.json @@ -86,6 +86,8 @@ "@streamdown/math": "^1.0.2", "@streamdown/mermaid": "^1.0.2", "@supabase/supabase-js": "^2.101.1", + "@tanstack/pacer": "^0.20.1", + "@tanstack/react-pacer": "^0.21.1", "@tanstack/react-query": "^5.96.1", "@tanstack/react-query-devtools": "^5.96.1", "@tiptap/core": "3.22.2", @@ -108,7 +110,7 @@ "@vercel/otel": "^2.1.1", "ai": "^6.0.146", "assistant-stream": "^0.3.10", - "better-auth": "^1.5.6", + "better-auth": "^1.6.2", "class-variance-authority": "^0.7.1", "clsx": "^2.1.1", "cmdk": "1.1.1", @@ -160,6 +162,7 @@ }, "devDependencies": { "@tailwindcss/postcss": "^4.2.1", + "@types/bun": "^1.3.11", "@types/lodash.throttle": "^4.1.9", "@types/node": "^25.3.0", "@types/pg": "^8.16.0", diff --git a/src/app/api/audio/process/route.ts b/src/app/api/audio/process/route.ts index d38cd1b6..062fbd75 100644 --- a/src/app/api/audio/process/route.ts +++ b/src/app/api/audio/process/route.ts @@ -1,8 +1,12 @@ import { NextRequest, NextResponse } from "next/server"; -import { auth } from "@/lib/auth"; -import { headers } from "next/headers"; import { start } from "workflow/api"; import { audioTranscribeWorkflow } from "@/workflows/audio-transcribe"; +import { + requireAuth, + verifyWorkspaceAccess, + withErrorHandling, +} from "@/lib/api/workspace-helpers"; +import { isAllowedAssetUrl } from "@/lib/tasks/validate-asset-url"; export const dynamic = "force-dynamic"; @@ -11,106 +15,70 @@ export const dynamic = "force-dynamic"; * Receives an audio file URL, runs a durable workflow to download, upload to Gemini, * and transcribe. Returns structured transcript + summary. */ -export async function POST(req: NextRequest) { - try { - const session = await auth.api.getSession({ - headers: await headers(), - }); - if (!session) { - return NextResponse.json({ error: "Unauthorized" }, { status: 401 }); - } +async function handlePOST(req: NextRequest) { + const userId = await requireAuth(); - if (!process.env.GOOGLE_GENERATIVE_AI_API_KEY) { - return NextResponse.json( - { error: "GOOGLE_GENERATIVE_AI_API_KEY is not set" }, - { status: 500 } - ); - } - - const body = await req.json(); - const { fileUrl, filename, mimeType, itemId, workspaceId } = body; - - if (!fileUrl) { - return NextResponse.json( - { error: "fileUrl is required" }, - { status: 400 } - ); - } + if (!process.env.GOOGLE_GENERATIVE_AI_API_KEY) { + return NextResponse.json( + { error: "GOOGLE_GENERATIVE_AI_API_KEY is not set" }, + { status: 500 }, + ); + } - if (!itemId || typeof itemId !== "string") { - return NextResponse.json( - { error: "itemId is required for polling" }, - { status: 400 } - ); - } + const body = await req.json(); + const { fileUrl, filename, mimeType, itemId, workspaceId } = body; - if (!workspaceId || typeof workspaceId !== "string") { - return NextResponse.json( - { error: "workspaceId is required" }, - { status: 400 } - ); - } + if (!fileUrl || typeof fileUrl !== "string") { + return NextResponse.json( + { error: "fileUrl is required" }, + { status: 400 }, + ); + } - // Validate URL origin to prevent SSRF - const allowedHosts: string[] = []; - if (process.env.NEXT_PUBLIC_SUPABASE_URL) { - allowedHosts.push(new URL(process.env.NEXT_PUBLIC_SUPABASE_URL).hostname); - } - if (process.env.NEXT_PUBLIC_APP_URL) { - allowedHosts.push(new URL(process.env.NEXT_PUBLIC_APP_URL).hostname); - } - if (process.env.NODE_ENV === "development") { - allowedHosts.push("localhost"); // local storage in dev - } + if (!itemId || typeof itemId !== "string") { + return NextResponse.json( + { error: "itemId is required for polling" }, + { status: 400 }, + ); + } - let parsedUrl: URL; - try { - parsedUrl = new URL(fileUrl); - } catch { - return NextResponse.json({ error: "Invalid fileUrl" }, { status: 400 }); - } + if (!workspaceId || typeof workspaceId !== "string") { + return NextResponse.json( + { error: "workspaceId is required" }, + { status: 400 }, + ); + } - if ( - !allowedHosts.some( - (host) => - parsedUrl.hostname === host || parsedUrl.hostname.endsWith(`.${host}`) - ) - ) { - return NextResponse.json( - { error: "fileUrl origin is not allowed" }, - { status: 400 } - ); - } + if (!isAllowedAssetUrl(fileUrl)) { + return NextResponse.json( + { error: "fileUrl origin is not allowed" }, + { status: 400 }, + ); + } - const audioMimeType = mimeType || guessMimeType(filename || fileUrl); + await verifyWorkspaceAccess(workspaceId, userId, "editor"); - const userId = session.user.id; + const audioMimeType = mimeType || guessMimeType(filename || fileUrl); - // Start durable workflow; return immediately for client to poll - const run = await start(audioTranscribeWorkflow, [ - fileUrl, - audioMimeType, - workspaceId, - itemId, - userId, - ]); + const run = await start(audioTranscribeWorkflow, [ + fileUrl, + audioMimeType, + workspaceId, + itemId, + userId, + ]); - return NextResponse.json({ - runId: run.runId, - itemId, - }); - } catch (error: unknown) { - console.error("[AUDIO_PROCESS] Error:", error); - return NextResponse.json( - { - error: - error instanceof Error ? error.message : "Failed to process audio", - }, - { status: 500 } - ); - } + return NextResponse.json({ + runId: run.runId, + itemId, + }); } +export const POST = withErrorHandling( + handlePOST, + "POST /api/audio/process", +); + function guessMimeType(filenameOrUrl: string): string { const lower = filenameOrUrl.toLowerCase(); if (lower.endsWith(".mp3")) return "audio/mp3"; diff --git a/src/app/api/workspaces/[id]/events/route.ts b/src/app/api/workspaces/[id]/events/route.ts index aa5fd235..b562ce66 100644 --- a/src/app/api/workspaces/[id]/events/route.ts +++ b/src/app/api/workspaces/[id]/events/route.ts @@ -11,7 +11,7 @@ import { hasDuplicateName } from "@/lib/workspace/unique-name"; import { db, workspaceEvents } from "@/lib/db/client"; import { eq, sql } from "drizzle-orm"; import { - requireAuth, + requireAuthWithUserInfo, verifyWorkspaceAccess, withErrorHandling, } from "@/lib/api/workspace-helpers"; @@ -25,11 +25,11 @@ async function handleGET( { params }: { params: Promise<{ id: string }> }, ) { const paramsPromise = params; - const authPromise = requireAuth(); + const authPromise = requireAuthWithUserInfo(); const paramsResolved = await paramsPromise; const id = paramsResolved.id; - const userId = await authPromise; + const { userId } = await authPromise; await verifyWorkspaceAccess(id, userId, "viewer"); const eventCountResult = await db .select({ count: sql`count(*)::int` }) @@ -153,36 +153,45 @@ async function handlePOST( const startTime = Date.now(); const timings: Record = {}; - // Start independent operations in parallel const paramsPromise = params; - const authPromise = requireAuth(); + const authPromise = requireAuthWithUserInfo(); const bodyPromise = request.json(); const paramsResolved = await paramsPromise; const id = paramsResolved.id; const authStart = Date.now(); - const userId = await authPromise; + const { userId, name: userName } = await authPromise; timings.auth = Date.now() - authStart; const bodyStart = Date.now(); const body = await bodyPromise; timings.bodyParse = Date.now() - bodyStart; - const { event, baseVersion } = body; - if (!event || baseVersion === undefined || isNaN(baseVersion)) { + const { event: rawEvent, baseVersion } = body ?? {}; + if ( + !rawEvent || + typeof rawEvent.type !== "string" || + typeof rawEvent.id !== "string" || + baseVersion === undefined || + isNaN(baseVersion) + ) { return NextResponse.json( { error: "Event and valid baseVersion are required" }, { status: 400 }, ); } - // Check if user has editor access (owner or editor collaborator) + const event: WorkspaceEvent = { + ...rawEvent, + userId, + userName, + }; + const workspaceCheckStart = Date.now(); await verifyWorkspaceAccess(id, userId, "editor"); timings.workspaceCheck = Date.now() - workspaceCheckStart; - // Validate unique name for ITEM_CREATED and ITEM_UPDATED (when name changes) if (event.type === "ITEM_CREATED") { const item = event.payload?.item; if (item?.name != null && item?.type) { @@ -201,13 +210,13 @@ async function handlePOST( } if (event.type === "ITEM_UPDATED" && event.payload?.changes?.name != null) { const itemId = event.payload?.id; - const newName = event.payload.changes.name; + const newName = event.payload.changes.name as string; if (itemId && newName) { const state = await loadWorkspaceState(id, { userId }); const existingItem = state.find((i: { id: string }) => i.id === itemId); if (existingItem) { const newFolderId = - event.payload.changes.folderId ?? existingItem.folderId ?? null; + (event.payload.changes as Record).folderId as string | undefined ?? existingItem.folderId ?? null; if ( hasDuplicateName( state, diff --git a/src/hooks/ai/use-create-card-from-message.ts b/src/hooks/ai/use-create-card-from-message.ts index 53e6fa34..2d16b980 100644 --- a/src/hooks/ai/use-create-card-from-message.ts +++ b/src/hooks/ai/use-create-card-from-message.ts @@ -1,6 +1,7 @@ "use client"; -import { useState, useCallback, useRef } from "react"; +import { useCallback } from "react"; +import { useAsyncDebouncer } from "@tanstack/react-pacer/async-debouncer"; import { useMessage, useThread } from "@assistant-ui/react"; import { toast } from "sonner"; import { useWorkspaceStore } from "@/lib/stores/workspace-store"; @@ -20,9 +21,6 @@ interface CreateCardOptions { */ export function useCreateCardFromMessage(options: CreateCardOptions = {}) { const { debounceMs = 300 } = options; // Reduced from 1000ms to 300ms - const [isCreating, setIsCreating] = useState(false); - const debounceTimerRef = useRef(null); - const message = useMessage(); const thread = useThread(); // Access the full thread to find sources @@ -66,21 +64,8 @@ export function useCreateCardFromMessage(options: CreateCardOptions = {}) { }); }; - const createCard = useCallback(async () => { - // Clear any existing debounce timer - if (debounceTimerRef.current) { - clearTimeout(debounceTimerRef.current); - } - - // Set up debounced execution - debounceTimerRef.current = setTimeout(async () => { - // Prevent creation if already in progress - if (isCreating) { - toast.error("Card creation already in progress"); - return; - } - - // Get message content + const createCardDebouncer = useAsyncDebouncer( + useCallback(async () => { const content = message.content .filter((part) => part.type === "text") .map((part) => (part as any).text) @@ -96,21 +81,16 @@ export function useCreateCardFromMessage(options: CreateCardOptions = {}) { return; } - setIsCreating(true); const toastId = toast.loading("Creating card..."); try { - // Get the current active folder ID const activeFolderId = useUIStore.getState().activeFolderId; - // Extract sources from the thread history - // We look backwards from the current message to find the relevant webSearch result let sources: | Array<{ title: string; url: string; favicon?: string }> | undefined; try { - // Use any cast since standard types might differ const messages = (thread as any).messages || []; const currentIndex = messages.findIndex( (m: any) => m.id === message.id, @@ -167,20 +147,14 @@ export function useCreateCardFromMessage(options: CreateCardOptions = {}) { const result = await response.json(); - // Invalidate React Query cache to refresh the UI immediately - if (currentWorkspaceId) { - logger.debug("🔄 [CREATE-CARD-BUTTON] Invalidating workspace cache", { - workspaceId: currentWorkspaceId.substring(0, 8), - }); - - // Force refetch workspace events to show the new card - queryClient.invalidateQueries({ - queryKey: ["workspace", currentWorkspaceId, "events"], - }); - } + logger.debug("🔄 [CREATE-CARD-BUTTON] Invalidating workspace cache", { + workspaceId: currentWorkspaceId.substring(0, 8), + }); + queryClient.invalidateQueries({ + queryKey: ["workspace", currentWorkspaceId, "events"], + }); toast.success("Card created successfully!", { id: toastId }); - return result; } catch (error) { console.error("Error creating card:", error); @@ -188,19 +162,32 @@ export function useCreateCardFromMessage(options: CreateCardOptions = {}) { error instanceof Error ? error.message : "Failed to create card", { id: toastId }, ); - throw error; - } finally { - setIsCreating(false); } - }, debounceMs); - }, [message, thread, currentWorkspaceId, isCreating, debounceMs]); + }, [message, thread, currentWorkspaceId, queryClient]), + { wait: debounceMs }, + (state) => ({ + isExecuting: state.isExecuting, + isPending: state.isPending, + }), + ); - // Cleanup on unmount - const cleanup = useCallback(() => { - if (debounceTimerRef.current) { - clearTimeout(debounceTimerRef.current); + const createCard = useCallback(() => { + if (createCardDebouncer.state.isExecuting) { + toast.error("Card creation already in progress"); + return; } - }, []); + + void createCardDebouncer.maybeExecute(); + }, [createCardDebouncer]); + + const cleanup = useCallback(() => { + createCardDebouncer.cancel(); + createCardDebouncer.abort(); + }, [createCardDebouncer]); + + const isCreating = + createCardDebouncer.state.isExecuting || + createCardDebouncer.state.isPending; return { createCard, diff --git a/src/hooks/ui/use-folder-url.ts b/src/hooks/ui/use-folder-url.ts index 0ec315fd..bd7c8fd7 100644 --- a/src/hooks/ui/use-folder-url.ts +++ b/src/hooks/ui/use-folder-url.ts @@ -1,6 +1,7 @@ "use client"; -import { useEffect, useRef } from "react"; +import { useCallback, useEffect, useRef } from "react"; +import { useDebouncer } from "@tanstack/react-pacer/debouncer"; import { useSearchParams, useRouter, usePathname } from "next/navigation"; import { useUIStore } from "@/lib/stores/ui-store"; @@ -38,11 +39,39 @@ export function useFolderUrl() { const activeFolderId = useUIStore((state) => state.activeFolderId); const openItems = useUIStore((state) => state.openItems); - const setActiveFolderIdDirect = useUIStore((state) => state._setActiveFolderIdDirect); + const setActiveFolderIdDirect = useUIStore( + (state) => state._setActiveFolderIdDirect, + ); const setOpenItemsFromUrl = useUIStore((state) => state._setOpenItemsFromUrl); const isSyncingFromUrl = useRef(false); const lastPushedState = useRef(undefined); + const pushUrlState = useCallback( + (next: UrlState) => { + lastPushedState.current = next; + + const params = new URLSearchParams(searchParams.toString()); + + if (next.folder) params.set("folder", next.folder); + else params.delete("folder"); + + if (next.items.length > 0) { + params.set("items", next.items.join(",")); + } else { + params.delete("items"); + } + + params.delete("focus"); + params.delete("item"); + + const qs = params.toString(); + router.push(qs ? `${pathname}?${qs}` : pathname, { scroll: false }); + }, + [pathname, router, searchParams], + ); + const pushUrlDebouncer = useDebouncer(pushUrlState, { + wait: DEBOUNCE_MS, + }); const parseUrlState = (): UrlState => { const folder = searchParams.get("folder") || null; @@ -85,6 +114,7 @@ export function useFolderUrl() { const storeItems = openItemsToUrlList(openItems); if (isSyncingFromUrl.current) { + pushUrlDebouncer.cancel(); lastPushedState.current = { folder: activeFolderId, items: storeItems, @@ -107,28 +137,12 @@ export function useFolderUrl() { return; } - const tid = setTimeout(() => { - lastPushedState.current = next; - - const params = new URLSearchParams(searchParams.toString()); - - if (next.folder) params.set("folder", next.folder); - else params.delete("folder"); - - if (next.items.length > 0) { - params.set("items", next.items.join(",")); - } else { - params.delete("items"); - } - - params.delete("focus"); - params.delete("item"); - - const qs = params.toString(); - router.push(qs ? `${pathname}?${qs}` : pathname, { scroll: false }); - }, DEBOUNCE_MS); - - return () => clearTimeout(tid); + pushUrlDebouncer.maybeExecute(next); // eslint-disable-next-line react-hooks/exhaustive-deps - }, [activeFolderId, openItems.primary, openItems.secondary]); + }, [ + activeFolderId, + openItems.primary, + openItems.secondary, + pushUrlDebouncer, + ]); } diff --git a/src/hooks/workspace/use-workspace-mutation.ts b/src/hooks/workspace/use-workspace-mutation.ts index 2fe15257..d43e9bfa 100644 --- a/src/hooks/workspace/use-workspace-mutation.ts +++ b/src/hooks/workspace/use-workspace-mutation.ts @@ -1,42 +1,27 @@ import { useMutation, useQueryClient } from "@tanstack/react-query"; import type { WorkspaceEvent, EventResponse } from "@/lib/workspace/events"; +import { + computeBaseVersion, + removeOptimisticEvent, + confirmOptimisticEvent, +} from "@/lib/workspace/version-helpers"; import { workspaceEventsQueryKey } from "./use-workspace-events"; import { applyConfirmedWorkspaceEventToStateQuery } from "./workspace-state-cache"; -import { logger } from "@/lib/utils/logger"; import { useRef } from "react"; import { toast } from "sonner"; -/** - * Maximum number of automatic retries for version conflicts - * This prevents infinite retry loops while allowing concurrent edits to succeed - */ const MAX_RETRY_ATTEMPTS = 3; -interface AppendEventParams { +type AppendEventResponse = + | { success: true; version: number; conflict: false } + | { conflict: true; version: number }; + +async function appendWorkspaceEvent(params: { workspaceId: string; event: WorkspaceEvent; baseVersion: number; -} - -/** - * Append event to workspace event log - */ -async function appendWorkspaceEvent({ - workspaceId, - event, - baseVersion, -}: AppendEventParams): Promise<{ - success: boolean; - version: number; - conflict?: boolean; -}> { - logger.debug("📤 [API] Appending event:", { - workspaceId, - eventType: event.type, - baseVersion, - userId: event.userId, - }); - +}): Promise { + const { workspaceId, event, baseVersion } = params; const response = await fetch(`/api/workspaces/${workspaceId}/events`, { method: "POST", headers: { "Content-Type": "application/json" }, @@ -48,260 +33,105 @@ async function appendWorkspaceEvent({ throw new Error("PERMISSION_DENIED"); } const errorText = await response.text(); - logger.error("❌ [API] Response error:", response.status, errorText); throw new Error( `Failed to append event: ${response.statusText} - ${errorText}`, ); } - const result = await response.json(); - logger.debug("✅ [API] Event appended successfully:", result); - return result; + return response.json(); } /** - * Hook to mutate workspace by appending events - * Implements optimistic updates with automatic rollback on error + * Hook to mutate workspace by appending events. + * Implements optimistic updates with automatic conflict retry. */ export function useWorkspaceMutation(workspaceId: string | null) { const queryClient = useQueryClient(); - - // Track retry attempts per event ID to prevent infinite loops const retryAttemptsRef = useRef>(new Map()); return useMutation({ mutationFn: (event: WorkspaceEvent) => { if (!workspaceId) { - logger.error("❌ [MUTATION] No workspace ID!"); throw new Error("No workspace ID provided"); } - // Get current version from cache at mutation time - const currentData = queryClient.getQueryData( + const cacheData = queryClient.getQueryData( workspaceEventsQueryKey(workspaceId), ); - // Calculate the actual max version from all events in cache - // This is critical because after tool completions, events may have versions - // that are higher than the cache's version field - const events = currentData?.events ?? []; - const optimisticEventsCount = events.filter( - (e) => typeof e.version !== "number", - ).length; - - // Find the max version from all events that have versions - // This accounts for tool events that were added with versions - const maxEventVersion = events - .filter((e) => typeof e.version === "number") - .reduce( - (max, e) => Math.max(max, e.version!), - currentData?.version ?? 0, - ); - - // Use the higher of: cache version or max event version - // This ensures we account for tool events that updated individual event versions - // but might not have updated the cache version field - const currentVersion = Math.max( - currentData?.version ?? 0, - maxEventVersion, - ); - - // CRITICAL FIX: Account for optimistic events when calculating baseVersion - // If there are optimistic events (pending mutations), they will increment the server version - // So we need to use a baseVersion that accounts for those pending increments - // Note: onMutate runs before mutationFn, so optimisticEventsCount includes our own - // optimistic event plus any other pending mutations. We subtract 1 to exclude our - // own event and only account for other mutations that will complete before this one. - // Example: If currentVersion=169 and optimisticEventsCount=2 (ours + 1 other), - // then 1 other mutation will complete first, making server version 170 - // So we should use baseVersion = 169 + (2-1) = 170 - const adjustedBaseVersion = - currentVersion + Math.max(0, optimisticEventsCount - 1); - - logger.debug("🚀 [MUTATION] Starting mutation:", { - workspaceId, - cacheVersion: currentData?.version ?? 0, - maxEventVersion, - currentVersion, - adjustedBaseVersion, - optimisticEventsCount, - eventType: event.type, - eventId: event.id, - }); + const { baseVersion } = computeBaseVersion(cacheData); return appendWorkspaceEvent({ workspaceId, event, - baseVersion: adjustedBaseVersion, + baseVersion, }); }, - // Optimistic update - apply event immediately to UI onMutate: async (event: WorkspaceEvent) => { if (!workspaceId) return; - // Cancel any outgoing refetches to avoid overwriting optimistic update await queryClient.cancelQueries({ queryKey: workspaceEventsQueryKey(workspaceId), }); - // Snapshot the previous value for rollback const previous = queryClient.getQueryData( workspaceEventsQueryKey(workspaceId), ); - logger.debug("⚡ [OPTIMISTIC] Applying event optimistically:", { - eventType: event.type, - currentVersion: previous?.version, - }); - - // Optimistically update to the new value - // NOTE: Don't increment version here - let server confirm the version queryClient.setQueryData( workspaceEventsQueryKey(workspaceId), (old) => { if (!old) { - const newState = { - events: [{ ...event }], // Don't assign version to optimistic events - version: 0, // Keep server version, don't increment optimistically - }; - return newState; + return { events: [{ ...event }], version: 0 }; } - - const newState = { + return { ...old, - events: [...old.events, { ...event }], // Don't assign version to optimistic events - // Keep the same version - server will tell us the new version + events: [...old.events, { ...event }], version: old.version, }; - return newState; }, ); - // Return context with previous value for rollback return { previous }; }, - // Rollback on error onError: (err, event, context) => { - logger.error("❌ [MUTATION] Event mutation failed:", err); - logger.error("❌ [MUTATION] Failed event:", event); - logger.error("❌ [MUTATION] Error details:", { - message: err instanceof Error ? err.message : String(err), - workspaceId, - }); - if (!workspaceId || !context?.previous) return; - // Show toast for permission errors if (err.message === "PERMISSION_DENIED") { toast.error("You don't have permission to edit this workspace"); - } else { - // Only show generic error for non-permission issues to avoid noise - // toast.error("Failed to save changes"); } - // Rollback to previous state queryClient.setQueryData( workspaceEventsQueryKey(workspaceId), context.previous, ); }, - // Refetch to ensure consistency on success onSuccess: (data, event, context) => { if (!workspaceId) return; - logger.debug("✅ [SUCCESS] Mutation succeeded:", { - conflict: data.conflict, - newVersion: data.version, - eventId: event.id, - }); - if (event.type === "ITEM_UPDATED") { - const payload = event.payload; - const changesData = payload?.changes?.data; - if ( - changesData && - "ocrStatus" in changesData && - changesData.ocrStatus - ) { - logger.debug("[OCR/UPDATE] ITEM_UPDATED persisted", { - itemId: payload.id, - ocrStatus: changesData.ocrStatus, - version: data.version, - }); - } - } - if (event.type === "BULK_ITEMS_PATCHED") { - const ocrUpdates = event.payload.updates.filter((update) => { - const dataChanges = update.changes?.data as - | { ocrStatus?: string } - | undefined; - return !!dataChanges?.ocrStatus; - }); - if (ocrUpdates.length > 0) { - logger.debug("[OCR/UPDATE] BULK_ITEMS_PATCHED persisted", { - itemIds: ocrUpdates.map((update) => update.id), - statuses: ocrUpdates.map( - (update) => - (update.changes?.data as { ocrStatus?: string } | undefined) - ?.ocrStatus, - ), - version: data.version, - }); - } - } - - // Handle conflicts with automatic retry if (data.conflict) { - // Check retry count for this event const currentRetries = retryAttemptsRef.current.get(event.id) || 0; if (currentRetries < MAX_RETRY_ATTEMPTS) { - logger.warn( - `⚠️ [CONFLICT] Version conflict detected, auto-retrying (attempt ${currentRetries + 1}/${MAX_RETRY_ATTEMPTS})...`, - ); - - // Increment retry counter retryAttemptsRef.current.set(event.id, currentRetries + 1); - // First, remove this specific optimistic event from cache queryClient.setQueryData( workspaceEventsQueryKey(workspaceId), - (old) => { - if (!old) return old; - - // Only remove the conflicting event, preserve other pending optimistic events - const filteredEvents = old.events.filter( - (e) => e.id !== event.id, - ); - - return { - ...old, - events: filteredEvents, - version: data.version, // Update to current server version - }; - }, + (old) => removeOptimisticEvent(old, event.id, data.version), ); - // Refetch to get latest events, then automatically retry queryClient .invalidateQueries({ queryKey: workspaceEventsQueryKey(workspaceId), }) .then(() => { - // After cache is updated with latest events, retry the mutation - logger.debug( - "🔄 [RETRY] Retrying event after refetch:", - event.type, - ); - - // Re-apply optimistic update queryClient.setQueryData( workspaceEventsQueryKey(workspaceId), (old) => { if (!old) return old; - return { ...old, events: [...old.events, { ...event }], @@ -309,178 +139,81 @@ export function useWorkspaceMutation(workspaceId: string | null) { }, ); - // Retry the mutation with updated base version const currentData = queryClient.getQueryData([ "workspace", workspaceId, "events", ]); - if (currentData) { - const events = currentData.events ?? []; - const maxEventVersion = events - .filter((e) => typeof e.version === "number") - .reduce( - (max, e) => Math.max(max, e.version!), - currentData.version ?? 0, - ); - const currentVersion = Math.max( - currentData.version ?? 0, - maxEventVersion, - ); - - logger.debug("🔄 [RETRY] Using base version:", currentVersion); + if (!currentData) return; - // Call API directly to retry - appendWorkspaceEvent({ - workspaceId, - event, - baseVersion: currentVersion, - }) - .then((retryResult) => { - if (retryResult.conflict) { - // Still conflicting after retry - treat as final failure - logger.error( - "❌ [RETRY] Retry failed with conflict, giving up", - ); - - // Remove optimistic event - queryClient.setQueryData( - workspaceEventsQueryKey(workspaceId), - (old) => { - if (!old) return old; - return { - ...old, - events: old.events.filter((e) => e.id !== event.id), - }; - }, - ); - - // Clean up retry counter - retryAttemptsRef.current.delete(event.id); - - // Force full refetch - queryClient.invalidateQueries({ - queryKey: workspaceEventsQueryKey(workspaceId), - }); - } else { - // Retry succeeded! - logger.debug( - "✅ [RETRY] Retry succeeded, version:", - retryResult.version, - ); - - // Update event with version - queryClient.setQueryData( - workspaceEventsQueryKey(workspaceId), - (old) => { - if (!old) return old; + const { baseVersion: retryBase } = + computeBaseVersion(currentData); - const updatedEvents = old.events.map((e) => - e.id === event.id - ? { ...e, version: retryResult.version } - : e, - ); - - return { - ...old, - events: updatedEvents, - version: retryResult.version, - }; - }, - ); - - // Clean up retry counter - retryAttemptsRef.current.delete(event.id); - applyConfirmedWorkspaceEventToStateQuery( - queryClient, - workspaceId, - { - ...event, - version: retryResult.version, - }, - ); - } - }) - .catch((err) => { - logger.error("❌ [RETRY] Retry failed with error:", err); - - // Remove optimistic event + appendWorkspaceEvent({ + workspaceId, + event, + baseVersion: retryBase, + }) + .then((retryResult) => { + if (retryResult.conflict) { queryClient.setQueryData( workspaceEventsQueryKey(workspaceId), - (old) => { - if (!old) return old; - return { - ...old, - events: old.events.filter((e) => e.id !== event.id), - }; - }, + (old) => removeOptimisticEvent(old, event.id), ); - - // Clean up retry counter retryAttemptsRef.current.delete(event.id); - - // Restore from context if available - if (context?.previous) { - queryClient.setQueryData( - workspaceEventsQueryKey(workspaceId), - context.previous, - ); - } - }); - } + queryClient.invalidateQueries({ + queryKey: workspaceEventsQueryKey(workspaceId), + }); + } else { + queryClient.setQueryData( + workspaceEventsQueryKey(workspaceId), + (old) => + confirmOptimisticEvent( + old, + event.id, + retryResult.version, + ), + ); + retryAttemptsRef.current.delete(event.id); + applyConfirmedWorkspaceEventToStateQuery( + queryClient, + workspaceId, + { ...event, version: retryResult.version }, + ); + } + }) + .catch(() => { + queryClient.setQueryData( + workspaceEventsQueryKey(workspaceId), + (old) => removeOptimisticEvent(old, event.id), + ); + retryAttemptsRef.current.delete(event.id); + if (context?.previous) { + queryClient.setQueryData( + workspaceEventsQueryKey(workspaceId), + context.previous, + ); + } + }); }); } else { - // Max retries exceeded - logger.error( - `❌ [CONFLICT] Max retries (${MAX_RETRY_ATTEMPTS}) exceeded for event ${event.id}`, - ); - - // Remove optimistic event queryClient.setQueryData( workspaceEventsQueryKey(workspaceId), - (old) => { - if (!old) return old; - return { - ...old, - events: old.events.filter((e) => e.id !== event.id), - }; - }, + (old) => removeOptimisticEvent(old, event.id), ); - - // Clean up retry counter retryAttemptsRef.current.delete(event.id); - - // Force full refetch to get true server state queryClient.invalidateQueries({ queryKey: workspaceEventsQueryKey(workspaceId), }); } } else { - // Success! No conflict - // Clean up retry counter if it exists retryAttemptsRef.current.delete(event.id); - // Update the version in cache without full refetch queryClient.setQueryData( workspaceEventsQueryKey(workspaceId), - (old) => { - if (!old) return old; - - // Find the specific optimistic event by ID and assign the server version - const updatedEvents = old.events.map((e) => - e.id === event.id ? { ...e, version: data.version } : e, - ); - - const newState = { - ...old, - events: updatedEvents, - version: data.version, // Update to server-confirmed version - }; - return newState; - }, + (old) => confirmOptimisticEvent(old, event.id, data.version), ); - logger.debug("✅ [SUCCESS] Version updated to:", data.version); applyConfirmedWorkspaceEventToStateQuery(queryClient, workspaceId, { ...event, version: data.version, diff --git a/src/hooks/workspace/use-workspace-operations.ts b/src/hooks/workspace/use-workspace-operations.ts index 86df17df..16ab9ab8 100644 --- a/src/hooks/workspace/use-workspace-operations.ts +++ b/src/hooks/workspace/use-workspace-operations.ts @@ -1,4 +1,5 @@ import { useCallback, useRef, useEffect } from "react"; +import { Debouncer } from "@tanstack/pacer/debouncer"; import { useSession } from "@/lib/auth-client"; import { toast } from "sonner"; import { useQueryClient } from "@tanstack/react-query"; @@ -31,6 +32,8 @@ import { } from "@/lib/workspace/unique-name"; import { filterItemIdsForFolderCreation } from "@/lib/workspace-state/search"; +const ITEM_UPDATE_DEBOUNCE_MS = 500; + function getAllDescendantIds(folderId: string, items: Item[]): string[] { const directChildren = items.filter((item) => item.folderId === folderId); const descendantIds: string[] = []; @@ -108,10 +111,11 @@ export function useWorkspaceOperations( const userId = user?.id || "anonymous"; const userName = user?.name || user?.email || undefined; - // Debounce refs for updateItem and updateItemData - // Maps item ID to timeout ID - const updateItemDebounceRef = useRef>(new Map()); - const updateItemDataDebounceRef = useRef>( + // Debouncer refs for updateItem and updateItemData keyed by item ID. + const updateItemDebouncerRef = useRef void>>>( + new Map(), + ); + const updateItemDataDebouncerRef = useRef void>>>( new Map(), ); // Store pending changes to merge them @@ -167,18 +171,169 @@ export function useWorkspaceOperations( [getLatestItemFromState], ); + const commitPendingItemUpdate = useCallback( + (itemId: string) => { + const finalChanges = pendingItemChangesRef.current.get(itemId); + if (!finalChanges) { + updateItemDebouncerRef.current.delete(itemId); + return; + } + + const latestItems = getLatestItemsFromState(); + const item = latestItems.find((candidate) => candidate.id === itemId); + if (!item) { + pendingItemChangesRef.current.delete(itemId); + updateItemDebouncerRef.current.delete(itemId); + return; + } + + const newName = finalChanges.name ?? item.name; + const newType = finalChanges.type ?? item.type; + const folderId = finalChanges.folderId ?? item.folderId ?? null; + + if (newName && newType && "name" in finalChanges) { + if (hasDuplicateName(latestItems, newName, newType, folderId, itemId)) { + toast.error( + `A ${newType} named "${newName}" already exists in this folder`, + ); + pendingItemChangesRef.current.delete(itemId); + updateItemDebouncerRef.current.delete(itemId); + return; + } + } + + mutation.mutate( + createEvent( + "ITEM_UPDATED", + { id: itemId, changes: finalChanges, name: newName }, + userId, + userName, + ), + ); + + pendingItemChangesRef.current.delete(itemId); + updateItemDebouncerRef.current.delete(itemId); + }, + [getLatestItemsFromState, mutation, userId, userName], + ); + + const commitPendingItemDataUpdate = useCallback( + (itemId: string) => { + const finalUpdater = pendingItemDataUpdatersRef.current.get(itemId); + if (!finalUpdater) { + updateItemDataDebouncerRef.current.delete(itemId); + return; + } + + // Read latest state from cache instead of relying on closure state so + // newly-created items remain addressable while optimistic events settle. + let latestItem: Item | undefined; + if (workspaceId) { + const cacheData = queryClient.getQueryData( + workspaceEventsQueryKey(workspaceId), + ); + if (cacheData?.events) { + const latestState = deriveWorkspaceStateFromCaches({ + workspaceId, + stateData: getCachedWorkspaceState(queryClient, workspaceId), + eventLog: cacheData, + }).state; + latestItem = latestState.find((item) => item.id === itemId); + } + } + if (!latestItem) { + latestItem = currentItems.find((item) => item.id === itemId); + } + if (!latestItem) { + const cacheData = workspaceId + ? queryClient.getQueryData( + workspaceEventsQueryKey(workspaceId), + ) + : null; + const itemCount = cacheData?.events + ? deriveWorkspaceStateFromCaches({ + workspaceId, + stateData: getCachedWorkspaceState(queryClient, workspaceId), + eventLog: cacheData, + }).state.length + : 0; + logger.warn( + `[OCR/UPDATE] updateItemData: Item ${itemId} not found. Item may have been deleted. Cache has ${itemCount} items.`, + { + itemId, + workspaceId, + }, + ); + pendingItemDataUpdatersRef.current.delete(itemId); + updateItemDataDebouncerRef.current.delete(itemId); + return; + } + + const newData = finalUpdater(latestItem.data); + + mutation.mutate( + createEvent( + "ITEM_UPDATED", + { + id: itemId, + changes: { data: newData }, + name: latestItem.name, + }, + userId, + userName, + ), + ); + + pendingItemDataUpdatersRef.current.delete(itemId); + updateItemDataDebouncerRef.current.delete(itemId); + }, + [workspaceId, queryClient, currentItems, mutation, userId, userName], + ); + + const getOrCreateUpdateItemDebouncer = useCallback( + (itemId: string) => { + const existing = updateItemDebouncerRef.current.get(itemId); + if (existing) { + return existing; + } + + const debouncer = new Debouncer(() => commitPendingItemUpdate(itemId), { + wait: ITEM_UPDATE_DEBOUNCE_MS, + }); + updateItemDebouncerRef.current.set(itemId, debouncer); + return debouncer; + }, + [commitPendingItemUpdate], + ); + + const getOrCreateUpdateItemDataDebouncer = useCallback( + (itemId: string) => { + const existing = updateItemDataDebouncerRef.current.get(itemId); + if (existing) { + return existing; + } + + const debouncer = new Debouncer( + () => commitPendingItemDataUpdate(itemId), + { wait: ITEM_UPDATE_DEBOUNCE_MS }, + ); + updateItemDataDebouncerRef.current.set(itemId, debouncer); + return debouncer; + }, + [commitPendingItemDataUpdate], + ); + // Cleanup timeouts on unmount useEffect(() => { - const updateItemDebounces = updateItemDebounceRef.current; - const updateItemDataDebounces = updateItemDataDebounceRef.current; + const updateItemDebouncers = updateItemDebouncerRef.current; + const updateItemDataDebouncers = updateItemDataDebouncerRef.current; const pendingItemChanges = pendingItemChangesRef.current; const pendingItemDataUpdaters = pendingItemDataUpdatersRef.current; return () => { - // Clear all pending timeouts - updateItemDebounces.forEach((timeout) => clearTimeout(timeout)); - updateItemDataDebounces.forEach((timeout) => clearTimeout(timeout)); - updateItemDebounces.clear(); - updateItemDataDebounces.clear(); + updateItemDebouncers.forEach((debouncer) => debouncer.cancel()); + updateItemDataDebouncers.forEach((debouncer) => debouncer.cancel()); + updateItemDebouncers.clear(); + updateItemDataDebouncers.clear(); pendingItemChanges.clear(); pendingItemDataUpdaters.clear(); }; @@ -191,16 +346,6 @@ export function useWorkspaceOperations( initialData?: Partial, initialLayout?: { w: number; h: number }, ) => { - // Validate type is a valid CardType - logger.debug( - "🔧 [CREATE-ITEM] Received type:", - type, - "typeof:", - typeof type, - "value:", - JSON.stringify(type), - ); - const validTypes: CardType[] = [ "pdf", "flashcard", @@ -214,19 +359,6 @@ export function useWorkspaceOperations( ]; const validType = validTypes.includes(type) ? type : "document"; - if (validType !== type) { - logger.warn( - `🔧 [CREATE-ITEM] Invalid type "${type}" (typeof ${typeof type}), using "document" instead`, - ); - } - - logger.debug("🔧 [CREATE-ITEM] Creating item:", { - type: validType, - name, - userId, - workspaceId, - hasInitialData: !!initialData, - }); const id = generateItemId(); // Merge default data with initial data @@ -235,13 +367,7 @@ export function useWorkspaceOperations( ? { ...baseData, ...initialData } : baseData; - // Get active folder - auto-assign new items to the currently viewed folder const activeFolderId = useUIStore.getState().activeFolderId; - logger.debug("🔧 [CREATE-ITEM] Active folder:", { activeFolderId }); - - logger.debug("🔧 [CREATE-ITEM] Base data:", baseData); - logger.debug("🔧 [CREATE-ITEM] Initial data:", initialData); - logger.debug("🔧 [CREATE-ITEM] Merged data:", mergedData); const folderId = activeFolderId ?? null; const finalName = @@ -296,13 +422,6 @@ export function useWorkspaceOperations( return []; } - logger.debug("🔧 [CREATE-ITEMS] Creating items:", { - count: items.length, - userId, - workspaceId, - }); - - // Get active folder - auto-assign new items to the currently viewed folder const activeFolderId = useUIStore.getState().activeFolderId; // Get items in current view for layout calculation @@ -334,12 +453,6 @@ export function useWorkspaceOperations( ]; const validType = validTypes.includes(type) ? type : "document"; - if (validType !== type) { - logger.warn( - `🔧 [CREATE-ITEMS] Invalid type "${type}", using "document" instead`, - ); - } - const id = generateItemId(); const folderId = activeFolderId ?? null; const allItemsSoFar = [...itemsInCurrentView, ...itemsSoFar]; @@ -352,9 +465,6 @@ export function useWorkspaceOperations( name && hasDuplicateName(allItemsSoFar, finalName, validType, folderId) ) { - logger.warn( - `🔧 [CREATE-ITEMS] Skipping duplicate: ${finalName} (${validType})`, - ); return null; } @@ -442,71 +552,16 @@ export function useWorkspaceOperations( const updateItem = useCallback( (id: string, changes: Partial) => { - // Merge with any pending changes for this item const existingPending = pendingItemChangesRef.current.get(id) || {}; const mergedChanges = { ...existingPending, ...changes }; pendingItemChangesRef.current.set(id, mergedChanges); - - // Clear existing debounce for this item - const existingTimeout = updateItemDebounceRef.current.get(id); - if (existingTimeout) { - clearTimeout(existingTimeout); - } - - // Set new debounce (500ms delay) - const timeout = setTimeout(() => { - const finalChanges = pendingItemChangesRef.current.get(id); - if (finalChanges) { - const latestItems = getLatestItemsFromState(); - const item = latestItems.find((i) => i.id === id); - if (!item) { - pendingItemChangesRef.current.delete(id); - updateItemDebounceRef.current.delete(id); - return; - } - - const newName = (finalChanges as Partial).name ?? item?.name; - const newType = (finalChanges as Partial).type ?? item?.type; - const folderId = - (finalChanges as Partial).folderId ?? item?.folderId ?? null; - - if (newName && newType && "name" in finalChanges) { - if (hasDuplicateName(latestItems, newName, newType, folderId, id)) { - toast.error( - `A ${newType} named "${newName}" already exists in this folder`, - ); - pendingItemChangesRef.current.delete(id); - updateItemDebounceRef.current.delete(id); - return; - } - } - - logger.debug("⏱️ [DEBOUNCE] updateItem firing after 500ms:", { - id, - changes: finalChanges, - }); - const event = createEvent( - "ITEM_UPDATED", - { id, changes: finalChanges, name: newName }, - userId, - userName, - ); - mutation.mutate(event); - // Clean up - pendingItemChangesRef.current.delete(id); - updateItemDebounceRef.current.delete(id); - } - }, 500); - - updateItemDebounceRef.current.set(id, timeout); + getOrCreateUpdateItemDebouncer(id).maybeExecute(); }, - [mutation, userId, userName, getLatestItemsFromState], + [getOrCreateUpdateItemDebouncer], ); const deleteItem = useCallback( async (id: string) => { - logger.debug("🗑️ [DELETE-ITEM] Deleting item:", { id, userId, userName }); - // If this is a PDF card, delete the file from Supabase storage first const itemToDelete = currentItems.find((item) => item.id === id); if (itemToDelete && itemToDelete.type === "pdf") { @@ -516,9 +571,6 @@ export function useWorkspaceOperations( }; if (pdfData?.fileUrl) { try { - logger.debug("🗑️ [DELETE-ITEM] Deleting PDF file from Supabase:", { - fileUrl: pdfData.fileUrl, - }); const deleteResponse = await fetch( `/api/delete-file?url=${encodeURIComponent(pdfData.fileUrl)}`, { @@ -530,18 +582,10 @@ export function useWorkspaceOperations( const errorData = await deleteResponse .json() .catch(() => ({ error: "Failed to delete file" })); - logger.warn( - "🗑️ [DELETE-ITEM] Failed to delete PDF file from Supabase:", - errorData.error, - ); - // Continue with card deletion even if file deletion fails (file might not exist) - } else { - logger.debug( - "🗑️ [DELETE-ITEM] Successfully deleted PDF file from Supabase", - ); + logger.warn("Failed to delete PDF file:", errorData.error); } } catch (error) { - logger.error("🗑️ [DELETE-ITEM] Error deleting PDF file:", error); + logger.error("Error deleting PDF file:", error); // Continue with card deletion even if file deletion fails } } @@ -554,7 +598,6 @@ export function useWorkspaceOperations( userId, userName, ); - logger.debug("🗑️ [DELETE-ITEM] Created event:", event); mutation.mutate(event); }, [mutation, userId, userName, currentItems], @@ -563,113 +606,24 @@ export function useWorkspaceOperations( // Helper for updating item data (used by field actions) const updateItemData = useCallback( (itemId: string, updater: (prev: Item["data"]) => Item["data"]) => { - // Chain updaters if there's already a pending one const existingUpdater = pendingItemDataUpdatersRef.current.get(itemId); if (existingUpdater) { - // Chain the updaters: apply existing, then new pendingItemDataUpdatersRef.current.set(itemId, (prev: Item["data"]) => updater(existingUpdater(prev)), ); } else { pendingItemDataUpdatersRef.current.set(itemId, updater); } - - // Clear existing debounce for this item - const existingTimeout = updateItemDataDebounceRef.current.get(itemId); - if (existingTimeout) { - clearTimeout(existingTimeout); - } - - // Set new debounce (500ms delay) - const timeout = setTimeout(() => { - const finalUpdater = pendingItemDataUpdatersRef.current.get(itemId); - if (finalUpdater) { - // CRITICAL: Read latest state from cache (not currentState closure) so we get - // items created after this callback was invoked (e.g. OCR completing after PDF create) - let latestItem: Item | undefined; - if (workspaceId) { - const cacheData = queryClient.getQueryData( - workspaceEventsQueryKey(workspaceId), - ); - if (cacheData?.events) { - const latestState = deriveWorkspaceStateFromCaches({ - workspaceId, - stateData: getCachedWorkspaceState(queryClient, workspaceId), - eventLog: cacheData, - }).state; - latestItem = latestState.find((item) => item.id === itemId); - } - } - if (!latestItem) { - latestItem = currentItems.find((item) => item.id === itemId); - } - if (!latestItem) { - const cacheData = workspaceId - ? queryClient.getQueryData( - workspaceEventsQueryKey(workspaceId), - ) - : null; - const itemCount = cacheData?.events - ? deriveWorkspaceStateFromCaches({ - workspaceId, - stateData: getCachedWorkspaceState(queryClient, workspaceId), - eventLog: cacheData, - }).state.length - : 0; - logger.warn( - `[OCR/UPDATE] updateItemData: Item ${itemId} not found. Item may have been deleted. Cache has ${itemCount} items.`, - { - itemId, - workspaceId, - }, - ); - pendingItemDataUpdatersRef.current.delete(itemId); - updateItemDataDebounceRef.current.delete(itemId); - return; - } - - // Apply the final updater to get new data - const newData = finalUpdater(latestItem.data); - - const ocrStatus = (newData as { ocrStatus?: string })?.ocrStatus; - logger.debug("[OCR/UPDATE] updateItemData firing:", { - itemId, - itemName: latestItem.name, - ocrStatus, - dataKeys: Object.keys(newData), - }); - - // Emit update event with new data and item name - const event = createEvent( - "ITEM_UPDATED", - { - id: itemId, - changes: { data: newData }, - name: latestItem.name, - }, - userId, - userName, - ); - mutation.mutate(event); - // Clean up - pendingItemDataUpdatersRef.current.delete(itemId); - updateItemDataDebounceRef.current.delete(itemId); - } - }, 500); - - updateItemDataDebounceRef.current.set(itemId, timeout); + getOrCreateUpdateItemDataDebouncer(itemId).maybeExecute(); }, - [workspaceId, queryClient, currentItems, mutation, userId, userName], + [getOrCreateUpdateItemDataDebouncer], ); // Update all items at once (used for layout changes, reordering, and bulk delete) const updateAllItems = useCallback( (items: Item[]) => { - logger.debug("🔧 [BULK-UPDATE] Updating all items:", { - count: items.length, - }); - - // CRITICAL FIX: Read the latest state from cache (including optimistic updates) + // Read the latest state from cache (including optimistic updates) + // instead of using the potentially stale currentState prop (including optimistic updates) // instead of using the potentially stale currentState prop // This ensures we compare against the most recent state even when a previous mutation is pending let latestState: Item[]; @@ -702,9 +656,6 @@ export function useWorkspaceOperations( const deletedIds = latestState .filter((i) => !newIds.has(i.id)) .map((i) => i.id); - logger.debug("🔧 [BULK-UPDATE] Bulk delete, sending deletedIds only", { - deletedIds, - }); mutation.mutate( createEvent( "BULK_ITEMS_UPDATED", @@ -719,9 +670,6 @@ export function useWorkspaceOperations( // Count increased: items added – send only the new items (typically 1–2) if (items.length > previousItemCount) { const addedItems = items.filter((i) => !currentIds.has(i.id)); - logger.debug("🔧 [BULK-UPDATE] Items added, sending addedItems only", { - count: addedItems.length, - }); mutation.mutate( createEvent( "BULK_ITEMS_UPDATED", @@ -783,10 +731,6 @@ export function useWorkspaceOperations( userName, ); mutation.mutate(event); - } else { - logger.debug( - "🔧 [BULK-UPDATE] No layout changes detected, skipping event", - ); } }, [workspaceId, queryClient, currentItems, mutation, userId, userName], @@ -805,10 +749,6 @@ export function useWorkspaceOperations( if (color) { updateItem(folderId, { color }); } - logger.debug("📁 [FOLDER-CREATE] Created folder item:", { - folderId, - name, - }); return folderId; }, [createItem, updateItem], @@ -827,21 +767,12 @@ export function useWorkspaceOperations( ); if (safeItemIds.length === 0) { - logger.warn( - "📁 [FOLDER-CREATE-WITH-ITEMS] No valid items after cycle filter, skipping", - ); toast.error( "Cannot create folder: selected items would create a circular reference.", ); return ""; // Return empty string since no folder was created } - logger.debug("📁 [FOLDER-CREATE-WITH-ITEMS] Active folder:", { - activeFolderId, - safeItemIds, - }); - - // Generate folder ID const folderId = generateItemId(); const baseData = defaultDataFor("folder"); @@ -866,12 +797,6 @@ export function useWorkspaceOperations( mutation.mutate(event); - logger.debug("📁 [FOLDER-CREATE-WITH-ITEMS] Created folder with items:", { - folderId, - name, - itemCount: safeItemIds.length, - }); - // Show success toast toast.success( `Folder created with ${safeItemIds.length} item${safeItemIds.length === 1 ? "" : "s"}`, @@ -885,10 +810,6 @@ export function useWorkspaceOperations( // updateFolder now just calls updateItem (folders are items) const updateFolder = useCallback( (folderId: string, changes: Partial) => { - logger.debug("📁 [FOLDER-UPDATE] Updating folder item:", { - folderId, - changes, - }); updateItem(folderId, changes); }, [updateItem], @@ -900,10 +821,6 @@ export function useWorkspaceOperations( const folder = currentItems?.find( (i) => i.id === folderId && i.type === "folder", ); - logger.debug("📁 [FOLDER-DELETE] Deleting folder item:", { - folderId, - folderName: folder?.name, - }); deleteItem(folderId); toast.success( folder ? `Folder "${folder.name}" deleted` : "Folder deleted", @@ -940,11 +857,6 @@ export function useWorkspaceOperations( const folder = latestItems.find( (i) => i.id === folderId && i.type === "folder", ); - logger.debug( - "📁 [FOLDER-DELETE-WITH-CONTENTS] Deleting folder and contents:", - { folderId, folderName: folder?.name }, - ); - // Find all descendant items recursively (handles nested folders) const allDescendantIds = getAllDescendantIds(folderId, latestItems); @@ -952,11 +864,6 @@ export function useWorkspaceOperations( const idsToDelete = new Set([...allDescendantIds, folderId]); const itemCount = allDescendantIds.length; - logger.debug("📁 [FOLDER-DELETE-WITH-CONTENTS] Found items to delete:", { - itemCount, - itemIds: [...idsToDelete], - }); - // Delete PDF files from storage (fire-and-forget, non-blocking) // This is best-effort cleanup - files may become orphaned if this fails const itemsToDelete = latestItems.filter((item) => @@ -971,12 +878,7 @@ export function useWorkspaceOperations( { method: "DELETE", }, - ).catch((err) => - logger.warn( - "📁 [FOLDER-DELETE-WITH-CONTENTS] Failed to delete PDF file:", - err, - ), - ); + ).catch(() => {}); } } } @@ -999,10 +901,6 @@ export function useWorkspaceOperations( const moveItemToFolder = useCallback( (itemId: string, folderId: string | null) => { const item = currentItems.find((i) => i.id === itemId); - logger.debug("📁 [ITEM-MOVE] Moving item to folder:", { - itemId, - folderId, - }); const event = createEvent( "ITEM_MOVED_TO_FOLDER", { itemId, folderId, itemName: item?.name }, @@ -1019,10 +917,6 @@ export function useWorkspaceOperations( const itemNames = itemIds .map((id) => currentItems.find((i) => i.id === id)?.name) .filter((n): n is string => n != null); - logger.debug("📁 [ITEMS-MOVE] Moving items to folder:", { - itemIds, - folderId, - }); const event = createEvent( "ITEMS_MOVED_TO_FOLDER", { itemIds, folderId, itemNames }, @@ -1035,73 +929,10 @@ export function useWorkspaceOperations( ); // Flush pending debounced changes for an item (called when modal closes) - const flushPendingChanges = useCallback( - (itemId: string) => { - logger.debug("💾 [FLUSH] Flushing pending changes for item:", { itemId }); - - // Flush updateItem pending changes - const pendingItemTimeout = updateItemDebounceRef.current.get(itemId); - if (pendingItemTimeout) { - clearTimeout(pendingItemTimeout); - updateItemDebounceRef.current.delete(itemId); - - const pendingChanges = pendingItemChangesRef.current.get(itemId); - if (pendingChanges) { - const item = getLatestItemFromState(itemId); - const name = (pendingChanges as Partial).name ?? item?.name; - logger.debug("💾 [FLUSH] Sending pending updateItem changes:", { - itemId, - changes: pendingChanges, - }); - const event = createEvent( - "ITEM_UPDATED", - { id: itemId, changes: pendingChanges, name }, - userId, - userName, - ); - mutation.mutate(event); - pendingItemChangesRef.current.delete(itemId); - } - } - - // Flush updateItemData pending changes - const pendingDataTimeout = updateItemDataDebounceRef.current.get(itemId); - if (pendingDataTimeout) { - clearTimeout(pendingDataTimeout); - updateItemDataDebounceRef.current.delete(itemId); - - const pendingUpdater = pendingItemDataUpdatersRef.current.get(itemId); - if (pendingUpdater) { - // Get the latest item state - const latestItem = getLatestItemFromState(itemId); - if (latestItem) { - logger.debug("💾 [FLUSH] Sending pending updateItemData changes:", { - itemId, - }); - const newData = pendingUpdater(latestItem.data); - const event = createEvent( - "ITEM_UPDATED", - { - id: itemId, - changes: { data: newData }, - name: latestItem.name, - }, - userId, - userName, - ); - mutation.mutate(event); - pendingItemDataUpdatersRef.current.delete(itemId); - } else { - logger.warn( - `💾 [FLUSH] Item ${itemId} not found when flushing pending data changes`, - ); - pendingItemDataUpdatersRef.current.delete(itemId); - } - } - } - }, - [getLatestItemFromState, mutation, userId, userName], - ); + const flushPendingChanges = useCallback((itemId: string) => { + updateItemDebouncerRef.current.get(itemId)?.flush(); + updateItemDataDebouncerRef.current.get(itemId)?.flush(); + }, []); const getDocumentMarkdownForExport = useCallback( (itemId: string) => { diff --git a/src/hooks/workspace/use-workspace-realtime.ts b/src/hooks/workspace/use-workspace-realtime.ts index 5d06eb00..baf73233 100644 --- a/src/hooks/workspace/use-workspace-realtime.ts +++ b/src/hooks/workspace/use-workspace-realtime.ts @@ -87,7 +87,6 @@ export function useWorkspaceRealtime( const cleanup = useCallback(() => { if (!channelRef.current) return; - logger.debug("[REALTIME] Cleaning up workspace events channel"); removeChannel(channelRef.current); channelRef.current = null; setIsConnected(false); @@ -106,7 +105,6 @@ export function useWorkspaceRealtime( const supabase = getSupabaseClient(); const channelName = `workspace:${workspaceId}:events`; onStatusChange?.("connecting"); - logger.debug("[REALTIME] Subscribing to workspace channel", channelName); const channel = supabase.channel(channelName, { config: { @@ -162,7 +160,6 @@ export function useWorkspaceRealtime( }); channel.subscribe((status) => { - logger.debug("[REALTIME] Workspace channel status", status); switch (status) { case "SUBSCRIBED": setIsConnected(true); diff --git a/src/hooks/workspace/use-workspace-state.ts b/src/hooks/workspace/use-workspace-state.ts index 83a8d3d9..c638f2a3 100644 --- a/src/hooks/workspace/use-workspace-state.ts +++ b/src/hooks/workspace/use-workspace-state.ts @@ -16,7 +16,7 @@ async function fetchWorkspaceState( throw new Error(`Failed to fetch workspace state: ${response.statusText}`); } - return response.json(); + return (await response.json()) as WorkspaceStateResponse; } export function useWorkspaceState(workspaceId: string | null) { diff --git a/src/hooks/workspace/workspace-state-cache.ts b/src/hooks/workspace/workspace-state-cache.ts index 62d42e07..604f52fb 100644 --- a/src/hooks/workspace/workspace-state-cache.ts +++ b/src/hooks/workspace/workspace-state-cache.ts @@ -25,8 +25,10 @@ export function applyConfirmedWorkspaceEventToState( return stateData; } + const nextState = eventReducer(stateData.state, event); + return { - state: eventReducer(stateData.state, event), + state: nextState, version: event.version, }; } diff --git a/src/lib/ai/workers/workspace-worker.ts b/src/lib/ai/workers/workspace-worker.ts index 9195df3b..bc2dd25f 100644 --- a/src/lib/ai/workers/workspace-worker.ts +++ b/src/lib/ai/workers/workspace-worker.ts @@ -1,17 +1,11 @@ -import { headers } from "next/headers"; import { createPatch, diffLines } from "diff"; -import { auth } from "@/lib/auth"; -import { db, workspaces } from "@/lib/db/client"; import { appendWorkspaceEventOrThrow, appendWorkspaceEventUsingCurrentVersionWithRetry, } from "@/lib/workspace/workspace-event-store"; -import { workspaceCollaborators } from "@/lib/db/schema"; -import { eq, and } from "drizzle-orm"; import { createEvent } from "@/lib/workspace/events"; import { generateItemId } from "@/lib/workspace-state/item-helpers"; import { getRandomCardColor } from "@/lib/workspace-state/colors"; -import { normalizeWorkspaceItems } from "@/lib/workspace-state/state"; import { logger } from "@/lib/utils/logger"; import type { Item, @@ -23,10 +17,13 @@ import type { FlashcardItem, DocumentData, } from "@/lib/workspace-state/types"; -import { getOcrPagesTextContent } from "@/lib/utils/ocr-pages"; import { executeWorkspaceOperation } from "./common"; -import { loadWorkspaceState } from "@/lib/workspace/state-loader"; -import { hasDuplicateName } from "@/lib/workspace/unique-name"; +import { + requireMutationIdentity, + requireWorkspaceEditor, + loadWorkspaceItemsForValidation, + checkDuplicateName, +} from "@/lib/workspace/mutation-helpers"; import type { WorkspaceEvent } from "@/lib/workspace/events"; import { replace as applyReplace, @@ -305,68 +302,16 @@ export async function workspaceWorker( params.workspaceId, async () => { try { - logger.debug("📝 [WORKSPACE-WORKER] Action:", action, params); - - // Get current user - const session = await auth.api.getSession({ - headers: await headers(), - }); - if (!session) { - throw new Error("User not authenticated"); - } - const userId = session.user.id; - const userName = session.user.name || session.user.email || undefined; - - // Verify workspace access - owner OR editor collaborator - const workspace = await db - .select({ userId: workspaces.userId }) - .from(workspaces) - .where(eq(workspaces.id, params.workspaceId)) - .limit(1); - - if (!workspace[0]) { - throw new Error("Workspace not found"); - } - - // Check if user is owner - const isOwner = workspace[0].userId === userId; - - // If not owner, check if user is an editor collaborator - if (!isOwner) { - const [collaborator] = await db - .select({ permissionLevel: workspaceCollaborators.permissionLevel }) - .from(workspaceCollaborators) - .where( - and( - eq(workspaceCollaborators.workspaceId, params.workspaceId), - eq(workspaceCollaborators.userId, userId), - ), - ) - .limit(1); - - if (!collaborator || collaborator.permissionLevel !== "editor") { - throw new Error("Access denied - editor permission required"); - } - } + const { userId, userName } = await requireMutationIdentity(); + await requireWorkspaceEditor(params.workspaceId, userId); // Handle different actions if (action === "create") { const item = await buildItemFromCreateParams(params); - const currentState = normalizeWorkspaceItems( - await loadWorkspaceState(params.workspaceId, { userId }), - ); - if ( - hasDuplicateName( - currentState, - item.name, - item.type, - item.folderId ?? null, - ) - ) { - return { - success: false, - message: `A ${item.type} named "${item.name}" already exists in this folder`, - }; + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); + const dupError = checkDuplicateName(currentState, item.name, item.type, item.folderId ?? null); + if (dupError) { + return { success: false, message: dupError }; } const event = createEvent( "ITEM_CREATED", @@ -381,8 +326,6 @@ export async function workspaceWorker( 2, ); - logger.info(`📝 [WORKSPACE-WORKER] Created ${item.type}:`, item.name); - // Include card count for flashcard decks (use created item.data.cards, not params) const flashcardCards = item.type === "flashcard" && item.data && "cards" in item.data @@ -412,6 +355,17 @@ export async function workspaceWorker( const items = await Promise.all( params.items.map((p) => buildItemFromCreateParams(p)), ); + + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); + for (let i = 0; i < items.length; i++) { + const item = items[i]; + const preceding = [...currentState, ...items.slice(0, i)]; + const dupError = checkDuplicateName(preceding, item.name, item.type, item.folderId ?? null); + if (dupError) { + return { success: false, message: dupError }; + } + } + const event = createEvent( "BULK_ITEMS_CREATED", { items }, @@ -424,9 +378,6 @@ export async function workspaceWorker( event, ); - logger.info( - `📝 [WORKSPACE-WORKER] Bulk created ${items.length} items`, - ); return { success: true, message: `Bulk created ${items.length} items successfully`, @@ -446,10 +397,7 @@ export async function workspaceWorker( throw new Error("Cards to add required for flashcard update"); } - // Use helper to load current state (duplicated logic removed) - const currentState = normalizeWorkspaceItems( - await loadWorkspaceState(params.workspaceId, { userId }), - ); + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); const existingItem = currentState.find( (i: any) => i.id === params.itemId, @@ -486,20 +434,9 @@ export async function workspaceWorker( // Handle title update if provided if (params.title) { - logger.debug("🎴 [UPDATE-FLASHCARD] Updating title:", params.title); - if ( - hasDuplicateName( - currentState, - params.title, - existingItem.type, - existingItem.folderId ?? null, - params.itemId, - ) - ) { - return { - success: false, - message: `A ${existingItem.type} named "${params.title}" already exists in this folder`, - }; + const dupError = checkDuplicateName(currentState, params.title, existingItem.type, existingItem.folderId ?? null, params.itemId); + if (dupError) { + return { success: false, message: dupError }; } changes.name = params.title; } @@ -520,13 +457,6 @@ export async function workspaceWorker( event, ); - logger.info("🎴 [WORKSPACE-WORKER] Updated flashcard deck:", { - itemId: params.itemId, - cardsAdded: newCards.length, - totalCards: updatedData.cards.length, - newTitle: params.title, - }); - return { success: true, itemId: params.itemId, @@ -553,15 +483,7 @@ export async function workspaceWorker( ); } - logger.debug("🎯 [WORKSPACE-WORKER] Updating quiz:", { - itemId: params.itemId, - questionsToAdd: questionsToAdd?.length ?? 0, - titleUpdate: hasTitle, - }); - - const currentState = normalizeWorkspaceItems( - await loadWorkspaceState(params.workspaceId, { userId }), - ); + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); const existingItem = currentState.find( (i: any) => i.id === params.itemId, @@ -589,20 +511,9 @@ export async function workspaceWorker( const changes: any = hasQuestions ? { data: updatedData } : {}; if (params.title) { - logger.debug("🎯 [UPDATE-QUIZ] Updating title:", params.title); - if ( - hasDuplicateName( - currentState, - params.title, - existingItem.type, - existingItem.folderId ?? null, - params.itemId, - ) - ) { - return { - success: false, - message: `A ${existingItem.type} named "${params.title}" already exists in this folder`, - }; + const dupError = checkDuplicateName(currentState, params.title, existingItem.type, existingItem.folderId ?? null, params.itemId); + if (dupError) { + return { success: false, message: dupError }; } changes.name = params.title; } @@ -618,27 +529,10 @@ export async function workspaceWorker( userName, ); - logger.debug("📝 [UPDATE-QUIZ-DB] Created event:", { - eventId: event.id, - eventType: event.type, - payloadId: event.payload.id, - questionsInPayload: (event.payload.changes?.data as any)?.questions - ?.length, - }); - const appendResult = await persistWorkspaceEvent( params.workspaceId, event, ); - logger.info("📝 [UPDATE-QUIZ-DB] Persisted result:", { - version: appendResult.version, - }); - - logger.info("🎯 [WORKSPACE-WORKER] Updated quiz:", { - itemId: params.itemId, - questionsAdded: questionsToAdd?.length ?? 0, - totalQuestions: updatedData.questions.length, - }); return { success: true, @@ -669,9 +563,7 @@ export async function workspaceWorker( ); } - const currentState = normalizeWorkspaceItems( - await loadWorkspaceState(params.workspaceId, { userId }), - ); + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); const existingItem = currentState.find( (i: any) => i.id === params.itemId, ); @@ -697,19 +589,9 @@ export async function workspaceWorker( const changes: Partial = { data: updatedData }; if (params.title) { - if ( - hasDuplicateName( - currentState, - params.title, - existingItem.type, - existingItem.folderId ?? null, - params.itemId, - ) - ) { - return { - success: false, - message: `A ${existingItem.type} named "${params.title}" already exists in this folder`, - }; + const dupError = checkDuplicateName(currentState, params.title, existingItem.type, existingItem.folderId ?? null, params.itemId); + if (dupError) { + return { success: false, message: dupError }; } changes.name = params.title; } @@ -730,16 +612,11 @@ export async function workspaceWorker( event, ); - const contentLen = getOcrPagesTextContent(params.pdfOcrPages).length; - logger.info("📄 [WORKSPACE-WORKER] Updated PDF OCR content:", { - itemId: params.itemId, - contentLength: contentLen, - }); return { success: true, itemId: params.itemId, - message: `Cached OCR content for PDF "${existingItem.name}" (${contentLen} chars)`, + message: `Cached OCR content for PDF "${existingItem.name}"`, event, version: appendResult.version, }; @@ -759,9 +636,7 @@ export async function workspaceWorker( ); } - const currentState = normalizeWorkspaceItems( - await loadWorkspaceState(params.workspaceId, { userId }), - ); + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); const existingItem = currentState.find( (i: any) => i.id === params.itemId, ); @@ -832,9 +707,7 @@ export async function workspaceWorker( throw new Error("oldString and newString required for edit"); } - const currentState = normalizeWorkspaceItems( - await loadWorkspaceState(params.workspaceId, { userId }), - ); + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); const existingItem = currentState.find( (i: Item) => i.id === params.itemId, ); @@ -898,19 +771,9 @@ export async function workspaceWorker( data: { ...data, cards: newCards } as FlashcardData, }; if (rename) { - if ( - hasDuplicateName( - currentState, - rename, - "flashcard", - existingItem.folderId ?? null, - params.itemId, - ) - ) { - return { - success: false, - message: `A flashcard deck named "${rename}" already exists in this folder`, - }; + const dupError = checkDuplicateName(currentState, rename, "flashcard", existingItem.folderId ?? null, params.itemId); + if (dupError) { + return { success: false, message: dupError }; } changes.name = rename; } @@ -1013,19 +876,9 @@ export async function workspaceWorker( const changes: Partial = { data: updatedData }; if (rename) { - if ( - hasDuplicateName( - currentState, - rename, - "quiz", - existingItem.folderId ?? null, - params.itemId, - ) - ) { - return { - success: false, - message: `A quiz named "${rename}" already exists in this folder`, - }; + const dupError = checkDuplicateName(currentState, rename, "quiz", existingItem.folderId ?? null, params.itemId); + if (dupError) { + return { success: false, message: dupError }; } changes.name = rename; } @@ -1059,19 +912,9 @@ export async function workspaceWorker( const changes: Partial = {}; const docData = existingItem.data as DocumentData; if (rename) { - if ( - hasDuplicateName( - currentState, - rename, - "document", - existingItem.folderId ?? null, - params.itemId, - ) - ) { - return { - success: false, - message: `A document named "${rename}" already exists in this folder`, - }; + const dupError = checkDuplicateName(currentState, rename, "document", existingItem.folderId ?? null, params.itemId); + if (dupError) { + return { success: false, message: dupError }; } changes.name = rename; } @@ -1169,19 +1012,9 @@ export async function workspaceWorker( "PDFs can only be renamed. Use oldString='', newString='', and newName='new name'.", }; } - if ( - hasDuplicateName( - currentState, - rename, - "pdf", - existingItem.folderId ?? null, - params.itemId, - ) - ) { - return { - success: false, - message: `A PDF named "${rename}" already exists in this folder`, - }; + const dupError = checkDuplicateName(currentState, rename, "pdf", existingItem.folderId ?? null, params.itemId); + if (dupError) { + return { success: false, message: dupError }; } const changes: Partial = { name: rename }; const event = createEvent( @@ -1211,9 +1044,7 @@ export async function workspaceWorker( throw new Error("Item ID required for delete"); } - const currentState = normalizeWorkspaceItems( - await loadWorkspaceState(params.workspaceId, { userId }), - ); + const currentState = await loadWorkspaceItemsForValidation(params.workspaceId, userId); const existingItem = currentState.find( (i: any) => i.id === params.itemId, ); @@ -1229,8 +1060,6 @@ export async function workspaceWorker( event, ); - logger.info("📝 [WORKSPACE-WORKER] Deleted item:", params.itemId); - return { success: true, itemId: params.itemId, @@ -1248,7 +1077,7 @@ export async function workspaceWorker( } catch (error) { const errorMessage = error instanceof Error ? error.message : "Unknown error"; - logger.error("📝 [WORKSPACE-WORKER] Error:", errorMessage); + logger.error("[WORKSPACE-WORKER] Error:", errorMessage); return { success: false, message: `Failed: ${errorMessage}`, diff --git a/src/lib/audio/poll-audio-processing.ts b/src/lib/audio/poll-audio-processing.ts index bbfe5437..6fbf8f12 100644 --- a/src/lib/audio/poll-audio-processing.ts +++ b/src/lib/audio/poll-audio-processing.ts @@ -1,71 +1,44 @@ -const POLL_INTERVAL_MS = 10_000; +import { pollTask } from "@/lib/tasks/poll-task"; + +const AUDIO_COMPLETE_EVENT = "audio-processing-complete"; +const POLL_INTERVAL_MS = 5_000; +const MAX_POLL_ATTEMPTS = 120; /** * Polls the audio processing status endpoint until the workflow completes or fails. * Dispatches audio-processing-complete when done. */ -export async function pollAudioProcessing(runId: string, itemId: string): Promise { - const dispatchError = (error: string) => { +export async function pollAudioProcessing( + runId: string, + itemId: string, + signal?: AbortSignal, +): Promise { + const dispatchComplete = (detail: Record) => { + if (typeof window === "undefined") return; window.dispatchEvent( - new CustomEvent("audio-processing-complete", { - detail: { itemId, error }, - }) + new CustomEvent(AUDIO_COMPLETE_EVENT, { detail: { itemId, ...detail } }), ); }; - while (true) { - try { - const res = await fetch(`/api/audio/process/status?runId=${encodeURIComponent(runId)}`); - if (!res.ok) { - dispatchError(`Status check failed: ${res.status}`); - return; - } - const data = await res.json(); - - if (data.status === "completed") { - window.dispatchEvent( - new CustomEvent("audio-processing-complete", { - detail: { - itemId, - summary: data.result.summary, - segments: data.result.segments, - duration: data.result.duration, - }, - }) - ); - return; - } - - if (data.status === "failed") { - window.dispatchEvent( - new CustomEvent("audio-processing-complete", { - detail: { - itemId, - error: data.error || "Processing failed", - }, - }) - ); - return; - } + const result = await pollTask({ + statusUrl: `/api/audio/process/status?runId=${encodeURIComponent(runId)}`, + intervalMs: POLL_INTERVAL_MS, + maxAttempts: MAX_POLL_ATTEMPTS, + signal, + }); - if (data.status === "cancelled") { - window.dispatchEvent( - new CustomEvent("audio-processing-complete", { - detail: { - itemId, - error: data.error || "Processing cancelled", - }, - }) - ); - return; - } - - await new Promise((r) => setTimeout(r, POLL_INTERVAL_MS)); - } catch (err) { - dispatchError( - err instanceof Error ? err.message : "Network or parse error during polling" - ); - return; - } + if (result.status === "completed" && result.data) { + const resultData = result.data.result as + | { summary?: string; segments?: unknown[]; duration?: number } + | undefined; + dispatchComplete({ + summary: resultData?.summary, + segments: resultData?.segments, + duration: resultData?.duration, + }); + } else { + dispatchComplete({ error: result.error ?? "Processing failed" }); } } + +export { AUDIO_COMPLETE_EVENT }; diff --git a/src/lib/ocr/client.ts b/src/lib/ocr/client.ts index d23313c9..a66c7081 100644 --- a/src/lib/ocr/client.ts +++ b/src/lib/ocr/client.ts @@ -1,9 +1,9 @@ +import { pollTask } from "@/lib/tasks/poll-task"; import type { OcrCandidate } from "./types"; const OCR_COMPLETE_EVENT = "ocr-processing-complete"; -const POLL_INTERVAL_MS = 1000; -const MAX_POLL_MS = 10 * 60 * 1000; -const MAX_POLL_ATTEMPTS = Math.ceil(MAX_POLL_MS / POLL_INTERVAL_MS); +const POLL_INTERVAL_MS = 1_000; +const MAX_POLL_ATTEMPTS = 600; interface OcrProcessingEventDetail { itemIds: string[]; @@ -19,112 +19,42 @@ function emitOcrProcessingComplete(detail: OcrProcessingEventDetail) { export async function pollOcrRun( runId: string, itemIds: string[], - signal?: AbortSignal + signal?: AbortSignal, ): Promise { - const emitFailure = (error: string) => { + const result = await pollTask({ + statusUrl: `/api/ocr/status?runId=${encodeURIComponent(runId)}`, + intervalMs: POLL_INTERVAL_MS, + maxAttempts: MAX_POLL_ATTEMPTS, + signal, + }); + + if (result.status === "completed") { + emitOcrProcessingComplete({ itemIds, status: "completed" }); + } else { emitOcrProcessingComplete({ itemIds, status: "failed", - error, - }); - }; - - for (let attempt = 0; attempt < MAX_POLL_ATTEMPTS; attempt += 1) { - if (signal?.aborted) { - emitFailure("OCR polling canceled"); - return; - } - - let res: Response; - let data: { status?: string; error?: string }; - - try { - res = await fetch(`/api/ocr/status?runId=${encodeURIComponent(runId)}`, { - signal, - }); - data = (await res.json().catch(() => { - throw new Error("Failed to parse OCR status response"); - })) as { status?: string; error?: string }; - } catch (error) { - const aborted = - signal?.aborted || - (error instanceof DOMException && error.name === "AbortError"); - emitFailure( - aborted - ? "OCR polling canceled" - : error instanceof Error - ? error.message - : "Failed to poll OCR status" - ); - return; - } - - if (data.status === "not_found" || res.status === 404) { - emitFailure(data.error ?? "OCR run not found or expired"); - return; - } - - if (!res.ok) { - emitFailure(data.error || `Failed to fetch OCR status: ${res.status}`); - return; - } - - if (data.status === "completed") { - emitOcrProcessingComplete({ - itemIds, - status: "completed", - }); - return; - } - - if (data.status === "failed") { - emitFailure(data.error || "OCR failed"); - return; - } - - if (data.status === "cancelled") { - emitFailure(data.error || "OCR cancelled"); - return; - } - - await new Promise((resolve) => { - const timeoutId = window.setTimeout(() => { - signal?.removeEventListener("abort", handleAbort); - resolve(); - }, POLL_INTERVAL_MS); - - const handleAbort = () => { - window.clearTimeout(timeoutId); - signal?.removeEventListener("abort", handleAbort); - resolve(); - }; - - signal?.addEventListener("abort", handleAbort, { once: true }); + error: result.error ?? "OCR failed", }); } - - emitFailure(`OCR status polling timed out after ${MAX_POLL_MS / 1000} seconds`); } export async function startOcrProcessing( workspaceId: string, - candidates: OcrCandidate[] + candidates: OcrCandidate[], ): Promise { if (!workspaceId || candidates.length === 0) return; const res = await fetch("/api/ocr/start", { method: "POST", headers: { "Content-Type": "application/json" }, - body: JSON.stringify({ - workspaceId, - candidates, - }), + body: JSON.stringify({ workspaceId, candidates }), }); const data = await res.json().catch(() => ({})); if (!res.ok) { emitOcrProcessingComplete({ - itemIds: candidates.map((candidate) => candidate.itemId), + itemIds: candidates.map((c) => c.itemId), status: "failed", error: data.error || `Failed to start OCR: ${res.status}`, }); @@ -135,7 +65,7 @@ export async function startOcrProcessing( data.runId, Array.isArray(data.itemIds) ? (data.itemIds as string[]) - : candidates.map((candidate) => candidate.itemId) + : candidates.map((c) => c.itemId), ); } diff --git a/src/lib/tasks/__tests__/poll-task.test.ts b/src/lib/tasks/__tests__/poll-task.test.ts new file mode 100644 index 00000000..20df036d --- /dev/null +++ b/src/lib/tasks/__tests__/poll-task.test.ts @@ -0,0 +1,143 @@ +import { describe, expect, it, vi, beforeEach, afterEach } from "vitest"; +import { pollTask } from "../poll-task"; + +describe("pollTask", () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); + }); + + it("returns completed immediately when first poll is terminal", async () => { + vi.stubGlobal( + "fetch", + vi.fn().mockResolvedValue({ + ok: true, + json: async () => ({ status: "completed" }), + }), + ); + + const result = await pollTask({ + statusUrl: "/api/test/status", + intervalMs: 100, + }); + + expect(result.status).toBe("completed"); + expect(fetch).toHaveBeenCalledTimes(1); + }); + + it("returns failed when status endpoint returns error", async () => { + vi.stubGlobal( + "fetch", + vi.fn().mockResolvedValue({ + ok: false, + status: 500, + json: async () => ({ error: "internal" }), + }), + ); + + const result = await pollTask({ + statusUrl: "/api/test/status", + intervalMs: 100, + }); + + expect(result.status).toBe("failed"); + expect(result.error).toContain("500"); + }); + + it("times out after maxAttempts", async () => { + vi.stubGlobal( + "fetch", + vi.fn().mockResolvedValue({ + ok: true, + json: async () => ({ status: "running" }), + }), + ); + + const promise = pollTask({ + statusUrl: "/api/test/status", + intervalMs: 10, + maxAttempts: 3, + }); + + for (let i = 0; i < 5; i++) { + await vi.advanceTimersByTimeAsync(50); + } + + const result = await promise; + expect(result.status).toBe("failed"); + expect(result.error).toContain("timed out"); + expect(fetch).toHaveBeenCalledTimes(3); + }); + + it("returns cancelled when signal is aborted", async () => { + const controller = new AbortController(); + vi.stubGlobal( + "fetch", + vi.fn().mockRejectedValue( + Object.assign(new DOMException("Aborted", "AbortError")), + ), + ); + + controller.abort(); + + const result = await pollTask({ + statusUrl: "/api/test/status", + signal: controller.signal, + }); + + expect(result.status).toBe("cancelled"); + }); + + it("propagates error string from terminal status", async () => { + vi.stubGlobal( + "fetch", + vi.fn().mockResolvedValue({ + ok: true, + json: async () => ({ status: "failed", error: "OCR engine down" }), + }), + ); + + const result = await pollTask({ + statusUrl: "/api/test/status", + intervalMs: 100, + }); + + expect(result.status).toBe("failed"); + expect(result.error).toBe("OCR engine down"); + }); + + it("calls onStatus callback on each poll", async () => { + let callCount = 0; + vi.stubGlobal( + "fetch", + vi.fn().mockImplementation(async () => { + callCount++; + return { + ok: true, + json: async () => + callCount >= 2 + ? { status: "completed" } + : { status: "running" }, + }; + }), + ); + + const statuses: string[] = []; + + const promise = pollTask({ + statusUrl: "/api/test/status", + intervalMs: 10, + onStatus: (s) => statuses.push(s), + }); + + await vi.advanceTimersByTimeAsync(50); + await promise; + + expect(statuses).toContain("running"); + expect(statuses).toContain("completed"); + }); +}); diff --git a/src/lib/tasks/poll-task.ts b/src/lib/tasks/poll-task.ts new file mode 100644 index 00000000..2f429a5b --- /dev/null +++ b/src/lib/tasks/poll-task.ts @@ -0,0 +1,98 @@ +import type { TaskStatus } from "./task-types"; + +const TERMINAL_STATUSES: Set = new Set([ + "completed", + "failed", + "cancelled", + "not_found", +]); + +export interface PollTaskOptions { + statusUrl: string; + intervalMs?: number; + maxAttempts?: number; + signal?: AbortSignal; + onStatus?: (status: TaskStatus, data: Record) => void; +} + +export interface PollTaskResult { + status: TaskStatus; + error?: string; + data?: Record; +} + +/** + * Generic polling loop for long-running task status endpoints. + * Returns when the task reaches a terminal status or polling is aborted/exhausted. + */ +export async function pollTask(opts: PollTaskOptions): Promise { + const { + statusUrl, + intervalMs = 2_000, + maxAttempts = 300, + signal, + onStatus, + } = opts; + + for (let attempt = 0; attempt < maxAttempts; attempt++) { + if (signal?.aborted) { + return { status: "cancelled", error: "Polling aborted" }; + } + + let data: Record; + + try { + const res = await fetch(statusUrl, { signal }); + if (!res.ok) { + return { + status: "failed", + error: `Status check failed: ${res.status}`, + }; + } + data = (await res.json()) as Record; + } catch (err) { + const aborted = + signal?.aborted || + (err instanceof DOMException && err.name === "AbortError"); + return { + status: aborted ? "cancelled" : "failed", + error: aborted + ? "Polling aborted" + : err instanceof Error + ? err.message + : "Network error during polling", + }; + } + + const status = data.status as TaskStatus; + onStatus?.(status, data); + + if (TERMINAL_STATUSES.has(status)) { + return { + status, + error: typeof data.error === "string" ? data.error : undefined, + data, + }; + } + + await new Promise((resolve) => { + const timeoutId = setTimeout(() => { + signal?.removeEventListener("abort", handleAbort); + resolve(); + }, intervalMs); + + const handleAbort = () => { + clearTimeout(timeoutId); + signal?.removeEventListener("abort", handleAbort); + resolve(); + }; + + signal?.addEventListener("abort", handleAbort, { once: true }); + }); + } + + return { + status: "failed", + error: `Polling timed out after ${maxAttempts} attempts`, + }; +} diff --git a/src/lib/tasks/task-types.ts b/src/lib/tasks/task-types.ts new file mode 100644 index 00000000..7422318e --- /dev/null +++ b/src/lib/tasks/task-types.ts @@ -0,0 +1,22 @@ +import { z } from "zod"; + +export const taskStatusSchema = z.enum([ + "running", + "completed", + "failed", + "cancelled", + "not_found", +]); + +export type TaskStatus = z.infer; + +export interface TaskStartResult { + runId: string; + itemIds: string[]; +} + +export interface TaskStatusResult { + status: TaskStatus; + error?: string; + result?: unknown; +} diff --git a/src/lib/tasks/validate-asset-url.ts b/src/lib/tasks/validate-asset-url.ts new file mode 100644 index 00000000..ab54dca0 --- /dev/null +++ b/src/lib/tasks/validate-asset-url.ts @@ -0,0 +1,10 @@ +import { isAllowedOcrFileUrl } from "@/lib/ocr/url-validation"; + +/** + * Validate that a file URL points to an allowed asset host. + * Delegates to the shared OCR allowlist which covers Supabase, app URL, + * and localhost in dev -- the same hosts that are valid for any asset. + */ +export function isAllowedAssetUrl(url: string): boolean { + return isAllowedOcrFileUrl(url); +} diff --git a/src/lib/workspace/__tests__/event-reducer.lightweight-ocr.test.ts b/src/lib/workspace/__tests__/event-reducer.lightweight-ocr.test.ts new file mode 100644 index 00000000..a0522b07 --- /dev/null +++ b/src/lib/workspace/__tests__/event-reducer.lightweight-ocr.test.ts @@ -0,0 +1,105 @@ +import { describe, expect, it } from "vitest"; +import { eventReducer } from "../event-reducer"; +import type { Item, PdfData, ImageData } from "@/lib/workspace-state/types"; +import type { WorkspaceEvent } from "../events"; + +describe("eventReducer lightweight OCR events", () => { + const basePdf: Item = { + id: "pdf-1", + type: "pdf", + name: "Doc", + subtitle: "", + data: { + fileUrl: "https://example.com/a.pdf", + filename: "a.pdf", + ocrStatus: "processing", + } as PdfData, + }; + + const baseImage: Item = { + id: "img-1", + type: "image", + name: "Image", + subtitle: "", + data: { + url: "https://example.com/b.png", + ocrStatus: "processing", + } as ImageData, + }; + + it("applies status-only BULK_ITEMS_PATCHED (no ocrPages in event)", () => { + const event: WorkspaceEvent = { + id: "evt-1", + type: "BULK_ITEMS_PATCHED", + timestamp: Date.now(), + userId: "user-1", + payload: { + updates: [ + { + id: "pdf-1", + changes: { + data: { ocrStatus: "complete" } as Item["data"], + }, + }, + { + id: "img-1", + changes: { + data: { + ocrStatus: "failed", + ocrError: "bad image", + } as Item["data"], + }, + }, + ], + }, + }; + + const next = eventReducer([basePdf, baseImage], event); + + const pdfData = next[0].data as PdfData; + expect(pdfData.ocrStatus).toBe("complete"); + expect(pdfData.fileUrl).toBe("https://example.com/a.pdf"); + expect(pdfData.filename).toBe("a.pdf"); + + const imgData = next[1].data as ImageData; + expect(imgData.ocrStatus).toBe("failed"); + expect(imgData.ocrError).toBe("bad image"); + expect(imgData.url).toBe("https://example.com/b.png"); + }); + + it("applies status-only ITEM_UPDATED for audio completion", () => { + const audioItem: Item = { + id: "audio-1", + type: "audio", + name: "Recording", + subtitle: "", + data: { + fileUrl: "https://example.com/a.mp3", + filename: "a.mp3", + processingStatus: "processing", + }, + }; + + const event: WorkspaceEvent = { + id: "evt-2", + type: "ITEM_UPDATED", + timestamp: Date.now(), + userId: "user-1", + payload: { + id: "audio-1", + changes: { + data: { + summary: "Meeting notes", + processingStatus: "complete", + } as Item["data"], + }, + }, + }; + + const next = eventReducer([audioItem], event); + const data = next[0].data as Record; + expect(data.processingStatus).toBe("complete"); + expect(data.summary).toBe("Meeting notes"); + expect(data.fileUrl).toBe("https://example.com/a.mp3"); + }); +}); diff --git a/src/lib/workspace/__tests__/version-helpers.test.ts b/src/lib/workspace/__tests__/version-helpers.test.ts new file mode 100644 index 00000000..05e5cc07 --- /dev/null +++ b/src/lib/workspace/__tests__/version-helpers.test.ts @@ -0,0 +1,99 @@ +import { describe, expect, it } from "vitest"; +import { + computeBaseVersion, + removeOptimisticEvent, + confirmOptimisticEvent, +} from "../version-helpers"; +import type { EventResponse } from "../events"; + +function makeCache( + version: number, + events: Array<{ id: string; version?: number }>, +): EventResponse { + return { + version, + events: events.map((e) => ({ + id: e.id, + type: "ITEM_DELETED" as const, + timestamp: Date.now(), + userId: "u", + payload: { id: "x" }, + ...(e.version !== undefined && { version: e.version }), + })), + }; +} + +describe("computeBaseVersion", () => { + it("returns 0 for undefined cache", () => { + const result = computeBaseVersion(undefined); + expect(result.baseVersion).toBe(0); + expect(result.currentVersion).toBe(0); + expect(result.optimisticCount).toBe(0); + }); + + it("returns cache version when no optimistic events", () => { + const cache = makeCache(10, [{ id: "a", version: 10 }]); + const result = computeBaseVersion(cache); + expect(result.baseVersion).toBe(10); + expect(result.currentVersion).toBe(10); + }); + + it("adjusts for optimistic events (subtracts own)", () => { + const cache = makeCache(10, [ + { id: "a", version: 10 }, + { id: "b" }, + { id: "c" }, + ]); + const result = computeBaseVersion(cache); + expect(result.optimisticCount).toBe(2); + expect(result.baseVersion).toBe(11); + }); + + it("uses max event version over cache version", () => { + const cache = makeCache(5, [ + { id: "a", version: 5 }, + { id: "b", version: 8 }, + ]); + const result = computeBaseVersion(cache); + expect(result.currentVersion).toBe(8); + expect(result.baseVersion).toBe(8); + }); +}); + +describe("removeOptimisticEvent", () => { + it("removes event by ID", () => { + const cache = makeCache(5, [ + { id: "a", version: 5 }, + { id: "b" }, + ]); + const result = removeOptimisticEvent(cache, "b"); + expect(result?.events).toHaveLength(1); + expect(result?.events[0].id).toBe("a"); + }); + + it("updates version when provided", () => { + const cache = makeCache(5, [{ id: "a", version: 5 }]); + const result = removeOptimisticEvent(cache, "x", 10); + expect(result?.version).toBe(10); + }); + + it("returns undefined for undefined input", () => { + expect(removeOptimisticEvent(undefined, "a")).toBeUndefined(); + }); +}); + +describe("confirmOptimisticEvent", () => { + it("stamps version on the matching event", () => { + const cache = makeCache(5, [ + { id: "a", version: 5 }, + { id: "b" }, + ]); + const result = confirmOptimisticEvent(cache, "b", 6); + expect(result?.events[1].version).toBe(6); + expect(result?.version).toBe(6); + }); + + it("returns undefined for undefined input", () => { + expect(confirmOptimisticEvent(undefined, "a", 1)).toBeUndefined(); + }); +}); diff --git a/src/lib/workspace/mutation-helpers.ts b/src/lib/workspace/mutation-helpers.ts new file mode 100644 index 00000000..4b958956 --- /dev/null +++ b/src/lib/workspace/mutation-helpers.ts @@ -0,0 +1,96 @@ +import { headers } from "next/headers"; +import { auth } from "@/lib/auth"; +import { db, workspaces } from "@/lib/db/client"; +import { workspaceCollaborators } from "@/lib/db/schema"; +import { eq, and } from "drizzle-orm"; +import { loadWorkspaceState } from "@/lib/workspace/state-loader"; +import { hasDuplicateName } from "@/lib/workspace/unique-name"; +import { normalizeWorkspaceItems } from "@/lib/workspace-state/state"; +import type { Item } from "@/lib/workspace-state/types"; + +export interface MutationIdentity { + userId: string; + userName?: string; +} + +/** + * Resolve authenticated user for server-side mutation contexts (API routes, + * workers, workflows). Returns userId and displayName from the session. + */ +export async function requireMutationIdentity(): Promise { + const session = await auth.api.getSession({ + headers: await headers(), + }); + if (!session) { + throw new Error("User not authenticated"); + } + return { + userId: session.user.id, + userName: session.user.name || session.user.email || undefined, + }; +} + +/** + * Verify the caller has editor-level access to a workspace (owner or editor + * collaborator). Throws on missing workspace or insufficient permissions. + */ +export async function requireWorkspaceEditor( + workspaceId: string, + userId: string, +): Promise { + const workspace = await db + .select({ userId: workspaces.userId }) + .from(workspaces) + .where(eq(workspaces.id, workspaceId)) + .limit(1); + + if (!workspace[0]) { + throw new Error("Workspace not found"); + } + + if (workspace[0].userId === userId) return; + + const [collaborator] = await db + .select({ permissionLevel: workspaceCollaborators.permissionLevel }) + .from(workspaceCollaborators) + .where( + and( + eq(workspaceCollaborators.workspaceId, workspaceId), + eq(workspaceCollaborators.userId, userId), + ), + ) + .limit(1); + + if (!collaborator || collaborator.permissionLevel !== "editor") { + throw new Error("Access denied - editor permission required"); + } +} + +/** + * Load workspace items for invariant checks (normalized). + */ +export async function loadWorkspaceItemsForValidation( + workspaceId: string, + userId: string, +): Promise { + return normalizeWorkspaceItems( + await loadWorkspaceState(workspaceId, { userId }), + ); +} + +/** + * Assert that no sibling item shares the same name + type in the given folder. + * Returns the error message if a duplicate exists, or null if the name is available. + */ +export function checkDuplicateName( + items: Item[], + name: string, + type: Item["type"], + folderId: string | null, + excludeItemId?: string, +): string | null { + if (hasDuplicateName(items, name, type, folderId, excludeItemId)) { + return `A ${type} named "${name}" already exists in this folder`; + } + return null; +} diff --git a/src/lib/workspace/state-loader.ts b/src/lib/workspace/state-loader.ts index 71fc9af5..b4966094 100644 --- a/src/lib/workspace/state-loader.ts +++ b/src/lib/workspace/state-loader.ts @@ -1,4 +1,5 @@ import { db } from "@/lib/db/client"; +import { logger } from "@/lib/utils/logger"; import type { Item } from "@/lib/workspace-state/types"; import { loadWorkspaceProjectionState } from "./workspace-items-projector"; diff --git a/src/lib/workspace/version-helpers.ts b/src/lib/workspace/version-helpers.ts new file mode 100644 index 00000000..7df54ba4 --- /dev/null +++ b/src/lib/workspace/version-helpers.ts @@ -0,0 +1,61 @@ +import type { EventResponse } from "./events"; + +/** + * Compute the effective base version for a new mutation, accounting for: + * 1. The cache's stored version field + * 2. The max version from individual events (may be higher after tool events) + * 3. Pending optimistic events that will increment the server version + */ +export function computeBaseVersion( + cacheData: EventResponse | undefined, +): { baseVersion: number; currentVersion: number; optimisticCount: number } { + const events = cacheData?.events ?? []; + + const optimisticCount = events.filter( + (e) => typeof e.version !== "number", + ).length; + + const maxEventVersion = events + .filter((e) => typeof e.version === "number") + .reduce((max, e) => Math.max(max, e.version!), cacheData?.version ?? 0); + + const currentVersion = Math.max(cacheData?.version ?? 0, maxEventVersion); + + const baseVersion = currentVersion + Math.max(0, optimisticCount - 1); + + return { baseVersion, currentVersion, optimisticCount }; +} + +/** + * Remove a single optimistic event from the cache by ID. + */ +export function removeOptimisticEvent( + cache: EventResponse | undefined, + eventId: string, + newVersion?: number, +): EventResponse | undefined { + if (!cache) return cache; + return { + ...cache, + events: cache.events.filter((e) => e.id !== eventId), + ...(newVersion !== undefined && { version: newVersion }), + }; +} + +/** + * Confirm an optimistic event by stamping its server-assigned version. + */ +export function confirmOptimisticEvent( + cache: EventResponse | undefined, + eventId: string, + serverVersion: number, +): EventResponse | undefined { + if (!cache) return cache; + return { + ...cache, + events: cache.events.map((e) => + e.id === eventId ? { ...e, version: serverVersion } : e, + ), + version: serverVersion, + }; +} diff --git a/src/lib/workspace/workspace-event-store.ts b/src/lib/workspace/workspace-event-store.ts index d094975a..5baf5d46 100644 --- a/src/lib/workspace/workspace-event-store.ts +++ b/src/lib/workspace/workspace-event-store.ts @@ -69,7 +69,6 @@ export async function getWorkspaceVersion( const result = await db.execute(sql` SELECT get_workspace_version(${workspaceId}::uuid) as version `); - return Number(result[0]?.version ?? 0); } @@ -183,6 +182,11 @@ export async function appendWorkspaceEventUsingCurrentVersionWithRetry(params: { return result; } + logger.warn("[EVENT-STORE] Conflict, retrying", { + workspaceId: params.workspaceId, + attempt: attempt + 1, + nextBaseVersion: result.version, + }); baseVersion = result.version; attempt += 1; } diff --git a/src/lib/workspace/workspace-items-projector.ts b/src/lib/workspace/workspace-items-projector.ts index 4a73237c..b01624d4 100644 --- a/src/lib/workspace/workspace-items-projector.ts +++ b/src/lib/workspace/workspace-items-projector.ts @@ -209,13 +209,15 @@ async function readProjectedWorkspaceState( getWorkspaceProjectionCheckpoint(client, workspaceId), ]); + const items = await loadProjectedItemsByShellRows( + client, + workspaceId, + shellRows as Array, + userId, + ); + return { - items: await loadProjectedItemsByShellRows( - client, - workspaceId, - shellRows as Array, - userId, - ), + items, version, }; } diff --git a/src/workflows/audio-transcribe/steps/persist-result.ts b/src/workflows/audio-transcribe/steps/persist-result.ts index 65d9b3fd..ec42b8fd 100644 --- a/src/workflows/audio-transcribe/steps/persist-result.ts +++ b/src/workflows/audio-transcribe/steps/persist-result.ts @@ -1,5 +1,6 @@ import { createEvent } from "@/lib/workspace/events"; import { appendWorkspaceEventOrThrow } from "@/lib/workspace/workspace-event-store"; +import { db, workspaceItemExtracted } from "@/lib/db/client"; import type { AudioData, Item } from "@/lib/workspace-state/types"; import type { TranscribeResult } from "./transcribe"; @@ -22,6 +23,33 @@ export async function persistAudioResult( ): Promise { "use step"; + const transcriptText = result.segments + ?.map((s) => s.content) + .join("\n") || null; + + await db + .insert(workspaceItemExtracted) + .values({ + workspaceId, + itemId, + searchText: transcriptText ?? "", + transcriptText, + transcriptSegments: result.segments as unknown as Record, + updatedAt: new Date().toISOString(), + }) + .onConflictDoUpdate({ + target: [ + workspaceItemExtracted.workspaceId, + workspaceItemExtracted.itemId, + ], + set: { + transcriptText, + transcriptSegments: result.segments as unknown as Record, + searchText: transcriptText ?? "", + updatedAt: new Date().toISOString(), + }, + }); + const event = createEvent( "ITEM_UPDATED", { @@ -29,7 +57,6 @@ export async function persistAudioResult( changes: { data: { summary: result.summary, - segments: result.segments, ...(typeof result.duration === "number" && result.duration > 0 && { duration: result.duration }), processingStatus: "complete" as const, diff --git a/src/workflows/ocr-dispatch/steps/persist-results.ts b/src/workflows/ocr-dispatch/steps/persist-results.ts index 664925de..8b599ba8 100644 --- a/src/workflows/ocr-dispatch/steps/persist-results.ts +++ b/src/workflows/ocr-dispatch/steps/persist-results.ts @@ -1,5 +1,7 @@ import { createEvent } from "@/lib/workspace/events"; import { appendWorkspaceEventOrThrow } from "@/lib/workspace/workspace-event-store"; +import { db, workspaceItemExtracted } from "@/lib/db/client"; +import { getOcrPagesTextContent } from "@/lib/utils/ocr-pages"; import type { OcrItemResult } from "@/lib/ocr/types"; import type { ImageData, Item, PdfData } from "@/lib/workspace-state/types"; @@ -10,19 +12,45 @@ export async function persistOcrResults( ): Promise { "use step"; + for (const result of results) { + const ocrPages = result.ok ? result.pages : []; + const ocrText = getOcrPagesTextContent(ocrPages) || null; + + await db + .insert(workspaceItemExtracted) + .values({ + workspaceId, + itemId: result.itemId, + searchText: ocrText ?? "", + ocrText, + ocrPages: ocrPages as unknown as Record, + updatedAt: new Date().toISOString(), + }) + .onConflictDoUpdate({ + target: [ + workspaceItemExtracted.workspaceId, + workspaceItemExtracted.itemId, + ], + set: { + ocrText, + ocrPages: ocrPages as unknown as Record, + searchText: ocrText ?? "", + updatedAt: new Date().toISOString(), + }, + }); + } + const event = createEvent( "BULK_ITEMS_PATCHED", { updates: results.map((result) => { - const dataPatch = ( + const statusPatch = ( result.ok ? { - ocrPages: result.pages, ocrStatus: "complete" as const, ocrError: undefined, } : { - ocrPages: [], ocrStatus: "failed" as const, ocrError: result.error, } @@ -31,7 +59,7 @@ export async function persistOcrResults( return { id: result.itemId, changes: { - data: dataPatch as Item["data"], + data: statusPatch as Item["data"], }, }; }),