Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions pkg/storage/internal/sqlstore/plan_comments.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@ type planCommentStore struct {

// Insert stores a newly posted plan comment and sets comment.ID.
func (s *planCommentStore) Insert(ctx context.Context, comment *storage.PlanComment) error {
canonicalizePlanCommentIdentity(comment)

id, err := s.identity.InsertID(ctx, s.db, `
INSERT INTO plan_comments (repository, pull_request, database_name, database_type, environment_scope, head_sha, github_comment_id, github_node_id)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
Expand All @@ -39,6 +41,9 @@ func (s *planCommentStore) Insert(ctx context.Context, comment *storage.PlanComm
// ListUnminimizedForSlot returns the not-yet-minimized comments for a
// (repository, pull_request, database) slot, ordered by id ascending.
func (s *planCommentStore) ListUnminimizedForSlot(ctx context.Context, repo string, pr int, database, databaseType string) ([]*storage.PlanComment, error) {
repo = storage.CanonicalKey(repo)
database = storage.CanonicalKey(database)
databaseType = storage.CanonicalKey(databaseType)
rows, err := s.db.QueryContext(ctx, `
SELECT `+planCommentColumns+`
FROM plan_comments
Expand All @@ -58,6 +63,7 @@ func (s *planCommentStore) ListUnminimizedForSlot(ctx context.Context, repo stri
// database has no other way to ask. The (repository, pull_request) prefix of
// the slot index serves it.
func (s *planCommentStore) ListUnminimizedForRepoPR(ctx context.Context, repo string, pr int) ([]*storage.PlanComment, error) {
repo = storage.CanonicalKey(repo)
rows, err := s.db.QueryContext(ctx, `
SELECT `+planCommentColumns+`
FROM plan_comments
Expand Down Expand Up @@ -101,6 +107,17 @@ func (s *planCommentStore) MarkMinimized(ctx context.Context, id int64) error {
return nil
}

// canonicalizePlanCommentIdentity folds the identity keys that appear in this
// store's SQL predicates. EnvironmentScope is deliberately not folded: no
// query filters on it, and its consumers compare it in Go against values built
// fresh from configured environment names, so folding only the stored side
// would break those comparisons.
func canonicalizePlanCommentIdentity(comment *storage.PlanComment) {
comment.Repository = storage.CanonicalKey(comment.Repository)
comment.DatabaseName = storage.CanonicalKey(comment.DatabaseName)
comment.DatabaseType = storage.CanonicalKey(comment.DatabaseType)
}

// scanPlanComment scans plan comment data from any scanner (Row or Rows).
func scanPlanComment(s scanner) (*storage.PlanComment, error) {
var comment storage.PlanComment
Expand Down
15 changes: 15 additions & 0 deletions pkg/storage/internal/sqlstore/plans.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,8 @@ type planStore struct {

// Create stores a new plan and returns its ID.
func (s *planStore) Create(ctx context.Context, plan *storage.Plan) (int64, error) {
canonicalizePlanIdentity(plan)

planDataJSON, err := json.Marshal(namespacesWithShardPlans(plan))
if err != nil {
return 0, fmt.Errorf("marshal plan data: %w", err)
Expand Down Expand Up @@ -98,6 +100,7 @@ func (s *planStore) GetByLock(ctx context.Context, lockID int64) ([]*storage.Pla
// latestPlanForTarget take the first matching row as "the newest plan" and rely
// on this deterministic order to avoid picking an older SHA on ties.
func (s *planStore) GetByPR(ctx context.Context, repo string, pr int) ([]*storage.Plan, error) {
repo = storage.CanonicalKey(repo)
rows, err := s.db.QueryContext(ctx, `
SELECT `+planColumns+`
FROM plans
Expand All @@ -119,6 +122,10 @@ func (s *planStore) GetByPR(ctx context.Context, repo string, pr int) ([]*storag
// plans.created_at is datetime (second precision), so the id tiebreaker keeps
// same-second plans in a deterministic newest-first order.
func (s *planStore) List(ctx context.Context, opts storage.ListPlansOptions) ([]*storage.Plan, error) {
opts.Database = storage.CanonicalKey(opts.Database)
opts.Environment = storage.CanonicalKey(opts.Environment)
opts.Repository = storage.CanonicalKey(opts.Repository)

if opts.Limit <= 0 {
return nil, fmt.Errorf("list plans for database %q environment %q: limit must be positive, got %d", opts.Database, opts.Environment, opts.Limit)
}
Expand Down Expand Up @@ -176,10 +183,18 @@ func (s *planStore) Delete(ctx context.Context, id int64) error {

// DeleteByPR removes all plans for a PR.
func (s *planStore) DeleteByPR(ctx context.Context, repo string, pr int) error {
repo = storage.CanonicalKey(repo)
_, err := s.db.ExecContext(ctx, `DELETE FROM plans WHERE repository = ? AND pull_request = ?`, repo, pr)
return err
}

func canonicalizePlanIdentity(plan *storage.Plan) {
plan.Database = storage.CanonicalKey(plan.Database)
plan.DatabaseType = storage.CanonicalKey(plan.DatabaseType)
plan.Repository = storage.CanonicalKey(plan.Repository)
plan.Environment = storage.CanonicalKey(plan.Environment)
}

// scanPlan scans a single plan row, returning nil if not found.
func scanPlan(row *sql.Row) (*storage.Plan, error) {
plan, err := scanPlanInto(row)
Expand Down
20 changes: 14 additions & 6 deletions pkg/storage/internal/sqlstore/webhook_events.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ func (s *webhookEventStore) Create(ctx context.Context, event *storage.WebhookEv
if event.Event == "" {
return false, fmt.Errorf("webhook event type is required")
}
event.Repository = storage.CanonicalKey(event.Repository)

provider := event.Provider
if provider == "" {
provider = storage.WebhookProviderGitHub
Expand Down Expand Up @@ -182,6 +184,7 @@ func webhookClaimableArgs() []any {
}

func (s *webhookEventStore) HasEventForHead(ctx context.Context, provider, repository string, pullRequest int, headSHA string) (bool, error) {
repository = storage.CanonicalKey(repository)
if provider == "" {
provider = storage.WebhookProviderGitHub
}
Expand Down Expand Up @@ -267,21 +270,23 @@ func (s *webhookEventStore) coveringSuccessorQuery(provider string, event *stora
// the live PR is worth a GitHub call at all — the common no-successor claim
// stays storage-only.
func (s *webhookEventStore) HasCoveringSuccessor(ctx context.Context, event *storage.WebhookEvent) (bool, error) {
if event.Repository == "" || event.PullRequest == 0 {
queryEvent := *event
queryEvent.Repository = storage.CanonicalKey(queryEvent.Repository)
if queryEvent.Repository == "" || queryEvent.PullRequest == 0 {
return false, fmt.Errorf("check covering successor for webhook event %d: repository and pull request are required for coalescing", event.ID)
}
provider := event.Provider
if provider == "" {
provider = storage.WebhookProviderGitHub
}
query, args := s.coveringSuccessorQuery(provider, event)
query, args := s.coveringSuccessorQuery(provider, &queryEvent)
var one int
err := s.db.QueryRowContext(ctx, query, args...).Scan(&one)
if errors.Is(err, sql.ErrNoRows) {
return false, nil
}
if err != nil {
return false, fmt.Errorf("check covering successor for webhook event %d (repo=%s, pr=%d): %w", event.ID, event.Repository, event.PullRequest, err)
return false, fmt.Errorf("check covering successor for webhook event %d (repo=%s, pr=%d): %w", event.ID, queryEvent.Repository, event.PullRequest, err)
}
return true, nil
}
Expand Down Expand Up @@ -311,18 +316,21 @@ func (s *webhookEventStore) HasCoveringSuccessor(ctx context.Context, event *sto
// terminally failed / superseded rows never run — none of those may justify
// discarding older work.
func (s *webhookEventStore) SupersedeIfCovered(ctx context.Context, event *storage.WebhookEvent) (bool, error) {
repository := storage.CanonicalKey(event.Repository)
if event.LeaseToken == "" {
return false, fmt.Errorf("webhook event lease token is required")
}
if event.Repository == "" || event.PullRequest == 0 {
if repository == "" || event.PullRequest == 0 {
return false, fmt.Errorf("supersede webhook event %d: repository and pull request are required for coalescing", event.ID)
}
provider := event.Provider
if provider == "" {
provider = storage.WebhookProviderGitHub
}
autoPlanIn := placeholders(len(storage.AutoPlanPullRequestActions))
successorQuery, successorArgs := s.coveringSuccessorQuery(provider, event)
queryEvent := *event
queryEvent.Repository = repository
successorQuery, successorArgs := s.coveringSuccessorQuery(provider, &queryEvent)
args := []any{storage.WebhookEventSuperseded, event.ID, event.LeaseToken, storage.WebhookEventProcessing}
args = append(args, stringArgs(storage.AutoPlanPullRequestActions)...)
args = append(args, successorArgs...)
Expand All @@ -336,7 +344,7 @@ func (s *webhookEventStore) SupersedeIfCovered(ctx context.Context, event *stora
)
`, args...)
if err != nil {
return false, fmt.Errorf("supersede webhook event %d (repo=%s, pr=%d): %w", event.ID, event.Repository, event.PullRequest, err)
return false, fmt.Errorf("supersede webhook event %d (repo=%s, pr=%d): %w", event.ID, repository, event.PullRequest, err)
}
rows, err := result.RowsAffected()
if err != nil {
Expand Down
9 changes: 9 additions & 0 deletions pkg/storage/storage.go
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,10 @@ type SettingsStore interface {
// primitive behind fast webhook acknowledgement: handlers can persist a delivery
// before returning 2xx, and drivers can claim/retry the stored event after the
// HTTP request has finished.
// Create canonicalizes the provided event's repository in place before
// persisting. The coalescing reads (HasCoveringSuccessor, SupersedeIfCovered)
// fold the repository only inside their SQL predicates and leave the caller's
// event untouched.
type WebhookEventStore interface {
// Create records a webhook delivery in the pending state. Returns
// inserted=false when provider + delivery GUID already exists, so callers
Expand Down Expand Up @@ -477,6 +481,7 @@ type ListPlansOptions struct {
// PlanStore manages schema change plans.
// Plans are created by Plan() and stored for Apply() and staleness detection.
// Both GRPCClient and LocalClient are stateless - SchemaBot owns plan storage.
// Create canonicalizes the provided plan's repository, environment, database, and database type in place before persisting.
type PlanStore interface {
// Create stores a new plan and returns its ID. Returns error if plan_identifier already exists.
Create(ctx context.Context, plan *Plan) (int64, error)
Expand Down Expand Up @@ -1082,6 +1087,10 @@ const SummaryClaimStaleAfter = 2 * time.Minute
// for comments actually posted; minimized_at is set only after the GitHub
// minimize call succeeded, so an unminimized row is always retried by the next
// supersede.
// Insert canonicalizes the provided comment's repository, database, and
// database type in place before persisting. EnvironmentScope is stored as
// given: no query predicate filters on it, and its consumers compare it in Go
// against a scope built from the configured environment names.
type PlanCommentStore interface {
// Insert stores a newly posted plan comment and sets comment.ID.
Insert(ctx context.Context, comment *PlanComment) error
Expand Down
24 changes: 24 additions & 0 deletions pkg/storage/storagetest/plan_comments.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,30 @@ func InsertPlanComment(t *testing.T, store storage.Storage, repo string, pr int,
// behavior — a minimized comment drops out of the unminimized listings and
// repeat or missing-id marks are no-ops.
func TestPlanComments(t *testing.T, h Harness) {
t.Run("CanonicalizesIdentityKeys", func(t *testing.T) {
ctx := t.Context()
store := h.NewStorage(t)

inserted := InsertPlanComment(t, store, "MixedCase/Sample-Repo", 42, "OrdersDB", "MySQL", "Production,Staging", "sha1", 100)
// Stored-value equality is the cross-dialect check that identity keys are canonicalized before persistence.
assert.Equal(t, "mixedcase/sample-repo", inserted.Repository)
assert.Equal(t, "ordersdb", inserted.DatabaseName)
assert.Equal(t, "mysql", inserted.DatabaseType)
// EnvironmentScope is not a query predicate and is compared in Go against
// a scope built from configured environment names, so it is stored as given.
assert.Equal(t, "Production,Staging", inserted.EnvironmentScope)

comments, err := store.PlanComments().ListUnminimizedForSlot(ctx, "MIXEDCASE/SAMPLE-REPO", 42, "ORDERSDB", "MYSQL")
require.NoError(t, err)
require.Len(t, comments, 1)
assert.Equal(t, inserted.ID, comments[0].ID)

comments, err = store.PlanComments().ListUnminimizedForRepoPR(ctx, "MIXEDCASE/SAMPLE-REPO", 42)
require.NoError(t, err)
require.Len(t, comments, 1)
assert.Equal(t, inserted.ID, comments[0].ID)
})

t.Run("Insert_And_ListUnminimizedForSlot", func(t *testing.T) {
ctx := t.Context()
store := h.NewStorage(t)
Expand Down
40 changes: 40 additions & 0 deletions pkg/storage/storagetest/plans.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,46 @@ import (
// storage.ErrNotImplemented. When an implementation lands, it joins this
// family.
func TestPlans(t *testing.T, h Harness) {
t.Run("CanonicalizesIdentityKeys", func(t *testing.T) {
ctx := t.Context()
store := h.NewStorage(t)
plan := &storage.Plan{
PlanIdentifier: "plan_mixed_case",
Database: "OrdersDB",
DatabaseType: "MySQL",
Repository: "MixedCase/Sample-Repo",
PullRequest: 42,
Environment: "Staging",
CreatedAt: time.Now().UTC().Truncate(time.Second),
}
_, err := store.Plans().Create(ctx, plan)
require.NoError(t, err)

// Stored-value equality is the cross-dialect check that identity keys are canonicalized before persistence.
assert.Equal(t, "ordersdb", plan.Database)
assert.Equal(t, "mysql", plan.DatabaseType)
assert.Equal(t, "mixedcase/sample-repo", plan.Repository)
assert.Equal(t, "staging", plan.Environment)

byPR, err := store.Plans().GetByPR(ctx, "MIXEDCASE/SAMPLE-REPO", 42)
require.NoError(t, err)
require.Len(t, byPR, 1)
assert.Equal(t, "plan_mixed_case", byPR[0].PlanIdentifier)

listed, err := store.Plans().List(ctx, storage.ListPlansOptions{
Database: "ORDERSDB", Environment: "STAGING",
Repository: "MIXEDCASE/SAMPLE-REPO", PullRequest: 42, Limit: 10,
})
require.NoError(t, err)
require.Len(t, listed, 1)
assert.Equal(t, "plan_mixed_case", listed[0].PlanIdentifier)

require.NoError(t, store.Plans().DeleteByPR(ctx, "MIXEDCASE/SAMPLE-REPO", 42))
deleted, err := store.Plans().Get(ctx, "plan_mixed_case")
require.NoError(t, err)
assert.Nil(t, deleted)
})

t.Run("Create_And_Get", func(t *testing.T) {
ctx := t.Context()
store := h.NewStorage(t)
Expand Down
30 changes: 30 additions & 0 deletions pkg/storage/storagetest/webhook_events.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,36 @@ func TestWebhookEvents(t *testing.T, h Harness) {
return claimExpecting(t, store, "old")
}

t.Run("CanonicalizesRepositoryKey", func(t *testing.T) {
ctx := t.Context()
store := h.NewStorage(t)

event := newPullRequestEvent("delivery-mixed-case", "synchronize", 42, "mixed-head", time.Time{})
event.Repository = "MixedCase/Sample-Repo"
createEvent(t, store, event)
// Stored-value equality is the cross-dialect check that the repository is canonicalized before persistence.
assert.Equal(t, "mixedcase/sample-repo", event.Repository)

found, err := store.WebhookEvents().HasEventForHead(ctx, storage.WebhookProviderGitHub, "MIXEDCASE/SAMPLE-REPO", 42, "mixed-head")
require.NoError(t, err)
assert.True(t, found)

claimed := claimExpecting(t, store, "delivery-mixed-case")
claimed.Repository = "MIXEDCASE/SAMPLE-REPO"
successor := newPullRequestEvent("delivery-mixed-successor", "synchronize", 42, "new-head", claimed.ReceivedAt.Add(time.Second))
successor.Repository = "mixedcase/sample-repo"
createEvent(t, store, successor)

covered, err := store.WebhookEvents().HasCoveringSuccessor(ctx, claimed)
require.NoError(t, err)
assert.True(t, covered)
assert.Equal(t, "MIXEDCASE/SAMPLE-REPO", claimed.Repository)
superseded, err := store.WebhookEvents().SupersedeIfCovered(ctx, claimed)
require.NoError(t, err)
assert.True(t, superseded)
assert.Equal(t, "MIXEDCASE/SAMPLE-REPO", claimed.Repository)
})

t.Run("Create_DeduplicatesDeliveryID", func(t *testing.T) {
ctx := t.Context()
store := h.NewStorage(t)
Expand Down
Loading