(
input: CreateConnectionInput,
)
| 3803 | // ------------------------------------------------------------------ |
| 3804 | |
| 3805 | const connectionsCreate = ( |
| 3806 | input: CreateConnectionInput, |
| 3807 | ): Effect.Effect< |
| 3808 | Connection, |
| 3809 | | IntegrationNotFoundError |
| 3810 | | ConnectionAlreadyExistsError |
| 3811 | | CredentialProviderNotRegisteredError |
| 3812 | | InvalidConnectionInputError |
| 3813 | | OrgWriteDeniedError |
| 3814 | | StorageFailure |
| 3815 | > => |
| 3816 | Effect.gen(function* () { |
| 3817 | yield* guardOrgWrite(input.owner); |
| 3818 | const name = connectionIdentifier(String(input.name)); |
| 3819 | // Typed (not StorageError) so the HTTP edge can answer 400 with the |
| 3820 | // reason instead of an opaque 500 — callers can act on it. |
| 3821 | if (input.owner === "user" && subject == null) { |
| 3822 | return yield* new InvalidConnectionInputError({ |
| 3823 | message: |
| 3824 | 'Cannot create a personal connection: this context has no user subject. Create it with owner "org", or connect as a signed-in user.', |
| 3825 | }); |
| 3826 | } |
| 3827 | const integrationRow = yield* findIntegrationRow(input.integration); |
| 3828 | if (!integrationRow) { |
| 3829 | return yield* new IntegrationNotFoundError({ |
| 3830 | slug: input.integration, |
| 3831 | }); |
| 3832 | } |
| 3833 | |
| 3834 | // Create is never a replace. This early check answers the common case |
| 3835 | // with a typed 409 before any other work, but it is NOT the guard |
| 3836 | // against concurrent creates — the row insert below is: the |
| 3837 | // transaction re-checks, the primary key breaks the tie, and the |
| 3838 | // provider write happens only after the insert wins. |
| 3839 | const duplicate = yield* findConnectionRow({ |
| 3840 | owner: input.owner, |
| 3841 | integration: input.integration, |
| 3842 | name, |
| 3843 | }); |
| 3844 | let retryingRowId: string | null = null; |
| 3845 | let retryingItemIds: readonly string[] = []; |
| 3846 | if (duplicate) { |
| 3847 | const duplicateItemIds = Object.values(connectionItemIds(duplicate)); |
| 3848 | const duplicateProvider = credentialProviders.get(duplicate.provider); |
| 3849 | const duplicateAttempt = parseCredentialWriteAttempt(duplicate.credential_write); |
| 3850 | const belongsToCrashedRuntime = |
| 3851 | duplicateItemIds.length > 0 && |
| 3852 | duplicateAttempt !== null && |
| 3853 | duplicateAttempt.runtimeId !== credentialWriteRuntimeId; |
| 3854 | const hasMissingCredential = |
| 3855 | belongsToCrashedRuntime && duplicateProvider?.set |
| 3856 | ? (yield* Effect.forEach(duplicateItemIds, (itemId) => |
| 3857 | duplicateProvider |
| 3858 | .get(ProviderItemId.make(itemId)) |
| 3859 | .pipe(Effect.map((value) => value === null)), |
| 3860 | )).some(Boolean) |
| 3861 | : false; |
| 3862 | retryingRowId = hasMissingCredential ? storageRowId(duplicate) : null; |
no test coverage detected