package api import ( "bytes" "context" "crypto/rand" "encoding/base64" "encoding/json" "fmt" "net/http" "net/http/httptest" "path/filepath" "strings" "testing" "time" "github.com/bluesky-social/indigo/atproto/atcrypto" "github.com/bluesky-social/indigo/atproto/auth" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/api/org_tangled" "tangled.org/core/log" "tangled.org/core/migrator/config" "tangled.org/core/migrator/db" "tangled.org/core/xrpc/serviceauth" ) const ( alice = "did:plc:alice" bob = "did:plc:bob" ) type caller struct { did syntax.DID signer atcrypto.PrivateKey } type allowGrants struct{} func (allowGrants) HasSession(context.Context, string) bool { return true } type denyGrants struct{} func (denyGrants) HasSession(context.Context, string) bool { return false } type testServer struct { server *Server db *db.DB cfg *config.Config alice caller bob caller } func (ts *testServer) token(t *testing.T, from caller, lxm string) string { t.Helper() nsid, err := syntax.ParseNSID(lxm) if err != nil { t.Fatal(err) } token, err := auth.SignServiceAuth(from.did, ts.cfg.ServiceDid.String(), time.Minute, &nsid, from.signer) if err != nil { t.Fatal(err) } return token } func (ts *testServer) do(t *testing.T, method, path, token string, body any) *httptest.ResponseRecorder { t.Helper() var reader *bytes.Reader if body != nil { encoded, err := json.Marshal(body) if err != nil { t.Fatal(err) } reader = bytes.NewReader(encoded) } else { reader = bytes.NewReader(nil) } req := httptest.NewRequest(method, path, reader) req.Header.Set("Content-Type", "application/json") if token != "" { req.Header.Set("Authorization", "Bearer "+token) } rec := httptest.NewRecorder() ts.server.Routes().ServeHTTP(rec, req) return rec } func setupTestServer(t *testing.T) *testServer { t.Helper() // the local test pds dials on loopback, which the oauth client's ssrf guard refuses t.Setenv("ATPROTO_OAUTH_DEV", "1") dir := t.TempDir() database, err := db.Make(context.Background(), filepath.Join(dir, "api_test.db")) if err != nil { t.Fatal(err) } t.Cleanup(func() { database.Close() }) rawKey := make([]byte, 32) rand.Read(rawKey) cfg := &config.Config{ Hostname: "migrator.example.com", PrivateKey: "private-multibase", MasterKey: base64.StdEncoding.EncodeToString(rawKey), WorkDir: filepath.Join(dir, "work"), } if err := cfg.Validate(); err != nil { t.Fatal(err) } logger := log.New("test") directory := identity.NewMockDirectory() newCaller := func(did string) caller { priv, err := atcrypto.GeneratePrivateKeyP256() if err != nil { t.Fatal(err) } pub, err := priv.PublicKey() if err != nil { t.Fatal(err) } directory.Insert(identity.Identity{ DID: syntax.DID(did), Keys: map[string]identity.VerificationMethod{ "atproto": {Type: "Multikey", PublicKeyMultibase: pub.Multibase()}, }, }) return caller{did: syntax.DID(did), signer: priv} } serviceAuth := serviceauth.NewServiceAuth(logger, directory, cfg.ServiceDid.String()) return &testServer{ server: NewServer(database, cfg, nil, serviceAuth, allowGrants{}, logger, nil, ""), db: database, cfg: cfg, alice: newCaller(alice), bob: newCaller(bob), } } func job(name string, private bool) *org_tangled.TempMigratorDefs_NewJob { return &org_tangled.TempMigratorDefs_NewJob{ Name: name, KnotDid: "did:web:knot.example.com", SourceUrl: "https://github.com/alice/" + name + ".git", Private: &private, } } func decodeTask(t *testing.T, rec *httptest.ResponseRecorder) *org_tangled.TempMigratorDefs_Task { t.Helper() var task org_tangled.TempMigratorDefs_Task if err := json.Unmarshal(rec.Body.Bytes(), &task); err != nil { t.Fatalf("decoding task: %v (body: %s)", err, rec.Body.String()) } return &task } func TestPreflightIsAnsweredWithoutAToken(t *testing.T) { ts := setupTestServer(t) createPath := "/xrpc/" + org_tangled.TempMigratorCreateTaskNSID req := httptest.NewRequest(http.MethodOptions, createPath, nil) req.Header.Set("Origin", "https://web.tngl.boltless.dev") req.Header.Set("Access-Control-Request-Method", http.MethodPost) rec := httptest.NewRecorder() ts.server.Routes().ServeHTTP(rec, req) if rec.Code != http.StatusNoContent { t.Fatalf("preflight: want %d, got %d", http.StatusNoContent, rec.Code) } if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "*" { t.Fatalf("preflight: allow-origin = %q", got) } if got := rec.Header().Get("Access-Control-Allow-Headers"); !strings.Contains(got, "Authorization") { t.Fatalf("preflight: allow-headers = %q, and the caller sends a bearer token", got) } body := &org_tangled.TempMigratorCreateTask_Input{RequestId: "r1", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false)}} rec = ts.do(t, http.MethodPost, createPath, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), body) if rec.Code != http.StatusAccepted { t.Fatalf("create with a token: want %d, got %d (%s)", http.StatusAccepted, rec.Code, rec.Body.String()) } if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "*" { t.Fatalf("create: allow-origin = %q", got) } } func TestXrpcRequiresATokenScopedToTheMethod(t *testing.T) { ts := setupTestServer(t) createPath := "/xrpc/" + org_tangled.TempMigratorCreateTaskNSID body := &org_tangled.TempMigratorCreateTask_Input{RequestId: "r1", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false)}} if rec := ts.do(t, http.MethodPost, createPath, "", body); rec.Code != http.StatusForbidden { t.Fatalf("no token: want 403, got %d", rec.Code) } otherToken := ts.token(t, ts.alice, org_tangled.TempMigratorListTasksNSID) if rec := ts.do(t, http.MethodPost, createPath, otherToken, body); rec.Code != http.StatusForbidden { t.Fatalf("wrong lxm: want 403, got %d", rec.Code) } nsid, _ := syntax.ParseNSID(org_tangled.TempMigratorCreateTaskNSID) wrongAudience, err := auth.SignServiceAuth(ts.alice.did, "did:web:elsewhere.example", time.Minute, &nsid, ts.alice.signer) if err != nil { t.Fatal(err) } if rec := ts.do(t, http.MethodPost, createPath, wrongAudience, body); rec.Code != http.StatusForbidden { t.Fatalf("wrong audience: want 403, got %d", rec.Code) } stranger, err := atcrypto.GeneratePrivateKeyP256() if err != nil { t.Fatal(err) } unknown, err := auth.SignServiceAuth("did:plc:nobody", ts.cfg.ServiceDid.String(), time.Minute, &nsid, stranger) if err != nil { t.Fatal(err) } if rec := ts.do(t, http.MethodPost, createPath, unknown, body); rec.Code != http.StatusForbidden { t.Fatalf("unknown issuer: want 403, got %d", rec.Code) } token := ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID) if rec := ts.do(t, http.MethodPost, createPath, token, body); rec.Code != http.StatusAccepted { t.Fatalf("valid token: want 202, got %d (body: %s)", rec.Code, rec.Body.String()) } } func TestCreateTaskAnswersWithTheTaskAndReplaysIt(t *testing.T) { ts := setupTestServer(t) token := ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID) path := "/xrpc/" + org_tangled.TempMigratorCreateTaskNSID body := &org_tangled.TempMigratorCreateTask_Input{RequestId: "req-1", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false), job("two", false)}} rec := ts.do(t, http.MethodPost, path, token, body) if rec.Code != http.StatusAccepted { t.Fatalf("want 202, got %d (body: %s)", rec.Code, rec.Body.String()) } if got := rec.Header().Get("Cache-Control"); !strings.Contains(got, "no-store") { t.Fatalf("tasks are per-caller state and must not be cached, got %q", got) } task := decodeTask(t, rec) if task.Id == "" { t.Fatal("task has no id") } if task.OwnerDid != alice { t.Fatalf("owner = %q, want the token issuer %q", task.OwnerDid, alice) } if len(task.Jobs) != 2 { t.Fatalf("want 2 jobs, got %d", len(task.Jobs)) } for _, j := range task.Jobs { if j.Id == "" { t.Fatal("job id is empty, so nothing can be retried by id later") } if j.Status != string(db.StatusQueued) { t.Fatalf("new job status = %q, want %q", j.Status, db.StatusQueued) } if j.RepoDid != nil { t.Fatalf("new job was answered with a repo did before the knot made one: %q", *j.RepoDid) } } if _, err := time.Parse(time.RFC3339, task.CreatedAt); err != nil { t.Fatalf("createdAt %q is not a datetime: %v", task.CreatedAt, err) } replay := ts.do(t, http.MethodPost, path, token, body) if replay.Code != http.StatusOK { t.Fatalf("replay: want 200, got %d (body: %s)", replay.Code, replay.Body.String()) } if again := decodeTask(t, replay); again.Id != task.Id { t.Fatalf("replay made a second task: %q then %q", task.Id, again.Id) } conflict := &org_tangled.TempMigratorCreateTask_Input{RequestId: "req-1", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("three", false)}} rec = ts.do(t, http.MethodPost, path, token, conflict) if rec.Code != http.StatusConflict { t.Fatalf("want 409, got %d (body: %s)", rec.Code, rec.Body.String()) } if !strings.Contains(rec.Body.String(), "TaskConflict") { t.Fatalf("conflict body = %s, want the TaskConflict tag", rec.Body.String()) } } func TestCreateTaskRejectsBadJobs(t *testing.T) { ts := setupTestServer(t) token := ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID) path := "/xrpc/" + org_tangled.TempMigratorCreateTaskNSID bad := []struct { name string body *org_tangled.TempMigratorCreateTask_Input }{ {"no requestId", &org_tangled.TempMigratorCreateTask_Input{Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false)}}}, {"no jobs", &org_tangled.TempMigratorCreateTask_Input{RequestId: "r"}}, {"private without a token", &org_tangled.TempMigratorCreateTask_Input{RequestId: "r", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", true)}}}, } for _, tc := range bad { rec := ts.do(t, http.MethodPost, path, token, tc.body) if rec.Code != http.StatusBadRequest { t.Fatalf("%s: want 400, got %d (body: %s)", tc.name, rec.Code, rec.Body.String()) } if !strings.Contains(rec.Body.String(), "InvalidRequest") { t.Fatalf("%s: body = %s, want the InvalidRequest tag", tc.name, rec.Body.String()) } } tooMany := &org_tangled.TempMigratorCreateTask_Input{RequestId: "r"} for i := 0; i < maxJobsPerTask+1; i++ { tooMany.Jobs = append(tooMany.Jobs, job("repo", false)) } if rec := ts.do(t, http.MethodPost, path, token, tooMany); rec.Code != http.StatusBadRequest { t.Fatalf("too many jobs: want 400, got %d", rec.Code) } } func TestTasksAreOnlyVisibleToTheirOwner(t *testing.T) { ts := setupTestServer(t) created := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: "r", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false)}}) task := decodeTask(t, created) aliceGet := ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorGetTaskNSID+"?taskId="+task.Id, ts.token(t, ts.alice, org_tangled.TempMigratorGetTaskNSID), nil) if aliceGet.Code != http.StatusOK { t.Fatalf("owner read: want 200, got %d (body: %s)", aliceGet.Code, aliceGet.Body.String()) } bobGet := ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorGetTaskNSID+"?taskId="+task.Id, ts.token(t, ts.bob, org_tangled.TempMigratorGetTaskNSID), nil) if bobGet.Code != http.StatusNotFound { t.Fatalf("foreign read: want 404, got %d (body: %s)", bobGet.Code, bobGet.Body.String()) } var bobList org_tangled.TempMigratorListTasks_Output rec := ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorListTasksNSID, ts.token(t, ts.bob, org_tangled.TempMigratorListTasksNSID), nil) if err := json.Unmarshal(rec.Body.Bytes(), &bobList); err != nil { t.Fatal(err) } if len(bobList.Tasks) != 0 { t.Fatalf("bob sees %d tasks, want none", len(bobList.Tasks)) } var aliceList org_tangled.TempMigratorListTasks_Output rec = ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorListTasksNSID, ts.token(t, ts.alice, org_tangled.TempMigratorListTasksNSID), nil) if err := json.Unmarshal(rec.Body.Bytes(), &aliceList); err != nil { t.Fatal(err) } if len(aliceList.Tasks) != 1 || aliceList.Tasks[0].Id != task.Id { t.Fatalf("owner list = %+v, want the one task %q", aliceList.Tasks, task.Id) } if len(aliceList.Tasks[0].Jobs) != 1 { t.Fatalf("listed task lost its jobs: %+v", aliceList.Tasks[0]) } } func TestGetTaskAndListTasksRejectBadInput(t *testing.T) { ts := setupTestServer(t) if rec := ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorGetTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorGetTaskNSID), nil); rec.Code != http.StatusBadRequest { t.Fatalf("missing taskId: want 400, got %d", rec.Code) } rec := ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorGetTaskNSID+"?taskId=missing", ts.token(t, ts.alice, org_tangled.TempMigratorGetTaskNSID), nil) if rec.Code != http.StatusNotFound || !strings.Contains(rec.Body.String(), "TaskNotFound") { t.Fatalf("unknown task: want 404 TaskNotFound, got %d %s", rec.Code, rec.Body.String()) } rec = ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorListTasksNSID+"?limit=soon", ts.token(t, ts.alice, org_tangled.TempMigratorListTasksNSID), nil) if rec.Code != http.StatusBadRequest { t.Fatalf("unparseable limit: want 400, got %d", rec.Code) } rec = ts.do(t, http.MethodGet, "/xrpc/"+org_tangled.TempMigratorListTasksNSID+"?limit=100000", ts.token(t, ts.alice, org_tangled.TempMigratorListTasksNSID), nil) if rec.Code != http.StatusOK { t.Fatalf("huge limit: want 200, got %d (body: %s)", rec.Code, rec.Body.String()) } } func TestRetryJobRequeuesOnlyTheCallersOwnJob(t *testing.T) { ts := setupTestServer(t) created := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: "r", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false)}}) task := decodeTask(t, created) jobID := task.Jobs[0].Id retryBody := &org_tangled.TempMigratorRetryJob_Input{TaskId: task.Id, JobId: jobID} rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorRetryJobNSID, ts.token(t, ts.alice, org_tangled.TempMigratorRetryJobNSID), retryBody) if rec.Code != http.StatusBadRequest || !strings.Contains(rec.Body.String(), "JobNotRetryable") { t.Fatalf("retrying a queued job: want 400 JobNotRetryable, got %d %s", rec.Code, rec.Body.String()) } foreign := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorRetryJobNSID, ts.token(t, ts.bob, org_tangled.TempMigratorRetryJobNSID), retryBody) if foreign.Code != http.StatusNotFound { t.Fatalf("foreign retry: want 404, got %d (body: %s)", foreign.Code, foreign.Body.String()) } mismatch := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorRetryJobNSID, ts.token(t, ts.alice, org_tangled.TempMigratorRetryJobNSID), &org_tangled.TempMigratorRetryJob_Input{TaskId: "somewhere-else", JobId: jobID}) if mismatch.Code != http.StatusNotFound { t.Fatalf("task mismatch: want 404, got %d", mismatch.Code) } id, err := parseJobID(jobID) if err != nil { t.Fatal(err) } reason := "the source said no" if err := ts.db.UpdateJobStatus(context.Background(), id, db.StatusFailed, &reason); err != nil { t.Fatal(err) } rec = ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorRetryJobNSID, ts.token(t, ts.alice, org_tangled.TempMigratorRetryJobNSID), retryBody) if rec.Code != http.StatusOK { t.Fatalf("retry: want 200, got %d (body: %s)", rec.Code, rec.Body.String()) } retried := decodeTask(t, rec) if len(retried.Jobs) != 1 { t.Fatalf("retry answered with %d jobs, want the task's one", len(retried.Jobs)) } if retried.Jobs[0].Status != string(db.StatusQueued) { t.Fatalf("retried job status = %q, want %q", retried.Jobs[0].Status, db.StatusQueued) } if retried.Jobs[0].Error != nil { t.Fatalf("retried job kept its error: %q", *retried.Jobs[0].Error) } if retried.Jobs[0].Attempts != 0 { t.Fatalf("retried job attempts = %d, want a full budget", retried.Jobs[0].Attempts) } } func TestCreateTaskEnforcesTheSchemasBounds(t *testing.T) { ts := setupTestServer(t) long := strings.Repeat("r", 129) rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: long, Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false)}}) if rec.Code != http.StatusBadRequest { t.Fatalf("requestId over the schema bound: want 400, got %d (body: %s)", rec.Code, rec.Body.String()) } twice := []*org_tangled.TempMigratorDefs_NewJob{job("same", false), job("Same", false)} rec = ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: "dupes", Jobs: twice}) if rec.Code != http.StatusBadRequest { t.Fatalf("two jobs for one name collide on the knot: want 400, got %d (body: %s)", rec.Code, rec.Body.String()) } description := strings.Repeat("d", 141) withDescription := job("described", false) withDescription.Description = &description rec = ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: "long-desc", Jobs: []*org_tangled.TempMigratorDefs_NewJob{withDescription}}) if rec.Code != http.StatusBadRequest { t.Fatalf("description over the schema bound: want 400, got %d (body: %s)", rec.Code, rec.Body.String()) } } func TestTheSameRequestIdWithADifferentDescriptionConflicts(t *testing.T) { ts := setupTestServer(t) token := ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID) first := job("described", false) one := "a small tool" first.Description = &one rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, token, &org_tangled.TempMigratorCreateTask_Input{RequestId: "same", Jobs: []*org_tangled.TempMigratorDefs_NewJob{first}}) if rec.Code != http.StatusAccepted { t.Fatalf("first: want 202, got %d (body: %s)", rec.Code, rec.Body.String()) } second := job("described", false) other := "a different tool" second.Description = &other rec = ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, token, &org_tangled.TempMigratorCreateTask_Input{RequestId: "same", Jobs: []*org_tangled.TempMigratorDefs_NewJob{second}}) if rec.Code != http.StatusConflict { t.Fatalf("want 409 for the same requestId with another description, got %d (body: %s)", rec.Code, rec.Body.String()) } } func TestRetryRefusesAnExpiredStoredCredential(t *testing.T) { ts := setupTestServer(t) ghToken := "ghp_secret" rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: "expired", GithubToken: &ghToken, Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("private-one", true)}}) task := decodeTask(t, rec) id, err := parseJobID(task.Jobs[0].Id) if err != nil { t.Fatal(err) } reason := "the source said no" if err := ts.db.UpdateJobStatus(context.Background(), id, db.StatusFailed, &reason); err != nil { t.Fatal(err) } expired := time.Now().Add(-time.Minute).UTC().Format(time.RFC3339) if _, err := ts.db.ExecContext(context.Background(), "update batches set credential_expires_at = ? where id = ?", expired, task.Id); err != nil { t.Fatal(err) } rec = ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorRetryJobNSID, ts.token(t, ts.alice, org_tangled.TempMigratorRetryJobNSID), &org_tangled.TempMigratorRetryJob_Input{TaskId: task.Id, JobId: task.Jobs[0].Id}) if rec.Code != http.StatusBadRequest { t.Fatalf("want 400 for an expired stored credential, got %d (body: %s)", rec.Code, rec.Body.String()) } } func TestPrivateRetryNeedsAFreshTokenOnceTheStoredOneIsDropped(t *testing.T) { ts := setupTestServer(t) ghToken := "ghp_secret" created := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: "r", GithubToken: &ghToken, Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("private-one", true)}}) task := decodeTask(t, created) jobID := task.Jobs[0].Id id, err := parseJobID(jobID) if err != nil { t.Fatal(err) } if strings.Contains(created.Body.String(), ghToken) { t.Fatal("the github token was echoed back to the caller") } reason := "the source said no" if err := ts.db.UpdateJobStatus(context.Background(), id, db.StatusAuthorizationRequired, &reason); err != nil { t.Fatal(err) } if _, err := ts.db.CheckAndScrubBatchToken(context.Background(), task.Id); err != nil { t.Fatal(err) } rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorRetryJobNSID, ts.token(t, ts.alice, org_tangled.TempMigratorRetryJobNSID), &org_tangled.TempMigratorRetryJob_Input{TaskId: task.Id, JobId: jobID}) if rec.Code != http.StatusBadRequest || !strings.Contains(rec.Body.String(), "githubToken") { t.Fatalf("retry without a token: want 400 asking for one, got %d %s", rec.Code, rec.Body.String()) } rec = ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorRetryJobNSID, ts.token(t, ts.alice, org_tangled.TempMigratorRetryJobNSID), &org_tangled.TempMigratorRetryJob_Input{TaskId: task.Id, JobId: jobID, GithubToken: &ghToken}) if rec.Code != http.StatusOK { t.Fatalf("retry with a token: want 200, got %d (body: %s)", rec.Code, rec.Body.String()) } } func TestDIDDocumentIsServedForTheSkillBackedIssuer(t *testing.T) { ts := setupTestServer(t) rec := ts.do(t, http.MethodGet, "/.well-known/did.json", "", nil) if rec.Code != http.StatusServiceUnavailable { t.Fatalf("want 503 without a signer, got %d", rec.Code) } } func parseJobID(raw string) (int64, error) { var id int64 _, err := fmt.Sscan(raw, &id) return id, err } func TestCreateTaskRequiresTheCallersOwnOAuthGrant(t *testing.T) { ts := setupTestServer(t) ts.server.grants = denyGrants{} rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), &org_tangled.TempMigratorCreateTask_Input{RequestId: "grant", Jobs: []*org_tangled.TempMigratorDefs_NewJob{job("one", false)}}) if rec.Code != http.StatusPreconditionRequired || !strings.Contains(rec.Body.String(), "GrantRequired") { t.Fatalf("want 428 GrantRequired, got %d %s", rec.Code, rec.Body.String()) } var count int if err := ts.db.QueryRow("select count(*) from batches").Scan(&count); err != nil { t.Fatal(err) } if count != 0 { t.Fatal("task was created without the caller's grant") } } func TestCreateTaskRejectsLegacyRepoDidInput(t *testing.T) { ts := setupTestServer(t) body := map[string]any{ "requestId": "legacy-repo-did", "jobs": []map[string]any{{ "name": "one", "repoDid": "did:plc:caller-choice", "knotDid": "did:web:knot.example.com", "sourceUrl": "https://github.com/alice/one.git", }}, } rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID), body) if rec.Code != http.StatusBadRequest || !strings.Contains(rec.Body.String(), "unknown field") { t.Fatalf("legacy repoDid input: want 400 unknown field, got %d %s", rec.Code, rec.Body.String()) } } func TestCreateTaskAcceptsAnyHttpsSource(t *testing.T) { ts := setupTestServer(t) token := ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID) path := "/xrpc/" + org_tangled.TempMigratorCreateTaskNSID foreign := job("from-gitlab", false) foreign.SourceUrl = "https://gitlab.example.com/team/project.git" rec := ts.do(t, http.MethodPost, path, token, &org_tangled.TempMigratorCreateTask_Input{ RequestId: "r-gitlab", Jobs: []*org_tangled.TempMigratorDefs_NewJob{foreign}, }) if rec.Code != http.StatusAccepted { t.Fatalf("a non-github https source should be accepted, got %d (body: %s)", rec.Code, rec.Body.String()) } if task := decodeTask(t, rec); task.Jobs[0].SourceUrl != foreign.SourceUrl { t.Fatalf("the source url was rewritten: %q", task.Jobs[0].SourceUrl) } private := job("private-foreign", true) private.SourceUrl = "https://gitlab.example.com/team/private.git" rec = ts.do(t, http.MethodPost, path, token, &org_tangled.TempMigratorCreateTask_Input{ RequestId: "r-private-foreign", Jobs: []*org_tangled.TempMigratorDefs_NewJob{private}, }) if rec.Code != http.StatusBadRequest { t.Fatalf("a private non-github source has to be refused, got %d", rec.Code) } if !strings.Contains(rec.Body.String(), "github.com") { t.Fatalf("the refusal should name the only host a token reaches: %s", rec.Body.String()) } } func TestCreateTaskCarriesTheDescription(t *testing.T) { ts := setupTestServer(t) token := ts.token(t, ts.alice, org_tangled.TempMigratorCreateTaskNSID) description := "a small tool" named := job("described", false) named.Description = &description rec := ts.do(t, http.MethodPost, "/xrpc/"+org_tangled.TempMigratorCreateTaskNSID, token, &org_tangled.TempMigratorCreateTask_Input{ RequestId: "r-described", Jobs: []*org_tangled.TempMigratorDefs_NewJob{named}, }) if rec.Code != http.StatusAccepted { t.Fatalf("want 202, got %d (body: %s)", rec.Code, rec.Body.String()) } task := decodeTask(t, rec) _, jobs, err := ts.db.GetBatch(context.Background(), task.Id) if err != nil { t.Fatal(err) } if len(jobs) != 1 || jobs[0].Description != description { t.Fatalf("the stored job lost the description: %+v", jobs) } }