diff --git a/apps/web/src/lib/cloud-agent/bitbucket-integration-helpers.ts b/apps/web/src/lib/cloud-agent/bitbucket-integration-helpers.ts index b020bfe715..8d5e12790a 100644 --- a/apps/web/src/lib/cloud-agent/bitbucket-integration-helpers.ts +++ b/apps/web/src/lib/cloud-agent/bitbucket-integration-helpers.ts @@ -4,6 +4,10 @@ import { and, eq, isNull } from 'drizzle-orm'; import { z } from 'zod'; import { db } from '@/lib/drizzle'; import { PLATFORM } from '@/lib/integrations/core/constants'; +import { + withRepositoryReadDeadline, + type RepositoryReadOptions, +} from '@/lib/integrations/core/repository-read-limits'; import { BitbucketOrganizationRepositoryListResultSchema, type BitbucketOrganizationRepositoryListResult, @@ -15,7 +19,7 @@ import { } from '@/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache'; import { platform_integrations } from '@kilocode/db/schema'; -async function findBitbucketIntegrationType(organizationId: string) { +async function findBitbucketIntegration(organizationId: string) { const [integration] = await db .select({ integrationType: platform_integrations.integration_type }) .from(platform_integrations) @@ -27,38 +31,63 @@ async function findBitbucketIntegrationType(organizationId: string) { ) ) .limit(1); - return integration?.integrationType ?? null; + return integration ?? null; } export async function fetchBitbucketRepositoriesForOrganization( organizationId: string, kiloUserId: string, - forceRefresh = false + forceRefresh = false, + options?: RepositoryReadOptions ): Promise { const canonicalOrganizationId = z.uuid().safeParse(organizationId); if (!canonicalOrganizationId.success) return { status: 'invalid_request' }; - const integrationType = await findBitbucketIntegrationType(canonicalOrganizationId.data); - if (integrationType === 'workspace_access_token') { - if (forceRefresh) { - return refreshBitbucketWorkspaceAccessTokenRepositoriesForMember({ - organizationId: canonicalOrganizationId.data, - kiloUserId, - }); - } - return readCachedBitbucketWorkspaceAccessTokenRepositories({ - organizationId: canonicalOrganizationId.data, - }); - } - if (integrationType === 'oauth') { - return listBitbucketRepositories({ - owner: { type: 'org', id: canonicalOrganizationId.data }, - kiloUserId, - forceRefresh, - }); + let integrationFound = false; + try { + const result = await withRepositoryReadDeadline( + options, + async signal => { + const readOptions = options?.bounded ? { bounded: true, signal } : undefined; + const integration = await findBitbucketIntegration(canonicalOrganizationId.data); + signal?.throwIfAborted(); + if (!integration) return { status: 'not_connected' }; + integrationFound = true; + const integrationType = integration.integrationType; + if (integrationType === 'workspace_access_token') { + if (!forceRefresh) { + const cached = await readCachedBitbucketWorkspaceAccessTokenRepositories({ + organizationId: canonicalOrganizationId.data, + readOptions, + }); + if (!readOptions || cached.status !== 'temporarily_unavailable') return cached; + } + signal?.throwIfAborted(); + return refreshBitbucketWorkspaceAccessTokenRepositoriesForMember({ + organizationId: canonicalOrganizationId.data, + kiloUserId, + readOptions, + }); + } + if (integrationType === 'oauth') { + return listBitbucketRepositories({ + owner: { type: 'org', id: canonicalOrganizationId.data }, + kiloUserId, + forceRefresh, + readOptions, + }); + } + if (integrationType || readOptions) return { status: 'reconnect_required' }; + return { status: 'not_connected' }; + } + ); + return options?.bounded && integrationFound && result.status === 'not_connected' + ? { status: 'temporarily_unavailable' } + : result; + } catch (error) { + if (!options?.bounded) throw error; + return { status: 'temporarily_unavailable' }; } - if (integrationType) return { status: 'reconnect_required' }; - return { status: 'not_connected' }; } export { BitbucketOrganizationRepositoryListResultSchema }; diff --git a/apps/web/src/lib/cloud-agent/github-integration-helpers.test.ts b/apps/web/src/lib/cloud-agent/github-integration-helpers.test.ts index ae735ca83f..cc5d91e156 100644 --- a/apps/web/src/lib/cloud-agent/github-integration-helpers.test.ts +++ b/apps/web/src/lib/cloud-agent/github-integration-helpers.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it, jest, beforeEach } from '@jest/globals'; import type { PlatformIntegration } from '@kilocode/db/schema'; import type { Owner } from '@/lib/integrations/core/types'; +import type { RepositoryReadOptions } from '@/lib/integrations/core/repository-read-limits'; // Define mock functions at module level with proper typing const mockGetIntegrationForOrganization = @@ -14,7 +15,9 @@ const mockUpdateRepositoriesForIntegration = const mockGetIntegrationsByOrganization = jest.fn<(organizationId: string, platform: string) => Promise>(); const mockFetchGitHubRepositories = - jest.fn<(installationId: string, appType: string) => Promise>(); + jest.fn< + (installationId: string, appType: string, options?: RepositoryReadOptions) => Promise + >(); const mockGenerateGitHubInstallationToken = jest.fn<(installationId: string, appType: string) => Promise<{ token: string }>>(); const mockCheckExistingFork = @@ -321,3 +324,204 @@ describe('github-integration-helpers', () => { }); }); }); + +describe('bounded GitHub repository reads', () => { + beforeEach(() => { + jest.clearAllMocks(); + mockFetchGitHubRepositories.mockReset(); + mockUpdateRepositoriesForIntegration.mockReset(); + }); + + describe.each(['personal', 'organization'] as const)('%s', scope => { + async function read(forceRefresh = false) { + const helpers = await import('./github-integration-helpers'); + return scope === 'personal' + ? helpers.fetchGitHubRepositoriesForUser('oauth/member', forceRefresh, { bounded: true }) + : helpers.fetchAllGitHubRepositoriesForOrganization('org-123', forceRefresh, { + bounded: true, + }); + } + + function configure(integration: PlatformIntegration | null) { + mockGetIntegrationForOwner.mockResolvedValue(integration); + mockGetIntegrationsByOrganization.mockResolvedValue(integration ? [integration] : []); + } + + it.each([ + ['absent', null, 'not_connected'], + ['empty', buildIntegration({ repositories: [] }), 'available'], + ['suspended', buildIntegration({ suspended_at: '2026-06-25 18:00:00+00' }), 'suspended'], + [ + 'auth-invalid', + buildIntegration({ auth_invalid_at: '2026-06-25 18:00:00+00' }), + 'reconnect_required', + ], + ['misconfigured', buildIntegration({ platform_installation_id: null }), 'misconfigured'], + ] as const)('keeps %s distinct without fetching', async (_label, integration, status) => { + configure(integration); + await expect(read()).resolves.toMatchObject({ + status, + integrationInstalled: integration !== null, + repositories: [], + }); + expect(mockFetchGitHubRepositories).not.toHaveBeenCalled(); + }); + + it('bounds cached entries before validation and projection', async () => { + const repositories = Array.from({ length: 60 }, (_, id) => ({ + id, + name: `repo-${id}`, + full_name: `org/repo-${id}`, + private: false, + })); + Object.defineProperty(repositories[50], 'id', { + get() { + throw new Error('Past the bound'); + }, + }); + configure(buildIntegration({ repositories })); + const result = await read(); + expect(result.status).toBe('available'); + expect(result.repositories).toHaveLength(50); + expect(result.repositories.at(-1)?.fullName).toBe('org/repo-49'); + }); + + it.each([false, true])( + 'does not replace the complete cache on bounded refresh=%s', + async forceRefresh => { + const repositories = Array.from({ length: 60 }, (_, id) => ({ + id, + name: `repo-${id}`, + full_name: `org/repo-${id}`, + private: false, + })); + const integration = buildIntegration({ + repositories: forceRefresh ? repositories : null, + repositories_synced_at: forceRefresh ? '2024-01-01T00:00:00Z' : null, + }); + configure(integration); + mockUpdateRepositoriesForIntegration.mockImplementation(async (_id, value) => { + integration.repositories = value as PlatformIntegration['repositories']; + }); + mockFetchGitHubRepositories.mockImplementation(async (_id, _app, options) => { + if (!options?.bounded || !options.signal) throw new Error('Unbounded transport'); + return repositories.slice(0, 50); + }); + const result = await read(forceRefresh); + expect(result.status).toBe('available'); + expect(result.repositories).toHaveLength(50); + mockFetchGitHubRepositories.mockResolvedValue(repositories); + const helpers = await import('./github-integration-helpers'); + const legacy = + scope === 'personal' + ? await helpers.fetchGitHubRepositoriesForUser('oauth/member') + : await helpers.fetchAllGitHubRepositoriesForOrganization('org-123'); + expect(legacy.repositories).toHaveLength(60); + expect(legacy).not.toHaveProperty('status'); + } + ); + + it('reports provider failure without exposing its raw error', async () => { + configure(buildIntegration({ repositories: null })); + mockFetchGitHubRepositories.mockRejectedValue(new Error('secret provider response')); + await expect(read()).resolves.toEqual({ + status: 'temporarily_unavailable', + integrationInstalled: true, + repositories: [], + syncedAt: null, + }); + }); + }); + + it('rejects more than ten configured installations before network work', async () => { + mockGetIntegrationsByOrganization.mockResolvedValue( + Array.from({ length: 11 }, () => buildIntegration()) + ); + const { fetchAllGitHubRepositoriesForOrganization } = + await import('./github-integration-helpers'); + await expect( + fetchAllGitHubRepositoriesForOrganization('org-123', true, { bounded: true }) + ).resolves.toMatchObject({ status: 'integration_limit_exceeded', repositories: [] }); + expect(mockFetchGitHubRepositories).not.toHaveBeenCalled(); + }); + + it('rejects an unavailable sibling instead of hiding it behind healthy repositories', async () => { + mockGetIntegrationsByOrganization.mockResolvedValue([ + buildIntegration(), + buildIntegration({ auth_invalid_at: '2026-06-25 18:00:00+00' }), + ]); + const { fetchAllGitHubRepositoriesForOrganization } = + await import('./github-integration-helpers'); + await expect( + fetchAllGitHubRepositoriesForOrganization('org-123', false, { bounded: true }) + ).resolves.toMatchObject({ status: 'reconnect_required', repositories: [] }); + }); + + it('rejects a later provider failure even after collecting fifty repositories', async () => { + mockGetIntegrationsByOrganization.mockResolvedValue([ + buildIntegration({ repositories: Array.from({ length: 50 }, () => cachedRepositories[0]) }), + buildIntegration({ id: 'integration-2', repositories: null }), + ]); + mockFetchGitHubRepositories.mockRejectedValue(new Error('unavailable')); + const { fetchAllGitHubRepositoriesForOrganization } = + await import('./github-integration-helpers'); + await expect( + fetchAllGitHubRepositoriesForOrganization('org-123', false, { bounded: true }) + ).resolves.toMatchObject({ status: 'temporarily_unavailable', repositories: [] }); + }); + + it('fetches sequentially and caps the aggregate while preserving installation provenance', async () => { + mockGetIntegrationsByOrganization.mockResolvedValue([ + buildIntegration({ repositories: null }), + buildIntegration({ + id: 'integration-2', + platform_installation_id: 'installation-2', + repositories: null, + }), + ]); + let active = 0; + let peak = 0; + mockFetchGitHubRepositories.mockImplementation(async installationId => { + peak = Math.max(peak, ++active); + await Promise.resolve(); + active--; + return Array.from({ length: 40 }, (_, id) => ({ + id, + name: 'repo', + full_name: `${installationId}/repo-${id}`, + private: false, + })); + }); + const { fetchAllGitHubRepositoriesForOrganization } = + await import('./github-integration-helpers'); + const result = await fetchAllGitHubRepositoriesForOrganization('org-123', true, { + bounded: true, + }); + expect(result.status).toBe('available'); + expect(peak).toBe(1); + expect(result.repositories).toHaveLength(50); + expect(result.repositories.at(-1)).toMatchObject({ + fullName: 'installation-2/repo-9', + platformIntegrationId: 'integration-2', + }); + }); + + it('uses one deadline and starts no later installation after timeout', async () => { + jest.useFakeTimers(); + try { + mockGetIntegrationsByOrganization.mockResolvedValue([buildIntegration(), buildIntegration()]); + mockFetchGitHubRepositories.mockImplementation(() => new Promise(() => {})); + const { fetchAllGitHubRepositoriesForOrganization } = + await import('./github-integration-helpers'); + const result = fetchAllGitHubRepositoriesForOrganization('org-123', true, { bounded: true }); + await jest.advanceTimersByTimeAsync(30_000); + await expect(result).resolves.toMatchObject({ + status: 'temporarily_unavailable', + repositories: [], + }); + expect(mockFetchGitHubRepositories).toHaveBeenCalledTimes(1); + } finally { + jest.useRealTimers(); + } + }); +}); diff --git a/apps/web/src/lib/cloud-agent/github-integration-helpers.ts b/apps/web/src/lib/cloud-agent/github-integration-helpers.ts index 91329a3f82..61cc0bdbd3 100644 --- a/apps/web/src/lib/cloud-agent/github-integration-helpers.ts +++ b/apps/web/src/lib/cloud-agent/github-integration-helpers.ts @@ -1,4 +1,10 @@ import { TRPCError } from '@trpc/server'; +import type { Owner } from '@/lib/integrations/core/types'; +import { + REPOSITORY_READ_LIMITS, + withRepositoryReadDeadline, + type RepositoryReadOptions, +} from '@/lib/integrations/core/repository-read-limits'; import { getIntegrationsByOrganization, getIntegrationForOrganization, @@ -23,6 +29,14 @@ import { } from '@/lib/integrations/core/types'; type GitHubRepositoriesResult = { + status?: + | 'not_connected' + | 'available' + | 'suspended' + | 'reconnect_required' + | 'misconfigured' + | 'temporarily_unavailable' + | 'integration_limit_exceeded'; integrationInstalled: boolean; repositories: { id: number; @@ -187,10 +201,89 @@ export async function fetchGitHubRepositoriesForOrganization( } } +async function fetchBoundedGitHubRepositories( + owner: Owner, + forceRefresh: boolean, + options: RepositoryReadOptions +): Promise { + const unavailable = { integrationInstalled: true, repositories: [], syncedAt: null }; + try { + return await withRepositoryReadDeadline(options, async signal => { + const candidates = + owner.type === 'org' + ? await getIntegrationsByOrganization(owner.id, PLATFORM.GITHUB) + : [await getIntegrationForOwner(owner, PLATFORM.GITHUB)]; + const integrations = candidates.filter(integration => integration !== null); + if (!integrations.length) { + return { ...unavailable, integrationInstalled: false, status: 'not_connected' }; + } + if (integrations.length > 10) { + return { ...unavailable, status: 'integration_limit_exceeded' }; + } + // Check every configured installation before fetching or applying the result cap. + for (const integration of integrations) { + if (isPlatformIntegrationSuspended(integration)) { + return { ...unavailable, status: 'suspended' }; + } + if (integration.auth_invalid_at) { + return { ...unavailable, status: 'reconnect_required' }; + } + if (!isPlatformIntegrationHealthy(integration) || !integration.platform_installation_id) { + return { ...unavailable, status: 'misconfigured' }; + } + } + const repositories: GitHubRepositoriesResult['repositories'] = []; + let syncedAt: string | null = null; + for (const integration of integrations) { + signal?.throwIfAborted(); + const installationId = integration.platform_installation_id; + if (!installationId) return { ...unavailable, status: 'misconfigured' }; + const cached = requireNumericPlatformRepositories( + integration.repositories?.slice(0, REPOSITORY_READ_LIMITS.repositories) ?? null + ); + const refresh = forceRefresh || !cached || !integration.repositories_synced_at; + const selected = refresh + ? await fetchGitHubRepositories( + installationId, + integration.github_app_type || 'standard', + { + bounded: true, + signal, + } + ) + : cached; + signal?.throwIfAborted(); + repositories.push( + ...mapRepositories( + selected.slice(0, REPOSITORY_READ_LIMITS.repositories - repositories.length), + owner.type === 'org' ? integration : undefined + ) + ); + const nextSyncedAt = refresh + ? new Date().toISOString() + : integration.repositories_synced_at; + if (nextSyncedAt && (!syncedAt || nextSyncedAt < syncedAt)) syncedAt = nextSyncedAt; + } + // Bounded transport results are not a complete shared cache. + return { status: 'available', integrationInstalled: true, repositories, syncedAt }; + }); + } catch { + return { ...unavailable, status: 'temporarily_unavailable' }; + } +} + export async function fetchAllGitHubRepositoriesForOrganization( organizationId: string, - forceRefresh: boolean = false + forceRefresh: boolean = false, + options?: RepositoryReadOptions ): Promise { + if (options?.bounded) { + return fetchBoundedGitHubRepositories( + { type: 'org', id: organizationId }, + forceRefresh, + options + ); + } const integrations = ( await getIntegrationsByOrganization(organizationId, PLATFORM.GITHUB) ).filter(isPlatformIntegrationHealthy); @@ -252,8 +345,12 @@ async function fetchRepositoriesForIntegrations( export async function fetchGitHubRepositoriesForUser( userId: string, - forceRefresh: boolean = false + forceRefresh: boolean = false, + options?: RepositoryReadOptions ): Promise { + if (options?.bounded) { + return fetchBoundedGitHubRepositories({ type: 'user', id: userId }, forceRefresh, options); + } const integration = await getIntegrationForOwner({ type: 'user', id: userId }, PLATFORM.GITHUB); if (!integration) { diff --git a/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.test.ts b/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.test.ts index cec31a2f19..8c96b0a363 100644 --- a/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.test.ts +++ b/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.test.ts @@ -1,7 +1,10 @@ -import { describe, expect, it, jest, beforeEach } from '@jest/globals'; +import { describe, expect, it, jest, beforeAll, beforeEach } from '@jest/globals'; import type { PlatformIntegration } from '@kilocode/db/schema'; import type { Owner } from '@/lib/integrations/core/types'; -import { buildGitLabCloneUrl } from './gitlab-integration-helpers'; +import type { RepositoryReadOptions } from '@/lib/integrations/core/repository-read-limits'; +import type { buildGitLabCloneUrl as BuildGitLabCloneUrl } from './gitlab-integration-helpers'; + +let buildGitLabCloneUrl: typeof BuildGitLabCloneUrl; // Define mock functions at module level with proper typing const mockGetGitLabIntegration = jest.fn<(owner: Owner) => Promise>(); @@ -19,7 +22,13 @@ const mockGetIntegrationForOwner = const mockUpdateRepositoriesForIntegration = jest.fn<(integrationId: string, repositories: unknown[]) => Promise>(); const mockFetchGitLabProjects = - jest.fn<(accessToken: string, instanceUrl: string) => Promise>(); + jest.fn< + ( + accessToken: string, + instanceUrl: string, + options?: RepositoryReadOptions + ) => Promise + >(); // Wire up the mocks jest.mock('@/lib/integrations/gitlab-service', () => ({ @@ -37,6 +46,14 @@ jest.mock('@/lib/integrations/platforms/gitlab/adapter', () => ({ fetchGitLabProjects: mockFetchGitLabProjects, })); +jest.mock('dns/promises', () => ({ + lookup: jest.fn(), +})); + +beforeAll(async () => { + ({ buildGitLabCloneUrl } = await import('./gitlab-integration-helpers')); +}); + describe('gitlab-integration-helpers', () => { beforeEach(() => { jest.clearAllMocks(); @@ -583,3 +600,241 @@ describe('gitlab-integration-helpers', () => { }); }); }); + +describe.each(['personal', 'organization'] as const)('bounded GitLab %s reads', scope => { + const repositories = Array.from({ length: 60 }, (_, id) => ({ + id, + name: `repo-${id}`, + full_name: `group/repo-${id}`, + private: false, + })); + const integration = (overrides: Partial = {}) => + ({ + id: 'integration-1', + platform: 'gitlab', + integration_status: 'active', + suspended_at: null, + auth_invalid_at: null, + metadata: {}, + repositories, + repositories_synced_at: '2026-06-25 18:00:00+00', + ...overrides, + }) as PlatformIntegration; + + beforeEach(() => { + jest.clearAllMocks(); + mockGetValidGitLabToken.mockReset().mockResolvedValue('valid-token'); + mockFetchGitLabProjects.mockReset(); + mockUpdateRepositoriesForIntegration.mockReset(); + }); + + function configure(value: PlatformIntegration | null) { + mockGetIntegrationForOwner.mockResolvedValue(value); + mockGetIntegrationForOrganization.mockResolvedValue(value); + } + + async function read( + forceRefresh = false, + options: RepositoryReadOptions | undefined = { bounded: true } + ) { + const helpers = await import('./gitlab-integration-helpers'); + return scope === 'personal' + ? helpers.fetchGitLabRepositoriesForUser('oauth/member', forceRefresh, options) + : helpers.fetchGitLabRepositoriesForOrganization( + 'org-123', + 'oauth/member', + forceRefresh, + options + ); + } + + it.each([ + ['absent', null, 'not_connected'], + ['empty', integration({ repositories: [] }), 'available'], + ['suspended', integration({ suspended_at: '2026-06-25 18:00:00+00' }), 'suspended'], + [ + 'auth-invalid', + integration({ auth_invalid_at: '2026-06-25 18:00:00+00' }), + 'reconnect_required', + ], + ['inactive', integration({ integration_status: 'pending' }), 'misconfigured'], + [ + 'malformed URL', + integration({ metadata: { gitlab_instance_url: 'not a URL' } }), + 'misconfigured', + ], + [ + 'private URL', + integration({ metadata: { gitlab_instance_url: 'https://127.0.0.1' } }), + 'misconfigured', + ], + [ + 'unsafe hostname', + integration({ metadata: { gitlab_instance_url: 'https://gitlab.local' } }), + 'misconfigured', + ], + [ + 'invalid URL', + integration({ metadata: { gitlab_instance_url: 'https://user:password@gitlab.com' } }), + 'misconfigured', + ], + ] as const)('keeps %s distinct without fetching', async (_label, value, status) => { + configure(value); + await expect(read()).resolves.toMatchObject({ + status, + integrationInstalled: value !== null, + repositories: [], + }); + expect(mockGetValidGitLabToken).not.toHaveBeenCalled(); + expect(mockFetchGitLabProjects).not.toHaveBeenCalled(); + }); + + it('caps a cached list before checking repository IDs', async () => { + const cached = [...repositories]; + cached[50] = { + ...cached[50], + get id(): number { + throw new Error('Past the bound'); + }, + }; + configure(integration({ repositories: cached })); + const result = await read(); + expect(result.status).toBe('available'); + expect(result.repositories).toHaveLength(50); + expect(result.repositories.at(-1)?.fullName).toBe('group/repo-49'); + }); + + it.each([false, true])( + 'preserves the complete cache and legacy behavior on refresh=%s', + async forceRefresh => { + const value = integration({ + repositories: forceRefresh ? repositories : null, + repositories_synced_at: forceRefresh ? '2026-06-25 18:00:00+00' : null, + }); + configure(value); + mockUpdateRepositoriesForIntegration.mockImplementation(async (_id, saved) => { + value.repositories = saved as PlatformIntegration['repositories']; + }); + mockFetchGitLabProjects.mockImplementation(async (_token, _url, options) => { + if (!options?.bounded || !options.signal) throw new Error('Unbounded transport'); + return repositories.slice(0, 50); + }); + const result = await read(forceRefresh); + expect(result.status).toBe('available'); + expect(result.repositories).toHaveLength(50); + mockFetchGitLabProjects.mockResolvedValue(repositories); + const legacy = await read(false, { bounded: false }); + expect(legacy.repositories).toHaveLength(60); + expect(legacy).not.toHaveProperty('status'); + } + ); + + it.each([ + ['UNAUTHORIZED', 'reconnect_required'], + ['INTERNAL_SERVER_ERROR', 'temporarily_unavailable'], + ] as const)( + 'keeps credential rejection separate from temporary failure', + async (code, status) => { + const { TRPCError } = await import('@trpc/server'); + configure(integration({ repositories: null })); + mockGetValidGitLabToken.mockRejectedValue( + new TRPCError({ code, message: 'provider details' }) + ); + await expect(read()).resolves.toEqual({ + status, + integrationInstalled: true, + repositories: [], + syncedAt: null, + }); + } + ); + + it('does not expose a raw project-fetch failure', async () => { + configure(integration({ repositories: null })); + mockFetchGitLabProjects.mockRejectedValue(new Error('secret provider response')); + await expect(read()).resolves.toEqual({ + status: 'temporarily_unavailable', + integrationInstalled: true, + repositories: [], + syncedAt: null, + }); + }); + + describe('self-hosted DNS failures', () => { + beforeEach(async () => { + const { lookup } = await import('dns/promises'); + const { buildGitLabUrl, resolveGitLabUrlSafely } = + await import('@/lib/integrations/platforms/gitlab/instance-url'); + jest + .mocked(lookup) + .mockReset() + .mockRejectedValue( + Object.assign(new Error('getaddrinfo EAI_AGAIN gitlab.example.com'), { + code: 'EAI_AGAIN', + }) + ); + configure( + integration({ + metadata: { gitlab_instance_url: 'https://gitlab.example.com' }, + repositories: null, + }) + ); + mockFetchGitLabProjects.mockImplementation(async (_token, instanceUrl) => { + await resolveGitLabUrlSafely(buildGitLabUrl(instanceUrl, '/api/v4/projects')); + return []; + }); + }); + + it('reports transient lookup failure without exposing resolver details', async () => { + await expect(read()).resolves.toEqual({ + status: 'temporarily_unavailable', + integrationInstalled: true, + repositories: [], + syncedAt: null, + }); + }); + + it('keeps private DNS answers misconfigured', async () => { + const { lookup } = await import('dns/promises'); + jest.mocked(lookup).mockResolvedValueOnce([{ address: '192.168.1.10', family: 4 }]); + await expect(read()).resolves.toEqual({ + status: 'misconfigured', + integrationInstalled: true, + repositories: [], + syncedAt: null, + }); + }); + + it('preserves the legacy failure when bounded options are omitted', async () => { + const helpers = await import('./gitlab-integration-helpers'); + const result = + scope === 'personal' + ? helpers.fetchGitLabRepositoriesForUser('oauth/member') + : helpers.fetchGitLabRepositoriesForOrganization('org-123', 'oauth/member'); + await expect(result).rejects.toMatchObject({ + code: 'INTERNAL_SERVER_ERROR', + message: 'Failed to fetch GitLab repositories', + }); + }); + }); + + it('stops before project fetching when credential retrieval exceeds the deadline', async () => { + jest.useFakeTimers(); + try { + configure(integration({ repositories: null })); + const token = Promise.withResolvers(); + mockGetValidGitLabToken.mockReturnValue(token.promise); + const result = read(); + await jest.advanceTimersByTimeAsync(30_000); + await expect(result).resolves.toMatchObject({ + status: 'temporarily_unavailable', + repositories: [], + }); + token.resolve('late-token'); + await Promise.resolve(); + expect(mockFetchGitLabProjects).not.toHaveBeenCalled(); + } finally { + jest.useRealTimers(); + } + }); +}); diff --git a/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.ts b/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.ts index 510c91b50d..f923cfeb1d 100644 --- a/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.ts +++ b/apps/web/src/lib/cloud-agent/gitlab-integration-helpers.ts @@ -1,4 +1,15 @@ import { TRPCError } from '@trpc/server'; +import { z } from 'zod'; +import type { Owner } from '@/lib/integrations/core/types'; +import { + REPOSITORY_READ_LIMITS, + withRepositoryReadDeadline, + type RepositoryReadOptions, +} from '@/lib/integrations/core/repository-read-limits'; +import { + GitLabInstanceUrlError, + normalizeGitLabInstanceUrl, +} from '@/lib/integrations/platforms/gitlab/instance-url'; import { getIntegrationForOrganization, getIntegrationForOwner, @@ -16,6 +27,13 @@ import { const DEFAULT_GITLAB_URL = 'https://gitlab.com'; type GitLabRepositoriesResult = { + status?: + | 'not_connected' + | 'available' + | 'suspended' + | 'reconnect_required' + | 'misconfigured' + | 'temporarily_unavailable'; integrationInstalled: boolean; repositories: { id: number; @@ -105,6 +123,70 @@ export async function getGitLabTokenForUser(userId: string): Promise { + const unavailable = { integrationInstalled: true, repositories: [], syncedAt: null }; + try { + return await withRepositoryReadDeadline(options, async signal => { + const integration = + owner.type === 'org' + ? await getIntegrationForOrganization(owner.id, PLATFORM.GITLAB) + : await getIntegrationForOwner(owner, PLATFORM.GITLAB); + if (!integration) { + return { ...unavailable, integrationInstalled: false, status: 'not_connected' }; + } + if (isPlatformIntegrationSuspended(integration)) + return { ...unavailable, status: 'suspended' }; + if (integration.auth_invalid_at) return { ...unavailable, status: 'reconnect_required' }; + const metadata = z + .object({ gitlab_instance_url: z.string().optional() }) + .safeParse(integration.metadata ?? {}); + if (integration.integration_status !== 'active' || !metadata.success) { + return { ...unavailable, status: 'misconfigured' }; + } + const instanceUrl = normalizeGitLabInstanceUrl(metadata.data.gitlab_instance_url); + const cached = requireNumericPlatformRepositories( + integration.repositories?.slice(0, REPOSITORY_READ_LIMITS.repositories) ?? null + ); + let repositories = cached; + let syncedAt = integration.repositories_synced_at; + if (forceRefresh || !cached || !syncedAt) { + signal?.throwIfAborted(); + const token = await getValidGitLabToken(integration, { + userId: actorUserId, + ...(owner.type === 'org' ? { organizationId: owner.id } : {}), + }); + signal?.throwIfAborted(); + repositories = await fetchGitLabProjects(token, instanceUrl, { bounded: true, signal }); + syncedAt = new Date().toISOString(); + } + signal?.throwIfAborted(); + return { + status: 'available', + integrationInstalled: true, + repositories: mapRepositories( + (repositories ?? []).slice(0, REPOSITORY_READ_LIMITS.repositories) + ), + syncedAt, + instanceUrl, + }; + }); + } catch (error) { + const status = + error instanceof GitLabInstanceUrlError && error.reason !== 'resolution_failed' + ? 'misconfigured' + : error instanceof TRPCError && + (error.code === 'UNAUTHORIZED' || error.code === 'FORBIDDEN') + ? 'reconnect_required' + : 'temporarily_unavailable'; + return { ...unavailable, status }; + } +} + /** * Fetch GitLab repositories for an organization * Returns cached repositories by default, fetches fresh from GitLab when forceRefresh is true @@ -112,8 +194,17 @@ export async function getGitLabTokenForUser(userId: string): Promise { + if (options?.bounded) { + return fetchBoundedGitLabRepositories( + { type: 'org', id: organizationId }, + actorUserId, + forceRefresh, + options + ); + } const integration = await getIntegrationForOrganization(organizationId, PLATFORM.GITLAB); if (!integration) { @@ -166,8 +257,17 @@ export async function fetchGitLabRepositoriesForOrganization( */ export async function fetchGitLabRepositoriesForUser( userId: string, - forceRefresh: boolean = false + forceRefresh: boolean = false, + options?: RepositoryReadOptions ): Promise { + if (options?.bounded) { + return fetchBoundedGitLabRepositories( + { type: 'user', id: userId }, + userId, + forceRefresh, + options + ); + } const integration = await getIntegrationForOwner({ type: 'user', id: userId }, PLATFORM.GITLAB); if (!integration) { diff --git a/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.test.ts b/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.test.ts index d680d7d7d7..ac7203fe6b 100644 --- a/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.test.ts +++ b/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.test.ts @@ -10,6 +10,8 @@ import { jest, } from '@jest/globals'; import type { Organization, User } from '@kilocode/db/schema'; +import { randomUUID } from 'node:crypto'; +import type { RepositoryReadOptions } from '@/lib/integrations/core/repository-read-limits'; import { kilocode_users, organizations, @@ -28,7 +30,11 @@ import type { const mockFetchBitbucketRepositoriesFromTokenService = jest.fn< - (kiloUserId: string, organizationId?: string) => Promise + ( + kiloUserId: string, + organizationId?: string, + options?: RepositoryReadOptions + ) => Promise >(); jest.mock('./token-service-client', () => ({ @@ -403,4 +409,171 @@ describe('Bitbucket repository cache', () => { expect(unchanged?.repositories).toEqual([CACHED_REPOSITORY]); expect(new Date(unchanged?.syncedAt ?? '').toISOString()).toBe(CACHED_AT); }); + + it('distinguishes bounded OAuth absence from an available empty cache', async () => { + const input = { + owner: { type: 'user' as const, id: user.id }, + kiloUserId: user.id, + readOptions: { bounded: true }, + }; + await expect(listBitbucketRepositories(input)).resolves.toEqual({ status: 'not_connected' }); + await insertActiveIntegration(user.id, { repositories: [] }); + await expect(listBitbucketRepositories(input)).resolves.toEqual({ + status: 'available', + repositories: [], + syncedAt: CACHED_AT, + }); + }); + + it.each(['reconnect_required', 'insufficient_permissions', 'temporarily_unavailable'] as const)( + 'retains bounded OAuth %s without changing the complete cache', + async status => { + const integration = await insertActiveIntegration(user.id); + mockFetchBitbucketRepositoriesFromTokenService.mockResolvedValue({ status }); + await expect( + listBitbucketRepositories({ + owner: { type: 'user', id: user.id }, + kiloUserId: user.id, + forceRefresh: true, + readOptions: { bounded: true }, + }) + ).resolves.toEqual({ status }); + const [saved] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, integration.id)); + expect(saved?.repositories).toEqual([CACHED_REPOSITORY]); + } + ); + + it('does not report a configured OAuth provider as absent after a remote disconnect', async () => { + const integration = await insertActiveIntegration(user.id, {}, organization.id); + mockFetchBitbucketRepositoriesFromTokenService.mockResolvedValue({ status: 'not_connected' }); + const { fetchBitbucketRepositoriesForOrganization } = + await import('@/lib/cloud-agent/bitbucket-integration-helpers'); + await expect( + fetchBitbucketRepositoriesForOrganization(organization.id, user.id, true, { bounded: true }) + ).resolves.toEqual({ status: 'temporarily_unavailable' }); + const [saved] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, integration.id)); + expect(saved?.repositories).toEqual([CACHED_REPOSITORY]); + }); + + it.each(['suspended', 'auth-invalid'] as const)( + 'rejects a bounded OAuth cache hit when the integration is %s', + async state => { + const integration = await insertActiveIntegration(user.id); + const invalidation = + state === 'suspended' + ? { suspended_at: CACHED_AT } + : { auth_invalid_at: CACHED_AT, auth_invalid_reason: 'provider_rejected' }; + await db + .update(platform_integrations) + .set(invalidation) + .where(eq(platform_integrations.id, integration.id)); + await expect( + listBitbucketRepositories({ + owner: { type: 'user', id: user.id }, + kiloUserId: user.id, + readOptions: { bounded: true }, + }) + ).resolves.toEqual({ status: 'reconnect_required' }); + expect(mockFetchBitbucketRepositoriesFromTokenService).not.toHaveBeenCalled(); + } + ); + + it('bounds a complete OAuth cache without changing legacy reads', async () => { + const repositories = Array.from({ length: 60 }, (_, index) => ({ + ...CACHED_REPOSITORY, + id: randomUUID(), + full_name: `acme/repo-${index}`, + })); + await insertActiveIntegration(user.id, { repositories }); + const input = { owner: { type: 'user' as const, id: user.id }, kiloUserId: user.id }; + const bounded = await listBitbucketRepositories({ ...input, readOptions: { bounded: true } }); + expect(bounded.status).toBe('available'); + if (bounded.status !== 'available') throw new Error('Expected bounded repositories'); + expect(bounded.repositories).toHaveLength(50); + expect(bounded.repositories.at(-1)?.fullName).toBe('acme/repo-49'); + const legacy = await listBitbucketRepositories(input); + expect(legacy.status).toBe('available'); + if (legacy.status !== 'available') throw new Error('Expected complete repositories'); + expect(legacy.repositories).toHaveLength(60); + }); + + it.each([false, true])( + 'does not persist a partial OAuth result on refresh=%s', + async forceRefresh => { + const integration = await insertActiveIntegration( + user.id, + forceRefresh ? {} : { repositories: null, syncedAt: null } + ); + mockFetchBitbucketRepositoriesFromTokenService.mockImplementation( + async (_userId, _organizationId, options) => { + if (!options?.bounded) throw new Error('Unbounded transport'); + return { status: 'available', repositories: [REFRESHED_REPOSITORY] }; + } + ); + await expect( + listBitbucketRepositories({ + owner: { type: 'user', id: user.id }, + kiloUserId: user.id, + forceRefresh, + readOptions: { bounded: true }, + }) + ).resolves.toMatchObject({ status: 'available', repositories: [REFRESHED_REPOSITORY] }); + const [unchanged] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, integration.id)); + expect(unchanged?.repositories).toEqual(forceRefresh ? [CACHED_REPOSITORY] : null); + expect( + unchanged?.repositories_synced_at + ? new Date(unchanged.repositories_synced_at).toISOString() + : null + ).toBe(forceRefresh ? CACHED_AT : null); + } + ); + + it.each(['available', 'reconnect_required', 'temporarily_unavailable'] as const)( + 'ignores stale OAuth %s after credential rotation and returns the bounded winner', + async status => { + const integration = await insertActiveIntegration(user.id); + const winner = Array.from({ length: 60 }, (_, index) => ({ + ...CACHED_REPOSITORY, + id: randomUUID(), + full_name: `acme/winner-${index}`, + })); + mockFetchBitbucketRepositoriesFromTokenService.mockImplementation(async () => { + await db + .update(platform_oauth_credentials) + .set({ credential_version: 2 }) + .where(eq(platform_oauth_credentials.platform_integration_id, integration.id)); + await db + .update(platform_integrations) + .set({ repositories: winner, repositories_synced_at: '2026-06-24T08:00:00.000Z' }) + .where(eq(platform_integrations.id, integration.id)); + return status === 'available' + ? { status, repositories: [REFRESHED_REPOSITORY] } + : { status }; + }); + const result = await listBitbucketRepositories({ + owner: { type: 'user', id: user.id }, + kiloUserId: user.id, + forceRefresh: true, + readOptions: { bounded: true }, + }); + expect(result.status).toBe('available'); + if (result.status !== 'available') throw new Error('Expected the winner cache'); + expect(result.repositories).toHaveLength(50); + expect(result.repositories[0].fullName).toBe('acme/winner-0'); + const [saved] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, integration.id)); + expect(saved?.repositories).toEqual(winner); + } + ); }); diff --git a/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.ts b/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.ts index ef76c3ffc8..95f94934bb 100644 --- a/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.ts +++ b/apps/web/src/lib/integrations/platforms/bitbucket/repository-cache.ts @@ -9,6 +9,10 @@ import { after } from 'next/server'; import { db } from '@/lib/drizzle'; import { INTEGRATION_STATUS, PLATFORM } from '@/lib/integrations/core/constants'; import type { Owner } from '@/lib/integrations/core/types'; +import { + REPOSITORY_READ_LIMITS, + type RepositoryReadOptions, +} from '@/lib/integrations/core/repository-read-limits'; import { BitbucketIntegrationMetadataSchema, type BitbucketWorkspace } from './metadata'; import { BitbucketRepositorySchema, @@ -50,6 +54,7 @@ type ListBitbucketRepositoriesInput = { kiloUserId: string; forceRefresh?: boolean; expectedIntegrationId?: string; + readOptions?: RepositoryReadOptions; }; type PrimeBitbucketRepositoryCacheInput = { @@ -70,10 +75,15 @@ function ownerCondition(owner: Owner) { function readCachedRepositories( value: unknown, syncedAt: string | null, - workspace: BitbucketWorkspace + workspace: BitbucketWorkspace, + readOptions?: RepositoryReadOptions ): Extract | null { if (value === null || !syncedAt) return null; - const repositories = z.array(CachedBitbucketRepositorySchema).safeParse(value); + const selected = + readOptions?.bounded && Array.isArray(value) + ? value.slice(0, REPOSITORY_READ_LIMITS.repositories) + : value; + const repositories = z.array(CachedBitbucketRepositorySchema).safeParse(selected); if (!repositories.success) return null; return { status: 'available', @@ -94,26 +104,32 @@ export async function listBitbucketRepositories({ kiloUserId, forceRefresh = false, expectedIntegrationId, + readOptions, }: ListBitbucketRepositoriesInput): Promise { - const [row] = await db - .select({ - integrationId: platform_integrations.id, - integrationStatus: platform_integrations.integration_status, - installationId: platform_integrations.platform_installation_id, - accountId: platform_integrations.platform_account_id, - accountLogin: platform_integrations.platform_account_login, - metadata: platform_integrations.metadata, - repositories: platform_integrations.repositories, - repositoriesSyncedAt: platform_integrations.repositories_synced_at, - credential: platform_oauth_credentials, - }) - .from(platform_integrations) - .leftJoin( - platform_oauth_credentials, - eq(platform_oauth_credentials.platform_integration_id, platform_integrations.id) - ) - .where(and(ownerCondition(owner), eq(platform_integrations.platform, PLATFORM.BITBUCKET))) - .limit(1); + readOptions?.signal?.throwIfAborted(); + const load = () => + db + .select({ + integrationId: platform_integrations.id, + integrationStatus: platform_integrations.integration_status, + suspendedAt: platform_integrations.suspended_at, + authInvalidAt: platform_integrations.auth_invalid_at, + installationId: platform_integrations.platform_installation_id, + accountId: platform_integrations.platform_account_id, + accountLogin: platform_integrations.platform_account_login, + metadata: platform_integrations.metadata, + repositories: platform_integrations.repositories, + repositoriesSyncedAt: platform_integrations.repositories_synced_at, + credential: platform_oauth_credentials, + }) + .from(platform_integrations) + .leftJoin( + platform_oauth_credentials, + eq(platform_oauth_credentials.platform_integration_id, platform_integrations.id) + ) + .where(and(ownerCondition(owner), eq(platform_integrations.platform, PLATFORM.BITBUCKET))) + .limit(1); + const [row] = await load(); if (!row) return { status: 'not_connected' }; if (expectedIntegrationId && row.integrationId !== expectedIntegrationId) { @@ -123,6 +139,9 @@ export async function listBitbucketRepositories({ if (!credential.success || credential.data.revoked_at) { return { status: 'reconnect_required' }; } + if (readOptions?.bounded && (row.suspendedAt || row.authInvalidAt)) { + return { status: 'reconnect_required' }; + } const metadata = BitbucketIntegrationMetadataSchema.safeParse(row.metadata); if (!metadata.success) return { status: 'reconnect_required' }; @@ -145,14 +164,51 @@ export async function listBitbucketRepositories({ const cachedResult = readCachedRepositories( row.repositories, row.repositoriesSyncedAt, - metadata.data.workspace + metadata.data.workspace, + readOptions ); if (!forceRefresh && cachedResult) return cachedResult; - const result = await fetchBitbucketRepositoriesFromTokenService( - kiloUserId, - owner.type === 'org' ? owner.id : undefined - ); + readOptions?.signal?.throwIfAborted(); + const result = readOptions?.bounded + ? await fetchBitbucketRepositoriesFromTokenService( + kiloUserId, + owner.type === 'org' ? owner.id : undefined, + readOptions + ).catch(() => ({ status: 'temporarily_unavailable' as const })) + : await fetchBitbucketRepositoriesFromTokenService( + kiloUserId, + owner.type === 'org' ? owner.id : undefined + ); + if (readOptions?.bounded) { + readOptions.signal?.throwIfAborted(); + const [current] = await load(); + const currentCredential = BitbucketOAuthCredentialRowSchema.safeParse(current?.credential); + if ( + !current || + !currentCredential.success || + current.integrationId !== row.integrationId || + current.integrationStatus !== row.integrationStatus || + current.suspendedAt !== row.suspendedAt || + current.authInvalidAt !== row.authInvalidAt || + current.installationId !== row.installationId || + current.accountId !== row.accountId || + current.accountLogin !== row.accountLogin || + JSON.stringify(current.metadata) !== JSON.stringify(row.metadata) || + current.repositoriesSyncedAt !== row.repositoriesSyncedAt || + currentCredential.data.id !== credential.data.id || + currentCredential.data.credential_version !== credential.data.credential_version || + currentCredential.data.revoked_at !== credential.data.revoked_at + ) { + return listBitbucketRepositories({ owner, kiloUserId, expectedIntegrationId, readOptions }); + } + if (result.status !== 'available') return result; + return { + status: 'available', + repositories: result.repositories.slice(0, REPOSITORY_READ_LIMITS.repositories), + syncedAt: new Date().toISOString(), + }; + } if (result.status !== 'available') return result; const repositories = result.repositories.map(repository => ({ diff --git a/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.test.ts b/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.test.ts index b57eb3aa55..41a88aef34 100644 --- a/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.test.ts +++ b/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.test.ts @@ -10,6 +10,8 @@ import { type User, } from '@kilocode/db/schema'; import { db } from '@/lib/drizzle'; +import { randomUUID } from 'node:crypto'; +import type { RepositoryReadOptions } from '@/lib/integrations/core/repository-read-limits'; import { eq } from 'drizzle-orm'; import { createTestOrganization } from '@/tests/helpers/organization.helper'; import { insertTestUser } from '@/tests/helpers/user.helper'; @@ -23,7 +25,13 @@ import type { } from './workspace-access-token-repository-cache'; const mockFetchBitbucketRepositoriesFromTokenService = - jest.fn<(kiloUserId: string, organizationId: string) => Promise>(); + jest.fn< + ( + kiloUserId: string, + organizationId: string, + options?: RepositoryReadOptions + ) => Promise + >(); jest.mock('./token-service-client', () => ({ BitbucketRepositorySchema: @@ -864,4 +872,174 @@ describe('Bitbucket Workspace Access Token repository cache', () => { expect(new Date(invalidated?.invalidAt ?? '').toISOString()).toBe('2026-06-24T09:00:00.000Z'); expect(invalidated?.invalidReason).toBe('provider_rejected'); }); + + it('distinguishes bounded absence from an available empty workspace-token cache', async () => { + await expect( + fetchBitbucketRepositoriesForOrganization(organization.id, user.id, false, { bounded: true }) + ).resolves.toEqual({ status: 'not_connected' }); + await insertStaticIntegration(organization.id, user.id, { repositories: [] }); + await expect( + fetchBitbucketRepositoriesForOrganization(organization.id, user.id, false, { bounded: true }) + ).resolves.toEqual({ status: 'available', repositories: [], syncedAt: CACHED_AT }); + }); + + it.each(['suspended', 'auth-invalid'] as const)( + 'rejects a bounded workspace-token cache hit when the integration is %s', + async state => { + const integration = await insertStaticIntegration(organization.id, user.id); + const invalidation = + state === 'suspended' + ? { suspended_at: CACHED_AT } + : { auth_invalid_at: CACHED_AT, auth_invalid_reason: 'provider_rejected' }; + await db + .update(platform_integrations) + .set(invalidation) + .where(eq(platform_integrations.id, integration.id)); + await expect( + fetchBitbucketRepositoriesForOrganization(organization.id, user.id, false, { + bounded: true, + }) + ).resolves.toEqual({ status: 'reconnect_required' }); + expect(mockFetchBitbucketRepositoriesFromTokenService).not.toHaveBeenCalled(); + } + ); + + it('bounds workspace-token cache projection while preserving complete legacy reads', async () => { + const repositories = Array.from({ length: 60 }, (_, index) => ({ + ...CACHED_REPOSITORY, + id: randomUUID(), + full_name: `acme/repo-${index}`, + })); + await insertStaticIntegration(organization.id, user.id, { repositories }); + const bounded = await fetchBitbucketRepositoriesForOrganization( + organization.id, + user.id, + false, + { bounded: true } + ); + expect(bounded.status).toBe('available'); + if (bounded.status !== 'available') throw new Error('Expected bounded repositories'); + expect(bounded.repositories).toHaveLength(50); + expect(bounded.repositories.at(-1)?.fullName).toBe('acme/repo-49'); + const legacy = await fetchBitbucketRepositoriesForOrganization(organization.id, user.id); + expect(legacy.status).toBe('available'); + if (legacy.status !== 'available') throw new Error('Expected complete repositories'); + expect(legacy.repositories).toHaveLength(60); + }); + + it.each([false, true])( + 'uses bounded transport without persisting partial results on refresh=%s', + async forceRefresh => { + const integration = await insertStaticIntegration( + organization.id, + user.id, + forceRefresh ? {} : { repositories: null, syncedAt: null } + ); + mockFetchBitbucketRepositoriesFromTokenService.mockImplementation( + async (_user, _org, options) => { + if (!options?.bounded || !options.signal) throw new Error('Unbounded transport'); + return { status: 'available', repositories: [REFRESHED_REPOSITORY] }; + } + ); + await expect( + fetchBitbucketRepositoriesForOrganization(organization.id, user.id, forceRefresh, { + bounded: true, + }) + ).resolves.toMatchObject({ status: 'available', repositories: [REFRESHED_REPOSITORY] }); + const [unchanged] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, integration.id)); + expect(unchanged?.repositories).toEqual(forceRefresh ? [CACHED_REPOSITORY] : null); + expect( + unchanged?.repositories_synced_at + ? new Date(unchanged.repositories_synced_at).toISOString() + : null + ).toBe(forceRefresh ? CACHED_AT : null); + } + ); + + it.each(['rotation', 'reconnection', 'cache refresh'] as const)( + 'returns a bounded winner after %s during provider work', + async race => { + const integration = await insertStaticIntegration(organization.id, user.id); + let winnerId = integration.id; + const winner = Array.from({ length: 60 }, (_, index) => ({ + ...CACHED_REPOSITORY, + id: randomUUID(), + full_name: `acme/winner-${index}`, + })); + mockFetchBitbucketRepositoriesFromTokenService.mockImplementation(async () => { + if (race === 'reconnection') { + winnerId = await replaceWithWinnerIntegration(organization.id, user.id, integration.id); + } else if (race === 'rotation') { + await db + .update(platform_access_token_credentials) + .set({ credential_version: 2 }) + .where(eq(platform_access_token_credentials.platform_integration_id, integration.id)); + } + await db + .update(platform_integrations) + .set({ repositories: winner, repositories_synced_at: WINNER_CACHED_AT }) + .where(eq(platform_integrations.id, winnerId)); + return { status: 'reconnect_required' }; + }); + const result = await fetchBitbucketRepositoriesForOrganization( + organization.id, + user.id, + true, + { bounded: true } + ); + expect(result.status).toBe('available'); + if (result.status !== 'available') throw new Error('Expected the winner cache'); + expect(result.repositories).toHaveLength(50); + expect(result.repositories[0].fullName).toBe('acme/winner-0'); + const [saved] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, winnerId)); + expect(saved?.repositories).toEqual(winner); + } + ); + + it('does not return bounded provider data after parent invalidation', async () => { + const integration = await insertStaticIntegration(organization.id, user.id); + mockFetchBitbucketRepositoriesFromTokenService.mockImplementation(async () => { + await db + .update(platform_integrations) + .set({ auth_invalid_at: WINNER_CACHED_AT, auth_invalid_reason: 'provider_rejected' }) + .where(eq(platform_integrations.id, integration.id)); + return { status: 'available', repositories: [REFRESHED_REPOSITORY] }; + }); + await expect( + fetchBitbucketRepositoriesForOrganization(organization.id, user.id, true, { bounded: true }) + ).resolves.toEqual({ status: 'reconnect_required' }); + const [saved] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, integration.id)); + expect(saved?.repositories).toEqual([CACHED_REPOSITORY]); + expect(saved?.auth_invalid_reason).toBe('provider_rejected'); + }); + + it.each([ + ['reconnect_required', 'reconnect_required'], + ['insufficient_permissions', 'insufficient_permissions'], + ['temporarily_unavailable', 'temporarily_unavailable'], + ['not_connected', 'temporarily_unavailable'], + ] as const)( + 'reports configured provider %s without cache writes', + async (status, expectedStatus) => { + const integration = await insertStaticIntegration(organization.id, user.id); + mockFetchBitbucketRepositoriesFromTokenService.mockResolvedValue({ status }); + await expect( + fetchBitbucketRepositoriesForOrganization(organization.id, user.id, true, { bounded: true }) + ).resolves.toEqual({ status: expectedStatus }); + const [saved] = await db + .select() + .from(platform_integrations) + .where(eq(platform_integrations.id, integration.id)); + expect(saved?.repositories).toEqual([CACHED_REPOSITORY]); + } + ); }); diff --git a/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.ts b/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.ts index a3c3458671..80b76043f0 100644 --- a/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.ts +++ b/apps/web/src/lib/integrations/platforms/bitbucket/workspace-access-token-repository-cache.ts @@ -15,6 +15,10 @@ import { and, eq, exists, isNull } from 'drizzle-orm'; import { z } from 'zod'; import { db, type DrizzleTransaction } from '@/lib/drizzle'; import { INTEGRATION_STATUS } from '@/lib/integrations/core/constants'; +import { + REPOSITORY_READ_LIMITS, + type RepositoryReadOptions, +} from '@/lib/integrations/core/repository-read-limits'; import { BitbucketWorkspaceAccessTokenMetadataSchema } from './metadata'; import { BitbucketRepositorySchema, @@ -69,6 +73,7 @@ type AvailableRepositoryCache = Extract< type ReadCachedRepositoriesInput = { organizationId: string; expectedIntegrationId?: string; + readOptions?: RepositoryReadOptions; }; type RefreshRepositoriesInput = { @@ -80,6 +85,7 @@ type RefreshRepositoriesInput = { type RefreshRepositoriesForMemberInput = { organizationId: string; kiloUserId: string; + readOptions?: RepositoryReadOptions; }; type ListRepositoriesInput = { @@ -118,6 +124,7 @@ async function loadIntegration(organizationId: string) { .select({ integrationId: platform_integrations.id, integrationStatus: platform_integrations.integration_status, + suspendedAt: platform_integrations.suspended_at, installationId: platform_integrations.platform_installation_id, workspaceUuid: platform_integrations.platform_account_id, workspaceSlug: platform_integrations.platform_account_login, @@ -167,10 +174,15 @@ function repositoriesHaveUniqueIdentity( function parseCachedRepositories( repositoriesValue: unknown, repositoriesSyncedAt: string | null, - workspace: WorkspaceIdentity + workspace: WorkspaceIdentity, + readOptions?: RepositoryReadOptions ): AvailableRepositoryCache | null { if (repositoriesValue === null || repositoriesSyncedAt === null) return null; - const repositories = z.array(CachedRepositorySchema).safeParse(repositoriesValue); + const selected = + readOptions?.bounded && Array.isArray(repositoriesValue) + ? repositoriesValue.slice(0, REPOSITORY_READ_LIMITS.repositories) + : repositoriesValue; + const repositories = z.array(CachedRepositorySchema).safeParse(selected); if (!repositories.success) return null; const projected = repositories.data.map(repository => ({ @@ -201,7 +213,7 @@ function toIsoTimestamp(value: string | null): string | null { return Number.isFinite(parsed.getTime()) ? parsed.toISOString() : null; } -function parseIntegration(row: LoadedIntegration) { +function parseIntegration(row: LoadedIntegration, readOptions?: RepositoryReadOptions) { const metadata = BitbucketWorkspaceAccessTokenMetadataSchema.safeParse(row.metadata); const workspaceUuid = z.uuid().safeParse(row.workspaceUuid); const workspaceSlug = WorkspaceSlugSchema.safeParse(row.workspaceSlug); @@ -214,7 +226,7 @@ function parseIntegration(row: LoadedIntegration) { ? { ...workspaceIdentity, displayName: metadata.data.displayName } : null; const cache = workspace - ? parseCachedRepositories(row.repositories, row.repositoriesSyncedAt, workspace) + ? parseCachedRepositories(row.repositories, row.repositoriesSyncedAt, workspace, readOptions) : null; const parsedCredential = BitbucketWorkspaceAccessTokenCredentialRowSchema.safeParse( row.credential @@ -235,6 +247,7 @@ function parseIntegration(row: LoadedIntegration) { row.installationId === null && hasValidCredentialEvidence && row.integrationStatus === INTEGRATION_STATUS.ACTIVE && + (!readOptions?.bounded || row.suspendedAt === null) && row.authInvalidAt === null; const rotatable = row.integrationStatus === INTEGRATION_STATUS.ACTIVE && @@ -259,9 +272,10 @@ function parseIntegration(row: LoadedIntegration) { }; } -async function loadParsedIntegration(organizationId: string) { +async function loadParsedIntegration(organizationId: string, readOptions?: RepositoryReadOptions) { + readOptions?.signal?.throwIfAborted(); const row = await loadIntegration(organizationId); - return row ? parseIntegration(row) : null; + return row ? parseIntegration(row, readOptions) : null; } function isRefreshableIntegration( @@ -331,6 +345,7 @@ export async function getBitbucketWorkspaceAccessTokenStatus(organizationId: str export async function readCachedBitbucketWorkspaceAccessTokenRepositories({ organizationId, expectedIntegrationId, + readOptions, }: ReadCachedRepositoriesInput): Promise { const canonicalOrganizationId = canonicalizeUuid(organizationId); const canonicalExpectedIntegrationId = expectedIntegrationId @@ -340,7 +355,7 @@ export async function readCachedBitbucketWorkspaceAccessTokenRepositories({ return { status: 'invalid_request' }; } - const integration = await loadParsedIntegration(canonicalOrganizationId); + const integration = await loadParsedIntegration(canonicalOrganizationId, readOptions); if (!integration) return { status: 'not_connected' }; if ( canonicalExpectedIntegrationId && @@ -492,13 +507,14 @@ export async function refreshBitbucketWorkspaceAccessTokenRepositories({ export async function refreshBitbucketWorkspaceAccessTokenRepositoriesForMember({ organizationId, kiloUserId, + readOptions, }: RefreshRepositoriesForMemberInput): Promise { const canonicalOrganizationId = canonicalizeUuid(organizationId); if (!canonicalOrganizationId) { return { status: 'invalid_request' }; } - const integration = await loadParsedIntegration(canonicalOrganizationId); + const integration = await loadParsedIntegration(canonicalOrganizationId, readOptions); if (!integration) return { status: 'not_connected' }; if (!isRefreshableIntegration(integration)) { return { status: 'reconnect_required' }; @@ -509,6 +525,7 @@ export async function refreshBitbucketWorkspaceAccessTokenRepositoriesForMember( kiloUserId, integration, requireOrganizationManager: false, + readOptions, }); } @@ -517,18 +534,28 @@ async function refreshLoadedBitbucketWorkspaceAccessTokenRepositories({ kiloUserId, integration, requireOrganizationManager, + readOptions, }: { organizationId: string; kiloUserId: string; integration: RefreshableParsedIntegration; requireOrganizationManager: boolean; + readOptions?: RepositoryReadOptions; }): Promise { const workspaceIdentity = integration.workspaceIdentity; const credential = integration.credential; - const providerResult = await fetchBitbucketWorkspaceAccessTokenRepositoriesFromTokenService( - kiloUserId, - organizationId - ); + readOptions?.signal?.throwIfAborted(); + const providerResult = readOptions?.bounded + ? await fetchBitbucketWorkspaceAccessTokenRepositoriesFromTokenService( + kiloUserId, + organizationId, + readOptions + ).catch(() => ({ status: 'temporarily_unavailable' as const })) + : await fetchBitbucketWorkspaceAccessTokenRepositoriesFromTokenService( + kiloUserId, + organizationId + ); + readOptions?.signal?.throwIfAborted(); const stillCurrent = await db.transaction(async tx => { await lockBitbucketWorkspaceAccessTokenOrganization(tx, organizationId); return isObservedCredentialGenerationCurrent(tx, organizationId, { @@ -540,8 +567,24 @@ async function refreshLoadedBitbucketWorkspaceAccessTokenRepositories({ if (!stillCurrent) { return readCachedBitbucketWorkspaceAccessTokenRepositories({ organizationId, + readOptions, }); } + if (readOptions?.bounded) { + const current = await loadParsedIntegration(organizationId, readOptions); + if (!current) return { status: 'not_connected' }; + if (!isRefreshableIntegration(current)) return { status: 'reconnect_required' }; + if ( + current.row.integrationId !== integration.row.integrationId || + current.credential.id !== credential.id || + current.credential.version !== credential.version || + current.workspaceIdentity.uuid !== workspaceIdentity.uuid || + current.workspaceIdentity.slug !== workspaceIdentity.slug || + current.row.repositoriesSyncedAt !== integration.row.repositoriesSyncedAt + ) { + return current.cache ?? { status: 'temporarily_unavailable' }; + } + } if (providerResult.status !== 'available') return providerResult; if ( providerResult.repositories.some( @@ -554,6 +597,14 @@ async function refreshLoadedBitbucketWorkspaceAccessTokenRepositories({ return { status: 'invalid_request' }; } + if (readOptions?.bounded) { + return { + status: 'available', + repositories: providerResult.repositories.slice(0, REPOSITORY_READ_LIMITS.repositories), + syncedAt: new Date().toISOString(), + }; + } + const repositories = providerResult.repositories.map(repository => ({ id: repository.id, name: repository.name, diff --git a/apps/web/src/lib/integrations/platforms/gitlab/instance-url.test.ts b/apps/web/src/lib/integrations/platforms/gitlab/instance-url.test.ts index c61d7c30a3..5115e9190c 100644 --- a/apps/web/src/lib/integrations/platforms/gitlab/instance-url.test.ts +++ b/apps/web/src/lib/integrations/platforms/gitlab/instance-url.test.ts @@ -7,8 +7,10 @@ import { assertGitLabUrlResolvesSafely, buildGitLabPlatformRepositoryId, buildGitLabUrl, + GitLabInstanceUrlError, isDefaultGitLabInstanceUrl, normalizeGitLabInstanceUrl, + resolveGitLabUrlSafely, } from './instance-url'; const mockLookup = lookup as jest.Mock; @@ -51,6 +53,7 @@ describe('GitLab instance URL safety', () => { urlWithCredentials.password = 'pass'; it.each([ + ['not a URL', 'Invalid URL format'], ['ftp://gitlab.example.com', 'Invalid URL protocol'], ['http://gitlab.example.com', 'must use https'], [urlWithCredentials.toString(), 'must not include credentials'], @@ -66,7 +69,13 @@ describe('GitLab instance URL safety', () => { ['http://192.168.0.1', 'host is not allowed'], ['https://gitlab.local', 'host is not allowed'], ])('rejects unsafe instance URL %p', (url, message) => { - expect(() => normalizeGitLabInstanceUrl(url)).toThrow(message); + expect(() => normalizeGitLabInstanceUrl(url)).toThrow( + expect.objectContaining({ + name: 'GitLabInstanceUrlError', + message: expect.stringContaining(message), + reason: 'invalid_url', + }) + ); }); it('accepts hostnames that resolve to public addresses', async () => { @@ -80,12 +89,41 @@ describe('GitLab instance URL safety', () => { }); }); + it('distinguishes transient lookup failures while preserving the legacy error', async () => { + mockLookup.mockRejectedValueOnce( + Object.assign(new Error('getaddrinfo EAI_AGAIN gitlab.example.com'), { code: 'EAI_AGAIN' }) + ); + + const result = resolveGitLabUrlSafely('https://gitlab.example.com/api/v4/user'); + await expect(result).rejects.toBeInstanceOf(GitLabInstanceUrlError); + await expect(result).rejects.toMatchObject({ + name: 'GitLabInstanceUrlError', + message: 'GitLab instance URL host could not be resolved.', + reason: 'resolution_failed', + }); + }); + + it('distinguishes empty DNS answers from unsafe addresses', async () => { + mockLookup.mockResolvedValueOnce([]); + + await expect( + resolveGitLabUrlSafely('https://gitlab.example.com/api/v4/user') + ).rejects.toMatchObject({ + name: 'GitLabInstanceUrlError', + message: 'GitLab instance URL host could not be resolved.', + reason: 'resolution_failed', + }); + }); + it('rejects hostnames that resolve to unsafe addresses', async () => { mockLookup.mockResolvedValueOnce([{ address: '192.168.1.10', family: 4 }]); await expect( assertGitLabUrlResolvesSafely('https://gitlab.example.com/api/v4/user') - ).rejects.toThrow('resolves to an address that is not allowed'); + ).rejects.toMatchObject({ + message: 'GitLab instance URL host resolves to an address that is not allowed.', + reason: 'invalid_url', + }); }); it('rejects hostnames that resolve to deprecated IPv6 site-local addresses', async () => { @@ -93,13 +131,19 @@ describe('GitLab instance URL safety', () => { await expect( assertGitLabUrlResolvesSafely('https://gitlab.example.com/api/v4/user') - ).rejects.toThrow('resolves to an address that is not allowed'); + ).rejects.toMatchObject({ + message: 'GitLab instance URL host resolves to an address that is not allowed.', + reason: 'invalid_url', + }); }); it('rejects unsafe literal IP URLs during fetch-time validation', async () => { - await expect(assertGitLabUrlResolvesSafely('http://127.0.0.1/api/v4/user')).rejects.toThrow( - 'host is not allowed' - ); + await expect( + assertGitLabUrlResolvesSafely('http://127.0.0.1/api/v4/user') + ).rejects.toMatchObject({ + message: 'GitLab instance URL host is not allowed.', + reason: 'invalid_url', + }); expect(mockLookup).not.toHaveBeenCalled(); }); diff --git a/apps/web/src/lib/integrations/platforms/gitlab/instance-url.ts b/apps/web/src/lib/integrations/platforms/gitlab/instance-url.ts index 9ccc6ecf85..618ca8cf6a 100644 --- a/apps/web/src/lib/integrations/platforms/gitlab/instance-url.ts +++ b/apps/web/src/lib/integrations/platforms/gitlab/instance-url.ts @@ -4,7 +4,10 @@ import { isIP } from 'net'; export const DEFAULT_GITLAB_INSTANCE_URL = 'https://gitlab.com'; export class GitLabInstanceUrlError extends Error { - constructor(message: string) { + constructor( + message: string, + readonly reason: 'invalid_url' | 'resolution_failed' = 'invalid_url' + ) { super(message); this.name = 'GitLabInstanceUrlError'; } @@ -197,11 +200,17 @@ async function resolveHostnameSafely( try { addresses = await lookup(hostname, { all: true, verbatim: true }); } catch { - throw new GitLabInstanceUrlError('GitLab instance URL host could not be resolved.'); + throw new GitLabInstanceUrlError( + 'GitLab instance URL host could not be resolved.', + 'resolution_failed' + ); } if (addresses.length === 0) { - throw new GitLabInstanceUrlError('GitLab instance URL host could not be resolved.'); + throw new GitLabInstanceUrlError( + 'GitLab instance URL host could not be resolved.', + 'resolution_failed' + ); } for (const { address } of addresses) { diff --git a/apps/web/src/routers/cloud-agent-next-router.test.ts b/apps/web/src/routers/cloud-agent-next-router.test.ts index cece37ede7..0e1e2ccff7 100644 --- a/apps/web/src/routers/cloud-agent-next-router.test.ts +++ b/apps/web/src/routers/cloud-agent-next-router.test.ts @@ -1,5 +1,6 @@ import { describe, expect, it, jest, beforeAll, beforeEach } from '@jest/globals'; -import { createCallerFactory } from '@/lib/trpc/init'; +import type * as GitHubIntegrationHelpers from '@/lib/cloud-agent/github-integration-helpers'; +import type * as GitLabIntegrationHelpers from '@/lib/cloud-agent/gitlab-integration-helpers'; import type { User } from '@kilocode/db/schema'; import type { z } from 'zod'; import type { personalPrepareSessionNextSchema } from '@/routers/cloud-agent-next-schemas'; @@ -83,27 +84,10 @@ const mockIsFeatureFlagEnabledOrDevelopment = const mockVerifyUserOwnsSessionV2ByCloudAgentId = jest.fn<() => Promise<{ kiloSessionId: string } | null>>(); const mockGetBalanceForUser = jest.fn<(user: User) => Promise<{ balance: number }>>(); -const mockFetchGitHubRepositoriesForUser = jest.fn< - ( - userId: string, - forceRefresh: boolean - ) => Promise<{ - repositories: unknown[]; - integrationInstalled: boolean; - syncedAt: null; - }> ->(); -const mockFetchGitLabRepositoriesForUser = jest.fn< - ( - userId: string, - forceRefresh: boolean - ) => Promise<{ - repositories: unknown[]; - integrationInstalled: boolean; - syncedAt: null; - instanceUrl?: string; - }> ->(); +const mockFetchGitHubRepositoriesForUser = + jest.fn(); +const mockFetchGitLabRepositoriesForUser = + jest.fn(); const mockOrderRepositoriesByUsage = jest.fn< (params: { @@ -115,6 +99,22 @@ const mockOrderRepositoriesByUsage = }) => Promise >(); +// Caller contexts supply authentication; these unit tests do not open external services. +jest.mock('@/lib/user/server', () => ({ getUserFromAuth: jest.fn() })); +jest.mock('@/lib/admin/admin-access-log', () => ({ + authViaTokenFromHeaders: jest.fn(), + clientIpFromHeaders: jest.fn(), + emitAdminAccessEvent: jest.fn(), +})); +jest.mock('@/lib/drizzle', () => ({ db: {}, readDb: {} })); +jest.mock('@/lib/r2/cloud-agent-pending-uploads', () => ({ + linkPendingUploads: jest.fn(async () => undefined), + releasePendingUploads: jest.fn(async () => undefined), +})); +jest.mock('@/lib/cloud-agent/stream-ticket', () => ({ + signStreamTicket: jest.fn(() => ({ ticket: 'test-stream-ticket', expiresAt: 1_900_000_000 })), +})); + jest.mock('@/lib/tokens', () => ({ generateCloudAgentToken: jest.fn(() => 'cloud-agent-token'), })); @@ -188,11 +188,18 @@ let createCaller: (ctx: { user: User }) => { isEligible: boolean; accessLevel: 'full' | 'limited' | 'blocked'; }>; - listGitHubRepositories: (input: { forceRefresh: boolean }) => Promise; - listGitLabRepositories: (input: { forceRefresh: boolean }) => Promise; + listGitHubRepositories: (input: { + forceRefresh?: boolean; + bounded?: boolean; + }) => Promise; + listGitLabRepositories: (input: { + forceRefresh?: boolean; + bounded?: boolean; + }) => Promise; }; beforeAll(async () => { + const { createCallerFactory } = await import('@/lib/trpc/init'); const mod = await import('./cloud-agent-next-router'); createCaller = createCallerFactory(mod.cloudAgentNextRouter); }); @@ -539,6 +546,109 @@ describe('cloudAgentNextRouter helper procedures', () => { }); }); +describe.each([ + ['listGitHubRepositories', mockFetchGitHubRepositoriesForUser], + ['listGitLabRepositories', mockFetchGitLabRepositoriesForUser], +] as const)('bounded personal %s', (method, fetchRepositories) => { + const user = { id: 'oauth/member', is_admin: false } as User; + beforeEach(() => { + jest.clearAllMocks(); + fetchRepositories.mockReset(); + mockOrderRepositoriesByUsage.mockImplementation(async ({ repositories }) => + [...repositories].reverse() + ); + }); + + it.each([ + 'not_connected', + 'available', + 'suspended', + 'reconnect_required', + 'misconfigured', + 'temporarily_unavailable', + ] as const)('preserves the %s state and derives the actor from authentication', async status => { + const expected = { + status, + repositories: [], + integrationInstalled: status !== 'not_connected', + syncedAt: null, + }; + fetchRepositories.mockImplementation(async (actor, forceRefresh, options) => { + if (actor !== user.id || forceRefresh !== false || !options?.bounded || !options.signal) { + throw new Error('Incorrect bounded authority or options'); + } + return expected; + }); + await expect(createCaller({ user })[method]({ bounded: true })).resolves.toEqual(expected); + }); + + it('preserves omitted options, refresh defaults, output fields, and usage ordering', async () => { + const repositories = [ + { id: 1, name: 'one', fullName: 'org/one', private: false }, + { id: 2, name: 'two', fullName: 'org/two', private: true }, + ]; + fetchRepositories.mockImplementation(async (actor, refresh, options) => { + if (actor !== user.id || refresh !== false || options !== undefined) + throw new Error('Legacy call changed'); + return { status: 'available', repositories, integrationInstalled: true, syncedAt: null }; + }); + await expect(createCaller({ user })[method]({})).resolves.toEqual({ + repositories: [repositories[1], repositories[0]], + integrationInstalled: true, + syncedAt: null, + }); + }); + + it.each([ + ['limit', 51], + ['pages', 3], + ['responseBytes', 1_048_577], + ['timeoutMs', 30_001], + ['endpoint', 'https://example.com/repositories'], + ['userId', 'oauth/other-user'], + ['credentials', 'client-supplied-token'], + ['signal', { aborted: false }], + ['organizationId', '9a283301-b75d-4375-a1ba-e319a02e18b7'], + ] as const)('rejects client-supplied %s on bounded reads', async (key, value) => { + const input = { bounded: true, [key]: value }; + await expect(createCaller({ user })[method](input)).rejects.toMatchObject({ + code: 'BAD_REQUEST', + }); + expect(fetchRepositories).not.toHaveBeenCalled(); + }); + + it('returns temporary unavailability without a raw provider error', async () => { + fetchRepositories.mockRejectedValue(new Error('raw credential response')); + await expect(createCaller({ user })[method]({ bounded: true })).resolves.toEqual({ + status: 'temporarily_unavailable', + integrationInstalled: true, + repositories: [], + syncedAt: null, + }); + }); + + it('bounds the complete procedure, including usage ordering', async () => { + jest.useFakeTimers(); + try { + fetchRepositories.mockResolvedValue({ + status: 'available', + repositories: [], + integrationInstalled: true, + syncedAt: null, + }); + mockOrderRepositoriesByUsage.mockImplementation(() => new Promise(() => {})); + const result = createCaller({ user })[method]({ bounded: true }); + await jest.advanceTimersByTimeAsync(30_000); + await expect(result).resolves.toMatchObject({ + status: 'temporarily_unavailable', + repositories: [], + }); + } finally { + jest.useRealTimers(); + } + }); +}); + describe('cloudAgentNextRouter.prepareSession', () => { beforeEach(() => { jest.clearAllMocks(); diff --git a/apps/web/src/routers/cloud-agent-next-router.ts b/apps/web/src/routers/cloud-agent-next-router.ts index a3e7ea328d..a30fb0ecce 100644 --- a/apps/web/src/routers/cloud-agent-next-router.ts +++ b/apps/web/src/routers/cloud-agent-next-router.ts @@ -10,6 +10,7 @@ import { rethrowAsTerminalError } from '@/lib/cloud-agent-next/terminal-errors'; import { generateCloudAgentToken } from '@/lib/tokens'; import { isFeatureFlagEnabledOrDevelopment } from '@/lib/posthog-feature-flags'; import { fetchGitHubRepositoriesForUser } from '@/lib/cloud-agent/github-integration-helpers'; +import { withRepositoryReadDeadline } from '@/lib/integrations/core/repository-read-limits'; import { getGitLabInstanceUrlForUser, buildGitLabCloneUrl, @@ -59,6 +60,11 @@ import { generateMessageId } from '@kilocode/cloud-agent-sdk/message-id'; import { getBalanceForUser } from '@/lib/user/balance'; import { buildCloudAgentNextEligibility } from './cloud-agent-next-eligibility'; +const ListRepositoriesInput = z.union([ + z.object({ forceRefresh: z.boolean().default(false), bounded: z.literal(true) }).strict(), + z.object({ forceRefresh: z.boolean().default(false), bounded: z.literal(false).optional() }), +]); + function buildTerminalUrl(params: { cloudAgentSessionId: string; ptyId: string; @@ -540,11 +546,7 @@ export const cloudAgentNextRouter = createTRPCRouter({ * List GitHub repositories available for cloud agent sessions. */ listGitHubRepositories: baseProcedure - .input( - z.object({ - forceRefresh: z.boolean().optional().default(false), - }) - ) + .input(ListRepositoriesInput) .output( z.object({ repositories: z.array( @@ -559,32 +561,61 @@ export const cloudAgentNextRouter = createTRPCRouter({ integrationInstalled: z.boolean(), syncedAt: z.string().nullish(), errorMessage: z.string().optional(), + status: z + .enum([ + 'not_connected', + 'available', + 'suspended', + 'reconnect_required', + 'misconfigured', + 'temporarily_unavailable', + 'integration_limit_exceeded', + ]) + .optional(), }) ) .query(async ({ ctx, input }) => { - const result = await fetchGitHubRepositoriesForUser(ctx.user.id, input.forceRefresh); - return { - repositories: await orderRepositoriesByUsage({ - userId: ctx.user.id, - organizationId: null, - platform: 'github', - repositories: result.repositories, - }), - integrationInstalled: result.integrationInstalled, - syncedAt: result.syncedAt, - errorMessage: result.errorMessage, - }; + try { + return await withRepositoryReadDeadline( + input.bounded ? { bounded: true } : undefined, + async signal => { + const result = input.bounded + ? await fetchGitHubRepositoriesForUser(ctx.user.id, input.forceRefresh, { + bounded: true, + signal, + }) + : await fetchGitHubRepositoriesForUser(ctx.user.id, input.forceRefresh); + signal?.throwIfAborted(); + return { + repositories: await orderRepositoriesByUsage({ + userId: ctx.user.id, + organizationId: null, + platform: 'github', + repositories: result.repositories, + }), + integrationInstalled: result.integrationInstalled, + syncedAt: result.syncedAt, + errorMessage: result.errorMessage, + ...(input.bounded ? { status: result.status } : {}), + }; + } + ); + } catch (error) { + if (!input.bounded) throw error; + return { + status: 'temporarily_unavailable' as const, + integrationInstalled: true, + repositories: [], + syncedAt: null, + }; + } }), /** * List GitLab repositories available for cloud agent sessions. */ listGitLabRepositories: baseProcedure - .input( - z.object({ - forceRefresh: z.boolean().optional().default(false), - }) - ) + .input(ListRepositoriesInput) .output( z.object({ repositories: z.array( @@ -598,21 +629,54 @@ export const cloudAgentNextRouter = createTRPCRouter({ integrationInstalled: z.boolean(), syncedAt: z.string().nullish(), errorMessage: z.string().optional(), + status: z + .enum([ + 'not_connected', + 'available', + 'suspended', + 'reconnect_required', + 'misconfigured', + 'temporarily_unavailable', + 'integration_limit_exceeded', + ]) + .optional(), }) ) .query(async ({ ctx, input }) => { - const result = await fetchGitLabRepositoriesForUser(ctx.user.id, input.forceRefresh); - return { - repositories: await orderRepositoriesByUsage({ - userId: ctx.user.id, - organizationId: null, - platform: 'gitlab', - repositories: result.repositories, - gitlabInstanceUrl: result.instanceUrl, - }), - integrationInstalled: result.integrationInstalled, - syncedAt: result.syncedAt, - errorMessage: result.errorMessage, - }; + try { + return await withRepositoryReadDeadline( + input.bounded ? { bounded: true } : undefined, + async signal => { + const result = input.bounded + ? await fetchGitLabRepositoriesForUser(ctx.user.id, input.forceRefresh, { + bounded: true, + signal, + }) + : await fetchGitLabRepositoriesForUser(ctx.user.id, input.forceRefresh); + signal?.throwIfAborted(); + return { + repositories: await orderRepositoriesByUsage({ + userId: ctx.user.id, + organizationId: null, + platform: 'gitlab', + repositories: result.repositories, + gitlabInstanceUrl: result.instanceUrl, + }), + integrationInstalled: result.integrationInstalled, + syncedAt: result.syncedAt, + errorMessage: result.errorMessage, + ...(input.bounded ? { status: result.status } : {}), + }; + } + ); + } catch (error) { + if (!input.bounded) throw error; + return { + status: 'temporarily_unavailable' as const, + integrationInstalled: true, + repositories: [], + syncedAt: null, + }; + } }), }); diff --git a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts index d0f13f8c3d..788c31a641 100644 --- a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts +++ b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts @@ -1,5 +1,6 @@ import { describe, expect, it, jest, beforeAll, beforeEach } from '@jest/globals'; -import { createCallerFactory } from '@/lib/trpc/init'; +import type * as GitHubIntegrationHelpers from '@/lib/cloud-agent/github-integration-helpers'; +import type * as GitLabIntegrationHelpers from '@/lib/cloud-agent/gitlab-integration-helpers'; import type * as TrpcInitModule from '@/lib/trpc/init'; import type * as ZodModule from 'zod'; import type { z } from 'zod'; @@ -87,37 +88,13 @@ const mockIsFeatureFlagEnabledOrDevelopment = const mockVerifyOrgOwnsSessionV2ByCloudAgentId = jest.fn<() => Promise<{ kiloSessionId: string } | null>>(); const mockFetchBitbucketRepositoriesForOrganization = - jest.fn< - ( - organizationId: string, - kiloUserId: string, - forceRefresh?: boolean - ) => Promise - >(); + jest.fn(); const mockGetBalanceForOrganizationUser = jest.fn<(organizationId: string, userId: string) => Promise<{ balance: number }>>(); -const mockFetchGitHubRepositoriesForOrganization = jest.fn< - ( - organizationId: string, - forceRefresh: boolean - ) => Promise<{ - repositories: unknown[]; - integrationInstalled: boolean; - syncedAt: null; - }> ->(); -const mockFetchGitLabRepositoriesForOrganization = jest.fn< - ( - organizationId: string, - actorUserId: string, - forceRefresh: boolean - ) => Promise<{ - repositories: unknown[]; - integrationInstalled: boolean; - syncedAt: null; - instanceUrl?: string; - }> ->(); +const mockFetchGitHubRepositoriesForOrganization = + jest.fn(); +const mockFetchGitLabRepositoriesForOrganization = + jest.fn(); const mockOrderRepositoriesByUsage = jest.fn< (params: { @@ -130,6 +107,23 @@ const mockOrderRepositoriesByUsage = >(); const mockEnsureOrganizationAccess = jest.fn<(userId: string, organizationId: string) => void>(); +// Caller contexts supply authentication; these unit tests do not open external services. +jest.mock('@/lib/user/server', () => ({ getUserFromAuth: jest.fn() })); +jest.mock('@/lib/admin/admin-access-log', () => ({ + authViaTokenFromHeaders: jest.fn(), + clientIpFromHeaders: jest.fn(), + emitAdminAccessEvent: jest.fn(), +})); +jest.mock('@/lib/drizzle', () => ({ db: {}, readDb: {} })); +jest.mock('@/lib/config.server', () => ({})); +jest.mock('@/lib/r2/cloud-agent-pending-uploads', () => ({ + linkPendingUploads: jest.fn(async () => undefined), + releasePendingUploads: jest.fn(async () => undefined), +})); +jest.mock('@/lib/cloud-agent/stream-ticket', () => ({ + signStreamTicket: jest.fn(() => ({ ticket: 'test-stream-ticket', expiresAt: 1_900_000_000 })), +})); + jest.mock('@/lib/tokens', () => ({ generateCloudAgentToken: jest.fn(() => 'cloud-agent-token'), })); @@ -196,6 +190,8 @@ jest.mock('@/routers/organizations/utils', () => { return { organizationMemberProcedure: organizationProcedure, organizationMemberMutationProcedure: organizationProcedure, + ensureOrganizationAccess: (ctx: { user: User }, organizationId: string) => + mockEnsureOrganizationAccess(ctx.user.id, organizationId), }; }); @@ -230,6 +226,7 @@ let createCaller: (ctx: { user: User }) => { listBitbucketRepositories: (input: { organizationId: string; forceRefresh?: boolean; + bounded?: boolean; }) => Promise; checkEligibility: (input: { organizationId: string }) => Promise<{ balance: number; @@ -239,11 +236,13 @@ let createCaller: (ctx: { user: User }) => { }>; listGitHubRepositories: (input: { organizationId: string; - forceRefresh: boolean; + forceRefresh?: boolean; + bounded?: boolean; }) => Promise; listGitLabRepositories: (input: { organizationId: string; - forceRefresh: boolean; + forceRefresh?: boolean; + bounded?: boolean; }) => Promise; refreshTerminalTicket: (input: { organizationId: string; @@ -269,6 +268,7 @@ let createCaller: (ctx: { user: User }) => { }; beforeAll(async () => { + const { createCallerFactory } = await import('@/lib/trpc/init'); const mod = await import('./organization-cloud-agent-next-router'); createCaller = createCallerFactory(mod.organizationCloudAgentNextRouter); }); @@ -659,6 +659,200 @@ describe('organizationCloudAgentNextRouter helper procedures', () => { }); }); +describe.each([ + 'listGitHubRepositories', + 'listGitLabRepositories', + 'listBitbucketRepositories', +] as const)('bounded organization %s', method => { + const user = { id: 'oauth/member', is_admin: false } as User; + const numeric = { + status: 'available' as const, + repositories: [], + integrationInstalled: true, + syncedAt: null, + }; + const bitbucket = { + status: 'available' as const, + repositories: [], + syncedAt: '2026-06-25T18:00:00.000Z', + }; + const expected = method === 'listBitbucketRepositories' ? bitbucket : numeric; + + beforeEach(() => { + jest.clearAllMocks(); + mockEnsureOrganizationAccess.mockReset().mockImplementation(() => undefined); + mockOrderRepositoriesByUsage.mockImplementation(async ({ repositories }) => repositories); + mockFetchGitHubRepositoriesForOrganization + .mockReset() + .mockImplementation(async (organizationId, refresh, options) => { + if ( + organizationId !== ORGANIZATION_ID || + refresh !== false || + !options?.bounded || + !options.signal + ) + throw new Error('Incorrect bounded request'); + return numeric; + }); + mockFetchGitLabRepositoriesForOrganization + .mockReset() + .mockImplementation(async (organizationId, actor, refresh, options) => { + if ( + organizationId !== ORGANIZATION_ID || + actor !== user.id || + refresh !== false || + !options?.bounded || + !options.signal + ) + throw new Error('Incorrect bounded request'); + return numeric; + }); + mockFetchBitbucketRepositoriesForOrganization + .mockReset() + .mockImplementation(async (organizationId, actor, refresh, options) => { + if ( + organizationId !== ORGANIZATION_ID || + actor !== user.id || + refresh !== false || + !options?.bounded || + !options.signal + ) + throw new Error('Incorrect bounded request'); + return bitbucket; + }); + }); + + it('uses fixed bounds, authenticated ownership, and the existing refresh default', async () => { + await expect( + createCaller({ user })[method]({ organizationId: ORGANIZATION_ID, bounded: true }) + ).resolves.toEqual(expected); + }); + + it.each([false, true])( + 'denies membership removed during bounded refresh=%s', + async forceRefresh => { + const revoke = () => + mockEnsureOrganizationAccess.mockImplementation(() => { + throw new TRPCError({ code: 'UNAUTHORIZED', message: 'Membership removed' }); + }); + mockFetchGitHubRepositoriesForOrganization.mockImplementation(async () => { + revoke(); + return numeric; + }); + mockFetchGitLabRepositoriesForOrganization.mockImplementation(async () => { + revoke(); + return numeric; + }); + mockFetchBitbucketRepositoriesForOrganization.mockImplementation(async () => { + revoke(); + return bitbucket; + }); + await expect( + createCaller({ user })[method]({ + organizationId: ORGANIZATION_ID, + forceRefresh, + bounded: true, + }) + ).rejects.toMatchObject({ code: 'UNAUTHORIZED', message: 'Membership removed' }); + } + ); + + it('denies removed membership before reading provider data', async () => { + mockEnsureOrganizationAccess.mockImplementation(() => { + throw new TRPCError({ code: 'UNAUTHORIZED', message: 'Membership removed' }); + }); + await expect( + createCaller({ user })[method]({ organizationId: ORGANIZATION_ID, bounded: true }) + ).rejects.toMatchObject({ code: 'UNAUTHORIZED' }); + expect(mockFetchGitHubRepositoriesForOrganization).not.toHaveBeenCalled(); + expect(mockFetchGitLabRepositoriesForOrganization).not.toHaveBeenCalled(); + expect(mockFetchBitbucketRepositoriesForOrganization).not.toHaveBeenCalled(); + }); + + it.each([ + ['limit', 51], + ['pages', 3], + ['responseBytes', 1_048_577], + ['timeoutMs', 30_001], + ['endpoint', 'https://example.com/repositories'], + ['actorUserId', 'oauth/other-user'], + ['credentials', 'client-supplied-token'], + ['signal', { aborted: false }], + ] as const)('rejects client-supplied %s', async (key, value) => { + const input = { organizationId: ORGANIZATION_ID, bounded: true, [key]: value }; + await expect(createCaller({ user })[method](input)).rejects.toMatchObject({ + code: 'BAD_REQUEST', + }); + }); + + it.each(['not_connected', 'reconnect_required', 'temporarily_unavailable'] as const)( + 'preserves the %s state', + async status => { + const result = { ...numeric, status, integrationInstalled: status !== 'not_connected' }; + mockFetchGitHubRepositoriesForOrganization.mockResolvedValue(result); + mockFetchGitLabRepositoriesForOrganization.mockResolvedValue(result); + mockFetchBitbucketRepositoriesForOrganization.mockResolvedValue({ status }); + await expect( + createCaller({ user })[method]({ organizationId: ORGANIZATION_ID, bounded: true }) + ).resolves.toEqual(method === 'listBitbucketRepositories' ? { status } : result); + } + ); + + it('preserves old output and default refresh when bounded options are omitted', async () => { + mockFetchGitHubRepositoriesForOrganization.mockImplementation( + async (_org, refresh, options) => { + if (refresh !== false || options !== undefined) throw new Error('Legacy call changed'); + return numeric; + } + ); + mockFetchGitLabRepositoriesForOrganization.mockImplementation( + async (_org, _actor, refresh, options) => { + if (refresh !== false || options !== undefined) throw new Error('Legacy call changed'); + return numeric; + } + ); + mockFetchBitbucketRepositoriesForOrganization.mockImplementation( + async (_org, _actor, refresh, options) => { + if (refresh !== false || options !== undefined) throw new Error('Legacy call changed'); + return bitbucket; + } + ); + await expect( + createCaller({ user })[method]({ organizationId: ORGANIZATION_ID }) + ).resolves.toEqual( + method === 'listBitbucketRepositories' + ? bitbucket + : { repositories: [], integrationInstalled: true, syncedAt: null } + ); + }); + + it('returns an explicit temporary state at the operation deadline', async () => { + jest.useFakeTimers(); + try { + mockFetchGitHubRepositoriesForOrganization.mockImplementation(() => new Promise(() => {})); + mockFetchGitLabRepositoriesForOrganization.mockImplementation(() => new Promise(() => {})); + mockFetchBitbucketRepositoriesForOrganization.mockImplementation(() => new Promise(() => {})); + const result = createCaller({ user })[method]({ + organizationId: ORGANIZATION_ID, + bounded: true, + }); + await jest.advanceTimersByTimeAsync(30_000); + await expect(result).resolves.toEqual( + method === 'listBitbucketRepositories' + ? { status: 'temporarily_unavailable' } + : { + status: 'temporarily_unavailable', + integrationInstalled: true, + repositories: [], + syncedAt: null, + } + ); + } finally { + jest.useRealTimers(); + } + }); +}); + describe('organizationCloudAgentNextRouter terminal ownership', () => { const organizationCloudAgentSessionId = 'agent_terminal_ticket_org_owned'; const personalCloudAgentSessionId = 'agent_terminal_ticket_org_personal'; diff --git a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts index fb351c3e1b..91d7df5b5b 100644 --- a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts +++ b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts @@ -12,7 +12,9 @@ import { isFeatureFlagEnabledOrDevelopment } from '@/lib/posthog-feature-flags'; import { organizationMemberProcedure, organizationMemberMutationProcedure, + ensureOrganizationAccess, } from '@/routers/organizations/utils'; +import { withRepositoryReadDeadline } from '@/lib/integrations/core/repository-read-limits'; import { fetchAllGitHubRepositoriesForOrganization } from '@/lib/cloud-agent/github-integration-helpers'; import { BitbucketOrganizationRepositoryListResultSchema, @@ -194,20 +196,20 @@ const AnswerPermissionInput = baseAnswerPermissionNextSchema.extend({ organizationId: z.uuid(), }); -const ListGitHubRepositoriesInput = z.object({ - organizationId: z.uuid(), - forceRefresh: z.boolean().optional().default(false), -}); - -const ListGitLabRepositoriesInput = z.object({ - organizationId: z.uuid(), - forceRefresh: z.boolean().optional().default(false), -}); - -const ListBitbucketRepositoriesInput = z.object({ - organizationId: z.uuid(), - forceRefresh: z.boolean().optional().default(false), -}); +const ListRepositoriesInput = z.union([ + z + .object({ + organizationId: z.uuid(), + forceRefresh: z.boolean().default(false), + bounded: z.literal(true), + }) + .strict(), + z.object({ + organizationId: z.uuid(), + forceRefresh: z.boolean().default(false), + bounded: z.literal(false).optional(), + }), +]); /** * Cloud Agent Next Router (Organization Context) @@ -727,7 +729,7 @@ export const organizationCloudAgentNextRouter = createTRPCRouter({ * List GitHub repositories available for cloud agent sessions (organization context). */ listGitHubRepositories: organizationMemberProcedure - .input(ListGitHubRepositoriesInput) + .input(ListRepositoriesInput) .output( z.object({ repositories: z.array( @@ -744,31 +746,75 @@ export const organizationCloudAgentNextRouter = createTRPCRouter({ integrationInstalled: z.boolean(), syncedAt: z.string().nullish(), errorMessage: z.string().optional(), + status: z + .enum([ + 'not_connected', + 'available', + 'suspended', + 'reconnect_required', + 'misconfigured', + 'temporarily_unavailable', + 'integration_limit_exceeded', + ]) + .optional(), }) ) .query(async ({ ctx, input }) => { - const result = await fetchAllGitHubRepositoriesForOrganization( - input.organizationId, - input.forceRefresh - ); - return { - repositories: await orderRepositoriesByUsage({ - userId: ctx.user.id, - organizationId: input.organizationId, - platform: 'github', - repositories: result.repositories, - }), - integrationInstalled: result.integrationInstalled, - syncedAt: result.syncedAt, - errorMessage: result.errorMessage, - }; + try { + return await withRepositoryReadDeadline( + input.bounded ? { bounded: true } : undefined, + async signal => { + const result = input.bounded + ? await fetchAllGitHubRepositoriesForOrganization( + input.organizationId, + input.forceRefresh, + { bounded: true, signal } + ) + : await fetchAllGitHubRepositoriesForOrganization( + input.organizationId, + input.forceRefresh + ); + signal?.throwIfAborted(); + const repositories = await orderRepositoriesByUsage({ + userId: ctx.user.id, + organizationId: input.organizationId, + platform: 'github', + repositories: result.repositories, + }); + if (input.bounded) { + signal?.throwIfAborted(); + await ensureOrganizationAccess(ctx, input.organizationId); + } + return { + repositories, + integrationInstalled: result.integrationInstalled, + syncedAt: result.syncedAt, + errorMessage: result.errorMessage, + ...(input.bounded ? { status: result.status } : {}), + }; + } + ); + } catch (error) { + if ( + !input.bounded || + (error instanceof TRPCError && + (error.code === 'UNAUTHORIZED' || error.code === 'FORBIDDEN')) + ) + throw error; + return { + status: 'temporarily_unavailable' as const, + integrationInstalled: true, + repositories: [], + syncedAt: null, + }; + } }), /** * List GitLab repositories available for cloud agent sessions (organization context). */ listGitLabRepositories: organizationMemberProcedure - .input(ListGitLabRepositoriesInput) + .input(ListRepositoriesInput) .output( z.object({ repositories: z.array( @@ -782,48 +828,120 @@ export const organizationCloudAgentNextRouter = createTRPCRouter({ integrationInstalled: z.boolean(), syncedAt: z.string().nullish(), errorMessage: z.string().optional(), + status: z + .enum([ + 'not_connected', + 'available', + 'suspended', + 'reconnect_required', + 'misconfigured', + 'temporarily_unavailable', + ]) + .optional(), }) ) .query(async ({ ctx, input }) => { - const result = await fetchGitLabRepositoriesForOrganization( - input.organizationId, - ctx.user.id, - input.forceRefresh - ); - return { - repositories: await orderRepositoriesByUsage({ - userId: ctx.user.id, - organizationId: input.organizationId, - platform: 'gitlab', - repositories: result.repositories, - gitlabInstanceUrl: result.instanceUrl, - }), - integrationInstalled: result.integrationInstalled, - syncedAt: result.syncedAt, - errorMessage: result.errorMessage, - }; + try { + return await withRepositoryReadDeadline( + input.bounded ? { bounded: true } : undefined, + async signal => { + const result = input.bounded + ? await fetchGitLabRepositoriesForOrganization( + input.organizationId, + ctx.user.id, + input.forceRefresh, + { bounded: true, signal } + ) + : await fetchGitLabRepositoriesForOrganization( + input.organizationId, + ctx.user.id, + input.forceRefresh + ); + signal?.throwIfAborted(); + const repositories = await orderRepositoriesByUsage({ + userId: ctx.user.id, + organizationId: input.organizationId, + platform: 'gitlab', + repositories: result.repositories, + gitlabInstanceUrl: result.instanceUrl, + }); + if (input.bounded) { + signal?.throwIfAborted(); + await ensureOrganizationAccess(ctx, input.organizationId); + } + return { + repositories, + integrationInstalled: result.integrationInstalled, + syncedAt: result.syncedAt, + errorMessage: result.errorMessage, + ...(input.bounded ? { status: result.status } : {}), + }; + } + ); + } catch (error) { + if ( + !input.bounded || + (error instanceof TRPCError && + (error.code === 'UNAUTHORIZED' || error.code === 'FORBIDDEN')) + ) + throw error; + return { + status: 'temporarily_unavailable' as const, + integrationInstalled: true, + repositories: [], + syncedAt: null, + }; + } }), listBitbucketRepositories: organizationMemberProcedure - .input(ListBitbucketRepositoriesInput) + .input(ListRepositoriesInput) .output(BitbucketOrganizationRepositoryListResultSchema) .query(async ({ ctx, input }) => { - const result = await fetchBitbucketRepositoriesForOrganization( - input.organizationId, - ctx.user.id, - input.forceRefresh - ); - if (result.status !== 'available') { - return result; + try { + return await withRepositoryReadDeadline( + input.bounded ? { bounded: true } : undefined, + async signal => { + const result = input.bounded + ? await fetchBitbucketRepositoriesForOrganization( + input.organizationId, + ctx.user.id, + input.forceRefresh, + { bounded: true, signal } + ) + : await fetchBitbucketRepositoriesForOrganization( + input.organizationId, + ctx.user.id, + input.forceRefresh + ); + signal?.throwIfAborted(); + const output = + result.status === 'available' + ? { + ...result, + repositories: await orderRepositoriesByUsage({ + userId: ctx.user.id, + organizationId: input.organizationId, + platform: 'bitbucket', + repositories: result.repositories, + }), + } + : result; + if (input.bounded) { + signal?.throwIfAborted(); + await ensureOrganizationAccess(ctx, input.organizationId); + } + return output; + } + ); + } catch (error) { + if ( + !input.bounded || + (error instanceof TRPCError && + (error.code === 'UNAUTHORIZED' || error.code === 'FORBIDDEN')) + ) + throw error; + return { status: 'temporarily_unavailable' as const }; } - return { - ...result, - repositories: await orderRepositoriesByUsage({ - userId: ctx.user.id, - organizationId: input.organizationId, - platform: 'bitbucket', - repositories: result.repositories, - }), - }; }), });