From cce5f0bb7e4db7f85eace3d7c5ba5cc73c05fb6b Mon Sep 17 00:00:00 2001 From: dawn Date: Tue, 15 Sep 2026 07:10:02 +0300 Subject: [PATCH] spindle/knotfeed: run the push when the old ref state is missing spindle writes down where a branch pointed after it handles a push, and reads that back as the old commit for the next push on the same branch. the knot's own feed hands us the old commit in the event, so the note only matters where it does not. what was wrong: - spindle built the note for every repo before it started serving, and gave up on the whole process when one knot was down. on production that is 183 knots - it took an old saved position as proof the note was built. it was not built, so the first update after a fallback was refused, not run - it asked whether a branch was new only after reading the current branches, so a branch made in that gap looked like it already existed and the changed files came out wrong - it wrote the note before making the run. if it died in between, the push was counted as handled and never ran - after wiping a note it did not build it again what it does now: - nothing is built at startup, so no knot can keep spindle from starting - an op that reports no old commit falls back to the note; if there is no note either, this is the first push we have seen for that branch and it runs with no old commit, then we write the new one down - a push whose changed files we could not work out is no longer excluded by a paths filter. the filter still skips when we know the files and none match, and every other trigger keeps its old behaviour - the note is written after the run exists, so a push can run twice but is not skipped - the repo row is written before we listen to its knot, so a ref op for a brand new repo is not dropped as someone else's the cost is the first push per branch after a note is lost: it runs without an old commit, so its changed-file list is empty and a paths filter cannot skip it. Signed-off-by: dawn --- spindle/db/feed.go | 9 ------ spindle/knotfeed.go | 43 ++++++++----------------- spindle/knotfeed_test.go | 69 +++++++++++++++++++++++----------------- spindle/server.go | 53 +++++++++++++++--------------- spindle/tapclient.go | 7 ++-- workflow/def.go | 7 ++-- workflow/def_test.go | 26 +++++++++++++++ 7 files changed, 113 insertions(+), 101 deletions(-) diff --git a/spindle/db/feed.go b/spindle/db/feed.go index d79fb2a60..57b1cbe83 100644 --- a/spindle/db/feed.go +++ b/spindle/db/feed.go @@ -65,15 +65,6 @@ func (d *DB) PutFeedRef(repoDid syntax.DID, rkey syntax.RecordKey, sha knotfeed. return err } -func (d *DB) SeedFeedRef(repoDid syntax.DID, rkey syntax.RecordKey, sha knotfeed.ObjectID) error { - _, err := d.Exec( - `insert into feed_refs (repo_did, rkey, sha) values (?, ?, ?) - on conflict(repo_did, rkey) do nothing`, - repoDid.String(), rkey.String(), sha.String(), - ) - return err -} - func (d *DB) DeleteFeedRef(repoDid syntax.DID, rkey syntax.RecordKey) error { _, err := d.Exec( `delete from feed_refs where repo_did = ? and rkey = ?`, diff --git a/spindle/knotfeed.go b/spindle/knotfeed.go index 39eac961b..a5eb6c679 100644 --- a/spindle/knotfeed.go +++ b/spindle/knotfeed.go @@ -125,7 +125,7 @@ func (s *Spindle) handleRefOp(ctx context.Context, knot string, repoDid syntax.D return fmt.Errorf("decoding ref record: %w", err) } - oldSha, err := s.priorSha(ctx, knot, repoDid, op) + oldSha, err := s.priorSha(repoDid, op) if err != nil { return err } @@ -150,10 +150,6 @@ func (s *Spindle) handleRefOp(ctx context.Context, knot string, repoDid syntax.D changedFiles := changedFilesUnderBudget(l, repoPath, oldSha, record.Sha) - if err := s.db.PutFeedRef(repoDid, op.Rkey, record.Sha); err != nil { - return fmt.Errorf("recording ref state: %w", err) - } - triggerRepo, err := s.buildTriggerRepo(ctx, repo) if err != nil { return fmt.Errorf("building trigger repo: %w", err) @@ -173,6 +169,9 @@ func (s *Spindle) handleRefOp(ctx context.Context, knot string, repoDid syntax.D if err != nil { return err } + if err := s.db.PutFeedRef(repoDid, op.Rkey, record.Sha); err != nil { + return fmt.Errorf("recording ref state: %w", err) + } if pipelineId == "" { l.Info("no workflow matched 'push' trigger, skipping the event") return nil @@ -185,36 +184,21 @@ func isMaterializedRef(refname string) bool { return strings.HasPrefix(refname, "refs/heads/") || strings.HasPrefix(refname, "refs/tags/") } -func (s *Spindle) priorSha(ctx context.Context, knot string, repoDid syntax.DID, op knotfeed.RecordOp) (knotfeed.ObjectID, error) { +// priorSha is where the ref pointed before this op: reported by the feed, recorded +// from the previous op, or zero when we have never seen the ref, which is not an +// error. +func (s *Spindle) priorSha(repoDid syntax.DID, op knotfeed.RecordOp) (knotfeed.ObjectID, error) { if sha, ok := op.Prior.Sha(); ok { return sha, nil } - if err := s.ensureRefState(ctx, knot, repoDid); err != nil { - return knotfeed.ObjectID{}, fmt.Errorf("bootstrapping ref state: %w", err) - } - sha, _, err := s.db.FeedRefSha(repoDid, op.Rkey) + sha, found, err := s.db.FeedRefSha(repoDid, op.Rkey) if err != nil { return knotfeed.ObjectID{}, fmt.Errorf("reading ref state: %w", err) } - return sha, nil -} - -func (s *Spindle) ensureRefState(ctx context.Context, knot string, repoDid syntax.DID) error { - if _, seeded := s.refStateSeeded.LoadOrStore(repoDid.String(), struct{}{}); seeded { - return nil - } - refs, err := s.refRecords(ctx, knot, repoDid) - if err != nil { - s.refStateSeeded.Delete(repoDid.String()) - return err + if !found { + return knotfeed.ObjectID{}, nil } - for _, ref := range refs { - if err := s.db.SeedFeedRef(repoDid, ref.Rkey, ref.Sha); err != nil { - s.refStateSeeded.Delete(repoDid.String()) - return fmt.Errorf("seeding ref state: %w", err) - } - } - return nil + return sha, nil } func changedFilesUnderBudget(l *slog.Logger, repoPath string, oldSha, newSha knotfeed.ObjectID) []string { @@ -349,11 +333,10 @@ func (s *Spindle) feedOutdatedReplay(ctx context.Context, knot string, feed knot s.l.Error("failed to reset ref state, old shas may be stale until the next push", "repo", repo.RepoDid, "err", err) continue } - s.refStateSeeded.Delete(repo.RepoDid.String()) reset++ } if reset > 0 { - s.l.Warn("knot cannot replay from our cursor, ref state will be re-seeded live", "knot", knot, "repos", reset) + s.l.Warn("knot cannot replay from our cursor, ref state reset", "knot", knot, "repos", reset) } return live } diff --git a/spindle/knotfeed_test.go b/spindle/knotfeed_test.go index 48da3b806..a258f82ab 100644 --- a/spindle/knotfeed_test.go +++ b/spindle/knotfeed_test.go @@ -4,7 +4,6 @@ import ( "bytes" "context" "encoding/json" - "errors" "io" "log/slog" "net/http" @@ -248,37 +247,49 @@ func TestKnotFeedRefOps(t *testing.T) { }) } -func TestKnotFeedSeedsRefStateOnce(t *testing.T) { - s := newTestFeedSpindle(t) - rkeyTag := syntax.RecordKey("refs~2ftags~2fv1") - fetched := 0 - s.listRefRecords = func(ctx context.Context, knot string, repoDid syntax.DID) ([]refRecord, error) { - if fetched++; fetched == 1 { - return nil, errors.New("knot unreachable") +func TestKnotFeedPriorState(t *testing.T) { + t.Run("a ref we have never seen has no prior and asks the knot nothing", func(t *testing.T) { + s := newTestFeedSpindle(t) + s.listRefRecords = func(context.Context, string, syntax.DID) ([]refRecord, error) { + t.Fatal("resolving a prior sha fetched the knot's refs") + return nil, nil } - return []refRecord{{Rkey: testRkeyMain, Sha: testShaOld}, {Rkey: rkeyTag, Sha: testShaNew}}, nil - } - if err := s.ensureRefState(context.Background(), testKnot, testRepoDid); err == nil { - t.Fatal("expected the first bootstrap attempt to surface the fetch error") - } - if err := s.ensureRefState(context.Background(), testKnot, testRepoDid); err != nil { - t.Fatalf("ensureRefState: %v", err) - } - for _, want := range []struct { - rkey syntax.RecordKey - sha knotfeed.ObjectID - }{ - {testRkeyMain, testShaOld}, - {rkeyTag, testShaNew}, - } { - if sha, ok, err := s.db.FeedRefSha(testRepoDid, want.rkey); err != nil || !ok || sha != want.sha { - t.Errorf("ref state (%s, %s) = (%q, %v, err %v)", want.rkey, want.sha, sha, ok, err) + for _, action := range []tapc.RecordAction{tapc.RecordCreateAction, tapc.RecordUpdateAction} { + sha, err := s.priorSha(testRepoDid, knotfeed.RecordOp{Action: action, Rkey: testRkeyMain}) + if err != nil || !sha.IsZero() { + t.Errorf("%s prior sha = (%q, %v), want zero and nil", action, sha, err) + } } - } - if err := s.ensureRefState(context.Background(), testKnot, testRepoDid); err != nil || fetched != 2 { - t.Fatalf("re-bootstrap fetched %d times (err %v), want 2 fetches and no error", fetched, err) - } + }) + + t.Run("the recorded state answers an op that reports none", func(t *testing.T) { + s := newTestFeedSpindle(t) + if err := s.db.PutFeedRef(testRepoDid, testRkeyMain, testShaOld); err != nil { + t.Fatalf("record ref state: %v", err) + } + + sha, err := s.priorSha(testRepoDid, knotfeed.RecordOp{Action: tapc.RecordUpdateAction, Rkey: testRkeyMain}) + if err != nil || sha != testShaOld { + t.Errorf("prior sha = (%q, %v), want the recorded %q", sha, err, testShaOld) + } + }) + + t.Run("the reported prior wins over the recorded state", func(t *testing.T) { + s := newTestFeedSpindle(t) + if err := s.db.PutFeedRef(testRepoDid, testRkeyMain, testShaOld); err != nil { + t.Fatalf("record ref state: %v", err) + } + + sha, err := s.priorSha(testRepoDid, knotfeed.RecordOp{ + Action: tapc.RecordUpdateAction, + Rkey: testRkeyMain, + Prior: knotfeed.ParsePriorSha(testHexNew), + }) + if err != nil || sha != testShaNew { + t.Errorf("prior sha = (%q, %v), want the reported %q", sha, err, testShaNew) + } + }) } func TestAdmitChangedFilesUnderBudget(t *testing.T) { diff --git a/spindle/server.go b/spindle/server.go index 048119cff..60ec0180e 100644 --- a/spindle/server.go +++ b/spindle/server.go @@ -77,33 +77,32 @@ type executorClient interface { } type Spindle struct { - jc *jetstream.JetstreamClient - tap *Tap - embedTap *embeddedTap - db *db.DB - e *rbac.Enforcer - l *slog.Logger - n *notifier.Notifier - wh *webhook.Service - metrics *observability.Metrics - engs map[string]models.Engine - jobWake chan struct{} - jobWorkers sync.WaitGroup - cfg *config.Config - feed *feed.Feed - refStateSeeded sync.Map - res *idresolver.Resolver - verify repoverify.Verifier - vault secrets.Manager - cache storage.Storage - motd []byte - motdMu sync.RWMutex - rootCtx context.Context - rootCancel context.CancelFunc - store artifactstore.Store - stores *artifactstore.Stores - reader artifactstore.Reader - qm *quota.Manager + jc *jetstream.JetstreamClient + tap *Tap + embedTap *embeddedTap + db *db.DB + e *rbac.Enforcer + l *slog.Logger + n *notifier.Notifier + wh *webhook.Service + metrics *observability.Metrics + engs map[string]models.Engine + jobWake chan struct{} + jobWorkers sync.WaitGroup + cfg *config.Config + feed *feed.Feed + res *idresolver.Resolver + verify repoverify.Verifier + vault secrets.Manager + cache storage.Storage + motd []byte + motdMu sync.RWMutex + rootCtx context.Context + rootCancel context.CancelFunc + store artifactstore.Store + stores *artifactstore.Stores + reader artifactstore.Reader + qm *quota.Manager // set only when this spindle hosts the mill or joins one as an executor mill *mill.Mill exec executorClient diff --git a/spindle/tapclient.go b/spindle/tapclient.go index 29b81880e..c28ca0adf 100644 --- a/spindle/tapclient.go +++ b/spindle/tapclient.go @@ -211,10 +211,6 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error return fmt.Errorf("add repo policy: %w", err) } - if t.spindle.feed != nil { - t.spindle.feed.Subscribe(t.spindle.rootCtx, record.Knot) - } - repo := db.Repo{ Knot: record.Knot, Owner: ownerDid, @@ -228,6 +224,9 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error l.Error("failed to add repo row", "err", err) return fmt.Errorf("add repo: %w", err) } + if t.spindle.feed != nil { + t.spindle.feed.Subscribe(t.spindle.rootCtx, record.Knot) + } repoCloneUri := t.spindle.newRepoCloneUrl(repo.Knot, repo.RepoDid) repoPath := t.spindle.newRepoPath(repo.RepoDid) diff --git a/workflow/def.go b/workflow/def.go index be6ace864..c9311c40b 100644 --- a/workflow/def.go +++ b/workflow/def.go @@ -181,8 +181,11 @@ func (c *Constraint) Match(trigger tangled.Pipeline_TriggerMetadata, changedFile } } - // apply paths filter: if specified, at least one changed file must match - if len(c.Paths) > 0 { + // apply paths filter: if specified, at least one changed file must match. the + // push trigger is the only one we compute a diff for, so a nil list there means + // the diff could not be computed, and that must not skip the workflow quietly. + // every other trigger keeps the old behaviour of a list that does not match + if len(c.Paths) > 0 && (trigger.Push == nil || changedFiles != nil) { matched, err := matchesAnyFile(changedFiles, c.Paths) if err != nil { return false, err diff --git a/workflow/def_test.go b/workflow/def_test.go index 559aa12ec..e011ad83c 100644 --- a/workflow/def_test.go +++ b/workflow/def_test.go @@ -551,6 +551,32 @@ func TestConstraintMatchTypes(t *testing.T) { } } +func TestConstraintMatch_PathsFilterWithUnknownChangedFiles(t *testing.T) { + constraint := func(event string) *Constraint { + return &Constraint{Event: []string{event}, Branch: []string{"main"}, Paths: []string{"src/**"}} + } + + t.Run("a push we could not diff runs", func(t *testing.T) { + trigger := tangled.Pipeline_TriggerMetadata{ + Kind: string(TriggerKindPush), + Push: &tangled.Pipeline_PushTriggerData{Ref: "refs/heads/main"}, + } + matched, err := constraint("push").Match(trigger, nil) + assert.NoError(t, err) + assert.True(t, matched) + }) + + t.Run("a pull request keeps skipping on a list that does not match", func(t *testing.T) { + trigger := tangled.Pipeline_TriggerMetadata{ + Kind: string(TriggerKindPullRequest), + PullRequest: &tangled.Pipeline_PullRequestTriggerData{TargetBranch: "main"}, + } + matched, err := constraint("pull_request").Match(trigger, nil) + assert.NoError(t, err) + assert.False(t, matched) + }) +} + func TestConstraintMatch_PullRequestTypes(t *testing.T) { prTrigger := func(action string) tangled.Pipeline_TriggerMetadata { return tangled.Pipeline_TriggerMetadata{ -- 2.51.2