From 20deab15676131cb456dfb7c57d7a4dac87b3445 Mon Sep 17 00:00:00 2001 From: "supabase-cli-releaser[bot]" <246109035+supabase-cli-releaser[bot]@users.noreply.github.com> Date: Thu, 24 Sep 2026 15:02:39 +0000 Subject: [PATCH 1/2] chore: sync API types from infrastructure --- apps/cli-go/api/v1-openapi.yaml | 511 +++++++++++++++++++++++------- apps/cli-go/pkg/api/client.gen.go | 76 +---- apps/cli-go/pkg/api/types.gen.go | 14 +- 3 files changed, 415 insertions(+), 186 deletions(-) diff --git a/apps/cli-go/api/v1-openapi.yaml b/apps/cli-go/api/v1-openapi.yaml index d1ac70b88e..8c64234393 100644 --- a/apps/cli-go/api/v1-openapi.yaml +++ b/apps/cli-go/api/v1-openapi.yaml @@ -400,7 +400,10 @@ paths: x-oauth-scope: environment:write /v1/branches/{branch_id_or_ref}/diff: get: - description: Diffs the specified database branch + description: |- + Diffs the specified database branch + + This endpoint is currently in its **Beta** stage. operationId: v1-diff-a-branch parameters: - name: branch_id_or_ref @@ -457,18 +460,21 @@ paths: description: Failed to diff database branch security: - bearer: [] - summary: '[Beta] Diffs a database branch' + summary: Diffs a database branch tags: - Environments x-badges: - name: 'OAuth scope: environment:write' position: after + - name: Beta + position: before x-endpoint-owners: - dev-workflows x-fga-permissions: - - branching_development_write - - branching_production_write x-oauth-scope: environment:write + x-scalar-stability: experimental /v1/projects: get: description: Returns a list of all projects you've previously created. @@ -540,6 +546,13 @@ paths: x-oauth-scope: projects:write /v1/projects/available-regions: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-available-regions parameters: - name: organization_slug @@ -604,17 +617,18 @@ paths: $ref: '#/components/schemas/RegionsInfo_Output' security: - bearer: [] - summary: >- - [Beta] Gets the list of available regions that can be used for a new - project + summary: Gets the list of available regions that can be used for a new project tags: - Projects x-badges: - name: 'OAuth scope: organizations:read' position: after + - name: Beta + position: before x-endpoint-owners: - infra x-oauth-scope: organizations:read + x-scalar-stability: experimental /v1/organizations: get: description: Returns a list of organizations that you currently belong to. @@ -686,6 +700,13 @@ paths: - - organizations_create /v1/oauth/authorize: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-authorize-user parameters: - name: client_id @@ -777,12 +798,16 @@ paths: description: Forbidden action '429': description: Rate limit exceeded - summary: '[Beta] Authorize user through oauth' + summary: Authorize user through oauth tags: - OAuth + x-badges: + - name: Beta + position: before x-endpoint-owners: - auth - control-plane + x-scalar-stability: experimental /v1/oauth/token: post: description: >- @@ -790,6 +815,9 @@ paths: `urn:ietf:params:oauth:grant-type:jwt-bearer` grant types. The `jwt-bearer` grant type (IDJAG — identity-directed JWT assertion) is in beta and available on Team and Enterprise plans only. + + + This endpoint is currently in its **Beta** stage. operationId: v1-exchange-oauth-token parameters: [] requestBody: @@ -811,14 +839,25 @@ paths: description: Forbidden action '429': description: Rate limit exceeded - summary: '[Beta] Exchange auth code for user''s access and refresh token' + summary: Exchange auth code for user's access and refresh token tags: - OAuth + x-badges: + - name: Beta + position: before x-endpoint-owners: - auth - control-plane + x-scalar-stability: experimental /v1/oauth/revoke: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-revoke-token parameters: [] requestBody: @@ -836,12 +875,16 @@ paths: description: Forbidden action '429': description: Rate limit exceeded - summary: '[Beta] Revoke oauth app authorization and it''s corresponding tokens' + summary: Revoke oauth app authorization and it's corresponding tokens tags: - OAuth + x-badges: + - name: Beta + position: before x-endpoint-owners: - auth - control-plane + x-scalar-stability: experimental /v1/oauth/authorize/project-claim: get: description: >- @@ -1916,6 +1959,13 @@ paths: x-oauth-scope: environment:read /v1/projects/{ref}/custom-hostname: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-hostname-config parameters: - name: ref @@ -1945,19 +1995,29 @@ paths: description: Failed to retrieve project's custom hostname config security: - bearer: [] - summary: '[Beta] Gets project''s custom hostname config' + summary: Gets project's custom hostname config tags: - Domains x-badges: - name: 'OAuth scope: domains:read' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - custom_domain_read x-oauth-scope: domains:read + x-scalar-stability: experimental delete: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-Delete hostname config parameters: - name: ref @@ -1991,20 +2051,30 @@ paths: description: Failed to delete project custom hostname configuration security: - bearer: [] - summary: '[Beta] Deletes a project''s custom hostname configuration' + summary: Deletes a project's custom hostname configuration tags: - Domains x-badges: - name: 'OAuth scope: domains:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - custom_domain_write x-oauth-scope: domains:write + x-scalar-stability: experimental /v1/projects/{ref}/custom-hostname/initialize: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-update-hostname-config parameters: - name: ref @@ -2040,20 +2110,30 @@ paths: description: Failed to update project custom hostname configuration security: - bearer: [] - summary: '[Beta] Updates project''s custom hostname configuration' + summary: Updates project's custom hostname configuration tags: - Domains x-badges: - name: 'OAuth scope: domains:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - custom_domain_write x-oauth-scope: domains:write + x-scalar-stability: experimental /v1/projects/{ref}/custom-hostname/reverify: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-verify-dns-config parameters: - name: ref @@ -2084,21 +2164,31 @@ paths: security: - bearer: [] summary: >- - [Beta] Attempts to verify the DNS configuration for project's custom - hostname configuration + Attempts to verify the DNS configuration for project's custom hostname + configuration tags: - Domains x-badges: - name: 'OAuth scope: domains:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - custom_domain_write x-oauth-scope: domains:write + x-scalar-stability: experimental /v1/projects/{ref}/custom-hostname/activate: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-activate-custom-hostname parameters: - name: ref @@ -2128,20 +2218,30 @@ paths: description: Failed to activate project custom hostname configuration security: - bearer: [] - summary: '[Beta] Activates a custom hostname for a project.' + summary: Activates a custom hostname for a project. tags: - Domains x-badges: - name: 'OAuth scope: domains:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - custom_domain_write x-oauth-scope: domains:write + x-scalar-stability: experimental /v1/projects/{ref}/jit-access: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-jit-access-config parameters: - name: ref @@ -2198,19 +2298,29 @@ paths: description: Failed to retrieve project's temporary access configuration. security: - bearer: [] - summary: '[Beta] Get project''s temporary access configuration.' + summary: Get project's temporary access configuration. tags: - Database x-badges: - name: 'OAuth scope: database:read' position: after + - name: Beta + position: before x-endpoint-owners: - security - control-plane x-fga-permissions: - - project_admin_read x-oauth-scope: database:read + x-scalar-stability: experimental put: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-update-jit-access-config parameters: - name: ref @@ -2273,20 +2383,30 @@ paths: description: Failed to update project's temporary access configuration. security: - bearer: [] - summary: '[Beta] Update project''s temporary access configuration.' + summary: Update project's temporary access configuration. tags: - Database x-badges: - name: 'OAuth scope: database:write' position: after + - name: Beta + position: before x-endpoint-owners: - security - control-plane x-fga-permissions: - - project_admin_write x-oauth-scope: database:write + x-scalar-stability: experimental /v1/projects/{ref}/network-bans/retrieve: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-list-all-network-bans parameters: - name: ref @@ -2316,20 +2436,30 @@ paths: description: Failed to retrieve project's network bans security: - bearer: [] - summary: '[Beta] Gets project''s network bans' + summary: Gets project's network bans tags: - Projects x-badges: - name: 'OAuth scope: projects:read' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane - infra x-fga-permissions: - - database_network_bans_read x-oauth-scope: projects:read + x-scalar-stability: experimental /v1/projects/{ref}/network-bans/retrieve/enriched: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-list-all-network-bans-enriched parameters: - name: ref @@ -2360,21 +2490,31 @@ paths: security: - bearer: [] summary: >- - [Beta] Gets project's network bans with additional information about - which databases they affect + Gets project's network bans with additional information about which + databases they affect tags: - Projects x-badges: - name: 'OAuth scope: projects:read' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane - infra x-fga-permissions: - - database_network_bans_read x-oauth-scope: projects:read + x-scalar-stability: experimental /v1/projects/{ref}/network-bans: delete: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-delete-network-bans parameters: - name: ref @@ -2406,20 +2546,30 @@ paths: description: Failed to remove network bans. security: - bearer: [] - summary: '[Beta] Remove network bans.' + summary: Remove network bans. tags: - Projects x-badges: - name: 'OAuth scope: projects:write' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane - infra x-fga-permissions: - - database_network_bans_write x-oauth-scope: projects:write + x-scalar-stability: experimental /v1/projects/{ref}/network-restrictions: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-network-restrictions parameters: - name: ref @@ -2449,19 +2599,29 @@ paths: description: Failed to retrieve project's network restrictions security: - bearer: [] - summary: '[Beta] Gets project''s network restrictions' + summary: Gets project's network restrictions tags: - Projects x-badges: - name: 'OAuth scope: projects:read' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - database_network_restrictions_read x-oauth-scope: projects:read + x-scalar-stability: experimental patch: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Alpha** stage. operationId: v1-patch-network-restrictions parameters: - name: ref @@ -2497,22 +2657,30 @@ paths: description: Failed to update project network restrictions security: - bearer: [] - summary: >- - [Alpha] Updates project's network restrictions by adding or removing - CIDRs + summary: Updates project's network restrictions by adding or removing CIDRs tags: - Projects x-badges: - name: 'OAuth scope: projects:write' position: after + - name: Alpha + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - database_network_restrictions_write x-oauth-scope: projects:write + x-scalar-stability: experimental /v1/projects/{ref}/network-restrictions/apply: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-update-network-restrictions parameters: - name: ref @@ -2548,20 +2716,30 @@ paths: description: Failed to update project network restrictions security: - bearer: [] - summary: '[Beta] Updates project''s network restrictions' + summary: Updates project's network restrictions tags: - Projects x-badges: - name: 'OAuth scope: projects:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - database_network_restrictions_write x-oauth-scope: projects:write + x-scalar-stability: experimental /v1/projects/{ref}/pgsodium: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-pgsodium-config parameters: - name: ref @@ -2591,18 +2769,28 @@ paths: description: Failed to retrieve project's pgsodium config security: - bearer: [] - summary: '[Beta] Gets project''s pgsodium config' + summary: Gets project's pgsodium config tags: - Secrets x-badges: - name: 'OAuth scope: secrets:read' position: after + - name: Beta + position: before x-endpoint-owners: - infra x-fga-permissions: - - project_admin_write x-oauth-scope: secrets:read + x-scalar-stability: experimental put: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-update-pgsodium-config parameters: - name: ref @@ -2639,18 +2827,21 @@ paths: security: - bearer: [] summary: >- - [Beta] Updates project's pgsodium config. Updating the root_key can - cause all data encrypted with the older key to become inaccessible. + Updates project's pgsodium config. Updating the root_key can cause all + data encrypted with the older key to become inaccessible. tags: - Secrets x-badges: - name: 'OAuth scope: secrets:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra x-fga-permissions: - - project_admin_write x-oauth-scope: secrets:write + x-scalar-stability: experimental /v1/projects/{ref}/postgrest: get: operationId: v1-get-postgrest-service-config @@ -3010,6 +3201,13 @@ paths: x-oauth-scope: secrets:write /v1/projects/{ref}/ssl-enforcement: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-ssl-enforcement-config parameters: - name: ref @@ -3039,19 +3237,29 @@ paths: description: Failed to retrieve project's SSL enforcement config security: - bearer: [] - summary: '[Beta] Get project''s SSL enforcement configuration.' + summary: Get project's SSL enforcement configuration. tags: - Database x-badges: - name: 'OAuth scope: database:read' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - database_ssl_config_read x-oauth-scope: database:read + x-scalar-stability: experimental put: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-update-ssl-enforcement-config parameters: - name: ref @@ -3087,18 +3295,21 @@ paths: description: Failed to update project's SSL enforcement configuration. security: - bearer: [] - summary: '[Beta] Update project''s SSL enforcement configuration.' + summary: Update project's SSL enforcement configuration. tags: - Database x-badges: - name: 'OAuth scope: database:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - database_ssl_config_write x-oauth-scope: database:write + x-scalar-stability: experimental /v1/projects/{ref}/types/typescript: get: description: Returns the TypeScript types of your schema for use with supabase-js. @@ -3151,6 +3362,13 @@ paths: x-oauth-scope: database:read /v1/projects/{ref}/vanity-subdomain: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-vanity-subdomain-config parameters: - name: ref @@ -3188,7 +3406,7 @@ paths: description: Failed to get project vanity subdomain configuration security: - bearer: [] - summary: '[Beta] Gets current vanity subdomain config' + summary: Gets current vanity subdomain config tags: - Domains x-allowed-plans: @@ -3198,6 +3416,8 @@ paths: x-badges: - name: 'OAuth scope: domains:read' position: after + - name: Beta + position: before - name: Only available on Pro, Team, Enterprise position: before x-endpoint-owners: @@ -3206,7 +3426,15 @@ paths: x-fga-permissions: - - vanity_subdomain_read x-oauth-scope: domains:read + x-scalar-stability: experimental delete: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-deactivate-vanity-subdomain-config parameters: - name: ref @@ -3232,20 +3460,30 @@ paths: description: Failed to delete project vanity subdomain configuration security: - bearer: [] - summary: '[Beta] Deletes a project''s vanity subdomain configuration' + summary: Deletes a project's vanity subdomain configuration tags: - Domains x-badges: - name: 'OAuth scope: domains:write' position: after + - name: Beta + position: before x-endpoint-owners: - infra - control-plane x-fga-permissions: - - vanity_subdomain_write x-oauth-scope: domains:write + x-scalar-stability: experimental /v1/projects/{ref}/vanity-subdomain/check-availability: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-check-vanity-subdomain-availability parameters: - name: ref @@ -3289,7 +3527,7 @@ paths: description: Failed to check project vanity subdomain configuration security: - bearer: [] - summary: '[Beta] Checks vanity subdomain availability' + summary: Checks vanity subdomain availability tags: - Domains x-allowed-plans: @@ -3299,6 +3537,8 @@ paths: x-badges: - name: 'OAuth scope: domains:write' position: after + - name: Beta + position: before - name: Only available on Pro, Team, Enterprise position: before x-endpoint-owners: @@ -3307,8 +3547,16 @@ paths: x-fga-permissions: - - vanity_subdomain_write x-oauth-scope: domains:write + x-scalar-stability: experimental /v1/projects/{ref}/vanity-subdomain/activate: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-activate-vanity-subdomain-config parameters: - name: ref @@ -3352,7 +3600,7 @@ paths: description: Failed to activate project vanity subdomain configuration security: - bearer: [] - summary: '[Beta] Activates a vanity subdomain for a project.' + summary: Activates a vanity subdomain for a project. tags: - Domains x-allowed-plans: @@ -3362,6 +3610,8 @@ paths: x-badges: - name: 'OAuth scope: domains:write' position: after + - name: Beta + position: before - name: Only available on Pro, Team, Enterprise position: before x-endpoint-owners: @@ -3370,8 +3620,16 @@ paths: x-fga-permissions: - - vanity_subdomain_write x-oauth-scope: domains:write + x-scalar-stability: experimental /v1/projects/{ref}/upgrade: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-upgrade-postgres-version parameters: - name: ref @@ -3407,12 +3665,14 @@ paths: description: Failed to initiate project upgrade security: - bearer: [] - summary: '[Beta] Upgrades the project''s Postgres version' + summary: Upgrades the project's Postgres version tags: - Projects x-badges: - name: 'OAuth scope: projects:write' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane - infra @@ -3420,8 +3680,16 @@ paths: - - project_admin_write - database_write x-oauth-scope: projects:write + x-scalar-stability: experimental /v1/projects/{ref}/upgrade/eligibility: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-postgres-upgrade-eligibility parameters: - name: ref @@ -3451,12 +3719,14 @@ paths: description: Failed to determine project upgrade eligibility security: - bearer: [] - summary: '[Beta] Returns the project''s eligibility for upgrades' + summary: Returns the project's eligibility for upgrades tags: - Projects x-badges: - name: 'OAuth scope: projects:read' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane - infra @@ -3464,8 +3734,16 @@ paths: - - project_admin_read - database_read x-oauth-scope: projects:read + x-scalar-stability: experimental /v1/projects/{ref}/upgrade/status: get: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-get-postgres-upgrade-status parameters: - name: ref @@ -3501,12 +3779,14 @@ paths: description: Failed to retrieve project upgrade status security: - bearer: [] - summary: '[Beta] Gets the latest status of the project''s upgrade' + summary: Gets the latest status of the project's upgrade tags: - Projects x-badges: - name: 'OAuth scope: projects:read' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane - infra @@ -3514,6 +3794,7 @@ paths: - - project_admin_read - database_read x-oauth-scope: projects:read + x-scalar-stability: experimental /v1/projects/{ref}/readonly: get: operationId: v1-get-readonly-mode-status @@ -3600,6 +3881,13 @@ paths: x-oauth-scope: database:write /v1/projects/{ref}/read-replicas/setup: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-setup-a-read-replica parameters: - name: ref @@ -3639,7 +3927,7 @@ paths: description: Failed to set up read replica security: - bearer: [] - summary: '[Beta] Set up a read replica' + summary: Set up a read replica tags: - Database x-allowed-plans: @@ -3647,6 +3935,8 @@ paths: - Team - Enterprise x-badges: + - name: Beta + position: before - name: Only available on Pro, Team, Enterprise position: before x-endpoint-owners: @@ -3654,8 +3944,16 @@ paths: - infra x-fga-permissions: - - infra_read_replicas_write + x-scalar-stability: experimental /v1/projects/{ref}/read-replicas/remove: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-remove-a-read-replica parameters: - name: ref @@ -3687,14 +3985,18 @@ paths: description: Failed to remove read replica security: - bearer: [] - summary: '[Beta] Remove a read replica' + summary: Remove a read replica tags: - Database + x-badges: + - name: Beta + position: before x-endpoint-owners: - control-plane - infra x-fga-permissions: - - infra_read_replicas_write + x-scalar-stability: experimental /v1/projects/{ref}/health: get: operationId: v1-get-services-health @@ -4837,7 +5139,6 @@ paths: x-internal: true /v1/projects/{ref}/advisors/performance: get: - deprecated: true description: >- This is an **experimental** endpoint. It is subject to change or removal in future versions. Use it with caution, as it may not remain supported @@ -4880,9 +5181,9 @@ paths: x-fga-permissions: - - advisors_read x-oauth-scope: database:read + x-scalar-stability: experimental /v1/projects/{ref}/advisors/security: get: - deprecated: true description: >- This is an **experimental** endpoint. It is subject to change or removal in future versions. Use it with caution, as it may not remain supported @@ -4933,28 +5234,15 @@ paths: x-fga-permissions: - - advisors_read x-oauth-scope: database:read + x-scalar-stability: experimental /v1/projects/{ref}/analytics/endpoints/logs.all: get: deprecated: true - description: > - Executes a SQL query on the project's logs. - - - Either the `iso_timestamp_start` and `iso_timestamp_end` parameters must - be provided. - - If both are not provided, only the last 1 minute of logs will be - queried. - - The timestamp range must be no more than 24 hours and is rounded to the - nearest minute. If the range is more than 24 hours, a validation error - will be thrown. - - - Note: Unless the `sql` parameter is provided, only edge_logs will be - queried. See the [log query - docs](https://supabase.com/docs/guides/monitoring-and-debugging/logs#logs-explorer) - for all available sources. + description: >- + This endpoint has been removed and always responds with `410 Gone`. Use + `GET /v1/projects/{ref}/analytics/endpoints/logs` instead. See the + [migration + guide](https://supabase.com/changelog/48235-migration-of-supabase-management-api-logs-all-analytics-endpoint-to-logs-endpoint). operationId: v1-get-project-logs-all parameters: - name: ref @@ -4967,47 +5255,13 @@ paths: pattern: ^[a-z]+$ example: abcdefghijklmnopqrst type: string - - name: sql - required: false - in: query - description: >- - Custom SQL query to execute on the logs. See [querying - logs](https://supabase.com/docs/guides/monitoring-and-debugging/logs#querying-with-the-logs-explorer) - for more details. - schema: - example: select event_message from edge_logs limit 10 - type: string - - name: iso_timestamp_start - required: false - in: query - schema: - format: date-time - pattern: >- - ^(?:(?:\d\d[2468][048]|\d\d[13579][26]|\d\d0[48]|[02468][048]00|[13579][26]00)-02-29|\d{4}-(?:(?:0[13578]|1[02])-(?:0[1-9]|[12]\d|3[01])|(?:0[469]|11)-(?:0[1-9]|[12]\d|30)|(?:02)-(?:0[1-9]|1\d|2[0-8])))T(?:(?:[01]\d|2[0-3]):[0-5]\d(?::[0-5]\d(?:\.\d+)?)?(?:Z))$ - example: '2025-03-01T00:00:00Z' - type: string - - name: iso_timestamp_end - required: false - in: query - schema: - format: date-time - pattern: >- - ^(?:(?:\d\d[2468][048]|\d\d[13579][26]|\d\d0[48]|[02468][048]00|[13579][26]00)-02-29|\d{4}-(?:(?:0[13578]|1[02])-(?:0[1-9]|[12]\d|3[01])|(?:0[469]|11)-(?:0[1-9]|[12]\d|30)|(?:02)-(?:0[1-9]|1\d|2[0-8])))T(?:(?:[01]\d|2[0-3]):[0-5]\d(?::[0-5]\d(?:\.\d+)?)?(?:Z))$ - example: '2025-03-01T23:59:59Z' - type: string responses: - '200': - description: '' - content: - application/json: - schema: - $ref: '#/components/schemas/AnalyticsResponse_Output' '401': description: Unauthorized - '402': - description: Usage exceeded. Enable additional usage to continue querying '403': description: Forbidden action + '410': + description: This endpoint has been removed '429': description: Rate limit exceeded security: @@ -5313,6 +5567,13 @@ paths: x-oauth-scope: analytics:read /v1/projects/{ref}/cli/login-role: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-create-login-role parameters: - name: ref @@ -5348,18 +5609,28 @@ paths: description: Failed to create login role security: - bearer: [] - summary: '[Beta] Create a login role for CLI with temporary password' + summary: Create a login role for CLI with temporary password tags: - Database x-badges: - name: 'OAuth scope: database:write' position: after + - name: Beta + position: before x-endpoint-owners: - dev-workflows x-fga-permissions: - - database_write x-oauth-scope: database:write + x-scalar-stability: experimental delete: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-delete-login-roles parameters: - name: ref @@ -5389,17 +5660,20 @@ paths: description: Failed to delete login roles security: - bearer: [] - summary: '[Beta] Delete existing login roles used by CLI' + summary: Delete existing login roles used by CLI tags: - Database x-badges: - name: 'OAuth scope: database:write' position: after + - name: Beta + position: before x-endpoint-owners: - dev-workflows x-fga-permissions: - - database_write x-oauth-scope: database:write + x-scalar-stability: experimental /v1/projects/{ref}/database/migrations: get: operationId: v1-list-migration-history @@ -5686,6 +5960,13 @@ paths: x-oauth-scope: database:write /v1/projects/{ref}/database/query: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-run-a-query parameters: - name: ref @@ -5717,21 +5998,27 @@ paths: description: Failed to run sql query security: - bearer: [] - summary: '[Beta] Run sql query' + summary: Run sql query tags: - Database x-badges: - name: 'OAuth scope: database:write' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane x-fga-permissions: - - database_read - - database_write x-oauth-scope: database:write + x-scalar-stability: experimental /v1/projects/{ref}/database/query/read-only: post: - description: All entity references must be schema qualified. + description: |- + All entity references must be schema qualified. + + This endpoint is currently in its **Beta** stage. operationId: v1-read-only-query parameters: - name: ref @@ -5763,19 +6050,29 @@ paths: description: Failed to run read-only sql query security: - bearer: [] - summary: '[Beta] Run a sql query as supabase_read_only_user' + summary: Run a sql query as supabase_read_only_user tags: - Database x-badges: - name: 'OAuth scope: database:read' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane x-fga-permissions: - - database_read x-oauth-scope: database:read + x-scalar-stability: experimental /v1/projects/{ref}/database/webhooks/enable: post: + description: >- + This is an **experimental** endpoint. It is subject to change or removal + in future versions. Use it with caution, as it may not remain supported + or stable. + + + This endpoint is currently in its **Beta** stage. operationId: v1-enable-database-webhook parameters: - name: ref @@ -5801,21 +6098,23 @@ paths: description: Failed to enable Database Webhooks on the project security: - bearer: [] - summary: '[Beta] Enables Database Webhooks on the project' + summary: Enables Database Webhooks on the project tags: - Database x-badges: - name: 'OAuth scope: database:write' position: after + - name: Beta + position: before x-endpoint-owners: - control-plane - infra x-fga-permissions: - - database_webhooks_config_write x-oauth-scope: database:write + x-scalar-stability: experimental /v1/projects/{ref}/database/context: get: - deprecated: true description: >- This is an **experimental** endpoint. It is subject to change or removal in future versions. Use it with caution, as it may not remain supported @@ -5858,6 +6157,7 @@ paths: x-fga-permissions: - - database_read x-oauth-scope: projects:read + x-scalar-stability: experimental /v1/projects/{ref}/database/password: patch: operationId: v1-update-database-password @@ -9381,7 +9681,6 @@ components: sql: type: string required: - - schema_version - sql required: - id @@ -15383,7 +15682,7 @@ components: max_concurrent_users: type: integer minimum: 1 - maximum: 50000 + maximum: 300000 description: Sets maximum number of concurrent users rate limit nullable: true max_events_per_second: @@ -15465,7 +15764,7 @@ components: max_concurrent_users: type: integer minimum: 1 - maximum: 50000 + maximum: 300000 description: Sets maximum number of concurrent users rate limit max_events_per_second: type: integer diff --git a/apps/cli-go/pkg/api/client.gen.go b/apps/cli-go/pkg/api/client.gen.go index a5bdc8b06f..179c078bc5 100644 --- a/apps/cli-go/pkg/api/client.gen.go +++ b/apps/cli-go/pkg/api/client.gen.go @@ -219,7 +219,7 @@ type ClientInterface interface { V1GetProjectLogs(ctx context.Context, ref string, params *V1GetProjectLogsParams, reqEditors ...RequestEditorFn) (*http.Response, error) // V1GetProjectLogsAll request - V1GetProjectLogsAll(ctx context.Context, ref string, params *V1GetProjectLogsAllParams, reqEditors ...RequestEditorFn) (*http.Response, error) + V1GetProjectLogsAll(ctx context.Context, ref string, reqEditors ...RequestEditorFn) (*http.Response, error) // V1ScrapeProjectMetrics request V1ScrapeProjectMetrics(ctx context.Context, ref string, reqEditors ...RequestEditorFn) (*http.Response, error) @@ -1273,8 +1273,8 @@ func (c *Client) V1GetProjectLogs(ctx context.Context, ref string, params *V1Get return c.Client.Do(req) } -func (c *Client) V1GetProjectLogsAll(ctx context.Context, ref string, params *V1GetProjectLogsAllParams, reqEditors ...RequestEditorFn) (*http.Response, error) { - req, err := NewV1GetProjectLogsAllRequest(c.Server, ref, params) +func (c *Client) V1GetProjectLogsAll(ctx context.Context, ref string, reqEditors ...RequestEditorFn) (*http.Response, error) { + req, err := NewV1GetProjectLogsAllRequest(c.Server, ref) if err != nil { return nil, err } @@ -5357,7 +5357,7 @@ func NewV1GetProjectLogsRequest(server string, ref string, params *V1GetProjectL } // NewV1GetProjectLogsAllRequest generates requests for V1GetProjectLogsAll -func NewV1GetProjectLogsAllRequest(server string, ref string, params *V1GetProjectLogsAllParams) (*http.Request, error) { +func NewV1GetProjectLogsAllRequest(server string, ref string) (*http.Request, error) { var err error var pathParam0 string @@ -5382,57 +5382,6 @@ func NewV1GetProjectLogsAllRequest(server string, ref string, params *V1GetProje return nil, err } - if params != nil { - // queryValues collects non-styled parameters (passthrough, JSON) - // that are safe to round-trip through url.Values.Encode(). - queryValues := queryURL.Query() - // rawQueryFragments collects pre-encoded query fragments from - // styled parameters, preserving literal commas as delimiters - // per the OpenAPI spec (e.g. "color=blue,black,brown"). - var rawQueryFragments []string - - if params.Sql != nil { - - if queryFrag, err := runtime.StyleParamWithOptions("form", true, "sql", *params.Sql, runtime.StyleParamOptions{ParamLocation: runtime.ParamLocationQuery, Type: "string", Format: ""}); err != nil { - return nil, err - } else { - for _, qp := range strings.Split(queryFrag, "&") { - rawQueryFragments = append(rawQueryFragments, qp) - } - } - - } - - if params.IsoTimestampStart != nil { - - if queryFrag, err := runtime.StyleParamWithOptions("form", true, "iso_timestamp_start", *params.IsoTimestampStart, runtime.StyleParamOptions{ParamLocation: runtime.ParamLocationQuery, Type: "string", Format: "date-time"}); err != nil { - return nil, err - } else { - for _, qp := range strings.Split(queryFrag, "&") { - rawQueryFragments = append(rawQueryFragments, qp) - } - } - - } - - if params.IsoTimestampEnd != nil { - - if queryFrag, err := runtime.StyleParamWithOptions("form", true, "iso_timestamp_end", *params.IsoTimestampEnd, runtime.StyleParamOptions{ParamLocation: runtime.ParamLocationQuery, Type: "string", Format: "date-time"}); err != nil { - return nil, err - } else { - for _, qp := range strings.Split(queryFrag, "&") { - rawQueryFragments = append(rawQueryFragments, qp) - } - } - - } - - if encoded := queryValues.Encode(); encoded != "" { - rawQueryFragments = append(rawQueryFragments, encoded) - } - queryURL.RawQuery = strings.Join(rawQueryFragments, "&") - } - req, err := http.NewRequest(http.MethodGet, queryURL.String(), nil) if err != nil { return nil, err @@ -11664,7 +11613,7 @@ type ClientWithResponsesInterface interface { V1GetProjectLogsWithResponse(ctx context.Context, ref string, params *V1GetProjectLogsParams, reqEditors ...RequestEditorFn) (*V1GetProjectLogsResponse, error) // V1GetProjectLogsAllWithResponse request - V1GetProjectLogsAllWithResponse(ctx context.Context, ref string, params *V1GetProjectLogsAllParams, reqEditors ...RequestEditorFn) (*V1GetProjectLogsAllResponse, error) + V1GetProjectLogsAllWithResponse(ctx context.Context, ref string, reqEditors ...RequestEditorFn) (*V1GetProjectLogsAllResponse, error) // V1ScrapeProjectMetricsWithResponse request V1ScrapeProjectMetricsWithResponse(ctx context.Context, ref string, reqEditors ...RequestEditorFn) (*V1ScrapeProjectMetricsResponse, error) @@ -13242,7 +13191,6 @@ func (r V1GetProjectLogsResponse) ContentType() string { type V1GetProjectLogsAllResponse struct { Body []byte HTTPResponse *http.Response - JSON200 *AnalyticsResponseOutput } // Status returns HTTPResponse.Status @@ -17636,8 +17584,8 @@ func (c *ClientWithResponses) V1GetProjectLogsWithResponse(ctx context.Context, } // V1GetProjectLogsAllWithResponse request returning *V1GetProjectLogsAllResponse -func (c *ClientWithResponses) V1GetProjectLogsAllWithResponse(ctx context.Context, ref string, params *V1GetProjectLogsAllParams, reqEditors ...RequestEditorFn) (*V1GetProjectLogsAllResponse, error) { - rsp, err := c.V1GetProjectLogsAll(ctx, ref, params, reqEditors...) +func (c *ClientWithResponses) V1GetProjectLogsAllWithResponse(ctx context.Context, ref string, reqEditors ...RequestEditorFn) (*V1GetProjectLogsAllResponse, error) { + rsp, err := c.V1GetProjectLogsAll(ctx, ref, reqEditors...) if err != nil { return nil, err } @@ -20114,16 +20062,6 @@ func ParseV1GetProjectLogsAllResponse(rsp *http.Response) (*V1GetProjectLogsAllR HTTPResponse: rsp, } - switch { - case strings.Contains(rsp.Header.Get("Content-Type"), "json") && rsp.StatusCode == 200: - var dest AnalyticsResponseOutput - if err := json.Unmarshal(bodyBytes, &dest); err != nil { - return nil, err - } - response.JSON200 = &dest - - } - return response, nil } diff --git a/apps/cli-go/pkg/api/types.gen.go b/apps/cli-go/pkg/api/types.gen.go index 959ea7ed62..844fc4fa66 100644 --- a/apps/cli-go/pkg/api/types.gen.go +++ b/apps/cli-go/pkg/api/types.gen.go @@ -7692,9 +7692,9 @@ type SnippetResponseOutput struct { Content struct { // Favorite Deprecated: Rely on root-level favorite property instead. // Deprecated: this property has been marked as deprecated upstream, but no `x-deprecated-reason` was set - Favorite *bool `json:"favorite,omitempty"` - SchemaVersion string `json:"schema_version"` - Sql string `json:"sql"` + Favorite *bool `json:"favorite,omitempty"` + SchemaVersion *string `json:"schema_version,omitempty"` + Sql string `json:"sql"` } `json:"content"` Description nullable.Nullable[string] `json:"description"` Favorite bool `json:"favorite"` @@ -9158,14 +9158,6 @@ type V1GetProjectLogsParams struct { IsoTimestampEnd *time.Time `form:"iso_timestamp_end,omitempty" json:"iso_timestamp_end,omitempty"` } -// V1GetProjectLogsAllParams defines parameters for V1GetProjectLogsAll. -type V1GetProjectLogsAllParams struct { - // Sql Custom SQL query to execute on the logs. See [querying logs](https://supabase.com/docs/guides/monitoring-and-debugging/logs#querying-with-the-logs-explorer) for more details. - Sql *string `form:"sql,omitempty" json:"sql,omitempty"` - IsoTimestampStart *time.Time `form:"iso_timestamp_start,omitempty" json:"iso_timestamp_start,omitempty"` - IsoTimestampEnd *time.Time `form:"iso_timestamp_end,omitempty" json:"iso_timestamp_end,omitempty"` -} - // V1GetProjectUsageApiCountParams defines parameters for V1GetProjectUsageApiCount. type V1GetProjectUsageApiCountParams struct { Interval *V1GetProjectUsageApiCountParamsInterval `form:"interval,omitempty" json:"interval,omitempty"` From 8d19c5c4a6693a01c2adbe5e74b32b7c5b15b593 Mon Sep 17 00:00:00 2001 From: Julien Goux Date: Thu, 24 Sep 2026 20:38:21 +0200 Subject: [PATCH 2/2] fix(stack): recover native pooler startup port collisions --- packages/stack/src/services/Pooler.ts | 13 + .../stack/src/services/Pooler.unit.test.ts | 34 + .../ProcessRecipe.integration.test.ts | 619 +++++++++++++++++- packages/stack/src/services/ProcessRecipe.ts | 479 +++++++++++--- 4 files changed, 1049 insertions(+), 96 deletions(-) create mode 100644 packages/stack/src/services/Pooler.unit.test.ts diff --git a/packages/stack/src/services/Pooler.ts b/packages/stack/src/services/Pooler.ts index abe1122295..53b224b279 100644 --- a/packages/stack/src/services/Pooler.ts +++ b/packages/stack/src/services/Pooler.ts @@ -79,6 +79,18 @@ const nativeStartupEnvironment: NonNullable["nativeS })), ); +const nativeReadinessOutput: NonNullable["nativeReadinessOutput"]> = ( + line, + endpoints, +) => { + const http = endpoints.get("http"); + if (http === undefined || /\bfailed\b/i.test(line)) return false; + return ( + line.includes("Running SupavisorWeb.Endpoint") && + new RegExp(`:${http.port}\\s+\\(http\\)(?:\\s|$)`).test(line) + ); +}; + export const makeSpec = (): ProcessRecipeSpec => ({ service: "pooler", executable: "bin/server", @@ -86,6 +98,7 @@ export const makeSpec = (): ProcessRecipeSpec => ({ healthPath: "/api/health", env: environment, nativeStartupEnv: nativeStartupEnvironment, + nativeReadinessOutput, args: (_creation, _endpoints, context) => Effect.succeed(context.container ? ["-s", "-g", "--", "/app/bin/server"] : ["start"]), mounts: () => Effect.succeed([]), diff --git a/packages/stack/src/services/Pooler.unit.test.ts b/packages/stack/src/services/Pooler.unit.test.ts new file mode 100644 index 0000000000..e813f5a6c9 --- /dev/null +++ b/packages/stack/src/services/Pooler.unit.test.ts @@ -0,0 +1,34 @@ +import { describe, expect, it } from "@effect/vitest"; +import { Effect } from "effect"; +import type { ServiceEndpoint } from "./Recipe.ts"; +import * as Pooler from "./Pooler.ts"; + +describe("native Pooler readiness output", () => { + it.effect("requires the selected HTTP port on a successful endpoint bind line", () => + Effect.sync(() => { + const readinessOutput = Pooler.makeSpec().nativeReadinessOutput; + const endpoints: ReadonlyMap = new Map([ + ["http", { kind: "tcp" as const, host: "127.0.0.1", port: 4000 }], + ]); + + expect( + readinessOutput?.( + "Running SupavisorWeb.Endpoint at http://127.0.0.1:4000 (http)", + endpoints, + ), + ).toBe(true); + expect( + readinessOutput?.( + "Running SupavisorWeb.Endpoint at http://127.0.0.1:4000 (http) failed, port already in use", + endpoints, + ), + ).toBe(false); + expect( + readinessOutput?.( + "Running SupavisorWeb.Endpoint at http://127.0.0.1:40001 (http)", + endpoints, + ), + ).toBe(false); + }), + ); +}); diff --git a/packages/stack/src/services/ProcessRecipe.integration.test.ts b/packages/stack/src/services/ProcessRecipe.integration.test.ts index 5f9925193f..83068eb1a2 100644 --- a/packages/stack/src/services/ProcessRecipe.integration.test.ts +++ b/packages/stack/src/services/ProcessRecipe.integration.test.ts @@ -1,13 +1,38 @@ import { NodeHttpClient, NodeServices } from "@effect/platform-node"; import { describe, expect, it } from "@effect/vitest"; -import { Crypto, Deferred, Effect, Fiber, FileSystem, Layer, Path, Sink, Stream } from "effect"; +import { + Crypto, + Deferred, + Effect, + Exit, + Fiber, + FileSystem, + Layer, + Path, + PlatformError, + Ref, + Scope, + Schema, + Sink, + Stream, +} from "effect"; import { TestClock } from "effect/testing"; -import { HttpClient } from "effect/unstable/http"; +import { HttpClient, HttpClientRequest } from "effect/unstable/http"; import { ChildProcessSpawner } from "effect/unstable/process"; +import type { ChildProcessSpawner as ChildProcessSpawnerService } from "effect/unstable/process/ChildProcessSpawner"; +// oxlint-disable-next-line effecttsgo/node-builtin-import -- the collision fixture owns a local HTTP listener. +import * as NodeHttp from "node:http"; +import { + makeArtifactStore, + type ArtifactRequest, + type ArtifactSource, +} from "../preparation/ArtifactStore.ts"; +import { PreparationError } from "../preparation/Errors.ts"; import type { ContainerRuntime } from "../runtime/Container.ts"; import { makeService } from "../Service.ts"; import { makeProcessRecipe } from "./ProcessRecipe.ts"; import * as Realtime from "./Realtime.ts"; +import * as Pooler from "./Pooler.ts"; const encode = (text: string) => new TextEncoder().encode(text); @@ -69,6 +94,122 @@ const realtimeService = Effect.fn(function* (container: ContainerRuntime) { const platform = Layer.merge(NodeServices.layer, NodeHttpClient.layerNodeHttp); +const nativePoolerArtifact = Effect.fn(function* (cacheRoot: string) { + const platformName = `${process.platform}-${process.arch}`; + const target = + platformName === "darwin-arm64" + ? "darwin-arm64" + : platformName === "linux-x64" + ? "linux-amd64" + : platformName === "linux-arm64" + ? "linux-arm64" + : undefined; + if (target === undefined) return yield* Effect.fail(`Unsupported test platform: ${platformName}`); + + const request: ArtifactRequest = { + key: `slim-services/pooler/v2.9.12/${target}`, + requiredRuntimePaths: ["bin/server", "bin/prepare", "bin/provision-tenant"], + executablePath: "bin/server", + }; + const fs = yield* FileSystem.FileSystem; + const server = + `#!${process.execPath}\nconst http = require("node:http");\n` + + `const port = Number(process.env.PORT);\n` + + `if (process.env.TENANT_ID === "unrelated-failure") { console.error("unrelated startup failure"); process.exit(1); }\n` + + `if (process.env.TENANT_ID === "retry-exhaustion") {\n` + + `const blocker = http.createServer();\n` + + `blocker.listen(port, "127.0.0.1", () => {\n` + + `const failed = http.createServer();\n` + + `failed.on("error", () => { console.error("port already in use " + "x".repeat(2000)); process.exit(1); });\n` + + `failed.listen(port, "127.0.0.1");\n` + + `});\nreturn;\n}\n` + + `const server = http.createServer((_request, response) => response.end("owned-fixture:" + port));\n` + + `server.on("error", (error) => { console.error(error); process.exit(1); });\n` + + `server.listen(port, "127.0.0.1", () => {\n` + + `console.log("Running SupavisorWeb.Endpoint at 127.0.0.1:" + port + " (http)");\n` + + `});\n`; + const oneShot = `#!${process.execPath}\nprocess.exit(0);\n`; + const source: ArtifactSource = { + checksum: () => Effect.succeed("0".repeat(64)), + materialize: (_entry, destination) => + Effect.gen(function* () { + yield* fs.makeDirectory(`${destination}/bin`, { recursive: true }); + yield* fs.writeFileString(`${destination}/bin/server`, server); + yield* fs.writeFileString(`${destination}/bin/prepare`, oneShot); + yield* fs.writeFileString(`${destination}/bin/provision-tenant`, oneShot); + yield* fs.chmod(`${destination}/bin/prepare`, 0o755); + yield* fs.chmod(`${destination}/bin/provision-tenant`, 0o755); + }).pipe( + Effect.mapError( + (cause) => + new PreparationError({ + message: `Unable to write native Pooler fixture: ${cause.message}`, + cause, + }), + ), + ), + }; + const store = yield* makeArtifactStore({ cacheRoot, source }); + yield* store.prepare(request); +}); + +const NativeLaunchPayload = Schema.Struct({ + executable: Schema.String, + args: Schema.optionalKey(Schema.Array(Schema.String)), + env: Schema.optionalKey(Schema.Record(Schema.String, Schema.String)), +}); + +const countingSpawner = ( + spawner: ChildProcessSpawnerService["Service"], + mainLaunches: Ref.Ref, + startupLaunches: Ref.Ref>, + splitStderr = false, +): ChildProcessSpawnerService["Service"] => ({ + ...spawner, + spawn: (command) => + spawner.spawn(command).pipe( + Effect.map((handle) => ({ + ...handle, + stderr: splitStderr + ? handle.stderr.pipe( + Stream.flatMap((bytes) => + Stream.fromIterable( + Array.from({ length: Math.ceil(bytes.length / 128) }, (_, index) => + bytes.slice(index * 128, (index + 1) * 128), + ), + ), + ), + ) + : handle.stderr, + getInputFd: (fd: number) => { + const sink = handle.getInputFd(fd); + if (fd !== 4) return sink; + return Sink.mapInputEffect(sink, (bytes) => + Effect.gen(function* () { + const payload = yield* Schema.decodeEffect( + Schema.fromJsonString(NativeLaunchPayload), + )(new TextDecoder().decode(bytes)).pipe(Effect.orDie); + if ( + payload.executable.endsWith("/bin/server") && + (payload.args ?? []).includes("start") + ) + yield* Ref.update(mainLaunches, (count) => count + 1); + if ( + payload.executable.endsWith("/bin/prepare") || + payload.executable.endsWith("/bin/provision-tenant") + ) + yield* Ref.update(startupLaunches, (launches) => [ + ...launches, + payload.executable.endsWith("/bin/prepare") ? "prepare" : "provision-tenant", + ]); + return bytes; + }), + ); + }, + })), + ), +}); + describe("process recipe startup", () => { it.effect("reports the startup process's recent stdout and stderr when it exits non-zero", () => Effect.scoped( @@ -94,6 +235,480 @@ describe("process recipe startup", () => { ).pipe(Effect.provide(platform)), ); + it.live("recovers when a healthy competing listener claims Pooler's selected HTTP port", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const client = yield* HttpClient.HttpClient; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "process-recipe-port-race-" }); + const cacheRoot = path.join(root, "cache"); + yield* nativePoolerArtifact(cacheRoot); + const collisionPort = yield* Ref.make(undefined); + const collisionInstalled = yield* Ref.make(false); + const competitor = yield* Ref.make(undefined); + const testScope = yield* Scope.fork(yield* Effect.scope, "sequential"); + yield* Effect.addFinalizer(() => Scope.close(testScope, Exit.void)); + const creation: Pooler.Creation = { + service: "pooler", + config: { + databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", + jwtSecret: "pooler-collision-test-secret-with-more-than-32-characters", + tenant: "collision-test", + poolMode: "transaction", + }, + }; + const interceptingSpawner: typeof spawner = { + ...spawner, + spawn: (command) => + spawner.spawn(command).pipe( + Effect.map((handle) => ({ + ...handle, + getInputFd: (fd: number) => { + const sink = handle.getInputFd(fd); + if (fd !== 4) return sink; + return Sink.mapInputEffect(sink, (bytes) => + Effect.gen(function* () { + const payload = yield* Schema.decodeEffect( + Schema.fromJsonString(NativeLaunchPayload), + )(new TextDecoder().decode(bytes)).pipe(Effect.orDie); + if ( + !payload.executable.endsWith("/bin/server") || + !(payload.args ?? []).includes("start") || + (yield* Ref.get(collisionInstalled)) + ) + return bytes; + + const port = Number(payload.env?.PORT); + const server = NodeHttp.createServer((_request, response) => + response.end("competing-listener"), + ); + yield* Effect.addFinalizer(() => + server.listening + ? Effect.callback((resume) => { + server.close(() => resume(Effect.void)); + return Effect.void; + }) + : Effect.void, + ).pipe(Scope.provide(testScope)); + yield* Effect.callback((resume) => { + const onError = (cause: Error) => resume(Effect.die(cause)); + server.once("error", onError); + server.listen(port, "127.0.0.1", () => resume(Effect.void)); + return Effect.sync(() => server.off("error", onError)); + }); + yield* Ref.set(collisionPort, port); + yield* Ref.set(competitor, server); + yield* Ref.set(collisionInstalled, true); + return bytes; + }), + ); + }, + })), + ), + }; + const recipe = yield* makeProcessRecipe( + creation, + { + stackId: "process-recipe-port-race", + instanceId: "instance", + root, + cacheRoot, + runtime: "native", + platform: { os: process.platform, arch: process.arch }, + }, + { fs, path, crypto, client, spawner: interceptingSpawner, container: undefined }, + Pooler.makeSpec(), + ); + if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); + const scope = yield* Scope.fork(testScope, "sequential"); + const runtime = yield* recipe.definition.launch({ + id: "pooler", + config: creation, + scope, + }); + const endpoints = yield* Ref.get(recipe.endpoints); + const endpoint = endpoints.get("http"); + const collided = yield* Ref.get(collisionPort); + expect(collided).toBeDefined(); + expect(endpoint?.port).not.toBe(collided); + const oldPort = yield* Ref.get(collisionPort); + const competingServer = yield* Ref.get(competitor); + expect(competingServer?.listening).toBe(true); + const competingResponse = yield* client.execute( + HttpClientRequest.get(`http://127.0.0.1:${oldPort}/api/health`), + ); + expect(yield* competingResponse.text).toBe("competing-listener"); + const response = yield* client.execute( + HttpClientRequest.get(`http://${endpoint?.host}:${endpoint?.port}/api/health`), + ); + expect(yield* response.text).toBe(`owned-fixture:${endpoint?.port}`); + expect(yield* Ref.get(collisionInstalled)).toBe(true); + yield* runtime.stop; + }), + ).pipe(Effect.provide(platform)), + ); + + it.live("stops after three consecutive native Pooler port collisions", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const client = yield* HttpClient.HttpClient; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const root = yield* fs.makeTempDirectoryScoped({ + prefix: "process-recipe-port-exhaustion-", + }); + const cacheRoot = path.join(root, "cache"); + yield* nativePoolerArtifact(cacheRoot); + const mainLaunches = yield* Ref.make(0); + const startupLaunches = yield* Ref.make>([]); + const creation: Pooler.Creation = { + service: "pooler", + config: { + databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", + jwtSecret: "pooler-collision-test-secret-with-more-than-32-characters", + tenant: "retry-exhaustion", + poolMode: "transaction", + }, + }; + const recipe = yield* makeProcessRecipe( + creation, + { + stackId: "process-recipe-port-exhaustion", + instanceId: "instance", + root, + cacheRoot, + runtime: "native", + platform: { os: process.platform, arch: process.arch }, + }, + { + fs, + path, + crypto, + client, + spawner: countingSpawner(spawner, mainLaunches, startupLaunches, true), + container: undefined, + }, + Pooler.makeSpec(), + ); + const service = yield* makeService(recipe.definition, { id: "pooler", config: creation }); + const subscribed = yield* Deferred.make(); + const exitObserved = yield* Deferred.make(); + yield* service.observation.pipe( + Stream.tap(() => Deferred.succeed(subscribed, undefined)), + Stream.runForEach((observation) => + observation.exit === undefined + ? Effect.void + : Deferred.succeed(exitObserved, undefined), + ), + Effect.forkScoped, + ); + yield* Deferred.await(subscribed); + yield* service.start; + yield* Deferred.await(exitObserved); + expect(yield* Ref.get(mainLaunches)).toBe(3); + expect(yield* Ref.get(startupLaunches)).toEqual(["prepare", "provision-tenant"]); + expect(yield* Ref.get(recipe.endpoints)).toEqual(new Map()); + const ready = yield* Effect.exit(service.ready); + expect(Exit.isFailure(ready)).toBe(true); + if (Exit.isFailure(ready)) + expect(ready.cause.toString()).toContain("native port collision"); + const observation = yield* service.get; + expect(observation.error?.operation).toBe("launch"); + expect(observation.error?.message).toContain("native port collision"); + expect(observation.exit).toBeDefined(); + if (observation.exit !== undefined) { + expect(Exit.isFailure(observation.exit)).toBe(true); + if (Exit.isFailure(observation.exit)) + expect(observation.exit.cause.toString()).toContain("native port collision"); + } + yield* service.stop; + expect((yield* service.get).lifecycle).toBe("stopped"); + expect(yield* Ref.get(recipe.endpoints)).toEqual(new Map()); + }), + ).pipe(Effect.provide(platform)), + ); + + it.effect("preserves the collision diagnostic when retries cross the readiness deadline", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const client = yield* HttpClient.HttpClient; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "process-recipe-deadline-" }); + const cacheRoot = path.join(root, "cache"); + yield* nativePoolerArtifact(cacheRoot); + const exitReached = yield* Deferred.make(); + const clockAdvanced = yield* Deferred.make(); + const deadlineSpawner: typeof spawner = { + ...spawner, + spawn: (command) => + spawner.spawn(command).pipe( + Effect.map((handle) => { + let mainProcess = false; + return { + ...handle, + exitCode: handle.exitCode.pipe( + Effect.tap(() => + mainProcess ? Deferred.succeed(exitReached, undefined) : Effect.void, + ), + ), + stderr: Stream.suspend(() => + mainProcess + ? handle.stderr.pipe( + Stream.concat( + Stream.fromEffectDrain( + Deferred.await(exitReached).pipe( + Effect.andThen(TestClock.adjust("61 seconds")), + Effect.tap(() => Deferred.succeed(clockAdvanced, undefined)), + ), + ), + ), + ) + : handle.stderr, + ), + getInputFd: (fd: number) => { + const sink = handle.getInputFd(fd); + if (fd !== 4) return sink; + return Sink.mapInputEffect(sink, (bytes) => + Effect.gen(function* () { + const payload = yield* Schema.decodeEffect( + Schema.fromJsonString(NativeLaunchPayload), + )(new TextDecoder().decode(bytes)).pipe(Effect.orDie); + if ( + payload.executable.endsWith("/bin/server") && + (payload.args ?? []).includes("start") + ) + mainProcess = true; + return bytes; + }), + ); + }, + }; + }), + ), + }; + const creation: Pooler.Creation = { + service: "pooler", + config: { + databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", + jwtSecret: "pooler-collision-test-secret-with-more-than-32-characters", + tenant: "retry-exhaustion", + poolMode: "transaction", + }, + }; + const recipe = yield* makeProcessRecipe( + creation, + { + stackId: "process-recipe-deadline", + instanceId: "instance", + root, + cacheRoot, + runtime: "native", + platform: { os: process.platform, arch: process.arch }, + }, + { fs, path, crypto, client, spawner: deadlineSpawner, container: undefined }, + Pooler.makeSpec(), + ); + if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); + const scope = yield* Scope.fork(yield* Effect.scope, "sequential"); + const launched = recipe.definition + .launch({ id: "pooler", config: creation, scope }) + .pipe(Effect.forkChild); + const fiber = yield* launched; + yield* Deferred.await(clockAdvanced); + const runtime = yield* Fiber.join(fiber); + const failure = yield* Effect.flip(runtime.health); + expect(failure.operation).toBe("launch"); + expect(failure.message).toContain("native port collision"); + expect(failure.message).not.toContain("readiness timed out"); + const exit = yield* runtime.exit; + expect(Exit.isFailure(exit)).toBe(true); + if (Exit.isFailure(exit)) expect(exit.cause.toString()).toContain("native port collision"); + expect(yield* Ref.get(recipe.endpoints)).toEqual(new Map()); + yield* runtime.stop; + }), + ).pipe(Effect.provide(platform)), + ); + + it.live("does not retry an unrelated native Pooler startup failure", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const client = yield* HttpClient.HttpClient; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const root = yield* fs.makeTempDirectoryScoped({ + prefix: "process-recipe-unrelated-failure-", + }); + const cacheRoot = path.join(root, "cache"); + yield* nativePoolerArtifact(cacheRoot); + const mainLaunches = yield* Ref.make(0); + const startupLaunches = yield* Ref.make>([]); + const creation: Pooler.Creation = { + service: "pooler", + config: { + databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", + jwtSecret: "pooler-collision-test-secret-with-more-than-32-characters", + tenant: "unrelated-failure", + poolMode: "transaction", + }, + }; + const recipe = yield* makeProcessRecipe( + creation, + { + stackId: "process-recipe-unrelated-failure", + instanceId: "instance", + root, + cacheRoot, + runtime: "native", + platform: { os: process.platform, arch: process.arch }, + }, + { + fs, + path, + crypto, + client, + spawner: countingSpawner(spawner, mainLaunches, startupLaunches), + container: undefined, + }, + Pooler.makeSpec(), + ); + if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); + const scope = yield* Scope.fork(yield* Effect.scope, "sequential"); + const runtime = yield* recipe.definition.launch({ id: "pooler", config: creation, scope }); + expect(yield* Ref.get(mainLaunches)).toBe(1); + expect(yield* Ref.get(startupLaunches)).toEqual(["prepare", "provision-tenant"]); + expect(yield* Ref.get(recipe.endpoints)).toEqual(new Map()); + const health = yield* Effect.exit(runtime.health); + expect(Exit.isFailure(health)).toBe(true); + if (Exit.isFailure(health)) + expect(health.cause.toString()).toContain("unrelated startup failure"); + const exit = yield* runtime.exit; + expect(Exit.isFailure(exit)).toBe(true); + if (Exit.isFailure(exit)) + expect(exit.cause.toString()).toContain("unrelated startup failure"); + yield* runtime.stop; + }), + ).pipe(Effect.provide(platform)), + ); + + it.live("cleans up a native process when exit observation fails before output ends", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const crypto = yield* Crypto.Crypto; + const client = yield* HttpClient.HttpClient; + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "process-recipe-exit-failure-" }); + const cacheRoot = path.join(root, "cache"); + yield* nativePoolerArtifact(cacheRoot); + const port = yield* Ref.make(undefined); + let mainProcess = false; + const failingSpawner: typeof spawner = { + ...spawner, + spawn: (command) => + spawner.spawn(command).pipe( + Effect.map((handle) => ({ + ...handle, + stdout: Stream.suspend(() => (mainProcess ? Stream.never : handle.stdout)), + stderr: Stream.suspend(() => (mainProcess ? Stream.never : handle.stderr)), + exitCode: Effect.suspend(() => + mainProcess + ? Effect.fail( + PlatformError.systemError({ + _tag: "PermissionDenied", + module: "ChildProcess", + method: "exitCode", + description: "injected native exit observation failure", + cause: { code: "EIO" }, + }), + ) + : handle.exitCode, + ), + getInputFd: (fd: number) => { + const sink = handle.getInputFd(fd); + if (fd !== 4) return sink; + return Sink.mapInputEffect(sink, (bytes) => + Effect.gen(function* () { + const payload = yield* Schema.decodeEffect( + Schema.fromJsonString(NativeLaunchPayload), + )(new TextDecoder().decode(bytes)).pipe(Effect.orDie); + if ( + payload.executable.endsWith("/bin/server") && + (payload.args ?? []).includes("start") + ) { + mainProcess = true; + yield* Ref.set(port, Number(payload.env?.PORT)); + } + return bytes; + }), + ); + }, + })), + ), + }; + const creation: Pooler.Creation = { + service: "pooler", + config: { + databaseUrl: "postgresql://postgres:postgres@127.0.0.1:5432/postgres", + jwtSecret: "pooler-collision-test-secret-with-more-than-32-characters", + tenant: "exit-observation-failure", + poolMode: "transaction", + }, + }; + const recipe = yield* makeProcessRecipe( + creation, + { + stackId: "process-recipe-exit-failure", + instanceId: "instance", + root, + cacheRoot, + runtime: "native", + platform: { os: process.platform, arch: process.arch }, + }, + { fs, path, crypto, client, spawner: failingSpawner, container: undefined }, + Pooler.makeSpec(), + ); + if (recipe.definition.prepare !== undefined) yield* recipe.definition.prepare(creation); + const scope = yield* Scope.fork(yield* Effect.scope, "sequential"); + const runtime = yield* recipe.definition.launch({ id: "pooler", config: creation, scope }); + const health = yield* Effect.exit(runtime.health); + expect(Exit.isFailure(health)).toBe(true); + const exit = yield* runtime.exit; + expect(Exit.isFailure(exit)).toBe(true); + if (Exit.isFailure(exit)) + expect(exit.cause.toString()).toContain("injected native exit observation failure"); + expect(yield* Ref.get(recipe.endpoints)).toEqual(new Map()); + + const portToProbe = yield* Ref.get(port); + expect(portToProbe).toBeDefined(); + if (portToProbe !== undefined) { + const probe = NodeHttp.createServer(); + const bind = Effect.callback((resume) => { + const onError = (cause: Error) => resume(Effect.die(cause)); + probe.once("error", onError); + probe.listen(portToProbe, "127.0.0.1", () => resume(Effect.void)); + return Effect.sync(() => probe.off("error", onError)); + }); + const close = Effect.callback((resume) => { + probe.close(() => resume(Effect.void)); + return Effect.void; + }); + yield* bind.pipe(Effect.ensuring(close)); + } + }), + ).pipe(Effect.provide(platform)), + ); + it.effect("keeps the stderr error when later stdout exceeds the tail", () => Effect.scoped( Effect.gen(function* () { diff --git a/packages/stack/src/services/ProcessRecipe.ts b/packages/stack/src/services/ProcessRecipe.ts index 69019f438a..ba3a9540ca 100644 --- a/packages/stack/src/services/ProcessRecipe.ts +++ b/packages/stack/src/services/ProcessRecipe.ts @@ -1,6 +1,8 @@ import { Cause, + Clock, Crypto, + Deferred, Duration, Effect, Exit, @@ -71,6 +73,10 @@ export interface ProcessRecipeSpec, ) => Effect.Effect>, ServiceError>; + readonly nativeReadinessOutput?: ( + line: string, + endpoints: ReadonlyMap, + ) => boolean; readonly mounts: ( creation: C, context: { readonly container: boolean }, @@ -135,48 +141,47 @@ const runtimeFromContainer = (process: ContainerProcess): RuntimeSession => ({ remove: process.remove.pipe(Effect.mapError((cause) => serviceError("remove", cause))), }); +interface NativePortReservation { + readonly port: number; + readonly server: Net.Server; +} + +const closeNativePort = (server: Net.Server): Effect.Effect => + Effect.callback((resume) => { + if (!server.listening) { + resume(Effect.void); + return Effect.void; + } + server.close(() => resume(Effect.void)); + return Effect.void; + }); + const reserveNativePort = Effect.fn("ProcessRecipe.reserveNativePort")( - (requested: number): Effect.Effect => - Effect.scoped( - Effect.gen(function* () { - const allocation = yield* Effect.acquireRelease( - Effect.callback<{ readonly port: number; readonly server: Net.Server }, CatalogError>( - (resume) => { - const server = Net.createServer(); - const onError = (cause: Error) => - resume( - Effect.fail( - catalogError( - "launch", - "Unable to reserve native service port", - undefined, - cause, - ), - ), - ); - server.once("error", onError); - server.listen({ host: "127.0.0.1", port: requested }, () => { - const address = server.address(); - if (address === null || typeof address === "string") { - onError(new Error("Native service port reservation returned no address")); - } else { - resume(Effect.succeed({ port: address.port, server })); - } - }); - return Effect.sync(() => { - server.off("error", onError); - if (server.listening) server.close(); - }); - }, - ), - ({ server }) => - Effect.callback((resume) => { - server.close(() => resume(Effect.void)); - return Effect.void; - }), - ); - return allocation.port; + (requested: number): Effect.Effect => + Effect.acquireRelease( + Effect.callback((resume) => { + const server = Net.createServer((socket) => socket.destroy()); + const onError = (cause: Error) => + resume( + Effect.fail( + catalogError("launch", "Unable to reserve native service port", undefined, cause), + ), + ); + server.once("error", onError); + server.listen({ host: "127.0.0.1", port: requested }, () => { + const address = server.address(); + if (address === null || typeof address === "string") { + onError(new Error("Native service port reservation returned no address")); + } else { + resume(Effect.succeed({ port: address.port, server })); + } + }); + return Effect.sync(() => { + server.off("error", onError); + if (server.listening) server.close(); + }); }), + ({ server }) => closeNativePort(server), ), ); @@ -297,11 +302,100 @@ const startupFailure = ( withRecentOutput(`${service} startup exited with ${result.code}`, result.output), ); +const isAddressInUse = (line: string) => + /eaddrinuse|address already in use|port already in use/i.test(line); + +const terminalSession = ( + error: ServiceError, + endpoints: Ref.Ref>, +) => + ({ + health: Effect.fail(error), + exit: Effect.succeed(Exit.fail(error)), + stop: Effect.void, + remove: Ref.set(endpoints, new Map()), + }) satisfies RuntimeSession; + +const collectNativeOutput = Effect.fn("ProcessRecipe.collectNativeOutput")(function* ( + process: NativeProcess, + logs: PubSub.PubSub, + endpoints: ReadonlyMap, + readinessOutput: + | ((line: string, endpoints: ReadonlyMap) => boolean) + | undefined, + scope: Scope.Closeable, +) { + const bindReady = yield* Deferred.make(); + const bindError = yield* Ref.make(false); + const stdout = yield* Ref.make>([]); + const stderr = yield* Ref.make>([]); + const drained = yield* Deferred.make(); + + const consume = Effect.fnUntraced(function* ( + stream: Stream.Stream, + name: CatalogLog["stream"], + tail: Ref.Ref>, + ) { + const partial = yield* Ref.make(""); + const append = (lines: ReadonlyArray) => + Effect.forEach( + lines, + (line) => + Effect.gen(function* () { + if (readinessOutput?.(line, endpoints) === true) + yield* Deferred.succeed(bindReady, undefined); + if (isAddressInUse(line)) yield* Ref.set(bindError, true); + if (line.trim().length > 0) + yield* Ref.update(tail, (current) => + [...current, clipLine(line)].slice(-startupOutputTailLines), + ); + }), + { discard: true }, + ); + yield* stream.pipe( + Stream.tap((bytes) => PubSub.publish(logs, { stream: name, bytes })), + Stream.decodeText, + Stream.runForEach((text) => + Ref.modify( + partial, + ( + rest, + ): [{ readonly lines: ReadonlyArray; readonly combined: string }, string] => { + const combined = `${rest}${text}`; + const lines = combined.split(/\r?\n/); + const next = lines.pop() ?? ""; + return [{ lines, combined }, next.slice(-(startupOutputLineChars + 1))]; + }, + ).pipe( + Effect.flatMap(({ lines, combined }) => + Effect.gen(function* () { + if (isAddressInUse(combined)) yield* Ref.set(bindError, true); + yield* append(lines); + }), + ), + ), + ), + Effect.ensuring(Ref.get(partial).pipe(Effect.flatMap((rest) => append([rest])))), + ); + }); + + const drain = Effect.all( + [consume(process.stdout, "stdout", stdout), consume(process.stderr, "stderr", stderr)], + { concurrency: "unbounded", discard: true }, + ).pipe( + Effect.ensuring(Deferred.succeed(drained, undefined)), + Effect.catch((cause) => Effect.logError(cause)), + ); + yield* Effect.forkIn(drain, scope); + return { bindReady, bindError, stdout, stderr, drained }; +}); + const readiness = Effect.fn("ProcessRecipe.readiness")( ( client: HttpClient.HttpClient, endpoint: ServiceEndpoint, path: string, + timeout: Duration.Input = "60 seconds", ): Effect.Effect => client.execute(HttpClientRequest.get(`http://${endpoint.host}:${endpoint.port}${path}`)).pipe( Effect.flatMap((response) => @@ -312,7 +406,7 @@ const readiness = Effect.fn("ProcessRecipe.readiness")( ), ), Effect.retry({ schedule: Schedule.spaced("250 millis") }), - Effect.timeout("60 seconds"), + Effect.timeout(timeout), Effect.mapError((cause) => serviceError("health", cause)), Effect.asVoid, ), @@ -365,20 +459,9 @@ export const makeProcessRecipe = readonly config: C; readonly scope: Scope.Closeable; }) { - const desired = new Map(); - for (const [name, port] of Object.entries(spec.ports)) { - if (spec.enabledPort !== undefined && !spec.enabledPort(context.config, name)) continue; - desired.set(name, { - kind: "tcp", - host: "127.0.0.1", - port: - options.runtime === "native" - ? yield* reserveNativePort(0).pipe( - Effect.mapError((cause) => serviceError("launch", cause)), - ) - : port, - }); - } + const portNames = Object.entries(spec.ports).filter( + ([name]) => spec.enabledPort === undefined || spec.enabledPort(context.config, name), + ); if (options.runtime === "native") { const executable = yield* Ref.get(prepared); const artifactRoot = yield* Ref.get(preparedRoot); @@ -386,55 +469,263 @@ export const makeProcessRecipe = return yield* serviceError("launch", "Artifact was not prepared"); if (artifactRoot === undefined) return yield* serviceError("launch", "Artifact root was not prepared"); - for (const [index, process] of spec.startup.entries()) { - const startupProcess = yield* spawnNativeProcess( + const reserveEndpoints = Effect.fn("ProcessRecipe.reserveEndpoints")(function* ( + parent: Scope.Closeable, + ) { + const portScope = yield* Scope.fork(parent, "sequential"); + const reservations = yield* Effect.forEach(portNames, () => reserveNativePort(0), { + concurrency: 1, + }).pipe( + Scope.provide(portScope), + Effect.mapError((cause) => serviceError("launch", cause)), + ); + const selected = new Map(); + for (const [index, [name]] of portNames.entries()) { + const reservation = reservations[index]; + if (reservation === undefined) + return yield* serviceError("launch", `Native ${name} port was not reserved`); + selected.set(name, { + kind: "tcp", + host: "127.0.0.1", + port: reservation.port, + }); + } + return { portScope, endpoints: selected }; + }); + + let reusableReservation: + | { + readonly portScope: Scope.Closeable; + readonly endpoints: ReadonlyMap; + } + | undefined; + let reusableEndpoints: ReadonlyMap | undefined; + if (spec.startup.length > 0) { + const startupScope = yield* Scope.fork(context.scope, "sequential"); + const reservation = yield* reserveEndpoints( + spec.nativeStartupEnv === undefined ? startupScope : context.scope, + ); + if (spec.nativeStartupEnv === undefined) { + yield* Scope.close(reservation.portScope, Exit.void); + reusableEndpoints = reservation.endpoints; + } + for (const [index, process] of spec.startup.entries()) { + const startupProcess = yield* spawnNativeProcess( + { + executable: `${artifactRoot}/bin/${process.nativeExecutable ?? "prepare"}`, + args: process.args, + env: yield* (spec.nativeStartupEnv ?? spec.env)( + context.config, + reservation.endpoints, + false, + ), + cwd: artifactRoot, + }, + defaultNativeProcessLauncher(), + { + stackId: String(options.stackId), + workloadId: `${context.id}-${context.config.service}-startup-${index}`, + }, + ).pipe( + Scope.provide(startupScope), + Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, deps.spawner), + Effect.mapError((cause) => serviceError("launch", cause)), + ); + const result = yield* awaitStartup(context.config.service, startupProcess, logs); + if (result.code !== 0) { + yield* Scope.close(reservation.portScope, Exit.void); + yield* Scope.close(startupScope, Exit.void); + return yield* startupFailure(context.config.service, result); + } + } + yield* Scope.close(startupScope, Exit.void); + if (spec.nativeStartupEnv !== undefined) reusableReservation = reservation; + } + + const deadline = (yield* Clock.currentTimeMillis) + startupTimeoutSeconds * 1_000; + const lastCollision = yield* Ref.make(undefined); + let firstAttempt = true; + const launchAttempt = Effect.fn("ProcessRecipe.launchNativeAttempt")(function* () { + if ( + spec.nativeReadinessOutput !== undefined && + (yield* Clock.currentTimeMillis) >= deadline + ) { + const previousCollision = yield* Ref.get(lastCollision); + return terminalSession( + previousCollision === undefined + ? serviceError("health", `${context.config.service} native readiness timed out`) + : serviceError("launch", previousCollision.message), + endpoints, + ); + } + const attemptScope = yield* Scope.fork(context.scope, "sequential"); + const heldReservation = firstAttempt ? reusableReservation : undefined; + const reuse = firstAttempt + ? (heldReservation?.endpoints ?? reusableEndpoints) + : undefined; + firstAttempt = false; + const reservation = + reuse === undefined ? yield* reserveEndpoints(attemptScope) : undefined; + const selected = reuse ?? reservation?.endpoints ?? new Map(); + const args = yield* spec.args(context.config, selected, { + container: false, + artifactRoot, + }); + const env = yield* spec.env(context.config, selected, false); + if (reservation !== undefined) yield* Scope.close(reservation.portScope, Exit.void); + if (heldReservation !== undefined) + yield* Scope.close(heldReservation.portScope, Exit.void); + + const native: NativeProcess = yield* spawnNativeProcess( { - executable: `${artifactRoot}/bin/${process.nativeExecutable ?? "prepare"}`, - args: process.args, - env: yield* (spec.nativeStartupEnv ?? spec.env)(context.config, desired, false), - cwd: artifactRoot, + executable, + args, + env, + gracefulStopSignal: "SIGTERM", + gracefulStopTimeout: "5 seconds", }, defaultNativeProcessLauncher(), - { - stackId: String(options.stackId), - workloadId: `${context.id}-${context.config.service}-startup-${index}`, - }, + { stackId: String(options.stackId), workloadId: context.id }, ).pipe( - Scope.provide(context.scope), + Scope.provide(attemptScope), Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, deps.spawner), Effect.mapError((cause) => serviceError("launch", cause)), ); - const result = yield* awaitStartup(context.config.service, startupProcess, logs); - if (result.code !== 0) return yield* startupFailure(context.config.service, result); - } - const native: NativeProcess = yield* spawnNativeProcess( - { - executable, - args: yield* spec.args(context.config, desired, { container: false, artifactRoot }), - env: yield* spec.env(context.config, desired, false), - gracefulStopSignal: "SIGTERM", - gracefulStopTimeout: "5 seconds", - }, - defaultNativeProcessLauncher(), - { stackId: String(options.stackId), workloadId: context.id }, - ).pipe( - Scope.provide(context.scope), - Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, deps.spawner), - Effect.mapError((cause) => serviceError("launch", cause)), - ); - yield* Ref.set(endpoints, desired); - yield* publishLogs(native, logs, context.scope); - const ready = desired.get("http"); - return { - health: + + if (spec.nativeReadinessOutput === undefined) { + yield* Ref.set(endpoints, selected); + yield* publishLogs(native, logs, attemptScope); + const ready = selected.get("http"); + return { + health: + ready === undefined + ? Effect.fail(serviceError("health", "Recipe has no HTTP readiness endpoint")) + : readiness(deps.client, ready, spec.healthPath), + exit: processExit(native.exitCode), + stop: native.kill.pipe(Effect.mapError((cause) => serviceError("stop", cause))), + remove: Ref.set(endpoints, new Map()), + } satisfies RuntimeSession; + } + + const output = yield* collectNativeOutput( + native, + logs, + selected, + spec.nativeReadinessOutput, + attemptScope, + ); + const ready = selected.get("http"); + const timeLeft = Math.max(0, deadline - (yield* Clock.currentTimeMillis)); + const readyCondition = ready === undefined ? Effect.fail(serviceError("health", "Recipe has no HTTP readiness endpoint")) - : readiness(deps.client, ready, spec.healthPath), - exit: processExit(native.exitCode), - stop: native.kill.pipe(Effect.mapError((cause) => serviceError("stop", cause))), - remove: Ref.set(endpoints, new Map()), - } satisfies RuntimeSession; + : Effect.all( + [ + readiness(deps.client, ready, spec.healthPath, Duration.millis(timeLeft)), + Deferred.await(output.bindReady), + ], + { concurrency: "unbounded", discard: true }, + ).pipe(Effect.asVoid); + const boundedReady = readyCondition.pipe( + Effect.timeout(Duration.millis(timeLeft)), + Effect.mapError((cause) => serviceError("health", cause)), + ); + const observed = yield* Effect.race( + Effect.exit(boundedReady).pipe( + Effect.map((result) => ({ _tag: "ready" as const, result })), + ), + Effect.exit(native.exitCode).pipe( + Effect.map((result) => ({ _tag: "exit" as const, result })), + ), + ); + + if (observed._tag === "ready" && Exit.isSuccess(observed.result)) { + yield* Ref.set(endpoints, selected); + return { + health: Effect.void, + exit: processExit(native.exitCode), + stop: native.kill.pipe(Effect.mapError((cause) => serviceError("stop", cause))), + remove: Ref.set(endpoints, new Map()), + } satisfies RuntimeSession; + } + + if (observed._tag === "ready") { + const bindReady = yield* Deferred.isDone(output.bindReady); + const [stdoutLines, stderrLines] = yield* Effect.all([ + Ref.get(output.stdout), + Ref.get(output.stderr), + ]); + const failure = serviceError( + "health", + withRecentOutput( + bindReady + ? `${context.config.service} HTTP readiness timed out` + : `${context.config.service} native listener bind was not confirmed`, + { stdout: stdoutLines, stderr: stderrLines }, + ), + ); + yield* Ref.set(endpoints, selected); + return { + health: Effect.fail(failure), + exit: processExit(native.exitCode), + stop: native.kill.pipe(Effect.mapError((cause) => serviceError("stop", cause))), + remove: Ref.set(endpoints, new Map()), + } satisfies RuntimeSession; + } + + if (Exit.isFailure(observed.result)) yield* Scope.close(attemptScope, Exit.void); + else yield* Deferred.await(output.drained); + const [stdoutLines, stderrLines] = yield* Effect.all([ + Ref.get(output.stdout), + Ref.get(output.stderr), + ]); + const exitMessage = Exit.isSuccess(observed.result) + ? `${context.config.service} startup exited with ${observed.result.value}` + : `${context.config.service} startup process failed: ${Cause.pretty(observed.result.cause)}`; + const failure = serviceError( + "launch", + withRecentOutput(exitMessage, { stdout: stdoutLines, stderr: stderrLines }), + ); + const collision = yield* Ref.get(output.bindError); + const nonzeroExit = Exit.isSuccess(observed.result) && observed.result.value !== 0; + if (Exit.isSuccess(observed.result)) yield* Scope.close(attemptScope, Exit.void); + yield* Ref.set(endpoints, new Map()); + if (collision && nonzeroExit) { + const collisionFailure = serviceError( + "native-port-collision", + `${context.config.service} native port collision\n${failure.message}`, + ); + yield* Ref.set(lastCollision, collisionFailure); + return yield* collisionFailure; + } + return terminalSession(failure, endpoints); + }); + + if (spec.nativeReadinessOutput !== undefined) { + const collisionRetries = Schedule.recurs(2).pipe( + Schedule.setInputType(), + Schedule.while( + ({ input }) => + input instanceof ServiceError && input.operation === "native-port-collision", + ), + ); + return yield* launchAttempt().pipe( + Effect.retry(collisionRetries), + Effect.catchIf( + (error) => + error instanceof ServiceError && error.operation === "native-port-collision", + (collision) => + Effect.succeed( + terminalSession(serviceError("launch", collision.message), endpoints), + ), + ), + ); + } + return yield* launchAttempt(); } + const desired = new Map(); + for (const [name, port] of portNames) + desired.set(name, { kind: "tcp", host: "127.0.0.1", port }); if (deps.container === undefined) return yield* serviceError("launch", "Container runtime unavailable"); const resolved = yield* resolveArtifact({