From 8e6f1316c4f4b12aaf64029518e4e229b4c1a1e1 Mon Sep 17 00:00:00 2001 From: Waleed Date: Sun, 22 Mar 2026 03:41:45 -0700 Subject: [PATCH] fix(kb): store filename with .txt extension for connector documents (#3707) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(kb): store filename with .txt extension for connector documents Connector documents (e.g. Fireflies transcripts) have titles without file extensions. The DB stored the raw title as filename, but the processing pipeline extracts file extension from filename to determine the parser. On retry/reprocess, this caused "Unsupported file type" errors with the document title treated as the extension. Now stores processingFilename (which includes .txt) instead of the raw title, consistent with what was actually uploaded to storage. * fix(kb): guard stuck document retry against filenames without extension Existing DB rows may have connector document filenames stored without a .txt extension (raw meeting titles). The stuck-doc retry path reads filename from DB and passes it to parseHttpFile, which extracts the extension via split('.'). When there's no dot, the entire title becomes the "extension", causing "Unsupported file type" errors. Falls back to 'document.txt' when the stored filename has no extension. * fix(kb): fix race condition in stuck document retry during sync The stuck document retry at the end of each sync was querying for all documents with processingStatus 'pending' or 'failed'. This included documents added in the CURRENT sync that were still processing asynchronously, causing duplicate concurrent processing attempts. The race between the original (correct) processing and the retry (which reads the raw title from DB as filename) produced nondeterministic failures — some documents would succeed while others would fail with "Unsupported file type: ". Fixes: - Filter stuck doc query by uploadedAt < syncStartedAt to exclude documents from the current sync - Pass mimeType through to parseHttpFile so text/plain content can be decoded directly without requiring a file extension in the filename (matches parseDataURI which already handles this) - Restore filename as extDoc.title in DB (the display name, not the processing filename) * fix(kb): fix race condition in stuck document retry during sync The stuck document retry at the end of each sync was querying for all documents with processingStatus 'pending' or 'failed'. This included documents added in the CURRENT sync that were still processing asynchronously, causing duplicate concurrent processing attempts. The race between the original (correct) processing and the retry (which reads the raw title from DB as filename) produced nondeterministic failures — some documents would succeed while others would fail with "Unsupported file type: ". Fixes: - Filter stuck doc query by uploadedAt < syncStartedAt to exclude documents from the current sync - Pass mimeType through to parseHttpFile and use existing getExtensionFromMimeType utility as fallback when filename has no extension (e.g. Fireflies meeting titles) - Apply same mimeType fallback in parseDataURI for consistency * lint * fix(kb): handle empty extension edge case in parseDataURI When filename ends with a dot (e.g. "file."), split('.').pop() returns an empty string. Fall through to mimeType-based extension lookup instead of passing empty string to parseBuffer. Co-Authored-By: Claude Opus 4.6 --------- Co-authored-by: Claude Opus 4.6 --- apps/sim/hooks/queries/kb/connectors.ts | 19 +++++++++++++++++-- .../lib/knowledge/connectors/sync-engine.ts | 13 +++++++++---- .../knowledge/documents/document-processor.ts | 19 +++++++++++++------ 3 files changed, 39 insertions(+), 12 deletions(-) diff --git a/apps/sim/hooks/queries/kb/connectors.ts b/apps/sim/hooks/queries/kb/connectors.ts index fd86f306e0..1b737d8f2a 100644 --- a/apps/sim/hooks/queries/kb/connectors.ts +++ b/apps/sim/hooks/queries/kb/connectors.ts @@ -88,6 +88,21 @@ async function fetchConnectorDetail( return result.data } +/** Stop polling for initial sync after 2 minutes */ +const PENDING_SYNC_WINDOW_MS = 2 * 60 * 1000 + +/** + * Checks if a connector is syncing or awaiting its first sync within the allowed window + */ +export function isConnectorSyncingOrPending(connector: ConnectorData): boolean { + if (connector.status === 'syncing') return true + return ( + connector.status === 'active' && + !connector.lastSyncAt && + Date.now() - new Date(connector.createdAt).getTime() < PENDING_SYNC_WINDOW_MS + ) +} + export function useConnectorList(knowledgeBaseId?: string) { return useQuery({ queryKey: connectorKeys.list(knowledgeBaseId), @@ -97,8 +112,8 @@ export function useConnectorList(knowledgeBaseId?: string) { placeholderData: keepPreviousData, refetchInterval: (query) => { const connectors = query.state.data - const hasSyncing = connectors?.some((c) => c.status === 'syncing') - return hasSyncing ? 3000 : false + if (!connectors?.length) return false + return connectors.some(isConnectorSyncingOrPending) ? 3000 : false }, }) } diff --git a/apps/sim/lib/knowledge/connectors/sync-engine.ts b/apps/sim/lib/knowledge/connectors/sync-engine.ts index a07005a95d..efef605f52 100644 --- a/apps/sim/lib/knowledge/connectors/sync-engine.ts +++ b/apps/sim/lib/knowledge/connectors/sync-engine.ts @@ -6,7 +6,7 @@ import { knowledgeConnectorSyncLog, } from '@sim/db/schema' import { createLogger } from '@sim/logger' -import { and, eq, inArray, isNull, ne, sql } from 'drizzle-orm' +import { and, eq, inArray, isNull, lt, ne, sql } from 'drizzle-orm' import { decryptApiKey } from '@/lib/api-key/crypto' import { getInternalApiBaseUrl } from '@/lib/core/utils/urls' import { @@ -272,11 +272,12 @@ export async function executeSync( } const syncLogId = crypto.randomUUID() + const syncStartedAt = new Date() await db.insert(knowledgeConnectorSyncLog).values({ id: syncLogId, connectorId, status: 'started', - startedAt: new Date(), + startedAt: syncStartedAt, }) let syncExitedCleanly = false @@ -536,19 +537,23 @@ export async function executeSync( throw new Error(`Knowledge base ${connector.knowledgeBaseId} was deleted during sync`) } - // Retry stuck documents that failed or never completed processing + // Retry stuck documents that failed or never completed processing. + // Only retry docs uploaded BEFORE this sync — docs added in the current sync + // are still processing asynchronously and would cause a duplicate processing race. const stuckDocs = await db .select({ id: document.id, fileUrl: document.fileUrl, filename: document.filename, fileSize: document.fileSize, + mimeType: document.mimeType, }) .from(document) .where( and( eq(document.connectorId, connectorId), inArray(document.processingStatus, ['pending', 'failed']), + lt(document.uploadedAt, syncStartedAt), eq(document.userExcluded, false), isNull(document.archivedAt), isNull(document.deletedAt) @@ -565,7 +570,7 @@ export async function executeSync( filename: doc.filename ?? 'document.txt', fileUrl: doc.fileUrl ?? '', fileSize: doc.fileSize ?? 0, - mimeType: 'text/plain', + mimeType: doc.mimeType ?? 'text/plain', }, {} ).catch((error) => { diff --git a/apps/sim/lib/knowledge/documents/document-processor.ts b/apps/sim/lib/knowledge/documents/document-processor.ts index 0185de495b..72bf9007c9 100644 --- a/apps/sim/lib/knowledge/documents/document-processor.ts +++ b/apps/sim/lib/knowledge/documents/document-processor.ts @@ -7,7 +7,7 @@ import { parseBuffer, parseFile } from '@/lib/file-parsers' import type { FileParseMetadata } from '@/lib/file-parsers/types' import { retryWithExponentialBackoff } from '@/lib/knowledge/documents/utils' import { StorageService } from '@/lib/uploads' -import { isInternalFileUrl } from '@/lib/uploads/utils/file-utils' +import { getExtensionFromMimeType, isInternalFileUrl } from '@/lib/uploads/utils/file-utils' import { downloadFileFromUrl } from '@/lib/uploads/utils/file-utils.server' import { mistralParserTool } from '@/tools/mistral/parser' @@ -727,7 +727,7 @@ async function parseWithFileParser(fileUrl: string, filename: string, mimeType: if (fileUrl.startsWith('data:')) { content = await parseDataURI(fileUrl, filename, mimeType) } else if (fileUrl.startsWith('http')) { - const result = await parseHttpFile(fileUrl, filename) + const result = await parseHttpFile(fileUrl, filename, mimeType) content = result.content metadata = result.metadata || {} } else { @@ -759,7 +759,10 @@ async function parseDataURI(fileUrl: string, filename: string, mimeType: string) : decodeURIComponent(base64Data) } - const extension = filename.split('.').pop()?.toLowerCase() || 'txt' + const extension = + (filename.includes('.') ? filename.split('.').pop()?.toLowerCase() : undefined) || + getExtensionFromMimeType(mimeType) || + 'txt' const buffer = Buffer.from(base64Data, 'base64') const result = await parseBuffer(buffer, extension) return result.content @@ -767,13 +770,17 @@ async function parseDataURI(fileUrl: string, filename: string, mimeType: string) async function parseHttpFile( fileUrl: string, - filename: string + filename: string, + mimeType?: string ): Promise<{ content: string; metadata?: FileParseMetadata }> { const buffer = await downloadFileWithTimeout(fileUrl) - const extension = filename.split('.').pop()?.toLowerCase() + let extension = filename.includes('.') ? filename.split('.').pop()?.toLowerCase() : undefined + if (!extension && mimeType) { + extension = getExtensionFromMimeType(mimeType) ?? undefined + } if (!extension) { - throw new Error(`Could not determine file extension: ${filename}`) + throw new Error(`Could not determine file type for: ${filename}`) } const result = await parseBuffer(buffer, extension)