diff --git a/spindle/db/artifacts.go b/spindle/db/artifacts.go index 599cff819..389aade68 100644 --- a/spindle/db/artifacts.go +++ b/spindle/db/artifacts.go @@ -1,5 +1,11 @@ package db +import ( + "errors" + + "tangled.org/core/spindle/models" +) + type FinishedLog struct { LeaseID string Workflow string @@ -7,14 +13,17 @@ type FinishedLog struct { Hash string } -func (d *DB) GetFinishedLog(workflow string) (*FinishedLog, error) { +func (d *DB) GetFinishedLog(wid models.WorkflowId) (*FinishedLog, error) { + if err := validateWorkflowIdentity(wid); err != nil { + return nil, err + } var fl FinishedLog err := d.QueryRow( `select lease_id, workflow, ref, hash from mill_artifacts - where workflow = ? + where knot = ? and rkey = ? and workflow = ? order by id desc limit 1`, - workflow, + wid.Knot, wid.Rkey, wid.Name, ).Scan(&fl.LeaseID, &fl.Workflow, &fl.Ref, &fl.Hash) if err != nil { return nil, err @@ -22,11 +31,21 @@ func (d *DB) GetFinishedLog(workflow string) (*FinishedLog, error) { return &fl, nil } -func (d *DB) SaveArtifactRef(leaseID, repoDid, workflow, ref, hash string) error { +func (d *DB) SaveArtifactRef(leaseID, repoDid string, wid models.WorkflowId, ref, hash string) error { + if err := validateWorkflowIdentity(wid); err != nil { + return err + } _, err := d.Exec( - `insert into mill_artifacts (lease_id, repo_did, workflow, ref, hash) - values (?, ?, ?, ?, ?)`, - leaseID, repoDid, workflow, ref, hash, + `insert into mill_artifacts (lease_id, repo_did, knot, rkey, workflow, ref, hash) + values (?, ?, ?, ?, ?, ?, ?)`, + leaseID, repoDid, wid.Knot, wid.Rkey, wid.Name, ref, hash, ) return err } + +func validateWorkflowIdentity(wid models.WorkflowId) error { + if wid.Knot == "" || wid.Rkey == "" || wid.Name == "" { + return errors.New("incomplete workflow identity") + } + return nil +} diff --git a/spindle/db/artifacts_test.go b/spindle/db/artifacts_test.go index 73e6b2ecb..75d2bd4d5 100644 --- a/spindle/db/artifacts_test.go +++ b/spindle/db/artifacts_test.go @@ -5,8 +5,62 @@ import ( "database/sql" "path/filepath" "testing" + + "tangled.org/core/spindle/models" ) +func TestMakeAcceptsPrebackfilledArtifactIdentity(t *testing.T) { + path := filepath.Join(t.TempDir(), "spindle.db") + legacy, err := sql.Open("sqlite3", path) + if err != nil { + t.Fatalf("open legacy database: %v", err) + } + if _, err := legacy.Exec(` + create table mill_artifacts ( + id integer primary key autoincrement, + lease_id text not null, + repo_did text not null, + knot text not null default '', + rkey text not null default '', + workflow text not null, + ref text not null, + hash text not null + ); + insert into mill_artifacts ( + id, lease_id, repo_did, knot, rkey, workflow, ref, hash + ) values ( + 1, 'legacy-lease', 'did:plc:legacy-repo', 'knot', 'rkey', + 'build.yml', 'legacy-ref', 'legacy-hash' + ); + `); err != nil { + legacy.Close() + t.Fatalf("seed legacy database: %v", err) + } + if err := legacy.Close(); err != nil { + t.Fatalf("close legacy database: %v", err) + } + + d, err := Make(context.Background(), path) + if err != nil { + t.Fatalf("Make: %v", err) + } + t.Cleanup(func() { d.Close() }) + + wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot", Rkey: "rkey"}, Name: "build.yml"} + log, err := d.GetFinishedLog(wid) + if err != nil { + t.Fatalf("GetFinishedLog(legacy): %v", err) + } + if log.LeaseID != "legacy-lease" || log.Ref != "legacy-ref" { + t.Fatalf("GetFinishedLog(legacy) = %+v", log) + } + + newWID := models.WorkflowId{PipelineId: models.PipelineId{Knot: "new-knot", Rkey: "new-rkey"}, Name: "build.yml"} + if err := d.SaveArtifactRef("new-lease", "did:plc:repo", newWID, "new-ref", "new-hash"); err != nil { + t.Fatalf("SaveArtifactRef: %v", err) + } +} + func TestMakeMigratesLegacyMillArtifactsRepoDID(t *testing.T) { path := filepath.Join(t.TempDir(), "spindle.db") legacy, err := sql.Open("sqlite3", path) @@ -17,6 +71,8 @@ func TestMakeMigratesLegacyMillArtifactsRepoDID(t *testing.T) { create table mill_artifacts ( id integer primary key autoincrement, lease_id text not null, + knot text not null, + rkey text not null, workflow text not null, ref text not null, hash text not null @@ -38,8 +94,8 @@ func TestMakeMigratesLegacyMillArtifactsRepoDID(t *testing.T) { 'legacy-lease', 'node', 'epoch', 'engine', 'knot', 'rkey', 'legacy-workflow', 'done', 'did:plc:legacy-repo' ); - insert into mill_artifacts (lease_id, workflow, ref, hash) - values ('legacy-lease', 'legacy-workflow', 'legacy-ref', 'legacy-hash'); + insert into mill_artifacts (lease_id, knot, rkey, workflow, ref, hash) + values ('legacy-lease', 'knot', 'rkey', 'legacy-workflow', 'legacy-ref', 'legacy-hash'); `); err != nil { legacy.Close() t.Fatalf("seed legacy database: %v", err) @@ -54,22 +110,57 @@ func TestMakeMigratesLegacyMillArtifactsRepoDID(t *testing.T) { } t.Cleanup(func() { d.Close() }) - var legacyRepoDID string - if err := d.QueryRow(`select repo_did from mill_artifacts where lease_id = 'legacy-lease'`).Scan(&legacyRepoDID); err != nil { + var repoDID string + if err := d.QueryRow(`select repo_did from mill_artifacts where lease_id = 'legacy-lease'`).Scan(&repoDID); err != nil { t.Fatalf("query migrated artifact: %v", err) } - if legacyRepoDID != "did:plc:legacy-repo" { - t.Fatalf("legacy repo_did = %q, want did:plc:legacy-repo", legacyRepoDID) + if repoDID != "did:plc:legacy-repo" { + t.Fatalf("legacy repo_did = %q, want did:plc:legacy-repo", repoDID) } +} - if err := d.SaveArtifactRef("new-lease", "did:plc:repo", "build", "new-ref", "new-hash"); err != nil { - t.Fatalf("SaveArtifactRef: %v", err) +func TestGetFinishedLogMatchesPipelineAndWorkflow(t *testing.T) { + d := newTestDB(t) + wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example.com", Rkey: "pipeline-1"}, Name: "build.yml"} + otherPipeline := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example.com", Rkey: "pipeline-2"}, Name: wid.Name} + otherKnot := models.WorkflowId{PipelineId: models.PipelineId{Knot: "other.example.com", Rkey: wid.Rkey}, Name: wid.Name} + + for _, artifact := range []struct { + leaseID string + wid models.WorkflowId + ref string + }{ + {"wanted-old", wid, "wanted-old.log"}, + {"other-pipeline", otherPipeline, "other-pipeline.log"}, + {"other-knot", otherKnot, "other-knot.log"}, + {"wanted-new", wid, "wanted-new.log"}, + } { + if err := d.SaveArtifactRef(artifact.leaseID, "did:plc:repo", artifact.wid, artifact.ref, "hash"); err != nil { + t.Fatalf("SaveArtifactRef(%s): %v", artifact.leaseID, err) + } } - var repoDID string - if err := d.QueryRow(`select repo_did from mill_artifacts where lease_id = 'new-lease'`).Scan(&repoDID); err != nil { - t.Fatalf("query new artifact: %v", err) + + log, err := d.GetFinishedLog(wid) + if err != nil { + t.Fatalf("GetFinishedLog: %v", err) } - if repoDID != "did:plc:repo" { - t.Fatalf("new repo_did = %q, want did:plc:repo", repoDID) + if log.LeaseID != "wanted-new" || log.Ref != "wanted-new.log" { + t.Fatalf("GetFinishedLog = %+v, want newest artifact for exact workflow identity", log) + } +} + +func TestFinishedLogRejectsIncompleteWorkflowIdentity(t *testing.T) { + d := newTestDB(t) + for _, wid := range []models.WorkflowId{ + {PipelineId: models.PipelineId{Rkey: "rkey"}, Name: "build.yml"}, + {PipelineId: models.PipelineId{Knot: "knot"}, Name: "build.yml"}, + {PipelineId: models.PipelineId{Knot: "knot", Rkey: "rkey"}}, + } { + if _, err := d.GetFinishedLog(wid); err == nil { + t.Fatalf("GetFinishedLog(%+v) succeeded", wid) + } + if err := d.SaveArtifactRef("lease", "did:plc:repo", wid, "ref", "hash"); err == nil { + t.Fatalf("SaveArtifactRef(%+v) succeeded", wid) + } } } diff --git a/spindle/db/db.go b/spindle/db/db.go index 1d21279ce..e52437b84 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -186,6 +186,8 @@ func Make(ctx context.Context, dbPath string) (*DB, error) { id integer primary key autoincrement, lease_id text not null, repo_did text not null, + knot text not null, + rkey text not null, workflow text not null, ref text not null, hash text not null @@ -619,6 +621,34 @@ func runMigrations(_ context.Context, conn *sql.Conn, logger *slog.Logger) error return err } + if err := orm.RunMigration(conn, logger, "mill-artifacts-workflow-identity", func(tx *sql.Tx) error { + for _, column := range []string{"knot", "rkey"} { + var present int + if err := tx.QueryRow( + `select count(*) from pragma_table_info('mill_artifacts') where name = ?`, + column, + ).Scan(&present); err != nil { + return err + } + if present != 0 { + continue + } + if _, err := tx.Exec( + `alter table mill_artifacts add column ` + column + ` text not null default ''`, + ); err != nil { + return err + } + } + + _, err := tx.Exec(` + create index if not exists idx_mill_artifacts_workflow_identity + on mill_artifacts (knot, rkey, workflow, id desc) + `) + return err + }); err != nil { + return err + } + if err := orm.RunMigration(conn, logger, "bans-schema", func(tx *sql.Tx) error { _, err := tx.Exec(` create table if not exists bans ( diff --git a/spindle/db/delete_test.go b/spindle/db/delete_test.go index 7bc29821d..7839e9bb5 100644 --- a/spindle/db/delete_test.go +++ b/spindle/db/delete_test.go @@ -80,7 +80,7 @@ func TestDeleteMillLeasesByRepo(t *testing.T) { if err := d.SaveMillLease(l); err != nil { t.Fatal(err) } - if _, err := d.Exec(`insert into mill_artifacts (lease_id, repo_did, workflow, ref, hash) values (?, ?, ?, ?, ?)`, l.LeaseID, l.RepoDID, l.Workflow, "logs/"+l.LeaseID+".log", "h"); err != nil { + if _, err := d.Exec(`insert into mill_artifacts (lease_id, repo_did, knot, rkey, workflow, ref, hash) values (?, ?, ?, ?, ?, ?, ?)`, l.LeaseID, l.RepoDID, l.Knot, l.Rkey, l.Workflow, "logs/"+l.LeaseID+".log", "h"); err != nil { t.Fatal(err) } if _, err := d.Exec(`insert into executor_pending_artifacts (lease_id, workflow, status, ref, hash) values (?, ?, 'done', 'r', 'h')`, l.LeaseID, l.Workflow); err != nil { diff --git a/spindle/db/mill_state.go b/spindle/db/mill_state.go index 26c67ed05..32fba498d 100644 --- a/spindle/db/mill_state.go +++ b/spindle/db/mill_state.go @@ -6,6 +6,7 @@ import ( "tangled.org/core/eventstream" "tangled.org/core/notifier" + "tangled.org/core/spindle/models" ) // recovery uses persisted identity, fencing and quota state @@ -329,11 +330,14 @@ func (d *DB) ClearPendingArtifacts() error { return err } -func (tx *EventBatchTx) InsertArtifactRef(leaseID, repoDid, workflow, ref, hash string) error { +func (tx *EventBatchTx) InsertArtifactRef(leaseID, repoDid string, wid models.WorkflowId, ref, hash string) error { + if err := validateWorkflowIdentity(wid); err != nil { + return err + } _, err := tx.tx.Exec( - `insert into mill_artifacts (lease_id, repo_did, workflow, ref, hash) - values (?, ?, ?, ?, ?)`, - leaseID, repoDid, workflow, ref, hash, + `insert into mill_artifacts (lease_id, repo_did, knot, rkey, workflow, ref, hash) + values (?, ?, ?, ?, ?, ?, ?)`, + leaseID, repoDid, wid.Knot, wid.Rkey, wid.Name, ref, hash, ) return err } diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go index fbd2a4086..87095c2b5 100644 --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -676,7 +676,7 @@ func archiveWorkflowLog(l *slog.Logger, stores *artifactstore.Stores, database * return } digest := "sha256:" + hex.EncodeToString(hash.Sum(nil)) - if err := database.SaveArtifactRef(wid.String(), repoDID, wid.Name, ref, digest); err != nil { + if err := database.SaveArtifactRef(wid.String(), repoDID, wid, ref, digest); err != nil { l.Error("save workflow log artifact", "wid", wid, "err", err) } } diff --git a/spindle/logview/logview.go b/spindle/logview/logview.go index 2bffc06ed..e45d28153 100644 --- a/spindle/logview/logview.go +++ b/spindle/logview/logview.go @@ -3,7 +3,10 @@ package logview import ( "bufio" "context" + "errors" + "fmt" "io" + "os" "strings" "github.com/hpcloud/tail" @@ -16,10 +19,19 @@ import ( // the uploaded artifact once finished, same for every role. stop ends a // live follow early, the channel closes when the source drains or ctx ends func Follow(ctx context.Context, d *db.DB, reader artifactstore.Reader, logDir string, wid models.WorkflowId, finished bool) (<-chan *tail.Line, func(), error) { + logPath := models.LogFilePath(logDir, wid) + var artifactErr error if finished && reader != nil && d != nil { - if fl, err := d.GetFinishedLog(wid.Name); err == nil && fl.Ref != "" { + fl, err := d.GetFinishedLog(wid) + if err != nil { + artifactErr = fmt.Errorf("find archived workflow log: %w", err) + } else if fl.Ref == "" { + artifactErr = errors.New("archived workflow log has an empty ref") + } else { rc, err := reader.Open(ctx, fl.Ref) - if err == nil { + if err != nil { + artifactErr = fmt.Errorf("open archived workflow log: %w", err) + } else { ch := make(chan *tail.Line, 64) followCtx, cancel := context.WithCancel(ctx) go func() { @@ -39,7 +51,16 @@ func Follow(ctx context.Context, d *db.DB, reader artifactstore.Reader, logDir s } } - t, err := tail.TailFile(models.LogFilePath(logDir, wid), tail.Config{ + if artifactErr != nil { + if _, err := os.Stat(logPath); err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil, nil, fmt.Errorf("finished workflow log unavailable: %w", artifactErr) + } + return nil, nil, fmt.Errorf("stat local workflow log: %w", err) + } + } + + t, err := tail.TailFile(logPath, tail.Config{ Follow: !finished, ReOpen: !finished, MustExist: false, diff --git a/spindle/logview/logview_test.go b/spindle/logview/logview_test.go new file mode 100644 index 000000000..2b004424a --- /dev/null +++ b/spindle/logview/logview_test.go @@ -0,0 +1,82 @@ +package logview + +import ( + "context" + "fmt" + "io" + "strings" + "testing" + + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" +) + +type recordingReader struct { + contents map[string]string + opens []string +} + +func (r *recordingReader) Open(_ context.Context, ref string) (io.ReadCloser, error) { + r.opens = append(r.opens, ref) + content, ok := r.contents[ref] + if !ok { + return nil, fmt.Errorf("missing artifact %q", ref) + } + return io.NopCloser(strings.NewReader(content)), nil +} + +func TestFollowDoesNotUseAnotherPipelineArtifact(t *testing.T) { + d, err := db.Make(context.Background(), t.TempDir()+"/spindle.db") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { d.Close() }) + + target := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example.com", Rkey: "target"}, Name: "build.yml"} + other := models.WorkflowId{PipelineId: models.PipelineId{Knot: target.Knot, Rkey: "other"}, Name: target.Name} + if err := d.SaveArtifactRef("other-lease", "did:plc:other", other, "other.log", "hash"); err != nil { + t.Fatal(err) + } + + reader := &recordingReader{contents: map[string]string{"other.log": "wrong log\n"}} + if _, _, err := Follow(context.Background(), d, reader, t.TempDir(), target, true); err == nil { + t.Fatal("Follow succeeded without an artifact for the target pipeline") + } + if len(reader.opens) != 0 { + t.Fatalf("opened unrelated artifacts: %v", reader.opens) + } +} + +func TestFollowReadsExactPipelineArtifact(t *testing.T) { + d, err := db.Make(context.Background(), t.TempDir()+"/spindle.db") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { d.Close() }) + + target := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example.com", Rkey: "target"}, Name: "build.yml"} + other := models.WorkflowId{PipelineId: models.PipelineId{Knot: target.Knot, Rkey: "other"}, Name: target.Name} + if err := d.SaveArtifactRef("other-lease", "did:plc:other", other, "other.log", "hash"); err != nil { + t.Fatal(err) + } + if err := d.SaveArtifactRef("target-lease", "did:plc:target", target, "target.log", "hash"); err != nil { + t.Fatal(err) + } + + reader := &recordingReader{contents: map[string]string{ + "other.log": "wrong log\n", + "target.log": "target log\n", + }} + lines, stop, err := Follow(context.Background(), d, reader, t.TempDir(), target, true) + if err != nil { + t.Fatal(err) + } + defer stop() + line := <-lines + if line == nil || line.Text != "target log" { + t.Fatalf("line = %#v, want target log", line) + } + if len(reader.opens) != 1 || reader.opens[0] != "target.log" { + t.Fatalf("opened artifacts = %v, want target.log", reader.opens) + } +} diff --git a/spindle/mill/integration_test.go b/spindle/mill/integration_test.go index fa780c5ca..375fc506c 100644 --- a/spindle/mill/integration_test.go +++ b/spindle/mill/integration_test.go @@ -129,6 +129,14 @@ func TestEndToEndDummyJob(t *testing.T) { t.Fatalf("mill live log file was not removed after terminal artifact was recorded") } + archived, err := bdb.GetFinishedLog(wid) + if err != nil { + t.Fatalf("GetFinishedLog: %v", err) + } + if archived.Workflow != wid.Name || archived.Ref == "" { + t.Fatalf("archived log = %+v, want artifact for %s", archived, wid) + } + if !waitForStatus(t, bdb, wid, "running") { events, _ := bdb.GetEvents(0, 1000) t.Logf("mill events after completion: %+v", events) diff --git a/spindle/mill/mill.go b/spindle/mill/mill.go index 91266ada0..dfbab51c4 100644 --- a/spindle/mill/mill.go +++ b/spindle/mill/mill.go @@ -1404,12 +1404,16 @@ func (m *Mill) onEventBatch(sess *millSession, batch *millv1.EventBatch) error { if !strings.HasPrefix(a.GetHash(), "sha256:") { return protoErrf("invalid log artifact hash %q", a.GetHash()) } - if tx != nil { - if err := tx.InsertArtifactRef(lease.id, lease.repoDID, lease.wid.Name, a.GetRef(), a.GetHash()); err != nil { - return err + if lease.wid.Knot == "" || lease.wid.Rkey == "" || lease.wid.Name == "" { + m.l.Error("skipping log artifact with incomplete workflow identity", "lease", lease.id, "wid", lease.wid) + } else { + if tx != nil { + if err := tx.InsertArtifactRef(lease.id, lease.repoDID, lease.wid, a.GetRef(), a.GetHash()); err != nil { + return err + } } + artifactLeases = append(artifactLeases, lease) } - artifactLeases = append(artifactLeases, lease) } } pendingTerminals = append(pendingTerminals, pendingTerminal{ diff --git a/spindle/mill/mill_test.go b/spindle/mill/mill_test.go index f27b38219..04b8c3206 100644 --- a/spindle/mill/mill_test.go +++ b/spindle/mill/mill_test.go @@ -676,6 +676,50 @@ func TestTerminalBeforeACK(t *testing.T) { } } +func TestTerminalWithIncompleteIdentityAdvancesCursor(t *testing.T) { + m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) + sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(*millproto.Message) error { return nil }), discardLogger()) + m.attachSession(sess) + + lease := newLease("lease-1", sess.nodeID, sess.epoch, "dummy") + lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Rkey: "r"}, Name: "build"} + m.mu.Lock() + m.leases[lease.id] = lease + m.mu.Unlock() + + err := m.onEventBatch(sess, &millv1.EventBatch{ + Epoch: sess.epoch, + Events: []*millv1.Event{{ + Seqno: 1, + LeaseId: lease.id, + Payload: &millv1.Event_AttemptResult{AttemptResult: &millv1.AttemptResult{ + Status: millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS, + LogArtifact: &millv1.LogArtifact{ + Ref: "logs/" + lease.id + ".log", + Hash: "sha256:test", + }, + }}, + }}, + }) + if err != nil { + t.Fatalf("onEventBatch: %v", err) + } + + m.mu.Lock() + seqno := m.nodeSeqno[sess.nodeID+"/"+sess.epoch] + m.mu.Unlock() + if seqno != 1 { + t.Fatalf("cursor = %d, want 1", seqno) + } + var artifacts int + if err := bdb.QueryRow(`select count(*) from mill_artifacts`).Scan(&artifacts); err != nil { + t.Fatal(err) + } + if artifacts != 0 { + t.Fatalf("artifact rows = %d, want 0", artifacts) + } +} + func TestExecutorRestartEmptySnapshot(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) diff --git a/spindle/wipe_test.go b/spindle/wipe_test.go index 70e231664..e5c2d0c78 100644 --- a/spindle/wipe_test.go +++ b/spindle/wipe_test.go @@ -95,7 +95,7 @@ func seedWipeRepo(t *testing.T, s *Spindle, repoDid, owner syntax.DID) { if err := s.db.SaveMillLease(db.MillLease{LeaseID: "l" + sfx, NodeID: "n1", Epoch: "e1", Engine: "microvm", Knot: "knot.example.com", Rkey: "p" + sfx, Workflow: "w" + sfx, State: "active", RepoDID: repoDid.String()}); err != nil { t.Fatal(err) } - if _, err := s.db.Exec(`insert into mill_artifacts (lease_id, repo_did, workflow, ref, hash) values (?, ?, ?, ?, 'h')`, "l"+sfx, repoDid.String(), "w"+sfx, "out/l"+sfx+".bin"); err != nil { + if _, err := s.db.Exec(`insert into mill_artifacts (lease_id, repo_did, knot, rkey, workflow, ref, hash) values (?, ?, 'knot.example.com', ?, ?, ?, 'h')`, "l"+sfx, repoDid.String(), "p"+sfx, "w"+sfx, "out/l"+sfx+".bin"); err != nil { t.Fatal(err) } if _, err := s.db.Exec(`insert into quota_allocations (repo_did, resource, kind, key, amount) values (?, 'compute', 'generic', 'k', 1)`, repoDid.String()); err != nil { @@ -248,7 +248,7 @@ func TestWipeRepoRemovesLocalEngineArtifacts(t *testing.T) { wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.example.com", Rkey: "p1"}, Name: "w1"} localRef := "logs/" + wid.String() + ".log" - if err := s.db.SaveArtifactRef(wid.String(), repoDid.String(), wid.Name, localRef, "h"); err != nil { + if err := s.db.SaveArtifactRef(wid.String(), repoDid.String(), wid, localRef, "h"); err != nil { t.Fatal(err) } if err := s.stores.PutFile(ctx, localRef, writeTempFile(t, "log")); err != nil && len(err) > 0 { diff --git a/spindle/xrpc/ci_pipeline_subscribe_logs.go b/spindle/xrpc/ci_pipeline_subscribe_logs.go index 51218a303..5adf44353 100644 --- a/spindle/xrpc/ci_pipeline_subscribe_logs.go +++ b/spindle/xrpc/ci_pipeline_subscribe_logs.go @@ -2,7 +2,9 @@ package xrpc import ( "context" + "database/sql" "encoding/json" + "errors" "fmt" "net/http" "sync" @@ -167,7 +169,11 @@ func (x *Xrpc) handleSubscribeLogs(w http.ResponseWriter, r *http.Request, pipel lines, stop, err := logview.Follow(ctx, x.Db, x.ArtifactReader, x.Config.Server.LogDir, wid, isFinished) if err != nil { - l.ErrorContext(r.Context(), "failed to follow workflow log", "workflow", wfName, "err", err) + if errors.Is(err, sql.ErrNoRows) { + l.WarnContext(r.Context(), "archived workflow log unavailable", "workflow", wfName, "err", err) + } else { + l.ErrorContext(r.Context(), "failed to follow workflow log", "workflow", wfName, "err", err) + } return } defer stop()