Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207package db
import ( "database/sql" "encoding/json" "fmt" "os"
"tangled.org/core/spindle/models")
type pipelineLogRename struct { knot string pipelineID models.PipelineId workflow string}
func stagePipelineLogRenames(tx *sql.Tx) error { var renames []pipelineLogRename hasLegacyPipelineID, err := tableHasColumn(tx, "pipelines", "rkey") if err != nil { return err } if hasLegacyPipelineID { rows, err := tx.Query(`select rkey, knot, payload from pipelines`) if err != nil { return err } for rows.Next() { var rkey, knot, payload string if err := rows.Scan(&rkey, &knot, &payload); err != nil { rows.Close() return err } var pipeline models.PipelineRecord if err := json.Unmarshal([]byte(payload), &pipeline); err != nil { rows.Close() return fmt.Errorf("decode pipeline %s while staging log migration: %w", rkey, err) } for _, workflow := range pipeline.Workflows { if workflow != nil && workflow.Name != "" { renames = append(renames, pipelineLogRename{knot, models.PipelineId(rkey), workflow.Name}) } } } if err := rows.Err(); err != nil { rows.Close() return err } if err := rows.Close(); err != nil { return err } }
rows, err := tx.Query(` select knot, rkey, workflow from mill_leases union select knot, rkey, workflow from mill_artifacts union select knot, rkey, workflow from executor_pending_artifacts `) if err != nil { return err } for rows.Next() { var knot, rkey, workflow string if err := rows.Scan(&knot, &rkey, &workflow); err != nil { rows.Close() return err } renames = append(renames, pipelineLogRename{knot, models.PipelineId(rkey), workflow}) } if err := rows.Err(); err != nil { rows.Close() return err } if err := rows.Close(); err != nil { return err }
for _, rename := range renames { if _, err := tx.Exec( `insert or ignore into pipeline_log_renames (knot, pipeline_id, workflow) values (?, ?, ?)`, rename.knot, rename.pipelineID, rename.workflow, ); err != nil { return err } } return nil}
func (d *DB) MigratePipelineLogFiles(logDir string) error { var tableExists int if err := d.QueryRow(` select count(*) from sqlite_master where type = 'table' and name = 'pipeline_log_renames' `).Scan(&tableExists); err != nil { return err } if tableExists == 0 { return nil } if logDir == "" { return d.finishPipelineLogMigration() } if _, err := os.Stat(logDir); err != nil { if os.IsNotExist(err) { return d.finishPipelineLogMigration() } return err }
rows, err := d.Query(`select knot, pipeline_id, workflow from pipeline_log_renames`) if err != nil { return err } var renames []pipelineLogRename for rows.Next() { var rename pipelineLogRename if err := rows.Scan(&rename.knot, &rename.pipelineID, &rename.workflow); err != nil { rows.Close() return err } renames = append(renames, rename) } if err := rows.Err(); err != nil { rows.Close() return err } if err := rows.Close(); err != nil { return err }
type plannedRename struct { rename pipelineLogRename oldPath string newPath string destinationIsSameFile bool } var planned []plannedRename var stale []pipelineLogRename destinations := make(map[string]string) for _, rename := range renames { oldPath := models.LegacyLogFilePath(logDir, rename.knot, rename.pipelineID, rename.workflow) newPath := models.LogFilePath(logDir, models.WorkflowId{PipelineId: rename.pipelineID, Name: rename.workflow}) oldInfo, err := os.Stat(oldPath) if err != nil { if os.IsNotExist(err) { stale = append(stale, rename) continue } return fmt.Errorf("stat legacy pipeline log %s: %w", oldPath, err) } if !oldInfo.Mode().IsRegular() { return fmt.Errorf("legacy pipeline log is not a regular file: %s", oldPath) } if previous, ok := destinations[newPath]; ok && previous != oldPath { return fmt.Errorf("legacy pipeline logs %s and %s both map to %s", previous, oldPath, newPath) } destinations[newPath] = oldPath
sameFile := false if newInfo, err := os.Stat(newPath); err == nil { if !os.SameFile(oldInfo, newInfo) { return fmt.Errorf("refusing to overwrite pipeline log %s while migrating %s", newPath, oldPath) } sameFile = true } else if !os.IsNotExist(err) { return fmt.Errorf("stat pipeline log %s: %w", newPath, err) } planned = append(planned, plannedRename{rename, oldPath, newPath, sameFile}) }
for _, rename := range stale { if err := d.clearPipelineLogRename(rename); err != nil { return err } } for _, plan := range planned { if !plan.destinationIsSameFile { if err := os.Link(plan.oldPath, plan.newPath); err != nil { return fmt.Errorf("link pipeline log %s to %s: %w", plan.oldPath, plan.newPath, err) } } if err := os.Remove(plan.oldPath); err != nil { return fmt.Errorf("remove legacy pipeline log %s: %w", plan.oldPath, err) } if err := d.clearPipelineLogRename(plan.rename); err != nil { return err } } return d.finishPipelineLogMigration()}
func (d *DB) finishPipelineLogMigration() error { _, err := d.Exec(`drop table if exists pipeline_log_renames`) return err}
func (d *DB) clearPipelineLogRename(rename pipelineLogRename) error { _, err := d.Exec( `delete from pipeline_log_renames where knot = ? and pipeline_id = ? and workflow = ?`, rename.knot, rename.pipelineID, rename.workflow, ) return err}