207 lines
4.7 KiB
Go
207 lines
4.7 KiB
Go
package recovery
|
|||
|
|
|
||
|
|
import (
|
||
|
|
"encoding/json"
|
||
|
|
"fmt"
|
||
|
|
"os"
|
||
|
|
"path/filepath"
|
||
|
|
"sync"
|
||
|
|
"time"
|
||
|
|
)
|
||
|
|
|
||
|
|
// DeadletterItem represents a failed activity/task
|
||
|
|
type DeadletterItem struct {
|
||
|
|
ID string `json:"id"`
|
||
|
|
Type string `json:"type"` // "activity", "task", "workflow"
|
||
|
|
WorkflowID string `json:"workflow_id"`
|
||
|
|
Error string `json:"error"`
|
||
|
|
LastAttempt time.Time `json:"last_attempt"`
|
||
|
|
AttemptCount int `json:"attempt_count"`
|
||
|
|
MaxAttempts int `json:"max_attempts"`
|
||
|
|
Data any `json:"data"` // Original input
|
||
|
|
Recoverable bool `json:"recoverable"`
|
||
|
|
RecoveryNote string `json:"recovery_note"`
|
||
|
|
CreatedAt time.Time `json:"created_at"`
|
||
|
|
UpdatedAt time.Time `json:"updated_at"`
|
||
|
|
}
|
||
|
|
|
||
|
|
// DeadletterQueue manages deadlettered items
|
||
|
|
type DeadletterQueue struct {
|
||
|
|
mu sync.RWMutex
|
||
|
|
path string
|
||
|
|
items map[string]*DeadletterItem
|
||
|
|
}
|
||
|
|
|
||
|
|
// NewDeadletterQueue creates a new deadletter queue
|
||
|
|
func NewDeadletterQueue(path string) *DeadletterQueue {
|
||
|
|
return &DeadletterQueue{
|
||
|
|
path: path,
|
||
|
|
items: make(map[string]*DeadletterItem),
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// Add adds an item to the deadletter queue
|
||
|
|
func (dq *DeadletterQueue) Add(item *DeadletterItem) error {
|
||
|
|
if item.ID == "" {
|
||
|
|
return fmt.Errorf("deadletter item must have an ID")
|
||
|
|
}
|
||
|
|
|
||
|
|
dq.mu.Lock()
|
||
|
|
defer dq.mu.Unlock()
|
||
|
|
|
||
|
|
now := time.Now()
|
||
|
|
if item.CreatedAt.IsZero() {
|
||
|
|
item.CreatedAt = now
|
||
|
|
}
|
||
|
|
item.UpdatedAt = now
|
||
|
|
|
||
|
|
dq.items[item.ID] = item
|
||
|
|
|
||
|
|
// Persist to disk
|
||
|
|
return dq.persistLocked()
|
||
|
|
}
|
||
|
|
|
||
|
|
// Get retrieves an item from the deadletter queue
|
||
|
|
func (dq *DeadletterQueue) Get(id string) *DeadletterItem {
|
||
|
|
dq.mu.RLock()
|
||
|
|
defer dq.mu.RUnlock()
|
||
|
|
|
||
|
|
return dq.items[id]
|
||
|
|
}
|
||
|
|
|
||
|
|
// GetAll returns all deadletter items
|
||
|
|
func (dq *DeadletterQueue) GetAll() []*DeadletterItem {
|
||
|
|
dq.mu.RLock()
|
||
|
|
defer dq.mu.RUnlock()
|
||
|
|
|
||
|
|
items := make([]*DeadletterItem, 0, len(dq.items))
|
||
|
|
for _, item := range dq.items {
|
||
|
|
items = append(items, item)
|
||
|
|
}
|
||
|
|
return items
|
||
|
|
}
|
||
|
|
|
||
|
|
// GetRecoverable returns all recoverable items
|
||
|
|
func (dq *DeadletterQueue) GetRecoverable() []*DeadletterItem {
|
||
|
|
dq.mu.RLock()
|
||
|
|
defer dq.mu.RUnlock()
|
||
|
|
|
||
|
|
items := make([]*DeadletterItem, 0)
|
||
|
|
for _, item := range dq.items {
|
||
|
|
if item.Recoverable {
|
||
|
|
items = append(items, item)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
return items
|
||
|
|
}
|
||
|
|
|
||
|
|
// Remove removes an item from the deadletter queue
|
||
|
|
func (dq *DeadletterQueue) Remove(id string) error {
|
||
|
|
dq.mu.Lock()
|
||
|
|
defer dq.mu.Unlock()
|
||
|
|
|
||
|
|
delete(dq.items, id)
|
||
|
|
return dq.persistLocked()
|
||
|
|
}
|
||
|
|
|
||
|
|
// Resolve marks an item as resolved
|
||
|
|
func (dq *DeadletterQueue) Resolve(id string, note string) error {
|
||
|
|
dq.mu.Lock()
|
||
|
|
defer dq.mu.Unlock()
|
||
|
|
|
||
|
|
item, exists := dq.items[id]
|
||
|
|
if !exists {
|
||
|
|
return fmt.Errorf("item not found: %s", id)
|
||
|
|
}
|
||
|
|
|
||
|
|
item.RecoveryNote = note
|
||
|
|
item.UpdatedAt = time.Now()
|
||
|
|
|
||
|
|
// Don't actually delete, just mark as recovered
|
||
|
|
// This maintains audit trail
|
||
|
|
return dq.persistLocked()
|
||
|
|
}
|
||
|
|
|
||
|
|
// Load loads deadletter queue from disk
|
||
|
|
func (dq *DeadletterQueue) Load() error {
|
||
|
|
dq.mu.Lock()
|
||
|
|
defer dq.mu.Unlock()
|
||
|
|
|
||
|
|
// Create directory if it doesn't exist
|
||
|
|
if err := os.MkdirAll(filepath.Dir(dq.path), 0755); err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
// If file doesn't exist, that's OK (queue is empty)
|
||
|
|
data, err := os.ReadFile(dq.path)
|
||
|
|
if err != nil {
|
||
|
|
if os.IsNotExist(err) {
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
var items []*DeadletterItem
|
||
|
|
if err := json.Unmarshal(data, &items); err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
dq.items = make(map[string]*DeadletterItem)
|
||
|
|
for _, item := range items {
|
||
|
|
dq.items[item.ID] = item
|
||
|
|
}
|
||
|
|
|
||
|
|
return nil
|
||
|
|
}
|
||
|
|
|
||
|
|
// persistLocked persists the queue to disk (must be called with lock held)
|
||
|
|
func (dq *DeadletterQueue) persistLocked() error {
|
||
|
|
items := make([]*DeadletterItem, 0, len(dq.items))
|
||
|
|
for _, item := range dq.items {
|
||
|
|
items = append(items, item)
|
||
|
|
}
|
||
|
|
|
||
|
|
data, err := json.MarshalIndent(items, "", " ")
|
||
|
|
if err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
// Create directory if it doesn't exist
|
||
|
|
if err := os.MkdirAll(filepath.Dir(dq.path), 0755); err != nil {
|
||
|
|
return err
|
||
|
|
}
|
||
|
|
|
||
|
|
return os.WriteFile(dq.path, data, 0644)
|
||
|
|
}
|
||
|
|
|
||
|
|
// Count returns the number of items in the queue
|
||
|
|
func (dq *DeadletterQueue) Count() int {
|
||
|
|
dq.mu.RLock()
|
||
|
|
defer dq.mu.RUnlock()
|
||
|
|
return len(dq.items)
|
||
|
|
}
|
||
|
|
|
||
|
|
// IsEmpty checks if the queue is empty
|
||
|
|
func (dq *DeadletterQueue) IsEmpty() bool {
|
||
|
|
dq.mu.RLock()
|
||
|
|
defer dq.mu.RUnlock()
|
||
|
|
return len(dq.items) == 0
|
||
|
|
}
|
||
|
|
|
||
|
|
// CreateDeadletterItem creates a new deadletter item from an error
|
||
|
|
func CreateDeadletterItem(id, itemType, workflowID string, err error, data any, recoverable bool) *DeadletterItem {
|
||
|
|
return &DeadletterItem{
|
||
|
|
ID: id,
|
||
|
|
Type: itemType,
|
||
|
|
WorkflowID: workflowID,
|
||
|
|
Error: err.Error(),
|
||
|
|
LastAttempt: time.Now(),
|
||
|
|
AttemptCount: 1,
|
||
|
|
MaxAttempts: 3,
|
||
|
|
Data: data,
|
||
|
|
Recoverable: recoverable,
|
||
|
|
CreatedAt: time.Now(),
|
||
|
|
UpdatedAt: time.Now(),
|
||
|
|
}
|
||
|
|
}
|