Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion apps/sim/ee/workspace-forking/lib/mapping/cascade.ts
Original file line number Diff line number Diff line change
Expand Up @@ -165,7 +165,8 @@ export async function detectForkCascadeReferences(params: {
inArray(knowledgeConnector.knowledgeBaseId, Array.from(knowledgeBaseIds)),
eq(knowledgeBase.workspaceId, sourceWorkspaceId),
isNull(knowledgeBase.deletedAt),
isNull(knowledgeConnector.deletedAt)
isNull(knowledgeConnector.deletedAt),
isNull(knowledgeConnector.detachedAt)
)
)
for (const connector of connectors) {
Expand Down
1 change: 1 addition & 0 deletions apps/sim/lib/billing/storage/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ export {
applyStorageUsageDeltasInTx,
checkAndIncrementStorageUsageInTx,
decrementStorageUsageForBillingContextInTx,
incrementAdmittedStorageUsageForBillingContextInTx,
incrementStorageUsageForBillingContextInTx,
type LegacyStorageUsageDelta,
maybeNotifyStorageLimitForBillingContext,
Expand Down
11 changes: 10 additions & 1 deletion apps/sim/lib/billing/storage/payer-transfer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,12 @@ vi.mock('@sim/db/schema', () => ({
id: 'knowledgeBase.id',
workspaceId: 'knowledgeBase.workspaceId',
},
knowledgeConnector: {
__table: 'knowledgeConnector',
detachedAt: 'knowledgeConnector.detachedAt',
detachReservedBytes: 'knowledgeConnector.detachReservedBytes',
knowledgeBaseId: 'knowledgeConnector.knowledgeBaseId',
},
organization: {
__table: 'organization',
id: 'organization.id',
Expand Down Expand Up @@ -535,7 +541,10 @@ describe('changeWorkspaceStoragePayerInTx', () => {
expect(query.values).not.toContain('workspaceFiles.deletedAt')
expect(query.values).toContain('document.connectorId')
expect(query.values).toContain('document.deletedAt')
expect(query.values.filter((value) => value === 'workspace-1')).toHaveLength(3)
/** A detaching connector's reservation is already charged, so a payer move carries it. */
expect(query.values).toContain('knowledgeConnector.detachReservedBytes')
expect(query.values).toContain('knowledgeConnector.detachedAt')
expect(query.values.filter((value) => value === 'workspace-1')).toHaveLength(4)
})

it('fails closed when a billable file is missing canonical size metadata', async () => {
Expand Down
25 changes: 24 additions & 1 deletion apps/sim/lib/billing/storage/payer-transfer.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import {
document,
knowledgeBase,
knowledgeConnector,
organization,
userStats,
workspace,
Expand Down Expand Up @@ -77,7 +78,8 @@ function parseExactBytes(value: number | string, label: string): number {
* Computes one workspace's live billable bytes with two index-bounded scalar
* aggregates. Archived workspace files and documents remain billable while
* their objects are retained; mothership files, connector documents, and
* deleted documents are excluded.
* deleted documents are excluded. A detaching connector's reservation counts:
* it was charged when the connector was removed with its documents kept.
*/
async function getExactWorkspaceStorageBytes(tx: DbOrTx, workspaceId: string): Promise<number> {
const [row] = await tx.execute<ExactWorkspaceStorageRow>(sql`
Expand All @@ -103,6 +105,13 @@ async function getExactWorkspaceStorageBytes(tx: DbOrTx, workspaceId: string): P
WHERE ${knowledgeBase.workspaceId} = ${workspaceId}
AND ${document.connectorId} IS NULL
AND ${document.deletedAt} IS NULL
), 0)::bigint + COALESCE((
SELECT SUM(${knowledgeConnector.detachReservedBytes})
FROM ${knowledgeConnector}
INNER JOIN ${knowledgeBase}
ON ${knowledgeBase.id} = ${knowledgeConnector.knowledgeBaseId}
WHERE ${knowledgeBase.workspaceId} = ${workspaceId}
AND ${knowledgeConnector.detachedAt} IS NOT NULL
), 0)::bigint AS document_bytes
`)

Expand Down Expand Up @@ -209,6 +218,20 @@ async function getExactWorkspaceStorageBytesBatch(
AND ${document.connectorId} IS NULL
AND ${document.deletedAt} IS NULL
GROUP BY ${knowledgeBase.workspaceId}

UNION ALL

SELECT
${knowledgeBase.workspaceId} AS workspace_id,
0::bigint AS workspace_file_bytes,
SUM(${knowledgeConnector.detachReservedBytes}) AS document_bytes,
0::bigint AS workspace_file_missing_size_count
FROM ${knowledgeConnector}
INNER JOIN ${knowledgeBase}
ON ${knowledgeBase.id} = ${knowledgeConnector.knowledgeBaseId}
WHERE ${inArray(knowledgeBase.workspaceId, workspaceIds)}
AND ${knowledgeConnector.detachedAt} IS NOT NULL
GROUP BY ${knowledgeBase.workspaceId}
) storage_by_workspace
GROUP BY storage_by_workspace.workspace_id
ORDER BY storage_by_workspace.workspace_id
Expand Down
23 changes: 23 additions & 0 deletions apps/sim/lib/billing/storage/tracking.ts
Original file line number Diff line number Diff line change
Expand Up @@ -590,6 +590,29 @@ export async function incrementStorageUsageForBillingContextInTx(
return result.updatedUsage
}

/**
* Increments one workspace and its current payer for bytes whose admission was
* already decided, such as documents a connector removal accepted and a
* background job now releases page by page. Never refuses: a page that could
* cross the limit after admission would otherwise leave the release half done.
*/
export async function incrementAdmittedStorageUsageForBillingContextInTx(
tx: DbOrTx,
context: StorageBillingContext,
bytes: number
): Promise<number | undefined> {
if (bytes <= 0) return undefined
const result = await mutateWorkspaceStorageUsage(
tx,
context.workspaceId,
bytes,
'increment',
undefined,
context
)
return result.updatedUsage
}

/**
* Atomically check quota and increment a user's (or their org's) storage
* counter inside an existing transaction, using a pre-resolved subscription.
Expand Down
9 changes: 2 additions & 7 deletions apps/sim/lib/copilot/chat/workspace-context.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import {
getAccessibleOAuthCredentials,
} from '@/lib/credentials/environment'
import { listWorkspaceSandboxes } from '@/lib/execution/remote-sandbox/workspace-sandboxes'
import { connectorIsLive } from '@/lib/knowledge/connectors/sync-lock'
import { listCustomBlockSummariesForWorkspace } from '@/lib/workflows/custom-blocks/operations'
import { listCustomTools } from '@/lib/workflows/custom-tools/operations'
import { listSkillsForUser } from '@/lib/workflows/skills/operations'
Expand Down Expand Up @@ -457,13 +458,7 @@ async function buildWorkspaceMdData(
connectorType: knowledgeConnector.connectorType,
})
.from(knowledgeConnector)
.where(
and(
inArray(knowledgeConnector.knowledgeBaseId, kbIds),
isNull(knowledgeConnector.archivedAt),
isNull(knowledgeConnector.deletedAt)
)
)
.where(and(inArray(knowledgeConnector.knowledgeBaseId, kbIds), connectorIsLive()))
: []
const connectorTypesByKb = new Map<string, string[]>()
for (const row of connectorRows) {
Expand Down
27 changes: 27 additions & 0 deletions apps/sim/lib/knowledge/__integration__/drain-connector-event.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
import { db } from '@sim/db'
import { outboxEvent } from '@sim/db/schema'
import { and, eq, sql } from 'drizzle-orm'
import { expect } from 'vitest'
import { processOutboxEventById } from '@/lib/core/outbox/service'
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'

/** Runs a connector's queued removal event until it completes, as the outbox worker would. */
export async function drainConnectorEvent(connectorId: string, eventType: string): Promise<void> {
const [job] = await db
.select()
.from(outboxEvent)
.where(
and(
eq(outboxEvent.eventType, eventType),
sql`${outboxEvent.payload}->>'connectorId' = ${connectorId}`
)
)
.limit(1)
expect(job).toBeDefined()
let status = await processOutboxEventById(job.id, knowledgeDocumentProcessingOutboxHandlers)
for (let attempt = 0; status === 'pending' && attempt < 20; attempt++) {
await db.update(outboxEvent).set({ availableAt: new Date() }).where(eq(outboxEvent.id, job.id))
status = await processOutboxEventById(job.id, knowledgeDocumentProcessingOutboxHandlers)
}
expect(status).toBe('completed')
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import { and, eq, inArray } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
import type { ConnectorDocumentFilter } from '@/lib/api/contracts/knowledge/connectors'
import * as embeddings from '@/lib/embeddings'
import { drainConnectorEvent } from '@/lib/knowledge/__integration__/drain-connector-event'
import {
createKnowledgeAclFixtureIds,
seedKnowledgeAclFixture,
Expand All @@ -29,6 +30,7 @@ import {
} from '@/lib/knowledge/application/documents'
import { readSearchSourceProgress } from '@/lib/knowledge/application/search-source-progress'
import { listSearchSources } from '@/lib/knowledge/application/search-sources'
import { KNOWLEDGE_CONNECTOR_DETACH_EVENT } from '@/lib/knowledge/connectors/detachment'
import { createContentSyncLease } from '@/lib/knowledge/connectors/sync-lock'
import { persistSkippedDocuments } from '@/lib/knowledge/connectors/sync-persistence'
import * as documentProcessor from '@/lib/knowledge/documents/document-processor'
Expand Down Expand Up @@ -630,6 +632,7 @@ describe('intentional skips and genuine failures across document reads', () => {
input: { ...scope, deleteDocuments: false },
})
expect(result).toMatchObject({ documentsDeleted: 0, documentsKept: 6 })
await drainConnectorEvent(fixture.connectorId, KNOWLEDGE_CONNECTOR_DETACH_EVENT)
const retained = await db
.select()
.from(document)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import { db } from '@sim/db'
import {
document,
embedding,
embeddingSearch,
knowledgeBase,
knowledgeConnector,
organization,
Expand All @@ -14,19 +15,19 @@ import {
workspace,
} from '@sim/db/schema'
import { generateId } from '@sim/utils/id'
import { and, eq, inArray, isNull, sql } from 'drizzle-orm'
import { and, eq, inArray, isNotNull, isNull, sql } from 'drizzle-orm'
import { afterAll, describe, expect, it, vi } from 'vitest'
import { processOutboxEventById } from '@/lib/core/outbox/service'
import { drainConnectorEvent } from '@/lib/knowledge/__integration__/drain-connector-event'
import {
createKnowledgeAclFixtureIds,
seedKnowledgeAclFixture,
} from '@/lib/knowledge/__integration__/seed-source-access-fixture'
import { WORKSPACE_ACCESS_SCOPE } from '@/lib/knowledge/access/scope'
import { SYSTEM_ACCESS_SCOPE } from '@/lib/knowledge/access/types'
import { KNOWLEDGE_CONNECTOR_CLEANUP_EVENT } from '@/lib/knowledge/connectors/deletion'
import { KNOWLEDGE_CONNECTOR_DETACH_EVENT } from '@/lib/knowledge/connectors/detachment'
import { createContentSyncLease, SyncLockLostException } from '@/lib/knowledge/connectors/sync-lock'
import { persistSkippedDocuments } from '@/lib/knowledge/connectors/sync-persistence'
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
import {
createSingleDocument,
getKnowledgeDocument,
Expand Down Expand Up @@ -104,13 +105,23 @@ function disconnect(ids: Fixture, deleteDocuments = false) {
})
}

/** Removes the source keeping its documents, then runs the background release to completion. */
async function detach(ids: Fixture) {
const outcome = await disconnect(ids)
if (outcome.success) await drainConnectorEvent(ids.connectorId, KNOWLEDGE_CONNECTOR_DETACH_EVENT)
return outcome
}

afterAll(async () => {
for (const ids of fixtures) {
await db
.delete(outboxEvent)
.where(
and(
eq(outboxEvent.eventType, KNOWLEDGE_CONNECTOR_CLEANUP_EVENT),
inArray(outboxEvent.eventType, [
KNOWLEDGE_CONNECTOR_CLEANUP_EVENT,
KNOWLEDGE_CONNECTOR_DETACH_EVENT,
]),
sql`${outboxEvent.payload}->>'knowledgeBaseId' = ${ids.knowledgeBaseId}`
)
)
Expand All @@ -130,7 +141,7 @@ describe('knowledge document storage ledgers', () => {
const transaction = db.transaction.bind(db)
/** Interleave real operations at the lock boundary; every query and commit still uses PostgreSQL. */
const detachBeforeLock: typeof db.transaction = async (callback, config) => {
expect(await disconnect(ids)).toMatchObject({ success: true })
expect(await detach(ids)).toMatchObject({ success: true })
expect(await ledger(ids)).toEqual({ workspaceBytes: 37, payerBytes: 37 })
return transaction(callback, config)
}
Expand Down Expand Up @@ -167,7 +178,19 @@ describe('knowledge document storage ledgers', () => {
{ success: true, documentsKept: 6, documentsDeleted: 0 },
])
expect(outcomes.filter((result) => !result.success)).toHaveLength(1)
/** Kept bytes are reserved at removal; the documents stay attached and readable until released. */
expect(await ledger(ids)).toEqual({ workspaceBytes: 70, payerBytes: 70 })
expect(
await getKnowledgeDocument(ids.knowledgeBaseId, source[0].id, WORKSPACE_ACCESS_SCOPE)
).not.toBeNull()
await drainConnectorEvent(ids.connectorId, KNOWLEDGE_CONNECTOR_DETACH_EVENT)
expect(await ledger(ids)).toEqual({ workspaceBytes: 70, payerBytes: 70 })
expect(
await db
.select({ id: knowledgeConnector.id })
.from(knowledgeConnector)
.where(eq(knowledgeConnector.id, ids.connectorId))
).toHaveLength(0)
const retained = await db
.select({ id: document.id, connectorId: document.connectorId, deletedAt: document.deletedAt })
.from(document)
Expand Down Expand Up @@ -197,6 +220,59 @@ describe('knowledge document storage ledgers', () => {
expect(await ledger(ids)).toEqual({ workspaceBytes: 0, payerBytes: 0 })
})

it('settles the reservation of a kept document deleted before its release', async () => {
const ids = await seed()
const [kept, removed] = [sourceDocument(ids, 37), sourceDocument(ids, 5)]
await db.insert(document).values([kept, removed])

expect(await disconnect(ids)).toMatchObject({ success: true, documentsKept: 2 })
expect(await ledger(ids)).toEqual({ workspaceBytes: 42, payerBytes: 42 })
expect(await hardDeleteDocuments([removed.id], generateId())).toBe(1)
expect(await ledger(ids)).toEqual({ workspaceBytes: 42, payerBytes: 42 })

await drainConnectorEvent(ids.connectorId, KNOWLEDGE_CONNECTOR_DETACH_EVENT)
expect(await ledger(ids)).toEqual({ workspaceBytes: 37, payerBytes: 37 })
expect(await hardDeleteDocuments([kept.id], generateId())).toBe(1)
expect(await ledger(ids)).toEqual({ workspaceBytes: 0, payerBytes: 0 })
})

it('releases a document larger than one page without a search row still naming its source', async () => {
const ids = await seed()
const row = sourceDocument(ids, 5)
await db.insert(document).values(row)
await db.insert(embedding).values(
Array.from({ length: 600 }, (_, chunkIndex) => ({
id: generateId(),
knowledgeBaseId: ids.knowledgeBaseId,
documentId: row.id,
chunkIndex,
chunkHash: `hash-${chunkIndex}`,
content: 'test chunk',
contentLength: 10,
tokenCount: 2,
startOffset: 0,
endOffset: 10,
embedding384: Array(384).fill(0.1),
}))
)
const sourceRows = () =>
db
.select({ count: sql<number>`COUNT(*)::integer` })
.from(embeddingSearch)
.where(and(eq(embeddingSearch.documentId, row.id), isNotNull(embeddingSearch.connectorId)))
expect((await sourceRows())[0].count).toBe(600)

expect(await detach(ids)).toEqual({ success: true, documentsKept: 1, documentsDeleted: 0 })

expect((await sourceRows())[0].count).toBe(0)
const [released] = await db
.select({ connectorId: document.connectorId })
.from(document)
.where(eq(document.id, row.id))
expect(released.connectorId).toBeNull()
expect(await ledger(ids)).toEqual({ workspaceBytes: 5, payerBytes: 5 })
})

it('hides a source immediately and cleans bounded batches without debiting manual storage', async () => {
const ids = await seed()
await manualDocument(ids, 31)
Expand Down Expand Up @@ -240,17 +316,6 @@ describe('knowledge document storage ledgers', () => {
for (const access of [SYSTEM_ACCESS_SCOPE, WORKSPACE_ACCESS_SCOPE]) {
expect(await getKnowledgeDocument(ids.knowledgeBaseId, rows[0].id, access)).toBeNull()
}
const [job] = await db
.select()
.from(outboxEvent)
.where(
and(
eq(outboxEvent.eventType, KNOWLEDGE_CONNECTOR_CLEANUP_EVENT),
sql`${outboxEvent.payload}->>'connectorId' = ${ids.connectorId}`
)
)
.limit(1)
expect(job).toBeDefined()
await expect(
persistSkippedDocuments(
ids.knowledgeBaseId,
Expand All @@ -274,17 +339,7 @@ describe('knowledge document storage ledgers', () => {
createContentSyncLease(ids.connectorId, ids.lockId)
)
).rejects.toBeInstanceOf(SyncLockLostException)
const handlers = knowledgeDocumentProcessingOutboxHandlers
let status = await processOutboxEventById(job.id, handlers)
expect(status).toBe('pending')
for (let attempt = 0; status === 'pending' && attempt < 5; attempt++) {
await db
.update(outboxEvent)
.set({ availableAt: new Date() })
.where(eq(outboxEvent.id, job.id))
status = await processOutboxEventById(job.id, handlers)
}
expect(status).toBe('completed')
await drainConnectorEvent(ids.connectorId, KNOWLEDGE_CONNECTOR_CLEANUP_EVENT)
expect(
await db
.select({ id: knowledgeConnector.id })
Expand Down Expand Up @@ -398,7 +453,7 @@ describe('knowledge document storage ledgers', () => {
it('serializes ordinary uploads against source detachment without losing either charge', async () => {
const ids = await seed()
await db.insert(document).values([sourceDocument(ids, 37)])
const [detached, manual] = await Promise.all([disconnect(ids), manualDocument(ids, 41)])
const [detached, manual] = await Promise.all([detach(ids), manualDocument(ids, 41)])
expect(detached).toMatchObject({ success: true })
expect(await ledger(ids)).toEqual({ workspaceBytes: 78, payerBytes: 78 })
const [source] = await db
Expand Down
Loading
Loading