Files
Test 60f9ca2b1d feat(T1.1): implement error recovery, retry policies, and deadletter handling
- Add internal/recovery package with comprehensive error recovery infrastructure
- Implement RetryPolicy with exponential backoff
- Three predefined policies: DefaultRetryPolicy, ActivityRetryPolicy, LLMActivityRetryPolicy
- Integrate with Temporal SDK via ToTemporalRetryPolicy()
- Implement DeadletterQueue for tracking permanently failed activities
- Thread-safe deadletter operations with JSON persistence
- Mark items as recoverable or non-recoverable
- Support batch retrieval of recoverable items
- Implement CheckpointManager for periodic state snapshots
- Track workflow stages and task lifecycle (completed/pending/failed)
- Persist checkpoints to enable recovery after crashes
- Add OrchestratorWorkflowWithRecovery demonstrating recovery patterns
- Structured logging at each workflow step
- Retry policies applied to all activity types
- Extended ActivityTuning with retry configuration fields

Test Coverage:
- 8/8 retry policy tests passing
- 10/10 deadletter queue tests passing
- 10/10 checkpoint manager tests passing
- 40 total recovery tests, all passing
- All existing tests continue to pass

Key Features:
- Exponential backoff prevents thundering herd
- Deadletter audit trail with timestamps
- Checkpoint interval configurable (30s default)
- Thread-safe concurrent access
- No external dependencies added

Closes T1.1
2026-08-23 16:43:30 -07:00

201 lines
4.6 KiB
Go

package recovery
import (
"errors"
"path/filepath"
"testing"
"github.com/stretchr/testify/assert"
)
func TestDeadletterQueue(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
dq := NewDeadletterQueue(queuePath)
item := &DeadletterItem{
ID: "task-1",
Type: "activity",
WorkflowID: "wf-1",
Error: "test error",
AttemptCount: 1,
MaxAttempts: 3,
Recoverable: true,
}
// Add item
err := dq.Add(item)
assert.NoError(t, err)
assert.Equal(t, 1, dq.Count())
// Get item
retrieved := dq.Get("task-1")
assert.NotNil(t, retrieved)
assert.Equal(t, "task-1", retrieved.ID)
assert.NotZero(t, retrieved.CreatedAt)
assert.NotZero(t, retrieved.UpdatedAt)
// Remove item
err = dq.Remove("task-1")
assert.NoError(t, err)
assert.Equal(t, 0, dq.Count())
}
func TestDeadletterQueuePersistence(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
// Create and add item
dq1 := NewDeadletterQueue(queuePath)
item := &DeadletterItem{
ID: "task-1",
Type: "activity",
WorkflowID: "wf-1",
Error: "test error",
Recoverable: true,
}
err := dq1.Add(item)
assert.NoError(t, err)
// Create new queue instance and load
dq2 := NewDeadletterQueue(queuePath)
err = dq2.Load()
assert.NoError(t, err)
// Verify item was loaded
assert.Equal(t, 1, dq2.Count())
retrieved := dq2.Get("task-1")
assert.NotNil(t, retrieved)
assert.Equal(t, "task-1", retrieved.ID)
}
func TestDeadletterQueueGetAll(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
dq := NewDeadletterQueue(queuePath)
// Add multiple items
for i := 1; i <= 3; i++ {
item := &DeadletterItem{
ID: "task-" + string(rune(48+i)),
Type: "activity",
WorkflowID: "wf-1",
Error: "error",
}
dq.Add(item)
}
all := dq.GetAll()
assert.Equal(t, 3, len(all))
}
func TestDeadletterQueueGetRecoverable(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
dq := NewDeadletterQueue(queuePath)
// Add recoverable item
dq.Add(&DeadletterItem{
ID: "task-1",
Type: "activity",
WorkflowID: "wf-1",
Recoverable: true,
})
// Add non-recoverable item
dq.Add(&DeadletterItem{
ID: "task-2",
Type: "activity",
WorkflowID: "wf-1",
Recoverable: false,
})
recoverable := dq.GetRecoverable()
assert.Equal(t, 1, len(recoverable))
assert.Equal(t, "task-1", recoverable[0].ID)
}
func TestDeadletterQueueResolve(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
dq := NewDeadletterQueue(queuePath)
dq.Add(&DeadletterItem{
ID: "task-1",
Type: "activity",
WorkflowID: "wf-1",
})
// Resolve item
err := dq.Resolve("task-1", "manually recovered")
assert.NoError(t, err)
item := dq.Get("task-1")
assert.NotNil(t, item)
assert.Equal(t, "manually recovered", item.RecoveryNote)
}
func TestDeadletterQueueEmpty(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
dq := NewDeadletterQueue(queuePath)
assert.True(t, dq.IsEmpty())
assert.Equal(t, 0, dq.Count())
dq.Add(&DeadletterItem{ID: "task-1"})
assert.False(t, dq.IsEmpty())
assert.Equal(t, 1, dq.Count())
}
func TestCreateDeadletterItem(t *testing.T) {
err := errors.New("test error")
data := map[string]any{"key": "value"}
item := CreateDeadletterItem("task-1", "activity", "wf-1", err, data, true)
assert.Equal(t, "task-1", item.ID)
assert.Equal(t, "activity", item.Type)
assert.Equal(t, "wf-1", item.WorkflowID)
assert.Equal(t, "test error", item.Error)
assert.Equal(t, 1, item.AttemptCount)
assert.Equal(t, 3, item.MaxAttempts)
assert.True(t, item.Recoverable)
assert.NotZero(t, item.CreatedAt)
assert.NotZero(t, item.UpdatedAt)
}
func TestDeadletterQueueNoFile(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "nonexistent.json")
dq := NewDeadletterQueue(queuePath)
// Loading non-existent file should not error
err := dq.Load()
assert.NoError(t, err)
assert.True(t, dq.IsEmpty())
}
func TestDeadletterRemoveNonexistent(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
dq := NewDeadletterQueue(queuePath)
// Removing non-existent item should not error
err := dq.Remove("nonexistent")
assert.NoError(t, err)
}
func TestDeadletterResolveNonexistent(t *testing.T) {
tmpDir := t.TempDir()
queuePath := filepath.Join(tmpDir, "deadletter.json")
dq := NewDeadletterQueue(queuePath)
// Resolving non-existent item should error
err := dq.Resolve("nonexistent", "note")
assert.Error(t, err)
}