From e6fdb4252b2302e8574704d36510f455024eb011 Mon Sep 17 00:00:00 2001 From: Cristian Magherusan-Stanciu Date: Fri, 22 May 2026 15:00:45 +0200 Subject: [PATCH] fix(execution): recover from panics inside FanOutWithConcurrency goroutines (closes #669) Defence-in-depth for the per-account / per-rec fan-out helper. Before this change, any panic inside fn (nil deref, type assertion failure, slice OOB in a future refactor) propagated up the unrecovered goroutine and crashed the entire Lambda process. Consequences: - Purchase execution row stays at 'approved' with no transition to 'failed' (state machine never runs the aggregator that would record the failure). - Lambda invocation terminates abnormally; CloudWatch may not flush the final log lines. - User sees a generic Lambda invocation error rather than the actual panic detail. The fan-out helper now installs a deferred recover() in every goroutine that converts a panic into a structured Err on the per-item Result slot, logs the goroutine stack at Error level for post-mortem, and lets the parent aggregator process the failure exactly like a normal fn-returned error. Current call sites in internal/purchase/execution.go (per-account executeForAccount + per-rec processPurchaseRecommendations) have no obvious panic source today; the recover() is cheap insurance for the high-impact failure mode and covers future refactors that might introduce one. Regression test (TestFanOut_PanicInFn) panics inside fn for one of three items, asserts the panicking item carries a non-nil Err containing the panic value AND that the other two items still succeed (i.e. one goroutine's panic doesn't cascade across the fan-out). Closes #669 --- internal/execution/fanout.go | 23 +++++++++++++++++++++ internal/execution/fanout_test.go | 34 +++++++++++++++++++++++++++++++ 2 files changed, 57 insertions(+) diff --git a/internal/execution/fanout.go b/internal/execution/fanout.go index 3d3763dad..9982cb976 100644 --- a/internal/execution/fanout.go +++ b/internal/execution/fanout.go @@ -3,9 +3,13 @@ package execution import ( "context" + "fmt" "os" + "runtime" "strconv" "sync" + + "github.com/LeanerCloud/CUDly/pkg/logging" ) // ConcurrencyFromEnv reads the CUDLY_MAX_ACCOUNT_PARALLELISM env var and @@ -69,6 +73,25 @@ func FanOutWithConcurrency[T any]( go func(idx int, accountID string) { defer wg.Done() defer func() { <-sem }() // release slot + // recover() guard so a panic inside fn (nil deref, type + // assertion failure, slice OOB, etc.) doesn't crash the + // whole Lambda process and strand the purchase execution + // at 'approved' with no way to debug post-mortem (issue + // #669). Surface the panic as an Err on the result slot + // so the parent aggregator records it on the execution + // row exactly like a regular fn-returned error, and log + // the goroutine stack at Error level for diagnosis. + defer func() { + if r := recover(); r != nil { + buf := make([]byte, 4096) + n := runtime.Stack(buf, false) + logging.Errorf("fan-out goroutine panic (account=%s): %v\n%s", accountID, r, buf[:n]) + results[idx] = Result[T]{ + AccountID: accountID, + Err: fmt.Errorf("panic during fan-out (account=%s): %v", accountID, r), + } + } + }() val, err := fn(ctx, accountID) results[idx] = Result[T]{AccountID: accountID, Value: val, Err: err} }(i, id) diff --git a/internal/execution/fanout_test.go b/internal/execution/fanout_test.go index 993473957..485888797 100644 --- a/internal/execution/fanout_test.go +++ b/internal/execution/fanout_test.go @@ -81,3 +81,37 @@ func TestPartition(t *testing.T) { require.Len(t, failures, 1) assert.Equal(t, "b", failures[0].AccountID) } + +// TestFanOut_PanicInFn asserts that a panic inside the per-item function is +// caught by the goroutine's deferred recover, converted to an Err on that +// item's result slot, and does NOT propagate up and crash the whole process +// (which would strand the surrounding purchase execution at 'approved' and +// terminate the Lambda invocation abnormally — see #669). +func TestFanOut_PanicInFn(t *testing.T) { + ids := []string{"a", "panic-me", "c"} + results := FanOut(context.Background(), ids, func(ctx context.Context, id string) (string, error) { + if id == "panic-me" { + panic("synthetic panic for test") + } + return "ok:" + id, nil + }) + + require.Len(t, results, 3) + // Map by AccountID since FanOut order is non-deterministic. + byID := make(map[string]Result[string], len(results)) + for _, r := range results { + byID[r.AccountID] = r + } + + // Non-panicking items still succeed. + assert.NoError(t, byID["a"].Err) + assert.Equal(t, "ok:a", byID["a"].Value) + assert.NoError(t, byID["c"].Err) + assert.Equal(t, "ok:c", byID["c"].Value) + + // The panicking item is surfaced as an Err containing the panic value. + require.Error(t, byID["panic-me"].Err) + assert.Contains(t, byID["panic-me"].Err.Error(), "panic during fan-out") + assert.Contains(t, byID["panic-me"].Err.Error(), "synthetic panic for test") + assert.Contains(t, byID["panic-me"].Err.Error(), "panic-me") +}