Skip to content

Commit dde4920

Browse files
committed
fix(knowledge): claim a member only under a proved lease and read email conversation dates
1 parent 209f328 commit dde4920

3 files changed

Lines changed: 28 additions & 11 deletions

File tree

apps/sim/lib/knowledge/connectors/member-sync-engine.ts

Lines changed: 20 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ import {
5252
MEMBER_SYNC_SOFT_BUDGET_SECONDS,
5353
} from '@/lib/knowledge/connectors/sync-limits'
5454
import {
55+
assertSyncLeaseHeldInTx,
5556
createMemberSyncLease,
5657
holdsMemberSyncLockToken,
5758
MEMBER_LOCKABLE_CONNECTOR_STATUSES,
@@ -660,15 +661,22 @@ async function reconcileMembership(
660661
* member be aborted at the deadline without touching the others.
661662
*/
662663
async function claimNextMember(run: MemberSyncRun): Promise<MemberRow | null> {
663-
const [claimed] = await db
664-
.update(knowledgeConnectorMember)
665-
.set({ lastStartedAt: new Date(), updatedAt: new Date() })
666-
.where(
667-
and(
668-
eq(knowledgeConnectorMember.connectorId, run.connectorId),
669-
eq(
670-
knowledgeConnectorMember.id,
671-
sql`(
664+
/**
665+
* Proved under the lease: a run reclaimed while it slept must not stamp
666+
* `lastStartedAt`, which would hide the member from its replacement's
667+
* selection and defer that member's access updates to a later run.
668+
*/
669+
const [claimed] = await db.transaction(async (tx) => {
670+
await assertSyncLeaseHeldInTx(tx, run.connectorId, run.lease)
671+
return tx
672+
.update(knowledgeConnectorMember)
673+
.set({ lastStartedAt: new Date(), updatedAt: new Date() })
674+
.where(
675+
and(
676+
eq(knowledgeConnectorMember.connectorId, run.connectorId),
677+
eq(
678+
knowledgeConnectorMember.id,
679+
sql`(
672680
SELECT ${knowledgeConnectorMember.id} FROM ${knowledgeConnectorMember}
673681
WHERE ${knowledgeConnectorMember.connectorId} = ${run.connectorId}
674682
AND ${knowledgeConnectorMember.status} = 'active'
@@ -678,10 +686,11 @@ async function claimNextMember(run: MemberSyncRun): Promise<MemberRow | null> {
678686
LIMIT 1
679687
FOR UPDATE SKIP LOCKED
680688
)`
689+
)
681690
)
682691
)
683-
)
684-
.returning()
692+
.returning()
693+
})
685694
return claimed ?? null
686695
}
687696

apps/sim/lib/knowledge/connectors/source-modified-at.test.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,12 @@ describe('resolveSourceModifiedAt', () => {
5151
expect(resolveSourceModifiedAt({ modifiedTime: 1e20 })).toBeNull()
5252
})
5353

54+
it('reads the newest message time an email conversation reports', () => {
55+
expect(
56+
resolveSourceModifiedAt({ lastMessageDate: '2026-08-29T09:30:00Z' })?.toISOString()
57+
).toBe('2026-08-29T09:30:00.000Z')
58+
})
59+
5460
it('reads the last activity a chat space or channel reports', () => {
5561
expect(resolveSourceModifiedAt({ lastActivity: '2026-08-30T10:00:00Z' })?.toISOString()).toBe(
5662
'2026-08-30T10:00:00.000Z'

apps/sim/lib/knowledge/connectors/source-modified-at.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@ const SOURCE_MODIFIED_AT_KEYS = [
1818
'statusDate',
1919
/** Google Chat spaces and Teams channels: the latest message time, the listing's only change signal. */
2020
'lastActivity',
21+
/** Gmail and Outlook conversations: the newest message's time. */
22+
'lastMessageDate',
2123
] as const
2224

2325
/** Earlier than any plausible document; guards against epoch-zero placeholders. */

0 commit comments

Comments
 (0)