package 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 }