mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-10 03:08:10 +02:00
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - GitHub connections give agents an account with selected repository access. > - Fresh local instances enroll with production Paperclip Cloud. > - Enrollment could finish while the GitHub OAuth profile remained disabled. > - Setup then switched to a personal access token form without explanation. > - This change preserves sign-in intent and shows the connected account and repositories. ## Linked Issues or Issue Description Related: #12907, #12943, #12947. Existing open GitHub connection work was checked. No duplicate was found. **What happened?** After Cloud enrollment, a fresh test-drive asked for a GitHub key. Production did not advertise the managed GitHub profile. Staging did. The permissions page also omitted the authenticated username and repository names. **Expected behavior** Continue with GitHub OAuth when available. Explain unavailable sign-in and allow retry otherwise. Show the GitHub username and complete accessible repository list. **Steps to reproduce** Start a fresh test-drive. Choose GitHub and complete instance enrollment while the Cloud GitHub profile is disabled. Open an existing GitHub connection's permissions page. ## What Changed - Preserve managed sign-in intent when the gallery omits its profile. - Refresh the selected gallery entry on retry without resetting the audience. - Fetch all pages of GitHub installations and repositories. - Store only repository IDs, full names, and installation IDs in grant metadata. - Show the GitHub username, repository list, management link, and refresh action. - Discard the repository snapshot after newer installation lifecycle events. Preserve snapshots verified after delayed events. - Lock and re-read grant metadata when applying installation events or saving refreshed access. Patch only webhook fields for other events. Reject snapshots if access changed during the external fetch, using unique access revisions even when timestamps collide. - Show repository installation recovery for managed OAuth even when the app also offers an advanced PAT method. - Update tests and the GitHub connection runbook. No SQL migration is required. ## Verification - Local typecheck, build, and token gates passed. All latest-head CI gates passed, including the complete test matrix and browser suites. Greptile is 5/5 with no unresolved findings. - All 382 focused setup, permissions, metadata, service, and webhook tests passed across final runs. One socket-hang-up test passed on rerun with the full service suite. Final service, metadata, and webhook checks passed all 230 tests. - The broad local suite was stopped after failures. Seven workspace-runtime exposure and control-conflict failures reproduce on base commit `54a99d884`. The broad run also overlapped local iteration; final focused tests and clean-checkout CI are tracked separately. - Browser: a fresh production-backed instance completed enrollment, retried after profile enablement, reached GitHub consent, recovered from a missing installation, and completed OAuth. - Browser: the permissions page showed the authenticated username and the selected private test repository. A real `get_me` call returned the same account. Reading the selected repository passed; reading an unselected private repository failed with 404. - Browser: a second fresh instance completed enrollment and OAuth without a PAT form or unavailable state. Its username and repository list survived reload and refresh. A real get_me call on the final code returned the displayed account. ## Risks - Repository names are now stored in company-scoped grant metadata and shown with that credential. They are display data, not authorization data. - Large selections require more GitHub API calls. A failed later page rejects the refresh rather than reporting a partial list. - Older grants and webhook-invalidated snapshots require Refresh access to load the list. - Cloud profile enablement is separate deployment configuration. This PR does not change OAuth scopes or GitHub App permissions. ## Model Used OpenAI GPT-6 (`gpt-6-astra`) via Codex. Reasoning, code execution, and browser tools were used. The exact context window size was not exposed. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip <noreply@paperclip.ing>
534 lines
21 KiB
TypeScript
534 lines
21 KiB
TypeScript
import { randomUUID } from "node:crypto";
|
||
import {
|
||
connectionEventDeliveries,
|
||
connectionGrants,
|
||
externalObjects,
|
||
toolConnections,
|
||
type Db,
|
||
} from "@paperclipai/db";
|
||
import { and, eq, sql } from "drizzle-orm";
|
||
import {
|
||
logActivity,
|
||
publishActivity,
|
||
type ActivityPublication,
|
||
} from "./activity-log.js";
|
||
import {
|
||
createPaperclipCloudConnector,
|
||
paperclipCloudConnectorConfigFromEnv,
|
||
type PaperclipCloudConnector,
|
||
type SealedConnectorEvents,
|
||
} from "./paperclip-cloud-connector.js";
|
||
import { issueThreadInteractionService } from "./issue-thread-interactions.js";
|
||
import { logger } from "../middleware/logger.js";
|
||
|
||
type LeasedEvent = SealedConnectorEvents["events"][number];
|
||
type GitHubBinding = {
|
||
id: string;
|
||
companyId: string;
|
||
connectionId: string;
|
||
grantId: string;
|
||
subject: string;
|
||
installationId: string;
|
||
providerTenant: NonNullable<typeof connectionGrants.$inferSelect.providerTenant>;
|
||
};
|
||
|
||
export type GitHubConnectionEventPollResult = {
|
||
leased: number;
|
||
processed: number;
|
||
duplicate: number;
|
||
ignored: number;
|
||
failed: number;
|
||
};
|
||
|
||
function record(value: unknown): Record<string, unknown> {
|
||
return value && typeof value === "object" && !Array.isArray(value)
|
||
? value as Record<string, unknown>
|
||
: {};
|
||
}
|
||
|
||
function stringValue(value: unknown): string | null {
|
||
return typeof value === "string" && value.length > 0 ? value : null;
|
||
}
|
||
|
||
function positiveInteger(value: unknown): number | null {
|
||
return typeof value === "number" && Number.isSafeInteger(value) && value > 0 ? value : null;
|
||
}
|
||
|
||
function boundedString(value: unknown, maximum: number): string | null {
|
||
return typeof value === "string" && value.length > 0 && value.length <= maximum ? value : null;
|
||
}
|
||
|
||
function identifier(value: unknown): string | null {
|
||
return typeof value === "string" && /^[1-9][0-9]{0,30}$/.test(value) ? value : null;
|
||
}
|
||
|
||
function isoDate(value: unknown): string | null {
|
||
const candidate = boundedString(value, 100);
|
||
if (!candidate || Number.isNaN(Date.parse(candidate))) return null;
|
||
return new Date(candidate).toISOString();
|
||
}
|
||
|
||
function commitSha(value: unknown): string | null {
|
||
return typeof value === "string" && /^[0-9a-f]{40,64}$/i.test(value) ? value.toLowerCase() : null;
|
||
}
|
||
|
||
function githubUrl(value: unknown): string | null {
|
||
const candidate = boundedString(value, 2_000);
|
||
if (!candidate) return null;
|
||
try {
|
||
const url = new URL(candidate);
|
||
return url.protocol === "https:" && url.hostname.toLowerCase() === "github.com" ? url.toString() : null;
|
||
} catch {
|
||
return null;
|
||
}
|
||
}
|
||
|
||
function compact(values: Record<string, unknown>): Record<string, unknown> {
|
||
return Object.fromEntries(Object.entries(values).filter(([, value]) => value !== null && value !== undefined));
|
||
}
|
||
|
||
/**
|
||
* Treat the sealed Cloud batch as an untrusted boundary. Cloud already normalizes
|
||
* GitHub payloads, but the instance independently allowlists and bounds the small
|
||
* reconciliation record that it persists and processes.
|
||
*/
|
||
function normalizeLeasedPayload(event: LeasedEvent): Record<string, unknown> {
|
||
const payload = record(event.payload);
|
||
const base = {
|
||
event: boundedString(payload.event, 100),
|
||
action: boundedString(payload.action, 100),
|
||
installationId: identifier(payload.installationId),
|
||
repositoryId: identifier(payload.repositoryId),
|
||
repository: boundedString(payload.repository, 300),
|
||
senderId: identifier(payload.senderId),
|
||
senderLogin: boundedString(payload.senderLogin, 100),
|
||
};
|
||
if (event.event === "pull_request") {
|
||
return compact({
|
||
...base,
|
||
number: positiveInteger(payload.number),
|
||
url: githubUrl(payload.url),
|
||
state: boundedString(payload.state, 40),
|
||
merged: payload.merged === true,
|
||
mergedAt: isoDate(payload.mergedAt),
|
||
updatedAt: isoDate(payload.updatedAt),
|
||
headRef: boundedString(payload.headRef, 300),
|
||
headSha: commitSha(payload.headSha),
|
||
baseRef: boundedString(payload.baseRef, 300),
|
||
baseSha: commitSha(payload.baseSha),
|
||
});
|
||
}
|
||
if (event.event === "installation_repositories") {
|
||
const repositoryIds = (value: unknown) => Array.isArray(value)
|
||
? value.slice(0, 1_000).flatMap((item) => identifier(item) ?? [])
|
||
: [];
|
||
return compact({
|
||
...base,
|
||
repositorySelection: boundedString(payload.repositorySelection, 40),
|
||
repositoriesAdded: repositoryIds(payload.repositoriesAdded),
|
||
repositoriesRemoved: repositoryIds(payload.repositoriesRemoved),
|
||
});
|
||
}
|
||
if (event.event === "installation") {
|
||
return compact({
|
||
...base,
|
||
accountId: identifier(payload.accountId),
|
||
accountLogin: boundedString(payload.accountLogin, 100),
|
||
repositorySelection: boundedString(payload.repositorySelection, 40),
|
||
});
|
||
}
|
||
return compact(base);
|
||
}
|
||
|
||
function bindingRows(rows: Array<{
|
||
grant: typeof connectionGrants.$inferSelect;
|
||
connection: typeof toolConnections.$inferSelect;
|
||
}>): GitHubBinding[] {
|
||
return rows.flatMap(({ grant, connection }) => {
|
||
const config = record(connection.config);
|
||
const oauth = record(config.oauth);
|
||
if (config.sourceTemplateKey !== "github" || oauth.connectorProfile !== "github.code") return [];
|
||
const github = grant.providerTenant?.github;
|
||
if (!github || grant.status !== "active") return [];
|
||
const subject = grant.kind === "agent" && grant.subjectAgentId
|
||
? `agent:${grant.subjectAgentId}`
|
||
: grant.kind === "user" && grant.subjectUserId
|
||
? grant.subjectUserId
|
||
: null;
|
||
if (!subject) return [];
|
||
return github.installationIds.map((installationId) => ({
|
||
id: `${grant.id}_${installationId}`,
|
||
companyId: grant.companyId,
|
||
connectionId: connection.id,
|
||
grantId: grant.id,
|
||
subject,
|
||
installationId,
|
||
providerTenant: grant.providerTenant!,
|
||
}));
|
||
});
|
||
}
|
||
|
||
function githubSnapshotUpdate(payload: Record<string, unknown>) {
|
||
const repository = stringValue(payload.repository);
|
||
const number = positiveInteger(payload.number);
|
||
if (!repository || !number) return null;
|
||
const state = stringValue(payload.state) ?? "unknown";
|
||
const merged = payload.merged === true;
|
||
const [owner, repo, ...extra] = repository.split("/");
|
||
if (!owner || !repo || extra.length > 0) return null;
|
||
return {
|
||
repository,
|
||
owner,
|
||
repo,
|
||
number,
|
||
externalId: `${repository}#pull/${number}`,
|
||
state,
|
||
merged,
|
||
statusKey: merged ? "merged" : state === "closed" ? "closed" : "open",
|
||
statusLabel: merged ? "Merged" : state === "closed" ? "Closed" : "Open",
|
||
statusCategory: merged ? "succeeded" : state === "closed" ? "closed" : "open",
|
||
statusTone: merged ? "success" : state === "closed" ? "muted" : "info",
|
||
statusIconKey: merged ? "git-merge" : state === "closed" ? "x-circle" : "git-pull-request",
|
||
data: {
|
||
provider: "github",
|
||
owner,
|
||
repo,
|
||
number,
|
||
state,
|
||
merged,
|
||
...(stringValue(payload.url) ? { url: stringValue(payload.url) } : {}),
|
||
...(stringValue(payload.mergedAt) ? { mergedAt: stringValue(payload.mergedAt) } : {}),
|
||
...(stringValue(payload.headRef) ? { headRef: stringValue(payload.headRef) } : {}),
|
||
...(stringValue(payload.headSha) ? { headSha: stringValue(payload.headSha) } : {}),
|
||
...(stringValue(payload.baseRef) ? { baseRef: stringValue(payload.baseRef) } : {}),
|
||
...(stringValue(payload.baseSha) ? { baseSha: stringValue(payload.baseSha) } : {}),
|
||
},
|
||
remoteVersion: stringValue(payload.updatedAt),
|
||
} as const;
|
||
}
|
||
|
||
export function githubConnectionEventService(
|
||
db: Db,
|
||
options: {
|
||
connector?: PaperclipCloudConnector;
|
||
env?: NodeJS.ProcessEnv;
|
||
now?: () => Date;
|
||
wakeup?: NonNullable<Parameters<typeof issueThreadInteractionService>[1]>["wakeup"];
|
||
} = {},
|
||
) {
|
||
const now = options.now ?? (() => new Date());
|
||
let nextPollAt = 0;
|
||
let emptyPolls = 0;
|
||
|
||
async function activeBindings() {
|
||
const rows = await db.select({ grant: connectionGrants, connection: toolConnections })
|
||
.from(connectionGrants)
|
||
.innerJoin(toolConnections, and(
|
||
eq(toolConnections.id, connectionGrants.connectionId),
|
||
eq(toolConnections.companyId, connectionGrants.companyId),
|
||
))
|
||
.where(and(
|
||
eq(connectionGrants.status, "active"),
|
||
eq(toolConnections.status, "active"),
|
||
eq(toolConnections.enabled, true),
|
||
));
|
||
return bindingRows(rows);
|
||
}
|
||
|
||
async function applyPullRequestEvent(companyId: string, event: LeasedEvent) {
|
||
const snapshot = githubSnapshotUpdate(event.payload);
|
||
if (!snapshot) return;
|
||
const appliedAt = now();
|
||
await db.update(externalObjects).set({
|
||
statusKey: snapshot.statusKey,
|
||
statusLabel: snapshot.statusLabel,
|
||
statusCategory: snapshot.statusCategory,
|
||
statusTone: snapshot.statusTone,
|
||
statusIconKey: snapshot.statusIconKey,
|
||
isTerminal: snapshot.merged || snapshot.state === "closed",
|
||
data: sql`${externalObjects.data} || ${JSON.stringify(snapshot.data)}::jsonb`,
|
||
remoteVersion: snapshot.remoteVersion,
|
||
lastResolvedAt: appliedAt,
|
||
lastChangedAt: appliedAt,
|
||
nextRefreshAt: appliedAt,
|
||
updatedAt: appliedAt,
|
||
}).where(and(
|
||
eq(externalObjects.companyId, companyId),
|
||
eq(externalObjects.providerKey, "github"),
|
||
eq(externalObjects.objectType, "pull_request"),
|
||
sql`lower(${externalObjects.externalId}) = lower(${snapshot.externalId})`,
|
||
));
|
||
if (snapshot.merged && event.action === "closed") {
|
||
await issueThreadInteractionService(db, { wakeup: options.wakeup })
|
||
.sweepMergedPullRequestConfirmations([{
|
||
companyId,
|
||
owner: snapshot.owner,
|
||
repo: snapshot.repo,
|
||
number: snapshot.number,
|
||
}]);
|
||
}
|
||
}
|
||
|
||
async function applyInstallationEvent(database: Db, binding: GitHubBinding, event: LeasedEvent) {
|
||
// Bindings are loaded before the Cloud request. Lock and read the grant
|
||
// again so a refresh completed during that request cannot be overwritten.
|
||
const [currentGrant] = await database.select().from(connectionGrants).where(and(
|
||
eq(connectionGrants.id, binding.grantId),
|
||
eq(connectionGrants.companyId, binding.companyId),
|
||
eq(connectionGrants.status, "active"),
|
||
)).for("update").limit(1);
|
||
const currentProviderTenant = currentGrant?.providerTenant;
|
||
const github = currentProviderTenant?.github;
|
||
if (!github) return;
|
||
// A newly bound instance can receive installation events from before OAuth
|
||
// verified its repository list. Those events must not erase newer access.
|
||
if (Date.parse(github.lastAccessRefreshAt ?? "") > Date.parse(event.createdAt)) {
|
||
await database.update(connectionGrants).set({
|
||
providerTenant: {
|
||
...currentProviderTenant,
|
||
github: { ...github, lastWebhookAt: now().toISOString(), webhookHealth: "healthy" },
|
||
},
|
||
updatedAt: now(),
|
||
}).where(and(eq(connectionGrants.id, binding.grantId), eq(connectionGrants.companyId, binding.companyId)));
|
||
return;
|
||
}
|
||
const unavailable = event.event === "installation" && (event.action === "deleted" || event.action === "suspend");
|
||
const installationIds = unavailable
|
||
? github.installationIds.filter((id) => id !== binding.installationId)
|
||
: [...new Set([...github.installationIds, binding.installationId])];
|
||
const added = Array.isArray(event.payload.repositoriesAdded) ? event.payload.repositoriesAdded.length : 0;
|
||
const removed = Array.isArray(event.payload.repositoriesRemoved) ? event.payload.repositoriesRemoved.length : 0;
|
||
const repositoryCount = Math.max(
|
||
0,
|
||
unavailable && installationIds.length === 0
|
||
? 0
|
||
: unavailable
|
||
? github.repositoryCount
|
||
: github.repositoryCount + added - removed,
|
||
);
|
||
const providerTenant = {
|
||
...currentProviderTenant,
|
||
github: {
|
||
...github,
|
||
// Lifecycle webhooks carry IDs, not the user token’s complete repository view.
|
||
// Discard the snapshot until Refresh access verifies it again.
|
||
accessRevision: randomUUID(),
|
||
repositories: undefined,
|
||
installationIds,
|
||
installationCount: installationIds.length,
|
||
repositoryCount,
|
||
repositorySelection: repositoryCount === 0 ? "none" as const : github.repositorySelection,
|
||
lastWebhookAt: now().toISOString(),
|
||
webhookHealth: unavailable ? "unhealthy" as const : "healthy" as const,
|
||
},
|
||
};
|
||
await database.update(connectionGrants).set({ providerTenant, updatedAt: now() })
|
||
.where(and(eq(connectionGrants.id, binding.grantId), eq(connectionGrants.companyId, binding.companyId)));
|
||
await database.update(toolConnections).set({
|
||
healthStatus: unavailable ? "failed" : "ok",
|
||
healthMessage: unavailable
|
||
? "GitHub installation access was removed or suspended. Manage repository access on GitHub."
|
||
: "GitHub installation and repository access are available.",
|
||
healthCheckedAt: now(),
|
||
lastHealthAt: now(),
|
||
lastError: unavailable ? "GitHub installation unavailable" : null,
|
||
updatedAt: now(),
|
||
}).where(and(eq(toolConnections.id, binding.connectionId), eq(toolConnections.companyId, binding.companyId)));
|
||
}
|
||
|
||
async function processForCompany(companyId: string, bindings: GitHubBinding[], event: LeasedEvent) {
|
||
const receiptAt = now();
|
||
const normalizedEvent = { ...event, payload: normalizeLeasedPayload(event) };
|
||
const [receipt] = await db.insert(connectionEventDeliveries).values({
|
||
companyId,
|
||
provider: event.provider,
|
||
providerDeliveryId: event.id,
|
||
event: event.event,
|
||
action: event.action,
|
||
installationId: event.installationId,
|
||
repositoryId: event.repositoryId,
|
||
normalizedPayload: normalizedEvent.payload,
|
||
providerCreatedAt: new Date(event.createdAt),
|
||
status: "received",
|
||
attempts: 1,
|
||
updatedAt: receiptAt,
|
||
}).onConflictDoNothing().returning();
|
||
if (!receipt) {
|
||
const [existing] = await db.select().from(connectionEventDeliveries).where(and(
|
||
eq(connectionEventDeliveries.companyId, companyId),
|
||
eq(connectionEventDeliveries.provider, event.provider),
|
||
eq(connectionEventDeliveries.providerDeliveryId, event.id),
|
||
)).limit(1);
|
||
if (existing?.status === "processed") return "duplicate" as const;
|
||
await db.update(connectionEventDeliveries).set({
|
||
status: "received",
|
||
attempts: sql`${connectionEventDeliveries.attempts} + 1`,
|
||
lastError: null,
|
||
updatedAt: receiptAt,
|
||
}).where(eq(connectionEventDeliveries.id, existing!.id));
|
||
}
|
||
const postCommitPublications: ActivityPublication[] = [];
|
||
try {
|
||
const applyAndFinalize = async (database: Db) => {
|
||
if (event.event === "pull_request") await applyPullRequestEvent(companyId, normalizedEvent);
|
||
if (event.event === "installation" || event.event === "installation_repositories") {
|
||
for (const binding of bindings) await applyInstallationEvent(database, binding, normalizedEvent);
|
||
} else {
|
||
const touchedAt = now();
|
||
for (const binding of bindings) {
|
||
const github = binding.providerTenant.github;
|
||
if (!github) continue;
|
||
await database.update(connectionGrants).set({
|
||
// Update only webhook fields; a concurrent access/token refresh
|
||
// owns the remaining metadata and must not be replaced here.
|
||
providerTenant: sql`jsonb_set(${connectionGrants.providerTenant}, '{github}',
|
||
(${connectionGrants.providerTenant}->'github') || ${JSON.stringify({
|
||
lastWebhookAt: touchedAt.toISOString(), webhookHealth: "healthy",
|
||
})}::jsonb)`,
|
||
updatedAt: touchedAt,
|
||
}).where(and(
|
||
eq(connectionGrants.id, binding.grantId),
|
||
eq(connectionGrants.companyId, companyId),
|
||
eq(connectionGrants.status, "active"),
|
||
sql`${connectionGrants.providerTenant}->'github' is not null`,
|
||
));
|
||
}
|
||
}
|
||
const finishedAt = now();
|
||
await database.update(connectionEventDeliveries).set({
|
||
status: "processed",
|
||
processedAt: finishedAt,
|
||
lastError: null,
|
||
updatedAt: finishedAt,
|
||
}).where(and(
|
||
eq(connectionEventDeliveries.companyId, companyId),
|
||
eq(connectionEventDeliveries.provider, event.provider),
|
||
eq(connectionEventDeliveries.providerDeliveryId, event.id),
|
||
));
|
||
await logActivity(database, {
|
||
companyId,
|
||
actorType: "system",
|
||
actorId: "system:github-webhook",
|
||
action: "tool_connection.webhook_processed",
|
||
entityType: "tool_connection",
|
||
entityId: bindings[0]!.connectionId,
|
||
details: {
|
||
provider: "github",
|
||
event: event.event,
|
||
action: event.action,
|
||
deliveryId: event.id,
|
||
installationId: event.installationId,
|
||
repositoryId: event.repositoryId,
|
||
},
|
||
}, postCommitPublications);
|
||
};
|
||
if (event.event === "installation" || event.event === "installation_repositories") {
|
||
await db.transaction(async (tx) => applyAndFinalize(tx as unknown as Db));
|
||
} else {
|
||
await applyAndFinalize(db);
|
||
}
|
||
} catch (error) {
|
||
await db.update(connectionEventDeliveries).set({
|
||
status: "failed",
|
||
lastError: error instanceof Error ? error.message.slice(0, 500) : "GitHub event processing failed",
|
||
updatedAt: now(),
|
||
}).where(and(
|
||
eq(connectionEventDeliveries.companyId, companyId),
|
||
eq(connectionEventDeliveries.provider, event.provider),
|
||
eq(connectionEventDeliveries.providerDeliveryId, event.id),
|
||
));
|
||
throw error;
|
||
}
|
||
// Persistence is complete at this point (and installation deltas have
|
||
// committed). A synchronous live-event subscriber must not turn that
|
||
// durable success back into a retryable receipt and replay the delta.
|
||
for (const publication of postCommitPublications) {
|
||
try {
|
||
publishActivity(publication);
|
||
} catch (error) {
|
||
logger.warn({
|
||
err: error,
|
||
companyId,
|
||
providerDeliveryId: event.id,
|
||
}, "GitHub webhook activity publication failed after commit");
|
||
}
|
||
}
|
||
return "processed" as const;
|
||
}
|
||
|
||
return {
|
||
async pollOnce(): Promise<GitHubConnectionEventPollResult> {
|
||
if (now().getTime() < nextPollAt) {
|
||
return { leased: 0, processed: 0, duplicate: 0, ignored: 0, failed: 0 };
|
||
}
|
||
const bindings = await activeBindings();
|
||
if (bindings.length === 0) {
|
||
nextPollAt = now().getTime() + 5 * 60_000;
|
||
return { leased: 0, processed: 0, duplicate: 0, ignored: 0, failed: 0 };
|
||
}
|
||
const config = options.connector ? null : paperclipCloudConnectorConfigFromEnv(options.env);
|
||
const connector = options.connector ?? (config ? createPaperclipCloudConnector({ config }) : null);
|
||
if (!connector) {
|
||
nextPollAt = now().getTime() + 5 * 60_000;
|
||
return { leased: 0, processed: 0, duplicate: 0, ignored: 0, failed: 0 };
|
||
}
|
||
const first = bindings[0]!;
|
||
const lease = await connector.leaseEvents({ subject: first.subject, companyId: first.companyId });
|
||
if (!lease) {
|
||
emptyPolls += 1;
|
||
nextPollAt = now().getTime() + Math.min(5 * 60_000, 5_000 * (2 ** Math.min(emptyPolls, 6)));
|
||
return { leased: 0, processed: 0, duplicate: 0, ignored: 0, failed: 0 };
|
||
}
|
||
emptyPolls = 0;
|
||
nextPollAt = now().getTime() + 5_000;
|
||
const result: GitHubConnectionEventPollResult = {
|
||
leased: lease.events.length,
|
||
processed: 0,
|
||
duplicate: 0,
|
||
ignored: 0,
|
||
failed: 0,
|
||
};
|
||
const acknowledge: string[] = [];
|
||
for (const event of lease.events) {
|
||
const matched = bindings.filter((binding) => event.bindingIds.includes(binding.id));
|
||
if (matched.length === 0) {
|
||
result.ignored += 1;
|
||
acknowledge.push(event.id);
|
||
continue;
|
||
}
|
||
try {
|
||
const companies = new Map<string, GitHubBinding[]>();
|
||
for (const binding of matched) companies.set(binding.companyId, [...(companies.get(binding.companyId) ?? []), binding]);
|
||
for (const [companyId, companyBindings] of companies) {
|
||
const status = await processForCompany(companyId, companyBindings, event);
|
||
result[status] += 1;
|
||
}
|
||
acknowledge.push(event.id);
|
||
if (event.event === "installation" && (event.action === "deleted" || event.action === "suspend")) {
|
||
await Promise.all(matched.map((binding) => connector.setWebhookBinding({
|
||
subject: binding.subject,
|
||
companyId: binding.companyId,
|
||
id: binding.id,
|
||
installationId: binding.installationId,
|
||
connectionId: binding.connectionId,
|
||
grantId: binding.grantId,
|
||
active: false,
|
||
})));
|
||
}
|
||
} catch {
|
||
result.failed += 1;
|
||
}
|
||
}
|
||
if (acknowledge.length > 0) {
|
||
await connector.acknowledgeEvents({
|
||
subject: first.subject,
|
||
companyId: first.companyId,
|
||
leaseId: lease.leaseId,
|
||
deliveryIds: acknowledge,
|
||
});
|
||
}
|
||
return result;
|
||
},
|
||
};
|
||
}
|