diff --git a/.github/workflows/project-intake.yml b/.github/workflows/project-intake.yml index 3221f02..bcfe25f 100644 --- a/.github/workflows/project-intake.yml +++ b/.github/workflows/project-intake.yml @@ -1,6 +1,8 @@ name: Project Intake on: + schedule: + - cron: "17 * * * *" issues: types: [opened, reopened, closed] workflow_dispatch: @@ -15,14 +17,14 @@ permissions: issues: read concurrency: - group: ${{ github.workflow }}-${{ github.event.issue.number || inputs.issue_number || github.run_id }} - cancel-in-progress: true + group: project-intake-${{ github.event.issue.number || inputs.issue_number || github.run_id }} + cancel-in-progress: false jobs: sync: name: Sync issue Project fields runs-on: ubuntu-latest - timeout-minutes: 10 + timeout-minutes: 30 env: BASE_PROJECT_OWNER: ${{ github.repository_owner }} BASE_PROJECT_TITLE: ${{ github.event.repository.name }} @@ -33,244 +35,275 @@ jobs: BASE_PROJECT_DEFAULT_SIZE: S BASE_PROJECT_DEFAULT_AREA: Product BASE_PROJECT_DEFAULT_INITIATIVE: Adoption Polish + BASE_PROJECT_NUMBER: 14 GH_TOKEN: ${{ secrets.BASE_PROJECT_TOKEN }} steps: - name: Reconcile Project item - shell: bash + shell: python run: | - set -euo pipefail - - if [[ -z "${GH_TOKEN:-}" ]]; then - echo "::error::BASE_PROJECT_TOKEN secret is required for Project Intake." - echo "::error::Fix: gh auth token | gh secret set BASE_PROJECT_TOKEN --repo $GITHUB_REPOSITORY" - exit 1 - fi - - issue_number="${BASE_PROJECT_ISSUE_NUMBER:-}" - if [[ -z "$issue_number" ]]; then - echo "::error::Issue number was not provided by the event or workflow_dispatch input." - exit 1 - fi - - project_intake_is_auth_failure() { - local output="$1" - - printf '%s\n' "$output" | grep -Eiq 'Bad credentials|401 Unauthorized|HTTP 401|authentication failed' - } - - project_intake_is_retryable_api_failure() { - local output="$1" - - printf '%s\n' "$output" | grep -Eiq 'rate limit|secondary rate limit|abuse detection|Retry-After|x-ratelimit-reset' - } - - project_intake_retry_delay_seconds() { - local output="$1" - local reset_epoch - local now - local delay - - if [[ "$output" =~ Retry-After:?[[:space:]]*([0-9]+) ]]; then - printf '%s\n' "${BASH_REMATCH[1]}" - return 0 - fi - - if [[ "$output" =~ (x-ratelimit-reset|X-RateLimit-Reset):?[[:space:]]*([0-9]+) ]]; then - reset_epoch="${BASH_REMATCH[2]}" - now="$(date +%s)" - delay=$(( reset_epoch - now + 5 )) - if (( delay < 5 )); then - delay=5 - fi - if (( delay > 300 )); then - delay=300 - fi - printf '%s\n' "$delay" - return 0 - fi - - # Project V2 operations consume the shared GraphQL quota, while - # the REST rate-limit endpoint remains available for discovering - # the actual reset time. Avoid a misleading fixed 60-second retry - # when the quota will not recover within this workflow attempt. - reset_epoch="$(gh api rate_limit --jq '.resources.graphql.reset' 2>/dev/null || true)" - if [[ "$reset_epoch" =~ ^[0-9]+$ ]]; then - now="$(date +%s)" - delay=$(( reset_epoch - now + 5 )) - if (( delay < 5 )); then - delay=5 - fi - if (( delay > 300 )); then - delay=300 - fi - printf '%s\n' "$delay" - return 0 - fi - - printf '60\n' - } - - project_intake_gh() { - local operation="$1" - shift - local output - local error_output - local diagnostic_output - local stderr_file - local status - local retry_delay - - stderr_file="$(mktemp)" - if output="$("$@" 2>"$stderr_file")"; then - error_output="$(cat "$stderr_file")" - rm -f "$stderr_file" - [[ -z "$error_output" ]] || printf '%s\n' "$error_output" >&2 - [[ -z "$output" ]] || printf '%s\n' "$output" - return 0 - else - status=$? - error_output="$(cat "$stderr_file")" - rm -f "$stderr_file" - fi - - diagnostic_output="$output" - if [[ -n "$error_output" ]]; then - [[ -z "$diagnostic_output" ]] || diagnostic_output+=$'\n' - diagnostic_output+="$error_output" - fi - - if project_intake_is_auth_failure "$diagnostic_output"; then - echo "::error::GitHub authentication failed during Project Intake: $operation" >&2 - echo "::error::Rotate BASE_PROJECT_TOKEN and rerun this workflow_dispatch for issue #${issue_number:-unknown}." >&2 - echo "::error::Fix: gh auth token | gh secret set BASE_PROJECT_TOKEN --repo $GITHUB_REPOSITORY" >&2 - printf '%s\n' "$diagnostic_output" >&2 - return "$status" - fi - - if project_intake_is_retryable_api_failure "$diagnostic_output"; then - retry_delay="$(project_intake_retry_delay_seconds "$diagnostic_output")" - echo "::warning::GitHub API pressure during Project Intake: $operation" >&2 - echo "::warning::Waiting ${retry_delay}s and retrying once." >&2 - sleep "$retry_delay" - stderr_file="$(mktemp)" - if output="$("$@" 2>"$stderr_file")"; then - error_output="$(cat "$stderr_file")" - rm -f "$stderr_file" - [[ -z "$error_output" ]] || printf '%s\n' "$error_output" >&2 - [[ -z "$output" ]] || printf '%s\n' "$output" - return 0 - else - status=$? - error_output="$(cat "$stderr_file")" - rm -f "$stderr_file" - fi - diagnostic_output="$output" - if [[ -n "$error_output" ]]; then - [[ -z "$diagnostic_output" ]] || diagnostic_output+=$'\n' - diagnostic_output+="$error_output" - fi - echo "::error::GitHub API retry failed during Project Intake: $operation" >&2 - printf '%s\n' "$diagnostic_output" >&2 - return "$status" - fi - - echo "::error::GitHub API command failed during Project Intake: $operation" >&2 - printf '%s\n' "$diagnostic_output" >&2 - return "$status" - } - - # Issue state and URL are REST fields. Using `gh issue view --json` - # here needlessly consumes the same GraphQL quota required for the - # Project V2 mutations below. - issue_json="$(project_intake_gh "read issue" gh api "repos/$GITHUB_REPOSITORY/issues/$issue_number" --jq '{state: (.state | ascii_upcase), url: .html_url}')" - issue_state="$(jq -r '.state' <<<"$issue_json")" - issue_url="$(jq -r '.url' <<<"$issue_json")" - - project_number="$( - project_intake_gh "list Projects" gh project list --owner "$BASE_PROJECT_OWNER" --format json --limit 100 | - jq -r --arg title "$BASE_PROJECT_TITLE" \ - '.projects[] | select(.title == $title) | .number' | - head -n 1 - )" - if [[ -z "$project_number" ]]; then - echo "::error::GitHub Project '$BASE_PROJECT_TITLE' was not found for owner '$BASE_PROJECT_OWNER'." - echo "::error::If this Project exists, set BASE_PROJECT_TOKEN with user Project read/write access." - echo "::error::Fix: gh auth token | gh secret set BASE_PROJECT_TOKEN --repo $GITHUB_REPOSITORY" - exit 1 - fi - - project_id="$(project_intake_gh "view Project" gh project view "$project_number" --owner "$BASE_PROJECT_OWNER" --format json --jq '.id')" - item_id="$(project_intake_gh "add Project item" gh project item-add "$project_number" --owner "$BASE_PROJECT_OWNER" --url "$issue_url" --format json --jq '.id')" - item_json="$( - project_intake_gh "list Project items" gh project item-list "$project_number" --owner "$BASE_PROJECT_OWNER" --format json --limit 1000 | - jq --arg id "$item_id" '.items[] | select(.id == $id)' - )" - fields_json="$(project_intake_gh "list Project fields" gh project field-list "$project_number" --owner "$BASE_PROJECT_OWNER" --format json)" - - field_id_for() { - local field_name="$1" - - jq -r --arg name "$field_name" \ - '.fields[] | select(.name == $name) | .id' <<<"$fields_json" | - head -n 1 - } - - option_id_for() { - local field_name="$1" - local option_name="$2" - - jq -r --arg name "$field_name" --arg option "$option_name" \ - '.fields[] | select(.name == $name) | .options[]? | select(.name == $option) | .id' \ - <<<"$fields_json" | - head -n 1 - } - - set_single_select() { - local field_name="$1" - local option_name="$2" - local field_id - local option_id - - [[ -n "$option_name" ]] || return 0 - - field_id="$(field_id_for "$field_name")" - option_id="$(option_id_for "$field_name" "$option_name")" - if [[ -z "$field_id" || -z "$option_id" ]]; then - echo "::error::Project field '$field_name' option '$option_name' was not found." - exit 1 - fi - - project_intake_gh "set Project field $field_name" gh project item-edit \ - --id "$item_id" \ - --project-id "$project_id" \ - --field-id "$field_id" \ - --single-select-option-id "$option_id" \ - >/dev/null - } - - set_single_select_if_missing() { - local field_name="$1" - local item_key="$2" - local option_name="$3" - local current_value - - current_value="$(jq -r --arg key "$item_key" '.[$key] // ""' <<<"$item_json")" - if [[ -n "$current_value" ]]; then - return 0 - fi - - set_single_select "$field_name" "$option_name" - } - - status_value="$BASE_PROJECT_DEFAULT_OPEN_STATUS" - if [[ "$issue_state" == "CLOSED" ]]; then - status_value="$BASE_PROJECT_DEFAULT_CLOSED_STATUS" - fi - - set_single_select Status "$status_value" - set_single_select_if_missing Priority priority "$BASE_PROJECT_DEFAULT_PRIORITY" - set_single_select_if_missing Size size "$BASE_PROJECT_DEFAULT_SIZE" - set_single_select_if_missing Area area "$BASE_PROJECT_DEFAULT_AREA" - set_single_select_if_missing Initiative initiative "$BASE_PROJECT_DEFAULT_INITIATIVE" - - printf 'Synced issue #%s into Project %s.\n' "$issue_number" "$BASE_PROJECT_TITLE" + import json + import os + import re + import subprocess + import sys + import time + + + class DeferredReconciliation(RuntimeError): + pass + + + calls = 0 + sweep_cache = {} + deadline = time.monotonic() + 1500 + + + def api(method, endpoint, *, payload=None, params=None, paginate=False): + """Use only REST, keeping JSON stdout separate from gh diagnostics.""" + global calls + cache_key = (endpoint, json.dumps(params, sort_keys=True)) + cacheable = (os.environ.get("GITHUB_EVENT_NAME") == "schedule" and method == "GET" + and (endpoint.endswith("/fields") or endpoint.endswith("/items") + or re.fullmatch(r"orgs/[^/]+/projectsV2/[0-9]+", endpoint))) + if cacheable and cache_key in sweep_cache: + return sweep_cache[cache_key] + command = ["gh", "api", "--method", method, endpoint, + "-H", "Accept: application/vnd.github+json"] + if paginate: + command += ["--paginate", "--slurp"] + else: + command += ["--include"] + for key, value in (params or {}).items(): + command += ["-f", f"{key}={value}"] + if payload is not None: + command += ["--input", "-"] + for attempt in range(3): + if time.monotonic() + 65 >= deadline: + raise DeferredReconciliation("Workflow time budget reached; the hourly full sweep will retry.") + calls += 1 + try: + result = subprocess.run( + command, input=json.dumps(payload) if payload is not None else None, + capture_output=True, text=True, timeout=60, check=False, + ) + if result.returncode == 0: + body = result.stdout + if body.startswith("HTTP/"): + body = re.split(r"\r?\n\r?\n", body, maxsplit=1)[1] + data = json.loads(body) + value = [item for page in data for item in page] if paginate else data + if cacheable: + sweep_cache[cache_key] = value + return value + detail = result.stderr + "\n" + result.stdout + except subprocess.TimeoutExpired: + detail = "GitHub request timed out" + token = os.environ.get("GH_TOKEN", "") + if token: + detail = detail.replace(token, "[REDACTED]") + retryable = re.search( + r"HTTP (429|5\d\d)|rate limit|secondary rate|timed out|" + r"Could not resolve host|connection (reset|refused)", detail, re.I, + ) + delay = 2 ** attempt + retry_after = re.search(r"Retry-After:?\s*(\d+)", detail, re.I) + if retry_after: + delay = max(delay, int(retry_after[1])) + reset = re.search(r"x-ratelimit-reset:\s*(\d+)", detail, re.I) + if reset: + delay = max(delay, int(reset[1]) - int(time.time()) + 1) + # Never spin until an hourly quota reset or ignore a long Retry-After. + if not retryable: + raise RuntimeError(f"{method} {endpoint} failed: {detail.strip()}") + if attempt == 2 or delay > 30: + raise DeferredReconciliation( + f"{method} {endpoint} deferred; hourly full reconciliation will retry. " + f"Requested delay: {delay}s. {detail.strip()}" + ) + print(f"Retrying {method} {endpoint} in {delay}s.", file=sys.stderr) + time.sleep(delay) + + + def invalidate_items_cache(project): + prefix = f"{project}/items" + for key in list(sweep_cache): + if key[0] == prefix: + del sweep_cache[key] + + + def main(issue_override=None): + if not os.environ.get("GH_TOKEN"): + raise RuntimeError("BASE_PROJECT_TOKEN is required for organization Project writes.") + number = os.environ.get("BASE_PROJECT_ISSUE_NUMBER", "") + if not re.fullmatch(r"[1-9][0-9]*", number): + raise RuntimeError("Issue number must be a positive integer.") + repo = os.environ["GITHUB_REPOSITORY"] + owner = os.environ["BASE_PROJECT_OWNER"] + project_number = os.environ["BASE_PROJECT_NUMBER"] + if not re.fullmatch(r"[1-9][0-9]*", project_number): + raise RuntimeError("Project number must be a positive integer.") + project = f"orgs/{owner}/projectsV2/{project_number}" + metadata = api("GET", project) + if metadata.get("title") != os.environ["BASE_PROJECT_TITLE"]: + raise RuntimeError("Configured Project does not match the repository Project title.") + issue = issue_override or api("GET", f"repos/{repo}/issues/{number}") + if "pull_request" in issue or issue.get("state") not in ("open", "closed"): + raise RuntimeError("Intake requires an open or closed issue.") + + defaults = { + "Status": os.environ["BASE_PROJECT_DEFAULT_OPEN_STATUS"], + "Priority": os.environ["BASE_PROJECT_DEFAULT_PRIORITY"], + "Size": os.environ["BASE_PROJECT_DEFAULT_SIZE"], + "Area": os.environ["BASE_PROJECT_DEFAULT_AREA"], + "Initiative": os.environ["BASE_PROJECT_DEFAULT_INITIATIVE"], + } + fields = api("GET", f"{project}/fields", paginate=True, params={"per_page": "100"}) + managed = {} + for name in defaults: + matches = [field for field in fields if field["name"] == name] + if len(matches) != 1 or matches[0].get("data_type") != "single_select": + raise RuntimeError(f"Expected exactly one single-select field: {name}") + managed[name] = matches[0] + field_ids = ",".join(str(field["id"]) for field in managed.values()) + + def find_item(): + # Scan every Project page, avoiding search grammar or index + # visibility dependencies, then match immutable issue identity. + items = api("GET", f"{project}/items", paginate=True, + params={"per_page": "100"}) + matches = [item for item in items if item.get("content_type") == "Issue" + and (item.get("content") or {}).get("id") == issue["id"]] + if len(matches) > 1: + raise RuntimeError("Multiple Project items match the exact issue.") + return matches[0] if matches else None + + item = find_item() + if item is None: + try: + added = api("POST", f"{project}/items", payload={"type": "Issue", "id": issue["id"]}) + item = added.get("value", added) + if not item.get("id"): + raise RuntimeError("Project item add did not return an item id.") + except RuntimeError as error: + # Another intake/Project automation may win the add race. + if "content already exists in this project" not in str(error).lower(): + raise + for attempt in range(3): + # The first lookup is cached during scheduled + # sweeps. Invalidate it before retrying so a + # concurrent add can become visible. + invalidate_items_cache(project) + item = find_item() + if item is not None: + break + if attempt < 2: + time.sleep(attempt + 1) + if item is None: + raise RuntimeError("Existing Project item is not visible after bounded retries.") from error + + item_endpoint = f"{project}/items/{item['id']}" + + def read_item(): + for attempt in range(3): + try: + current = api("GET", item_endpoint, params={"fields": field_ids}) + except RuntimeError as error: + if "HTTP 404" not in str(error): + raise + if attempt == 2: + raise RuntimeError("Project item is not visible after bounded retries.") from error + else: + if (current.get("content_type") != "Issue" + or (current.get("content") or {}).get("id") != issue["id"]): + raise RuntimeError("Project item identity did not match the requested issue.") + return current + time.sleep(attempt + 1) + + def values(current): + return {field["id"]: (field.get("value") or {}).get("id") + for field in current.get("fields", [])} + + def option_id(field, label): + matches = [option["id"] for option in field.get("options", []) + if (option["name"].get("raw") if isinstance(option["name"], dict) + else option["name"]) == label] + if len(matches) != 1: + raise RuntimeError(f"Project field {field['name']!r} option {label!r} was not found uniquely.") + return matches[0] + + current_values = values(read_item()) + expected = {} + updates = [] + closed_status = os.environ["BASE_PROJECT_DEFAULT_CLOSED_STATUS"] + for name, field in managed.items(): + current = current_values.get(field["id"]) + desired = current + if name == "Status" and issue["state"] == "closed": + desired = option_id(field, closed_status) + elif name == "Status" and current == option_id(field, closed_status): + desired = option_id(field, defaults[name]) + elif not current: + desired = option_id(field, defaults[name]) + expected[field["id"]] = desired + if current != desired: + updates.append({"id": field["id"], "value": desired}) + if updates: + api("PATCH", item_endpoint, payload={"fields": updates}) + status_field = managed["Status"] + open_status_ids = {option["id"] for option in status_field["options"]} + open_status_ids.discard(option_id(status_field, closed_status)) + for attempt in range(3): + verified = values(read_item()) + verification_expected = dict(expected) + # Linked-PR automation or a maintainer can advance an open + # issue while intake initializes its other fields. Preserve + # that valid status just as we preserve it before the write. + observed_status = verified.get(status_field["id"]) + if issue["state"] == "open" and observed_status in open_status_ids: + verification_expected[status_field["id"]] = observed_status + if all(verified.get(field_id) == value for field_id, value in verification_expected.items()): + if verification_expected != expected: + print("Preserved a concurrent open-issue Project status change.") + print(f"Synced issue #{number} into Project {project_number} via REST; verified all five fields.") + return + if attempt < 2: + time.sleep(attempt + 1) + mismatches = [name for name, field in managed.items() + if verified.get(field["id"]) != verification_expected[field["id"]]] + raise RuntimeError("Project field readback did not match the intended values after bounded retries: " + + ", ".join(mismatches)) + + + try: + if os.environ.get("GITHUB_EVENT_NAME") == "schedule": + # A full sweep is a durable recovery path even if the original + # event exhausted quota before any Project item was created. + issues = api("GET", f"repos/{os.environ['GITHUB_REPOSITORY']}/issues", + paginate=True, params={"state": "all", "per_page": "100"}) + failures = [] + for issue in issues: + if "pull_request" not in issue: + os.environ["BASE_PROJECT_ISSUE_NUMBER"] = str(issue["number"]) + try: + main(issue) + except DeferredReconciliation: + raise + except (RuntimeError, ValueError, KeyError, TypeError, OSError) as error: + failures.append(f"#{issue['number']}: {error}") + if failures: + raise RuntimeError( + "Scheduled Project sweep completed with failures; successful issues were preserved: " + + "; ".join(failures) + ) + else: + main() + print(f"Project Intake REST calls: {calls}; GraphQL calls/points: 0.") + except DeferredReconciliation as error: + message = str(error).replace("%", "%25").replace("\r", "%0D").replace("\n", "%0A") + print(f"::notice::Reconciliation deferred: {message}", file=sys.stderr) + print(f"Project Intake REST calls: {calls}; GraphQL calls/points: 0; status: deferred.") + except (RuntimeError, ValueError, KeyError, TypeError, OSError) as error: + # Escape annotation control characters; never claim success on partial sync. + message = str(error).replace("%", "%25").replace("\r", "%0D").replace("\n", "%0A") + print(f"::error::{message}", file=sys.stderr) + sys.exit(1) diff --git a/docs/testing.md b/docs/testing.md index eb3f143..049fd2e 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -74,3 +74,13 @@ exclusions. The typing gate checks every Git-visible Python source outside `lib/ (checked separately), `scripts/` (validation tools), and `tests/` (test harnesses) with strict mypy. This includes example and compatibility consumer packages and new top-level source directories; untracked sources are included during development. +### Project Intake recovery + +Project Intake uses REST exclusively, verifies each managed field by independent +readback, and preserves existing planning metadata and active statuses. Primary +or secondary quota pressure beyond its bounded retry window reports `deferred`, +not a successful sync. An hourly full issue sweep retries missing and stale cards, +including events that failed before item creation; manual dispatch remains available. +The sweep caches stable Project metadata/fields/items during the job. Authentication +and permission errors fail explicitly and require operator repair. Logs report REST +calls and zero GraphQL calls/points; token values are redacted from diagnostics. diff --git a/tests/test_project_intake.py b/tests/test_project_intake.py new file mode 100644 index 0000000..472312a --- /dev/null +++ b/tests/test_project_intake.py @@ -0,0 +1,509 @@ +#!/usr/bin/env python3 +"""Execute the actual checkout-free workflow with a stateful REST boundary fake.""" + +import contextlib +import copy +import io +import json +import os +import subprocess +import textwrap +import unittest +from pathlib import Path +from unittest.mock import patch + +ROOT = Path(__file__).resolve().parents[1] +WORKFLOW = (ROOT / ".github/workflows/project-intake.yml").read_text() + + +def extract_reconcile_script(workflow): + """Read the named step's literal block without adding a YAML dependency.""" + lines = workflow.splitlines(keepends=True) + step_marker = " - name: Reconcile Project item" + matches = [index for index, line in enumerate(lines) if line.rstrip() == step_marker] + if len(matches) != 1: + raise AssertionError("Expected exactly one workflow step named Reconcile Project item.") + start = matches[0] + 1 + end = next( + ( + index + for index in range(start, len(lines)) + if lines[index].strip() and not lines[index].startswith(" ") + ), + len(lines), + ) + step = lines[start:end] + run_markers = [index for index, line in enumerate(step) if line.rstrip() == " run: |"] + if len(run_markers) != 1: + raise AssertionError("Reconcile Project item must contain exactly one literal run: | block.") + start = run_markers[0] + 1 + end = next( + ( + index + for index in range(start, len(step)) + if step[index].strip() and not step[index].startswith(" ") + ), + len(step), + ) + script = textwrap.dedent("".join(step[start:end])).strip() + if not script: + raise AssertionError("Reconcile Project item has an empty run block.") + return script + "\n" + + +SCRIPT = extract_reconcile_script(WORKFLOW) +DEFAULTS = {"Status": "Backlog", "Priority": "P2", "Size": "S", "Area": "Product", "Initiative": "Adoption Polish"} + + +class RestFixture: + def __init__(self): + self.fields = [] + for field_id, (name, default) in enumerate(DEFAULTS.items(), 1): + options = [default] + if name == "Status": + options += ["Done", "In Progress", "In Review"] + else: + options += ["Custom"] + self.fields.append( + { + "id": field_id, + "name": name, + "data_type": "single_select", + "options": [{"id": label, "name": {"raw": label}} for label in options], + } + ) + self.current = {} + self.exists = True + self.closed = False + self.duplicate = False + self.delayed = 0 + self.delayed_search = 0 + self.mismatch = False + self.wrong_identity = False + self.calls = [] + self.failures = [] + self.direct_add = False + self.stale_readback = 0 + self.patch_failure = False + self.concurrent_status = "" + self.filtered_search_empty = False + self.item_read_failure = "" + + def item(self): + return { + "id": 101, + "content_type": "Issue", + "content": {"id": 42 if not self.wrong_identity else 99}, + "fields": [{"id": key, "value": {"id": value}} for key, value in self.current.items()], + } + + def run(self, command, **kwargs): + self.calls.append((command, kwargs)) + # GraphQL is unavailable in every scenario, including the happy path. + if command[:2] != ["gh", "api"] or "graphql" in command: + return subprocess.CompletedProcess(command, 1, "", "GraphQL quota exhausted") + method, endpoint = command[3:5] + self.assert_safe_request(command, kwargs) + if self.failures: + failure = self.failures.pop(0) + if isinstance(failure, BaseException): + raise failure + if failure: + return subprocess.CompletedProcess(command, 1, "", failure) + result = None + if endpoint.endswith("projectsV2/14"): + result = {"title": "base-cli"} + elif endpoint == "repos/basefoundry/base-cli/issues": + result = [[{"number": 495, "id": 42, "state": "open"}, {"number": 496, "pull_request": {}}]] + elif endpoint == "repos/basefoundry/base-cli/issues/495": + result = {"id": 42, "state": "closed" if self.closed else "open"} + elif endpoint.endswith("/fields"): + # Both fields and item discovery must work beyond the first page. + result = [self.fields[:2], self.fields[2:]] + elif endpoint.endswith("/items"): + if method == "POST": + self.exists = True + if self.duplicate: + return subprocess.CompletedProcess( + command, 1, "", "Content already exists in this project (HTTP 422)" + ) + result = self.item() if self.direct_add else {"value": self.item()} + elif self.filtered_search_empty and any(arg.startswith("q=") for arg in command): + result = [[]] + elif not self.exists or self.delayed_search: + self.delayed_search = max(0, self.delayed_search - 1) + result = [[]] + else: + result = [ + [ + { + "id": 100, + "content_type": "Issue", + "content": { + "id": 41, + "number": 495, + "repository_url": "https://api.github.com/repos/another/repo", + }, + }, + {"id": 102, "content_type": "DraftIssue", "content": {"id": 42}}, + {"id": 103, "content_type": "Issue", "content": None}, + ], + [self.item()], + ] + elif endpoint.endswith("/items/101"): + if method == "PATCH": + if self.patch_failure: + return subprocess.CompletedProcess(command, 1, "", "Forbidden (HTTP 403)") + if not self.mismatch: + for field in json.loads(kwargs["input"])["fields"]: + self.current[field["id"]] = field["value"] + if self.concurrent_status: + self.current[1] = self.concurrent_status + result = self.item() + elif self.item_read_failure: + return subprocess.CompletedProcess(command, 1, "", self.item_read_failure) + elif self.delayed: + self.delayed -= 1 + return subprocess.CompletedProcess(command, 1, "", "Not Found (HTTP 404)") + elif self.stale_readback and any(c[0][3] == "PATCH" for c in self.calls): + self.stale_readback -= 1 + result = {**self.item(), "fields": []} + else: + result = self.item() + else: + raise AssertionError(f"Unexpected REST request: {command}") + return subprocess.CompletedProcess(command, 0, json.dumps(result), "nonfatal gh notice") + + @staticmethod + def assert_safe_request(command, kwargs): + assert kwargs["timeout"] == 60 + assert kwargs["capture_output"] is True + assert not kwargs.get("shell") + if command[3] == "GET" and command[4].endswith(("/fields", "/items")): + assert "--paginate" in command and "--slurp" in command + + +class ProjectIntakeTests(unittest.TestCase): + def execute(self, fixture=None, **overrides): + fixture = fixture or RestFixture() + env = { + "GH_TOKEN": "test-only", + "GITHUB_REPOSITORY": "basefoundry/base-cli", + "BASE_PROJECT_OWNER": "basefoundry", + "BASE_PROJECT_TITLE": "base-cli", + "BASE_PROJECT_NUMBER": "14", + "BASE_PROJECT_ISSUE_NUMBER": "495", + "BASE_PROJECT_DEFAULT_CLOSED_STATUS": "Done", + } + env.update( + {f"BASE_PROJECT_DEFAULT_{name.upper()}": value for name, value in DEFAULTS.items() if name != "Status"} + ) + env["BASE_PROJECT_DEFAULT_OPEN_STATUS"] = "Backlog" + env.update(overrides) + stdout, stderr = io.StringIO(), io.StringIO() + with ( + patch.dict(os.environ, env, clear=True), + patch("subprocess.run", fixture.run), + patch("time.sleep") as sleep, + contextlib.redirect_stdout(stdout), + contextlib.redirect_stderr(stderr), + ): + status = 0 + try: + exec(compile(SCRIPT, "project-intake.yml", "exec"), {"__name__": "__main__"}) + except SystemExit as error: + status = error.code + return status, stdout.getvalue(), stderr.getvalue(), sleep + + def test_hourly_sweep_recovers_an_event_without_an_item(self): + fixture = RestFixture() + fixture.exists = False + status, stdout, stderr, _ = self.execute(fixture, GITHUB_EVENT_NAME="schedule", BASE_PROJECT_ISSUE_NUMBER="") + self.assertEqual(status, 0, stderr) + self.assertIn("verified all five fields", stdout) + self.assertTrue(fixture.exists) + self.assertEqual(fixture.current, {i: v for i, v in enumerate(DEFAULTS.values(), 1)}) + + def test_primary_reset_headers_defer_and_redact_token(self): + fixture = RestFixture() + fixture.failures = ["rate limit (HTTP 403) x-ratelimit-reset: 9999999999 token-secret"] + status, stdout, stderr, sleep = self.execute(fixture, GH_TOKEN="token-secret") + self.assertEqual(status, 0) + self.assertIn("status: deferred", stdout) + self.assertNotIn("token-secret", stdout + stderr) + sleep.assert_not_called() + + def test_partial_patch_failure_recovers_on_rerun(self): + fixture = RestFixture() + fixture.current = {1: "In Progress", 2: "Custom"} + fixture.patch_failure = True + self.assertEqual(self.execute(fixture)[0], 1) + fixture.patch_failure = False + status, stdout, stderr, _ = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertEqual(fixture.current[1], "In Progress") + self.assertEqual(fixture.current[2], "Custom") + self.assertIn("verified all five fields", stdout) + + def test_extraction_ignores_neighboring_run_steps_and_metadata(self): + before = " - name: Setup\n run: |\n echo setup\n" + after = " env:\n EXAMPLE: value\n - name: Finish\n run: |\n echo finish\n" + workflow = WORKFLOW.replace(" steps:\n", " steps:\n" + before) + "\n" + after + self.assertEqual(extract_reconcile_script(workflow), SCRIPT) + + def test_extraction_reports_missing_or_duplicate_named_step(self): + for workflow in ( + WORKFLOW.replace("name: Reconcile Project item", "name: Renamed"), + WORKFLOW + "\n - name: Reconcile Project item\n", + ): + with self.subTest(workflow=workflow[-80:]): + with self.assertRaisesRegex(AssertionError, "exactly one workflow step"): + extract_reconcile_script(workflow) + + def test_extraction_reports_unsupported_or_empty_run_block(self): + for workflow in ( + WORKFLOW.replace(" run: |\n", " run: >\n"), + " - name: Reconcile Project item\n run: |\n", + ): + with self.subTest(workflow=workflow[-80:]): + with self.assertRaisesRegex(AssertionError, "Reconcile Project item.*run"): + extract_reconcile_script(workflow) + + def test_workflow_uses_trusted_inline_python_and_secret(self): + self.assertIn("shell: python", WORKFLOW) + self.assertIn("GH_TOKEN: ${{ secrets.BASE_PROJECT_TOKEN }}", WORKFLOW) + self.assertIn("timeout-minutes: 30", WORKFLOW) + self.assertIn("cancel-in-progress: false", WORKFLOW) + self.assertNotIn("uses:", WORKFLOW) + self.assertNotIn("github.token", SCRIPT) + + def test_graphql_unavailable_still_syncs_and_reads_back(self): + fixture = RestFixture() + status, stdout, stderr, _ = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertIn("verified all five fields", stdout) + self.assertEqual(fixture.current, dict(enumerate(DEFAULTS.values(), 1))) + self.assertEqual(sum(c[0][3] == "PATCH" for c in fixture.calls), 1) + self.assertEqual(fixture.calls[-1][0][3:5], ["GET", "orgs/basefoundry/projectsV2/14/items/101"]) + + def test_preserves_nonempty_metadata_and_active_status_without_writes(self): + for status_value in ("In Progress", "In Review"): + with self.subTest(status=status_value): + fixture = RestFixture() + fixture.current = {1: status_value, 2: "Custom", 3: "Custom", 4: "Custom", 5: "Custom"} + expected = copy.deepcopy(fixture.current) + status, _, stderr, _ = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertEqual(fixture.current, expected) + self.assertTrue(all(c[0][3] == "GET" for c in fixture.calls)) + + def test_close_and_reopen(self): + fixture = RestFixture() + fixture.closed = True + self.assertEqual(self.execute(fixture)[0], 0) + self.assertEqual(fixture.current[1], "Done") + fixture.closed = False + self.assertEqual(self.execute(fixture)[0], 0) + self.assertEqual(fixture.current[1], "Backlog") + + def test_add_response_shapes_and_delayed_visibility(self): + for direct in (False, True): + with self.subTest(direct=direct): + fixture = RestFixture() + fixture.exists = False + fixture.direct_add = direct + fixture.delayed = 2 + status, _, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertEqual(sleep.call_count, 2) + self.assertEqual(sum(c[0][3] == "POST" for c in fixture.calls), 1) + + def test_duplicate_add_race_recovers_exact_item(self): + fixture = RestFixture() + fixture.exists = False + fixture.duplicate = True + fixture.delayed_search = 2 + status, _, stderr, _ = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertEqual(sum(c[0][3] == "POST" for c in fixture.calls), 1) + + def test_item_lookup_does_not_depend_on_search_filter_grammar(self): + fixture = RestFixture() + fixture.filtered_search_empty = True + fixture.duplicate = True + status, stdout, stderr, _ = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertIn("verified all five fields", stdout) + self.assertFalse(any(c[0][3] == "POST" for c in fixture.calls)) + self.assertFalse(any(arg.startswith("q=") for c in fixture.calls for arg in c[0])) + + def test_transient_rest_failures_retry(self): + for error in ( + "rate limit (HTTP 429) Retry-After: 7", + "Bad Gateway (HTTP 502)", + subprocess.TimeoutExpired("gh", 60), + ): + with self.subTest(error=error): + fixture = RestFixture() + fixture.failures = [error] + status, _, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertEqual(sleep.call_count, 1) + if isinstance(error, str) and "Retry-After" in error: + sleep.assert_called_once_with(7) + + def test_auth_and_permission_failures_are_not_retried(self): + for error in ("Bad credentials (HTTP 401)", "Forbidden (HTTP 403)"): + with self.subTest(error=error): + fixture = RestFixture() + fixture.failures = [error] + status, stdout, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn(error, stderr) + sleep.assert_not_called() + self.assertEqual(len(fixture.calls), 1) + + def test_persistent_rate_limit_is_deferred_for_scheduled_recovery(self): + fixture = RestFixture() + fixture.failures = ["rate limit (HTTP 429)"] * 3 + status, stdout, _, sleep = self.execute(fixture) + self.assertEqual(status, 0) + self.assertNotIn("Synced issue", stdout) + self.assertEqual(len(fixture.calls), 3) + self.assertEqual(sleep.call_count, 2) + + def test_long_retry_after_defers_without_early_retry(self): + fixture = RestFixture() + fixture.failures = ["rate limit (HTTP 429) Retry-After: 3600"] + status, _, _, sleep = self.execute(fixture) + self.assertEqual(status, 0) + sleep.assert_not_called() + self.assertEqual(len(fixture.calls), 1) + + def test_missing_token_and_invalid_input_fail_before_api(self): + for overrides in ( + {"GH_TOKEN": ""}, + {"BASE_PROJECT_ISSUE_NUMBER": "0"}, + {"BASE_PROJECT_ISSUE_NUMBER": "495; echo unsafe"}, + ): + with self.subTest(overrides=overrides): + fixture = RestFixture() + self.assertEqual(self.execute(fixture, **overrides)[0], 1) + self.assertEqual(fixture.calls, []) + + def test_invalid_option_or_field_fails_before_patch(self): + for missing_field in (False, True): + fixture = RestFixture() + if missing_field: + fixture.fields.pop() + else: + fixture.fields[-1]["options"] = [] + status, stdout, _, _ = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertFalse(any(c[0][3] == "PATCH" for c in fixture.calls)) + + def test_wrong_item_identity_fails(self): + fixture = RestFixture() + fixture.exists = False + fixture.wrong_identity = True + status, stdout, stderr, _ = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn("identity did not match", stderr) + + def test_missing_item_after_add_fails_without_duplicate_writes(self): + fixture = RestFixture() + fixture.exists = False + fixture.delayed = 3 + status, stdout, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertEqual(sleep.call_count, 2) + self.assertEqual(sum(c[0][3] == "POST" for c in fixture.calls), 1) + self.assertIn("Project item is not visible after bounded retries.", stderr) + self.assertEqual( + sum(c[0][3:5] == ["GET", "orgs/basefoundry/projectsV2/14/items/101"] for c in fixture.calls), 3 + ) + + def test_non_404_item_read_failure_keeps_original_diagnostic(self): + fixture = RestFixture() + fixture.item_read_failure = "Forbidden (HTTP 403)" + status, stdout, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn("Forbidden (HTTP 403)", stderr) + self.assertNotIn("not visible after bounded retries", stderr) + sleep.assert_not_called() + + def test_readback_mismatch_never_claims_success(self): + fixture = RestFixture() + fixture.mismatch = True + status, stdout, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn("readback did not match", stderr) + self.assertEqual(sleep.call_count, 2) + + def test_failed_patch_never_claims_success(self): + fixture = RestFixture() + fixture.patch_failure = True + status, stdout, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn("Forbidden", stderr) + sleep.assert_not_called() + + def test_duplicate_item_that_never_becomes_visible_fails(self): + fixture = RestFixture() + fixture.exists = False + fixture.duplicate = True + fixture.delayed_search = 10 + status, stdout, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn("not visible after bounded retries", stderr) + self.assertEqual(sleep.call_count, 2) + + def test_stale_readback_recovers(self): + fixture = RestFixture() + fixture.stale_readback = 2 + status, _, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertEqual(sleep.call_count, 2) + + def test_preserves_concurrent_open_status_change(self): + for status_value in ("In Progress", "In Review"): + with self.subTest(status=status_value): + fixture = RestFixture() + fixture.concurrent_status = status_value + status, stdout, stderr, sleep = self.execute(fixture) + self.assertEqual(status, 0, stderr) + self.assertEqual(fixture.current[1], status_value) + self.assertIn("Preserved a concurrent open-issue Project status change", stdout) + self.assertEqual(sum(c[0][3] == "PATCH" for c in fixture.calls), 1) + sleep.assert_not_called() + + def test_open_status_readback_rejects_done_and_invalid_options(self): + for status_value in ("Done", "invalid-option"): + with self.subTest(status=status_value): + fixture = RestFixture() + fixture.concurrent_status = status_value + status, stdout, stderr, _ = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn("bounded retries: Status", stderr) + + def test_closed_status_readback_remains_strict(self): + fixture = RestFixture() + fixture.closed = True + fixture.concurrent_status = "In Progress" + status, stdout, stderr, _ = self.execute(fixture) + self.assertEqual(status, 1) + self.assertNotIn("Synced issue", stdout) + self.assertIn("bounded retries: Status", stderr) + + +if __name__ == "__main__": + unittest.main()