@@ -19,9 +19,11 @@ import {
1919 user ,
2020 workspace ,
2121} from '@sim/db/schema'
22+ import { toNumberOrNull } from '@sim/utils/coerce'
2223import { generateId } from '@sim/utils/id'
24+ import { toArray , toRecord } from '@sim/utils/object'
2325import { and , eq , inArray , sql } from 'drizzle-orm'
24- import { afterAll , beforeAll , describe , expect , it , vi } from 'vitest'
26+ import { afterAll , afterEach , beforeAll , beforeEach , describe , expect , it , vi } from 'vitest'
2527
2628const fixture = vi . hoisted ( ( ) => ( {
2729 storageRoot : '' ,
@@ -92,6 +94,7 @@ import {
9294 executeMemberSync ,
9395 resumeMembershipRewrites ,
9496} from '@/lib/knowledge/connectors/member-sync-engine'
97+ import { RECONCILIATION_WINDOW_SIZE } from '@/lib/knowledge/connectors/reconciliation-window'
9598import { runConnectorContentPass } from '@/lib/knowledge/connectors/sync-content-pass'
9699import { executeSync } from '@/lib/knowledge/connectors/sync-engine'
97100import {
@@ -898,6 +901,220 @@ describe('durable source and member cycles in PostgreSQL', () => {
898901 }
899902 } )
900903
904+ describe ( 'reconciliation windows' , ( ) => {
905+ const UNRELATED_FILENAME = 'reconciliation-window-unrelated'
906+ const acl = [ `u:${ ids . aliceId } @fixture.test` ]
907+ let connectorId = ''
908+ let runId = ''
909+ let startedAt = new Date ( )
910+ beforeEach ( ( ) => {
911+ connectorId = generateId ( )
912+ runId = generateId ( )
913+ startedAt = new Date ( )
914+ } )
915+ afterEach ( async ( ) => {
916+ await db . delete ( document ) . where ( eq ( document . filename , UNRELATED_FILENAME ) )
917+ await db . delete ( document ) . where ( eq ( document . connectorId , connectorId ) )
918+ await db . delete ( knowledgeConnector ) . where ( eq ( knowledgeConnector . id , connectorId ) )
919+ } )
920+ const row = ( id : string ) => ( {
921+ id,
922+ knowledgeBaseId : ids . knowledgeBaseId ,
923+ connectorId,
924+ externalId : id ,
925+ filename : id ,
926+ fileUrl : '' ,
927+ fileSize : 0 ,
928+ mimeType : 'text/plain' ,
929+ processingStatus : 'completed' ,
930+ contentHash : `hash-${ id } ` ,
931+ acl,
932+ aclVerifiedAt : startedAt ,
933+ } )
934+ /**
935+ * Seeds a connector whose listing already completed, so the pass goes straight to
936+ * reconciliation. Other documents spread across the id space keep the connector the minority
937+ * it is at scale; a connector that is nearly the whole table may be walked through the
938+ * primary key instead, one row at a time all the same, reading the other rows between its
939+ * own.
940+ */
941+ const seed = async ( rows : ReturnType < typeof row > [ ] , listedCount : number ) => {
942+ await db . execute ( sql `
943+ INSERT INTO ${ document } (id, knowledge_base_id, filename, file_url, file_size, mime_type, processing_status)
944+ SELECT md5(${ connectorId } || g), ${ ids . knowledgeBaseId } , ${ UNRELATED_FILENAME } , '', 0, 'text/plain', 'completed'
945+ FROM generate_series(1, ${ rows . length } ) g` )
946+ const checkpoint = beginListingCheckpoint ( {
947+ fingerprint : listingFingerprint ( { connectorId } ) ,
948+ generationId : runId ,
949+ startedAt,
950+ } )
951+ checkpoint . complete = true
952+ checkpoint . listedCount = listedCount
953+ await db . insert ( knowledgeConnector ) . values ( {
954+ id : connectorId ,
955+ knowledgeBaseId : ids . knowledgeBaseId ,
956+ connectorType : 'google_drive' ,
957+ sourceConfig : { } ,
958+ accessMode : 'admin' ,
959+ status : 'syncing' ,
960+ syncLockToken : runId ,
961+ listingCheckpoint : checkpoint ,
962+ } )
963+ for ( let offset = 0 ; offset < rows . length ; offset += 1_000 )
964+ await db . insert ( document ) . values ( rows . slice ( offset , offset + 1_000 ) )
965+ await db . execute ( sql `ANALYZE document` )
966+ }
967+ const isWindowScan = ( query : string ) => / ^ \s * w i t h " d o c u m e n t " a s m a t e r i a l i z e d / i. test ( query )
968+ /** Runs the pass, returning the window statements it issued outside transactions. */
969+ const reconcile = async ( ) => {
970+ const hardDelete = vi . spyOn ( documentService , 'hardDeleteDocuments' )
971+ const statements = vi . spyOn ( db . $client , 'unsafe' )
972+ try {
973+ const [ connector ] = await db
974+ . select ( )
975+ . from ( knowledgeConnector )
976+ . where ( eq ( knowledgeConnector . id , connectorId ) )
977+ const stats = result ( )
978+ const pass = await runConnectorContentPass ( {
979+ connectorId,
980+ connector,
981+ connectorConfig : CONNECTOR_REGISTRY . google_drive ,
982+ sourceConfig : { } ,
983+ syncContext : { } ,
984+ kbOwner : { userId : ids . aliceId , workspaceId : ids . workspaceId } ,
985+ billingAttribution : billing ,
986+ result : stats ,
987+ lease : createContentSyncLease ( connectorId , runId ) ,
988+ leaseKind : 'content' ,
989+ runId,
990+ fingerprint : listingFingerprint ( { connectorId } ) ,
991+ documentAccess : 'admin' ,
992+ getAccessToken : async ( ) => 'fixture' ,
993+ hydration : { getDocument : fixture . get } ,
994+ forceRehydrate : false ,
995+ deadlineAt : Date . now ( ) + 60_000 ,
996+ } )
997+ return {
998+ pass,
999+ stats,
1000+ hardDeleted : hardDelete . mock . calls . flatMap ( ( [ batch ] ) => batch ) ,
1001+ walked : statements . mock . calls
1002+ . map ( ( [ query , params ] ) => ( { query, params } ) )
1003+ . filter ( ( { query } ) => isWindowScan ( query ) ) ,
1004+ }
1005+ } finally {
1006+ statements . mockRestore ( )
1007+ hardDelete . mockRestore ( )
1008+ }
1009+ }
1010+ /** Plans a window statement from freshly analyzed statistics, returning the document rows it read. */
1011+ const explainWalk = async ( query : string , params : Parameters < typeof db . $client . unsafe > [ 1 ] ) => {
1012+ const [ explained ] = await db . $client . begin ( async ( tx ) => {
1013+ await tx . unsafe ( 'ANALYZE document' )
1014+ return tx . unsafe ( `EXPLAIN (ANALYZE, FORMAT JSON) ${ query } ` , params )
1015+ } )
1016+ const nodes = ( node : unknown ) : Record < string , unknown > [ ] => {
1017+ const plan = toRecord ( node )
1018+ return [ plan , ...toArray ( plan . Plans ) . flatMap ( nodes ) ]
1019+ }
1020+ const plan = nodes ( toRecord ( toArray ( explained [ 'QUERY PLAN' ] ) [ 0 ] ) . Plan )
1021+ return plan
1022+ . filter ( ( node ) => node [ 'Relation Name' ] === 'document' )
1023+ . reduce (
1024+ ( total , node ) =>
1025+ total +
1026+ ( ( toNumberOrNull ( node [ 'Actual Rows' ] ) ?? 0 ) +
1027+ ( toNumberOrNull ( node [ 'Rows Removed by Filter' ] ) ?? 0 ) ) *
1028+ ( toNumberOrNull ( node [ 'Actual Loops' ] ) ?? 1 ) ,
1029+ 0
1030+ )
1031+ }
1032+ /**
1033+ * Every window statement reads at most one window of documents, whichever connector index
1034+ * the planner walks: a read of the whole connector, or of every tombstone it has, is what
1035+ * must never happen.
1036+ */
1037+ const expectBounded = async (
1038+ walked : { query : string ; params : Parameters < typeof db . $client . unsafe > [ 1 ] } [ ]
1039+ ) => {
1040+ expect ( walked . length ) . toBeGreaterThan ( 0 )
1041+ for ( const { query, params } of walked ) {
1042+ expect ( await explainWalk ( query , params ) , query ) . toBeLessThanOrEqual (
1043+ RECONCILIATION_WINDOW_SIZE
1044+ )
1045+ }
1046+ }
1047+
1048+ it ( 'bounds every page by the ids it scans when absence is rare and late in id order' , async ( ) => {
1049+ /** Ids sort present rows first, so each absent row is found only after three full windows. */
1050+ const present = Array . from ( { length : 3 * RECONCILIATION_WINDOW_SIZE } , ( _ , index ) => ( {
1051+ ...row ( `${ connectorId } -a-${ String ( index ) . padStart ( 6 , '0' ) } ` ) ,
1052+ sourceSeenAt : startedAt ,
1053+ } ) )
1054+ const absent = Array . from ( { length : 3 } , ( _ , index ) => ( {
1055+ ...row ( `${ connectorId } -z-live-${ index } ` ) ,
1056+ sourceSeenAt : null ,
1057+ } ) )
1058+ const tombstoned = Array . from ( { length : 2 } , ( _ , index ) => ( {
1059+ ...row ( `${ connectorId } -z-tombstone-${ index } ` ) ,
1060+ sourceSeenAt : null ,
1061+ deletedAt : new Date ( startedAt . getTime ( ) - 60_000 ) ,
1062+ } ) )
1063+ await seed ( [ ...present , ...absent , ...tombstoned ] , present . length )
1064+ const { pass, stats, hardDeleted, walked } = await reconcile ( )
1065+ expect ( pass ) . toMatchObject ( { complete : true , holdNotice : null } )
1066+ expect ( stats . docsDeleted ) . toBe ( absent . length )
1067+ expect ( hardDeleted . sort ( ) ) . toEqual ( tombstoned . map ( ( item ) => item . id ) . sort ( ) )
1068+ const stored = await db
1069+ . select ( { id : document . id , acl : document . acl , deletedAt : document . deletedAt } )
1070+ . from ( document )
1071+ . where ( eq ( document . connectorId , connectorId ) )
1072+ expect ( stored ) . toHaveLength ( present . length + absent . length )
1073+ for ( const item of stored ) expect ( item . deletedAt !== null ) . toBe ( item . id . includes ( '-z-' ) )
1074+ expect (
1075+ stored . filter ( ( item ) => item . id . includes ( '-z-' ) ) . every ( ( item ) => item . acl . length === 0 )
1076+ ) . toBe ( true )
1077+ await expectBounded ( walked )
1078+ } , 120_000 )
1079+
1080+ it ( 'scans a dense window once and never reads every tombstone of the connector' , async ( ) => {
1081+ /**
1082+ * Three hard-delete pages of absent tombstones open the first window; the rest of the
1083+ * connector is tombstones the listing still sees, which the hard walk must pass over.
1084+ */
1085+ const dense = Array . from ( { length : 75 } , ( _ , index ) => ( {
1086+ ...row ( `${ connectorId } -a-${ String ( index ) . padStart ( 3 , '0' ) } ` ) ,
1087+ acl : [ ] ,
1088+ sourceSeenAt : null ,
1089+ deletedAt : new Date ( startedAt . getTime ( ) - 60_000 ) ,
1090+ } ) )
1091+ const seenTombstones = Array . from ( { length : 3 * RECONCILIATION_WINDOW_SIZE } , ( _ , index ) => ( {
1092+ ...row ( `${ connectorId } -b-${ String ( index ) . padStart ( 6 , '0' ) } ` ) ,
1093+ sourceSeenAt : startedAt ,
1094+ deletedAt : new Date ( startedAt . getTime ( ) - 60_000 ) ,
1095+ } ) )
1096+ await seed ( [ ...dense , ...seenTombstones ] , seenTombstones . length )
1097+ const { pass, hardDeleted, walked } = await reconcile ( )
1098+ expect ( pass ) . toMatchObject ( { complete : true , holdNotice : null } )
1099+ expect ( hardDeleted . sort ( ) ) . toEqual ( dense . map ( ( item ) => item . id ) . sort ( ) )
1100+ expect (
1101+ await db
1102+ . select ( { id : document . id } )
1103+ . from ( document )
1104+ . where ( eq ( document . connectorId , connectorId ) )
1105+ ) . toHaveLength ( seenTombstones . length )
1106+ /**
1107+ * The only walk is the hard one: one statement per window, the dense first one included,
1108+ * then the tail, however many pages of matches a window holds.
1109+ */
1110+ expect ( walked ) . toHaveLength (
1111+ Math . floor ( ( dense . length + seenTombstones . length ) / RECONCILIATION_WINDOW_SIZE ) + 1
1112+ )
1113+ /** `deleted_at < $1` implies the tombstone index, which would read every connector tombstone. */
1114+ await expectBounded ( walked )
1115+ } , 120_000 )
1116+ } )
1117+
9011118 it ( 'indexes a page, resumes under a new lease, and reconciles absence only after EOF' , async ( ) => {
9021119 await db
9031120 . update ( knowledgeConnector )
0 commit comments