From 17c10863a9b5c6c39ea8f44adfafc2d5832990c7 Mon Sep 17 00:00:00 2001 From: Konstantin Wohlwend Date: Fri, 17 Jul 2026 16:57:38 -0700 Subject: [PATCH] Various fixes before pushing to prod --- .github/actions/wait-for-check/action.yml | 88 +++---- .github/workflows/all-good.yaml | 220 ++++++++++++------ .github/workflows/e2e-api-tests.yaml | 4 +- .../e2e-custom-base-port-api-tests.yaml | 4 +- .github/workflows/e2e-fallback-tests.yaml | 4 +- .github/workflows/fast-fail-tests.yaml | 27 ++- .../setup-tests-with-custom-base-port.yaml | 4 +- .github/workflows/setup-tests.yaml | 4 +- AGENTS.md | 2 +- .../api/latest/internal/cli-auth/route.tsx | 32 ++- .../src/databases/bulldozer/index.test.ts | 44 ++++ .../src/databases/bulldozer/index.ts | 33 ++- apps/bulldozer-js/src/databases/index.ts | 4 + .../low-level/implementations/in-memory.ts | 3 + .../instant-availability.test.ts | 27 +++ .../implementations/instant-availability.ts | 17 ++ .../low-level/implementations/lmdb.test.ts | 23 ++ .../low-level/implementations/lmdb.ts | 38 ++- .../src/databases/piledriver/index.ts | 3 + apps/bulldozer-js/src/index.ts | 49 +++- .../[projectId]/cli-auth/page-client.tsx | 4 +- .../api/v1/internal/cli-auth.test.ts | 46 +++- scripts/reseed-bulldozer.sh | 37 ++- 23 files changed, 560 insertions(+), 157 deletions(-) diff --git a/.github/actions/wait-for-check/action.yml b/.github/actions/wait-for-check/action.yml index 442e27e50..56c500ae2 100644 --- a/.github/actions/wait-for-check/action.yml +++ b/.github/actions/wait-for-check/action.yml @@ -1,98 +1,102 @@ -name: "Wait for check" -description: "Polls the GitHub Checks API until a named check-run on this commit completes, then succeeds/fails accordingly." +name: "Wait for workflow" +description: "Polls the GitHub Actions API until a workflow run for this event completes, then succeeds/fails accordingly." inputs: - check-name: - description: "Name of the check-run (i.e. the job name) to wait for." + workflow-file: + description: "Repository-relative workflow filename to wait for." required: true github-token: - description: "Token used to query the Checks API." + description: "Token used to query the Actions API." required: true timeout-seconds: - description: "How long to wait for the check to appear and complete before giving up." + description: "How long to wait for the workflow run to appear and complete before giving up." required: false default: "3600" poll-interval-seconds: description: "How long to sleep between polls." required: false - default: "15" + default: "30" runs: using: "composite" steps: - - name: Wait for "${{ inputs.check-name }}" to succeed + - name: Wait for "${{ inputs.workflow-file }}" to succeed shell: bash env: GITHUB_TOKEN: ${{ inputs.github-token }} REPO: ${{ github.repository }} - # A check-run can be reported against either the commit the workflow ran - # on (`github.sha`, the merge commit for pull_request events) or the PR - # head commit, depending on the event. Poll both, like all-good.yaml. MERGE_SHA: ${{ github.sha }} HEAD_SHA: ${{ github.event.pull_request.head.sha }} - CHECK_NAME: ${{ inputs.check-name }} + EVENT_NAME: ${{ github.event_name }} + WORKFLOW_FILE: ${{ inputs.workflow-file }} TIMEOUT_SECONDS: ${{ inputs.timeout-seconds }} POLL_INTERVAL_SECONDS: ${{ inputs.poll-interval-seconds }} run: | set -euo pipefail + workflow_path=".github/workflows/${WORKFLOW_FILE}" shas="$MERGE_SHA" if [ -n "${HEAD_SHA:-}" ] && [ "$HEAD_SHA" != "$MERGE_SHA" ]; then shas="$shas $HEAD_SHA" fi - echo "Waiting for check \"${CHECK_NAME}\" on ${REPO} (commits: ${shas})..." + echo "Waiting for ${workflow_path} on ${REPO} for ${EVENT_NAME} (commits: ${shas})..." deadline=$(( $(date +%s) + TIMEOUT_SECONDS )) while true; do - # Gather matching check-runs across all candidate commits and dedupe by - # id, so the same run reported under two SHAs isn't counted twice. - # Responses can be large (up to 100 runs), so pipe them through files - # rather than passing them as jq argv (which hits ARG_MAX). - : > matching.jsonl + matching_file=$(mktemp) for sha in $shas; do - # Page through the check-runs API: a busy commit can have >100 runs, - # so the target might not be on the first page. page=1 while true; do - curl -s -f \ + if ! curl -sS --fail-with-body \ + --retry 3 \ + --retry-all-errors \ + --connect-timeout 10 \ + --max-time 30 \ -H "Authorization: Bearer ${GITHUB_TOKEN}" \ -H "Accept: application/vnd.github+json" \ - "https://api.github.com/repos/${REPO}/commits/${sha}/check-runs?per_page=100&page=${page}" \ - -o response.json || echo '{}' > response.json - jq -c --arg name "$CHECK_NAME" \ - '.check_runs[]? | select(.name == $name)' response.json >> matching.jsonl || true - page_count=$(jq '.check_runs | length' response.json 2>/dev/null || echo 0) + -H "X-GitHub-Api-Version: 2022-11-28" \ + "https://api.github.com/repos/${REPO}/actions/runs?event=${EVENT_NAME}&head_sha=${sha}&per_page=100&page=${page}" \ + -o response.json; then + rm -f "$matching_file" response.json + echo "Failed to fetch workflow runs; refusing to treat an API failure as an absent workflow." >&2 + exit 1 + fi + if ! jq -c --arg path "$workflow_path" \ + '.workflow_runs[]? | select(.path == $path)' response.json >> "$matching_file"; then + rm -f "$matching_file" response.json + echo "GitHub returned an invalid workflow-runs response." >&2 + exit 1 + fi + page_count=$(jq '.workflow_runs | length' response.json) [ "$page_count" -lt 100 ] && break page=$((page + 1)) done done - # Keep only the most recent check-run (highest id) so a stale failed - # rerun doesn't block when a newer run of the same check succeeded. - matching=$(jq -s -c 'unique_by(.id) | if length == 0 then [] else [max_by(.id)] end' matching.jsonl) - count=$(echo "$matching" | jq 'length') + rm -f response.json - if [ "$count" -eq 0 ]; then - echo "Check \"${CHECK_NAME}\" not reported yet; waiting..." + matching=$(jq -s -c 'unique_by(.id) | if length == 0 then null else max_by(.id) end' "$matching_file") + rm -f "$matching_file" + + if [ "$matching" = "null" ]; then + echo "Workflow ${workflow_path} not reported for this event yet; waiting..." else - pending=$(echo "$matching" | jq '[.[] | select(.status != "completed")] | length') - if [ "$pending" -gt 0 ]; then - echo "Check \"${CHECK_NAME}\" is still running; waiting..." + status=$(echo "$matching" | jq -r '.status') + if [ "$status" != "completed" ]; then + echo "Workflow ${workflow_path} is ${status}; waiting..." else - # Completed. Treat success/skipped/neutral as "did not block". - failed=$(echo "$matching" | jq '[.[] | select(.conclusion != "success" and .conclusion != "skipped" and .conclusion != "neutral")] | length') - if [ "$failed" -eq 0 ]; then - echo "Check \"${CHECK_NAME}\" passed. Proceeding." + conclusion=$(echo "$matching" | jq -r '.conclusion') + if [ "$conclusion" = "success" ] || [ "$conclusion" = "skipped" ] || [ "$conclusion" = "neutral" ]; then + echo "Workflow ${workflow_path} completed with ${conclusion}. Proceeding." exit 0 fi - echo "Check \"${CHECK_NAME}\" did not pass:" - echo "$matching" | jq -r '.[] | " - \(.name): \(.conclusion)"' + echo "$matching" | jq -r '"Workflow \(.html_url) completed with \(.conclusion)."' exit 1 fi fi if [ "$(date +%s)" -ge "$deadline" ]; then - echo "Timed out after ${TIMEOUT_SECONDS}s waiting for check \"${CHECK_NAME}\"." + echo "Timed out after ${TIMEOUT_SECONDS}s waiting for ${workflow_path}." exit 1 fi sleep "$POLL_INTERVAL_SECONDS" diff --git a/.github/workflows/all-good.yaml b/.github/workflows/all-good.yaml index 266358198..5bb57f96f 100644 --- a/.github/workflows/all-good.yaml +++ b/.github/workflows/all-good.yaml @@ -15,24 +15,106 @@ jobs: all-good: runs-on: ubuntu-latest permissions: + actions: read checks: read contents: read env: REPO: ${{ github.repository }} COMMIT: ${{ github.sha }} PR_HEAD_SHA: ${{ github.event.pull_request.head.sha }} + EVENT_NAME: ${{ github.event_name }} + REF_NAME: ${{ github.ref_name }} steps: - - name: Wait for 60 seconds - run: sleep 60 - - - name: Poll Checks API until complete + - name: Wait for every expected workflow and check env: GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} run: | set -euo pipefail - echo "Checking check runs for commit ${COMMIT} in repo ${REPO}..." - + # This is deliberately an explicit contract. Discovering only the + # workflow runs or check runs that already exist has an open-world + # race: a queued workflow may not have reported either one yet. + expected_workflows=( + "e2e-fallback-tests.yaml" + "lint-and-build.yaml" + "docker-server-build-run.yaml" + "setup-tests.yaml" + "check-prisma-migrations.yaml" + "e2e-custom-base-port-api-tests.yaml" + "table-of-contents.yaml" + "setup-tests-with-custom-base-port.yaml" + "fast-fail-tests.yaml" + "e2e-api-tests.yaml" + ) + case "$EVENT_NAME" in + pull_request) + expected_workflows+=("db-migration-backwards-compatibility.yaml") + ;; + push) + expected_workflows+=("docker-server-build-push.yaml") + if [ "$REF_NAME" = "dev" ]; then + expected_workflows+=("db-migration-backwards-compatibility.yaml") + elif [ "$REF_NAME" = "main" ]; then + expected_workflows+=("npm-publish.yaml") + else + echo "Unexpected push branch for all-good: ${REF_NAME}" >&2 + exit 1 + fi + ;; + *) + echo "Unexpected event for all-good: ${EVENT_NAME}" >&2 + exit 1 + ;; + esac + + function candidate_shas() { + printf '%s\n' "$COMMIT" + if [ -n "${PR_HEAD_SHA:-}" ] && [ "$PR_HEAD_SHA" != "$COMMIT" ]; then + printf '%s\n' "$PR_HEAD_SHA" + fi + } + + function get_workflow_runs() { + local responses_file + responses_file=$(mktemp) + + while IFS= read -r sha; do + local page=1 + while true; do + local response + if ! response=$(curl -sS --fail-with-body \ + --retry 3 \ + --retry-all-errors \ + --connect-timeout 10 \ + --max-time 30 \ + -H "Authorization: Bearer ${GITHUB_TOKEN}" \ + -H "Accept: application/vnd.github+json" \ + -H "X-GitHub-Api-Version: 2022-11-28" \ + "https://api.github.com/repos/${REPO}/actions/runs?event=${EVENT_NAME}&head_sha=${sha}&per_page=100&page=${page}"); then + echo "Error fetching workflow runs for ${sha}" >&2 + rm -f "$responses_file" + return 1 + fi + if ! jq -c '.workflow_runs[]?' <<< "$response" >> "$responses_file"; then + echo "Invalid Actions API response for ${sha}" >&2 + rm -f "$responses_file" + return 1 + fi + local page_count + if ! page_count=$(jq '.workflow_runs | length' <<< "$response"); then + echo "Invalid Actions API response for ${sha}" >&2 + rm -f "$responses_file" + return 1 + fi + [ "$page_count" -lt 100 ] && break + page=$((page + 1)) + done + done < <(candidate_shas) + + jq -s '{ workflow_runs: unique_by(.id) }' "$responses_file" + rm -f "$responses_file" + } + function get_check_runs() { local endpoint=$1 local page=1 @@ -88,75 +170,81 @@ jobs: printf '%s\n' "$combined" } - + function count_pending_checks() { local response=$1 echo "$response" | jq '([.check_runs[]? | select(.status != "completed" and .name != "all-good" and .name != "claude-review")] | length) // 0' } - + function count_failed_checks() { local response=$1 - echo "$response" | jq '([.check_runs[]? | select(.conclusion != "success" and .conclusion != "skipped" and .conclusion != "neutral" and .name != "all-good" and .name != "claude-review")] | length) // 0' + echo "$response" | jq '([.check_runs[]? | select(.status == "completed" and .conclusion != "success" and .conclusion != "skipped" and .conclusion != "neutral" and .name != "all-good" and .name != "claude-review")] | length) // 0' } - + + deadline=$(( $(date +%s) + 21600 )) while true; do - # Always check the current commit's checks - commit_response=$(get_check_runs "https://api.github.com/repos/${REPO}/commits/${COMMIT}/check-runs") - commit_total=$(echo "$commit_response" | jq -r '.total_count // 0') - - # If this is a PR, check the PR's head commit checks - pr_total=0 - pr_response='{"total_count":0,"check_runs":[]}' - if [ -n "$PR_HEAD_SHA" ] && [ "$PR_HEAD_SHA" != "$COMMIT" ]; then - pr_response=$(get_check_runs "https://api.github.com/repos/${REPO}/commits/${PR_HEAD_SHA}/check-runs") - pr_total=$(echo "$pr_response" | jq -r '.total_count // 0') - echo "Found ${commit_total} current commit checks and ${pr_total} PR head commit checks" - else - echo "Found ${commit_total} commit checks" - fi - - # If no checks found at all, wait and retry - if [ "$commit_total" -eq 0 ] && { [ -z "$PR_HEAD_SHA" ] || [ "$pr_total" -eq 0 ]; }; then - echo "No check runs found. Waiting..." - sleep 10 - continue - fi - - # Check for pending runs in both current and PR head commit checks - commit_pending=$(count_pending_checks "$commit_response") - pr_pending=0 - if [ -n "$PR_HEAD_SHA" ] && [ "$PR_HEAD_SHA" != "$COMMIT" ]; then - pr_pending=$(count_pending_checks "$pr_response") - fi - - total_pending=$((commit_pending + pr_pending)) - if [ "$total_pending" -gt 0 ]; then - echo "$total_pending check run(s) still in progress. Waiting..." - sleep 10 - continue - fi - - # Check for failures in both current and PR head commit checks - commit_failed=$(count_failed_checks "$commit_response") - pr_failed=0 - if [ -n "$PR_HEAD_SHA" ] && [ "$PR_HEAD_SHA" != "$COMMIT" ]; then - pr_failed=$(count_failed_checks "$pr_response") - fi - - total_failed=$((commit_failed + pr_failed)) - if [ "$total_failed" -eq 0 ]; then - echo "All check runs passed!" - exit 0 - else - echo "The following check run(s) failed:" - # Failed checks on the current commit (excluding claude-review) - echo "$commit_response" | jq -r '.check_runs[] | select(.conclusion != "success" and .conclusion != "skipped" and .conclusion != "neutral" and .name != "claude-review") | .name' | sed 's/^/ - /' - # Failed checks on the PR head commit (if different, excluding claude-review) - if [ -n "$PR_HEAD_SHA" ] && [ "$PR_HEAD_SHA" != "$COMMIT" ]; then - echo "$pr_response" | jq -r '.check_runs[] | select(.conclusion != "success" and .conclusion != "skipped" and .conclusion != "neutral" and .name != "claude-review") | .name' | sed 's/^/ - /' + workflow_response=$(get_workflow_runs) + missing_workflows=() + pending_workflows=() + failed_workflows=() + + for workflow_file in "${expected_workflows[@]}"; do + workflow_path=".github/workflows/${workflow_file}" + latest=$(jq -c --arg path "$workflow_path" ' + [.workflow_runs[]? | select(.path == $path)] + | if length == 0 then null else max_by(.id) end + ' <<< "$workflow_response") + if [ "$latest" = "null" ]; then + missing_workflows+=("$workflow_file") + continue fi - echo "$commit_response" - echo "$pr_response" + + status=$(jq -r '.status' <<< "$latest") + conclusion=$(jq -r '.conclusion // ""' <<< "$latest") + if [ "$status" != "completed" ]; then + pending_workflows+=("${workflow_file} (${status})") + elif [ "$conclusion" != "success" ] && [ "$conclusion" != "skipped" ] && [ "$conclusion" != "neutral" ]; then + failed_workflows+=("${workflow_file} (${conclusion})") + fi + done + + if [ "${#failed_workflows[@]}" -gt 0 ]; then + echo "Expected workflow failures:" + printf ' - %s\n' "${failed_workflows[@]}" exit 1 fi + + if [ "${#missing_workflows[@]}" -gt 0 ] || [ "${#pending_workflows[@]}" -gt 0 ]; then + if [ "${#missing_workflows[@]}" -gt 0 ]; then + echo "Expected workflows not reported yet:" + printf ' - %s\n' "${missing_workflows[@]}" + fi + if [ "${#pending_workflows[@]}" -gt 0 ]; then + echo "Expected workflows still running:" + printf ' - %s\n' "${pending_workflows[@]}" + fi + else + # Expected GitHub workflows are now a closed, complete set. Check + # every check run currently attached to this event's commit too, + # so already-reported external checks can still block the gate. + commit_response=$(get_check_runs "https://api.github.com/repos/${REPO}/commits/${COMMIT}/check-runs") + pending_checks=$(count_pending_checks "$commit_response") + failed_checks=$(count_failed_checks "$commit_response") + if [ "$failed_checks" -gt 0 ]; then + echo "The following check run(s) failed:" + echo "$commit_response" | jq -r '.check_runs[] | select(.status == "completed" and .conclusion != "success" and .conclusion != "skipped" and .conclusion != "neutral" and .name != "all-good" and .name != "claude-review") | .name' | sed 's/^/ - /' + exit 1 + fi + if [ "$pending_checks" -eq 0 ]; then + echo "Every expected workflow and reported check passed." + exit 0 + fi + echo "${pending_checks} reported check run(s) still in progress." + fi + + if [ "$(date +%s)" -ge "$deadline" ]; then + echo "Timed out after six hours waiting for the expected CI contract." >&2 + exit 1 + fi + sleep 30 done diff --git a/.github/workflows/e2e-api-tests.yaml b/.github/workflows/e2e-api-tests.yaml index a5a03d3f2..2dda01537 100644 --- a/.github/workflows/e2e-api-tests.yaml +++ b/.github/workflows/e2e-api-tests.yaml @@ -16,7 +16,7 @@ jobs: name: wait-for-fast-fail runs-on: ubuntu-latest permissions: - checks: read + actions: read contents: read steps: - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6 @@ -24,7 +24,7 @@ jobs: persist-credentials: false - uses: ./.github/actions/wait-for-check with: - check-name: fast-fail-tests + workflow-file: fast-fail-tests.yaml github-token: ${{ secrets.GITHUB_TOKEN }} build: diff --git a/.github/workflows/e2e-custom-base-port-api-tests.yaml b/.github/workflows/e2e-custom-base-port-api-tests.yaml index b4b14a309..ab49f6520 100644 --- a/.github/workflows/e2e-custom-base-port-api-tests.yaml +++ b/.github/workflows/e2e-custom-base-port-api-tests.yaml @@ -16,7 +16,7 @@ jobs: name: wait-for-fast-fail runs-on: ubuntu-latest permissions: - checks: read + actions: read contents: read steps: - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6 @@ -24,7 +24,7 @@ jobs: persist-credentials: false - uses: ./.github/actions/wait-for-check with: - check-name: fast-fail-tests + workflow-file: fast-fail-tests.yaml github-token: ${{ secrets.GITHUB_TOKEN }} build: diff --git a/.github/workflows/e2e-fallback-tests.yaml b/.github/workflows/e2e-fallback-tests.yaml index 7333d7957..4388f7a4a 100644 --- a/.github/workflows/e2e-fallback-tests.yaml +++ b/.github/workflows/e2e-fallback-tests.yaml @@ -18,7 +18,7 @@ jobs: name: wait-for-fast-fail runs-on: ubuntu-latest permissions: - checks: read + actions: read contents: read steps: - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6 @@ -26,7 +26,7 @@ jobs: persist-credentials: false - uses: ./.github/actions/wait-for-check with: - check-name: fast-fail-tests + workflow-file: fast-fail-tests.yaml github-token: ${{ secrets.GITHUB_TOKEN }} build: diff --git a/.github/workflows/fast-fail-tests.yaml b/.github/workflows/fast-fail-tests.yaml index ee3618736..62aa35d10 100644 --- a/.github/workflows/fast-fail-tests.yaml +++ b/.github/workflows/fast-fail-tests.yaml @@ -111,21 +111,33 @@ jobs: } JSON + - name: Check whether the AI selector is available + id: ai-selector + env: + OPENAI_API_KEY: ${{ secrets.OPENAI_API_KEY }} + run: | + if [ -z "${OPENAI_API_KEY:-}" ]; then + echo "::warning::OPENAI_API_KEY is unavailable; skipping the fast-fail subset and falling back to the full test suites." + echo "available=false" >> "$GITHUB_OUTPUT" + else + echo "available=true" >> "$GITHUB_OUTPUT" + fi + - name: Install Codex CLI + id: install-codex + if: ${{ steps.ai-selector.outputs.available == 'true' }} + continue-on-error: true run: npm install -g @openai/codex@0.142.5 - name: Ask Codex which tests to run id: codex + if: ${{ steps.ai-selector.outputs.available == 'true' && steps.install-codex.outcome == 'success' }} + continue-on-error: true env: OPENAI_API_KEY: ${{ secrets.OPENAI_API_KEY }} run: | set -euo pipefail - if [ -z "${OPENAI_API_KEY:-}" ]; then - echo "::error::OPENAI_API_KEY is not set. Add it as a repository Actions secret (Settings -> Secrets and variables -> Actions)." - exit 1 - fi - # This Codex version does not pick up OPENAI_API_KEY automatically for # `codex exec`; it must be logged in explicitly first, otherwise every # request is sent unauthenticated and fails with 401. @@ -145,6 +157,11 @@ jobs: echo "Codex raw output:" cat codex-out.json + - name: Report AI selector fallback + if: ${{ always() && (steps.ai-selector.outputs.available != 'true' || steps.install-codex.outcome != 'success' || steps.codex.outcome != 'success') }} + run: | + echo "::warning::The AI selector was unavailable. No fast-fail subset will run; the full test suites remain authoritative." + - name: Parse selected test files id: select run: | diff --git a/.github/workflows/setup-tests-with-custom-base-port.yaml b/.github/workflows/setup-tests-with-custom-base-port.yaml index 9c2a1e13f..8ce3ada23 100644 --- a/.github/workflows/setup-tests-with-custom-base-port.yaml +++ b/.github/workflows/setup-tests-with-custom-base-port.yaml @@ -20,7 +20,7 @@ jobs: if: ${{ (github.head_ref || github.ref_name) == 'dev' }} runs-on: ubuntu-latest permissions: - checks: read + actions: read contents: read steps: - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6 @@ -28,7 +28,7 @@ jobs: persist-credentials: false - uses: ./.github/actions/wait-for-check with: - check-name: fast-fail-tests + workflow-file: fast-fail-tests.yaml github-token: ${{ secrets.GITHUB_TOKEN }} setup-tests-with-custom-base-port: diff --git a/.github/workflows/setup-tests.yaml b/.github/workflows/setup-tests.yaml index 14a214928..45516d9a1 100644 --- a/.github/workflows/setup-tests.yaml +++ b/.github/workflows/setup-tests.yaml @@ -20,7 +20,7 @@ jobs: if: ${{ (github.head_ref || github.ref_name) == 'dev' }} runs-on: ubuntu-latest permissions: - checks: read + actions: read contents: read steps: - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6 @@ -28,7 +28,7 @@ jobs: persist-credentials: false - uses: ./.github/actions/wait-for-check with: - check-name: fast-fail-tests + workflow-file: fast-fail-tests.yaml github-token: ${{ secrets.GITHUB_TOKEN }} setup-tests: diff --git a/AGENTS.md b/AGENTS.md index e999f4ca7..2c814bf9f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -100,7 +100,7 @@ To see all development ports, refer to the index.html of `apps/dev-launchpad/pub - NEVER implement a hacky solution without EXPLICIT approval from the user. Always go the extra mile to make sure the solution is clean, maintainable, and robust. - Fail early, fail loud. Fail fast with an error instead of silently continuing. - Do NOT use `as`/`any`/type casts or anything else like that to bypass the type system unless you specifically asked the user about it. Most of the time a place where you would use type casts is not one where you actually need them. Avoid wherever possible. -- When writing database migration files, assume that we have >>>10,000,000 rows in every table (unless otherwise specified). This means you may have to use CONDITIONALLY_REPEAT_MIGRATION_SENTINEL to avoid running the migration and things like concurrent index builds; see the existing migrations for examples. One common pattern is to add a temporary index or extra boolean column marking whether the row has already been migrated (then deleting the column at the end). +- When writing database migration files, assume that we have >>>1,000,000 rows in every table (unless otherwise specified). This means you may have to use CONDITIONALLY_REPEAT_MIGRATION_SENTINEL to avoid running the migration and things like concurrent index builds; see the existing migrations for examples. One common pattern is to add a temporary index or extra boolean column marking whether the row has already been migrated (then deleting the column at the end). - Each migration file runs in its own transaction with a relatively short timeout. Split long-running operations into separate migration files to avoid timeouts. For example, when adding CHECK constraints, use `NOT VALID` in one migration, then `VALIDATE CONSTRAINT` in a separate migration file. - Note that each database migration file is executed in a single transaction. Even with the run-outside-transaction sentinel, the transaction will still continue during the entire migration file. If you want to split things up into multiple transactions, put it into their own migration files. - When writing database migration files, ALWAYS ALWAYS add tests for all the potential edge cases! See the folder structure of the other migrations to see how that works. diff --git a/apps/backend/src/app/api/latest/internal/cli-auth/route.tsx b/apps/backend/src/app/api/latest/internal/cli-auth/route.tsx index 8203df5ca..93deab1be 100644 --- a/apps/backend/src/app/api/latest/internal/cli-auth/route.tsx +++ b/apps/backend/src/app/api/latest/internal/cli-auth/route.tsx @@ -131,17 +131,21 @@ export const GET = createSmartRouteHandler({ }; }); - // Active-session discovery is also bounded: it considers refresh tokens - // attached to the most recently used CLI attempts, then checks which of - // those sessions still exist in the canonical refresh-token table. + // Apply LIMIT before the nullable-column predicates. Otherwise LIMIT only + // bounds returned rows and a tenant with sparse completed attempts could + // scan its entire history looking for 200 qualifying entries. const cliRefreshTokens = await prisma.$replica().$queryRaw<{ refreshToken: string }[]>(Prisma.sql` + WITH "recentAttempts" AS MATERIALIZED ( + SELECT "refreshToken", "usedAt" + FROM ${sqlQuoteIdent(schema)}."CliAuthAttempt" + WHERE "tenancyId" = ${tenancy.id}::UUID + ORDER BY "createdAt" DESC, "id" DESC + LIMIT ${activeTokenAttemptLimit} + ) SELECT "refreshToken" - FROM ${sqlQuoteIdent(schema)}."CliAuthAttempt" - WHERE "tenancyId" = ${tenancy.id}::UUID - AND "refreshToken" IS NOT NULL + FROM "recentAttempts" + WHERE "refreshToken" IS NOT NULL AND "usedAt" IS NOT NULL - ORDER BY "createdAt" DESC, "id" DESC - LIMIT ${activeTokenAttemptLimit} `); let activeCliUsers: Array<{ @@ -157,15 +161,6 @@ export const GET = createSmartRouteHandler({ if (cliRefreshTokens.length > 0) { const tokenValues = cliRefreshTokens.map((row) => row.refreshToken); - const countResult = await globalPrismaClient.$replica().$queryRaw<{ count: bigint }[]>(Prisma.sql` - SELECT COUNT(*) AS count - FROM "ProjectUserRefreshToken" - WHERE "tenancyId" = ${tenancy.id}::UUID - AND "refreshToken" = ANY(${tokenValues}) - AND ("expiresAt" IS NULL OR "expiresAt" >= ${now}) - `); - activeTokenCount = Number(countResult[0]?.count ?? 0n); - const activeTokens = await globalPrismaClient.$replica().$queryRaw(Prisma.sql` SELECT "id", @@ -177,8 +172,9 @@ export const GET = createSmartRouteHandler({ WHERE "tenancyId" = ${tenancy.id}::UUID AND "refreshToken" = ANY(${tokenValues}) ORDER BY "lastActiveAt" DESC - LIMIT 50 + LIMIT ${activeTokenAttemptLimit} `); + activeTokenCount = activeTokens.filter((token) => token.expiresAt == null || token.expiresAt >= now).length; if (activeTokens.length > 0) { const userIds = [...new Set(activeTokens.map((token) => token.projectUserId))]; diff --git a/apps/bulldozer-js/src/databases/bulldozer/index.test.ts b/apps/bulldozer-js/src/databases/bulldozer/index.test.ts index 5cdd3802b..946aea5b2 100644 --- a/apps/bulldozer-js/src/databases/bulldozer/index.test.ts +++ b/apps/bulldozer-js/src/databases/bulldozer/index.test.ts @@ -1,5 +1,10 @@ import { describe, expect, it } from "vitest"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; import { declareInMemoryLowLevelDatabase } from "../low-level/implementations/in-memory.js"; +import { declareInstantAvailabilityLowLevelDatabase } from "../low-level/implementations/instant-availability.js"; +import { declareLmdbLowLevelDatabase } from "../low-level/implementations/lmdb.js"; import { declarePiledriverDatabase, PiledriverObject } from "../piledriver/index.js"; import { ConcatTreeList } from "../piledriver/data-structures/concat-tree-list.js"; import { stringCompare } from "@hexclave/shared/dist/utils/strings"; @@ -87,6 +92,45 @@ describe("Bulldozer", () => { ]); }); + it("retains the latest write sequence after instant-availability cache eviction", async () => { + const path = await mkdtemp(join(tmpdir(), "bulldozer-durability-barrier-")); + const lmdb = declareLmdbLowLevelDatabase({ path, dbId: "durability-barrier" }); + const instant = declareInstantAvailabilityLowLevelDatabase(lmdb, { dbId: "instant-durability-barrier" }); + const db = declareBulldozerDatabase(declarePiledriverDatabase(instant), { + migrations: [[{ type: "initTable", tableId: "store", table: defineStoredTable(), inputTables: {} }]], + }); + let releaseDurability: (() => void) | undefined; + const durabilityGate = new Promise((resolve) => { + releaseDurability = resolve; + }); + try { + await db.applyRemainingMigrations(); + + const originalWaitUntilDurable = lmdb.waitUntilDurable.bind(lmdb); + lmdb.waitUntilDurable = async (seq) => { + if (seq !== lmdb.initialSeq) await durabilityGate; + await originalWaitUntilDurable(seq); + }; + + const write = await db.withSnapshot(async snapshot => await set(snapshot, "store", "a", 1)); + await instant.waitUntilUnderlyingAvailable(write.seq); + + const barrier = db.waitUntilCurrentStateDurable(); + expect(await Promise.race([ + barrier.then(() => "resolved"), + Promise.resolve("pending"), + ])).toBe("pending"); + + if (releaseDurability === undefined) throw new Error("Durability gate was not initialized"); + releaseDurability(); + await barrier; + } finally { + if (releaseDurability !== undefined) releaseDurability(); + await db.close(); + await rm(path, { recursive: true, force: true }); + } + }); + it("does not expose stored rows through non-null groups", async () => { let snapshot = await initializedSnapshot([[ { type: "initTable", tableId: "store", table: defineStoredTable(), inputTables: {} }, diff --git a/apps/bulldozer-js/src/databases/bulldozer/index.ts b/apps/bulldozer-js/src/databases/bulldozer/index.ts index f22109419..9451e0157 100644 --- a/apps/bulldozer-js/src/databases/bulldozer/index.ts +++ b/apps/bulldozer-js/src/databases/bulldozer/index.ts @@ -857,6 +857,8 @@ export type BulldozerDatabase = { getSnapshot(): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, withSnapshot(updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, withSnapshotReplicated(updateSnapshot: (snapshot: BulldozerDatabaseSnapshot) => Promise): Promise<{ snapshot: BulldozerDatabaseSnapshot, seq: DatabaseSeq }>, + waitUntilCurrentStateDurable(): Promise, + close(): Promise, applyRemainingMigrations(): Promise<{ seq: DatabaseSeq }>, }; @@ -865,10 +867,19 @@ type BulldozerDatabaseRootSerialized = { }; export function declareBulldozerDatabase(piledriverDatabase: PiledriverDatabase, options: { migrations: readonly BulldozerDatabaseMigration[] }): BulldozerDatabase { const rootKey = new TextEncoder().encode("bulldozer-database-root").buffer; + let latestRootWriteSeq = piledriverDatabase.initialSeq; const getRoot = async () => await piledriverDatabase.getRootObject(rootKey) as { object: BulldozerDatabaseRootSerialized, seq: DatabaseSeq }; - const setRoot = async (root: BulldozerDatabaseRootSerialized) => await piledriverDatabase.setRootObject(rootKey, root); + const setRoot = async (root: BulldozerDatabaseRootSerialized) => { + const result = await piledriverDatabase.setRootObject(rootKey, root); + // Keep the exact instant-availability sequence alive even after its cached + // value is evicted. Re-reading the root can return initialSeq before the + // corresponding LMDB flush has made this write durable. + latestRootWriteSeq = result.seq; + return result; + }; const tablesState = createTablesStateFromMigrations(options.migrations); let currentOperation: Promise | null = null; + let closePromise: Promise | null = null; const withWriteLock = async (operation: () => Promise): Promise => { while (currentOperation !== null) { @@ -968,8 +979,10 @@ export function declareBulldozerDatabase(piledriverDatabase: PiledriverDatabase, rootKey, getRoot, setRoot, + latestRootWriteSeq, tablesState, currentOperation, + closePromise, }; }, listTables: () => Object.entries(tablesState.tables).map(([tableId, tableState]) => ({ @@ -986,6 +999,24 @@ export function declareBulldozerDatabase(piledriverDatabase: PiledriverDatabase, getSnapshot, withSnapshot: async (updateSnapshot) => await withSnapshot(updateSnapshot, { replicated: false }), withSnapshotReplicated: async (updateSnapshot) => await withSnapshot(updateSnapshot, { replicated: true }), + waitUntilCurrentStateDurable: async () => await traceSpan("bulldozer-js.bulldozer.waitUntilCurrentStateDurable", async () => await withWriteLock(async () => { + // Taking the write lock makes this a barrier for every earlier mutation. + // Waiting on the retained write sequence avoids the eviction race where a + // fresh read sees initialSeq while the original LMDB flush is still pending. + await piledriverDatabase.waitUntilDurable(latestRootWriteSeq); + })), + close() { + if (closePromise === null) { + closePromise = traceSpan("bulldozer-js.bulldozer.close", async () => await withWriteLock(async () => { + try { + await piledriverDatabase.waitUntilDurable(latestRootWriteSeq); + } finally { + await piledriverDatabase.close(); + } + })); + } + return closePromise; + }, applyRemainingMigrations: async () => await traceSpan("bulldozer-js.bulldozer.applyRemainingMigrations", async () => await withWriteLock(async () => { let snapshot: BulldozerDatabaseSnapshotSerialized; let currentSeq = piledriverDatabase.initialSeq; diff --git a/apps/bulldozer-js/src/databases/index.ts b/apps/bulldozer-js/src/databases/index.ts index 68a48c417..11eb586d4 100644 --- a/apps/bulldozer-js/src/databases/index.ts +++ b/apps/bulldozer-js/src/databases/index.ts @@ -27,5 +27,9 @@ export type Database = { */ waitUntilReplicated(seq: DatabaseSeq): Promise, combineSeqs(...seqs: DatabaseSeq[]): DatabaseSeq, + /** + * Drains pending writes and releases resources. Calls are idempotent. + */ + close(): Promise, initialSeq: DatabaseSeq, }; diff --git a/apps/bulldozer-js/src/databases/low-level/implementations/in-memory.ts b/apps/bulldozer-js/src/databases/low-level/implementations/in-memory.ts index 6ff04ae12..a0519c2d8 100644 --- a/apps/bulldozer-js/src/databases/low-level/implementations/in-memory.ts +++ b/apps/bulldozer-js/src/databases/low-level/implementations/in-memory.ts @@ -154,6 +154,9 @@ export function declareInMemoryLowLevelDatabase(dbId: string): LowLevelDatabase combineSeqs(...seqs) { return this.initialSeq; }, + async close() { + // In-memory databases have no external resources to release. + }, async debugSnapshot() { return await traceSpanHot({ description: "bulldozer-js.low-level.in-memory.debugSnapshot", attributes: { "bulldozer.low_level.backend": "in-memory" } }, async () => { const stores: Record = {}; diff --git a/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.test.ts b/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.test.ts index 030d0c993..72b39c71d 100644 --- a/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.test.ts +++ b/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.test.ts @@ -11,6 +11,7 @@ const text = (value: ArrayBuffer | null) => value === null ? null : textDecoder. function createSlowSetDatabase() { const releaseSets: Array<() => void> = []; let setCallCount = 0; + let closeCallCount = 0; const committed = new Map(); const initialSeq = ["slow", "initial"] as unknown as DatabaseSeq; const seqToPromise = new Map>([[initialSeq, Promise.resolve()]]); @@ -64,6 +65,10 @@ function createSlowSetDatabase() { combineSeqs(...seqs) { return seqs[seqs.length - 1] ?? initialSeq; }, + async close() { + closeCallCount++; + await Promise.all(seqToPromise.values()); + }, initialSeq, }; return { @@ -74,6 +79,7 @@ function createSlowSetDatabase() { releaseSet(); }, setCallCount: () => setCallCount, + closeCallCount: () => closeCallCount, }; } @@ -127,6 +133,27 @@ describe("instant-availability low-level database", () => { slow.releaseSet(); }); + it("drains pending writes and closes the wrapped database exactly once", async () => { + const slow = createSlowSetDatabase(); + const db = declareInstantAvailabilityLowLevelDatabase(slow.db, { dbId: "instant-close-test" }); + const store = db.declareKvStore("store"); + await store.setAll([{ key: buffer("key"), value: buffer("pending") }]); + + const closing = db.close(); + expect(await Promise.race([ + closing.then(() => "closed"), + Promise.resolve("pending"), + ])).toBe("pending"); + expect(slow.closeCallCount()).toBe(0); + + slow.releaseSet(); + await closing; + expect(slow.closeCallCount()).toBe(1); + await db.close(); + expect(slow.closeCallCount()).toBe(1); + await expect(store.setAll([{ key: buffer("late"), value: buffer("rejected") }])).rejects.toThrow("closing"); + }); + it("only allows a single winner for concurrent compareAndSet on the same key", async () => { const slow = createSlowSetDatabase(); const db = declareInstantAvailabilityLowLevelDatabase(slow.db, { dbId: "instant-test" }); diff --git a/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.ts b/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.ts index 6ec71e68d..a289a4813 100644 --- a/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.ts +++ b/apps/bulldozer-js/src/databases/low-level/implementations/instant-availability.ts @@ -52,6 +52,8 @@ export function declareInstantAvailabilityLowLevelDatabase(wrapped: LowLevelData let underlyingAvailableSeqRecords = 1; const pendingSeqRecords = new Set(); let currentWriteGateOperation: Promise = Promise.resolve(); + let isClosing = false; + let closePromise: Promise | null = null; let pendingSeqRecordsChangedResolve: () => void; let pendingSeqRecordsChanged = new Promise(resolve => { pendingSeqRecordsChangedResolve = resolve; @@ -92,6 +94,7 @@ export function declareInstantAvailabilityLowLevelDatabase(wrapped: LowLevelData }; const withWriteGate = async (operation: () => Promise): Promise => { return await traceSpanHot({ description: "bulldozer-js.low-level.instant.withWriteGate", attributes: { "bulldozer.low_level.backend": "instant-availability" } }, async () => { + if (isClosing) throw new Error("Instant-availability database is closing and cannot accept writes"); const previousOperation = currentWriteGateOperation; let releaseCurrentOperation: () => void; currentWriteGateOperation = new Promise(resolve => { @@ -282,6 +285,20 @@ export function declareInstantAvailabilityLowLevelDatabase(wrapped: LowLevelData if (seqs.every(seq => getSeqId(seq) === initialSeqId)) return initialSeq; return createSeq((async () => wrapped.combineSeqs(...await Promise.all(seqs.map(async seq => await getUnderlyingSeq(seq)))))()); }, + close() { + if (closePromise === null) { + isClosing = true; + closePromise = traceSpanHot({ description: "bulldozer-js.low-level.instant.close", attributes: { "bulldozer.low_level.backend": "instant-availability" } }, async () => { + await currentWriteGateOperation; + try { + await Promise.all([...pendingSeqRecords].map(async record => await record.underlyingAvailable)); + } finally { + await wrapped.close(); + } + }); + } + return closePromise; + }, async debugSnapshot() { return await traceSpanHot({ description: "bulldozer-js.low-level.instant.debugSnapshot", attributes: { "bulldozer.low_level.backend": "instant-availability" } }, async () => await wrapped.debugSnapshot?.() ?? { stores: {}, dumps: {} }); }, diff --git a/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.test.ts b/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.test.ts index 5472d6a4c..bec5c83ce 100644 --- a/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.test.ts +++ b/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.test.ts @@ -37,6 +37,29 @@ describe("LMDB low-level database", () => { } }); + it("flushes delayed commits before close resolves", async () => { + const path = await tempLmdbPath(); + const db = declareLmdbLowLevelDatabase({ path, dbId: "close-drain" }); + try { + const store = db.declareKvStore("store"); + await store.setAll([{ key: buffer("key"), value: buffer("durable") }]); + // setAll intentionally returns before the 10ms commit batch is submitted. + // close must flush that application-level queue before closing LMDB. + await db.close(); + + const reopened = declareLmdbLowLevelDatabase({ path, dbId: "close-drain" }); + try { + const reopenedStore = reopened.declareKvStore("store"); + expect(text((await reopenedStore.get(buffer("key"))).buffer)).toBe("durable"); + } finally { + await reopened.close(); + } + } finally { + await db.close(); + await rm(path, { recursive: true, force: true }); + } + }); + it("supports compareAndSet without advancing seq on failed comparisons", async () => { const path = await tempLmdbPath(); try { diff --git a/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts b/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts index 747d9fd02..e0d5b3e29 100644 --- a/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts +++ b/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts @@ -146,6 +146,8 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri let pendingCommitOperations: PendingCommitOperation[] = []; let pendingCommitFlushTimer: ReturnType | null = null; let pendingCommitFlushPromise: Promise | null = null; + let isClosing = false; + let closePromise: Promise | null = null; let activityStats = emptyActivityStats(); let activityWindowStartedAt = performance.now(); if (!shouldSuppressPeriodicBulldozerLogs) { @@ -208,11 +210,12 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri return seqToDurability.get(seqId) ?? combinedSeqToDurability.get(seqId) ?? Promise.resolve(); }; // LMDB may reject with an opaque "Commit failed" wrapper whose real status - // lives on `.commitError` (a Promise). Unwrap before surfacing to callers. - const awaitLmdbPromise = (promise: Promise) => promise.catch(async (error) => { + // lives on `.commitError` (a Promise). LMDB's `committed` value is only a + // PromiseLike, so normalize it before using native Promise methods. + const awaitLmdbPromise = (promise: PromiseLike) => Promise.resolve(promise).catch(async (error) => { throw await unwrapLmdbCommitError(error); }); - const rememberAvailability = (seqId: string, promise: Promise) => { + const rememberAvailability = (seqId: string, promise: PromiseLike) => { const insertedAt = performance.now(); const availability = traceSpanHot({ description: "bulldozer-js.low-level.lmdb.availability", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => { await awaitLmdbPromise(promise); @@ -223,7 +226,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri availability.catch(() => {}); seqToAvailability.set(seqId, availability); }; - const rememberDurability = (seqId: string, promise: Promise) => { + const rememberDurability = (seqId: string, promise: PromiseLike) => { const insertedAt = performance.now(); const durability = traceSpanHot({ description: "bulldozer-js.low-level.lmdb.durability", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => { await awaitLmdbPromise(promise); @@ -235,7 +238,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri durability.catch(() => {}); seqToDurability.set(seqId, durability); }; - const rememberCombinedAvailability = (seqId: string, promise: Promise) => { + const rememberCombinedAvailability = (seqId: string, promise: PromiseLike) => { const insertedAt = performance.now(); const availability = traceSpanHot({ description: "bulldozer-js.low-level.lmdb.combinedAvailability", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => { await awaitLmdbPromise(promise); @@ -247,7 +250,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri availability.catch(() => {}); combinedSeqToAvailability.set(seqId, availability); }; - const rememberCombinedDurability = (seqId: string, promise: Promise) => { + const rememberCombinedDurability = (seqId: string, promise: PromiseLike) => { const insertedAt = performance.now(); const durability = traceSpanHot({ description: "bulldozer-js.low-level.lmdb.combinedDurability", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => { await awaitLmdbPromise(promise); @@ -258,7 +261,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri durability.catch(() => {}); combinedSeqToDurability.set(seqId, durability); }; - const trackCommit = (seqId: string, promise: Promise) => { + const trackCommit = (seqId: string, promise: PromiseLike) => { rememberAvailability(seqId, promise); rememberDurability(seqId, promise); return toSeq(seqId); @@ -321,6 +324,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri }, 10); }; const commit = (requiresSeq: DatabaseSeq, action: (version: number) => Promise) => { + if (isClosing) throw new Error("LMDB database is closing and cannot accept writes"); const seqId = nextSeqId(); const deferred = createVoidDeferred(); pendingCommitOperations.push({ seqId, requiresSeq, action, resolve: deferred.resolve, reject: deferred.reject }); @@ -328,6 +332,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri return trackCommit(seqId, deferred.promise); }; const commitIfVersion = async (db: BinaryDatabase, key: Buffer, version: number, action: (version: number) => Promise) => { + if (isClosing) throw new Error("LMDB database is closing and cannot accept writes"); const nextVersionRef: { value: number | null } = { value: null }; const seqId = nextSeqId(); const wasSet = await db.ifVersion(key, version, () => { @@ -524,6 +529,25 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri rememberCombinedDurability(seqId, Promise.all(seqs.map(seq => getDurabilityPromise(getSeqId(seq))))); return toSeq(seqId); }, + close() { + if (closePromise === null) { + isClosing = true; + closePromise = traceSpanHot({ description: "bulldozer-js.low-level.lmdb.close", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => { + try { + if (pendingCommitFlushPromise !== null) await pendingCommitFlushPromise; + await flushPendingCommits(); + if (pendingCommitFlushPromise !== null) await pendingCommitFlushPromise; + await Promise.all([ + ...seqToDurability.values(), + ...combinedSeqToDurability.values(), + ]); + } finally { + await root.close(); + } + }); + } + return closePromise; + }, async debugSnapshot() { return await traceSpanHot({ description: "bulldozer-js.low-level.lmdb.debugSnapshot", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => { const stores: Record = {}; diff --git a/apps/bulldozer-js/src/databases/piledriver/index.ts b/apps/bulldozer-js/src/databases/piledriver/index.ts index 7e6e3737d..ad71918fb 100644 --- a/apps/bulldozer-js/src/databases/piledriver/index.ts +++ b/apps/bulldozer-js/src/databases/piledriver/index.ts @@ -700,6 +700,9 @@ export function declarePiledriverDatabase(lowLevelDb: LowLevelDatabase, options: combineSeqs(...seqs) { return lowLevelDb.combineSeqs(...seqs); }, + async close() { + await lowLevelDb.close(); + }, waitUntilAvailable(seq) { return traceSpan("bulldozer-js.piledriver.waitUntilAvailable", async () => await lowLevelDb.waitUntilAvailable(seq)); }, diff --git a/apps/bulldozer-js/src/index.ts b/apps/bulldozer-js/src/index.ts index 04a3e92d1..77127eb3e 100644 --- a/apps/bulldozer-js/src/index.ts +++ b/apps/bulldozer-js/src/index.ts @@ -910,6 +910,9 @@ const bulldozerServerSecret = process.env.HEXCLAVE_BULLDOZER_SERVER_SECRET ?? "" if (bulldozerServerSecret.length === 0) { throw new HexclaveAssertionError("bulldozer-js refuses to start: HEXCLAVE_BULLDOZER_SERVER_SECRET is not set. This shared secret authenticates the backend to bulldozer-js and must be configured (it has a default in .env.development for local/self-host)."); } +let isShuttingDown = false; +let shutdownPromise: Promise | null = null; +let stopHttpServer: (() => void | Promise) | undefined; const app = new Elysia({ adapter: node() }) .use(instrumentation) @@ -924,6 +927,10 @@ const app = new Elysia({ adapter: node() }) } }) .get("/health", () => ({ ok: true })) + .post("/internal/wait-until-durable", () => handler("wait-until-durable", async () => { + await bulldozerDb.waitUntilCurrentStateDurable(); + return ok(); + })) .post("/internal/payments/verify-data-integrity", () => handler("verify-data-integrity", async () => ok())) .get("/v1/:tenancyId/transactions", ({ params, query }) => handler("list-transactions", async () => { const parsedLimit = Number.parseInt(typeof query.limit === "string" ? query.limit : "50", 10); @@ -1030,9 +1037,47 @@ const app = new Elysia({ adapter: node() }) .post("/v1/:tenancyId/test-mode/subscriptions/:subscriptionId/end", () => notImplemented("end-test-mode-subscription")) .post("/v1/:tenancyId/test-mode/one-time-purchases", () => notImplemented("create-test-mode-one-time-purchase")) .post("/v1/:tenancyId/test-mode/subscriptions/:subscriptionId/switch", () => notImplemented("switch-test-mode-subscription")) - .listen(port); + .listen(port, (server) => { + // @elysiajs/node 1.4.5 does not assign the server to app.server, so + // Elysia.stop() throws even though the callback's server is running. + stopHttpServer = () => server.stop(); + }); console.log(`Bulldozer JS server listening on http://localhost:${app.server?.port ?? port}`); + +function shutDown(): Promise { + if (shutdownPromise === null) { + isShuttingDown = true; + shutdownPromise = traceSpan("bulldozer-js.shutdown", async () => { + try { + // Do not force-close active requests; database.close() serializes behind + // any in-flight write or tick before draining the retained write seq. + if (stopHttpServer === undefined) throw new Error("Bulldozer HTTP server started without providing a stop handle"); + await stopHttpServer(); + } finally { + await bulldozerDb.close(); + } + }); + } + return shutdownPromise; +} + +function requestShutdown(signal: "SIGINT" | "SIGTERM") { + logBulldozerService("service-shutdown-requested", { signal }); + runAsynchronously(async () => { + await shutDown(); + process.exitCode = signal === "SIGINT" ? 130 : 0; + logBulldozerService("service-shutdown-complete", { signal, exitCode: process.exitCode }); + }, { + onError: () => { + process.exitCode = 1; + }, + }); +} + +process.on("SIGINT", () => requestShutdown("SIGINT")); +process.on("SIGTERM", () => requestShutdown("SIGTERM")); + const startupFields = { port: app.server?.port ?? port, pid: process.pid, @@ -1065,7 +1110,7 @@ if (sentryEnabled) { // backwards wall-clock jump can't rewind it. runAsynchronously(async () => { let lastTickMillis = 0; - while (true) { + while (!isShuttingDown) { await traceSpan("bulldozer-js-tick-loop-iteration", async () => { const tickStartedAt = performance.now(); try { diff --git a/apps/dashboard/src/app/(main)/(protected)/projects/[projectId]/cli-auth/page-client.tsx b/apps/dashboard/src/app/(main)/(protected)/projects/[projectId]/cli-auth/page-client.tsx index c103eb1c7..afe6163ad 100644 --- a/apps/dashboard/src/app/(main)/(protected)/projects/[projectId]/cli-auth/page-client.tsx +++ b/apps/dashboard/src/app/(main)/(protected)/projects/[projectId]/cli-auth/page-client.tsx @@ -175,7 +175,7 @@ function CliAuthDashboard({ data }: { data: CliAuthData }) { CLI Auth allows users to authenticate from command-line tools using a browser-based login flow. The CLI initiates a session, the user confirms in a browser, and a refresh token is issued to the CLI. Metrics below are bounded snapshots: attempt counts cover the newest {summary.attempt_window_limit} attempts, - and active sessions are discovered from the newest {summary.active_token_lookup_window_limit} used attempts. + and CLI sessions are discovered from the newest {summary.active_token_lookup_window_limit} attempts. } /> @@ -189,7 +189,7 @@ function CliAuthDashboard({ data }: { data: CliAuthData }) { diff --git a/apps/e2e/tests/backend/endpoints/api/v1/internal/cli-auth.test.ts b/apps/e2e/tests/backend/endpoints/api/v1/internal/cli-auth.test.ts index 66badb6cb..c98e72c89 100644 --- a/apps/e2e/tests/backend/endpoints/api/v1/internal/cli-auth.test.ts +++ b/apps/e2e/tests/backend/endpoints/api/v1/internal/cli-auth.test.ts @@ -1,5 +1,5 @@ import { it } from "../../../../../helpers"; -import { Project, niceBackendFetch } from "../../../../backend-helpers"; +import { Auth, Project, niceBackendFetch } from "../../../../backend-helpers"; it("returns bounded CLI authentication metrics without exposing login secrets", async ({ expect }) => { await Project.createAndSwitch(); @@ -39,3 +39,47 @@ it("returns bounded CLI authentication metrics without exposing login secrets", expect(JSON.stringify(response.body)).not.toContain(createdAttempt.body.polling_code); expect(JSON.stringify(response.body)).not.toContain(createdAttempt.body.login_code); }); + +it("returns completed CLI sessions with a count matching the active rows", async ({ expect }) => { + await Project.createAndSwitch(); + const user = await Auth.fastSignUp(); + + const createdAttempt = await niceBackendFetch("/api/latest/auth/cli", { + method: "POST", + accessType: "server", + body: {}, + }); + const completedAttempt = await niceBackendFetch("/api/latest/auth/cli/complete", { + method: "POST", + accessType: "server", + body: { + login_code: createdAttempt.body.login_code, + mode: "complete", + refresh_token: user.refreshToken, + }, + }); + expect(completedAttempt.status).toBe(200); + + const polledAttempt = await niceBackendFetch("/api/latest/auth/cli/poll", { + method: "POST", + accessType: "server", + body: { polling_code: createdAttempt.body.polling_code }, + }); + expect(polledAttempt.status).toBe(201); + + const response = await niceBackendFetch("/api/latest/internal/cli-auth", { + method: "GET", + accessType: "admin", + }); + + expect(response.status).toBe(200); + expect(response.body.summary.active_tokens_in_lookup_window).toBe(1); + expect(response.body.active_cli_users).toHaveLength(1); + expect(response.body.active_cli_users[0]).toMatchObject({ + user_id: user.userId, + is_expired: false, + }); + expect(response.body.active_cli_users.filter((session: { is_expired: boolean }) => !session.is_expired)).toHaveLength( + response.body.summary.active_tokens_in_lookup_window, + ); +}); diff --git a/scripts/reseed-bulldozer.sh b/scripts/reseed-bulldozer.sh index ecec4314b..3bc72c1f3 100755 --- a/scripts/reseed-bulldozer.sh +++ b/scripts/reseed-bulldozer.sh @@ -18,6 +18,7 @@ BULLDOZER_DIR="$REPO_ROOT/apps/bulldozer-js" DATA_DIR="$BULLDOZER_DIR/.data" TEMP_DATA_DIR="" BACKUP_DATA_DIR="$BULLDOZER_DIR/.data-backup.untracked.$$" +BULLDOZER_SECRET="$(node -e 'console.log(require("node:crypto").randomUUID())')" BULLDOZER_PID="" port_is_open() { @@ -48,6 +49,23 @@ stop_bulldozer() { BULLDOZER_PID="" } +stop_bulldozer_cleanly() { + if [[ -z "$BULLDOZER_PID" ]]; then + echo "Temporary bulldozer-js is not running; refusing to install its store." >&2 + return 1 + fi + local pid="$BULLDOZER_PID" + echo "Gracefully shutting down temporary bulldozer-js (pid $pid) ..." + kill "$pid" 2>/dev/null || true + local process_status=0 + wait "$pid" || process_status=$? + BULLDOZER_PID="" + if [[ "$process_status" -ne 0 ]]; then + echo "Temporary bulldozer-js did not shut down cleanly (status $process_status); keeping the existing store." >&2 + return 1 + fi +} + cleanup() { stop_bulldozer if [[ -n "$TEMP_DATA_DIR" && -d "$TEMP_DATA_DIR" ]]; then @@ -78,6 +96,8 @@ echo "Starting temporary bulldozer-js on port $BULLDOZER_PORT ..." ( cd "$BULLDOZER_DIR" && \ NODE_ENV=development \ + BULLDOZER_JS_PORT="$BULLDOZER_PORT" \ + HEXCLAVE_BULLDOZER_SERVER_SECRET="$BULLDOZER_SECRET" \ HEXCLAVE_BULLDOZER_JS_LMDB_PATH="$TEMP_LMDB_PATH" \ exec node --import tsx --expose-gc src/index.ts ) & @@ -105,9 +125,22 @@ if [[ "$ready" -ne 1 ]]; then fi echo "Running Postgres->bulldozer backfill ..." -pnpm run db:backfill-bulldozer-from-prisma +HEXCLAVE_BULLDOZER_SERVER_SECRET="$BULLDOZER_SECRET" \ + HEXCLAVE_BULLDOZER_SERVER_URL="http://127.0.0.1:$BULLDOZER_PORT" \ + pnpm run db:backfill-bulldozer-from-prisma -stop_bulldozer +echo "Waiting for the replacement store to become durable ..." +curl --silent --show-error --fail-with-body \ + --retry 3 \ + --retry-all-errors \ + --connect-timeout 10 \ + --max-time 300 \ + --request POST \ + --header "Authorization: Bearer $BULLDOZER_SECRET" \ + "http://127.0.0.1:$BULLDOZER_PORT/internal/wait-until-durable" \ + --output /dev/null + +stop_bulldozer_cleanly echo "Replacing bulldozer-js LMDB store at $DATA_DIR ..." if [[ -e "$DATA_DIR" ]]; then