Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596package db
import ( "context" "path/filepath" "testing" "time"
coresqlite "tangled.org/core/sqlite")
func newTestDB(t *testing.T) *DB { t.Helper() dir := t.TempDir() dbPath := filepath.Join(dir, "test.db") db, err := Make(context.Background(), dbPath) if err != nil { t.Fatalf("Make db failed: %v", err) } t.Cleanup(func() { db.Close() }) return db}
func TestDBSchemaAndIdempotency(t *testing.T) { ctx := context.Background() db := newTestDB(t)
token := "enc_token_123" in := CreateBatchInput{ ID: "batch-1", OwnerDid: "did:plc:alice", RequestID: "req-1", RequestDigest: "digest-aaa", EncryptedToken: &token, Jobs: []CreateJobInput{ { Name: "repo1", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/repo1", Private: true, }, }, }
batch, jobs, reused, err := db.CreateBatch(ctx, in) if err != nil { t.Fatalf("CreateBatch failed: %v", err) } if reused { t.Fatal("expected new batch, got reused") } if batch.ID != "batch-1" || len(jobs) != 1 { t.Fatalf("unexpected batch/jobs: %+v, %+v", batch, jobs) } if jobs[0].RepoDid != "" { t.Fatalf("queued job accepted a pre-minted repo did: %q", jobs[0].RepoDid) }
batch2, jobs2, reused2, err := db.CreateBatch(ctx, in) if err != nil { t.Fatalf("CreateBatch idempotent replay failed: %v", err) } if !reused2 { t.Fatal("expected reused batch on idempotent replay") } if batch2.ID != batch.ID || len(jobs2) != len(jobs) { t.Fatalf("mismatched replay: %+v vs %+v", batch2, batch) }
conflictIn := in conflictIn.RequestDigest = "digest-different" _, _, _, err = db.CreateBatch(ctx, conflictIn) if err != ErrConflict { t.Fatalf("expected ErrConflict, got %v", err) }}
func TestDBClaimAndResetRecovery(t *testing.T) { ctx := context.Background() db := newTestDB(t)
in := CreateBatchInput{ ID: "batch-1", OwnerDid: "did:plc:alice", RequestID: "req-1", RequestDigest: "digest-1", Jobs: []CreateJobInput{ { Name: "repo1", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/repo1", Private: false, }, }, } _, jobs, _, err := db.CreateBatch(ctx, in) if err != nil { t.Fatal(err) }
job, batch, err := db.ClaimNextQueuedJob(ctx) if err != nil { t.Fatal(err) } if job == nil || batch == nil { t.Fatal("expected claimed job") } if job.ID != jobs[0].ID || job.Status != StatusCloning || job.Attempts != 1 { t.Fatalf("unexpected claimed job state: %+v", job) }
err = db.UpdateJobStatus(ctx, job.ID, StatusImporting, nil) if err != nil { t.Fatal(err) }
resetCount, err := db.ResetInflightJobsToQueued(ctx) if err != nil { t.Fatal(err) } if resetCount != 1 { t.Fatalf("expected 1 reset job, got %d", resetCount) }
reloaded, _, err := db.GetJob(ctx, job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != StatusQueued { t.Fatalf("expected queued, got %s", reloaded.Status) }}
func TestDBScrubTokenOnTerminalBatch(t *testing.T) { ctx := context.Background() db := newTestDB(t)
token := "secret_token" in := CreateBatchInput{ ID: "batch-1", OwnerDid: "did:plc:alice", RequestID: "req-1", RequestDigest: "digest-1", EncryptedToken: &token, Jobs: []CreateJobInput{ { Name: "repo1", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/repo1", Private: true, }, { Name: "repo2", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/repo2", Private: true, }, }, } _, jobs, _, err := db.CreateBatch(ctx, in) if err != nil { t.Fatal(err) }
err = db.UpdateJobStatus(ctx, jobs[0].ID, StatusCompleted, nil) if err != nil { t.Fatal(err) }
scrubbed, err := db.CheckAndScrubBatchToken(ctx, "batch-1") if err != nil { t.Fatal(err) } if scrubbed { t.Fatal("should not have scrubbed token while job 2 is active") }
err = db.UpdateJobStatus(ctx, jobs[1].ID, StatusFailed, nil) if err != nil { t.Fatal(err) }
scrubbed, err = db.CheckAndScrubBatchToken(ctx, "batch-1") if err != nil { t.Fatal(err) } if !scrubbed { t.Fatal("expected token to be scrubbed when all jobs are terminal") }
b, _, err := db.GetBatch(ctx, "batch-1") if err != nil { t.Fatal(err) } if b.EncryptedToken != nil { t.Fatalf("expected nil EncryptedToken, got %v", *b.EncryptedToken) }}
func TestDBRequeueJob(t *testing.T) { ctx := context.Background() db := newTestDB(t)
in := CreateBatchInput{ ID: "batch-1", OwnerDid: "did:plc:alice", RequestID: "req-1", RequestDigest: "digest-1", Jobs: []CreateJobInput{ { Name: "repo1", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/repo1", Private: true, }, }, } _, jobs, _, err := db.CreateBatch(ctx, in) if err != nil { t.Fatal(err) }
_, _, err = db.RequeueJob(ctx, jobs[0].ID, "did:plc:alice", nil) if err == nil { t.Fatal("expected error retrying queued job") }
if _, _, err = db.ClaimNextQueuedJob(ctx); err != nil { t.Fatal(err) } if err := db.UpdateJobStatus(ctx, jobs[0].ID, StatusImporting, nil); err != nil { t.Fatal(err) } const capability = "cap-token-000000000000000000000000000" if err := db.SetJobCapability(ctx, jobs[0].ID, capability); err != nil { t.Fatal(err) } if err := db.SetJobRepoDid(ctx, jobs[0].ID, "did:plc:created"); err != nil { t.Fatal(err) }
errMsg := "git error" err = db.UpdateJobStatus(ctx, jobs[0].ID, StatusFailed, &errMsg) if err != nil { t.Fatal(err) }
_, _, err = db.RequeueJob(ctx, jobs[0].ID, "did:plc:bob", nil) if err != ErrNotFound { t.Fatalf("expected ErrNotFound for foreign ownerDid, got %v", err) }
newToken := "new_encrypted_token" parentBatch, batchJobs, err := db.RequeueJob(ctx, jobs[0].ID, "did:plc:alice", &newToken) if err != nil { t.Fatalf("RequeueJob failed: %v", err) } if parentBatch.ID != "batch-1" || len(batchJobs) != 1 { t.Fatalf("unexpected parent batch: %+v", parentBatch) } if batchJobs[0].Status != StatusQueued { t.Fatalf("expected queued, got %s", batchJobs[0].Status) } if batchJobs[0].Error != nil { t.Fatalf("expected error cleared, got %v", *batchJobs[0].Error) } if batchJobs[0].RepoDid != "did:plc:created" { t.Fatalf("repo_did = %q, want the knot identity preserved for state-aware resume", batchJobs[0].RepoDid) } if batchJobs[0].CapabilityToken == nil || *batchJobs[0].CapabilityToken != capability { t.Fatalf("capability = %v, want it alive until a terminal state", batchJobs[0].CapabilityToken) }
b, _, err := db.GetBatch(ctx, "batch-1") if err != nil { t.Fatal(err) } if b.EncryptedToken == nil || *b.EncryptedToken != newToken { t.Fatalf("expected batch token updated to %q, got %v", newToken, b.EncryptedToken) }}
func TestDBQueuedCredentialExpiry(t *testing.T) { ctx := context.Background() db := newTestDB(t)
token := "enc_token_expiry" in := CreateBatchInput{ ID: "batch-exp", OwnerDid: "did:plc:alice", RequestID: "req-exp", RequestDigest: "digest-exp", EncryptedToken: &token, Jobs: []CreateJobInput{ { Name: "exp1", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/exp1", Private: true, }, }, } _, jobs, _, err := db.CreateBatch(ctx, in) if err != nil { t.Fatal(err) }
past := time.Now().UTC().Add(-2 * time.Hour).Format(time.RFC3339) _, err = db.ExecContext(ctx, "update batches set credential_expires_at = ? where id = 'batch-exp'", past) if err != nil { t.Fatal(err) }
scrubbed, err := db.ScrubExpiredCredentials(ctx) if err != nil { t.Fatal(err) } if scrubbed != 1 { t.Fatalf("expected 1 scrubbed batch, got %d", scrubbed) }
b, _, err := db.GetBatch(ctx, "batch-exp") if err != nil { t.Fatal(err) } if b.EncryptedToken != nil { t.Fatalf("expected token scrubbed, got %v", *b.EncryptedToken) }
j, _, err := db.GetJob(ctx, jobs[0].ID) if err != nil { t.Fatal(err) } if j.Status != StatusAuthorizationRequired { t.Fatalf("expected status authorization_required, got %s", j.Status) }}
func TestExistingJobsTableGainsCapabilityColumn(t *testing.T) { path := filepath.Join(t.TempDir(), "existing.db") old, err := coresqlite.Open(path) if err != nil { t.Fatal(err) } _, err = old.Exec(`create table jobs (id integer primary key, batch_id text, repo_did text, name text, knot_did text, source_url text, private integer, status text, attempts integer, next_attempt_at text, error text, created_at text, updated_at text)`) if err != nil { t.Fatal(err) } _ = old.Close()
database, err := Make(context.Background(), path) if err != nil { t.Fatalf("additive migration failed: %v", err) } defer database.Close() rows, err := database.Query("pragma table_info(jobs)") if err != nil { t.Fatal(err) } defer rows.Close() found := false for rows.Next() { var cid, notNull, pk int var name, typ string var defaultValue any if err := rows.Scan(&cid, &name, &typ, ¬Null, &defaultValue, &pk); err != nil { t.Fatal(err) } found = found || name == "capability_token" } if !found { t.Fatal("existing jobs table did not gain capability_token") }}
func TestCapabilityResolutionRequiresTokenAndName(t *testing.T) { database := newTestDB(t) _, jobs, _, err := database.CreateBatch(context.Background(), CreateBatchInput{ ID: "cap", OwnerDid: "did:plc:alice", RequestID: "request", RequestDigest: "digest", Jobs: []CreateJobInput{{Name: "repo", KnotDid: "did:web:knot.example", SourceURL: "https://github.com/alice/repo.git"}}, }) if err != nil { t.Fatal(err) } if err := database.UpdateJobStatus(context.Background(), jobs[0].ID, StatusImporting, nil); err != nil { t.Fatal(err) } token := "secret-capability" if err := database.SetJobCapability(context.Background(), jobs[0].ID, token); err != nil { t.Fatal(err) } if job, err := database.ResolveJobCapability(context.Background(), token, "repo"); err != nil || job.ID != jobs[0].ID { t.Fatalf("exact mapping: job=%v err=%v", job, err) } if _, err := database.ResolveJobCapability(context.Background(), token, "other"); err != ErrNotFound { t.Fatalf("wrong name: %v", err) } if err := database.RevokeJobCapability(context.Background(), jobs[0].ID); err != nil { t.Fatal(err) } if _, err := database.ResolveJobCapability(context.Background(), token, "repo"); err != ErrNotFound { t.Fatalf("revoked token: %v", err) }}
func TestReleaseJobRefundsTheClaimedAttempt(t *testing.T) { ctx := context.Background() database := newTestDB(t)
if _, _, _, err := database.CreateBatch(ctx, CreateBatchInput{ ID: "batch-release", OwnerDid: "did:plc:alice", RequestID: "req-release", RequestDigest: "digest", Jobs: []CreateJobInput{{Name: "one", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/one"}}, }); err != nil { t.Fatal(err) }
job, _, err := database.ClaimNextQueuedJob(ctx) if err != nil || job == nil { t.Fatalf("claim: job=%v err=%v", job, err) } if job.Attempts != 1 || job.Status != StatusCloning { t.Fatalf("claim left %q at %d attempts", job.Status, job.Attempts) }
reason := "daemon stopping" if err := database.ReleaseJob(ctx, job.ID, &reason); err != nil { t.Fatal(err) }
reloaded, _, err := database.GetJob(ctx, job.ID) if err != nil { t.Fatal(err) } if reloaded.Status != StatusQueued || reloaded.Attempts != 0 { t.Fatalf("released job is %q at %d attempts, want queued at 0", reloaded.Status, reloaded.Attempts) } if reloaded.NextAttemptAt.After(time.Now().Add(time.Second)) { t.Fatalf("a released job waits for %s, want it claimable now", reloaded.NextAttemptAt) }
again, _, err := database.ClaimNextQueuedJob(ctx) if err != nil || again == nil || again.ID != job.ID { t.Fatalf("released job was not claimable: job=%v err=%v", again, err) } if again.Attempts != 1 { t.Fatalf("second claim left %d attempts, want 1", again.Attempts) }}
func TestClaimRespectsTheRetrySchedule(t *testing.T) { ctx := context.Background() database := newTestDB(t)
if _, _, _, err := database.CreateBatch(ctx, CreateBatchInput{ ID: "batch-backoff", OwnerDid: "did:plc:alice", RequestID: "req-backoff", RequestDigest: "digest", Jobs: []CreateJobInput{{Name: "one", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/one"}}, }); err != nil { t.Fatal(err) } job, _, err := database.ClaimNextQueuedJob(ctx) if err != nil || job == nil { t.Fatalf("claim: job=%v err=%v", job, err) } reason := "the source said no" if err := database.ScheduleJobRetry(ctx, job.ID, time.Minute, &reason); err != nil { t.Fatal(err) }
next, _, err := database.ClaimNextQueuedJob(ctx) if err != nil { t.Fatal(err) } if next != nil { t.Fatalf("a job held back for a minute was claimed after %d attempt(s)", next.Attempts) }}
func TestClaimCapsEachOwnerAndSharesFairly(t *testing.T) { ctx := context.Background() database := newTestDB(t) database.MaxActivePerOwner = 2
_, _, _, err := database.CreateBatch(ctx, CreateBatchInput{ ID: "batch-alice", OwnerDid: "did:plc:alice", RequestID: "req-a", RequestDigest: "digest-a", Jobs: []CreateJobInput{ {Name: "a1", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/a1"}, {Name: "a2", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/a2"}, {Name: "a3", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/a3"}, }, }) if err != nil { t.Fatal(err) } _, bobJobs, _, err := database.CreateBatch(ctx, CreateBatchInput{ ID: "batch-bob", OwnerDid: "did:plc:bob", RequestID: "req-b", RequestDigest: "digest-b", Jobs: []CreateJobInput{{Name: "b1", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/bob/b1"}}, }) if err != nil { t.Fatal(err) } b1 := bobJobs[0].ID
name := func(jobID int64) string { job, _, err := database.GetJob(ctx, jobID) if err != nil { t.Fatal(err) } return job.Name }
// alice never takes a third slot while her two claims run; bob's single job // claims ahead of alice's older queue instead of waiting behind it want := []struct { jobID int64 failed bool }{ {1, false}, // neither owner has an active job, so the oldest wins {4, false}, // alice has one active, bob none, so bob goes first {2, false}, // alice has one active, bob one, so alice's next {0, true}, // alice at the cap, bob has nothing queued } for _, step := range want { job, _, err := database.ClaimNextQueuedJob(ctx) if err != nil { t.Fatal(err) } if step.failed { if job != nil { t.Fatalf("claimed job %d while every owner sat at its cap", job.ID) } continue } if job == nil || job.ID != step.jobID { got := "nil" if job != nil { got = name(job.ID) } t.Fatalf("claim = %s, want job %d", got, step.jobID) } }
if err := database.UpdateJobStatus(ctx, 1, StatusCompleted, nil); err != nil { t.Fatal(err) } if err := database.UpdateJobStatus(ctx, b1, StatusCompleted, nil); err != nil { t.Fatal(err) } job, _, err := database.ClaimNextQueuedJob(ctx) if err != nil || job == nil || job.ID != 3 { t.Fatalf("a finished job freed no slot: job=%v err=%v", job, err) }}
func TestAScheduledRetryKeepsTheRepoDidAndCapability(t *testing.T) { ctx := context.Background() database := newTestDB(t)
if _, _, _, err := database.CreateBatch(ctx, CreateBatchInput{ ID: "batch-retry", OwnerDid: "did:plc:alice", RequestID: "req-r", RequestDigest: "digest", Jobs: []CreateJobInput{{Name: "one", KnotDid: "did:web:knot.test", SourceURL: "https://github.com/alice/one"}}, }); err != nil { t.Fatal(err) } job, _, err := database.ClaimNextQueuedJob(ctx) if err != nil || job == nil { t.Fatalf("claim: job=%v err=%v", job, err) } if err := database.UpdateJobStatus(ctx, job.ID, StatusImporting, nil); err != nil { t.Fatal(err) } const capability = "cap-token-111111111111111111111111111" if err := database.SetJobCapability(ctx, job.ID, capability); err != nil { t.Fatal(err) } if err := database.SetJobRepoDid(ctx, job.ID, "did:plc:halfstaged"); err != nil { t.Fatal(err) }
reason := "the create stalled halfway" if err := database.ScheduleJobRetry(ctx, job.ID, time.Second, &reason); err != nil { t.Fatal(err) } reloaded, _, err := database.GetJob(ctx, job.ID) if err != nil { t.Fatal(err) } if reloaded.RepoDid != "did:plc:halfstaged" { t.Fatalf("repo_did = %q, want state-aware resume of the existing knot repo", reloaded.RepoDid) } if reloaded.CapabilityToken == nil || *reloaded.CapabilityToken != capability { t.Fatalf("capability = %v, want it alive so the knot keeps fetching during backoff", reloaded.CapabilityToken) }}