diff --git a/.gitattributes b/.gitattributes index b359f866..9ad90d86 100644 --- a/.gitattributes +++ b/.gitattributes @@ -2,4 +2,5 @@ api/tangled/** linguist-generated -diff api/tangled/*_ext.go -linguist-generated diff flake.lock -diff web/src/lib/api/lexicons/** linguist-generated -diff +spindle/mill/proto/gen/** linguist-generated -diff Cargo.lock -diff diff --git a/spindle/db/db.go b/spindle/db/db.go index 766b52bd..bd1a3349 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -13,6 +13,7 @@ import ( "tangled.org/core/api/tangled" "tangled.org/core/log" "tangled.org/core/orm" + "tangled.org/core/spindle/models" ) type DB struct { @@ -381,6 +382,17 @@ func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error return err } + if err := orm.RunMigration(conn, logger, "drop-pipeline-knot", func(tx *sql.Tx) error { + _, err := tx.Exec(` + alter table pipelines drop column knot; + alter table jobs drop column pipeline_id_knot; + alter table mill_leases drop column knot; + `) + return err + }); err != nil { + return err + } + return nil } @@ -488,7 +500,7 @@ func migratePipelines(tx *sql.Tx, logger *slog.Logger) error { continue } - p, kind := mapToCiPipeline(rkey, time.Unix(0, created), raw) + p, kind := mapToCiPipeline(models.PipelineId(rkey), time.Unix(0, created), raw) payload, err := json.Marshal(p) if err != nil { skipped++ diff --git a/spindle/db/events.go b/spindle/db/events.go index 44dcb2d7..f60f32d0 100644 --- a/spindle/db/events.go +++ b/spindle/db/events.go @@ -23,7 +23,7 @@ func insertStatusTx(tx DBTX, wid models.WorkflowId, kind models.StatusKind, work _, err := tx.Exec( `insert into workflow_statuses (rkey, workflow, status, error, exit_code, created_at) values (?, ?, ?, ?, ?, ?)`, - wid.PipelineId.Rkey, wid.Name, string(kind), workflowError, exitCode, + wid.PipelineId, wid.Name, string(kind), workflowError, exitCode, time.Now().Format(time.RFC3339), ) return err @@ -93,7 +93,7 @@ func (d *DB) CompleteMillLease( } func (d *DB) GetStatus(workflowId models.WorkflowId) (models.StatusKind, error) { - pipelineId := workflowId.PipelineId.Rkey + pipelineId := workflowId.PipelineId var status string err := d.QueryRow( diff --git a/spindle/db/jobs.go b/spindle/db/jobs.go index e98bc173..05494ca0 100644 --- a/spindle/db/jobs.go +++ b/spindle/db/jobs.go @@ -9,12 +9,11 @@ import ( ) type JobRow struct { - Id int64 - RepoDid string - PipelineIdKnot string - PipelineIdRkey string - SourceRepo *tangled.Pipeline_TriggerRepo - Tpl tangled.Pipeline + Id int64 + RepoDid string + PipelineId models.PipelineId + SourceRepo *tangled.Pipeline_TriggerRepo + Tpl tangled.Pipeline } func (d *DB) EnqueueJob(ctx context.Context, repoDid string, pipelineId models.PipelineId, sourceRepo *tangled.Pipeline_TriggerRepo, tpl tangled.Pipeline) error { @@ -23,9 +22,9 @@ func (d *DB) EnqueueJob(ctx context.Context, repoDid string, pipelineId models.P return err } _, err = d.ExecContext(ctx, ` - insert into jobs (repo_did, pipeline_id_knot, pipeline_id_rkey, source_repo, tpl) - values (?, ?, ?, ?, ?) - `, repoDid, pipelineId.Knot, pipelineId.Rkey, string(sourceRepoJson(sourceRepo)), string(tplJson)) + insert into jobs (repo_did, pipeline_id_rkey, source_repo, tpl) + values (?, ?, ?, ?) + `, repoDid, pipelineId, string(sourceRepoJson(sourceRepo)), string(tplJson)) return err } func (d *DB) DequeueJob(ctx context.Context) (*JobRow, error) { @@ -39,8 +38,8 @@ func (d *DB) DequeueJob(ctx context.Context) (*JobRow, error) { order by id asc limit 1 ) - returning id, repo_did, pipeline_id_knot, pipeline_id_rkey, source_repo, tpl - `).Scan(&row.Id, &row.RepoDid, &row.PipelineIdKnot, &row.PipelineIdRkey, &sourceRepoStr, &tplJson) + returning id, repo_did, pipeline_id_rkey, source_repo, tpl + `).Scan(&row.Id, &row.RepoDid, &row.PipelineId, &sourceRepoStr, &tplJson) if err != nil { if err == sql.ErrNoRows { return nil, nil diff --git a/spindle/db/mill_state.go b/spindle/db/mill_state.go index fe74c310..8e0cd330 100644 --- a/spindle/db/mill_state.go +++ b/spindle/db/mill_state.go @@ -14,7 +14,6 @@ type MillLease struct { NodeID string Epoch string Engine string - Knot string Rkey string Workflow string State string @@ -41,10 +40,10 @@ type OutboxDeletion struct { func (d *DB) SaveMillLease(l MillLease) error { _, err := d.Exec( `insert into mill_leases ( - lease_id, node_id, epoch, engine, knot, rkey, workflow, state - ) values (?, ?, ?, ?, ?, ?, ?, ?) + lease_id, node_id, epoch, engine, rkey, workflow, state + ) values (?, ?, ?, ?, ?, ?, ?) on conflict(lease_id) do update set state = excluded.state`, - l.LeaseID, l.NodeID, l.Epoch, l.Engine, l.Knot, l.Rkey, l.Workflow, l.State, + l.LeaseID, l.NodeID, l.Epoch, l.Engine, l.Rkey, l.Workflow, l.State, ) return err } @@ -56,7 +55,7 @@ func (d *DB) DeleteMillLease(leaseID string) error { func (d *DB) ListMillLeases() ([]MillLease, error) { rows, err := d.Query(` - select lease_id, node_id, epoch, engine, knot, rkey, workflow, state + select lease_id, node_id, epoch, engine, rkey, workflow, state from mill_leases `) if err != nil { @@ -68,7 +67,7 @@ func (d *DB) ListMillLeases() ([]MillLease, error) { for rows.Next() { var l MillLease if err := rows.Scan( - &l.LeaseID, &l.NodeID, &l.Epoch, &l.Engine, &l.Knot, &l.Rkey, &l.Workflow, &l.State, + &l.LeaseID, &l.NodeID, &l.Epoch, &l.Engine, &l.Rkey, &l.Workflow, &l.State, ); err != nil { return nil, err } diff --git a/spindle/db/mill_state_test.go b/spindle/db/mill_state_test.go index 0d6fe42f..2fc1bd1a 100644 --- a/spindle/db/mill_state_test.go +++ b/spindle/db/mill_state_test.go @@ -18,7 +18,6 @@ func TestMillLeaseRoundTrip(t *testing.T) { NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", - Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: "reserved", @@ -87,7 +86,7 @@ func TestCompleteMillLeaseIsAtomic(t *testing.T) { d := newTestDB(t) lease := MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", - Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: "running", + Rkey: "rkey1", Workflow: "build", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) @@ -107,7 +106,7 @@ func TestCompleteMillLeaseIsAtomic(t *testing.T) { defer n.Unsubscribe(notifications) err := d.CompleteMillLease( "lease-1", - models.WorkflowId{PipelineId: models.PipelineId{Rkey: "rkey1"}, Name: "build"}, + models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"}, "failed", nil, nil, @@ -137,7 +136,7 @@ func TestCompleteMillLeaseIsAtomic(t *testing.T) { } if err := d.CompleteMillLease( "lease-1", - models.WorkflowId{PipelineId: models.PipelineId{Rkey: "rkey1"}, Name: "build"}, + models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"}, "failed", nil, nil, @@ -170,8 +169,7 @@ func TestRestartPersistence(t *testing.T) { } lease := MillLease{ - LeaseID: "lease-p", NodeID: "node-p", Epoch: "inc-p", Engine: "dummy", - Knot: "k", Rkey: "r", Workflow: "w", State: "running", + LeaseID: "lease-p", NodeID: "node-p", Epoch: "inc-p", Engine: "dummy", Rkey: "r", Workflow: "w", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) @@ -304,8 +302,7 @@ func TestBatchRollback(t *testing.T) { d := newTestDB(t) lease := MillLease{ - LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", - Knot: "k", Rkey: "r", Workflow: "w", State: "running", + LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Rkey: "r", Workflow: "w", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) @@ -347,8 +344,7 @@ func TestTerminalCursorAtomicity(t *testing.T) { d := newTestDB(t) lease := MillLease{ - LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", - Knot: "k", Rkey: "r", Workflow: "w", State: "running", + LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Rkey: "r", Workflow: "w", State: "running", } if err := d.SaveMillLease(lease); err != nil { t.Fatalf("SaveMillLease: %v", err) diff --git a/spindle/db/pipelines.go b/spindle/db/pipelines.go index 7ceb5f6c..29c0669f 100644 --- a/spindle/db/pipelines.go +++ b/spindle/db/pipelines.go @@ -91,30 +91,23 @@ func (d *DB) QueryPipelines(ctx context.Context, repoDid string, commits []strin return pipelines, nextCursor, total, nil } -// GetPipelineWithKnot also returns the knot the pipeline was created on, which -// subscribePipelineLogs needs to locate the workflow log file. -func (d *DB) GetPipelineWithKnot(ctx context.Context, rkey string) (*tangled.CiPipeline, string, error) { - var payload, knot string +func (d *DB) GetPipeline(ctx context.Context, rkey models.PipelineId) (*tangled.CiPipeline, error) { + var payload string if err := d.QueryRowContext(ctx, - `select payload, knot from pipelines where rkey = ?`, rkey, - ).Scan(&payload, &knot); err != nil { - return nil, "", err + `select payload from pipelines where rkey = ?`, rkey, + ).Scan(&payload); err != nil { + return nil, err } var p tangled.CiPipeline if err := json.Unmarshal([]byte(payload), &p); err != nil { - return nil, "", err + return nil, err } if err := d.applyStatuses(ctx, []*tangled.CiPipeline{&p}); err != nil { - return nil, "", err + return nil, err } - return &p, knot, nil -} - -func (d *DB) GetPipeline(ctx context.Context, rkey string) (*tangled.CiPipeline, error) { - p, _, err := d.GetPipelineWithKnot(ctx, rkey) - return p, err + return &p, nil } // mapToCiPipeline converts a compiled legacy pipeline into the sh.tangled.ci.pipeline @@ -125,7 +118,7 @@ func (d *DB) GetPipeline(ctx context.Context, rkey string) (*tangled.CiPipeline, // kind is taken straight from raw.TriggerMetadata.Kind, which already holds the // workflow.TriggerKind vocabulary the queryPipelines `kinds` param uses. Producing it // in the same switch as the trigger union keeps the two from disagreeing. -func mapToCiPipeline(rkey string, createdAt time.Time, raw tangled.Pipeline) (*tangled.CiPipeline, workflow.TriggerKind) { +func mapToCiPipeline(rkey models.PipelineId, createdAt time.Time, raw tangled.Pipeline) (*tangled.CiPipeline, workflow.TriggerKind) { createdAtStr := createdAt.Format(time.RFC3339) var repoDidStr string @@ -201,7 +194,7 @@ func mapToCiPipeline(rkey string, createdAt time.Time, raw tangled.Pipeline) (*t } return &tangled.CiPipeline{ - Id: rkey, + Id: string(rkey), Commit: commitSha, Repo: repoDidStr, CreatedAt: &createdAtStr, @@ -228,17 +221,15 @@ func pipelinePairsToCiTriggerPairs(inputs []*tangled.Pipeline_Pair) []*tangled.C return pairs } -// knot is stored only because the workflow log path still embeds it; it is not -// part of pipeline identity and nothing resolves or dials it. func (d *DB) CreatePipeline(id models.PipelineId, raw tangled.Pipeline) error { - p, kind := mapToCiPipeline(id.Rkey, time.Now(), raw) + p, kind := mapToCiPipeline(id, time.Now(), raw) payload, err := json.Marshal(p) if err != nil { return err } _, err = d.Exec( - `insert into pipelines (rkey, knot, repo_did, commit_sha, kind, payload) values (?, ?, ?, ?, ?, ?)`, - id.Rkey, id.Knot, p.Repo, p.Commit, string(kind), string(payload), + `insert into pipelines (rkey, repo_did, commit_sha, kind, payload) values (?, ?, ?, ?, ?)`, + id, p.Repo, p.Commit, string(kind), string(payload), ) return err } diff --git a/spindle/db/pipelines_test.go b/spindle/db/pipelines_test.go index 20274f87..3c6b72cb 100644 --- a/spindle/db/pipelines_test.go +++ b/spindle/db/pipelines_test.go @@ -32,7 +32,7 @@ func seedPipeline(t *testing.T, d *DB, rkey, repoDid, kind string) { TriggerMetadata: tm, Workflows: []*tangled.Pipeline_Workflow{{Name: "ci.yml"}}, } - if err := d.CreatePipeline(models.PipelineId{Knot: "knot.test", Rkey: rkey}, raw); err != nil { + if err := d.CreatePipeline(models.PipelineId(rkey), raw); err != nil { t.Fatalf("seed pipeline %s: %v", rkey, err) } } @@ -125,12 +125,12 @@ func TestQueryPipelines_WorkflowStatuses(t *testing.T) { }, Workflows: []*tangled.Pipeline_Workflow{{Name: "a"}, {Name: "b"}}, } - if err := d.CreatePipeline(models.PipelineId{Knot: "knot.test", Rkey: "pl1"}, raw); err != nil { + if err := d.CreatePipeline(models.PipelineId("pl1"), raw); err != nil { t.Fatalf("CreatePipeline: %v", err) } - widA := models.WorkflowId{PipelineId: models.PipelineId{Rkey: "pl1"}, Name: "a"} - widB := models.WorkflowId{PipelineId: models.PipelineId{Rkey: "pl1"}, Name: "b"} + widA := models.WorkflowId{PipelineId: models.PipelineId("pl1"), Name: "a"} + widB := models.WorkflowId{PipelineId: models.PipelineId("pl1"), Name: "b"} for _, step := range []func() error{ func() error { return d.StatusPending(widA, &n) }, diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go index e011fea9..77c8d22f 100644 --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -103,7 +103,7 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s for _, w := range wfs { wid := models.WorkflowId{ PipelineId: pipelineId, - Name: w.Name, + Name: w.Name, } wfCounts[wid.String()]++ } @@ -116,7 +116,7 @@ func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, s for _, w := range wfs { wid := models.WorkflowId{ PipelineId: pipelineId, - Name: w.Name, + Name: w.Name, } if wfCounts[wid.String()] > 1 { diff --git a/spindle/engine/engine_test.go b/spindle/engine/engine_test.go index dedcbb68..bea8bc00 100644 --- a/spindle/engine/engine_test.go +++ b/spindle/engine/engine_test.go @@ -88,10 +88,7 @@ func TestStartWorkflows_CollisionRejection(t *testing.T) { logger := slog.New(slog.NewTextHandler(os.Stderr, nil)) eng := &mockEngine{} - pipelineId := models.PipelineId{ - Knot: "test-knot", - Rkey: "test-rkey", - } + pipelineId := models.PipelineId("test-rkey") // two names that normalize to the same wid must not both run wfColliding1 := models.Workflow{ @@ -171,14 +168,11 @@ func TestCancelWorkflow_NotOverwritten(t *testing.T) { }, } - pipelineId := models.PipelineId{ - Knot: "test-knot", - Rkey: "test-rkey", - } + pipelineId := models.PipelineId("test-rkey") wid := models.WorkflowId{ PipelineId: pipelineId, - Name: "cancel_test_job", + Name: "cancel_test_job", } pipeline := &models.Pipeline{ @@ -240,7 +234,7 @@ func TestSetupTimeout_ReportsTimeout(t *testing.T) { }, } - pipelineId := models.PipelineId{Knot: "test-knot", Rkey: "test-rkey"} + pipelineId := models.PipelineId("test-rkey") wid := models.WorkflowId{PipelineId: pipelineId, Name: "timeout_job"} pipeline := &models.Pipeline{ diff --git a/spindle/mill/auth_test.go b/spindle/mill/auth_test.go index 218e6f16..2d9549da 100644 --- a/spindle/mill/auth_test.go +++ b/spindle/mill/auth_test.go @@ -249,9 +249,9 @@ func TestOnStatusEventOwnership(t *testing.T) { m.Attach(bdb, &n) foreign := newLease("lease-foreign", "node-x", "inc-x", "dummy") - foreign.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "foreign"}, Name: "build"} + foreign.wid = models.WorkflowId{PipelineId: models.PipelineId("foreign"), Name: "build"} owned := newLease("lease-owned", "node-z", "inc-z", "dummy") - owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "owned"}, Name: "build"} + owned.wid = models.WorkflowId{PipelineId: models.PipelineId("owned"), Name: "build"} m.mu.Lock() m.leases[foreign.id] = foreign m.leases[owned.id] = owned diff --git a/spindle/mill/executor/capability_test.go b/spindle/mill/executor/capability_test.go index a9b4cde7..9cba0d1d 100644 --- a/spindle/mill/executor/capability_test.go +++ b/spindle/mill/executor/capability_test.go @@ -49,7 +49,6 @@ func TestHandleReserveRejectsMissingTriggerMetadata(t *testing.T) { TargetEngine: "microvm", RawWorkflowJson: string(twf), RawPipelineJson: string(tpl), - Knot: "k", Rkey: "r", }) @@ -93,7 +92,6 @@ func TestHandleReserveValidatesPlacementBeforeAcquiringSlot(t *testing.T) { TargetEngine: "microvm", RawWorkflowJson: string(twf), RawPipelineJson: string(tpl), - Knot: "k", Rkey: "r", }) diff --git a/spindle/mill/executor/executor.go b/spindle/mill/executor/executor.go index 014e740a..5351c69a 100644 --- a/spindle/mill/executor/executor.go +++ b/spindle/mill/executor/executor.go @@ -357,7 +357,7 @@ func (e *Executor) handleReserve(ctx context.Context, rs *millv1.ReserveSeat) { return } - pipelineId := models.PipelineId{Knot: rs.GetKnot(), Rkey: rs.GetRkey()} + pipelineId := models.PipelineId(rs.GetRkey()) wid := models.WorkflowId{PipelineId: pipelineId, Name: twf.Name} wf, err := realEngine.InitWorkflow(twf, tpl) diff --git a/spindle/mill/executor/observe.go b/spindle/mill/executor/observe.go index df9ed09b..e4e06232 100644 --- a/spindle/mill/executor/observe.go +++ b/spindle/mill/executor/observe.go @@ -151,7 +151,7 @@ func (e *Executor) reservationFor(rkey, workflow string) *reservation { e.mu.Lock() defer e.mu.Unlock() for _, res := range e.active { - if res.wid.PipelineId.Rkey == rkey && res.wid.Name == workflow { + if res.wid.PipelineId == models.PipelineId(rkey) && res.wid.Name == workflow { return res } } diff --git a/spindle/mill/executor/reserved_test.go b/spindle/mill/executor/reserved_test.go index d13da68e..e5a9b57d 100644 --- a/spindle/mill/executor/reserved_test.go +++ b/spindle/mill/executor/reserved_test.go @@ -206,7 +206,7 @@ func TestHandleCommitPreservesPreauthorizedSecrets(t *testing.T) { } res := &reservation{ leaseID: "lease-1", - wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, + wid: models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"}, realEngine: inner, slot: slot, wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}}, @@ -315,7 +315,7 @@ func TestFinishJobReportsCancelledReservationAsCancelled(t *testing.T) { d := testDB(t) res := &reservation{ leaseID: "lease-1", - wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, + wid: models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"}, cancelled: true, } e := &Executor{ @@ -329,7 +329,7 @@ func TestFinishJobReportsCancelledReservationAsCancelled(t *testing.T) { } e.finishJob(res, &db.StatusRow{ - Pipeline: res.wid.PipelineId.Rkey, + Pipeline: string(res.wid.PipelineId), Workflow: res.wid.Name, Status: string(models.StatusKindFailed), }) @@ -397,7 +397,7 @@ func TestSocketCancellationIndependence(t *testing.T) { } res := &reservation{ leaseID: "lease-1", - wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, + wid: models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"}, realEngine: inner, slot: slot, wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}}, @@ -522,7 +522,7 @@ func TestStructuredShutdown(t *testing.T) { } res := &reservation{ leaseID: "lease-1", - wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, + wid: models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"}, realEngine: inner, slot: slot, wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}}, diff --git a/spindle/mill/integration_test.go b/spindle/mill/integration_test.go index 553706b2..8e92ef07 100644 --- a/spindle/mill/integration_test.go +++ b/spindle/mill/integration_test.go @@ -75,7 +75,7 @@ func TestEndToEndDummyJob(t *testing.T) { if err != nil { t.Fatalf("InitWorkflow: %v", err) } - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.test", Rkey: "rkey1"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"} placeCtx, placeCancel := context.WithTimeout(ctx, 10*time.Second) defer placeCancel() @@ -248,7 +248,7 @@ func TestEndToEndDummyJobUsesRequiredLabelsAcrossExecutors(t *testing.T) { if err != nil { t.Fatalf("InitWorkflow: %v", err) } - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.test", Rkey: "rkey-arm"}, Name: "build-arm"} + wid := models.WorkflowId{PipelineId: models.PipelineId("rkey-arm"), Name: "build-arm"} placeCtx, placeCancel := context.WithTimeout(ctx, 10*time.Second) defer placeCancel() @@ -300,7 +300,7 @@ func waitForStatus(t *testing.T, d *db.DB, wid models.WorkflowId, want string) b t.Fatalf("GetEvents: %v", err) } for _, ev := range evs { - if ev.Pipeline == wid.PipelineId.Rkey && ev.Workflow == wid.Name && ev.Status == want { + if ev.Pipeline == string(wid.PipelineId) && ev.Workflow == wid.Name && ev.Status == want { return true } } diff --git a/spindle/mill/mill.go b/spindle/mill/mill.go index 319e55a9..7bd12d78 100644 --- a/spindle/mill/mill.go +++ b/spindle/mill/mill.go @@ -429,8 +429,7 @@ func (m *Mill) bid(ctx context.Context, engineName string, wid models.WorkflowId TargetEngine: engineName, RawPipelineJson: rawPipeline, RawWorkflowJson: rawWorkflow, - Knot: wid.PipelineId.Knot, - Rkey: wid.PipelineId.Rkey, + Rkey: string(wid.PipelineId), TtlSeconds: uint32(m.cfg.ReconnectGrace / time.Second), }} resp, err := sess.request(bidCtx, leaseID, msg) diff --git a/spindle/mill/mill_test.go b/spindle/mill/mill_test.go index e4a44461..791644f4 100644 --- a/spindle/mill/mill_test.go +++ b/spindle/mill/mill_test.go @@ -139,7 +139,7 @@ func TestCommitRetriesAfterSessionCloseBeforeCommitted(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{BidTimeout: 25 * time.Millisecond, ReconnectGrace: time.Second}) wf := testWorkflow("build") - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = wid wf.Data.(*millWorkflowState).Lease = lease @@ -209,7 +209,7 @@ func TestCommitRetriesAfterSessionCloseBeforeCommitted(t *testing.T) { func TestDestroyRunningLeaseDoesNotDropCancelledTerminal(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{}) - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = wid lease.setState(leaseRunning) @@ -236,7 +236,7 @@ func TestPlaceBlocksWhenNoCapacity(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{}) wf := testWorkflow("build") - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} // with no executors at all, place must block until ctx expires and the user sees pending ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond) @@ -307,7 +307,7 @@ func TestPlaceWithMissingRequiredLabelsStaysPendingWithoutReserve(t *testing.T) return nil })) wf := testWorkflowWithRunsOn("build", []string{"linux", "arm64"}) - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} ctx, cancel := context.WithTimeout(context.Background(), 120*time.Millisecond) defer cancel() @@ -339,7 +339,7 @@ func TestMaxPendingRejects(t *testing.T) { func TestCancelledRunningLeaseSurvivesReleaseForReconnectReplay(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = wid lease.setState(leaseRunning) @@ -500,7 +500,7 @@ func TestPlaceReleasesRemoteReservationWhenInitialPersistenceFails(t *testing.T) slot, err := m.place( ctx, "dummy", - models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, + models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"}, testWorkflow("build"), ) if err == nil { @@ -531,7 +531,7 @@ func TestGapsAndDuplicates(t *testing.T) { m.attachSession(sess) owned := newLease("lease-1", "node-1", "inc-1", "dummy") - owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + owned.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock() @@ -575,7 +575,7 @@ func TestAtomicBatchRollback(t *testing.T) { m.attachSession(sess) owned := newLease("lease-1", "node-1", "inc-1", "dummy") - owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + owned.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock() @@ -621,7 +621,7 @@ func TestTerminalBeforeACK(t *testing.T) { ackSent := make(chan struct{}) sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error { if msg.GetAck() != nil { - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} st, err := bdb.GetStatus(wid) if err != nil || st != "success" { t.Errorf("expected terminal status success at ACK time, got status: %v, err: %v", st, err) @@ -633,7 +633,7 @@ func TestTerminalBeforeACK(t *testing.T) { m.attachSession(sess) owned := newLease("lease-1", "node-1", "inc-1", "dummy") - owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + owned.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock() @@ -665,7 +665,7 @@ func TestExecutorRestartEmptySnapshot(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) lease := newLease("lease-1", "node-1", "inc-old", "dummy") - lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock() @@ -695,7 +695,7 @@ func TestExecutorRestartEmptySnapshot(t *testing.T) { func TestReplacementLostBeforeSnapshotFailsOldEpochLease(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) lease := newLease("lease-1", "node-1", "inc-old", "dummy") - lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) @@ -743,7 +743,7 @@ func TestClaimedSweep(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) lease := newLease("lease-1", "node-1", "inc-1", "dummy") - lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} lease.orphaned = true lease.claimed = false m.mu.Lock() @@ -786,7 +786,7 @@ func TestCancelDeadline(t *testing.T) { m.attachSession(sess) lease := newLease("lease-1", "node-1", "inc-1", "dummy") - lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} lease.setState(leaseRunning) m.mu.Lock() m.leases[lease.id] = lease @@ -805,7 +805,7 @@ func TestCleanupRetry(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) lease := newLease("lease-1", "node-1", "inc-1", "dummy") - lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + lease.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persist lease: %v", err) } diff --git a/spindle/mill/proto/gen/mill.pb.go b/spindle/mill/proto/gen/mill.pb.go index e31758f2..1ddef549 100644 --- a/spindle/mill/proto/gen/mill.pb.go +++ b/spindle/mill/proto/gen/mill.pb.go @@ -419,12 +419,10 @@ type ReserveSeat struct { TargetEngine string `protobuf:"bytes,2,opt,name=target_engine,json=targetEngine,proto3" json:"target_engine,omitempty"` RawPipelineJson string `protobuf:"bytes,3,opt,name=raw_pipeline_json,json=rawPipelineJson,proto3" json:"raw_pipeline_json,omitempty"` RawWorkflowJson string `protobuf:"bytes,4,opt,name=raw_workflow_json,json=rawWorkflowJson,proto3" json:"raw_workflow_json,omitempty"` - // pipeline id, the executor reconstructs the exact WorkflowId from it - Knot string `protobuf:"bytes,5,opt,name=knot,proto3" json:"knot,omitempty"` - Rkey string `protobuf:"bytes,6,opt,name=rkey,proto3" json:"rkey,omitempty"` - TtlSeconds uint32 `protobuf:"varint,7,opt,name=ttl_seconds,json=ttlSeconds,proto3" json:"ttl_seconds,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + Rkey string `protobuf:"bytes,6,opt,name=rkey,proto3" json:"rkey,omitempty"` + TtlSeconds uint32 `protobuf:"varint,7,opt,name=ttl_seconds,json=ttlSeconds,proto3" json:"ttl_seconds,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *ReserveSeat) Reset() { @@ -485,13 +483,6 @@ func (x *ReserveSeat) GetRawWorkflowJson() string { return "" } -func (x *ReserveSeat) GetKnot() string { - if x != nil { - return x.Knot - } - return "" -} - func (x *ReserveSeat) GetRkey() string { if x != nil { return x.Rkey @@ -1470,16 +1461,15 @@ const file_spindle_mill_v1_mill_proto_rawDesc = "" + "\x10active_lease_ids\x18\x03 \x03(\tR\x0eactiveLeaseIds\x1a_\n" + "\fEnginesEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x129\n" + - "\x05value\x18\x02 \x01(\v2#.spindle.mill.v1.EngineAvailabilityR\x05value:\x028\x01\"\x80\x02\n" + + "\x05value\x18\x02 \x01(\v2#.spindle.mill.v1.EngineAvailabilityR\x05value:\x028\x01\"\xf2\x01\n" + "\vReserveSeat\x12\"\n" + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12,\n" + "\rtarget_engine\x18\x02 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\ftargetEngine\x12*\n" + "\x11raw_pipeline_json\x18\x03 \x01(\tR\x0frawPipelineJson\x12*\n" + "\x11raw_workflow_json\x18\x04 \x01(\tR\x0frawWorkflowJson\x12\x12\n" + - "\x04knot\x18\x05 \x01(\tR\x04knot\x12\x12\n" + "\x04rkey\x18\x06 \x01(\tR\x04rkey\x12\x1f\n" + "\vttl_seconds\x18\a \x01(\rR\n" + - "ttlSeconds\"\xbf\x01\n" + + "ttlSecondsJ\x04\b\x05\x10\x06\"\xbf\x01\n" + "\rReserveResult\x12\"\n" + "\blease_id\x18\x01 \x01(\tB\a\xbaH\x04r\x02\x10\x01R\aleaseId\x12\x1a\n" + "\baccepted\x18\x02 \x01(\bR\baccepted\x12#\n" + diff --git a/spindle/mill/proto/protocol_test.go b/spindle/mill/proto/protocol_test.go index 0e0b9637..683bac0a 100644 --- a/spindle/mill/proto/protocol_test.go +++ b/spindle/mill/proto/protocol_test.go @@ -17,7 +17,6 @@ func TestEncodeDecodeRoundTrip(t *testing.T) { LeaseId: "lease-1", TargetEngine: "microvm", RawWorkflowJson: `{"name":"build"}`, - Knot: "knot.example", Rkey: "abc123", TtlSeconds: 30, }, diff --git a/spindle/mill/proto/spindle/mill/v1/mill.proto b/spindle/mill/proto/spindle/mill/v1/mill.proto index df71b493..e580bf95 100644 --- a/spindle/mill/proto/spindle/mill/v1/mill.proto +++ b/spindle/mill/proto/spindle/mill/v1/mill.proto @@ -43,8 +43,9 @@ message ReserveSeat { string target_engine = 2 [(buf.validate.field).string.min_len = 1]; string raw_pipeline_json = 3; string raw_workflow_json = 4; - // pipeline id, the executor reconstructs the exact WorkflowId from it - string knot = 5; + // pipeline id, the executor reconstructs the exact WorkflowId from it. + // 5 was `knot`, dropped when knot left pipeline identity. + reserved 5; string rkey = 6; uint32 ttl_seconds = 7; } diff --git a/spindle/mill/restore.go b/spindle/mill/restore.go index d5470072..bcea9b5b 100644 --- a/spindle/mill/restore.go +++ b/spindle/mill/restore.go @@ -23,8 +23,7 @@ func (m *Mill) persistLease(lease *RemoteLease, state string) error { NodeID: lease.nodeID, Epoch: lease.epoch, Engine: lease.engine, - Knot: lease.wid.PipelineId.Knot, - Rkey: lease.wid.PipelineId.Rkey, + Rkey: string(lease.wid.PipelineId), Workflow: lease.wid.Name, State: state, }) @@ -50,8 +49,8 @@ func (m *Mill) RestoreState() error { for _, r := range rows { lease := newLease(r.LeaseID, r.NodeID, r.Epoch, r.Engine) lease.wid = models.WorkflowId{ - PipelineId: models.PipelineId{Knot: r.Knot, Rkey: r.Rkey}, - Name: r.Workflow, + PipelineId: models.PipelineId(r.Rkey), + Name: r.Workflow, } // restored leases start as orphans, an executor must reclaim it via // its first snapshot, or the sweep will fail it diff --git a/spindle/mill/restore_test.go b/spindle/mill/restore_test.go index 806dc8dd..d8a217bb 100644 --- a/spindle/mill/restore_test.go +++ b/spindle/mill/restore_test.go @@ -42,7 +42,7 @@ func TestRestoreStateRebuildsLeasesAndCursors(t *testing.T) { if err := bdb.SaveMillLease(db.MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", - Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, + Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, }); err != nil { t.Fatalf("SaveMillLease: %v", err) } @@ -66,7 +66,7 @@ func TestRestoreStateRebuildsLeasesAndCursors(t *testing.T) { if lease.getState() != leaseRunning { t.Fatalf("restored lease state = %v, want leaseRunning", lease.getState()) } - wantWid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"} + wantWid := models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"} if lease.wid != wantWid { t.Fatalf("restored lease wid = %+v, want %+v", lease.wid, wantWid) } @@ -79,7 +79,7 @@ func TestOrphanTerminalAuthorsStatusRow(t *testing.T) { _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) if err := bdb.SaveMillLease(db.MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", - Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, + Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, }); err != nil { t.Fatalf("SaveMillLease: %v", err) } @@ -113,7 +113,7 @@ func TestOrphanTerminalAuthorsStatusRow(t *testing.T) { }, }) - wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "rkey1"}, Name: "build"} + wid := models.WorkflowId{PipelineId: models.PipelineId("rkey1"), Name: "build"} st, err := bdb.GetStatus(wid) if err != nil { t.Fatalf("GetStatus after orphan terminal: %v", err) @@ -136,8 +136,8 @@ func TestOrphanTerminalAuthorsStatusRow(t *testing.T) { func TestSnapshotReconciliationFailsDroppedOrphans(t *testing.T) { _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) for _, l := range []db.MillLease{ - {LeaseID: "lease-kept", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowRunning}, - {LeaseID: "lease-gone", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Knot: "k", Rkey: "r2", Workflow: "w", State: leaseRowRunning}, + {LeaseID: "lease-kept", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Rkey: "r1", Workflow: "w", State: leaseRowRunning}, + {LeaseID: "lease-gone", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", Rkey: "r2", Workflow: "w", State: leaseRowRunning}, } { if err := bdb.SaveMillLease(l); err != nil { t.Fatalf("SaveMillLease(%s): %v", l.LeaseID, err) @@ -165,7 +165,7 @@ func TestSnapshotReconciliationFailsDroppedOrphans(t *testing.T) { t.Fatal("reconciliation kept a lease the executor no longer holds") } - st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r2"}, Name: "w"}) + st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId("r2"), Name: "w"}) if err != nil { t.Fatalf("GetStatus for dropped orphan: %v", err) } @@ -177,8 +177,8 @@ func TestSnapshotReconciliationPreservesRequestedCancellation(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) lease := newLease("lease-1", "node-1", "inc-old", "dummy") lease.wid = models.WorkflowId{ - PipelineId: models.PipelineId{Knot: "k", Rkey: "r1"}, - Name: "w", + PipelineId: models.PipelineId("r1"), + Name: "w", } lease.setState(leaseRunning) lease.requestCancel("workflow destroyed") @@ -235,7 +235,7 @@ func TestSweepFailsOrphansOfAbsentExecutors(t *testing.T) { _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) if err := bdb.SaveMillLease(db.MillLease{ LeaseID: "lease-1", NodeID: "node-absent", Epoch: "inc-absent", Engine: "dummy", - Knot: "k", Rkey: "r1", Workflow: "w", State: leaseRowReserved, + Rkey: "r1", Workflow: "w", State: leaseRowReserved, }); err != nil { t.Fatalf("SaveMillLease: %v", err) } @@ -249,7 +249,7 @@ func TestSweepFailsOrphansOfAbsentExecutors(t *testing.T) { if still { t.Fatal("sweep kept an orphan whose executor never reconnected") } - st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r1"}, Name: "w"}) + st, err := bdb.GetStatus(models.WorkflowId{PipelineId: models.PipelineId("r1"), Name: "w"}) if err != nil { t.Fatalf("GetStatus after sweep: %v", err) } @@ -267,7 +267,7 @@ func TestAckSeqnoPersistsCursor(t *testing.T) { m.attachSession(sess) owned := newLease("lease-1", "node-1", "inc-1", "dummy") - owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} + owned.wid = models.WorkflowId{PipelineId: models.PipelineId("r"), Name: "build"} m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock() @@ -303,7 +303,7 @@ func TestOrphanTerminalFailureKeepsLeaseAndSeqnoRetryable(t *testing.T) { _, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) if err := bdb.SaveMillLease(db.MillLease{ LeaseID: "lease-1", NodeID: "node-1", Epoch: "inc-1", Engine: "dummy", - Knot: "knot.example", Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, + Rkey: "rkey1", Workflow: "build", State: leaseRowRunning, }); err != nil { t.Fatalf("SaveMillLease: %v", err) } diff --git a/spindle/models/logger_test.go b/spindle/models/logger_test.go index 79cec1ef..93332018 100644 --- a/spindle/models/logger_test.go +++ b/spindle/models/logger_test.go @@ -9,7 +9,7 @@ import ( ) func testWorkflowId(name string) WorkflowId { - return WorkflowId{PipelineId: PipelineId{Knot: "knot1", Rkey: "rkey1"}, Name: name} + return WorkflowId{PipelineId: PipelineId("rkey1"), Name: name} } func readDataContents(t *testing.T, path string) []string { diff --git a/spindle/models/models.go b/spindle/models/models.go index 77bbbb02..d5b27925 100644 --- a/spindle/models/models.go +++ b/spindle/models/models.go @@ -6,8 +6,6 @@ import ( "slices" "time" - "tangled.org/core/api/tangled" - "github.com/bluesky-social/indigo/atproto/syntax" ) @@ -15,14 +13,9 @@ var ( re = regexp.MustCompile(`[^a-zA-Z0-9_.-]`) ) -type PipelineId struct { - Knot string - Rkey string -} +type PipelineId syntax.RecordKey -func (p *PipelineId) AtUri() syntax.ATURI { - return syntax.ATURI(fmt.Sprintf("at://did:web:%s/%s/%s", p.Knot, tangled.PipelineNSID, p.Rkey)) -} +func (p PipelineId) String() string { return string(p) } type WorkflowId struct { PipelineId @@ -30,7 +23,7 @@ type WorkflowId struct { } func (wid WorkflowId) String() string { - return fmt.Sprintf("%s-%s-%s", normalize(wid.PipelineId.Knot), wid.PipelineId.Rkey, normalize(wid.Name)) + return fmt.Sprintf("%s-%s", normalize(string(wid.PipelineId)), normalize(wid.Name)) } func normalize(name string) string { diff --git a/spindle/models/pipeline_env.go b/spindle/models/pipeline_env.go index 23534395..9935f54d 100644 --- a/spindle/models/pipeline_env.go +++ b/spindle/models/pipeline_env.go @@ -23,7 +23,7 @@ func PipelineEnvVarsForSource(tr *tangled.Pipeline_TriggerMetadata, pipelineId P // standard CI env vars env["CI"] = "true" - env["TANGLED_PIPELINE_ID"] = pipelineId.AtUri().String() + env["TANGLED_PIPELINE_ID"] = pipelineId.String() env["TANGLED_PIPELINE_KIND"] = tr.Kind if tr.SourceRepo != nil && *tr.SourceRepo != "" { diff --git a/spindle/models/pipeline_env_test.go b/spindle/models/pipeline_env_test.go index 54fc5f5e..fd57b69d 100644 --- a/spindle/models/pipeline_env_test.go +++ b/spindle/models/pipeline_env_test.go @@ -23,10 +23,7 @@ func TestPipelineEnvVars_PushBranch(t *testing.T) { DefaultBranch: "main", }, } - id := PipelineId{ - Knot: "example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVars(tr, id) // Check standard CI variable @@ -86,10 +83,7 @@ func TestPipelineEnvVars_PushTag(t *testing.T) { RepoDid: sp("did:plc:boltless"), }, } - id := PipelineId{ - Knot: "example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVars(tr, id) if env["TANGLED_REF"] != "refs/tags/v1.2.3" { @@ -118,10 +112,7 @@ func TestPipelineEnvVars_PullRequest(t *testing.T) { RepoDid: sp("did:plc:boltless"), }, } - id := PipelineId{ - Knot: "example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVars(tr, id) // Check ref variables for PR @@ -185,10 +176,7 @@ func TestPipelineEnvVars_SourceRepo(t *testing.T) { RepoDid: &sourceRepoDid, DefaultBranch: "feature-branch", } - id := PipelineId{ - Knot: "target.example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVarsForSource(tr, id, sourceRepo) @@ -220,10 +208,7 @@ func TestPipelineEnvVars_ManualWithInputs(t *testing.T) { RepoDid: sp("did:plc:boltless"), }, } - id := PipelineId{ - Knot: "example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVars(tr, id) // Check manual input variables @@ -261,10 +246,7 @@ func TestPipelineEnvVars_DevMode(t *testing.T) { RepoDid: sp("did:plc:boltless"), }, } - id := PipelineId{ - Knot: "example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVars(tr, id) expectedURL := "http://localhost:3000/did:plc:boltless" @@ -274,10 +256,7 @@ func TestPipelineEnvVars_DevMode(t *testing.T) { } func TestPipelineEnvVars_NilTrigger(t *testing.T) { - id := PipelineId{ - Knot: "example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVars(nil, id) if env != nil { @@ -296,10 +275,7 @@ func TestPipelineEnvVars_NilPushData(t *testing.T) { RepoDid: sp("did:plc:boltless"), }, } - id := PipelineId{ - Knot: "example.com", - Rkey: "123123", - } + id := PipelineId("123123") env := PipelineEnvVars(tr, id) // Should still have repo variables diff --git a/spindle/server.go b/spindle/server.go index 9a699aa8..3bc98600 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -556,11 +556,11 @@ func (s *Spindle) processKnotStream(ctx context.Context, src eventconsumer.Sourc if err != nil { return err } - if pipelineId.Rkey == "" { + if pipelineId == "" { l.Info("no workflow matched 'push' trigger, skipping the event") return nil } - l.Info("pipeline triggered", "pipeline", pipelineId.AtUri()) + l.Info("pipeline triggered", "pipeline", pipelineId) } return nil @@ -699,10 +699,10 @@ func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger t rawPipeline, err := s.loadPipeline(ctx, repoCloneUri, repoPath, rev) if err != nil { - return models.PipelineId{}, fmt.Errorf("loading pipeline: %w", err) + return "", fmt.Errorf("loading pipeline: %w", err) } if len(rawPipeline) == 0 { - return models.PipelineId{}, nil + return "", nil } tpl := compiler.Compile(compiler.Parse(rawPipeline)) @@ -718,15 +718,12 @@ func (s *Spindle) runPipeline(ctx context.Context, repoDid syntax.DID, trigger t tpl.Workflows = filterWorkflows(tpl.Workflows, only) } if len(tpl.Workflows) == 0 { - return models.PipelineId{}, nil + return "", nil } - pipelineId := models.PipelineId{ - Knot: trigger.Repo.Knot, - Rkey: tid.TID(), - } + pipelineId := models.PipelineId(tid.TID()) if err := s.db.CreatePipeline(pipelineId, tpl); err != nil { - return models.PipelineId{}, fmt.Errorf("creating pipeline: %w", err) + return "", fmt.Errorf("creating pipeline: %w", err) } err = s.processPipeline(repoDid, tpl, pipelineId, sourceRepo) return pipelineId, err @@ -752,7 +749,7 @@ func filterWorkflows(workflows []*tangled.Pipeline_Workflow, only []string) []*t // TriggerManual dispatches a pipeline at sha, authorized against and recorded // under repoDid. sourceRepo, pull, and inputs are optional trigger payload. -func (s *Spindle) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull xrpc.PullContext, inputs []*tangled.Pipeline_Pair) (syntax.ATURI, error) { +func (s *Spindle) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull xrpc.PullContext, inputs []*tangled.Pipeline_Pair) (models.PipelineId, error) { repo, err := s.db.GetRepoByDid(repoDid) if err != nil { return "", fmt.Errorf("unknown repoDid %s: %w", repoDid, err) @@ -805,10 +802,10 @@ func (s *Spindle) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, re if err != nil { return "", err } - if pipelineId.Rkey == "" { + if pipelineId == "" { return "", xrpc.ErrNoMatchingWorkflows } - return pipelineId.AtUri(), nil + return pipelineId, nil } // sourceInfo is nil when the checkout comes from the target repo. @@ -996,10 +993,7 @@ func (s *Spindle) StartJobWorkers(ctx context.Context) { } func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) { - pipelineId := models.PipelineId{ - Knot: job.PipelineIdKnot, - Rkey: job.PipelineIdRkey, - } + pipelineId := job.PipelineId pipelineEnv := models.PipelineEnvVarsForSource(job.Tpl.TriggerMetadata, pipelineId, job.SourceRepo) trustedSource := true @@ -1024,7 +1018,7 @@ func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) { if !ok { _ = s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, - Name: w.Name, + Name: w.Name, }, fmt.Sprintf("unknown engine %#v", w.Engine), -1, s.n) continue } @@ -1033,7 +1027,7 @@ func (s *Spindle) runJob(ctx context.Context, job *db.JobRow) { if err != nil { _ = s.db.StatusFailed(models.WorkflowId{ PipelineId: pipelineId, - Name: w.Name, + Name: w.Name, }, fmt.Sprintf("init workflow: %s", err), -1, s.n) continue } @@ -1073,7 +1067,7 @@ func (s *Spindle) processPipeline(repoDid syntax.DID, tpl tangled.Pipeline, pipe } if err := s.db.StatusPending(models.WorkflowId{ PipelineId: pipelineId, - Name: w.Name, + Name: w.Name, }, s.n); err != nil { return fmt.Errorf("db.StatusPending: %w", err) } diff --git a/spindle/tapclient.go b/spindle/tapclient.go index 4bfeb52d..799df2bf 100644 --- a/spindle/tapclient.go +++ b/spindle/tapclient.go @@ -597,10 +597,7 @@ func (s *Spindle) triggerPullRequestPipeline(ctx context.Context, l *slog.Logger return nil } - pipelineId := models.PipelineId{ - Knot: tpl.TriggerMetadata.Repo.Knot, - Rkey: tid.TID(), - } + pipelineId := models.PipelineId(tid.TID()) if err := s.db.CreatePipeline(pipelineId, tpl); err != nil { l.Error("failed to create pipeline event", "err", err) return nil diff --git a/spindle/xrpc/ci_pipeline_subscribe_logs.go b/spindle/xrpc/ci_pipeline_subscribe_logs.go index 3627b645..780719f4 100644 --- a/spindle/xrpc/ci_pipeline_subscribe_logs.go +++ b/spindle/xrpc/ci_pipeline_subscribe_logs.go @@ -9,7 +9,6 @@ import ( "time" "github.com/bluesky-social/indigo/atproto/atclient" - "github.com/bluesky-social/indigo/atproto/syntax" "github.com/gorilla/websocket" "tangled.org/core/api/tangled" "tangled.org/core/spindle/logview" @@ -22,13 +21,14 @@ func (x *Xrpc) HandleCiSubscribePipelineLogs(w http.ResponseWriter, r *http.Requ workflows = r.URL.Query()["workflows"] ) - pipeline, err := syntax.ParseTID(pipelineQuery) - if err != nil { - writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: fmt.Sprintf("pipeline parameter invalid: %s", pipelineQuery)}) + // pipeline ids are opaque, any format. WorkflowId.String() normalizes before + // it reaches a filesystem path, so only emptiness is rejected here. + if pipelineQuery == "" { + writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: "pipeline parameter is required"}) return } - x.handleSubscribeLogs(w, r, pipeline, workflows) + x.handleSubscribeLogs(w, r, models.PipelineId(pipelineQuery), workflows) } var wsUpgrader = websocket.Upgrader{ @@ -36,11 +36,11 @@ var wsUpgrader = websocket.Upgrader{ WriteBufferSize: 10_000, } -func (x *Xrpc) handleSubscribeLogs(w http.ResponseWriter, r *http.Request, pipeline syntax.TID, workflows []string) { +func (x *Xrpc) handleSubscribeLogs(w http.ResponseWriter, r *http.Request, pipeline models.PipelineId, workflows []string) { l := x.Logger.With("pipeline", pipeline, "workflows", workflows) - // 1. query the pipeline from database to get the knot. knot is used to locate the workflow log files. - tpl, knot, err := x.Db.GetPipelineWithKnot(r.Context(), pipeline.String()) + // 1. the pipeline supplies the workflow names to stream + tpl, err := x.Db.GetPipeline(r.Context(), pipeline) if err != nil { l.Error("failed to find pipeline event", "err", err) writeJson(w, http.StatusNotFound, atclient.ErrorBody{Name: "NotFound", Message: fmt.Sprintf("pipeline not found: %s", pipeline.String())}) @@ -132,10 +132,7 @@ func (x *Xrpc) handleSubscribeLogs(w http.ResponseWriter, r *http.Request, pipel defer wg.Done() wid := models.WorkflowId{ - PipelineId: models.PipelineId{ - Knot: knot, - Rkey: pipeline.String(), - }, + PipelineId: pipeline, Name: wfName, } diff --git a/spindle/xrpc/ci_pipeline_trigger_pipeline.go b/spindle/xrpc/ci_pipeline_trigger_pipeline.go index 7cf1a74e..a4430486 100644 --- a/spindle/xrpc/ci_pipeline_trigger_pipeline.go +++ b/spindle/xrpc/ci_pipeline_trigger_pipeline.go @@ -110,7 +110,7 @@ func (x *Xrpc) TriggerPipeline(w http.ResponseWriter, r *http.Request) { return } - pipelineAt, err := x.Trigger.TriggerManual(r.Context(), repoDid, sha, ref, input.Workflows, sourceRepo, pull, inputs) + pipelineId, err := x.Trigger.TriggerManual(r.Context(), repoDid, sha, ref, input.Workflows, sourceRepo, pull, inputs) if errors.Is(err, ErrNoMatchingWorkflows) { fail(xrpcerr.GenericError(err)) return @@ -121,7 +121,7 @@ func (x *Xrpc) TriggerPipeline(w http.ResponseWriter, r *http.Request) { } if err := writeJson(w, http.StatusOK, tangled.CiTriggerPipeline_Output{ - Pipeline: pipelineAt.String(), + Pipeline: pipelineId.String(), }); err != nil { l.Error("failed to write response", "err", err) } diff --git a/spindle/xrpc/ci_query_pipelines.go b/spindle/xrpc/ci_query_pipelines.go index 2c0b71e9..d0337447 100644 --- a/spindle/xrpc/ci_query_pipelines.go +++ b/spindle/xrpc/ci_query_pipelines.go @@ -6,6 +6,7 @@ import ( "strconv" "tangled.org/core/api/tangled" + "tangled.org/core/spindle/models" xrpcerr "tangled.org/core/xrpc/errors" ) @@ -66,7 +67,7 @@ func (x *Xrpc) HandleCiGetPipeline(w http.ResponseWriter, r *http.Request) { return } - p, err := x.Db.GetPipeline(r.Context(), pipeline) + p, err := x.Db.GetPipeline(r.Context(), models.PipelineId(pipeline)) if err != nil { fail(xrpcerr.GenericError(err), http.StatusInternalServerError) return diff --git a/spindle/xrpc/pipeline_cancel_pipeline.go b/spindle/xrpc/pipeline_cancel_pipeline.go index 5ec26e6e..6abecd15 100644 --- a/spindle/xrpc/pipeline_cancel_pipeline.go +++ b/spindle/xrpc/pipeline_cancel_pipeline.go @@ -32,9 +32,10 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { return } - pipelineTid, err := syntax.ParseTID(input.Pipeline) - if err != nil { - fail(xrpcerr.GenericError(fmt.Errorf("invalid pipeline TID %q: %w", input.Pipeline, err))) + // pipeline ids are opaque, any format. the GetPipeline + repo-ownership check + // below is the real gate; it also covers existence. + if input.Pipeline == "" { + fail(xrpcerr.GenericError(fmt.Errorf("pipeline is required"))) return } @@ -43,15 +44,9 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { fail(xerr) return } - repo, err := x.Db.GetRepoByDid(repoDid) - if err != nil { - fail(xrpcerr.GenericError(fmt.Errorf("failed to get repo: %w", err))) - return - } - // the actor is only authorized against input.Repo, so make sure the // pipeline actually belongs to it before cancelling anything - p, err := x.Db.GetPipeline(r.Context(), pipelineTid.String()) + p, err := x.Db.GetPipeline(r.Context(), models.PipelineId(input.Pipeline)) if err != nil { fail(xrpcerr.GenericError(fmt.Errorf("failed to get pipeline: %w", err))) return @@ -61,11 +56,8 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { return } - pipelineId := models.PipelineId{ - Knot: repo.Knot, - Rkey: pipelineTid.String(), - } - l = l.With("input.pipeline", pipelineTid, "input.workflows", input.Workflows) + pipelineId := models.PipelineId(input.Pipeline) + l = l.With("input.pipeline", input.Pipeline, "input.workflows", input.Workflows) workflows := input.Workflows if len(workflows) == 0 { @@ -83,7 +75,7 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { for _, wName := range workflows { wid := models.WorkflowId{ PipelineId: pipelineId, - Name: wName, + Name: wName, } l.Debug("cancel pipeline", "wid", wid) diff --git a/spindle/xrpc/xrpc.go b/spindle/xrpc/xrpc.go index 4d390d96..1e97126d 100644 --- a/spindle/xrpc/xrpc.go +++ b/spindle/xrpc/xrpc.go @@ -39,7 +39,7 @@ func requireSha(sha string) error { // this is to break an import cycle. spindle imports this package for Xrpc, // so this package can't import *spindle.Spindle back. type PipelineTrigger interface { - TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull PullContext, inputs []*tangled.Pipeline_Pair) (syntax.ATURI, error) + TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull PullContext, inputs []*tangled.Pipeline_Pair) (models.PipelineId, error) DescribeWorkflowDefinition(ctx context.Context, repoDid syntax.DID, sha string, sourceRepo syntax.DID) (*tangled.CiDescribeWorkflowDefinition_Output, error) } diff --git a/spindle/xrpc/xrpc_test.go b/spindle/xrpc/xrpc_test.go index 8255ff66..7778b946 100644 --- a/spindle/xrpc/xrpc_test.go +++ b/spindle/xrpc/xrpc_test.go @@ -27,9 +27,9 @@ type mockTrigger struct { triggered bool } -func (m *mockTrigger) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull PullContext, inputs []*tangled.Pipeline_Pair) (syntax.ATURI, error) { +func (m *mockTrigger) TriggerManual(ctx context.Context, repoDid syntax.DID, sha, ref string, workflows []string, sourceRepo syntax.DID, pull PullContext, inputs []*tangled.Pipeline_Pair) (models.PipelineId, error) { m.triggered = true - return syntax.ParseATURI("at://did:plc:repoowner/sh.tangled.ci.pipeline/testrkey") + return models.PipelineId("testrkey"), nil } func (m *mockTrigger) DescribeWorkflowDefinition(context.Context, syntax.DID, string, syntax.DID) (*tangled.CiDescribeWorkflowDefinition_Output, error) { @@ -191,7 +191,7 @@ func TestCancelPipeline_RBAC(t *testing.T) { {Name: "test-workflow"}, }, } - err = d.CreatePipeline(models.PipelineId{Knot: "knot.test", Rkey: pipelineTid}, tpl) + err = d.CreatePipeline(models.PipelineId(pipelineTid), tpl) if err != nil { t.Fatalf("CreatePipeline: %v", err) }