From bdeefe2a11a2602b93aad885563894637f80e3b3 Mon Sep 17 00:00:00 2001 From: Guilherme Caulada Date: Wed, 23 Sep 2026 11:14:07 -0300 Subject: [PATCH 1/3] feat: add stateless registration cleanup with multi-App support --- .../termination-watcher/package.json | 3 +- .../src/github-app-client.test.ts | 161 +++++++ .../src/github-app-client.ts | 167 +++++++ .../termination-watcher/src/lambda.ts | 2 + .../src/registration-janitor.test.ts | 434 ++++++++++++++++++ .../src/registration-janitor.ts | 230 ++++++++++ .../aws/ec2/registration-cleanup.ts | 45 ++ lambdas/libs/compute-providers/package.json | 3 +- .../providers.config.registration-cleanup.ts | 7 + .../registration-cleanup.test.ts | 53 +++ .../compute-providers/registration-cleanup.ts | 28 ++ lambdas/yarn.lock | 1 + modules/registration-janitor/README.md | 117 +++++ modules/registration-janitor/compute-ec2.tf | 27 ++ modules/registration-janitor/main.tf | 116 +++++ modules/registration-janitor/outputs.tf | 9 + .../tests/registration-janitor.tftest.hcl | 175 +++++++ modules/registration-janitor/variables.tf | 53 +++ modules/registration-janitor/versions.tf | 10 + 19 files changed, 1639 insertions(+), 2 deletions(-) create mode 100644 lambdas/functions/termination-watcher/src/github-app-client.test.ts create mode 100644 lambdas/functions/termination-watcher/src/github-app-client.ts create mode 100644 lambdas/functions/termination-watcher/src/registration-janitor.test.ts create mode 100644 lambdas/functions/termination-watcher/src/registration-janitor.ts create mode 100644 lambdas/libs/compute-providers/aws/ec2/registration-cleanup.ts create mode 100644 lambdas/libs/compute-providers/providers.config.registration-cleanup.ts create mode 100644 lambdas/libs/compute-providers/registration-cleanup.test.ts create mode 100644 lambdas/libs/compute-providers/registration-cleanup.ts create mode 100644 modules/registration-janitor/README.md create mode 100644 modules/registration-janitor/compute-ec2.tf create mode 100644 modules/registration-janitor/main.tf create mode 100644 modules/registration-janitor/outputs.tf create mode 100644 modules/registration-janitor/tests/registration-janitor.tftest.hcl create mode 100644 modules/registration-janitor/variables.tf create mode 100644 modules/registration-janitor/versions.tf diff --git a/lambdas/functions/termination-watcher/package.json b/lambdas/functions/termination-watcher/package.json index 3c58b5c272..d3d8cd1820 100644 --- a/lambdas/functions/termination-watcher/package.json +++ b/lambdas/functions/termination-watcher/package.json @@ -32,7 +32,8 @@ "@octokit/core": "7.0.6", "@octokit/plugin-throttling": "11.0.3", "@octokit/request": "^9.2.2", - "@octokit/rest": "22.0.1" + "@octokit/rest": "22.0.1", + "@aws-github-runner/compute-providers": "*" }, "nx": { "includedScripts": [ diff --git a/lambdas/functions/termination-watcher/src/github-app-client.test.ts b/lambdas/functions/termination-watcher/src/github-app-client.test.ts new file mode 100644 index 0000000000..17d278221e --- /dev/null +++ b/lambdas/functions/termination-watcher/src/github-app-client.test.ts @@ -0,0 +1,161 @@ +import { createCommonStorage, type GitHubAppCredentialsStore } from '@aws-github-runner/storage-providers'; +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import { createThrottleOptions, resetAppCredentialsCache, createRunnerInstallationClient } from './github-app-client'; +import type { EndpointDefaults } from '@octokit/types'; + +vi.mock('@aws-github-runner/storage-providers', () => ({ + createCommonStorage: vi.fn(), +})); + +const mockedCreateCommonStorage = vi.mocked(createCommonStorage); +const mockGetCredentials = vi.fn(); +const credentialsStore = { get: mockGetCredentials } satisfies GitHubAppCredentialsStore; + +const mockCreateAppAuth = vi.fn(); +vi.mock('@octokit/auth-app', () => ({ + createAppAuth: (...args: unknown[]) => mockCreateAppAuth(...args), +})); + +const mockPaginate = { + iterator: vi.fn(), +}; + +const mockActions = { + listSelfHostedRunnersForOrg: vi.fn(), + listSelfHostedRunnersForRepo: vi.fn(), + deleteSelfHostedRunnerFromOrg: vi.fn(), + deleteSelfHostedRunnerFromRepo: vi.fn(), +}; + +const mockApps = { + getOrgInstallation: vi.fn(), + getRepoInstallation: vi.fn(), +}; + +const mockHookAfter = vi.fn(); + +function MockOctokit() { + return { + hook: { after: mockHookAfter }, + actions: mockActions, + apps: mockApps, + paginate: mockPaginate, + }; +} +MockOctokit.plugin = vi.fn().mockReturnValue(MockOctokit); + +vi.mock('@octokit/rest', () => ({ + Octokit: MockOctokit, +})); + +vi.mock('@octokit/plugin-throttling', () => ({ + throttling: vi.fn(), +})); + +vi.mock('@octokit/request', () => ({ + request: { + defaults: vi.fn().mockReturnValue(vi.fn()), + }, +})); + +function setupAuthMocks() { + mockGetCredentials.mockResolvedValue([{ appId: 12345, privateKey: 'fake-private-key' }]); + + // App auth returns app token + const mockAuth = vi.fn(); + mockAuth.mockImplementation((opts: { type: string }) => { + if (opts.type === 'app') { + return Promise.resolve({ token: 'app-token' }); + } + return Promise.resolve({ token: 'installation-token' }); + }); + mockCreateAppAuth.mockReturnValue(mockAuth); +} + +describe('multi-App installation clients', () => { + beforeEach(() => { + vi.clearAllMocks(); + resetAppCredentialsCache(); + mockedCreateCommonStorage.mockReturnValue({ githubAppCredentials: credentialsStore }); + setupAuthMocks(); + }); + + it('creates an organization installation client for scheduled registration cleanup', async () => { + mockApps.getOrgInstallation.mockResolvedValue({ data: { id: 999 } }); + const client = await createRunnerInstallationClient('test-org', 'Org', ''); + expect(client.actions).toBe(mockActions); + expect(mockApps.getOrgInstallation).toHaveBeenCalledWith({ org: 'test-org' }); + expect(mockCreateAppAuth).toHaveBeenCalledWith({ + appId: 12345, + privateKey: 'fake-private-key', + installationId: 999, + }); + }); + + it('uses additional credentials and keeps App and installation auth paired', async () => { + mockGetCredentials.mockResolvedValue([ + { appId: 1, privateKey: 'one' }, + { appId: 2, privateKey: 'two' }, + ]); + mockApps.getOrgInstallation.mockResolvedValue({ data: { id: 222 } }); + const random = vi.spyOn(Math, 'random').mockReturnValue(0.99); + try { + await createRunnerInstallationClient('test-org', 'Org', ''); + expect(mockCreateAppAuth).toHaveBeenNthCalledWith(1, { appId: 2, privateKey: 'two' }); + expect(mockCreateAppAuth).toHaveBeenNthCalledWith(2, { appId: 2, privateKey: 'two', installationId: 222 }); + } finally { + random.mockRestore(); + } + }); + + it('selects another App after an exhausted installation response', async () => { + mockGetCredentials.mockResolvedValue([ + { appId: 1, privateKey: 'one' }, + { appId: 2, privateKey: 'two' }, + ]); + mockApps.getOrgInstallation.mockResolvedValue({ data: { id: 222 } }); + const random = vi.spyOn(Math, 'random').mockReturnValue(0); + try { + await createRunnerInstallationClient('test-org', 'Org', ''); + mockHookAfter.mock.calls[0][1]({ headers: { 'x-ratelimit-remaining': '0' } }); + mockCreateAppAuth.mockClear(); + await createRunnerInstallationClient('test-org', 'Org', ''); + expect(mockCreateAppAuth).toHaveBeenNthCalledWith(1, { appId: 2, privateKey: 'two' }); + } finally { + random.mockRestore(); + } + }); + + it('skips an App in secondary-limit cooldown', async () => { + mockGetCredentials.mockResolvedValue([ + { appId: 1, privateKey: 'one' }, + { appId: 2, privateKey: 'two' }, + ]); + mockApps.getOrgInstallation.mockResolvedValue({ data: { id: 222 } }); + createThrottleOptions(1).onSecondaryRateLimit(60, { method: 'GET', url: '/runners' } as Required); + const random = vi.spyOn(Math, 'random').mockReturnValue(0); + try { + await createRunnerInstallationClient('test-org', 'Org', ''); + expect(mockCreateAppAuth).toHaveBeenNthCalledWith(1, { appId: 2, privateKey: 'two' }); + } finally { + random.mockRestore(); + } + }); + + it('tries another configured App when installation authentication fails', async () => { + mockGetCredentials.mockResolvedValue([ + { appId: 1, privateKey: 'one' }, + { appId: 2, privateKey: 'two' }, + ]); + mockApps.getOrgInstallation + .mockRejectedValueOnce(new Error('not installed')) + .mockResolvedValue({ data: { id: 222 } }); + const random = vi.spyOn(Math, 'random').mockReturnValue(0); + try { + await createRunnerInstallationClient('test-org', 'Org', ''); + expect(mockCreateAppAuth).toHaveBeenLastCalledWith({ appId: 2, privateKey: 'two', installationId: 222 }); + } finally { + random.mockRestore(); + } + }); +}); diff --git a/lambdas/functions/termination-watcher/src/github-app-client.ts b/lambdas/functions/termination-watcher/src/github-app-client.ts new file mode 100644 index 0000000000..ab10f73643 --- /dev/null +++ b/lambdas/functions/termination-watcher/src/github-app-client.ts @@ -0,0 +1,167 @@ +import { createAppAuth } from '@octokit/auth-app'; +import { Octokit } from '@octokit/rest'; +import { throttling } from '@octokit/plugin-throttling'; +import { request } from '@octokit/request'; +import { createChildLogger } from '@aws-github-runner/aws-powertools-util'; +import { createCommonStorage, type GitHubAppCredential } from '@aws-github-runner/storage-providers'; +import type { EndpointDefaults } from '@octokit/types'; + +const logger = createChildLogger('github-app-client'); + +let appCredentialsPromise: Promise | undefined; + +const appBudgets = new Map(); + +function coolDown(appId?: number): void { + if (appId !== undefined) appBudgets.set(appId, { remaining: 0, cooldownUntil: Date.now() + 60000 }); +} + +function selectCredential(credentials: GitHubAppCredential[]): GitHubAppCredential { + const offset = Math.floor(Math.random() * credentials.length); + const rotated = [...credentials.slice(offset), ...credentials.slice(0, offset)]; + const available = rotated.filter( + (credential) => (appBudgets.get(credential.appId)?.cooldownUntil ?? 0) <= Date.now(), + ); + return (available.length ? available : rotated).reduce((best, credential) => + (appBudgets.get(credential.appId)?.remaining ?? Infinity) > (appBudgets.get(best.appId)?.remaining ?? Infinity) + ? credential + : best, + ); +} + +export function createThrottleOptions(appId?: number) { + return { + onRateLimit: (_retryAfter: number, options: Required) => { + coolDown(appId); + logger.warn(`Rate limit hit for ${options.method} ${options.url}`); + return false; + }, + onSecondaryRateLimit: (_retryAfter: number, options: Required) => { + coolDown(appId); + logger.warn(`Secondary rate limit hit for ${options.method} ${options.url}`); + return false; + }, + }; +} + +async function loadAppCredentials(): Promise { + const credentials = await createCommonStorage().githubAppCredentials.get(); + if (credentials.length === 0) { + throw new Error('No GitHub App credentials found'); + } + return credentials; +} + +function getAppCredentials(): Promise { + if (!appCredentialsPromise) { + appCredentialsPromise = loadAppCredentials().catch((error: unknown) => { + appCredentialsPromise = undefined; + throw error; + }); + } + return appCredentialsPromise; +} + +export function resetAppCredentialsCache(): void { + appCredentialsPromise = undefined; + appBudgets.clear(); +} + +function createOctokitInstance(token: string, ghesApiUrl: string, appId?: number): Octokit { + const CustomOctokit = Octokit.plugin(throttling); + const octokitOptions: ConstructorParameters[0] = { + auth: token, + }; + if (ghesApiUrl) { + octokitOptions.baseUrl = ghesApiUrl; + } + const client = new CustomOctokit({ + ...octokitOptions, + userAgent: 'github-aws-runners-termination-watcher', + throttle: createThrottleOptions(appId), + }); + if (appId !== undefined) + client.hook.after('request', (response) => { + const remaining = Number.parseInt(String(response.headers['x-ratelimit-remaining']), 10); + if (Number.isFinite(remaining)) + appBudgets.set(appId, { + remaining, + cooldownUntil: remaining === 0 ? Date.now() + 60000 : 0, + }); + }); + return client; +} + +async function createAuthenticatedClient(ghesApiUrl: string, credential: GitHubAppCredential): Promise { + const { appId, privateKey } = credential; + const authOptions: { appId: number; privateKey: string; request?: typeof request } = { + appId, + privateKey, + }; + if (ghesApiUrl) { + authOptions.request = request.defaults({ baseUrl: ghesApiUrl }); + } + const auth = createAppAuth(authOptions); + const appAuth = await auth({ type: 'app' }); + return createOctokitInstance(appAuth.token, ghesApiUrl); +} + +async function getInstallationId(octokit: Octokit, owner: string): Promise { + const { data: installation } = await octokit.apps.getOrgInstallation({ org: owner }); + return installation.id; +} + +async function getInstallationIdForRepo(octokit: Octokit, owner: string, repo: string): Promise { + const { data: installation } = await octokit.apps.getRepoInstallation({ owner, repo }); + return installation.id; +} + +async function createInstallationClient( + appOctokit: Octokit, + owner: string, + runnerType: string, + ghesApiUrl: string, + credential: GitHubAppCredential, +): Promise { + let installationId: number; + if (runnerType === 'Repo') { + const [repoOwner, repo] = owner.split('/'); + installationId = await getInstallationIdForRepo(appOctokit, repoOwner, repo); + } else { + installationId = await getInstallationId(appOctokit, owner); + } + + const { appId, privateKey } = credential; + const authOptions: { appId: number; privateKey: string; installationId: number; request?: typeof request } = { + appId, + privateKey, + installationId, + }; + if (ghesApiUrl) { + authOptions.request = request.defaults({ baseUrl: ghesApiUrl }); + } + const auth = createAppAuth(authOptions); + const installationAuth = await auth({ type: 'installation' }); + return createOctokitInstance(installationAuth.token, ghesApiUrl, appId); +} + +export async function createRunnerInstallationClient( + owner: string, + runnerType: string, + ghesApiUrl: string, +): Promise { + const remaining = [...(await getAppCredentials())]; + while (remaining.length) { + const credential = selectCredential(remaining); + remaining.splice(remaining.indexOf(credential), 1); + try { + const appClient = await createAuthenticatedClient(ghesApiUrl, credential); + return await createInstallationClient(appClient, owner, runnerType, ghesApiUrl, credential); + } catch (error) { + coolDown(credential.appId); + if (!remaining.length) throw error; + logger.warn('GitHub App authentication failed; trying another configured app', { appId: credential.appId }); + } + } + throw new Error('No GitHub App credentials found'); +} diff --git a/lambdas/functions/termination-watcher/src/lambda.ts b/lambdas/functions/termination-watcher/src/lambda.ts index cb493d12a4..7d68746a59 100644 --- a/lambdas/functions/termination-watcher/src/lambda.ts +++ b/lambdas/functions/termination-watcher/src/lambda.ts @@ -75,3 +75,5 @@ const addMiddleware = () => { }; addMiddleware(); + +export { registrationJanitor } from './registration-janitor'; diff --git a/lambdas/functions/termination-watcher/src/registration-janitor.test.ts b/lambdas/functions/termination-watcher/src/registration-janitor.test.ts new file mode 100644 index 0000000000..aaaa33c4c0 --- /dev/null +++ b/lambdas/functions/termination-watcher/src/registration-janitor.test.ts @@ -0,0 +1,434 @@ +import { EC2Client, DescribeInstancesCommand } from '@aws-sdk/client-ec2'; +import { SQSClient, SendMessageCommand } from '@aws-sdk/client-sqs'; +import { mockClient } from 'aws-sdk-client-mock'; +import 'aws-sdk-client-mock-jest/vitest'; +import { beforeEach, afterEach, describe, expect, it, vi } from 'vitest'; +import type { Context, SQSEvent } from 'aws-lambda'; + +import { registrationJanitor, type RegistrationJanitorConfig } from './registration-janitor'; + +const github = vi.hoisted(() => ({ + request: vi.fn(), + actions: { getSelfHostedRunnerForOrg: vi.fn(), deleteSelfHostedRunnerFromOrg: vi.fn() }, +})); +vi.mock('./github-app-client', () => ({ createRunnerInstallationClient: vi.fn().mockResolvedValue(github) })); +const customProvider = vi.hoisted(() => ({ + scope: 'test-scope', + resourceIdFromRunnerName: vi.fn(), + exists: vi.fn(), +})); +vi.mock('@aws-github-runner/compute-providers/registration-cleanup', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + createRegistrationCleanupProvider: (config: { type: string; options: unknown }, prefix: string) => + config.type === 'test-provider' ? customProvider : actual.createRegistrationCleanupProvider(config, prefix), + }; +}); +const ec2 = mockClient(EC2Client); +const sqs = mockClient(SQSClient); +const context = { getRemainingTimeInMillis: () => 60000 } as Context; +const instanceId = 'i-0123456789abcdef0'; +const runner = { id: 10, name: `account-prod_${instanceId}`, status: 'offline', busy: false }; +let config: RegistrationJanitorConfig; + +function configure(overrides: Partial = {}) { + config = { ...config, ...overrides }; + process.env.REGISTRATION_JANITOR_CONFIG = JSON.stringify(config); +} + +async function confirmationEvent(): Promise { + await registrationJanitor({}, context); + const body = sqs.commandCalls(SendMessageCommand)[0].args[0].input.MessageBody!; + vi.setSystemTime(Date.now() + 900000); + return { Records: [{ body, messageId: 'candidate-1' }] } as SQSEvent; +} + +beforeEach(() => { + vi.useFakeTimers({ toFake: ['Date'] }); + vi.setSystemTime(new Date('2026-09-23T00:00:00Z')); + vi.clearAllMocks(); + ec2.reset(); + sqs.reset(); + ec2.on(DescribeInstancesCommand).resolves({ Reservations: [] }); + sqs.on(SendMessageCommand).resolves({}); + github.request.mockImplementation(async () => ({ data: { runners: [runner] } })); + github.actions.getSelfHostedRunnerForOrg.mockResolvedValue({ data: runner }); + github.actions.deleteSelfHostedRunnerFromOrg.mockResolvedValue({ status: 204 }); + config = { + organization: 'example', + runnerGroupIds: [1], + runnerNamePrefix: 'account-prod_', + computeProvider: { type: 'ec2', options: { regions: ['eu-west-1', 'us-east-1'] } }, + dryRun: false, + maxCandidates: 100, + ghesApiUrl: '', + }; + configure(); + process.env.REGISTRATION_JANITOR_QUEUE_URL = 'https://sqs.eu-west-1.amazonaws.com/123456789012/confirmation'; +}); +afterEach(() => { + vi.useRealTimers(); + vi.unstubAllEnvs(); +}); + +describe('registration janitor discovery', () => { + it('discovers remaining inventory after earlier candidates have been removed', async () => { + configure({ maxCandidates: 1 }); + const second = { ...runner, id: 20, name: 'account-prod_i-11111111111111111' }; + github.request.mockResolvedValue({ data: { runners: [runner, second] } }); + const event = await confirmationEvent(); + await registrationJanitor(event, context); + expect(github.actions.deleteSelfHostedRunnerFromOrg).toHaveBeenCalledWith({ org: 'example', runner_id: 10 }); + github.request.mockClear(); + github.request.mockResolvedValue({ data: { runners: [second] } }); + sqs.resetHistory(); + await registrationJanitor({}, context); + expect(github.request).toHaveBeenCalledWith(expect.anything(), expect.objectContaining({ page: 1 })); + expect(JSON.parse(sqs.commandCalls(SendMessageCommand)[0].args[0].input.MessageBody!).runnerId).toBe(20); + }); + + it('continues with another candidate when one EC2 lookup fails', async () => { + github.request.mockResolvedValue({ + data: { runners: [runner, { ...runner, id: 20, name: 'account-prod_i-11111111111111111' }] }, + }); + ec2.on(DescribeInstancesCommand).rejectsOnce(new Error('EC2 unavailable')).resolves({ Reservations: [] }); + await registrationJanitor({}, context); + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + expect(JSON.parse(sqs.commandCalls(SendMessageCommand)[0].args[0].input.MessageBody!).runnerId).toBe(20); + }); + + it('queues early-page candidates before a later GitHub page fails', async () => { + const page = [ + runner, + ...Array.from({ length: 99 }, (_, index) => ({ ...runner, id: index + 100, status: 'online' })), + ]; + github.request + .mockResolvedValueOnce({ data: { runners: page } }) + .mockImplementationOnce(() => { + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + throw new Error('later page unavailable'); + }) + .mockResolvedValue({ data: { runners: [] } }); + await registrationJanitor({}, context); + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + expect(github.request).toHaveBeenCalledWith(expect.anything(), expect.objectContaining({ page: 3 })); + }); + + it('queues a delayed check only for scoped offline, non-busy EC2 runners', async () => { + github.request.mockImplementation(async () => ({ + data: { + runners: [ + runner, + { ...runner, id: 11, status: 'online' }, + { ...runner, id: 12, busy: true }, + { ...runner, id: 13, name: `other_${instanceId}` }, + { ...runner, id: 14, name: 'account-prod_i-123' }, + { ...runner, id: 15, name: `account-prod_nested_${instanceId}` }, + ], + }, + })); + await registrationJanitor({}, context); + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + expect(sqs).toHaveReceivedCommandWith(SendMessageCommand, { DelaySeconds: 900 }); + expect(ec2).toHaveReceivedCommandTimes(DescribeInstancesCommand, 2); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it.each(['pending', 'running', 'stopped', 'stopping', 'shutting-down', undefined])( + 'preserves %s instances in another configured region', + async (state) => { + ec2 + .on(DescribeInstancesCommand) + .resolvesOnce({ Reservations: [] }) + .resolves({ + Reservations: [ + { Instances: [{ InstanceId: instanceId, State: state ? { Name: state as 'running' } : undefined }] }, + ], + }); + await registrationJanitor({}, context); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + }, + ); + + it('protects instances on later EC2 pages', async () => { + ec2 + .on(DescribeInstancesCommand) + .resolvesOnce({ NextToken: 'next' }) + .resolves({ + Reservations: [{ Instances: [{ InstanceId: instanceId, State: { Name: 'running' } }] }], + }); + await registrationJanitor({}, context); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + expect(ec2).toHaveReceivedCommandWith(DescribeInstancesCommand, { NextToken: 'next' }); + }); + + it('allows terminated instances to become confirmation candidates', async () => { + ec2.on(DescribeInstancesCommand).resolves({ + Reservations: [{ Instances: [{ InstanceId: instanceId, State: { Name: 'terminated' } }] }], + }); + await registrationJanitor({}, context); + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + }); + + it('enqueues nothing if any region cannot be checked', async () => { + ec2.on(DescribeInstancesCommand).resolvesOnce({ Reservations: [] }).rejects(new Error('Access denied')); + await registrationJanitor({}, context); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + }); + + it('enqueues nothing if group listing fails', async () => { + github.request.mockRejectedValue(new Error('GitHub unavailable')); + await registrationJanitor({}, context); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + }); + + it('does not mutate in dry-run mode', async () => { + configure({ dryRun: true }); + await registrationJanitor({}, context); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it('shares the candidate budget across groups', async () => { + configure({ runnerGroupIds: [1, 2], maxCandidates: 1 }); + github.request.mockImplementation(async () => ({ + data: { runners: [runner, { ...runner, id: 20, name: 'account-prod_i-11111111111111111' }] }, + })); + await registrationJanitor({}, context); + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + }); + + it('does not spend the candidate limit on retained instances or failed lookups', async () => { + configure({ maxCandidates: 1 }); + github.request.mockResolvedValue({ + data: { + runners: [ + runner, + { ...runner, id: 20, name: 'account-prod_i-11111111111111111' }, + { ...runner, id: 30, name: 'account-prod_i-22222222222222222' }, + ], + }, + }); + ec2 + .on(DescribeInstancesCommand) + .resolvesOnce({ Reservations: [{ Instances: [{ InstanceId: instanceId, State: { Name: 'running' } }] }] }) + .resolvesOnce({ Reservations: [] }) + .rejectsOnce(new Error('unavailable')) + .resolves({ Reservations: [] }); + await registrationJanitor({}, context); + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + expect(JSON.parse(sqs.commandCalls(SendMessageCommand)[0].args[0].input.MessageBody!).runnerId).toBe(30); + }); + + it('requires the confirmation queue before sending candidates', async () => { + delete process.env.REGISTRATION_JANITOR_QUEUE_URL; + await registrationJanitor({}, context); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + }); + + it('handles empty EC2 pages and reservations', async () => { + ec2 + .on(DescribeInstancesCommand) + .resolvesOnce({}) + .resolves({ Reservations: [{}] }); + await registrationJanitor({}, context); + expect(sqs).toHaveReceivedCommandTimes(SendMessageCommand, 1); + }); + + it('rejects absent configuration', async () => { + delete process.env.REGISTRATION_JANITOR_CONFIG; + await expect(registrationJanitor({}, context)).rejects.toThrow('Invalid REGISTRATION_JANITOR_CONFIG'); + }); + + it('stops queueing before the Lambda deadline', async () => { + await registrationJanitor({}, { getRemainingTimeInMillis: () => 1000 } as Context); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + }); + + it.each([{ runnerNamePrefix: '' }, { runnerGroupIds: [] }, { dryRun: undefined }, { maxCandidates: 0 }])( + 'rejects unsafe config %s', + async (invalid) => { + configure(invalid); + await expect(registrationJanitor({}, context)).rejects.toThrow('Invalid REGISTRATION_JANITOR_CONFIG'); + }, + ); +}); + +describe('registration janitor confirmation', () => { + it('continues a confirmation batch after one record fails and lists each group once', async () => { + const event = await confirmationEvent(); + event.Records.unshift({ body: 'invalid JSON', messageId: 'invalid' } as SQSEvent['Records'][number]); + github.request.mockClear(); + expect(await registrationJanitor(event, context)).toEqual({ batchItemFailures: [{ itemIdentifier: 'invalid' }] }); + expect(github.actions.deleteSelfHostedRunnerFromOrg).toHaveBeenCalledWith({ org: 'example', runner_id: 10 }); + expect(github.request).toHaveBeenCalledTimes(1); + }); + + it('checks membership page by page without storing a continuation', async () => { + const event = await confirmationEvent(); + sqs.resetHistory(); + github.request + .mockResolvedValueOnce({ + data: { runners: Array.from({ length: 100 }, (_, index) => ({ ...runner, id: 100 + index })) }, + }) + .mockResolvedValue({ data: { runners: [runner] } }); + await registrationJanitor(event, context); + expect(github.request).toHaveBeenLastCalledWith(expect.anything(), expect.objectContaining({ page: 2 })); + expect(github.actions.deleteSelfHostedRunnerFromOrg).toHaveBeenCalledWith({ org: 'example', runner_id: 10 }); + expect(sqs).not.toHaveReceivedCommand(SendMessageCommand); + }); + + it('preserves a candidate when the membership check reaches the deadline', async () => { + const event = await confirmationEvent(); + expect(await registrationJanitor(event, { getRemainingTimeInMillis: () => 1000 } as Context)).toEqual({ + batchItemFailures: [{ itemIdentifier: 'candidate-1' }], + }); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it('deletes only after delayed EC2 and GitHub rechecks', async () => { + const event = await confirmationEvent(); + ec2.resetHistory(); + expect(await registrationJanitor(event, context)).toEqual({ batchItemFailures: [] }); + expect(ec2).toHaveReceivedCommandTimes(DescribeInstancesCommand, 2); + expect(github.actions.getSelfHostedRunnerForOrg).toHaveBeenCalledWith({ org: 'example', runner_id: 10 }); + expect(github.actions.deleteSelfHostedRunnerFromOrg).toHaveBeenCalledWith({ org: 'example', runner_id: 10 }); + }); + + it.each([{ status: 'online' }, { busy: true }, { name: 'renamed' }])( + 'preserves registrations whose state changed: %s', + async (change) => { + const event = await confirmationEvent(); + github.actions.getSelfHostedRunnerForOrg.mockResolvedValue({ data: { ...runner, ...change } }); + await registrationJanitor(event, context); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }, + ); + + it('preserves instances that appear after discovery', async () => { + const event = await confirmationEvent(); + ec2 + .on(DescribeInstancesCommand) + .resolves({ Reservations: [{ Instances: [{ InstanceId: instanceId, State: { Name: 'stopped' } }] }] }); + await registrationJanitor(event, context); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it('preserves runners moved out of the configured groups', async () => { + const event = await confirmationEvent(); + github.request.mockImplementation(async () => ({ data: { runners: [] } })); + await registrationJanitor(event, context); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it('discards messages after scope configuration changes', async () => { + const event = await confirmationEvent(); + configure({ computeProvider: { type: 'ec2', options: { regions: ['eu-west-1'] } } }); + await registrationJanitor(event, context); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it('discards messages older than one day', async () => { + const event = await confirmationEvent(); + vi.setSystemTime(Date.now() + 86400000); + await registrationJanitor(event, context); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it('does not delete if dry-run is enabled while messages are pending', async () => { + const event = await confirmationEvent(); + configure({ dryRun: true }); + await registrationJanitor(event, context); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it('retries confirmation delivered before the observation delay', async () => { + const event = await confirmationEvent(); + vi.setSystemTime(Date.now() - 1000); + expect(await registrationJanitor(event, context)).toEqual({ + batchItemFailures: [{ itemIdentifier: 'candidate-1' }], + }); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); + + it.each(['ec2', 'github', 'delete'])( + 'retries %s errors without treating unavailable state as absence', + async (failure) => { + const event = await confirmationEvent(); + if (failure === 'ec2') ec2.on(DescribeInstancesCommand).rejects(new Error('EC2 unavailable')); + if (failure === 'github') + github.actions.getSelfHostedRunnerForOrg.mockRejectedValue(new Error('GitHub unavailable')); + if (failure === 'delete') + github.actions.deleteSelfHostedRunnerFromOrg.mockRejectedValue(new Error('GitHub unavailable')); + expect(await registrationJanitor(event, context)).toEqual({ + batchItemFailures: [{ itemIdentifier: 'candidate-1' }], + }); + if (failure !== 'delete') expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }, + ); + + it.each(['get', 'delete'])('treats a %s 404 as already removed', async (operation) => { + const event = await confirmationEvent(); + const method = + operation === 'get' ? github.actions.getSelfHostedRunnerForOrg : github.actions.deleteSelfHostedRunnerFromOrg; + method.mockRejectedValue({ status: 404 }); + expect(await registrationJanitor(event, context)).toEqual({ batchItemFailures: [] }); + }); +}); + +describe('provider-independent registration cleanup', () => { + it('discovers and confirms a non-EC2 resource without AWS inventory calls', async () => { + configure({ computeProvider: { type: 'test-provider', options: {} } }); + const customRunner = { ...runner, name: 'account-prod_container-42' }; + github.request.mockResolvedValue({ data: { runners: [customRunner] } }); + github.actions.getSelfHostedRunnerForOrg.mockResolvedValue({ data: customRunner }); + customProvider.resourceIdFromRunnerName.mockReturnValue('container-42'); + customProvider.exists.mockResolvedValue(false); + const event = await confirmationEvent(); + expect(JSON.parse(event.Records[0].body).resourceId).toBe('container-42'); + await registrationJanitor(event, context); + expect(customProvider.exists).toHaveBeenNthCalledWith(1, 'container-42'); + expect(customProvider.exists).toHaveBeenNthCalledWith(2, 'container-42'); + expect(ec2).not.toHaveReceivedCommand(DescribeInstancesCommand); + expect(github.actions.deleteSelfHostedRunnerFromOrg).toHaveBeenCalledWith({ org: 'example', runner_id: 10 }); + }); + + it('retains a candidate when the provider cannot establish absence at confirmation', async () => { + configure({ computeProvider: { type: 'test-provider', options: {} } }); + customProvider.resourceIdFromRunnerName.mockReturnValue('resource-42'); + customProvider.exists.mockResolvedValueOnce(false).mockRejectedValueOnce(new Error('inventory unavailable')); + const event = await confirmationEvent(); + expect(await registrationJanitor(event, context)).toEqual({ + batchItemFailures: [{ itemIdentifier: 'candidate-1' }], + }); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }); +}); + +describe('confirmation deadline guards', () => { + it.each([ + { stage: 'authentication', checks: 0, membership: 0, inventory: 0, recheck: 0 }, + { stage: 'inventory after membership', checks: 2, membership: 1, inventory: 0, recheck: 0 }, + { stage: 'GitHub recheck after inventory', checks: 3, membership: 1, inventory: 2, recheck: 0 }, + { stage: 'deletion after GitHub recheck', checks: 4, membership: 1, inventory: 2, recheck: 1 }, + ])( + 'defers work before $stage and leaves later records retryable', + async ({ checks, membership, inventory, recheck }) => { + const event = await confirmationEvent(); + event.Records.push({ ...event.Records[0], messageId: 'candidate-2' }); + ec2.resetHistory(); + github.request.mockClear(); + github.actions.getSelfHostedRunnerForOrg.mockClear(); + const remaining = vi.fn().mockReturnValue(9000); + for (let i = 0; i < checks; i++) remaining.mockReturnValueOnce(60000); + expect(await registrationJanitor(event, { getRemainingTimeInMillis: remaining } as unknown as Context)).toEqual({ + batchItemFailures: [{ itemIdentifier: 'candidate-1' }, { itemIdentifier: 'candidate-2' }], + }); + expect(github.request).toHaveBeenCalledTimes(membership); + expect(ec2.commandCalls(DescribeInstancesCommand)).toHaveLength(inventory); + expect(github.actions.getSelfHostedRunnerForOrg).toHaveBeenCalledTimes(recheck); + expect(github.actions.deleteSelfHostedRunnerFromOrg).not.toHaveBeenCalled(); + }, + ); +}); diff --git a/lambdas/functions/termination-watcher/src/registration-janitor.ts b/lambdas/functions/termination-watcher/src/registration-janitor.ts new file mode 100644 index 0000000000..ca9a761ab3 --- /dev/null +++ b/lambdas/functions/termination-watcher/src/registration-janitor.ts @@ -0,0 +1,230 @@ +import { createHash } from 'node:crypto'; +import { + createRegistrationCleanupProvider, + type RegistrationCleanupProvider, + type RegistrationCleanupProviderConfig, +} from '@aws-github-runner/compute-providers/registration-cleanup'; +import { SQSClient, SendMessageCommand } from '@aws-sdk/client-sqs'; +import { createChildLogger } from '@aws-github-runner/aws-powertools-util'; +import type { Context, SQSEvent, SQSBatchResponse } from 'aws-lambda'; +import type { Octokit } from '@octokit/rest'; + +import { createRunnerInstallationClient } from './github-app-client'; + +const logger = createChildLogger('registration-janitor'); +const confirmationSeconds = 900; + +export interface RegistrationJanitorConfig { + organization: string; + runnerGroupIds: number[]; + runnerNamePrefix: string; + computeProvider: RegistrationCleanupProviderConfig; + dryRun: boolean; + maxCandidates: number; + ghesApiUrl: string; +} + +interface Candidate { + runnerId: number; + runnerName: string; + resourceId: string; + observedAt: number; + scope: string; + groupId: number; +} + +function loadConfig(): RegistrationJanitorConfig { + const config = JSON.parse(process.env.REGISTRATION_JANITOR_CONFIG ?? '{}') as RegistrationJanitorConfig; + if ( + !config.organization || + !config.runnerNamePrefix || + !Array.isArray(config.runnerGroupIds) || + config.runnerGroupIds.length === 0 || + !config.runnerGroupIds.every((id) => Number.isSafeInteger(id) && id > 0) || + typeof config.dryRun !== 'boolean' || + !Number.isSafeInteger(config.maxCandidates) || + config.maxCandidates < 1 || + config.maxCandidates > 1000 + ) { + throw new Error('Invalid REGISTRATION_JANITOR_CONFIG'); + } + return { + ...config, + runnerGroupIds: [...config.runnerGroupIds].sort((a, b) => a - b), + }; +} + +function scopeKey(config: RegistrationJanitorConfig, provider: RegistrationCleanupProvider): string { + return createHash('sha256') + .update( + JSON.stringify({ + organization: config.organization, + runnerNamePrefix: config.runnerNamePrefix, + computeProvider: { type: config.computeProvider.type, scope: provider.scope }, + runnerGroupIds: [...config.runnerGroupIds].sort((a, b) => a - b), + ghesApiUrl: config.ghesApiUrl, + }), + ) + .digest('hex'); +} + +async function groupPage(client: Octokit, config: RegistrationJanitorConfig, groupId: number, page: number) { + return ( + await client.request('GET /orgs/{org}/actions/runner-groups/{runner_group_id}/runners', { + org: config.organization, + runner_group_id: groupId, + page, + per_page: 100, + }) + ).data.runners; +} + +async function enqueue(candidate: Candidate, delaySeconds: number): Promise { + const queueUrl = process.env.REGISTRATION_JANITOR_QUEUE_URL; + if (!queueUrl) throw new Error('REGISTRATION_JANITOR_QUEUE_URL is required'); + await new SQSClient({}).send( + new SendMessageCommand({ + QueueUrl: queueUrl, + DelaySeconds: delaySeconds, + MessageBody: JSON.stringify(candidate), + }), + ); +} + +async function discover( + config: RegistrationJanitorConfig, + context: Context, + provider: RegistrationCleanupProvider, +): Promise { + let candidates = 0; + for (const groupId of config.runnerGroupIds) { + let failures = 0; + for (let page = 1; ; page++) { + if (context.getRemainingTimeInMillis() < 10000 || candidates >= config.maxCandidates) return; + let runners; + try { + const client = await createRunnerInstallationClient(config.organization, 'Org', config.ghesApiUrl); + runners = await groupPage(client, config, groupId, page); + failures = 0; + } catch (error) { + logger.warn('Skipping unavailable discovery page', { error, groupId, page }); + // Isolate a broken group while still trying later pages after a transient failure. + if (++failures >= 3) break; + continue; + } + for (const runner of runners) { + if (context.getRemainingTimeInMillis() < 10000 || candidates >= config.maxCandidates) return; + if (runner.status !== 'offline' || runner.busy !== false) continue; + try { + const resourceId = provider.resourceIdFromRunnerName(runner.name); + if (!resourceId) continue; + if (await provider.exists(resourceId)) continue; + const candidate: Candidate = { + runnerId: runner.id, + runnerName: runner.name, + resourceId, + groupId, + observedAt: Date.now(), + scope: scopeKey(config, provider), + }; + if (config.dryRun) logger.info('Would confirm stale registration', { candidate }); + else await enqueue(candidate, confirmationSeconds); + candidates++; + } catch (error) { + logger.warn('Skipping unverified candidate', { error, runnerId: runner.id }); + } + } + if (runners.length < 100) break; + } + } +} + +function isNotFound(error: unknown): boolean { + return typeof error === 'object' && error !== null && 'status' in error && error.status === 404; +} + +function requireConfirmationTime(context: Context): void { + if (context.getRemainingTimeInMillis() < 10000) throw new Error('Insufficient time to confirm registration'); +} + +async function confirm( + client: Octokit, + config: RegistrationJanitorConfig, + candidate: Candidate, + context: Context, + provider: RegistrationCleanupProvider, +): Promise { + const elapsed = Date.now() - candidate.observedAt; + if ( + candidate.scope !== scopeKey(config, provider) || + !config.runnerGroupIds.includes(candidate.groupId) || + typeof candidate.resourceId !== 'string' || + provider.resourceIdFromRunnerName(candidate.runnerName) !== candidate.resourceId || + !Number.isSafeInteger(candidate.runnerId) || + candidate.runnerId < 1 || + !Number.isFinite(elapsed) || + elapsed > 24 * 60 * 60 * 1000 + ) { + logger.warn('Discarding stale or out-of-scope registration candidate'); + return; + } + if (elapsed < confirmationSeconds * 1000) throw new Error('Registration confirmation arrived too early'); + if (config.dryRun) return; + + // A runner moved out of the configured groups is no longer ours to remove. + for (let pageNumber = 1; ; pageNumber++) { + requireConfirmationTime(context); + const members = await groupPage(client, config, candidate.groupId, pageNumber); + if (members.some((runner) => runner.id === candidate.runnerId)) break; + if (members.length < 100) return; + } + requireConfirmationTime(context); + if (await provider.exists(candidate.resourceId)) return; + requireConfirmationTime(context); + let runner; + try { + runner = ( + await client.actions.getSelfHostedRunnerForOrg({ + org: config.organization, + runner_id: candidate.runnerId, + }) + ).data; + } catch (error) { + if (isNotFound(error)) return; + throw error; + } + if (runner.name !== candidate.runnerName || runner.status !== 'offline' || runner.busy !== false) return; + requireConfirmationTime(context); + try { + await client.actions.deleteSelfHostedRunnerFromOrg({ org: config.organization, runner_id: candidate.runnerId }); + } catch (error) { + if (!isNotFound(error)) throw error; + } + logger.info('Removed stale GitHub registration', { runnerId: candidate.runnerId, resourceId: candidate.resourceId }); +} + +export async function registrationJanitor( + event: Partial, + context: Context, +): Promise { + const config = loadConfig(); + const provider = createRegistrationCleanupProvider(config.computeProvider, config.runnerNamePrefix); + if (event.Records) { + const batchItemFailures: SQSBatchResponse['batchItemFailures'] = []; + for (const record of event.Records) { + try { + requireConfirmationTime(context); + const client = await createRunnerInstallationClient(config.organization, 'Org', config.ghesApiUrl); + await confirm(client, config, JSON.parse(record.body) as Candidate, context, provider); + } catch (error) { + logger.warn('Registration confirmation failed; leaving the runner registered', { + error, + messageId: record.messageId, + }); + batchItemFailures.push({ itemIdentifier: record.messageId }); + } + } + return { batchItemFailures }; + } + await discover(config, context, provider); +} diff --git a/lambdas/libs/compute-providers/aws/ec2/registration-cleanup.ts b/lambdas/libs/compute-providers/aws/ec2/registration-cleanup.ts new file mode 100644 index 0000000000..c1d9b16068 --- /dev/null +++ b/lambdas/libs/compute-providers/aws/ec2/registration-cleanup.ts @@ -0,0 +1,45 @@ +import { EC2Client, paginateDescribeInstances } from '@aws-sdk/client-ec2'; +import type { RegistrationCleanupProviderModule } from '../../registration-cleanup'; + +export const provider: RegistrationCleanupProviderModule = { + type: 'ec2', + create(config, runnerNamePrefix) { + const regions = (config as { regions?: unknown } | null)?.regions; + if ( + !Array.isArray(regions) || + regions.length === 0 || + !regions.every((region): region is string => typeof region === 'string' && region.trim().length > 0) + ) { + throw new Error('Invalid EC2 registration cleanup regions'); + } + const sortedRegions = [...new Set(regions)].sort(); + return { + scope: JSON.stringify(sortedRegions), + resourceIdFromRunnerName(name) { + if (!name.startsWith(runnerNamePrefix)) return undefined; + const id = name.slice(runnerNamePrefix.length); + return /^i-(?:[0-9a-f]{8}|[0-9a-f]{17})$/.test(id) ? id : undefined; + }, + async exists(resourceId) { + // No tag filters: stopped or retagged instances still protect registrations. + for (const region of sortedRegions) { + for await (const page of paginateDescribeInstances( + { client: new EC2Client({ region }) }, + { Filters: [{ Name: 'instance-id', Values: [resourceId] }] }, + )) { + for (const reservation of page.Reservations ?? []) { + if ( + reservation.Instances?.some( + (instance) => instance.InstanceId === resourceId && instance.State?.Name !== 'terminated', + ) + ) + return true; + } + } + } + // Absence requires a complete successful lookup in every configured region. + return false; + }, + }; + }, +}; diff --git a/lambdas/libs/compute-providers/package.json b/lambdas/libs/compute-providers/package.json index b62e261495..7780307265 100644 --- a/lambdas/libs/compute-providers/package.json +++ b/lambdas/libs/compute-providers/package.json @@ -13,7 +13,8 @@ "./aws/ec2/control-plane": "./aws/ec2/control-plane.ts", "./aws/ec2/scale-set": "./aws/ec2/scale-set.ts", "./aws/ec2/runners": "./aws/ec2/src/runners.ts", - "./aws/ec2/control-plane/runner-creation": "./aws/ec2/src/control-plane/runner-creation.ts" + "./aws/ec2/control-plane/runner-creation": "./aws/ec2/src/control-plane/runner-creation.ts", + "./registration-cleanup": "./registration-cleanup.ts" }, "type": "module", "license": "MIT", diff --git a/lambdas/libs/compute-providers/providers.config.registration-cleanup.ts b/lambdas/libs/compute-providers/providers.config.registration-cleanup.ts new file mode 100644 index 0000000000..da32e753a7 --- /dev/null +++ b/lambdas/libs/compute-providers/providers.config.registration-cleanup.ts @@ -0,0 +1,7 @@ +import { provider as ec2 } from './aws/ec2/registration-cleanup'; +import type { RegistrationCleanupProviderModule } from './registration-cleanup'; + +/** Provider plugins included in the registration cleanup bundle. */ +export const enabledRegistrationCleanupProviders = [ + ec2, +] as const satisfies readonly RegistrationCleanupProviderModule[]; diff --git a/lambdas/libs/compute-providers/registration-cleanup.test.ts b/lambdas/libs/compute-providers/registration-cleanup.test.ts new file mode 100644 index 0000000000..e8917d0133 --- /dev/null +++ b/lambdas/libs/compute-providers/registration-cleanup.test.ts @@ -0,0 +1,53 @@ +import { beforeEach, describe, expect, it } from 'vitest'; +import { EC2Client, DescribeInstancesCommand } from '@aws-sdk/client-ec2'; +import { mockClient } from 'aws-sdk-client-mock'; +import { createRegistrationCleanupProvider } from './registration-cleanup'; + +const ec2 = mockClient(EC2Client); +const id = 'i-0123456789abcdef0'; +const create = (regions = ['eu-west-1', 'us-east-1']) => + createRegistrationCleanupProvider({ type: 'ec2', options: { regions } }, 'owned_'); +beforeEach(() => ec2.reset()); + +describe('EC2 registration cleanup provider', () => { + it('rejects unknown providers and invalid inventory scope', () => { + expect(() => createRegistrationCleanupProvider({ type: 'missing', options: {} }, 'owned_')).toThrow('Unknown'); + for (const options of [null, {}, { regions: [] }, { regions: [''] }, { regions: [1] }]) { + expect(() => createRegistrationCleanupProvider({ type: 'ec2', options }, 'owned_')).toThrow('Invalid'); + } + }); + it('only claims exact prefixed EC2 names and normalizes scope', () => { + expect(create().resourceIdFromRunnerName(`owned_${id}`)).toBe(id); + expect(create().resourceIdFromRunnerName(`other_${id}`)).toBeUndefined(); + expect(create().resourceIdFromRunnerName('owned_container-42')).toBeUndefined(); + expect(create().scope).toBe(create(['us-east-1', 'eu-west-1', 'eu-west-1']).scope); + }); + it('requires every page and region to establish absence', async () => { + ec2 + .on(DescribeInstancesCommand) + .resolvesOnce({ Reservations: [], NextToken: 'next' }) + .resolvesOnce({ Reservations: [{ Instances: [{ InstanceId: id, State: { Name: 'terminated' } }] }] }) + .resolvesOnce({ Reservations: [] }); + expect(await create().exists(id)).toBe(false); + expect(ec2.commandCalls(DescribeInstancesCommand)).toHaveLength(3); + }); + it.each(['stopped', 'running', undefined])('protects a resource with state %s', async (state) => { + ec2.on(DescribeInstancesCommand).resolves({ + Reservations: [ + { + Instances: [ + { + InstanceId: id, + State: state ? { Name: state as 'running' } : undefined, + }, + ], + }, + ], + }); + expect(await create().exists(id)).toBe(true); + }); + it('does not convert an incomplete lookup to absence', async () => { + ec2.on(DescribeInstancesCommand).resolvesOnce({ Reservations: [] }).rejectsOnce(new Error('unavailable')); + await expect(create().exists(id)).rejects.toThrow('unavailable'); + }); +}); diff --git a/lambdas/libs/compute-providers/registration-cleanup.ts b/lambdas/libs/compute-providers/registration-cleanup.ts new file mode 100644 index 0000000000..bf61274f4d --- /dev/null +++ b/lambdas/libs/compute-providers/registration-cleanup.ts @@ -0,0 +1,28 @@ +import { enabledRegistrationCleanupProviders } from './providers.config.registration-cleanup'; + +/** A failed lookup must throw, never report absence. Resource IDs are provider-owned. */ +export interface RegistrationCleanupProvider { + /** Stable representation of all provider settings affecting ownership and lookup scope. */ + scope: string; + resourceIdFromRunnerName(name: string): string | undefined; + exists(resourceId: string): Promise; +} + +export interface RegistrationCleanupProviderModule { + type: string; + create(config: unknown, runnerNamePrefix: string): RegistrationCleanupProvider; +} + +export interface RegistrationCleanupProviderConfig { + type: string; + options: unknown; +} + +export function createRegistrationCleanupProvider( + config: RegistrationCleanupProviderConfig, + runnerNamePrefix: string, +): RegistrationCleanupProvider { + const provider = enabledRegistrationCleanupProviders.find((provider) => provider.type === config?.type); + if (!provider) throw new Error(`Unknown registration cleanup provider: ${config?.type}`); + return provider.create(config.options, runnerNamePrefix); +} diff --git a/lambdas/yarn.lock b/lambdas/yarn.lock index 372f01fe00..1a86fdbefa 100644 --- a/lambdas/yarn.lock +++ b/lambdas/yarn.lock @@ -245,6 +245,7 @@ __metadata: resolution: "@aws-github-runner/termination-watcher@workspace:functions/termination-watcher" dependencies: "@aws-github-runner/aws-powertools-util": "npm:*" + "@aws-github-runner/compute-providers": "npm:*" "@aws-github-runner/storage-providers": "npm:*" "@aws-sdk/client-ec2": "npm:^3.1009.0" "@aws-sdk/client-sqs": "npm:^3.1009.0" diff --git a/modules/registration-janitor/README.md b/modules/registration-janitor/README.md new file mode 100644 index 0000000000..82c863251b --- /dev/null +++ b/modules/registration-janitor/README.md @@ -0,0 +1,117 @@ +# GitHub registration janitor + +Opt-in reconciliation for offline GitHub organization runner registrations left behind after their EC2 instances disappear. This complements scale-down and termination-event handling: it can discover registrations even when there are no remaining EC2 instances or a termination event was missed. + +The module deploys a separate scheduled Lambda using the existing `termination-watcher.zip` archive. It does not terminate EC2 instances or delete SSM parameters. Existing module users receive no new resources or behavior unless they instantiate this module. + +## Scope and rollout + +Use a nonempty `runner_name_prefix` reserved exclusively for this AWS account in the configured organization and runner groups. Names must equal that prefix followed by an EC2 instance ID. Include **every AWS region** where that prefix is used. Do not use a prefix shared with another AWS account: absence in this account cannot establish that another account's instance is gone. EC2 tags alone cannot prove ownership after an instance disappears. + +Start with the default `dry_run = true` and inspect the candidate logs. Set it to `false` only after verifying the scope. A GitHub App needs organization self-hosted runner read/write permissions. The Lambda supports GHES and optional additional GitHub Apps through the credential manifest described below, including quota-based selection and authentication fallback. + +```hcl +module "registration_janitor" { + source = "github-aws-runners/github-runner/aws//modules/registration-janitor" + + config = { + prefix = "my-account-prod" + organization = "example" + runner_group_ids = [123] + runner_name_prefix = "my-account-prod_" + regions = ["eu-west-1", "us-east-1"] + dry_run = true + + github_app_parameters = { + id = { name = "/runner/app/id", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/runner/app/id" } + key_base64 = { name = "/runner/app/key", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/runner/app/key" } + } + # Set github_app_kms_key_arn when these parameters use a customer-managed key. + # Use an archive from a release containing this handler, or build it locally. + zip = "./termination-watcher.zip" + } +} +``` + +## Additional GitHub Apps + +Set `config.github_app_parameters.additional_apps_manifest` to the manifest's `{ name, arn }` reference and `additional_app_parameter_arns` to every ID, key, and optional installation-ID parameter ARN referenced by that manifest. The credential store uses the existing multi-App manifest format. All Apps must have the permissions and installation access required for the configured organization and runner groups. Include the corresponding KMS decryption access when using customer-managed keys. + +Discovery selects an App for each group page; confirmation selects one for each candidate. Selection prefers the greatest observed remaining quota and skips Apps in a one-minute cooldown after throttling. Unobserved Apps begin with equal priority and use a random starting offset. App JWT and installation authentication use the same credential, with installation lookup scoped to the requested owner. Failed authentication tries another configured App. Rate-limit errors on an authenticated operation remain subject to the existing per-page or per-candidate retry behavior. With no manifest, the primary App remains the only choice. + +## Confirmation and failures + +Every 30 minutes by default, discovery starts a fresh scan of the configured runner groups. It processes each page immediately and checks scoped, offline, non-busy registrations individually against EC2 in every configured region. Any non-terminated instance state protects a registration. Verified candidates are queued before later pages are fetched; one candidate failure does not block others. Unavailable pages are skipped, with advancement to another group after three consecutive page failures. + +Discovery is stateless: no scan cursor or completed-item list is saved. `max_candidates` bounds queued candidates (or reported candidates in dry-run), while retained instances and failed checks do not consume that limit. A deadline guard stops new work with ten seconds remaining. Later invocations list current inventory again; registrations already deleted are no longer listed. Busy, retained, or failed items may be rechecked. Inventory changes can shift pagination boundaries, and subsequent scans reconcile remaining registrations. + +Candidates wait at least 15 minutes in SQS for the existing two-observation safety check. These messages contain candidate identity and observation time, not scan progress. Confirmation searches group membership page by page, stopping when the candidate is found, then rechecks EC2 absence and current GitHub identity/offline/busy state before deletion. Incomplete or failed lookups never establish absence. Scope changes invalidate candidates, and candidates older than 24 hours are discarded. Duplicate candidates are safe because deletion treats 404 as already removed. + +SQS retries failed confirmations independently and routes repeated failures to the dead-letter queue. A verification failure preserves that registration while other records continue. Monitor Lambda errors and the dead-letter queue. + +Dry-run logs candidates and neither enqueues confirmations nor deletes registrations, including confirmations already queued. Disabling the discovery schedule does not stop pending confirmations; set `dry_run = true` to stop deletions as well. + +The delayed observations reduce launch/eventual-consistency races but do not provide an atomic transaction between GitHub and EC2. This first version supports organization runners backed by EC2 in one AWS account, across explicitly configured regions. Repository runners, other compute providers, and cross-account reconciliation are outside its scope. + + +## Requirements + +| Name | Version | +| ---- | ------- | +| [terraform](#requirement\_terraform) | >= 1.5.6 | +| [aws](#requirement\_aws) | >= 6.21 | + +## Providers + +| Name | Version | +| ---- | ------- | +| [aws](#provider\_aws) | 6.60.0 | + +## Modules + +| Name | Source | Version | +| ---- | ------ | ------- | +| [lambda](#module\_lambda) | ../lambda | n/a | + +## Resources + +| Name | Type | +| ---- | ---- | +| [aws_cloudwatch_event_rule.schedule](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/cloudwatch_event_rule) | resource | +| [aws_cloudwatch_event_target.schedule](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/cloudwatch_event_target) | resource | +| [aws_iam_role_policy.cleanup](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/iam_role_policy) | resource | +| [aws_iam_role_policy.compute](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/iam_role_policy) | resource | +| [aws_lambda_event_source_mapping.confirmation](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/lambda_event_source_mapping) | resource | +| [aws_lambda_permission.schedule](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/lambda_permission) | resource | +| [aws_sqs_queue.confirmation](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/sqs_queue) | resource | +| [aws_sqs_queue.dead_letter](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/sqs_queue) | resource | +| [aws_iam_policy_document.cleanup](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/data-sources/iam_policy_document) | data source | +| [aws_iam_policy_document.compute](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/data-sources/iam_policy_document) | data source | + +## Inputs + +| Name | Description | Type | Default | Required | +| ---- | ----------- | ---- | ------- | :------: | +| [config](#input\_config) | Opt-in GitHub registration cleanup for EC2 runners. The name prefix must be exclusive to this AWS account within the configured groups; regions must include every region using it. |
object({
prefix = string
organization = string
runner_group_ids = set(number)
runner_name_prefix = string
regions = set(string)
dry_run = optional(bool, true)
max_candidates = optional(number, 100)
schedule_expression = optional(string, "rate(30 minutes)")
schedule_state = optional(string, "ENABLED")
ghes_api_url = optional(string, "")
github_app_parameters = object({
id = object({ name = string, arn = string })
key_base64 = object({ name = string, arn = string })
additional_apps_manifest = optional(object({ name = string, arn = string }))
additional_app_parameter_arns = optional(list(string), [])
})
github_app_kms_key_arn = optional(string)
zip = optional(string)
s3_bucket = optional(string)
s3_key = optional(string)
s3_object_version = optional(string)
architecture = optional(string, "arm64")
runtime = optional(string, "nodejs24.x")
memory_size = optional(number, 512)
timeout = optional(number, 300)
log_level = optional(string, "info")
logging_retention_in_days = optional(number, 30)
logging_kms_key_id = optional(string)
role_path = optional(string)
role_permissions_boundary = optional(string)
tags = optional(map(string), {})
})
| n/a | yes | + +## Outputs + +| Name | Description | +| ---- | ----------- | +| [dead\_letter\_queue](#output\_dead\_letter\_queue) | Failed confirmation messages. Inspect before redriving; candidates expire after 24 hours. | +| [lambda](#output\_lambda) | Scheduled registration janitor Lambda resources. | + + +### Compute-provider integration + +The scheduled janitor complements termination-watcher deregistration: it discovers registrations left behind by missed termination events, earlier failures, or resources deleted before the watcher was enabled. It does not replace event-driven cleanup or its busy-runner retries. + +GitHub discovery and delayed confirmation use the `RegistrationCleanupProvider` interface in the compute-providers library. Providers interpret runner names, supply a stable inventory scope, and check one backing resource at a time. Lookup failures must throw; only a complete successful lookup may establish absence. The scope includes the provider type and settings, so changing providers or inventory scope invalidates queued observations. + +EC2 is the first bundled implementation, and this Terraform module configures its regions and IAM permissions. Additional providers can register an implementation in `providers.config.registration-cleanup.ts` and supply their configuration and permissions without changing the GitHub cleanup algorithm. No inventory-wide listing or persisted scan progress is required. + +### Deployment boundaries + +Use one janitor deployment per compute provider and ownership scope, with a distinct resource prefix and disjoint runner-name ownership. Multiple deployments can run concurrently for different providers without combining their IAM permissions, queues, or execution budgets. This module currently bundles EC2 only; additional providers require their implementation and Terraform configuration/permissions before deployment. + +`compute-ec2.tf` owns the EC2 provider configuration and region-restricted inventory policy. The shared cleanup policy covers GitHub credential access and confirmation messages only. Both policies use `aws_iam_policy_document` and attach separately to the Lambda role. diff --git a/modules/registration-janitor/compute-ec2.tf b/modules/registration-janitor/compute-ec2.tf new file mode 100644 index 0000000000..f91f480ca9 --- /dev/null +++ b/modules/registration-janitor/compute-ec2.tf @@ -0,0 +1,27 @@ +# One provider per janitor deployment. Keep provider configuration and read-only +# inventory permissions separate from shared GitHub/SQS cleanup infrastructure. +locals { + compute_provider = { + type = "ec2" + options = { regions = var.config.regions } + } +} + +data "aws_iam_policy_document" "compute" { + statement { + sid = "ReadEc2Inventory" + actions = ["ec2:DescribeInstances"] + resources = ["*"] + condition { + test = "StringEquals" + variable = "aws:RequestedRegion" + values = var.config.regions + } + } +} + +resource "aws_iam_role_policy" "compute" { + name = "registration-janitor-compute" + role = module.lambda.lambda.role.name + policy = data.aws_iam_policy_document.compute.json +} diff --git a/modules/registration-janitor/main.tf b/modules/registration-janitor/main.tf new file mode 100644 index 0000000000..1a16be5913 --- /dev/null +++ b/modules/registration-janitor/main.tf @@ -0,0 +1,116 @@ +locals { + name = "${var.config.prefix}-registration-janitor" +} + +module "lambda" { + source = "../lambda" + lambda = { + prefix = var.config.prefix + name = "registration-janitor" + handler = "index.registrationJanitor" + zip = var.config.s3_bucket == null ? coalesce(var.config.zip, "${path.module}/../../lambdas/functions/termination-watcher/termination-watcher.zip") : null + s3_bucket = var.config.s3_bucket + s3_key = var.config.s3_key + s3_object_version = var.config.s3_object_version + architecture = var.config.architecture + runtime = var.config.runtime + memory_size = var.config.memory_size + timeout = var.config.timeout + log_level = var.config.log_level + logging_retention_in_days = var.config.logging_retention_in_days + logging_kms_key_id = var.config.logging_kms_key_id + role_path = var.config.role_path + role_permissions_boundary = var.config.role_permissions_boundary + tags = var.config.tags + environment_variables = merge({ + PARAMETER_GITHUB_APP_ID_NAME = var.config.github_app_parameters.id.name + PARAMETER_GITHUB_APP_KEY_BASE64_NAME = var.config.github_app_parameters.key_base64.name + REGISTRATION_JANITOR_QUEUE_URL = aws_sqs_queue.confirmation.url + REGISTRATION_JANITOR_CONFIG = jsonencode({ + organization = var.config.organization + runnerGroupIds = var.config.runner_group_ids + runnerNamePrefix = var.config.runner_name_prefix + computeProvider = local.compute_provider + dryRun = var.config.dry_run + maxCandidates = var.config.max_candidates + ghesApiUrl = var.config.ghes_api_url + }) + }, var.config.github_app_parameters.additional_apps_manifest == null ? {} : { + PARAMETER_GITHUB_APPS_MANIFEST_NAME = var.config.github_app_parameters.additional_apps_manifest.name + }) + } +} + +resource "aws_sqs_queue" "dead_letter" { + name = "${local.name}-dlq" + message_retention_seconds = 1209600 + sqs_managed_sse_enabled = true + tags = var.config.tags +} + +resource "aws_sqs_queue" "confirmation" { + name = "${local.name}-confirmation" + delay_seconds = 900 + visibility_timeout_seconds = var.config.timeout * 6 + message_retention_seconds = 86400 + sqs_managed_sse_enabled = true + redrive_policy = jsonencode({ + deadLetterTargetArn = aws_sqs_queue.dead_letter.arn + maxReceiveCount = 3 + }) + tags = var.config.tags +} + +data "aws_iam_policy_document" "cleanup" { + statement { + sid = "ReadGitHubCredentials" + actions = concat(["ssm:GetParameters"], var.config.github_app_parameters.additional_apps_manifest == null ? [] : ["ssm:GetParameter"]) + resources = concat([var.config.github_app_parameters.id.arn, var.config.github_app_parameters.key_base64.arn], var.config.github_app_parameters.additional_apps_manifest == null ? [] : [var.config.github_app_parameters.additional_apps_manifest.arn], var.config.github_app_parameters.additional_app_parameter_arns) + } + statement { + sid = "ConfirmRegistrations" + actions = ["sqs:SendMessage", "sqs:ReceiveMessage", "sqs:DeleteMessage", "sqs:GetQueueAttributes"] + resources = [aws_sqs_queue.confirmation.arn] + } + dynamic "statement" { + for_each = var.config.github_app_kms_key_arn == null ? [] : [var.config.github_app_kms_key_arn] + content { + sid = "DecryptGitHubCredentials" + actions = ["kms:Decrypt"] + resources = [statement.value] + } + } +} + +resource "aws_iam_role_policy" "cleanup" { + name = "registration-janitor" + role = module.lambda.lambda.role.name + policy = data.aws_iam_policy_document.cleanup.json +} + +resource "aws_lambda_event_source_mapping" "confirmation" { + event_source_arn = aws_sqs_queue.confirmation.arn + function_name = module.lambda.lambda.function.arn + batch_size = 10 + function_response_types = ["ReportBatchItemFailures"] + depends_on = [aws_iam_role_policy.cleanup, aws_iam_role_policy.compute] +} + +resource "aws_cloudwatch_event_rule" "schedule" { + name = local.name + schedule_expression = var.config.schedule_expression + state = var.config.schedule_state + tags = var.config.tags +} + +resource "aws_cloudwatch_event_target" "schedule" { + rule = aws_cloudwatch_event_rule.schedule.name + arn = module.lambda.lambda.function.arn +} + +resource "aws_lambda_permission" "schedule" { + action = "lambda:InvokeFunction" + function_name = module.lambda.lambda.function.function_name + principal = "events.amazonaws.com" + source_arn = aws_cloudwatch_event_rule.schedule.arn +} diff --git a/modules/registration-janitor/outputs.tf b/modules/registration-janitor/outputs.tf new file mode 100644 index 0000000000..4900d5676d --- /dev/null +++ b/modules/registration-janitor/outputs.tf @@ -0,0 +1,9 @@ +output "lambda" { + description = "Scheduled registration janitor Lambda resources." + value = module.lambda.lambda +} + +output "dead_letter_queue" { + description = "Failed confirmation messages. Inspect before redriving; candidates expire after 24 hours." + value = aws_sqs_queue.dead_letter +} diff --git a/modules/registration-janitor/tests/registration-janitor.tftest.hcl b/modules/registration-janitor/tests/registration-janitor.tftest.hcl new file mode 100644 index 0000000000..148cfe469b --- /dev/null +++ b/modules/registration-janitor/tests/registration-janitor.tftest.hcl @@ -0,0 +1,175 @@ +mock_provider "aws" { + mock_data "aws_iam_policy_document" { + defaults = { json = "{\"Version\":\"2012-10-17\",\"Statement\":[]}" } + } + mock_resource "aws_iam_role" { + defaults = { arn = "arn:aws:iam::123456789012:role/janitor" } + } + mock_resource "aws_lambda_function" { + defaults = { arn = "arn:aws:lambda:eu-west-1:123456789012:function:janitor" } + } + mock_resource "aws_cloudwatch_event_rule" { + defaults = { arn = "arn:aws:events:eu-west-1:123456789012:rule/janitor" } + } + mock_resource "aws_cloudwatch_log_group" { + defaults = { arn = "arn:aws:logs:eu-west-1:123456789012:log-group:/aws/lambda/janitor" } + } +} + +override_resource { + target = aws_sqs_queue.dead_letter + values = { arn = "arn:aws:sqs:eu-west-1:123456789012:janitor-dlq" } +} + +override_resource { + target = aws_sqs_queue.confirmation + values = { + arn = "arn:aws:sqs:eu-west-1:123456789012:janitor-confirmation" + url = "https://sqs.eu-west-1.amazonaws.com/123456789012/janitor-confirmation" + } +} + +variables { + config = { + prefix = "test" + organization = "example" + runner_group_ids = [123] + runner_name_prefix = "account-prod_" + regions = ["eu-west-1", "us-east-1"] + s3_bucket = "artifacts" + s3_key = "termination-watcher.zip" + github_app_parameters = { + id = { name = "/app/id", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/id" } + key_base64 = { name = "/app/key", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/key" } + } + } +} + +run "safe_defaults_and_delayed_confirmation" { + assert { + condition = aws_iam_role_policy.cleanup.policy == data.aws_iam_policy_document.cleanup.json && aws_iam_role_policy.compute.policy == data.aws_iam_policy_document.compute.json + error_message = "Shared cleanup and compute permissions must attach their respective policy documents." + } + assert { + condition = alltrue([for s in data.aws_iam_policy_document.cleanup.statement : !contains(s.actions, "ec2:DescribeInstances")]) && one(one(data.aws_iam_policy_document.compute.statement).condition).test == "StringEquals" && one(one(data.aws_iam_policy_document.compute.statement).condition).variable == "aws:RequestedRegion" + error_message = "EC2 permissions must stay in the separate provider policy with the region restriction." + } + assert { + condition = jsondecode(output.lambda.function.environment[0].variables["REGISTRATION_JANITOR_CONFIG"]).computeProvider == { + type = "ec2" + options = { regions = ["eu-west-1", "us-east-1"] } + } + error_message = "The janitor must receive the EC2 provider and its complete region scope." + } + + assert { + condition = !contains(keys(output.lambda.function.environment[0].variables), "PARAMETER_GITHUB_APPS_MANIFEST_NAME") + error_message = "Single-App deployments must omit the optional environment key." + } + command = apply + assert { + condition = one([for s in data.aws_iam_policy_document.cleanup.statement : s.actions if s.sid == "ReadGitHubCredentials"]) == toset(["ssm:GetParameters"]) + error_message = "Credential loading uses the batched GetParameters API." + } + assert { + condition = output.lambda.function.handler == "index.registrationJanitor" && jsondecode(output.lambda.function.environment[0].variables["REGISTRATION_JANITOR_CONFIG"]).dryRun + error_message = "The Lambda must invoke the janitor handler in dry-run mode." + } + assert { + condition = var.config.dry_run && var.config.max_candidates == 100 + error_message = "The janitor must start in dry-run mode with bounded candidate batches." + } + assert { + condition = aws_sqs_queue.confirmation.delay_seconds == 900 && aws_sqs_queue.confirmation.visibility_timeout_seconds == 1800 + error_message = "Confirmations must be delayed 15 minutes and remain invisible during retries." + } + assert { + condition = aws_lambda_event_source_mapping.confirmation.function_response_types == toset(["ReportBatchItemFailures"]) + error_message = "Only failed confirmations should be retried." + } + assert { + condition = one(data.aws_iam_policy_document.compute.statement).actions == toset(["ec2:DescribeInstances"]) + error_message = "The janitor must not be able to terminate instances." + } + assert { + condition = toset(one(one(data.aws_iam_policy_document.compute.statement).condition).values) == toset(["eu-west-1", "us-east-1"]) + error_message = "EC2 access must be limited to the configured regions." + } + assert { + condition = length(data.aws_iam_policy_document.cleanup.statement) == 2 + error_message = "KMS access should be absent unless a customer-managed key is configured." + } +} + +run "reject_empty_scope" { + command = plan + variables { + config = { + prefix = "test" + organization = "example" + runner_group_ids = [123] + runner_name_prefix = "" + regions = ["eu-west-1"] + s3_bucket = "artifacts" + s3_key = "termination-watcher.zip" + github_app_parameters = { + id = { name = "/app/id", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/id" } + key_base64 = { name = "/app/key", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/key" } + } + } + } + expect_failures = [var.config] +} + +run "customer_managed_key" { + command = apply + variables { + config = { + github_app_kms_key_arn = "arn:aws:kms:eu-west-1:123456789012:key/11111111-1111-1111-1111-111111111111" + prefix = "test" + organization = "example" + runner_group_ids = [123] + runner_name_prefix = "account-prod_" + regions = ["eu-west-1"] + s3_bucket = "artifacts" + s3_key = "termination-watcher.zip" + github_app_parameters = { + id = { name = "/app/id", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/id" } + key_base64 = { name = "/app/key", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/key" } + } + } + } + assert { + condition = one([for s in data.aws_iam_policy_document.cleanup.statement : s.resources if s.sid == "DecryptGitHubCredentials"]) == toset([var.config.github_app_kms_key_arn]) + error_message = "App credential decryption must be limited to the configured KMS key." + } +} + +run "additional_apps_manifest" { + command = apply + variables { + config = { + prefix = "test" + organization = "example" + runner_group_ids = [123] + runner_name_prefix = "account-prod_" + regions = ["eu-west-1"] + s3_bucket = "artifacts" + s3_key = "termination-watcher.zip" + github_app_parameters = { + id = { name = "/app/id", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/id" } + key_base64 = { name = "/app/key", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/key" } + additional_apps_manifest = { name = "/app/manifest", arn = "arn:aws:ssm:eu-west-1:123456789012:parameter/app/manifest" } + additional_app_parameter_arns = ["arn:aws:ssm:eu-west-1:123456789012:parameter/app/extra/key"] + } + } + } + assert { + condition = output.lambda.function.environment[0].variables["PARAMETER_GITHUB_APPS_MANIFEST_NAME"] == "/app/manifest" + error_message = "The Lambda must receive the additional Apps manifest." + } + assert { + condition = one([for s in data.aws_iam_policy_document.cleanup.statement : s.actions if s.sid == "ReadGitHubCredentials"]) == toset(["ssm:GetParameter", "ssm:GetParameters"]) && one([for s in data.aws_iam_policy_document.cleanup.statement : s.resources if s.sid == "ReadGitHubCredentials"]) == toset([var.config.github_app_parameters.id.arn, var.config.github_app_parameters.key_base64.arn, var.config.github_app_parameters.additional_apps_manifest.arn, "arn:aws:ssm:eu-west-1:123456789012:parameter/app/extra/key"]) + error_message = "Read access must cover the manifest and only the configured credential parameters." + } +} diff --git a/modules/registration-janitor/variables.tf b/modules/registration-janitor/variables.tf new file mode 100644 index 0000000000..8f94fa5911 --- /dev/null +++ b/modules/registration-janitor/variables.tf @@ -0,0 +1,53 @@ +variable "config" { + description = "Opt-in GitHub registration cleanup for EC2 runners. The name prefix must be exclusive to this AWS account within the configured groups; regions must include every region using it." + type = object({ + prefix = string + organization = string + runner_group_ids = set(number) + runner_name_prefix = string + regions = set(string) + dry_run = optional(bool, true) + max_candidates = optional(number, 100) + schedule_expression = optional(string, "rate(30 minutes)") + schedule_state = optional(string, "ENABLED") + ghes_api_url = optional(string, "") + github_app_parameters = object({ + id = object({ name = string, arn = string }) + key_base64 = object({ name = string, arn = string }) + additional_apps_manifest = optional(object({ name = string, arn = string })) + additional_app_parameter_arns = optional(list(string), []) + }) + github_app_kms_key_arn = optional(string) + zip = optional(string) + s3_bucket = optional(string) + s3_key = optional(string) + s3_object_version = optional(string) + architecture = optional(string, "arm64") + runtime = optional(string, "nodejs24.x") + memory_size = optional(number, 512) + timeout = optional(number, 300) + log_level = optional(string, "info") + logging_retention_in_days = optional(number, 30) + logging_kms_key_id = optional(string) + role_path = optional(string) + role_permissions_boundary = optional(string) + tags = optional(map(string), {}) + }) + + validation { + condition = length(trimspace(var.config.organization)) > 0 && length(trimspace(var.config.runner_name_prefix)) > 0 + error_message = "An organization and a nonempty, account-exclusive runner name prefix are required." + } + validation { + condition = length(var.config.regions) > 0 && alltrue([for region in var.config.regions : length(trimspace(region)) > 0]) + error_message = "Include every AWS region used by runners with this prefix." + } + validation { + condition = length(var.config.runner_group_ids) > 0 && alltrue([for id in var.config.runner_group_ids : id > 0 && floor(id) == id]) + error_message = "At least one positive integer runner group ID is required." + } + validation { + condition = var.config.max_candidates >= 1 && var.config.max_candidates <= 1000 && floor(var.config.max_candidates) == var.config.max_candidates + error_message = "max_candidates must be an integer between 1 and 1000." + } +} diff --git a/modules/registration-janitor/versions.tf b/modules/registration-janitor/versions.tf new file mode 100644 index 0000000000..1238b79cc3 --- /dev/null +++ b/modules/registration-janitor/versions.tf @@ -0,0 +1,10 @@ +terraform { + required_version = ">= 1.5.6" + + required_providers { + aws = { + source = "hashicorp/aws" + version = ">= 6.21" + } + } +} From af6cc4be977becee30e05277c9daad8a4ed963f5 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Wed, 23 Sep 2026 19:08:23 +0000 Subject: [PATCH 2/3] docs: auto update terraform docs --- modules/registration-janitor/README.md | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/modules/registration-janitor/README.md b/modules/registration-janitor/README.md index 82c863251b..28a3089e9c 100644 --- a/modules/registration-janitor/README.md +++ b/modules/registration-janitor/README.md @@ -57,26 +57,26 @@ The delayed observations reduce launch/eventual-consistency races but do not pro ## Requirements | Name | Version | -| ---- | ------- | +|------|---------| | [terraform](#requirement\_terraform) | >= 1.5.6 | | [aws](#requirement\_aws) | >= 6.21 | ## Providers | Name | Version | -| ---- | ------- | -| [aws](#provider\_aws) | 6.60.0 | +|------|---------| +| [aws](#provider\_aws) | >= 6.21 | ## Modules | Name | Source | Version | -| ---- | ------ | ------- | +|------|--------|---------| | [lambda](#module\_lambda) | ../lambda | n/a | ## Resources | Name | Type | -| ---- | ---- | +|------|------| | [aws_cloudwatch_event_rule.schedule](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/cloudwatch_event_rule) | resource | | [aws_cloudwatch_event_target.schedule](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/cloudwatch_event_target) | resource | | [aws_iam_role_policy.cleanup](https://registry.terraform.io/providers/hashicorp/aws/latest/docs/resources/iam_role_policy) | resource | @@ -91,13 +91,13 @@ The delayed observations reduce launch/eventual-consistency races but do not pro ## Inputs | Name | Description | Type | Default | Required | -| ---- | ----------- | ---- | ------- | :------: | +|------|-------------|------|---------|:--------:| | [config](#input\_config) | Opt-in GitHub registration cleanup for EC2 runners. The name prefix must be exclusive to this AWS account within the configured groups; regions must include every region using it. |
object({
prefix = string
organization = string
runner_group_ids = set(number)
runner_name_prefix = string
regions = set(string)
dry_run = optional(bool, true)
max_candidates = optional(number, 100)
schedule_expression = optional(string, "rate(30 minutes)")
schedule_state = optional(string, "ENABLED")
ghes_api_url = optional(string, "")
github_app_parameters = object({
id = object({ name = string, arn = string })
key_base64 = object({ name = string, arn = string })
additional_apps_manifest = optional(object({ name = string, arn = string }))
additional_app_parameter_arns = optional(list(string), [])
})
github_app_kms_key_arn = optional(string)
zip = optional(string)
s3_bucket = optional(string)
s3_key = optional(string)
s3_object_version = optional(string)
architecture = optional(string, "arm64")
runtime = optional(string, "nodejs24.x")
memory_size = optional(number, 512)
timeout = optional(number, 300)
log_level = optional(string, "info")
logging_retention_in_days = optional(number, 30)
logging_kms_key_id = optional(string)
role_path = optional(string)
role_permissions_boundary = optional(string)
tags = optional(map(string), {})
})
| n/a | yes | ## Outputs | Name | Description | -| ---- | ----------- | +|------|-------------| | [dead\_letter\_queue](#output\_dead\_letter\_queue) | Failed confirmation messages. Inspect before redriving; candidates expire after 24 hours. | | [lambda](#output\_lambda) | Scheduled registration janitor Lambda resources. | From 5b3d69b8ee04ce42037e1726fa602bbf8641c0f2 Mon Sep 17 00:00:00 2001 From: Guilherme Caulada Date: Thu, 24 Sep 2026 12:48:42 -0300 Subject: [PATCH 3/3] perf(registration-janitor): reuse clients and report discovery progress --- .../src/github-app-client.test.ts | 47 ++++++++ .../src/github-app-client.ts | 13 +- .../src/registration-janitor.test.ts | 89 ++++++++++++++ .../src/registration-janitor.ts | 114 ++++++++++++------ modules/registration-janitor/README.md | 6 +- 5 files changed, 229 insertions(+), 40 deletions(-) diff --git a/lambdas/functions/termination-watcher/src/github-app-client.test.ts b/lambdas/functions/termination-watcher/src/github-app-client.test.ts index 17d278221e..1f18322f76 100644 --- a/lambdas/functions/termination-watcher/src/github-app-client.test.ts +++ b/lambdas/functions/termination-watcher/src/github-app-client.test.ts @@ -158,4 +158,51 @@ describe('multi-App installation clients', () => { random.mockRestore(); } }); + it('reuses a selected App client within one invocation, while still switching exhausted Apps', async () => { + mockGetCredentials.mockResolvedValue([ + { appId: 1, privateKey: 'one' }, + { appId: 2, privateKey: 'two' }, + ]); + mockApps.getOrgInstallation.mockResolvedValue({ data: { id: 222 } }); + const random = vi.spyOn(Math, 'random').mockReturnValue(0); + const clients = new Map(); + try { + const first = await createRunnerInstallationClient('test-org', 'Org', '', clients); + // Unknown budgets tie; with the same selection, the client is reused. + expect(await createRunnerInstallationClient('test-org', 'Org', '', clients)).toBe(first); + expect(mockCreateAppAuth).toHaveBeenCalledTimes(2); + mockHookAfter.mock.calls[0][1]({ headers: { 'x-ratelimit-remaining': '0' } }); + const second = await createRunnerInstallationClient('test-org', 'Org', '', clients); + expect(await createRunnerInstallationClient('test-org', 'Org', '', clients)).toBe(second); + expect(mockCreateAppAuth).toHaveBeenCalledTimes(4); + expect(mockApps.getOrgInstallation).toHaveBeenCalledTimes(2); + await createRunnerInstallationClient('test-org', 'Org', '', new Map()); + expect(mockCreateAppAuth).toHaveBeenCalledTimes(6); + } finally { + random.mockRestore(); + } + }); + + it('separates cached clients by owner, runner type, and GHES endpoint', async () => { + mockApps.getOrgInstallation.mockResolvedValue({ data: { id: 222 } }); + mockApps.getRepoInstallation.mockResolvedValue({ data: { id: 333 } }); + const clients = new Map(); + await createRunnerInstallationClient('owner', 'Org', '', clients); + await createRunnerInstallationClient('other', 'Org', '', clients); + await createRunnerInstallationClient('owner/repo', 'Repo', '', clients); + await createRunnerInstallationClient('owner', 'Org', 'https://ghe.example/api/v3', clients); + expect(clients.size).toBe(4); + expect(mockCreateAppAuth).toHaveBeenCalledTimes(8); + }); + + it('does not cache failed authentication', async () => { + const clients = new Map(); + mockApps.getOrgInstallation + .mockRejectedValueOnce(new Error('unavailable')) + .mockResolvedValue({ data: { id: 222 } }); + await expect(createRunnerInstallationClient('owner', 'Org', '', clients)).rejects.toThrow('unavailable'); + expect(clients.size).toBe(0); + await createRunnerInstallationClient('owner', 'Org', '', clients); + expect(clients.size).toBe(1); + }); }); diff --git a/lambdas/functions/termination-watcher/src/github-app-client.ts b/lambdas/functions/termination-watcher/src/github-app-client.ts index ab10f73643..86167d395f 100644 --- a/lambdas/functions/termination-watcher/src/github-app-client.ts +++ b/lambdas/functions/termination-watcher/src/github-app-client.ts @@ -145,18 +145,29 @@ async function createInstallationClient( return createOctokitInstance(installationAuth.token, ghesApiUrl, appId); } +/** Kept by the caller for one invocation; never stores scan progress. */ +export type InstallationClientCache = Map; + export async function createRunnerInstallationClient( owner: string, runnerType: string, ghesApiUrl: string, + clients?: InstallationClientCache, ): Promise { const remaining = [...(await getAppCredentials())]; while (remaining.length) { const credential = selectCredential(remaining); remaining.splice(remaining.indexOf(credential), 1); + // Select against current quota before consulting the cache, so a reused + // client never pins a batch to an exhausted or cooling-down App. + const key = JSON.stringify([credential.appId, runnerType, owner, ghesApiUrl]); + const cached = clients?.get(key); + if (cached) return cached; try { const appClient = await createAuthenticatedClient(ghesApiUrl, credential); - return await createInstallationClient(appClient, owner, runnerType, ghesApiUrl, credential); + const client = await createInstallationClient(appClient, owner, runnerType, ghesApiUrl, credential); + clients?.set(key, client); + return client; } catch (error) { coolDown(credential.appId); if (!remaining.length) throw error; diff --git a/lambdas/functions/termination-watcher/src/registration-janitor.test.ts b/lambdas/functions/termination-watcher/src/registration-janitor.test.ts index aaaa33c4c0..d8bbe0d428 100644 --- a/lambdas/functions/termination-watcher/src/registration-janitor.test.ts +++ b/lambdas/functions/termination-watcher/src/registration-janitor.test.ts @@ -6,6 +6,21 @@ import { beforeEach, afterEach, describe, expect, it, vi } from 'vitest'; import type { Context, SQSEvent } from 'aws-lambda'; import { registrationJanitor, type RegistrationJanitorConfig } from './registration-janitor'; +import { createRunnerInstallationClient } from './github-app-client'; + +const janitorLogs = vi.hoisted(() => ({ info: vi.fn(), setContext: vi.fn() })); +vi.mock('@aws-github-runner/aws-powertools-util', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + setContext: janitorLogs.setContext, + createChildLogger: (name: string) => { + const child = actual.createChildLogger(name); + if (name === 'registration-janitor') vi.spyOn(child, 'info').mockImplementation(janitorLogs.info); + return child; + }, + }; +}); const github = vi.hoisted(() => ({ request: vi.fn(), @@ -432,3 +447,77 @@ describe('confirmation deadline guards', () => { }, ); }); + +describe('invocation reuse and discovery observability', () => { + it('shares one client cache across discovery pages and resets it next invocation', async () => { + github.request + .mockResolvedValueOnce({ + data: { runners: Array.from({ length: 100 }, () => ({ ...runner, status: 'online' })) }, + }) + .mockResolvedValue({ data: { runners: [] } }); + await registrationJanitor({}, context); + const calls = vi.mocked(createRunnerInstallationClient).mock.calls; + expect(calls).toHaveLength(2); + expect(calls[0][3]).toBe(calls[1][3]); + await registrationJanitor({}, context); + expect(calls[2][3]).not.toBe(calls[0][3]); + }); + + it('shares one client cache within a confirmation batch and sets fresh invocation context', async () => { + const event = await confirmationEvent(); + event.Records.push({ ...event.Records[0], messageId: 'second' }); + vi.mocked(createRunnerInstallationClient).mockClear(); + const batchContext = { ...context, awsRequestId: 'confirmation-request', functionName: 'janitor' } as Context; + await registrationJanitor(event, batchContext); + const calls = vi.mocked(createRunnerInstallationClient).mock.calls; + expect(calls).toHaveLength(2); + expect(calls[0][3]).toBe(calls[1][3]); + expect(janitorLogs.setContext).toHaveBeenLastCalledWith(batchContext, 'registration-janitor'); + }); + + it('reports queued candidates and resources retained in live discovery', async () => { + github.request.mockResolvedValue({ data: { runners: [runner, { ...runner, id: 11 }] } }); + ec2 + .on(DescribeInstancesCommand) + .resolvesOnce({ Reservations: [] }) + .resolvesOnce({ Reservations: [] }) + .resolves({ Reservations: [{ Instances: [{ InstanceId: instanceId, State: { Name: 'stopped' } }] }] }); + await registrationJanitor({}, context); + expect(janitorLogs.info).toHaveBeenCalledWith( + 'Registration discovery finished.', + expect.objectContaining({ + pagesScanned: 1, + candidatesQueued: 1, + existingResources: 1, + dryRunCandidates: 0, + stopReason: 'completed', + }), + ); + }); + + it('reports dry-run candidates separately from queued work', async () => { + configure({ dryRun: true }); + await registrationJanitor({}, context); + expect(janitorLogs.info).toHaveBeenCalledWith( + 'Registration discovery finished.', + expect.objectContaining({ candidatesQueued: 0, dryRunCandidates: 1, dryRun: true }), + ); + }); + + it('reports failures without claiming failed queue sends succeeded', async () => { + sqs.on(SendMessageCommand).rejects(new Error('queue unavailable')); + await registrationJanitor({}, context); + expect(janitorLogs.info).toHaveBeenCalledWith( + 'Registration discovery finished.', + expect.objectContaining({ candidatesQueued: 0, candidateFailures: 1 }), + ); + }); + + it('reports partial discovery when the deadline stops scanning', async () => { + await registrationJanitor({}, { ...context, getRemainingTimeInMillis: () => 1000 }); + expect(janitorLogs.info).toHaveBeenCalledWith( + 'Registration discovery finished.', + expect.objectContaining({ pagesScanned: 0, candidatesQueued: 0, stopReason: 'deadline' }), + ); + }); +}); diff --git a/lambdas/functions/termination-watcher/src/registration-janitor.ts b/lambdas/functions/termination-watcher/src/registration-janitor.ts index ca9a761ab3..8b4914a847 100644 --- a/lambdas/functions/termination-watcher/src/registration-janitor.ts +++ b/lambdas/functions/termination-watcher/src/registration-janitor.ts @@ -5,11 +5,11 @@ import { type RegistrationCleanupProviderConfig, } from '@aws-github-runner/compute-providers/registration-cleanup'; import { SQSClient, SendMessageCommand } from '@aws-sdk/client-sqs'; -import { createChildLogger } from '@aws-github-runner/aws-powertools-util'; +import { createChildLogger, setContext } from '@aws-github-runner/aws-powertools-util'; import type { Context, SQSEvent, SQSBatchResponse } from 'aws-lambda'; import type { Octokit } from '@octokit/rest'; -import { createRunnerInstallationClient } from './github-app-client'; +import { createRunnerInstallationClient, type InstallationClientCache } from './github-app-client'; const logger = createChildLogger('registration-janitor'); const confirmationSeconds = 900; @@ -95,47 +95,83 @@ async function discover( config: RegistrationJanitorConfig, context: Context, provider: RegistrationCleanupProvider, + clients: InstallationClientCache, ): Promise { let candidates = 0; - for (const groupId of config.runnerGroupIds) { - let failures = 0; - for (let page = 1; ; page++) { - if (context.getRemainingTimeInMillis() < 10000 || candidates >= config.maxCandidates) return; - let runners; - try { - const client = await createRunnerInstallationClient(config.organization, 'Org', config.ghesApiUrl); - runners = await groupPage(client, config, groupId, page); - failures = 0; - } catch (error) { - logger.warn('Skipping unavailable discovery page', { error, groupId, page }); - // Isolate a broken group while still trying later pages after a transient failure. - if (++failures >= 3) break; - continue; - } - for (const runner of runners) { - if (context.getRemainingTimeInMillis() < 10000 || candidates >= config.maxCandidates) return; - if (runner.status !== 'offline' || runner.busy !== false) continue; + const summary = { + pagesScanned: 0, + candidatesQueued: 0, + dryRunCandidates: 0, + existingResources: 0, + pageFailures: 0, + candidateFailures: 0, + stopReason: 'completed', + }; + const shouldStop = () => { + if (context.getRemainingTimeInMillis() < 10000) { + summary.stopReason = 'deadline'; + return true; + } + if (candidates >= config.maxCandidates) { + summary.stopReason = 'candidate-limit'; + return true; + } + return false; + }; + try { + for (const groupId of config.runnerGroupIds) { + let failures = 0; + for (let page = 1; ; page++) { + if (shouldStop()) return; + let runners; try { - const resourceId = provider.resourceIdFromRunnerName(runner.name); - if (!resourceId) continue; - if (await provider.exists(resourceId)) continue; - const candidate: Candidate = { - runnerId: runner.id, - runnerName: runner.name, - resourceId, - groupId, - observedAt: Date.now(), - scope: scopeKey(config, provider), - }; - if (config.dryRun) logger.info('Would confirm stale registration', { candidate }); - else await enqueue(candidate, confirmationSeconds); - candidates++; + const client = await createRunnerInstallationClient(config.organization, 'Org', config.ghesApiUrl, clients); + runners = await groupPage(client, config, groupId, page); + summary.pagesScanned++; + failures = 0; } catch (error) { - logger.warn('Skipping unverified candidate', { error, runnerId: runner.id }); + summary.pageFailures++; + logger.warn('Skipping unavailable discovery page', { error, groupId, page }); + // Isolate a broken group while still trying later pages after a transient failure. + if (++failures >= 3) break; + continue; + } + for (const runner of runners) { + if (shouldStop()) return; + if (runner.status !== 'offline' || runner.busy !== false) continue; + try { + const resourceId = provider.resourceIdFromRunnerName(runner.name); + if (!resourceId) continue; + if (await provider.exists(resourceId)) { + summary.existingResources++; + continue; + } + const candidate: Candidate = { + runnerId: runner.id, + runnerName: runner.name, + resourceId, + groupId, + observedAt: Date.now(), + scope: scopeKey(config, provider), + }; + if (config.dryRun) { + logger.info('Would confirm stale registration', { candidate }); + summary.dryRunCandidates++; + } else { + await enqueue(candidate, confirmationSeconds); + summary.candidatesQueued++; + } + candidates++; + } catch (error) { + summary.candidateFailures++; + logger.warn('Skipping unverified candidate', { error, runnerId: runner.id }); + } } + if (runners.length < 100) break; } - if (runners.length < 100) break; } + } finally { + logger.info('Registration discovery finished.', { ...summary, dryRun: config.dryRun }); } } @@ -207,6 +243,8 @@ export async function registrationJanitor( event: Partial, context: Context, ): Promise { + setContext(context, 'registration-janitor'); + const clients: InstallationClientCache = new Map(); const config = loadConfig(); const provider = createRegistrationCleanupProvider(config.computeProvider, config.runnerNamePrefix); if (event.Records) { @@ -214,7 +252,7 @@ export async function registrationJanitor( for (const record of event.Records) { try { requireConfirmationTime(context); - const client = await createRunnerInstallationClient(config.organization, 'Org', config.ghesApiUrl); + const client = await createRunnerInstallationClient(config.organization, 'Org', config.ghesApiUrl, clients); await confirm(client, config, JSON.parse(record.body) as Candidate, context, provider); } catch (error) { logger.warn('Registration confirmation failed; leaving the runner registered', { @@ -226,5 +264,5 @@ export async function registrationJanitor( } return { batchItemFailures }; } - await discover(config, context, provider); + await discover(config, context, provider, clients); } diff --git a/modules/registration-janitor/README.md b/modules/registration-janitor/README.md index 28a3089e9c..14e5dce886 100644 --- a/modules/registration-janitor/README.md +++ b/modules/registration-janitor/README.md @@ -37,7 +37,7 @@ module "registration_janitor" { Set `config.github_app_parameters.additional_apps_manifest` to the manifest's `{ name, arn }` reference and `additional_app_parameter_arns` to every ID, key, and optional installation-ID parameter ARN referenced by that manifest. The credential store uses the existing multi-App manifest format. All Apps must have the permissions and installation access required for the configured organization and runner groups. Include the corresponding KMS decryption access when using customer-managed keys. -Discovery selects an App for each group page; confirmation selects one for each candidate. Selection prefers the greatest observed remaining quota and skips Apps in a one-minute cooldown after throttling. Unobserved Apps begin with equal priority and use a random starting offset. App JWT and installation authentication use the same credential, with installation lookup scoped to the requested owner. Failed authentication tries another configured App. Rate-limit errors on an authenticated operation remain subject to the existing per-page or per-candidate retry behavior. With no manifest, the primary App remains the only choice. +Discovery selects an App for each group page; confirmation selects one for each candidate. Within an invocation, clients are reused per App, owner, runner type, and GitHub endpoint, avoiding repeated JWT creation, installation lookup, and token minting. Selection still runs before cache lookup so quota and cooldown can move subsequent work to another App. Clients are discarded between invocations; failed authentication is not cached. Selection prefers the greatest observed remaining quota and skips Apps in a one-minute cooldown after throttling. Unobserved Apps begin with equal priority and use a random starting offset. App JWT and installation authentication use the same credential, with installation lookup scoped to the requested owner. Failed authentication tries another configured App. Rate-limit errors on an authenticated operation remain subject to the existing per-page or per-candidate retry behavior. With no manifest, the primary App remains the only choice. ## Confirmation and failures @@ -115,3 +115,7 @@ EC2 is the first bundled implementation, and this Terraform module configures it Use one janitor deployment per compute provider and ownership scope, with a distinct resource prefix and disjoint runner-name ownership. Multiple deployments can run concurrently for different providers without combining their IAM permissions, queues, or execution budgets. This module currently bundles EC2 only; additional providers require their implementation and Terraform configuration/permissions before deployment. `compute-ec2.tf` owns the EC2 provider configuration and region-restricted inventory policy. The shared cleanup policy covers GitHub credential access and confirmation messages only. Both policies use `aws_iam_policy_document` and attach separately to the Lambda role. + +### Discovery logs + +Each discovery invocation emits `Registration discovery finished.` with `pagesScanned`, `candidatesQueued`, `existingResources`, `dryRunCandidates`, `pageFailures`, `candidateFailures`, and `stopReason` (`completed`, `deadline`, or `candidate-limit`). Queue counts increase only after successful sends; dry-run observations are counted separately. Summaries are also emitted for partial scans, and failure counters distinguish skipped work from successful traversal. The handler sets request context for discovery and confirmation logs on every invocation. These counters and the client cache are invocation-local, not persisted cleanup progress.