Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271package db
import ( "context" "database/sql" "encoding/json" "os" "path/filepath" "slices" "testing"
"tangled.org/core/api/tangled" "tangled.org/core/spindle/models")
func TestPipelineIdentityFinalSchemaHasNoLegacyState(t *testing.T) { ctx := context.Background() database, err := Make(ctx, filepath.Join(t.TempDir(), "final-schema.db")) if err != nil { t.Fatal(err) } defer database.Close() if err := database.MigratePipelineLogFiles(""); err != nil { t.Fatal(err) }
for _, table := range []string{"events", "events_legacy", "pipeline_log_renames"} { var count int if err := database.QueryRow(`select count(*) from sqlite_master where type = 'table' and name = ?`, table).Scan(&count); err != nil || count != 0 { t.Fatalf("legacy table %s count = %d, %v", table, count, err) } } legacyColumns := map[string][]string{ "pipelines": {"knot", "rkey", "definitions"}, "workflow_statuses": {"rkey"}, "jobs": {"pipeline_id_knot", "pipeline_id_rkey"}, "mill_leases": {"knot", "rkey"}, "mill_artifacts": {"knot", "rkey"}, "executor_pending_artifacts": {"knot", "rkey"}, } for table, columns := range legacyColumns { for _, column := range columns { var count int if err := database.QueryRow(`select count(*) from pragma_table_info(?) where name = ?`, table, column).Scan(&count); err != nil || count != 0 { t.Fatalf("legacy column %s.%s count = %d, %v", table, column, count, err) } } }}
func TestPipelineMigrationSkipsMalformedLegacyRows(t *testing.T) { ctx := context.Background() dbPath := filepath.Join(t.TempDir(), "malformed.db") database, err := Make(ctx, dbPath) if err != nil { t.Fatal(err) } if err := database.Close(); err != nil { t.Fatal(err) }
rawDB, err := sql.Open("sqlite3", dbPath) if err != nil { t.Fatal(err) } if _, err := rawDB.Exec(` delete from migrations where name in ('pipelines-and-workflow-statuses', 'pipeline-id-identity', 'drop-legacy-pipeline-events');
alter table jobs rename column pipeline_id to pipeline_id_rkey; alter table jobs add column pipeline_id_knot text not null default '';
alter table mill_leases rename column pipeline_id to rkey; alter table mill_leases add column knot text not null default '';
drop index idx_mill_artifacts_workflow_identity; alter table mill_artifacts rename column pipeline_id to rkey; alter table mill_artifacts add column knot text not null default ''; create index idx_mill_artifacts_workflow_identity on mill_artifacts(knot, rkey, workflow, id desc);
alter table executor_pending_artifacts rename column pipeline_id to rkey; alter table executor_pending_artifacts add column knot text not null default '';
drop table pipeline_log_renames; drop table if exists events_legacy; drop table pipelines; drop table workflow_statuses; create table events ( rkey text not null, nsid text not null, event text not null, created integer not null ); insert into events (rkey, nsid, event, created) values ('3mu2xwiorc2xl', 'sh.tangled.pipeline', '{', 1); `); err != nil { rawDB.Close() t.Fatal(err) } repo := "did:plc:validrepo" validPipeline := tangled.Pipeline{ TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ Kind: "push", Repo: &tangled.Pipeline_TriggerRepo{Knot: "knot.test", RepoDid: &repo}, Push: &tangled.Pipeline_PushTriggerData{NewSha: "1111111111111111111111111111111111111111", Ref: "refs/heads/main"}, }, Workflows: []*tangled.Pipeline_Workflow{{Name: "ci.yml", Raw: "when:\n - event: push\n"}}, } validJSON, err := json.Marshal(validPipeline) if err != nil { t.Fatal(err) } if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "3mu2xwiorc2xm", tangled.PipelineNSID, validJSON, int64(2)); err != nil { t.Fatal(err) } statusJSON, err := json.Marshal(tangled.PipelineStatus{ Pipeline: "at://did:plc:validrepo/sh.tangled.pipeline/3mu2xwiorc2xm", Workflow: "ci.yml", Status: "success", CreatedAt: "2026-08-31T08:00:00Z", }) if err != nil { t.Fatal(err) } if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "status-valid", tangled.PipelineStatusNSID, statusJSON, int64(3)); err != nil { t.Fatal(err) } if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "status-broken", tangled.PipelineStatusNSID, "{", int64(4)); err != nil { t.Fatal(err) } invalidDefinitionPipeline := validPipeline invalidDefinitionPipeline.Workflows = []*tangled.Pipeline_Workflow{ {Name: "valid.yml", Raw: "when:\n - event: push\n"}, {Name: "broken.yml", Raw: "when: ["}, } invalidDefinitionJSON, err := json.Marshal(invalidDefinitionPipeline) if err != nil { t.Fatal(err) } if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "3mu2xwiorc2xn", tangled.PipelineNSID, invalidDefinitionJSON, int64(5)); err != nil { t.Fatal(err) }
orphanStatusJSON, err := json.Marshal(tangled.PipelineStatus{ Pipeline: "at://did:plc:validrepo/sh.tangled.pipeline/3mu2xwiorc2xn", Workflow: "broken.yml", Status: "failed", CreatedAt: "2026-08-31T08:01:00Z", }) if err != nil { t.Fatal(err) } if _, err := rawDB.Exec(`insert into events (rkey, nsid, event, created) values (?, ?, ?, ?)`, "status-orphan", tangled.PipelineStatusNSID, orphanStatusJSON, int64(6)); err != nil { t.Fatal(err) } if err := rawDB.Close(); err != nil { t.Fatal(err) } logDir := filepath.Join(t.TempDir(), "logs") if err := os.MkdirAll(logDir, 0o755); err != nil { t.Fatal(err) } pipelineID := models.PipelineId("3mu2xwiorc2xm") oldLog := models.LegacyLogFilePath(logDir, "knot.test", pipelineID, "ci.yml") if err := os.WriteFile(oldLog, []byte("historical log"), 0o600); err != nil { t.Fatal(err) }
upgraded, err := Make(ctx, dbPath) if err != nil { t.Fatal(err) } defer upgraded.Close() if err := upgraded.MigratePipelineLogFiles(logDir); err != nil { t.Fatal(err) } var archivedEventsTable int if err := upgraded.QueryRow(`select count(*) from sqlite_master where type = 'table' and name = 'events_legacy'`).Scan(&archivedEventsTable); err != nil || archivedEventsTable != 0 { t.Fatalf("events_legacy table count = %d, %v", archivedEventsTable, err) } var converted int if err := upgraded.QueryRow(`select count(*) from pipelines`).Scan(&converted); err != nil { t.Fatal(err) } if converted != 1 { t.Fatalf("valid legacy pipelines converted = %d", converted) }
var statuses int if err := upgraded.QueryRow(`select count(*) from workflow_statuses`).Scan(&statuses); err != nil || statuses != 1 { t.Fatalf("workflow statuses = %d, %v", statuses, err) } record, err := upgraded.GetPipeline(ctx, pipelineID) if err != nil { t.Fatal(err) } if len(record.Workflows) != 1 || record.Workflows[0].Definition == nil { t.Fatalf("canonical workflow = %#v", record.Workflows) } workflow := record.Workflows[0] definition := workflow.Definition if definition.Hash != nil || !slices.Equal(definition.Triggers, []string{"org.tangled.event.push"}) { t.Fatalf("canonical definition = %#v", definition) } file := definition.Source.File if file == nil || file.Repo != repo || file.Commit != record.Commit || file.Path != ".tangled/workflows/ci.yml" { t.Fatalf("canonical source = %#v", file) } if record.ID != pipelineID || record.Trigger.Push == nil || record.Trigger.Push.Ref != "refs/heads/main" || workflow.Status != "success" || workflow.FinishedAt == nil { t.Fatalf("canonical run = %#v", record) } newLog := models.LogFilePath(logDir, models.WorkflowId{PipelineId: pipelineID, Name: "ci.yml"}) if contents, err := os.ReadFile(newLog); err != nil || string(contents) != "historical log" { t.Fatalf("migrated log = %q, %v", contents, err) } if _, err := upgraded.GetPipeline(ctx, models.PipelineId("3mu2xwiorc2xn")); err != sql.ErrNoRows { t.Fatalf("pipeline with malformed definition survived migration: %v", err) }}
func TestMigratePipelineLogFilesDoesNotClobberAndResumesHardLink(t *testing.T) { ctx := context.Background() database, err := Make(ctx, filepath.Join(t.TempDir(), "logs.db")) if err != nil { t.Fatal(err) } defer database.Close() pipelineID := models.PipelineId("3mu2xwiorc2xl") if _, err := database.Exec(`insert into pipeline_log_renames (knot, pipeline_id, workflow) values (?, ?, ?)`, "knot.test", pipelineID, "ci.yml"); err != nil { t.Fatal(err) } logDir := filepath.Join(t.TempDir(), "logs") if err := os.MkdirAll(logDir, 0o755); err != nil { t.Fatal(err) } oldPath := models.LegacyLogFilePath(logDir, "knot.test", pipelineID, "ci.yml") newPath := models.LogFilePath(logDir, models.WorkflowId{PipelineId: pipelineID, Name: "ci.yml"}) if err := os.WriteFile(oldPath, []byte("old"), 0o600); err != nil { t.Fatal(err) } if err := os.WriteFile(newPath, []byte("new"), 0o600); err != nil { t.Fatal(err) } if err := database.MigratePipelineLogFiles(logDir); err == nil { t.Fatal("expected existing destination to reject migration") } if oldData, _ := os.ReadFile(oldPath); string(oldData) != "old" { t.Fatalf("old path changed: %q", oldData) } if newData, _ := os.ReadFile(newPath); string(newData) != "new" { t.Fatalf("new path changed: %q", newData) }
if err := os.Remove(newPath); err != nil { t.Fatal(err) } if err := os.Link(oldPath, newPath); err != nil { t.Fatal(err) } if err := database.MigratePipelineLogFiles(logDir); err != nil { t.Fatal(err) } if _, err := os.Stat(oldPath); !os.IsNotExist(err) { t.Fatalf("old hard link still exists: %v", err) } if data, err := os.ReadFile(newPath); err != nil || string(data) != "old" { t.Fatalf("new hard link = %q, %v", data, err) } var renameTable int if err := database.QueryRow(`select count(*) from sqlite_master where type = 'table' and name = 'pipeline_log_renames'`).Scan(&renameTable); err != nil || renameTable != 0 { t.Fatalf("pipeline_log_renames table count = %d, %v", renameTable, err) } if err := database.MigratePipelineLogFiles(logDir); err != nil { t.Fatalf("repeat log migration: %v", err) }}