diff --git a/docker-compose.yml b/docker-compose.yml index 5498035c4..5cbc7a9f9 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -925,6 +925,7 @@ services: restart: unless-stopped environment: SITESD_LISTEN: :8080 + SITESD_DEV: "true" SITESD_INTERNAL_ADDR: :8081 SITESD_WORKER_URL: http://sites:8787 SITESD_AUDIENCE: did:web:sites.tngl.boltless.dev @@ -947,18 +948,24 @@ services: condition: service_started sites: + dns: [11.0.0.254] build: context: . dockerfile: localinfra/sites.Dockerfile restart: unless-stopped init: true profiles: ["sites-worker"] + environment: + NODE_EXTRA_CA_CERTS: /usr/local/share/ca-certificates/caddy.crt volumes: - - ./sites:/sites:cached - sites-wrangler-state:/sites/.wrangler + - ./localinfra/certs/root.crt:/usr/local/share/ca-certificates/caddy.crt:ro ports: - "8787:8787" networks: [tngl] + depends_on: + dns: + condition: service_started volumes: caddy-data: diff --git a/localinfra/sites.Dockerfile b/localinfra/sites.Dockerfile index 92dd78c24..4efc2bed1 100644 --- a/localinfra/sites.Dockerfile +++ b/localinfra/sites.Dockerfile @@ -30,4 +30,4 @@ COPY --from=build /out/build ./build COPY sites/wrangler.dev.toml ./wrangler.dev.toml COPY sites/migrations ./migrations EXPOSE 8787 -CMD ["sh", "-c", "wrangler d1 migrations apply tangled-sites --local --config /sites/wrangler.dev.toml && exec wrangler dev --local --config /sites/wrangler.dev.toml --port 8787 --ip 0.0.0.0"] \ No newline at end of file +CMD ["sh", "-c", "update-ca-certificates && wrangler d1 migrations apply tangled-sites --local --config /sites/wrangler.dev.toml && exec wrangler dev --local --config /sites/wrangler.dev.toml --port 8787 --ip 0.0.0.0"] \ No newline at end of file diff --git a/sitesd/README.md b/sitesd/README.md new file mode 100644 index 000000000..ca6c193c0 --- /dev/null +++ b/sitesd/README.md @@ -0,0 +1,31 @@ +# sitesd cutover + +`sitesd` takes over site pushes after appview stops. do not restart appview +while sitesd is subscribed: its old sitefeed would also deploy pushes. sitesd +continues from the existing `sitefeed:cursor:*` Redis key; it does not call +appview or maintain a second cursor. start the optional sites worker with +`docker compose --profile sites-worker up`; without it, local site pushes +have no active deploy executor. compose sets `SITESD_DEV=true` so sitesd starts +without R2 credentials; the feed remains inactive until they are supplied. +production starts fail if any R2 credential is missing. + +for production, configure sitesd, the sites worker, Redis, and R2 before +stopping appview. start sitesd only after appview is down. set `SITESD_KNOTS`, +`SITESD_REDIS`, Cloudflare R2 credentials and bucket, and a private +`SITESD_INTERNAL_ADDR`. +Expose that private `/deploy` listener only to the sites worker. Set the +worker's `SITES_SD_URL` to that reachable private endpoint; the checked-in +production value is empty, so config writes save the row but report +`WorkerUnavailable` and cannot trigger a deploy. Do not deploy this cutover +as if the empty value were working. The worker verifies ownership using the repo DID document and +the knot's `describeRepo` answer, not a record-supplied owner. + +`sites-migrate -apply -yes -remote` imports only into an **empty** D1 target; +it refuses to overwrite live claims or deploy history. `-validate -apply` +checks parity after the import. Run the PDS claim backfill (`claimer`) after +import, and do not rerun the one-shot migration on a live D1. The standalone +count-parity check is for the initial import, before later claims and deploys +make the two databases intentionally different. Verify the new worker has the +site configs, then trigger initial deployments explicitly. sitesd resumes from +any saved sitefeed cursor; a fresh knot with no cursor starts live and does not +replay old pushes. diff --git a/sitesd/reposync/reposync.go b/sitesd/reposync/reposync.go index 502013586..63940df6f 100644 --- a/sitesd/reposync/reposync.go +++ b/sitesd/reposync/reposync.go @@ -14,9 +14,11 @@ type entry struct { } type Registry struct { - mu sync.Mutex - locks map[string]*entry - slots chan struct{} + mu sync.Mutex + locks map[string]*entry + slots chan struct{} + draining bool + active sync.WaitGroup } func New() *Registry { @@ -26,6 +28,23 @@ func New() *Registry { } } +func (r *Registry) Begin() (func(), bool) { + r.mu.Lock() + defer r.mu.Unlock() + if r.draining { + return nil, false + } + r.active.Add(1) + return r.active.Done, true +} + +func (r *Registry) Drain() { + r.mu.Lock() + r.draining = true + r.mu.Unlock() + r.active.Wait() +} + // Acquire takes a deploy slot; the returned func releases it. func (r *Registry) Acquire(ctx context.Context) (func(), error) { select { diff --git a/sitesd/reposync/reposync_test.go b/sitesd/reposync/reposync_test.go index 15c254625..f95bc97f3 100644 --- a/sitesd/reposync/reposync_test.go +++ b/sitesd/reposync/reposync_test.go @@ -89,3 +89,30 @@ func TestAcquireBoundsTheCeiling(t *testing.T) { release() } } + +func TestDrainWaitsForActiveDeploysAndRejectsNewWork(t *testing.T) { + r := New() + finish, ok := r.Begin() + if !ok { + t.Fatal("fresh registry refused a deploy") + } + drained := make(chan struct{}) + go func() { + r.Drain() + close(drained) + }() + select { + case <-drained: + t.Fatal("drained while a deploy was active") + case <-time.After(20 * time.Millisecond): + } + finish() + select { + case <-drained: + case <-time.After(time.Second): + t.Fatal("drain did not finish after the deploy") + } + if _, ok := r.Begin(); ok { + t.Fatal("accepted a deploy after drain") + } +} diff --git a/sitesd/sitefeed/feed.go b/sitesd/sitefeed/feed.go index f97fb6a88..41120e7e1 100644 --- a/sitesd/sitefeed/feed.go +++ b/sitesd/sitefeed/feed.go @@ -34,12 +34,13 @@ type feedSource struct { } type Feed struct { - rdb *cache.Cache - cfg *config.Config - cf *cloudflare.Client - wrk *worker.Client - locks *reposync.Registry - logger *slog.Logger + rdb *cache.Cache + cfg *config.Config + cf *cloudflare.Client + wrk *worker.Client + locks *reposync.Registry + deployCtx context.Context + logger *slog.Logger fetchCtx func(context.Context, string) (*worker.DeployContext, error) deploy func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error @@ -52,18 +53,19 @@ type Feed struct { sources map[string]*feedSource } -func New(rdb *cache.Cache, cfg *config.Config, cf *cloudflare.Client, wrk *worker.Client, locks *reposync.Registry, logger *slog.Logger) *Feed { +func New(rdb *cache.Cache, cfg *config.Config, cf *cloudflare.Client, wrk *worker.Client, locks *reposync.Registry, deployCtx context.Context, logger *slog.Logger) *Feed { return &Feed{ - rdb: rdb, - cfg: cfg, - cf: cf, - wrk: wrk, - locks: locks, - logger: log.SubLogger(logger, "sitefeed"), - fetchCtx: wrk.GetDeployContext, - deploy: sites.Deploy, - record: recordDeploy, - sources: make(map[string]*feedSource), + rdb: rdb, + cfg: cfg, + cf: cf, + wrk: wrk, + locks: locks, + deployCtx: deployCtx, + logger: log.SubLogger(logger, "sitefeed"), + fetchCtx: wrk.GetDeployContext, + deploy: sites.Deploy, + record: recordDeploy, + sources: make(map[string]*feedSource), } } @@ -217,6 +219,10 @@ func (f *Feed) refOp(ctx context.Context, host string, repoDid syntax.DID, op kn ctxInfo, err := f.fetchCtx(ctx, repoDid.String()) if err != nil { + if perRepoRefusal(err) { + logger.Warn("repo is not authorized for site deploys, skipping its ref", "err", err) + return nil + } return fmt.Errorf("fetching deploy context for %s: %w", repoDid, err) } if ctxInfo == nil { @@ -240,6 +246,19 @@ func (f *Feed) refOp(ctx context.Context, host string, repoDid syntax.DID, op kn return nil } +func perRepoRefusal(err error) bool { + var denied *worker.Error + if !errors.As(err, &denied) || denied.NSID != "org.tangled.temp.repo.getDeployContext" { + return false + } + switch denied.Name { + case "NotOwner", "RepoNotFound", "HandleNotOwned", "HandleUnknown", "SiteConfigDisputed": + return true + default: + return false + } +} + func (f *Feed) maybeDeploy(ctx context.Context, ctxInfo *worker.DeployContext, branch string, sha knotfeed.ObjectID) error { if f.cf == nil || !f.cf.Enabled() { return nil @@ -247,20 +266,27 @@ func (f *Feed) maybeDeploy(ctx context.Context, ctxInfo *worker.DeployContext, b if ctxInfo.Branch != branch { return nil } + done, ok := f.locks.Begin() + if !ok { + return context.Canceled + } release, ok := f.locks.TryAcquire() if !ok { + defer done() f.logger.Warn("deploy ceiling reached, dropping the deploy", "repo", ctxInfo.RepoDid, "branch", branch) f.recordDroppedDeploy(ctxInfo, sha.String(), "deploy ceiling reached") return nil } f.pendingSha.Store(ctxInfo.RepoDid, sha) go func() { + defer done() defer release() - deployCtx, cancel := context.WithTimeout(context.Background(), deployTimeout) + deployCtx, cancel := context.WithTimeout(f.deployCtx, deployTimeout) defer cancel() unlock, err := f.locks.Lock(deployCtx, ctxInfo.RepoDid) if err != nil { f.logger.Info("repo lock abandoned", "repo", ctxInfo.RepoDid, "branch", branch, "sha", sha, "err", err) + f.recordDroppedDeploy(ctxInfo, sha.String(), "repo lock abandoned: "+err.Error()) return } defer unlock() @@ -272,8 +298,9 @@ func (f *Feed) maybeDeploy(ctx context.Context, ctxInfo *worker.DeployContext, b f.logger.Info("duplicate push deploy skipped", "repo", ctxInfo.RepoDid, "sha", sha) return } - f.triggerDeploy(deployCtx, ctxInfo, sha.String()) - f.lastDeployed.Store(ctxInfo.RepoDid, sha) + if f.triggerDeploy(deployCtx, ctxInfo, sha.String()) { + f.lastDeployed.Store(ctxInfo.RepoDid, sha) + } }() return nil } @@ -288,12 +315,14 @@ func (f *Feed) recordDroppedDeploy(ctxInfo *worker.DeployContext, sha, reason st Status: models.SiteDeployStatusFailure, Error: reason, } - if err := f.record(context.Background(), f.wrk, deploy); err != nil { + recordCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + if err := f.record(recordCtx, f.wrk, deploy); err != nil { f.logger.Error("failed to record the dropped deploy", "repo", ctxInfo.RepoDid, "err", err) } } -func (f *Feed) triggerDeploy(ctx context.Context, ctxInfo *worker.DeployContext, sha string) { +func (f *Feed) triggerDeploy(ctx context.Context, ctxInfo *worker.DeployContext, sha string) bool { logger := f.logger.With("repo", ctxInfo.RepoDid) repo := &models.Repo{ @@ -323,11 +352,13 @@ func (f *Feed) triggerDeploy(ctx context.Context, ctxInfo *worker.DeployContext, recordCtx, recordCancel := context.WithTimeout(context.Background(), 30*time.Second) defer recordCancel() - if err := f.record(recordCtx, f.wrk, deploy); err != nil { - logger.Error("sites: failed to record deploy", "err", err) + recordErr := f.record(recordCtx, f.wrk, deploy) + if recordErr != nil { + logger.Error("sites: failed to record deploy", "err", recordErr) } - - if deployErr == nil { + if deployErr == nil && recordErr == nil { logger.Info("site deployed to r2") + return true } + return false } diff --git a/sitesd/sitefeed/feed_test.go b/sitesd/sitefeed/feed_test.go index e7d12c68d..fde40c2fb 100644 --- a/sitesd/sitefeed/feed_test.go +++ b/sitesd/sitefeed/feed_test.go @@ -159,7 +159,7 @@ func feedFor(t *testing.T, deployErr error) (*Feed, *deployRecorder, *recordReco t.Helper() d := &deployRecorder{} r := &recordRecorder{} - f := New(nil, &config.Config{}, enabledCf(t), &worker.Client{}, reposync.New(), log.New("test")) + f := New(nil, &config.Config{}, enabledCf(t), &worker.Client{}, reposync.New(), context.Background(), log.New("test")) f.deploy = d.hook(deployErr) f.record = r.hook() t.Cleanup(func() { @@ -169,6 +169,13 @@ func feedFor(t *testing.T, deployErr error) (*Feed, *deployRecorder, *recordReco return f, d, r } +func TestSitefeedCursorContinuesAppviewKey(t *testing.T) { + f, _, _ := feedFor(t, nil) + if got := f.cursorKey(feedTestHost); got != "sitefeed:cursor:"+feedTestHost { + t.Fatalf("sitefeed cursor = %q", got) + } +} + func TestFeedRefOp_BranchMatchTriggersDeploy(t *testing.T) { f, d, r := feedFor(t, nil) f.fetchCtx = func(_ context.Context, repoDid string) (*worker.DeployContext, error) { @@ -341,8 +348,34 @@ func TestFeedHandleFailsTheCommitWhenTheWorkerIsDown(t *testing.T) { } } +func TestFeedSkipsOneRepoRefusalWithoutStallingKnot(t *testing.T) { + f, d, _ := feedFor(t, nil) + msg := knotfeed.Message{Type: knotfeed.TypeCommit, Commit: &knotfeed.Commit{ + Repo: feedTestRepoDid, + Seq: 9, + Records: []knotfeed.RecordOp{ + refOpFor(t, "refs/heads/main", tapc.RecordCreateAction, feedTestSha), + }, + }} + f.fetchCtx = func(context.Context, string) (*worker.DeployContext, error) { + return nil, &worker.Error{StatusCode: 403, Name: "NotOwner", NSID: "org.tangled.temp.repo.getDeployContext"} + } + if err := f.handle(feedTestHost)(context.Background(), msg); err != nil { + t.Fatalf("one repo's permanent refusal stalled the knot feed: %v", err) + } + f.fetchCtx = func(context.Context, string) (*worker.DeployContext, error) { + return nil, &worker.Error{StatusCode: 503, Name: "OwnershipUnverifiable", NSID: "org.tangled.temp.repo.getDeployContext"} + } + if err := f.handle(feedTestHost)(context.Background(), msg); err == nil { + t.Fatal("transient resolver failure advanced the cursor") + } + if calls := d.got(); len(calls) != 0 { + t.Fatalf("unauthorized repo was deployed: %v", calls) + } +} + func TestFeedSubscribeDeduplicatesAndUnsubscribes(t *testing.T) { - f := New(nil, &config.Config{}, nil, &worker.Client{}, reposync.New(), log.New("test")) + f := New(nil, &config.Config{}, nil, &worker.Client{}, reposync.New(), context.Background(), log.New("test")) ctx, cancel := context.WithCancel(context.Background()) defer cancel() @@ -414,3 +447,122 @@ func TestFeedTrigger_DuplicateSameShaDeploysOnce(t *testing.T) { t.Errorf("deploy records = %d, want exactly 1 for a duplicated same-SHA frame", len(rows)) } } + +func TestFailedPushCanRetrySameCommit(t *testing.T) { + f, d, r := feedFor(t, errBoom) + sha := object(t, feedTestSha) + if err := f.maybeDeploy(context.Background(), ctxFor(), "main", sha); err != nil { + t.Fatal(err) + } + pollFor(t, "first failed push record", func() bool { return len(r.got()) == 1 }) + unlock, err := f.locks.Lock(context.Background(), feedTestRepoDid) + if err != nil { + t.Fatal(err) + } + unlock() + + if err := f.maybeDeploy(context.Background(), ctxFor(), "main", sha); err != nil { + t.Fatal(err) + } + pollFor(t, "retry of failed commit", func() bool { return len(d.got()) == 2 }) + pollFor(t, "second failed push record", func() bool { return len(r.got()) == 2 }) + for _, row := range r.got() { + if row.Status != models.SiteDeployStatusFailure { + t.Fatalf("retry row status = %q, want failure", row.Status) + } + } +} + +func TestUnrecordedPushCanRetrySameCommit(t *testing.T) { + f, d, r := feedFor(t, nil) + r.error = errBoom + sha := object(t, feedTestSha) + if err := f.maybeDeploy(context.Background(), ctxFor(), "main", sha); err != nil { + t.Fatal(err) + } + pollFor(t, "first failed history write", func() bool { return len(r.got()) == 1 }) + unlock, err := f.locks.Lock(context.Background(), feedTestRepoDid) + if err != nil { + t.Fatal(err) + } + unlock() + if err := f.maybeDeploy(context.Background(), ctxFor(), "main", sha); err != nil { + t.Fatal(err) + } + pollFor(t, "retry after failed history write", func() bool { return len(d.got()) == 2 }) +} + +func TestFeedShutdownCancelsPushAndKeepsItsFailureRecord(t *testing.T) { + f, _, r := feedFor(t, nil) + base, cancel := context.WithCancel(context.Background()) + f.deployCtx = base + started := make(chan struct{}) + f.deploy = func(ctx context.Context, _ *cloudflare.Client, _ *config.Config, _ *models.Repo, _, _ string) error { + close(started) + <-ctx.Done() + return ctx.Err() + } + if err := f.maybeDeploy(context.Background(), ctxFor(), "main", object(t, feedTestSha)); err != nil { + t.Fatal(err) + } + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("push never started") + } + cancel() + drained := make(chan struct{}) + go func() { f.locks.Drain(); close(drained) }() + select { + case <-drained: + case <-time.After(time.Second): + t.Fatal("cancelled push did not drain") + } + rows := r.got() + if len(rows) != 1 || rows[0].Status != models.SiteDeployStatusFailure { + t.Fatalf("cancelled push history = %v, want one failure", rows) + } +} + +func TestFeedDrainDoesNotAdvanceCursorPastPush(t *testing.T) { + f, _, _ := feedFor(t, nil) + f.locks.Drain() + if err := f.maybeDeploy(context.Background(), ctxFor(), "main", object(t, feedTestSha)); !errors.Is(err, context.Canceled) { + t.Fatalf("push during drain = %v, want retryable cancellation", err) + } +} + +func TestFeedDrainWaitsForPushDeploy(t *testing.T) { + f, _, _ := feedFor(t, nil) + started := make(chan struct{}) + release := make(chan struct{}) + f.deploy = func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(started) + <-release + return nil + } + if err := f.maybeDeploy(context.Background(), ctxFor(), "main", object(t, feedTestSha)); err != nil { + t.Fatalf("maybeDeploy: %v", err) + } + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("push deploy did not start") + } + drained := make(chan struct{}) + go func() { + f.locks.Drain() + close(drained) + }() + select { + case <-drained: + t.Fatal("drain returned during a push deploy") + case <-time.After(20 * time.Millisecond): + } + close(release) + select { + case <-drained: + case <-time.After(time.Second): + t.Fatal("drain never observed the push deploy finishing") + } +} diff --git a/sitesd/sitesd.go b/sitesd/sitesd.go index cb751ee92..470ab8476 100644 --- a/sitesd/sitesd.go +++ b/sitesd/sitesd.go @@ -9,7 +9,6 @@ import ( "net/http" "os" "strings" - "sync" "time" "github.com/bluesky-social/indigo/atproto/atcrypto" @@ -32,7 +31,9 @@ const ( // recordTimeout bounds the history write recordTimeout = 30 * time.Second // shutdownTimeout is the HTTP server drain budget at Stop. - shutdownTimeout = 15 * time.Second + shutdownTimeout = 15 * time.Second + preflightTimeout = 10 * time.Second + internalWriteTimeout = 30 * time.Second ) // Options are the runtime parameters of the site deploy executor, provided by @@ -91,28 +92,36 @@ func Serve(ctx context.Context, logger *slog.Logger, opts Options) error { if err != nil { return fmt.Errorf("building cloudflare client: %w", err) } + if !opts.Dev && (!cf.Enabled() || cfg.Cloudflare.R2.AccessKeyID == "" || cfg.Cloudflare.R2.SecretAccessKey == "") { + return fmt.Errorf("cloudflare R2 requires CLOUDFLARE_ACCOUNT_ID and CLOUDFLARE_R2_* credentials") + } + if !cf.Enabled() { + logger.Warn("cloudflare is not configured (CLOUDFLARE_ACCOUNT_ID, CLOUDFLARE_R2_*): every deploy is refused") + } wrk := worker.New(opts.WorkerURL, opts.Issuer, opts.Audience, priv) var rdb *cache.Cache if opts.RedisAddr != "" { rdb = cache.New(opts.RedisAddr) + } else { + logger.Warn("no redis address (SITESD_REDIS): the knot feed cursor is not persisted, so a restart resumes live and the pushes it missed never deploy") } + deployCtx, stopDeploys := context.WithCancel(context.Background()) + defer stopDeploys() + locks := reposync.New() - feed := sitefeed.New(rdb, cfg, cf, wrk, locks, logger) - for _, knot := range strings.Split(opts.Knots, ",") { - knot = strings.TrimSpace(knot) - if knot == "" { - continue + feed := sitefeed.New(rdb, cfg, cf, wrk, locks, deployCtx, logger) + if cf.Enabled() { + for _, knot := range strings.Split(opts.Knots, ",") { + knot = strings.TrimSpace(knot) + if knot != "" { + feed.Subscribe(runCtx, knot) + } } - feed.Subscribe(runCtx, knot) } - // deploys ride this, not a request: Stop cancels it after deployShutdownGrace - deployCtx, stopDeploys := context.WithCancel(context.Background()) - defer stopDeploys() - mux := http.NewServeMux() srv := &server{ cfg: cfg, @@ -154,7 +163,7 @@ func Serve(ctx context.Context, logger *slog.Logger, opts Options) error { Addr: opts.InternalListen, Handler: internalMux, ReadHeaderTimeout: 10 * time.Second, - WriteTimeout: 10 * time.Minute, + WriteTimeout: internalWriteTimeout, IdleTimeout: 2 * time.Minute, } go func() { @@ -163,6 +172,8 @@ func Serve(ctx context.Context, logger *slog.Logger, opts Options) error { serveErr <- fmt.Errorf("internal server: %w", err) } }() + } else { + logger.Warn("no internal listen address (SITESD_INTERNAL_ADDR): config-change and manual deploys have no endpoint") } var serveErrVal error @@ -187,12 +198,11 @@ func Serve(ctx context.Context, logger *slog.Logger, opts Options) error { return serveErrVal } -// drainDeploys lets in-flight deploys record, then cancels the rest. func (s *server) drainDeploys(stop context.CancelFunc) { done := make(chan struct{}) go func() { defer close(done) - s.deploys.Wait() + s.locks.Drain() }() select { case <-done: @@ -221,9 +231,7 @@ type server struct { locks *reposync.Registry logger *slog.Logger - // base of every deploy; never a request context deployCtx context.Context - deploys sync.WaitGroup getCtx func(context.Context, string) (*worker.DeployContext, error) deploy func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error @@ -236,6 +244,14 @@ type deployRequest struct { Trigger string `json:"trigger"` } +type acceptedDeploy struct { + dc *worker.DeployContext + trigger models.SiteDeployTrigger + release func() + done func() +} + +// ack 202 before deploy runs because worker wait_until expires in seconds while deploys take minutes func (s *server) handleInternalDeploy(w http.ResponseWriter, r *http.Request) { var req deployRequest if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<20)).Decode(&req); err != nil { @@ -254,53 +270,59 @@ func (s *server) handleInternalDeploy(w http.ResponseWriter, r *http.Request) { return } } + if !s.cf.Enabled() { + writeErr(w, http.StatusServiceUnavailable, "DeployUnavailable", "cloudflare is not configured") + return + } - // r.Context() dies ~30s after the worker's response (wait_until); run on deployCtx - status, deployErr := s.runDeploy(req.RepoDid, trigger) - if deployErr != nil { - s.logger.Error("internal deploy: failed", "repo", req.RepoDid, "err", deployErr) - writeErr(w, http.StatusBadGateway, "DeployFailed", deployErr.Error()) + preflightCtx, preflightCancel := context.WithTimeout(r.Context(), preflightTimeout) + defer preflightCancel() + dc, err := s.getCtx(preflightCtx, req.RepoDid) + if err != nil { + s.logger.Error("internal deploy: deploy context unavailable", "repo", req.RepoDid, "err", err) + writeErr(w, http.StatusBadGateway, "DeployContextUnavailable", err.Error()) + return + } + if dc == nil { + writeErr(w, http.StatusUnprocessableEntity, "NoDeployContext", "no site config or active claim for this repo") return } - w.Header().Set("Content-Type", "application/json") - w.WriteHeader(http.StatusOK) - json.NewEncoder(w).Encode(map[string]string{"status": status}) -} -// runDeploy runs one R2 sync on a deploy slot; the history write is the outcome record. -func (s *server) runDeploy(repoDid string, trigger models.SiteDeployTrigger) (string, error) { - if !s.cf.Enabled() { - return "", fmt.Errorf("cloudflare is not configured") + // register before ack so draining waits for accepted work + done, ok := s.locks.Begin() + if !ok { + writeErr(w, http.StatusServiceUnavailable, "ShuttingDown", "deploy executor is shutting down") + return + } + release, ok := s.locks.TryAcquire() + if !ok { + done() + writeErr(w, http.StatusTooManyRequests, "DeployAtCapacity", "every deploy slot is busy") + return } - s.deploys.Add(1) - defer s.deploys.Done() + go s.runDeploy(acceptedDeploy{dc: dc, trigger: trigger, release: release, done: done}) - runCtx, cancel := context.WithTimeout(s.deployCtx, deployTimeout) - defer cancel() + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusAccepted) + json.NewEncoder(w).Encode(map[string]string{"status": "accepted"}) +} - release, err := s.locks.Acquire(runCtx) - if err != nil { - return "", fmt.Errorf("waiting for a deploy slot: %w", err) - } - defer release() +func (s *server) runDeploy(a acceptedDeploy) { + defer a.done() + defer a.release() - dc, err := s.getCtx(runCtx, repoDid) - if err != nil { - return "", fmt.Errorf("fetching context: %w", err) - } - if dc == nil { - return "", fmt.Errorf("no site config or active claim for this repo") - } + runCtx, cancel := context.WithTimeout(s.deployCtx, deployTimeout) + defer cancel() + dc := a.dc deploy := &models.SiteDeploy{ RepoDid: syntax.DID(dc.RepoDid), Branch: dc.Branch, Dir: dc.Dir, CommitSHA: "", - Trigger: trigger, + Trigger: a.trigger, } - repo := &models.Repo{ Did: dc.OwnerDid, Name: dc.Name, @@ -309,30 +331,28 @@ func (s *server) runDeploy(repoDid string, trigger models.SiteDeployTrigger) (st RepoDid: dc.RepoDid, } - unlock, err := s.locks.Lock(runCtx, dc.RepoDid) - if err != nil { - return "", fmt.Errorf("waiting for the repo lock: %w", err) + // push feed shares this lock, serializing against trigger deploys for the same repo + unlock, lockErr := s.locks.Lock(runCtx, dc.RepoDid) + var deployErr error + if lockErr != nil { + deployErr = fmt.Errorf("waiting for the repo lock: %w", lockErr) + } else { + defer unlock() + deployErr = s.deploy(runCtx, s.cf, s.cfg, repo, dc.Branch, dc.Dir) } - defer unlock() - - deployErr := s.deploy(runCtx, s.cf, s.cfg, repo, dc.Branch, dc.Dir) + deploy.Status = models.SiteDeployStatusSuccess if deployErr != nil { deploy.Status = models.SiteDeployStatusFailure deploy.Error = deployErr.Error() - } else { - deploy.Status = models.SiteDeployStatusSuccess + s.logger.Error("internal deploy: failed", "repo", dc.RepoDid, "err", deployErr) } // outside the deploy's budget: timeouts and shutdowns still record recordCtx, recordCancel := context.WithTimeout(context.Background(), recordTimeout) defer recordCancel() if err := s.record(recordCtx, dc, deploy.CommitSHA, string(deploy.Status), string(deploy.Trigger), deploy.Error); err != nil { - return "", fmt.Errorf("deploy ran but recording history failed: %w", err) - } - if deployErr != nil { - return string(deploy.Status), deployErr + s.logger.Error("internal deploy: recording the outcome failed", "repo", dc.RepoDid, "status", deploy.Status, "err", err) } - return string(deploy.Status), nil } func (s *server) handleDidDoc(issuer string, priv *atcrypto.PrivateKeyP256) http.HandlerFunc { diff --git a/sitesd/sitesd_test.go b/sitesd/sitesd_test.go index b170da70e..2365ecd2d 100644 --- a/sitesd/sitesd_test.go +++ b/sitesd/sitesd_test.go @@ -10,6 +10,7 @@ import ( "net/http/httptest" "strings" "testing" + "time" "tangled.org/core/appview/models" "tangled.org/core/sitesd/cloudflare" @@ -22,10 +23,34 @@ const ( testOwnerDid = "did:plc:mockowner" testRepoDid = "did:plc:repo1" testKnot = "knot.invalid" + waitFor = 5 * time.Second ) var errBoom = errors.New("boom") +func TestServeRequiresR2OutsideDev(t *testing.T) { + for _, tc := range []struct { + name, account, bucket, access, secret string + }{ + {name: "account", bucket: "bucket", access: "access", secret: "secret"}, + {name: "bucket", account: "account", access: "access", secret: "secret"}, + {name: "access", account: "account", bucket: "bucket", secret: "secret"}, + {name: "secret", account: "account", bucket: "bucket", access: "access"}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Setenv("CLOUDFLARE_ACCOUNT_ID", tc.account) + t.Setenv("CLOUDFLARE_R2_BUCKET", tc.bucket) + t.Setenv("CLOUDFLARE_R2_ACCESS_KEY_ID", tc.access) + t.Setenv("CLOUDFLARE_R2_SECRET_ACCESS_KEY", tc.secret) + opts := Options{Knots: testKnot, KeyHex: strings.Repeat("0", 63) + "1"} + err := Serve(context.Background(), slog.New(slog.NewTextHandler(io.Discard, nil)), opts) + if err == nil || !strings.Contains(err.Error(), "cloudflare R2 requires") { + t.Fatalf("Serve with missing %s: %v", tc.name, err) + } + }) + } +} + type serverHooks struct { getCtx func(context.Context, string) (*worker.DeployContext, error) deploy func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error @@ -84,6 +109,68 @@ func decodeStatus(t *testing.T, w *httptest.ResponseRecorder) string { return out.Status } +func await(t *testing.T, ch <-chan struct{}, what string) { + t.Helper() + select { + case <-ch: + case <-time.After(waitFor): + t.Fatalf("timed out waiting for %s", what) + } +} + +func awaitQuiet(t *testing.T, ch <-chan struct{}, d time.Duration, why string) { + t.Helper() + select { + case <-ch: + t.Fatal(why) + case <-time.After(d): + } +} + +func requireAccepted(t *testing.T, w *httptest.ResponseRecorder) { + t.Helper() + if w.Code != http.StatusAccepted { + t.Fatalf("status = %d, want 202 (body %q)", w.Code, w.Body.String()) + } + if got := decodeStatus(t, w); got != "accepted" { + t.Errorf("status body = %q, want accepted", got) + } +} + +func requireRefused(t *testing.T, srv *server, w *httptest.ResponseRecorder, deployed, recorded <-chan struct{}, want int) { + t.Helper() + if w.Code != want { + t.Fatalf("status = %d, want %d (body %q)", w.Code, want, w.Body.String()) + } + for what, ch := range map[string]<-chan struct{}{"deploy": deployed, "history write": recorded} { + select { + case <-ch: + t.Errorf("a refused request still did the %s", what) + default: + } + } + release, ok := srv.locks.TryAcquire() + if !ok { + t.Error("a refused request kept a deploy slot") + return + } + release() +} + +func fillSlots(t *testing.T, r *reposync.Registry) []func() { + t.Helper() + var held []func() + for len(held) < 64 { + release, ok := r.TryAcquire() + if !ok { + return held + } + held = append(held, release) + } + t.Fatal("the deploy pool never filled") + return nil +} + func TestInternalDeploy_RequiresRepoDid(t *testing.T) { srv := newTestServer(serverHooks{}) w := httptest.NewRecorder() @@ -93,12 +180,17 @@ func TestInternalDeploy_RequiresRepoDid(t *testing.T) { } } -func TestInternalDeploy_Success(t *testing.T) { +func TestInternalDeploy_AcceptsBeforeTheDeployRuns(t *testing.T) { var repo models.Repo var recordCall *worker.DeployContext var recordArgs [4]string + var ctxCalls, recordCalls int + started := make(chan struct{}) + release := make(chan struct{}) + recorded := make(chan struct{}) srv := newTestServer(serverHooks{ getCtx: func(_ context.Context, repoDid string) (*worker.DeployContext, error) { + ctxCalls++ if repoDid != testRepoDid { t.Errorf("getCtx repoDid = %q, want %q", repoDid, testRepoDid) } @@ -109,21 +201,37 @@ func TestInternalDeploy_Success(t *testing.T) { if branch != "main" || dir != "/" { t.Errorf("deploy branch/dir = %q/%q, want main//", branch, dir) } + close(started) + <-release return nil }, record: func(_ context.Context, dc *worker.DeployContext, sha, status, trigger, errMsg string) error { + recordCalls++ recordCall = dc recordArgs = [4]string{sha, status, trigger, errMsg} + close(recorded) return nil }, }) + w := httptest.NewRecorder() srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) - if w.Code != http.StatusOK { - t.Fatalf("status = %d, want 200 (body %q)", w.Code, w.Body.String()) + requireAccepted(t, w) + select { + case <-recorded: + t.Fatal("the ack waited for the deploy to finish") + default: + } + + await(t, started, "the accepted deploy to run") + close(release) + await(t, recorded, "the history write") + + if ctxCalls != 1 { + t.Errorf("getDeployContext calls = %d, want the ack's one preflight", ctxCalls) } - if got := decodeStatus(t, w); got != "success" { - t.Errorf("status body = %q, want success", got) + if recordCalls != 1 { + t.Errorf("recordDeploy calls = %d, want 1", recordCalls) } if repo.Did != testOwnerDid || repo.Rkey != "abc1" || repo.Knot != testKnot || repo.RepoDid != testRepoDid { t.Errorf("deploy repo = %+v, want the owner's rkey/knot/repoDid", repo) @@ -150,34 +258,30 @@ func TestInternalDeploy_Trigger(t *testing.T) { for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { var got string - deployed := false + deployed := make(chan struct{}) + recorded := make(chan struct{}) srv := newTestServer(serverHooks{ getCtx: func(context.Context, string) (*worker.DeployContext, error) { return ctxFor(), nil }, deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { - deployed = true + close(deployed) return nil }, record: func(_ context.Context, _ *worker.DeployContext, _, _, trigger, _ string) error { got = trigger + close(recorded) return nil }, }) w := httptest.NewRecorder() srv.handleInternalDeploy(w, internalDeployRequest(tc.body)) if tc.want == "" { - if w.Code != http.StatusBadRequest { - t.Fatalf("status = %d, want 400 (body %q)", w.Code, w.Body.String()) - } - if deployed { - t.Error("a refused trigger must not start a deploy") - } + requireRefused(t, srv, w, deployed, recorded, http.StatusBadRequest) return } - if w.Code != http.StatusOK { - t.Fatalf("status = %d, want 200 (body %q)", w.Code, w.Body.String()) - } + requireAccepted(t, w) + await(t, recorded, "the history write") if got != tc.want { t.Errorf("recorded trigger = %q, want %q", got, tc.want) } @@ -185,8 +289,9 @@ func TestInternalDeploy_Trigger(t *testing.T) { } } -func TestInternalDeploy_DeployFailure(t *testing.T) { +func TestInternalDeploy_FailureIsRecordedNotReturned(t *testing.T) { var recordArgs [4]string + recorded := make(chan struct{}) srv := newTestServer(serverHooks{ getCtx: func(_ context.Context, _ string) (*worker.DeployContext, error) { return ctxFor(), nil @@ -196,15 +301,362 @@ func TestInternalDeploy_DeployFailure(t *testing.T) { }, record: func(_ context.Context, _ *worker.DeployContext, sha, status, trigger, errMsg string) error { recordArgs = [4]string{sha, status, trigger, errMsg} + close(recorded) return nil }, }) w := httptest.NewRecorder() srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) - if w.Code != http.StatusBadGateway { - t.Fatalf("status = %d, want 502 (body %q)", w.Code, w.Body.String()) - } + requireAccepted(t, w) + await(t, recorded, "the history write") if recordArgs[1] != "failure" || recordArgs[3] != "boom" { t.Errorf("record args = %v, want failure/boom", recordArgs) } } + +func TestInternalDeploy_AcksWhileTheDeployIsBlocked(t *testing.T) { + started := make(chan struct{}) + release := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + return ctxFor(), nil + }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(started) + <-release + return nil + }, + record: func(context.Context, *worker.DeployContext, string, string, string, string) error { + close(recorded) + return nil + }, + }) + + w := httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + requireAccepted(t, w) + select { + case <-recorded: + t.Fatal("the ack carried the deploy's outcome") + default: + } + + await(t, started, "the accepted deploy to run") + close(release) + await(t, recorded, "the history write") +} + +// negative control: synchronous handler blocks until deploy returns +func TestInternalDeploy_SynchronousShapeStillBlocks(t *testing.T) { + release := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + return ctxFor(), nil + }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + <-release + return nil + }, + record: func(context.Context, *worker.DeployContext, string, string, string, string) error { + close(recorded) + return nil + }, + }) + + w := httptest.NewRecorder() + served := make(chan struct{}) + go func() { + defer close(served) + synchronousDeployHandler(srv)(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + }() + + awaitQuiet(t, served, 100*time.Millisecond, "the synchronous handler answered while its deploy was blocked") + close(release) + await(t, served, "the synchronous handler's answer") + if w.Code != http.StatusOK { + t.Fatalf("status = %d, want the pre-fix 200 (body %q)", w.Code, w.Body.String()) + } +} + +func synchronousDeployHandler(s *server) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + var req deployRequest + if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<20)).Decode(&req); err != nil || req.RepoDid == "" { + writeErr(w, http.StatusBadRequest, "InvalidRequest", "request body must be JSON with a repoDid") + return + } + trigger := models.SiteDeployTriggerConfigChange + if req.Trigger != "" { + var ok bool + if trigger, ok = models.ParseSiteDeployTrigger(req.Trigger); !ok { + writeErr(w, http.StatusBadRequest, "InvalidRequest", "unknown trigger") + return + } + } + done, ok := s.locks.Begin() + if !ok { + writeErr(w, http.StatusServiceUnavailable, "ShuttingDown", "deploy executor is shutting down") + return + } + release, err := s.locks.Acquire(r.Context()) + if err != nil { + done() + writeErr(w, http.StatusServiceUnavailable, "DeployUnavailable", err.Error()) + return + } + dc, err := s.getCtx(r.Context(), req.RepoDid) + if err != nil || dc == nil { + done() + release() + writeErr(w, http.StatusBadGateway, "DeployFailed", "no deploy context") + return + } + s.runDeploy(acceptedDeploy{dc: dc, trigger: trigger, release: release, done: done}) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + json.NewEncoder(w).Encode(map[string]string{"status": "success"}) + } +} + +func TestInternalDeploy_DrainWaitsForAcceptedDeploy(t *testing.T) { + started := make(chan struct{}) + release := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + return ctxFor(), nil + }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(started) + <-release + return nil + }, + record: func(context.Context, *worker.DeployContext, string, string, string, string) error { + close(recorded) + return nil + }, + }) + + w := httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + requireAccepted(t, w) + await(t, started, "the accepted deploy to run") + + drained := make(chan struct{}) + go func() { + defer close(drained) + srv.locks.Drain() + }() + awaitQuiet(t, drained, 50*time.Millisecond, "Drain returned while an accepted deploy was still running") + + close(release) + await(t, drained, "Drain after the accepted deploy") + select { + case <-recorded: + default: + t.Error("Drain returned before the accepted deploy recorded its outcome") + } +} + +func TestInternalDeploy_DrainRefusesNewWork(t *testing.T) { + deployed := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + return ctxFor(), nil + }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(deployed) + return nil + }, + record: func(context.Context, *worker.DeployContext, string, string, string, string) error { + close(recorded) + return nil + }, + }) + srv.locks.Drain() + + w := httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + requireRefused(t, srv, w, deployed, recorded, http.StatusServiceUnavailable) +} + +func TestInternalDeploy_DrainDuringPreflightRefuses(t *testing.T) { + preflight := make(chan struct{}) + finishPreflight := make(chan struct{}) + deployed := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + close(preflight) + <-finishPreflight + return ctxFor(), nil + }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(deployed) + return nil + }, + record: func(context.Context, *worker.DeployContext, string, string, string, string) error { + close(recorded) + return nil + }, + }) + + w := httptest.NewRecorder() + served := make(chan struct{}) + go func() { + defer close(served) + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + }() + await(t, preflight, "the preflight") + + srv.locks.Drain() + close(finishPreflight) + await(t, served, "the refused answer") + + requireRefused(t, srv, w, deployed, recorded, http.StatusServiceUnavailable) +} + +func TestInternalDeploy_NoDeployContext(t *testing.T) { + cases := []struct { + name string + ctx func() (*worker.DeployContext, error) + want int + }{ + {"no site config or active claim", func() (*worker.DeployContext, error) { return nil, nil }, http.StatusUnprocessableEntity}, + {"the context lookup failed", func() (*worker.DeployContext, error) { return nil, errBoom }, http.StatusBadGateway}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + deployed := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { return tc.ctx() }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(deployed) + return nil + }, + record: func(context.Context, *worker.DeployContext, string, string, string, string) error { + close(recorded) + return nil + }, + }) + w := httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + requireRefused(t, srv, w, deployed, recorded, tc.want) + }) + } +} + +func TestInternalDeploy_RefusesWhenSlotsFull(t *testing.T) { + deployed := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + return ctxFor(), nil + }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(deployed) + return nil + }, + record: func(context.Context, *worker.DeployContext, string, string, string, string) error { + close(recorded) + return nil + }, + }) + held := fillSlots(t, srv.locks) + + w := httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + if w.Code != http.StatusTooManyRequests { + t.Fatalf("status = %d, want 429 (body %q)", w.Code, w.Body.String()) + } + select { + case <-deployed: + t.Error("a deploy ran without a slot") + default: + } + select { + case <-recorded: + t.Error("a deploy without a slot still wrote history") + default: + } + + held[0]() + w = httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + requireAccepted(t, w) + await(t, recorded, "the history write") +} + +func TestInternalDeploy_CloudflareNotConfigured(t *testing.T) { + var ctxCalls int + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + ctxCalls++ + return ctxFor(), nil + }, + }) + cf, err := cloudflare.New(&config.Config{}) + if err != nil { + t.Fatal(err) + } + if cf.Enabled() { + t.Fatal("an empty cloudflare config must not read as enabled") + } + srv.cf = cf + + w := httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + if w.Code != http.StatusServiceUnavailable { + t.Fatalf("status = %d, want 503 (body %q)", w.Code, w.Body.String()) + } + if ctxCalls != 0 { + t.Errorf("getDeployContext calls = %d, want none before the cloudflare check", ctxCalls) + } +} + +func TestInternalDeploy_RecordsAnAbandonedRepoLock(t *testing.T) { + var recordArgs [4]string + deployed := make(chan struct{}) + recorded := make(chan struct{}) + srv := newTestServer(serverHooks{ + getCtx: func(context.Context, string) (*worker.DeployContext, error) { + return ctxFor(), nil + }, + deploy: func(context.Context, *cloudflare.Client, *config.Config, *models.Repo, string, string) error { + close(deployed) + return nil + }, + record: func(_ context.Context, _ *worker.DeployContext, sha, status, trigger, errMsg string) error { + recordArgs = [4]string{sha, status, trigger, errMsg} + close(recorded) + return nil + }, + }) + + unlock, err := srv.locks.Lock(context.Background(), testRepoDid) + if err != nil { + t.Fatal(err) + } + defer unlock() + deployCtx, cancel := context.WithCancel(context.Background()) + cancel() + srv.deployCtx = deployCtx + + w := httptest.NewRecorder() + srv.handleInternalDeploy(w, internalDeployRequest(`{"repoDid":"did:plc:repo1"}`)) + requireAccepted(t, w) + await(t, recorded, "the abandoned deploy's history write") + + if recordArgs[1] != "failure" || !strings.Contains(recordArgs[3], "repo lock") { + t.Errorf("record args = %v, want a failure naming the abandoned repo lock", recordArgs) + } + select { + case <-deployed: + t.Error("the deploy ran without the repo lock") + default: + } +}