CreateAccessPolicy treated a Cloudflare "policy_already_exists" as idempotent success and kept whatever policy was there. On a re-run with a changed identity — or against a hand-made broader policy — op.console would stay guarded by something weaker than the fail-closed body this package builds and guards, while Setup reported success. The fail-closed validation only ever ran on the policy we built, never on the one that stayed live. Now it upserts by name: lookup, PUT the guarded body over the existing policy, POST only when absent (a racing POST re-looks up and PUTs). apiPost/apiPut share one apiWrite; three httptest cases pin update-over-existing, create-when- absent, and the race fallback.
431 lines
17 KiB
Go
431 lines
17 KiB
Go
package cfsetup
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"regexp"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// This file is INTEGRATION-ONLY. ExecRunner shells out to the real `cloudflared`
|
|
// binary and calls the live Cloudflare API; none of it can run — or be honestly
|
|
// faked — on a box without the operator's own Cloudflare account and the
|
|
// interactive `cloudflared tunnel login` consent already completed. The orchestration
|
|
// that uses it (Setup) and the request-body/guard logic are unit-verified in
|
|
// cfsetup.go; what lives here is exercised only against a real account.
|
|
|
|
const defaultAPIBase = "https://api.cloudflare.com/client/v4"
|
|
|
|
// DetectPreconditions inspects the local environment for the gating facts Setup
|
|
// needs: whether cloudflared is installed and whether the operator has logged in
|
|
// (cert.pem present). The API token is supplied by the caller (the TUI prompts
|
|
// for it); it is passed through so the returned value is ready to hand to Setup.
|
|
// This only READS the environment — it performs no Cloudflare side effects — but
|
|
// it touches the real filesystem/PATH, so it is integration-side.
|
|
func DetectPreconditions(apiToken string) Preconditions {
|
|
pre := Preconditions{APIToken: apiToken}
|
|
if path, err := exec.LookPath("cloudflared"); err == nil {
|
|
pre.CloudflaredPath = path
|
|
}
|
|
if home, err := os.UserHomeDir(); err == nil {
|
|
if _, err := os.Stat(filepath.Join(home, ".cloudflared", "cert.pem")); err == nil {
|
|
pre.CertExists = true
|
|
}
|
|
}
|
|
return pre
|
|
}
|
|
|
|
// ExecRunner is the production Runner: cloudflared via os/exec for the tunnel and
|
|
// the Cloudflare API via HTTP for Access.
|
|
type ExecRunner struct {
|
|
// Cloudflared is the resolved cloudflared binary path (Preconditions.CloudflaredPath).
|
|
Cloudflared string
|
|
// APIToken authenticates the Access API calls (Bearer).
|
|
APIToken string
|
|
// AccountID is the Cloudflare account the Access app/policy are created under.
|
|
AccountID string
|
|
// APIBase defaults to the public Cloudflare API; overridable for testing.
|
|
APIBase string
|
|
// HTTP is the client used for API calls; nil means a default with a timeout.
|
|
HTTP *http.Client
|
|
}
|
|
|
|
func (r *ExecRunner) httpClient() *http.Client {
|
|
if r.HTTP != nil {
|
|
return r.HTTP
|
|
}
|
|
return &http.Client{Timeout: 30 * time.Second}
|
|
}
|
|
|
|
func (r *ExecRunner) apiBase() string {
|
|
if r.APIBase != "" {
|
|
return r.APIBase
|
|
}
|
|
return defaultAPIBase
|
|
}
|
|
|
|
// VerifyAPIToken satisfies the apiTokenVerifier seam Setup probes for: it does a
|
|
// read-only Cloudflare API call so a bad, expired, wrong-account or
|
|
// under-permissioned token fails BEFORE any tunnel/DNS/config is created, instead
|
|
// of surfacing late at CreateAccessApplication with a half-built edge left behind.
|
|
//
|
|
// It lists Access apps (per_page=1 — the cheapest authenticated read) against the
|
|
// exact account and permission Setup will write to, so a green result means the
|
|
// credential that actually gates the side-effecting calls works. The tunnel and
|
|
// DNS authenticate via cert.pem, not this token, so the Access read is precisely
|
|
// the credential worth pre-checking. apiGet surfaces 401/403 with the actionable
|
|
// permission checklist; this method only adds the cheap pre-flight argument
|
|
// validation so an empty token/account never reaches the wire.
|
|
func (r *ExecRunner) VerifyAPIToken(ctx context.Context) error {
|
|
if strings.TrimSpace(r.APIToken) == "" {
|
|
return fmt.Errorf("cfsetup: Cloudflare API token is required")
|
|
}
|
|
if strings.TrimSpace(r.AccountID) == "" {
|
|
return fmt.Errorf("cfsetup: Cloudflare account ID is required")
|
|
}
|
|
return r.apiGet(ctx, fmt.Sprintf("/accounts/%s/access/apps?per_page=1", r.AccountID), nil)
|
|
}
|
|
|
|
// tunnelIDRE extracts the UUID cloudflared prints when a tunnel is created or
|
|
// already exists.
|
|
var tunnelIDRE = regexp.MustCompile(`[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}`)
|
|
|
|
// CreateTunnel runs `cloudflared tunnel create <name>`. cloudflared writes the
|
|
// credentials JSON under ~/.cloudflared/<id>.json and prints the id; we parse it
|
|
// out. If the tunnel already exists this returns its id (idempotent re-run).
|
|
func (r *ExecRunner) CreateTunnel(ctx context.Context, name string) (string, string, error) {
|
|
var id string
|
|
out, err := r.runCloudflared(ctx, "tunnel", "create", name)
|
|
if err != nil {
|
|
// An "already exists" is not fatal — recover the id via `tunnel list`.
|
|
lid, lerr := r.lookupTunnel(ctx, name)
|
|
if lerr != nil || lid == "" {
|
|
return "", "", err
|
|
}
|
|
id = lid
|
|
} else if id = tunnelIDRE.FindString(out); id == "" {
|
|
return "", "", fmt.Errorf("cfsetup: could not parse tunnel id from cloudflared output: %s", out)
|
|
}
|
|
cred := r.credentialsPath(id)
|
|
if err := r.ensureCredentials(ctx, id, cred); err != nil {
|
|
return "", "", err
|
|
}
|
|
return id, cred, nil
|
|
}
|
|
|
|
// ensureCredentials guarantees the tunnel credentials JSON exists at credPath.
|
|
// cloudflared writes that file only at `tunnel create`, so an idempotent re-run
|
|
// against a tunnel that already exists — or a reset+re-bootstrap where the old
|
|
// box's ~/.cloudflared was wiped but the Cloudflare-side tunnel survived — finds
|
|
// no local file, and the connector crash-loops with "Tunnel credentials file
|
|
// doesn't exist". `cloudflared tunnel token --cred-file` re-fetches the token into
|
|
// the file (authenticating with cert.pem, keeping the same id/DNS/Access), healing
|
|
// the re-run. The secret is written to the file, not stdout.
|
|
func (r *ExecRunner) ensureCredentials(ctx context.Context, id, credPath string) error {
|
|
// Any existing file counts as healthy; re-fetch only on absence
|
|
// (the failure actually seen). A truncated/zero-byte file would still
|
|
// crash-loop — validate the JSON here if that ever shows up.
|
|
if _, err := os.Stat(credPath); err == nil {
|
|
return nil
|
|
}
|
|
if err := os.MkdirAll(filepath.Dir(credPath), 0o700); err != nil {
|
|
return fmt.Errorf("cfsetup: preparing credentials dir for tunnel %s: %w", id, err)
|
|
}
|
|
if _, err := r.runCloudflared(ctx, "tunnel", "token", "--cred-file", credPath, id); err != nil {
|
|
return fmt.Errorf("cfsetup: tunnel %s credentials file %s is missing and could not be regenerated: %w", id, credPath, err)
|
|
}
|
|
// The credentials file is a secret sitting next to cert.pem; don't rely on
|
|
// cloudflared's umask to keep it owner-only.
|
|
if err := os.Chmod(credPath, 0o600); err != nil {
|
|
return fmt.Errorf("cfsetup: securing credentials file %s: %w", credPath, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// lookupTunnel finds an existing tunnel's id by name via `tunnel list`.
|
|
func (r *ExecRunner) lookupTunnel(ctx context.Context, name string) (string, error) {
|
|
out, err := r.runCloudflared(ctx, "tunnel", "list", "--name", name, "--output", "json")
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
var tunnels []struct {
|
|
ID string `json:"id"`
|
|
Name string `json:"name"`
|
|
}
|
|
if err := json.Unmarshal([]byte(out), &tunnels); err != nil {
|
|
return "", err
|
|
}
|
|
for _, t := range tunnels {
|
|
if t.Name == name {
|
|
return t.ID, nil
|
|
}
|
|
}
|
|
return "", fmt.Errorf("cfsetup: tunnel %q not found", name)
|
|
}
|
|
|
|
func (r *ExecRunner) credentialsPath(id string) string {
|
|
if home, err := os.UserHomeDir(); err == nil {
|
|
return filepath.Join(home, ".cloudflared", id+".json")
|
|
}
|
|
return id + ".json"
|
|
}
|
|
|
|
// RouteDNS runs `cloudflared tunnel route dns --overwrite-dns <tunnelID> <hostname>`,
|
|
// creating (or repointing) the proxied CNAME so hostname resolves to THIS tunnel.
|
|
//
|
|
// --overwrite-dns is load-bearing, not cosmetic. Without it, when a record for
|
|
// hostname already exists — most commonly a stale CNAME left by an earlier tunnel
|
|
// that was created and later deleted/recreated on the same box — cloudflared refuses
|
|
// with "record already exists" and changes nothing, leaving the name bound to the
|
|
// dead tunnel. The old code swallowed exactly that error as success, so a re-run
|
|
// reported "routed" while the hostname kept returning Cloudflare error 1033: the
|
|
// tunnel it still pointed at had no connector. Overwriting repoints the record at the
|
|
// tunnel we just created, making the route idempotent AND correct on every re-run.
|
|
func (r *ExecRunner) RouteDNS(ctx context.Context, tunnelID, hostname string) error {
|
|
_, err := r.runCloudflared(ctx, "tunnel", "route", "dns", "--overwrite-dns", tunnelID, hostname)
|
|
return err
|
|
}
|
|
|
|
// WriteTunnelConfig writes the rendered config.yml, creating its parent directory.
|
|
func (r *ExecRunner) WriteTunnelConfig(path string, contents []byte) error {
|
|
if dir := filepath.Dir(path); dir != "" {
|
|
if err := os.MkdirAll(dir, 0o755); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return os.WriteFile(path, contents, 0o644)
|
|
}
|
|
|
|
// CreateAccessApplication POSTs the self-hosted Access app and returns its id and
|
|
// issued aud (spec §14: the aud felis [auth] access_jwt_aud must adopt).
|
|
func (r *ExecRunner) CreateAccessApplication(ctx context.Context, app AccessApplication) (string, string, error) {
|
|
var resp struct {
|
|
Result struct {
|
|
ID string `json:"id"`
|
|
AUD string `json:"aud"`
|
|
} `json:"result"`
|
|
}
|
|
if err := r.apiPost(ctx, fmt.Sprintf("/accounts/%s/access/apps", r.AccountID), app, &resp); err != nil {
|
|
// If application already exists, look it up instead of failing (idempotency)
|
|
if strings.Contains(err.Error(), "application_already_exists") || strings.Contains(err.Error(), "11010") {
|
|
if id, aud, lerr := r.lookupAccessApplication(ctx, app.Domain); lerr == nil && id != "" {
|
|
return id, aud, nil
|
|
}
|
|
}
|
|
return "", "", err
|
|
}
|
|
return resp.Result.ID, resp.Result.AUD, nil
|
|
}
|
|
|
|
// CreateAccessPolicy makes the Access app carry exactly the recommended policy.
|
|
// It upserts by name instead of POST-then-swallow: when a policy of this name
|
|
// already exists — a re-run, or a previous hand-made setup — the swallow looked
|
|
// idempotent but left the OLD rule set in place. If that old policy is broader
|
|
// than the fail-closed one just built (a changed identity, a hand-made
|
|
// allow-everyone rule), op.console ends up guarded by something weaker while
|
|
// Setup reports success. So: find it and PUT our body over it; POST only when
|
|
// absent (a POST that races into "already exists" falls back to the PUT).
|
|
func (r *ExecRunner) CreateAccessPolicy(ctx context.Context, appID string, policy AccessPolicy) error {
|
|
pid, err := r.lookupAccessPolicy(ctx, appID, policy.Name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if pid != "" {
|
|
return r.apiPut(ctx, fmt.Sprintf("/accounts/%s/access/apps/%s/policies/%s", r.AccountID, appID, pid), policy, nil)
|
|
}
|
|
if err := r.apiPost(ctx, fmt.Sprintf("/accounts/%s/access/apps/%s/policies", r.AccountID, appID), policy, nil); err != nil {
|
|
if strings.Contains(err.Error(), "already_exists") || strings.Contains(err.Error(), "11015") {
|
|
if pid, lerr := r.lookupAccessPolicy(ctx, appID, policy.Name); lerr == nil && pid != "" {
|
|
return r.apiPut(ctx, fmt.Sprintf("/accounts/%s/access/apps/%s/policies/%s", r.AccountID, appID, pid), policy, nil)
|
|
}
|
|
}
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// lookupAccessPolicy finds the id of the policy named name on an Access app,
|
|
// returning "" when absent.
|
|
func (r *ExecRunner) lookupAccessPolicy(ctx context.Context, appID, name string) (string, error) {
|
|
var resp struct {
|
|
Result []struct {
|
|
ID string `json:"id"`
|
|
Name string `json:"name"`
|
|
} `json:"result"`
|
|
}
|
|
if err := r.apiGet(ctx, fmt.Sprintf("/accounts/%s/access/apps/%s/policies?per_page=100", r.AccountID, appID), &resp); err != nil {
|
|
return "", err
|
|
}
|
|
for _, p := range resp.Result {
|
|
if p.Name == name {
|
|
return p.ID, nil
|
|
}
|
|
}
|
|
return "", nil
|
|
}
|
|
|
|
// runCloudflared executes the cloudflared binary with the given args, returning
|
|
// combined output. The interactive `tunnel login` browser consent is NOT done
|
|
// here — it is a separate, operator-driven step the TUI suspends to run.
|
|
func (r *ExecRunner) runCloudflared(ctx context.Context, args ...string) (string, error) {
|
|
bin := r.Cloudflared
|
|
if bin == "" {
|
|
bin = "cloudflared"
|
|
}
|
|
cmd := exec.CommandContext(ctx, bin, args...)
|
|
var buf bytes.Buffer
|
|
cmd.Stdout = &buf
|
|
cmd.Stderr = &buf
|
|
if err := cmd.Run(); err != nil {
|
|
return buf.String(), fmt.Errorf("cfsetup: cloudflared %s: %w: %s", strings.Join(args, " "), err, buf.String())
|
|
}
|
|
return buf.String(), nil
|
|
}
|
|
|
|
// apiPost and apiPut send an authenticated JSON write to the Cloudflare API and,
|
|
// on a non-2xx or success:false body, return the error. out, when non-nil,
|
|
// receives the decoded response.
|
|
func (r *ExecRunner) apiPost(ctx context.Context, path string, body, out any) error {
|
|
return r.apiWrite(ctx, http.MethodPost, path, body, out)
|
|
}
|
|
|
|
func (r *ExecRunner) apiPut(ctx context.Context, path string, body, out any) error {
|
|
return r.apiWrite(ctx, http.MethodPut, path, body, out)
|
|
}
|
|
|
|
func (r *ExecRunner) apiWrite(ctx context.Context, method, path string, body, out any) error {
|
|
payload, err := json.Marshal(body)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, method, r.apiBase()+path, bytes.NewReader(payload))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
req.Header.Set("Authorization", "Bearer "+r.APIToken)
|
|
req.Header.Set("Content-Type", "application/json")
|
|
resp, err := r.httpClient().Do(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
raw, _ := io.ReadAll(resp.Body)
|
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
|
if resp.StatusCode == http.StatusUnauthorized {
|
|
return fmt.Errorf("cfsetup: Cloudflare API authentication failed (status 401). Please verify that:\n"+
|
|
" 1. The API Token is valid, active, and has not expired.\n"+
|
|
" 2. You did not enter a Global API Key (a Bearer API Token is required).\n"+
|
|
" 3. The token has the required permissions under the Account scope:\n"+
|
|
" - Account > Access Apps and Policies: Edit\n"+
|
|
" - Account > Cloudflare Tunnel: Edit\n"+
|
|
" - Zone > DNS: Edit\n"+
|
|
" Original error: %s", string(raw))
|
|
}
|
|
if resp.StatusCode == http.StatusForbidden {
|
|
return fmt.Errorf("cfsetup: Cloudflare API access forbidden (status 403). Please verify that:\n"+
|
|
" 1. The API Token has permission to access Account ID %q.\n"+
|
|
" 2. The token has the required permissions under the Account scope:\n"+
|
|
" - Account > Access Apps and Policies: Edit\n"+
|
|
" - Account > Cloudflare Tunnel: Edit\n"+
|
|
" - Zone > DNS: Edit\n"+
|
|
" Original error: %s", r.AccountID, string(raw))
|
|
}
|
|
return fmt.Errorf("cfsetup: Cloudflare API %s: status %d: %s", path, resp.StatusCode, string(raw))
|
|
}
|
|
// Cloudflare wraps every response in {success, errors, result}; surface a
|
|
// success:false even on a 200.
|
|
var envelope struct {
|
|
Success bool `json:"success"`
|
|
Errors []json.RawMessage `json:"errors"`
|
|
}
|
|
if err := json.Unmarshal(raw, &envelope); err == nil && !envelope.Success && len(envelope.Errors) > 0 {
|
|
return fmt.Errorf("cfsetup: Cloudflare API %s: %s", path, string(raw))
|
|
}
|
|
if out != nil {
|
|
if err := json.Unmarshal(raw, out); err != nil {
|
|
return fmt.Errorf("cfsetup: decode Cloudflare API %s response: %w", path, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// apiGet sends an authenticated JSON GET to the Cloudflare API and, on a
|
|
// non-2xx or success:false body, returns the error. out, when non-nil, receives
|
|
// the decoded response.
|
|
func (r *ExecRunner) apiGet(ctx context.Context, path string, out any) error {
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, r.apiBase()+path, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
req.Header.Set("Authorization", "Bearer "+r.APIToken)
|
|
resp, err := r.httpClient().Do(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
raw, _ := io.ReadAll(resp.Body)
|
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
|
if resp.StatusCode == http.StatusUnauthorized {
|
|
return fmt.Errorf("cfsetup: Cloudflare API authentication failed (status 401). Please verify that:\n"+
|
|
" 1. The API Token is valid, active, and has not expired.\n"+
|
|
" 2. You did not enter a Global API Key (a Bearer API Token is required).\n"+
|
|
" 3. The token has the required permissions under the Account scope:\n"+
|
|
" - Account > Access Apps and Policies: Edit\n"+
|
|
" - Account > Cloudflare Tunnel: Edit\n"+
|
|
" - Zone > DNS: Edit\n"+
|
|
" Original error: %s", string(raw))
|
|
}
|
|
if resp.StatusCode == http.StatusForbidden {
|
|
return fmt.Errorf("cfsetup: Cloudflare API access forbidden (status 403). Please verify that:\n"+
|
|
" 1. The API Token has permission to access Account ID %q.\n"+
|
|
" 2. The token has the required permissions under the Account scope:\n"+
|
|
" - Account > Access Apps and Policies: Edit\n"+
|
|
" - Account > Cloudflare Tunnel: Edit\n"+
|
|
" - Zone > DNS: Edit\n"+
|
|
" Original error: %s", r.AccountID, string(raw))
|
|
}
|
|
return fmt.Errorf("cfsetup: Cloudflare API %s: status %d: %s", path, resp.StatusCode, string(raw))
|
|
}
|
|
var envelope struct {
|
|
Success bool `json:"success"`
|
|
Errors []json.RawMessage `json:"errors"`
|
|
}
|
|
if err := json.Unmarshal(raw, &envelope); err == nil && !envelope.Success && len(envelope.Errors) > 0 {
|
|
return fmt.Errorf("cfsetup: Cloudflare API %s: %s", path, string(raw))
|
|
}
|
|
if out != nil {
|
|
if err := json.Unmarshal(raw, out); err != nil {
|
|
return fmt.Errorf("cfsetup: decode Cloudflare API %s response: %w", path, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// lookupAccessApplication finds an existing Access application's id and aud by domain.
|
|
func (r *ExecRunner) lookupAccessApplication(ctx context.Context, domain string) (string, string, error) {
|
|
var resp struct {
|
|
Result []struct {
|
|
ID string `json:"id"`
|
|
Domain string `json:"domain"`
|
|
AUD string `json:"aud"`
|
|
} `json:"result"`
|
|
}
|
|
if err := r.apiGet(ctx, fmt.Sprintf("/accounts/%s/access/apps?per_page=100", r.AccountID), &resp); err != nil {
|
|
return "", "", err
|
|
}
|
|
for _, app := range resp.Result {
|
|
if app.Domain == domain {
|
|
return app.ID, app.AUD, nil
|
|
}
|
|
}
|
|
return "", "", fmt.Errorf("cfsetup: access application for domain %q not found in list", domain)
|
|
}
|