Something went wrong. Try again.
Monorepo for Tangled
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922package 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) }}