package worker import ( "context" "crypto/rand" "encoding/base64" "errors" "fmt" "os" "path/filepath" "strings" "testing" "time" "github.com/bluesky-social/indigo/atproto/atcrypto" "tangled.org/core/log" "tangled.org/core/migrator/config" "tangled.org/core/migrator/crypto" "tangled.org/core/migrator/db" "tangled.org/core/migrator/git" migratoroauth "tangled.org/core/migrator/oauth" ) type fakeRunner struct { runs []string cloneErr error cloneOut string clonePackBytes int64 lsRemoteErr error lsRemoteOut string lfsFetchErr error hasGitLFS bool } func (f *fakeRunner) LookPath(file string) (string, error) { if file == "git-lfs" && f.hasGitLFS { return "/usr/bin/git-lfs", nil } return "", errors.New("not found") } func (f *fakeRunner) Run(_ context.Context, _ string, _ []string, name string, args ...string) (string, error) { f.runs = append(f.runs, name+" "+strings.Join(args, " ")) if len(args) > 0 && args[0] == "clone" { if f.cloneErr == nil && f.clonePackBytes > 0 { packDir := filepath.Join(args[len(args)-1], "objects", "pack") if err := os.MkdirAll(packDir, 0700); err != nil { return "", err } pack := filepath.Join(packDir, "pack-test.pack") if err := os.WriteFile(pack, nil, 0600); err != nil { return "", err } if err := os.Truncate(pack, f.clonePackBytes); err != nil { return "", err } } return f.cloneOut, f.cloneErr } if len(args) > 0 && args[0] == "ls-remote" { return f.lsRemoteOut, f.lsRemoteErr } if len(args) > 1 && args[0] == "lfs" && args[1] == "fetch" { return "", f.lfsFetchErr } return "", nil } type fakeKnot struct { createURL string createErr error repoDid string contents []string describeErrs []error describes int onCreate func(string) records []string recordErr error } func (f *fakeKnot) CreateRepo(_ context.Context, _, _, _, _, sourceURL string) (string, error) { f.createURL = sourceURL if f.onCreate != nil { f.onCreate(sourceURL) } if f.repoDid == "" { f.repoDid = "did:plc:created" } return f.repoDid, f.createErr } func (f *fakeKnot) PutRepoRecord(_ context.Context, ownerDid, rkey, name, description, knotDid, repoDid string) error { if f.recordErr != nil { return f.recordErr } // tracks describe count to verify content was held before record creation f.records = append(f.records, fmt.Sprintf("%s/%s %s desc=%q %s %s describes=%d", ownerDid, rkey, name, description, knotDid, repoDid, f.describes)) return nil } func (f *fakeKnot) Content(_ context.Context, _, _ string) (string, error) { content, _, err := f.DescribeRepo(context.Background(), "", "", "") return content, err } func (f *fakeKnot) DescribeRepo(_ context.Context, _, _, _ string) (string, *string, error) { if len(f.describeErrs) > 0 { err := f.describeErrs[0] f.describeErrs = f.describeErrs[1:] f.describes++ return "", nil, err } if len(f.contents) == 0 { return "present", nil, nil } index := min(f.describes, len(f.contents)-1) f.describes++ return f.contents[index], nil, nil } func setupTestWorker(t *testing.T, runner git.CommandRunner, knot KnotClient) (*WorkerPool, *db.DB, *config.Config) { t.Helper() dir := t.TempDir() database, err := db.Make(context.Background(), filepath.Join(dir, "worker.db")) if err != nil { t.Fatal(err) } t.Cleanup(func() { database.Close() }) rawKey := make([]byte, 32) _, _ = rand.Read(rawKey) privateKey, err := atcrypto.GeneratePrivateKeyP256() if err != nil { t.Fatal(err) } cfg := &config.Config{ Hostname: "migrator.example.com", MasterKey: base64.StdEncoding.EncodeToString(rawKey), PrivateKey: privateKey.Multibase(), WorkDir: filepath.Join(dir, "scratch-work"), JobTimeout: time.Second, } if err := cfg.Validate(); err != nil { t.Fatal(err) } if knot == nil { knot = &fakeKnot{} } pool, err := NewWorkerPool(database, cfg, nil, runner, knot, log.New("test-worker")) if err != nil { t.Fatal(err) } pool.pollInterval = time.Millisecond return pool, database, cfg } func batchInput(sealed *string, private bool) db.CreateBatchInput { return db.CreateBatchInput{ ID: "batch", OwnerDid: "did:plc:alice", RequestID: "request", RequestDigest: "digest", EncryptedToken: sealed, Jobs: []db.CreateJobInput{{ Name: "repo", KnotDid: "did:web:knot.example", SourceURL: "https://github.com/alice/repo.git", Description: "a small tool", Private: private, }}, } } func createAndClaim(t *testing.T, database *db.DB) (*db.Job, *db.Batch) { t.Helper() return claimOne(t, database, batchInput(nil, false)) } func createAndClaimPrivate(t *testing.T, database *db.DB, cfg *config.Config) (*db.Job, *db.Batch) { t.Helper() sealed, err := crypto.Encrypt(cfg.ParsedMasterKey, "ghp_credential", crypto.ComputeAAD("did:plc:alice", "batch", "request")) if err != nil { t.Fatal(err) } return claimOne(t, database, batchInput(&sealed, true)) } func claimOne(t *testing.T, database *db.DB, in db.CreateBatchInput) (*db.Job, *db.Batch) { t.Helper() if _, _, _, err := database.CreateBatch(context.Background(), in); err != nil { t.Fatal(err) } job, batch, err := database.ClaimNextQueuedJob(context.Background()) if err != nil || job == nil { t.Fatalf("claim: job=%v err=%v", job, err) } return job, batch } func TestJobPublishesCapabilityWaitsAndCleansUp(t *testing.T) { runner := &fakeRunner{} knot := &fakeKnot{contents: []string{"pending", "present"}} pool, database, cfg := setupTestWorker(t, runner, knot) job, batch := createAndClaimPrivate(t, database, cfg) knot.onCreate = func(sourceURL string) { parts := strings.Split(sourceURL, "/") token := parts[len(parts)-2] resolved, err := database.ResolveJobCapability(context.Background(), token, "repo") if err != nil || resolved.ID != job.ID { t.Fatalf("capability did not map to the job: resolved=%v err=%v", resolved, err) } } pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusCompleted || reloaded.RepoDid != "did:plc:created" { t.Fatalf("unexpected completed job: %+v", reloaded) } if reloaded.CapabilityToken != nil { t.Fatal("capability remained active after completion") } if _, err := os.Stat(filepath.Join(cfg.WorkDir, "job-1")); !os.IsNotExist(err) { t.Fatal("scratch mirror remained after completion") } if len(runner.runs) != 1 || !strings.Contains(runner.runs[0], "clone --mirror") { t.Fatalf("expected only a mirror clone, got %v", runner.runs) } if !strings.HasPrefix(knot.createURL, "https://migrator.example.com/git/") || !strings.HasSuffix(knot.createURL, "/repo.git") { t.Fatalf("unexpected capability url %q", knot.createURL) } if len(knot.records) != 1 { t.Fatalf("repo record writes = %v, want exactly one", knot.records) } if !strings.Contains(knot.records[0], "did:plc:alice/repo ") { t.Fatalf("record written for the wrong owner or rkey: %s", knot.records[0]) } if !strings.Contains(knot.records[0], `desc="a small tool"`) { t.Fatalf("the record lost the description: %s", knot.records[0]) } if strings.Contains(knot.records[0], "describes=0") { t.Fatalf("the record was written before the knot held the content: %s", knot.records[0]) } if knot.describes != 2 { t.Fatalf("describeRepo calls = %d, want 2", knot.describes) } } func TestAFailedMigrationWritesNoRepoRecord(t *testing.T) { knot := &fakeKnot{contents: []string{"failed"}} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) job, batch := createAndClaim(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusQueued { t.Fatalf("job status = %s, want it requeued for another attempt", reloaded.Status) } if len(knot.records) != 0 { t.Fatalf("an attempt that never landed wrote a record: %v", knot.records) } } func TestARetryKeepsTheCapabilityAndMirrorAlive(t *testing.T) { runner := &fakeRunner{} knot := &fakeKnot{createErr: errors.New("temporary knot failure")} calls := 0 knot.onCreate = func(string) { calls++ if calls > 1 { knot.createErr = nil // the knot recovered for the retry } } pool, database, cfg := setupTestWorker(t, runner, knot) job, batch := createAndClaimPrivate(t, database, cfg) pool.executeJob(context.Background(), job, batch) firstURL := knot.createURL reloaded, _, _ := database.GetJob(context.Background(), job.ID) if reloaded.Status != db.StatusQueued { t.Fatalf("status = %s, want queued retry", reloaded.Status) } if reloaded.CapabilityToken == nil { t.Fatal("the capability died with the attempt; the knot could not finish fetching") } if _, err := database.ResolveJobCapability(context.Background(), *reloaded.CapabilityToken, "repo"); err != nil { t.Fatalf("capability does not resolve during retry backoff: %v", err) } if _, err := os.Stat(filepath.Join(cfg.WorkDir, "job-1", "mirror-ready")); err != nil { t.Fatal("the completed mirror did not survive the retry backoff; the retry would pay for a second clone") } pool.executeJob(context.Background(), reloaded, batch) reloaded, _, _ = database.GetJob(context.Background(), job.ID) if reloaded.Status != db.StatusCompleted || reloaded.RepoDid == "" { t.Fatalf("retry did not complete: %+v", reloaded) } if reloaded.CapabilityToken != nil { t.Fatal("capability remained active after completion") } if len(runner.runs) != 1 { t.Fatalf("the retry cloned the source again: %v", runner.runs) } if knot.createURL == "" || knot.createURL != firstURL { t.Fatalf("the retry replaced capability url %q; the knot may still be fetching %q", knot.createURL, firstURL) } if _, err := os.Stat(filepath.Join(cfg.WorkDir, "job-1")); !os.IsNotExist(err) { t.Fatal("scratch mirror remained after completion") } } func existingRepoRetry(t *testing.T, database *db.DB, repoDid string) (*db.Job, *db.Batch) { t.Helper() job, _ := createAndClaim(t, database) if err := database.UpdateJobStatus(context.Background(), job.ID, db.StatusImporting, nil); err != nil { t.Fatal(err) } if err := database.SetJobRepoDid(context.Background(), job.ID, repoDid); err != nil { t.Fatal(err) } if err := database.ScheduleJobRetry(context.Background(), job.ID, 0, strPtr("retry")); err != nil { t.Fatal(err) } job, batch, err := database.ClaimNextQueuedJob(context.Background()) if err != nil || job == nil { t.Fatalf("claim retry: job=%v err=%v", job, err) } return job, batch } func TestPresentRepoRetryOnlyWritesTheMissingPDSRecord(t *testing.T) { runner := &fakeRunner{} knot := &fakeKnot{contents: []string{"present"}, repoDid: "did:plc:created"} creates := 0 knot.onCreate = func(string) { creates++ } pool, database, _ := setupTestWorker(t, runner, knot) job, batch := existingRepoRetry(t, database, "did:plc:created") pool.executeJob(context.Background(), job, batch) if creates != 0 || len(runner.runs) != 0 { t.Fatalf("present repo was re-created or cloned: creates=%d commands=%v", creates, runner.runs) } if len(knot.records) != 1 { t.Fatalf("PDS records = %v, want the missing record written", knot.records) } } func TestFetchingRepoRetryKeepsPollingWithoutRecreating(t *testing.T) { runner := &fakeRunner{} knot := &fakeKnot{contents: []string{"fetching", "present"}, repoDid: "did:plc:created"} creates := 0 knot.onCreate = func(string) { creates++ } pool, database, _ := setupTestWorker(t, runner, knot) job, batch := existingRepoRetry(t, database, "did:plc:created") pool.executeJob(context.Background(), job, batch) if creates != 0 || len(runner.runs) != 0 || len(knot.records) != 1 { t.Fatalf("fetching resume did extra work: creates=%d commands=%v records=%v", creates, runner.runs, knot.records) } } func TestFailedRepoRetryRestagesTheExistingIdentity(t *testing.T) { runner := &fakeRunner{} knot := &fakeKnot{contents: []string{"failed", "present"}, repoDid: "did:plc:created"} creates := 0 knot.onCreate = func(string) { creates++ } pool, database, _ := setupTestWorker(t, runner, knot) job, batch := existingRepoRetry(t, database, "did:plc:created") pool.executeJob(context.Background(), job, batch) if creates != 1 || len(runner.runs) != 1 || len(knot.records) != 1 { t.Fatalf("failed resume did not re-stage once: creates=%d commands=%v records=%v", creates, runner.runs, knot.records) } } func TestWaitUntilPresentStopsOnFailedContent(t *testing.T) { reason := "source rejected" knot := &reasonKnot{reason: &reason} pool, _, _ := setupTestWorker(t, &fakeRunner{}, knot) err := pool.waitUntilPresent(context.Background(), "did:plc:alice", "did:web:knot.example", "did:plc:repo") if err == nil || !strings.Contains(err.Error(), reason) { t.Fatalf("error = %v, want knot reason", err) } } type describeOnly struct{} func (describeOnly) CreateRepo(context.Context, string, string, string, string, string) (string, error) { return "", nil } func (describeOnly) PutRepoRecord(context.Context, string, string, string, string, string, string) error { return nil } type reasonKnot struct { describeOnly reason *string } func (r *reasonKnot) Content(context.Context, string, string) (string, error) { return "failed", nil } func (r *reasonKnot) DescribeRepo(context.Context, string, string, string) (string, *string, error) { return "failed", r.reason, nil } func TestWaitUntilPresentStopsOnPartialContent(t *testing.T) { reason := "one ref failed" pool, _, _ := setupTestWorker(t, &fakeRunner{}, &partialKnot{reason: &reason}) err := pool.waitUntilPresent(context.Background(), "did:plc:alice", "did:web:knot.example", "did:plc:repo") if err == nil || !strings.Contains(err.Error(), "partial") || !strings.Contains(err.Error(), reason) { t.Fatalf("error = %v, want partial reason", err) } } type partialKnot struct { describeOnly reason *string } func (r *partialKnot) Content(context.Context, string, string) (string, error) { return "partial", nil } func (r *partialKnot) DescribeRepo(context.Context, string, string, string) (string, *string, error) { return "partial", r.reason, nil } func TestCloneRejectsPackAboveKnotLimitBeforeLFS(t *testing.T) { runner := &fakeRunner{clonePackBytes: 4, hasGitLFS: true} knot := &fakeKnot{} pool, database, cfg := setupTestWorker(t, runner, knot) cfg.MaxPackBytes = 3 job, batch := createAndClaimPrivate(t, database, cfg) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusFailed || reloaded.Error == nil || !strings.Contains(*reloaded.Error, "git pack limit exceeded") { t.Fatalf("job = %+v, want terminal pack-limit failure", reloaded) } if len(runner.runs) != 1 || strings.Contains(strings.Join(runner.runs, " "), "lfs fetch") { t.Fatalf("oversize pack reached LFS fetch: %v", runner.runs) } if knot.createURL != "" { t.Fatalf("oversize pack reached the knot: %s", knot.createURL) } } func TestLFSFetchIsKeptButNoLFSPushRuns(t *testing.T) { runner := &fakeRunner{hasGitLFS: true} pool, database, cfg := setupTestWorker(t, runner, &fakeKnot{}) job, batch := createAndClaimPrivate(t, database, cfg) pool.executeJob(context.Background(), job, batch) if len(runner.runs) != 2 || !strings.Contains(runner.runs[1], "lfs fetch --all") { t.Fatalf("commands = %v", runner.runs) } for _, command := range runner.runs { if strings.Contains(command, "push") { t.Fatalf("obsolete push command ran: %s", command) } } } func TestAPublicRepoIsHandedToTheKnotUncloned(t *testing.T) { runner := &fakeRunner{} knot := &fakeKnot{contents: []string{"pending", "present"}} pool, database, cfg := setupTestWorker(t, runner, knot) job, batch := createAndClaim(t, database) knot.onCreate = func(string) { live, _, err := database.GetJob(context.Background(), job.ID) if err != nil || live.CapabilityToken != nil { t.Fatalf("a public import minted a capability: job=%+v err=%v", live, err) } } pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusCompleted || reloaded.RepoDid != "did:plc:created" { t.Fatalf("unexpected completed job: %+v", reloaded) } if knot.createURL != "https://github.com/alice/repo.git" { t.Fatalf("knot was sent %q, want the source origin", knot.createURL) } if len(runner.runs) != 1 || !strings.HasPrefix(runner.runs[0], "git ls-remote") { t.Fatalf("expected only a reachability probe, got %v", runner.runs) } if _, err := os.Stat(filepath.Join(cfg.WorkDir, "job-1")); !os.IsNotExist(err) { t.Fatal("a public import made a scratch directory") } if len(knot.records) != 1 { t.Fatalf("repo record writes = %v, want exactly one", knot.records) } } func TestAnUnreadablePublicSourceNeverReachesTheKnot(t *testing.T) { runner := &fakeRunner{lsRemoteErr: errors.New("exit status 128"), lsRemoteOut: "ERROR: Repository not found."} knot := &fakeKnot{} creates := 0 knot.onCreate = func(string) { creates++ } pool, database, _ := setupTestWorker(t, runner, knot) job, batch := createAndClaim(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if creates != 0 || knot.createURL != "" { t.Fatalf("an unreadable source reached the knot: creates=%d url=%q", creates, knot.createURL) } if reloaded.Error == nil || !strings.Contains(*reloaded.Error, "Repository not found") { t.Fatalf("job error = %v, want git's own answer", reloaded.Error) } } func TestStartupRecoveryRevokesInflightCapabilities(t *testing.T) { pool, database, cfg := setupTestWorker(t, &fakeRunner{}, &fakeKnot{}) job, _ := createAndClaim(t, database) if err := database.UpdateJobStatus(context.Background(), job.ID, db.StatusImporting, nil); err != nil { t.Fatal(err) } if err := database.SetJobCapability(context.Background(), job.ID, strings.Repeat("a", 43)); err != nil { t.Fatal(err) } jobDir := filepath.Join(cfg.WorkDir, "job-1") _ = os.MkdirAll(jobDir, 0700) if err := pool.startupRecovery(context.Background()); err != nil { t.Fatal(err) } reloaded, _, _ := database.GetJob(context.Background(), job.ID) if reloaded.Status != db.StatusQueued || reloaded.CapabilityToken != nil { t.Fatalf("recovered job = %+v", reloaded) } if _, err := os.Stat(jobDir); !os.IsNotExist(err) { t.Fatal("recovery left scratch directory") } } func TestTerminalFailureRevokesDeletesAndRedactsCapability(t *testing.T) { runner := &fakeRunner{} knot := &fakeKnot{} pool, database, cfg := setupTestWorker(t, runner, knot) job, _ := createAndClaimPrivate(t, database, cfg) if _, err := database.Exec("update jobs set attempts = ? where id = ?", MaxAttempts, job.ID); err != nil { t.Fatal(err) } job, batch, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } knot.onCreate = func(sourceURL string) { knot.createErr = errors.New("knot rejected " + sourceURL) } pool.executeJob(context.Background(), job, batch) reloaded, _, _ := database.GetJob(context.Background(), job.ID) if reloaded.Status != db.StatusFailed || reloaded.CapabilityToken != nil { t.Fatalf("terminal job = %+v", reloaded) } if reloaded.Error == nil || !strings.Contains(*reloaded.Error, "[REDACTED]") || strings.Contains(*reloaded.Error, knot.createURL) { t.Fatalf("capability leaked through job error: %v", reloaded.Error) } if _, err := os.Stat(filepath.Join(cfg.WorkDir, "job-1")); !os.IsNotExist(err) { t.Fatal("terminal failure left the mirror on disk") } } func TestWaitUntilPresentHonorsContext(t *testing.T) { knot := &fakeKnot{contents: []string{"pending"}} pool, _, _ := setupTestWorker(t, &fakeRunner{}, knot) ctx, cancel := context.WithTimeout(context.Background(), 5*time.Millisecond) defer cancel() err := pool.waitUntilPresent(ctx, "did:plc:alice", "did:web:knot.example", "did:plc:repo") if !errors.Is(err, context.DeadlineExceeded) { t.Fatalf("error = %v, want deadline exceeded", err) } } func TestWaitingSurvivesATransientKnotAnswer(t *testing.T) { knot := &fakeKnot{describeErrs: []error{errors.New("503 service unavailable")}} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) job, batch := createAndClaim(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusCompleted { t.Fatalf("status = %q, want %q (err: %v)", reloaded.Status, db.StatusCompleted, reloaded.Error) } } type conflictKnot struct { createErr error contents []string describeErr []error creates int } func (k *conflictKnot) CreateRepo(context.Context, string, string, string, string, string) (string, error) { k.creates++ return "", k.createErr } func (conflictKnot) PutRepoRecord(context.Context, string, string, string, string, string, string) error { return nil } func (k *conflictKnot) Content(_ context.Context, _, _ string) (string, error) { content, _, err := k.DescribeRepo(context.Background(), "", "", "") return content, err } func (k *conflictKnot) DescribeRepo(_ context.Context, _, _, _ string) (string, *string, error) { if len(k.describeErr) > 0 { err := k.describeErr[0] k.describeErr = k.describeErr[1:] return "", nil, err } if len(k.contents) == 0 { return "present", nil, nil } content := k.contents[0] k.contents = k.contents[1:] return content, nil, nil } func seedEarlierAttempt(t *testing.T, database *db.DB, repoDid string) { t.Helper() seedEarlierAttemptFrom(t, database, repoDid, "https://github.com/alice/repo.git") } func seedEarlierAttemptFrom(t *testing.T, database *db.DB, repoDid, sourceURL string) { t.Helper() in := batchInput(nil, false) in.Jobs[0].SourceURL = sourceURL in.ID, in.RequestID = "batch-earlier", "earlier" if _, _, _, err := database.CreateBatch(context.Background(), in); err != nil { t.Fatal(err) } job, _, err := database.ClaimNextQueuedJob(context.Background()) if err != nil || job == nil { t.Fatalf("claim earlier job: job=%v err=%v", job, err) } if err := database.SetJobRepoDid(context.Background(), job.ID, repoDid); err != nil { t.Fatal(err) } if err := database.UpdateJobStatus(context.Background(), job.ID, db.StatusFailed, strPtr("recording failed")); err != nil { t.Fatal(err) } } func conflictBatch(t *testing.T, database *db.DB) (*db.Job, *db.Batch) { t.Helper() in := batchInput(nil, false) in.ID, in.RequestID = "batch-conflict", "conflict" if _, _, _, err := database.CreateBatch(context.Background(), in); err != nil { t.Fatal(err) } job, batch, err := database.ClaimNextQueuedJob(context.Background()) if err != nil || job == nil { t.Fatalf("claim conflict job: job=%v err=%v", job, err) } return job, batch } func conflictErr() error { return fmt.Errorf("knot create returned HTTP 409: %w", migratoroauth.ErrRepoExists) } func TestAConflictWithAContentBearingRepoCompletesTheJob(t *testing.T) { knot := &conflictKnot{createErr: conflictErr(), contents: []string{"present"}} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) seedEarlierAttempt(t, database, "did:plc:candidate") job, batch := conflictBatch(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusCompleted || reloaded.RepoDid != "did:plc:candidate" { t.Fatalf("job = %+v, want completed on the existing repository", reloaded) } if reloaded.Attempts != 1 { t.Fatalf("attempts = %d, want the conflict settled on the first attempt", reloaded.Attempts) } if knot.creates != 0 { t.Fatalf("knot creates = %d, want existing repository resolved without knot create", knot.creates) } } func TestAConflictWithAnEmptyRepoIsADistinctFailure(t *testing.T) { knot := &conflictKnot{createErr: conflictErr(), contents: []string{""}} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) seedEarlierAttempt(t, database, "did:plc:candidate") job, batch := conflictBatch(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusFailed || reloaded.Error == nil || !strings.Contains(*reloaded.Error, "holds no import") { t.Fatalf("job = %+v, want the distinct empty-target failure", reloaded) } if reloaded.Attempts != 1 { t.Fatalf("attempts = %d, want no retry burn on a conflict", reloaded.Attempts) } } func TestAConflictFromADifferentSourceIsNotAdopted(t *testing.T) { knot := &conflictKnot{createErr: conflictErr(), contents: []string{"present"}} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) seedEarlierAttemptFrom(t, database, "did:plc:other", "https://github.com/someone-else/tool.git") job, batch := conflictBatch(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusFailed || reloaded.Error == nil || !strings.Contains(*reloaded.Error, "was not created from") { t.Fatalf("job = %+v, want a refusal instead of adopting a repository another import created", reloaded) } if reloaded.RepoDid != "" { t.Fatalf("repoDid = %q, want the job not to attach to another source's repository", reloaded.RepoDid) } } func TestAConflictWithNoEarlierAttemptIsDistinct(t *testing.T) { knot := &conflictKnot{createErr: conflictErr()} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) job, batch := conflictBatch(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusFailed || reloaded.Error == nil || !strings.Contains(*reloaded.Error, "was not created from") { t.Fatalf("job = %+v, want the distinct conflict failure", reloaded) } if reloaded.Attempts != 1 { t.Fatalf("attempts = %d, want no retry burn on a conflict", reloaded.Attempts) } } func TestAConflictSurvivesATransientContentCheck(t *testing.T) { knot := &conflictKnot{createErr: conflictErr(), describeErr: []error{errors.New("503 service unavailable")}} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) seedEarlierAttempt(t, database, "did:plc:candidate") job, batch := conflictBatch(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusQueued { t.Fatalf("status = %s, want the conflict requeued while the knot cannot be asked", reloaded.Status) } } func TestARefusedGrantIsTerminalNotRetried(t *testing.T) { knot := &fakeKnot{recordErr: fmt.Errorf("refreshing oauth session: %w", migratoroauth.ErrGrantRequired)} pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) job, batch := createAndClaim(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusAuthorizationRequired { t.Fatalf("status = %s, want %q (the grant can never recover on retry)", reloaded.Status, db.StatusAuthorizationRequired) } if len(knot.records) != 0 { t.Fatalf("a refused grant wrote a repo record: %v", knot.records) } } func TestAStoppingDaemonReleasesTheJob(t *testing.T) { pool, database, _ := setupTestWorker(t, &fakeRunner{}, &fakeKnot{}) job, batch := createAndClaim(t, database) if job.Attempts != 1 { t.Fatalf("the claim did not spend an attempt: %d", job.Attempts) } stopped, cancel := context.WithCancel(context.Background()) cancel() pool.handleJobFailure(stopped, job, batch, context.Canceled, "cloning", "", "") reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusQueued { t.Fatalf("status = %q, want %q", reloaded.Status, db.StatusQueued) } if reloaded.Attempts != 0 { t.Fatalf("attempts = %d, want the claim refunded", reloaded.Attempts) } } type strictKnot struct { repoExists bool createdDid string creates int contents []string records []string t *testing.T } func (s *strictKnot) CreateRepo(_ context.Context, _, _, _, _, _ string) (string, error) { s.creates++ if s.repoExists { s.t.Fatalf("CreateRepo called when repo already exists on knot") } return s.createdDid, nil } func (s *strictKnot) Content(_ context.Context, _, _ string) (string, error) { if len(s.contents) == 0 { return "present", nil } return s.contents[0], nil } func (s *strictKnot) DescribeRepo(ctx context.Context, _, _, _ string) (string, *string, error) { c, err := s.Content(ctx, "", "") return c, nil, err } func (s *strictKnot) PutRepoRecord(_ context.Context, ownerDid, rkey, name, description, knotDid, repoDid string) error { s.records = append(s.records, repoDid) return nil } func TestExistingRepoPathAvoidsSecondKnotCreate(t *testing.T) { // the repo already exists on the knot, so push into it instead of creating one // fake knot fails test if CreateRepo is called when repo already exists. knot := &strictKnot{ repoExists: true, contents: []string{"present"}, t: t, } pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) seedEarlierAttempt(t, database, "did:plc:pre-existing") job, batch := createAndClaim(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusCompleted || reloaded.RepoDid != "did:plc:pre-existing" { t.Fatalf("job = %+v, want completed on pre-existing repo did:plc:pre-existing", reloaded) } if knot.creates != 0 { t.Fatalf("knot creates = %d, want 0 (no second knot create)", knot.creates) } if len(knot.records) != 1 || knot.records[0] != "did:plc:pre-existing" { t.Fatalf("records = %v, want exactly 1 record naming did:plc:pre-existing", knot.records) } } func TestMissingRepoFallbackCallsCreateRepo(t *testing.T) { // the repo is missing on the knot, so the fallback create runs knot := &strictKnot{ repoExists: false, createdDid: "did:plc:fallback-created", contents: []string{"present"}, t: t, } pool, database, _ := setupTestWorker(t, &fakeRunner{}, knot) job, batch := createAndClaim(t, database) pool.executeJob(context.Background(), job, batch) reloaded, _, err := database.GetJob(context.Background(), job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != db.StatusCompleted || reloaded.RepoDid != "did:plc:fallback-created" { t.Fatalf("job = %+v, want completed on fallback-created repo did:plc:fallback-created", reloaded) } if knot.creates != 1 { t.Fatalf("knot creates = %d, want 1 (fallback create called)", knot.creates) } if len(knot.records) != 1 || knot.records[0] != "did:plc:fallback-created" { t.Fatalf("records = %v, want exactly 1 record naming did:plc:fallback-created", knot.records) } }