package delivery_test import ( "context" "net/http" "net/http/httptest" "sync/atomic" "testing" "github.com/google/uuid" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/delivery" ) // The two terminal-state gaps of // https://git.eeqj.de/sneak/webhooker/issues/107: a delivery failed // with nothing in its event log to say why, and a retrying delivery // whose target was deleted, which used to keep sending and then never // terminalise. Section 4 is the same deleted-target gap for a pending // delivery: https://git.eeqj.de/sneak/webhooker/issues/293. // tUnknownType is a target type no build implements. It stands in for // a target whose type was written by a build that knew a type this one // does not. const tUnknownType = database.TargetType("pubsub") // tSeedDeletedTarget creates a target, a delivery against it at the // given status with one recorded failed attempt, and then deletes the // target the way the source page does. // // It asserts the delete is soft, because that is the whole reason the // engine could not tell a deleted target from a target id that never // named a row: the surviving row is invisible to a scoped read. func tSeedDeletedTarget( t *testing.T, s iSetup, name, url string, status database.DeliveryStatus, ) string { t.Helper() targetID := uuid.New().String() iCreateTarget( t, s.MainDB, targetID, s.WebhookID, name, database.TargetTypeHTTP, iHTTPConfig(url), 5, ) event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"target":"deleted"}`, ) d := iSeedDelivery( t, s.WebhookDB, event.ID, targetID, status, ) iSeedFailedResult(t, s.WebhookDB, d.ID) require.NoError(t, s.MainDB.Delete( &database.Target{}, "id = ?", targetID, ).Error) var scoped, unscoped int64 require.NoError(t, s.MainDB. Model(&database.Target{}). Where("id = ?", targetID). Count(&scoped).Error) require.NoError(t, s.MainDB.Unscoped(). Model(&database.Target{}). Where("id = ?", targetID). Count(&unscoped).Error) require.Zero(t, scoped, "the deleted target is still visible to a scoped read", ) require.Equal(t, int64(1), unscoped, "the delete was hard, so this test proves nothing about "+ "the soft-delete case it exists for", ) return d.ID } // tLastResult returns a delivery's final recorded attempt, asserting // the expected number of them. func tLastResult( t *testing.T, s iSetup, deliveryID string, want int, ) database.DeliveryResult { t.Helper() results := iResults(t, s.WebhookDB, deliveryID) require.Len(t, results, want) return results[want-1] } // --- 1. A failure with nothing recorded --- func TestProcessDelivery_UnknownTargetType_RecordsWhy( t *testing.T, ) { t.Parallel() s := newISetup(t) targetID := uuid.New().String() event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"unknown":"type"}`, ) seeded := iSeedDelivery( t, s.WebhookDB, event.ID, targetID, database.DeliveryStatusPending, ) target := database.Target{ Name: "mystery", Type: tUnknownType, Config: iHTTPConfig("http://example.com/hook"), } target.ID = targetID d := database.Delivery{ EventID: event.ID, TargetID: targetID, Status: database.DeliveryStatusPending, Event: event, Target: target, } d.ID = seeded.ID body := event.Body task := iTask( seeded, event, s.WebhookID, targetID, "mystery", target.Config, 0, 1, &body, ) task.TargetType = tUnknownType s.Engine.ExportProcessDelivery( context.Background(), s.WebhookDB, &d, &task, ) iAssertStatus( t, s.WebhookDB, d.ID, database.DeliveryStatusFailed, ) last := tLastResult(t, s, d.ID, 1) assert.False(t, last.Success) assert.Equal(t, 1, last.AttemptNum) assert.Contains(t, last.Error, string(tUnknownType), "the recorded reason does not name the offending type", ) } // --- 2. A retrying delivery whose target is gone --- func TestRecoverSingleRetry_TargetDeleted(t *testing.T) { t.Parallel() s := newISetup(t) iCreateWebhook( t, s.MainDB, s.WebhookID, "deleted-target-recovery", ) deliveryID := tSeedDeletedTarget( t, s, "gone-on-recovery", "http://example.com/hook", database.DeliveryStatusRetrying, ) s.Engine.ExportRecoverWebhookDeliveries( context.Background(), s.WebhookID, ) iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusFailed, ) last := tLastResult(t, s, deliveryID, 2) assert.False(t, last.Success) assert.Equal(t, 2, last.AttemptNum) assert.Contains(t, last.Error, "gone-on-recovery") assert.Contains(t, last.Error, "was deleted") assert.Empty(t, s.Engine.ExportRetryCh(), "a delivery whose target is gone was rescheduled", ) assert.Zero(t, s.Engine.ExportInflightHeld(), "the terminal path leaked its ownership reference", ) } func TestSweepSingleRetry_TargetDeleted(t *testing.T) { t.Parallel() s := newISetup(t) iCreateWebhook( t, s.MainDB, s.WebhookID, "deleted-target-sweep", ) deliveryID := tSeedDeletedTarget( t, s, "gone-on-sweep", "http://example.com/hook", database.DeliveryStatusRetrying, ) // Twice, because the bug was an error the sweep repeated every // minute for the life of the database: the second sweep must // find nothing left to do. s.Engine.ExportSweepWebhookRetries( context.Background(), s.WebhookID, ) s.Engine.ExportSweepWebhookRetries( context.Background(), s.WebhookID, ) iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusFailed, ) last := tLastResult(t, s, deliveryID, 2) assert.Contains(t, last.Error, "gone-on-sweep") assert.Contains(t, last.Error, "was deleted") assert.Empty(t, s.Engine.ExportRetryCh()) assert.Zero(t, s.Engine.ExportInflightHeld()) } // TestSweepSingleRetry_TargetNeverExisted covers the other half of the // soft-delete distinction: an id with no row at all, deleted or // otherwise, must not be reported as something the operator deleted. func TestSweepSingleRetry_TargetNeverExisted(t *testing.T) { t.Parallel() s := newISetup(t) iCreateWebhook( t, s.MainDB, s.WebhookID, "target-never-existed", ) targetID := uuid.New().String() event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"target":"absent"}`, ) d := iSeedDelivery( t, s.WebhookDB, event.ID, targetID, database.DeliveryStatusRetrying, ) iSeedFailedResult(t, s.WebhookDB, d.ID) s.Engine.ExportSweepWebhookRetries( context.Background(), s.WebhookID, ) iAssertStatus( t, s.WebhookDB, d.ID, database.DeliveryStatusFailed, ) last := tLastResult(t, s, d.ID, 2) assert.Contains(t, last.Error, targetID) assert.Contains(t, last.Error, "no longer exists") assert.NotContains(t, last.Error, "was deleted", "an id that never named a row was reported as a deletion", ) } // TestFailMissingTarget_WritesNoTargetRow holds the new terminal path // to the same rule as the existing one: no target row, and so no // plaintext target config, may be written into the per-webhook event // database. See https://git.eeqj.de/sneak/webhooker/issues/206. func TestFailMissingTarget_WritesNoTargetRow( t *testing.T, ) { t.Parallel() s := newISetup(t) iCreateWebhook( t, s.MainDB, s.WebhookID, "no-target-row-deleted", ) hookURL := "https://hooks.slack.com/services/T00/B00/x" deliveryID := tSeedDeletedTarget( t, s, "credential-bearing", hookURL, database.DeliveryStatusRetrying, ) s.Engine.ExportSweepWebhookRetries( context.Background(), s.WebhookID, ) iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusFailed, ) var configs []string require.NoError(t, s.WebhookDB. Table("targets"). Pluck("config", &configs).Error) assert.Empty(t, configs, "the deleted-target terminal path wrote a target row "+ "into the per-webhook event database", ) } // --- 3. The scheduled retry chain --- // tRetryChainSetup wires a counting sink and a retrying delivery // against a live target pointing at it, and returns the task a // scheduled retry would carry — config and all, snapshotted as // ScheduleRetry snapshots it. func tRetryChainSetup( t *testing.T, s iSetup, name string, hits *atomic.Int64, ) (delivery.Task, string) { t.Helper() ts := httptest.NewServer(http.HandlerFunc( func(w http.ResponseWriter, _ *http.Request) { hits.Add(1) w.WriteHeader(http.StatusOK) }, )) t.Cleanup(ts.Close) iCreateWebhook(t, s.MainDB, s.WebhookID, name) targetID := uuid.New().String() cfg := iHTTPConfig(ts.URL) iCreateTarget( t, s.MainDB, targetID, s.WebhookID, name, database.TargetTypeHTTP, cfg, 5, ) event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"chain":"retry"}`, ) d := iSeedDelivery( t, s.WebhookDB, event.ID, targetID, database.DeliveryStatusRetrying, ) iSeedFailedResult(t, s.WebhookDB, d.ID) body := event.Body return iTask( d, event, s.WebhookID, targetID, name, cfg, 5, 2, &body, ), targetID } // TestProcessRetryTask_TargetDeleted_MakesNoAttempt is the half the // deployability audit found worse than filed: terminalising on // recovery and sweep alone leaves the already-scheduled timer chain // running, and it holds the target's configuration from before the // deletion, so it goes on sending to a destination that was removed. func TestProcessRetryTask_TargetDeleted_MakesNoAttempt( t *testing.T, ) { t.Parallel() s := newISetup(t) var hits atomic.Int64 task, targetID := tRetryChainSetup( t, s, "gone-mid-chain", &hits, ) require.NoError(t, s.MainDB.Delete( &database.Target{}, "id = ?", targetID, ).Error) s.Engine.ExportProcessRetryTask( context.Background(), &task, ) assert.Zero(t, hits.Load(), "a scheduled retry fired at a target the operator "+ "had already deleted", ) iAssertStatus( t, s.WebhookDB, task.DeliveryID, database.DeliveryStatusFailed, ) last := tLastResult(t, s, task.DeliveryID, 2) assert.False(t, last.Success) assert.Contains(t, last.Error, "was deleted") assert.Zero(t, s.Engine.ExportInflightHeld()) } // TestProcessRetryTask_TargetPresent_StillDelivers is the guard's // mutation check: a liveness check that refused every retry would pass // the test above and break every retry there is. func TestProcessRetryTask_TargetPresent_StillDelivers( t *testing.T, ) { t.Parallel() s := newISetup(t) var hits atomic.Int64 task, _ := tRetryChainSetup(t, s, "still-there", &hits) s.Engine.ExportProcessRetryTask( context.Background(), &task, ) assert.Equal(t, int64(1), hits.Load()) iAssertStatus( t, s.WebhookDB, task.DeliveryID, database.DeliveryStatusDelivered, ) } // TestProcessRetryTask_TargetUnreadable_StillDelivers pins the other // half of the guard: only a target that is confirmed gone stops a // retry. A main database that cannot be read is a transient fault, and // a guard that abandoned deliveries on one would be a worse bug than // the one it fixes. func TestProcessRetryTask_TargetUnreadable_StillDelivers( t *testing.T, ) { t.Parallel() s := newISetup(t) var hits atomic.Int64 task, _ := tRetryChainSetup(t, s, "unreadable-main", &hits) sqlDB, err := s.MainDB.DB() require.NoError(t, err) require.NoError(t, sqlDB.Close()) s.Engine.ExportProcessRetryTask( context.Background(), &task, ) assert.Equal(t, int64(1), hits.Load(), "a retry was abandoned because the main database "+ "could not be read, not because its target was gone", ) iAssertStatus( t, s.WebhookDB, task.DeliveryID, database.DeliveryStatusDelivered, ) } // TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone is the // same rule on the recovery path. A read failure that is not // "record not found" must leave every retrying delivery of every // webhook exactly as it was. func TestRecoverSingleRetry_TargetUnreadable_LeavesDeliveryAlone( t *testing.T, ) { t.Parallel() s := newISetup(t) iCreateWebhook( t, s.MainDB, s.WebhookID, "unreadable-on-recovery", ) targetID := uuid.New().String() iCreateTarget( t, s.MainDB, targetID, s.WebhookID, "healthy", database.TargetTypeHTTP, iHTTPConfig("http://example.com/hook"), 5, ) event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"still":"retrying"}`, ) d := iSeedDelivery( t, s.WebhookDB, event.ID, targetID, database.DeliveryStatusRetrying, ) iSeedFailedResult(t, s.WebhookDB, d.ID) sqlDB, err := s.MainDB.DB() require.NoError(t, err) require.NoError(t, sqlDB.Close()) s.Engine.ExportRecoverRetryingDeliveries( s.WebhookDB, s.WebhookID, ) iAssertStatus( t, s.WebhookDB, d.ID, database.DeliveryStatusRetrying, ) assert.Len(t, iResults(t, s.WebhookDB, d.ID), 1, "an unreadable main database produced a terminal "+ "failure row", ) assert.Zero(t, s.Engine.ExportInflightHeld()) } // --- 4. A pending delivery whose target is gone --- func TestRecoverPending_TargetDeleted(t *testing.T) { t.Parallel() s := newISetup(t) deliveryID := tSeedDeletedTarget( t, s, "gone-while-pending", "http://example.com/hook", database.DeliveryStatusPending, ) s.Engine.ExportRecoverWebhookDeliveries( context.Background(), s.WebhookID, ) iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusFailed, ) last := tLastResult(t, s, deliveryID, 2) assert.False(t, last.Success) assert.Equal(t, 2, last.AttemptNum) assert.Contains(t, last.Error, "gone-while-pending") assert.Contains(t, last.Error, "was deleted") assert.Empty(t, fDrain(s.Engine), "a delivery whose target is gone was sent", ) assert.Zero(t, s.Engine.ExportInflightHeld(), "the terminal path leaked its ownership reference", ) } // TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone: the // terminal write takes ownership like every other recovery write, so a // delivery the engine still holds is not failed underneath its worker. func TestRecoverPending_TargetDeleted_LeavesAnOwnedDeliveryAlone( t *testing.T, ) { t.Parallel() s := newISetup(t) deliveryID := tSeedDeletedTarget( t, s, "gone-but-owned", "http://example.com/hook", database.DeliveryStatusPending, ) require.True(t, s.Engine.ExportRetainDelivery(deliveryID)) s.Engine.ExportRecoverWebhookDeliveries( context.Background(), s.WebhookID, ) iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusPending, ) assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1, "a delivery the engine owns was failed underneath it", ) } // TestFailMissingTarget_LeavesASettledDeliveryAlone: the recovery paths // read their batch before taking ownership, and a worker may send a // delivery and let it go in between. The terminal write goes by the row // as it is now, not as the batch read it. func TestFailMissingTarget_LeavesASettledDeliveryAlone( t *testing.T, ) { t.Parallel() s := newISetup(t) deliveryID := tSeedDeletedTarget( t, s, "gone-after-sending", "http://example.com/hook", database.DeliveryStatusPending, ) var batch database.Delivery require.NoError(t, s.WebhookDB.First( &batch, "id = ?", deliveryID, ).Error) // A worker settles the delivery after the batch was read. require.NoError(t, s.WebhookDB.Model(&database.Delivery{}). Where("id = ?", deliveryID). Update("status", database.DeliveryStatusDelivered).Error) s.Engine.ExportFailMissingTarget( s.WebhookDB, s.WebhookID, &batch, ) iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusDelivered, ) assert.Len(t, iResults(t, s.WebhookDB, deliveryID), 1, "a delivery settled after the batch read was then failed", ) assert.Zero(t, s.Engine.ExportInflightHeld()) } // TestSweepPending_TargetDeleted sweeps twice over a batch that also // holds a healthy stranded delivery. The one whose target is gone is // failed once and then left alone; the healthy one is queued by the // first sweep and not again by the second. func TestSweepPending_TargetDeleted(t *testing.T) { t.Parallel() liveTargetID := uuid.New().String() s := fSweepSetup(t, liveTargetID, "still-there") deliveryID := tSeedDeletedTarget( t, s, "gone-on-pending-sweep", "http://example.com/hook", database.DeliveryStatusPending, ) rAgePending(t, s.WebhookDB, deliveryID) event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"target":"live"}`, ) healthy := iSeedDelivery( t, s.WebhookDB, event.ID, liveTargetID, database.DeliveryStatusPending, ) rAgePending(t, s.WebhookDB, healthy.ID) ctx := context.Background() s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID) tasks := fDrain(s.Engine) require.Len(t, tasks, 1, "the first sweep did not queue the healthy delivery", ) assert.Equal(t, healthy.ID, tasks[0].DeliveryID) s.Engine.ExportSweepWebhookRetries(ctx, s.WebhookID) assert.Empty(t, fDrain(s.Engine), "the second sweep queued a delivery again", ) iAssertStatus( t, s.WebhookDB, deliveryID, database.DeliveryStatusFailed, ) last := tLastResult(t, s, deliveryID, 2) assert.Contains(t, last.Error, "gone-on-pending-sweep") assert.Contains(t, last.Error, "was deleted") } // TestSendRecoveredDeliveries_TargetMissingFromMap: the batch's target // map is empty when its query failed, so every delivery in the batch is // looked up on its own. A healthy one is sent to the target that lookup // finds. func TestSendRecoveredDeliveries_TargetMissingFromMap( t *testing.T, ) { t.Parallel() s := newISetup(t) targetID := uuid.New().String() iCreateTarget( t, s.MainDB, targetID, s.WebhookID, "found-on-lookup", database.TargetTypeLog, "", 0, ) event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"map":"empty"}`, ) d := iSeedDelivery( t, s.WebhookDB, event.ID, targetID, database.DeliveryStatusPending, ) s.Engine.ExportSendRecoveredDeliveries( context.Background(), s.WebhookDB, []database.Delivery{d}, s.WebhookID, map[string]database.Target{}, nil, ) tasks := fDrain(s.Engine) require.Len(t, tasks, 1, "the healthy delivery was not queued exactly once", ) assert.Equal(t, d.ID, tasks[0].DeliveryID) assert.Equal(t, targetID, tasks[0].TargetID) assert.Equal(t, database.TargetTypeLog, tasks[0].TargetType) iAssertStatus( t, s.WebhookDB, d.ID, database.DeliveryStatusPending, ) } // TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone: a failed // read of the main database is not a deleted target. Restart recovery // holds every pending delivery of the webhook in one batch, so failing // on this would fail all of them. func TestRecoverPending_TargetUnreadable_LeavesDeliveryAlone( t *testing.T, ) { t.Parallel() s := newISetup(t) targetID := uuid.New().String() iCreateTarget( t, s.MainDB, targetID, s.WebhookID, "healthy", database.TargetTypeLog, "", 0, ) event := iSeedEvent( t, s.WebhookDB, s.WebhookID, `{"still":"pending"}`, ) d := iSeedDelivery( t, s.WebhookDB, event.ID, targetID, database.DeliveryStatusPending, ) sqlDB, err := s.MainDB.DB() require.NoError(t, err) require.NoError(t, sqlDB.Close()) s.Engine.ExportRecoverPendingDeliveries( context.Background(), s.WebhookDB, s.WebhookID, ) iAssertStatus( t, s.WebhookDB, d.ID, database.DeliveryStatusPending, ) assert.Empty(t, iResults(t, s.WebhookDB, d.ID), "an unreadable main database produced a terminal "+ "failure row", ) assert.Empty(t, fDrain(s.Engine)) assert.Zero(t, s.Engine.ExportInflightHeld()) }