From 5f5bcba2b20e252b144a542f5fe3fa6051dbece6 Mon Sep 17 00:00:00 2001 From: Max Date: Fri, 27 Feb 2026 18:19:57 +0800 Subject: [PATCH] Refactor delivery event handling and payload structure for improved clarity and functionality - Update DeliveryPayload to use structured types for Content and Preferences, enhancing type safety and readability. - Modify tests to reflect changes in payload structure, ensuring proper serialization and deserialization of delivery content. - Implement a new robotHandler for processing delivery events, streamlining the handling of different delivery channels (email, webhook, process). - Remove the deprecated DeliveryCenter, consolidating delivery logic within the new handler for better maintainability. - Enhance error handling and logging during delivery processing to improve observability and debugging capabilities. --- agent/robot/events/event_push_test.go | 21 +- agent/robot/events/events.go | 15 +- agent/robot/events/events_test.go | 23 +- agent/robot/events/handlers.go | 440 ++++++++- agent/robot/events/handlers_test.go | 187 +++- agent/robot/executor/standard/delivery.go | 138 +-- .../executor/standard/delivery_center.go | 444 --------- .../robot/executor/standard/delivery_test.go | 884 +----------------- 8 files changed, 586 insertions(+), 1566 deletions(-) delete mode 100644 agent/robot/executor/standard/delivery_center.go diff --git a/agent/robot/events/event_push_test.go b/agent/robot/events/event_push_test.go index 4ef496b8..473d0473 100644 --- a/agent/robot/events/event_push_test.go +++ b/agent/robot/events/event_push_test.go @@ -6,6 +6,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + robottypes "github.com/yaoapp/yao/agent/robot/types" ) // EP1: ExecPayload with all execution statuses @@ -76,18 +77,14 @@ func TestTaskPayloadErrorSerialization(t *testing.T) { // EP4: DeliveryPayload with nested content func TestDeliveryPayloadNestedContent(t *testing.T) { - content := map[string]interface{}{ - "report": map[string]interface{}{ - "title": "Daily Summary", - "sections": []interface{}{"intro", "body", "conclusion"}, - }, - } - payload := DeliveryPayload{ ExecutionID: "exec-ep4", MemberID: "member-ep4", TeamID: "team-ep4", - Content: content, + Content: &robottypes.DeliveryContent{ + Summary: "Daily Summary", + Body: "Full body with sections: intro, body, conclusion", + }, } data, err := json.Marshal(payload) @@ -97,11 +94,9 @@ func TestDeliveryPayloadNestedContent(t *testing.T) { err = json.Unmarshal(data, &parsed) require.NoError(t, err) - contentMap, ok := parsed.Content.(map[string]interface{}) - require.True(t, ok) - report, ok := contentMap["report"].(map[string]interface{}) - require.True(t, ok) - assert.Equal(t, "Daily Summary", report["title"]) + require.NotNil(t, parsed.Content) + assert.Equal(t, "Daily Summary", parsed.Content.Summary) + assert.Contains(t, parsed.Content.Body, "sections") } // EP5: Event constants follow naming convention diff --git a/agent/robot/events/events.go b/agent/robot/events/events.go index 492ddfcd..f1a7689a 100644 --- a/agent/robot/events/events.go +++ b/agent/robot/events/events.go @@ -1,5 +1,7 @@ package events +import robottypes "github.com/yaoapp/yao/agent/robot/types" + // Robot event type constants for event.Push integration. // Events are fire-and-forget; handlers are registered via event.Register(). const ( @@ -46,11 +48,10 @@ type TaskPayload struct { // DeliveryPayload is the event payload for Delivery events. type DeliveryPayload struct { - ExecutionID string `json:"execution_id"` - MemberID string `json:"member_id"` - TeamID string `json:"team_id"` - ChatID string `json:"chat_id,omitempty"` - Result interface{} `json:"result,omitempty"` - Content interface{} `json:"content,omitempty"` // DeliveryContent from agent - Preferences interface{} `json:"preferences,omitempty"` // DeliveryPreferences for routing + ExecutionID string `json:"execution_id"` + MemberID string `json:"member_id"` + TeamID string `json:"team_id"` + ChatID string `json:"chat_id,omitempty"` + Content *robottypes.DeliveryContent `json:"content,omitempty"` + Preferences *robottypes.DeliveryPreferences `json:"preferences,omitempty"` } diff --git a/agent/robot/events/events_test.go b/agent/robot/events/events_test.go index fb6f77a1..5da71539 100644 --- a/agent/robot/events/events_test.go +++ b/agent/robot/events/events_test.go @@ -6,6 +6,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + robottypes "github.com/yaoapp/yao/agent/robot/types" ) func TestEventConstants(t *testing.T) { @@ -121,8 +122,16 @@ func TestDeliveryPayloadMarshalling(t *testing.T) { MemberID: "member-d1", TeamID: "team-d1", ChatID: "chat-d1", - Content: map[string]interface{}{"summary": "done"}, - Preferences: map[string]interface{}{"channel": "email"}, + Content: &robottypes.DeliveryContent{ + Summary: "done", + Body: "full report", + }, + Preferences: &robottypes.DeliveryPreferences{ + Email: &robottypes.EmailPreference{ + Enabled: true, + Targets: []robottypes.EmailTarget{{To: []string{"a@b.com"}}}, + }, + }, } data, err := json.Marshal(payload) @@ -134,13 +143,7 @@ func TestDeliveryPayloadMarshalling(t *testing.T) { assert.Equal(t, "exec-d1", parsed.ExecutionID) assert.Equal(t, "member-d1", parsed.MemberID) assert.NotNil(t, parsed.Content) + assert.Equal(t, "done", parsed.Content.Summary) assert.NotNil(t, parsed.Preferences) - - contentMap, ok := parsed.Content.(map[string]interface{}) - require.True(t, ok) - assert.Equal(t, "done", contentMap["summary"]) - - prefMap, ok := parsed.Preferences.(map[string]interface{}) - require.True(t, ok) - assert.Equal(t, "email", prefMap["channel"]) + assert.NotNil(t, parsed.Preferences.Email) } diff --git a/agent/robot/events/handlers.go b/agent/robot/events/handlers.go index 8112e7c0..fb219965 100644 --- a/agent/robot/events/handlers.go +++ b/agent/robot/events/handlers.go @@ -1,54 +1,426 @@ package events import ( + "bytes" "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" "encoding/json" "fmt" + "io" + "net/http" + "strings" + "time" + "github.com/yaoapp/gou/process" + "github.com/yaoapp/gou/text" "github.com/yaoapp/kun/log" + robottypes "github.com/yaoapp/yao/agent/robot/types" + "github.com/yaoapp/yao/attachment" + "github.com/yaoapp/yao/event" eventtypes "github.com/yaoapp/yao/event/types" + "github.com/yaoapp/yao/messenger" + messengerTypes "github.com/yaoapp/yao/messenger/types" ) -// DeliveryHandler processes robot.delivery events asynchronously. -// It routes delivery content to configured channels (email, webhook, process). -type DeliveryHandler struct{} +func init() { + event.Register("robot", &robotHandler{ + httpClient: &http.Client{Timeout: 30 * time.Second}, + }) +} -// Handle processes a delivery event from the event bus. -func (h *DeliveryHandler) Handle(ctx context.Context, ev *eventtypes.Event, resp chan<- eventtypes.Result) { - var payload DeliveryPayload - if err := ev.Should(&payload); err != nil { - log.Error("delivery handler: invalid payload: %v", err) - return - } +// robotHandler processes all robot.* events. +type robotHandler struct { + httpClient *http.Client +} - log.Info( - "delivery handler: processing delivery for execution=%s member=%s", - payload.ExecutionID, payload.MemberID, - ) - - // Log delivery content summary for observability - if payload.Content != nil { - if data, err := json.Marshal(payload.Content); err == nil { - log.Debug("delivery handler: content=%s", string(data)) - } - } - - // Actual delivery routing is deferred to registered channel handlers. - // In the current implementation, the DeliveryCenter logic in delivery.go - // can be invoked here if needed. For now, this handler serves as the - // event-driven entry point for future channel-specific handlers. - - if ev.IsCall { - resp <- eventtypes.Result{Data: fmt.Sprintf("delivery processed for %s", payload.ExecutionID)} +// Handle dispatches robot events by type. +func (h *robotHandler) Handle(ctx context.Context, ev *eventtypes.Event, resp chan<- eventtypes.Result) { + switch ev.Type { + case Delivery: + h.handleDelivery(ctx, ev, resp) + default: + log.Debug("robot handler: unhandled event type=%s id=%s", ev.Type, ev.ID) } } -// Shutdown gracefully shuts down the delivery handler. -func (h *DeliveryHandler) Shutdown(ctx context.Context) error { +// Shutdown gracefully shuts down the robot handler. +func (h *robotHandler) Shutdown(ctx context.Context) error { return nil } -// NewDeliveryHandler creates a new DeliveryHandler. -func NewDeliveryHandler() *DeliveryHandler { - return &DeliveryHandler{} +// handleDelivery routes delivery content to configured channels (email, webhook, process). +func (h *robotHandler) handleDelivery(ctx context.Context, ev *eventtypes.Event, resp chan<- eventtypes.Result) { + var payload DeliveryPayload + if err := ev.Should(&payload); err != nil { + log.Error("delivery handler: invalid payload: %v", err) + if ev.IsCall { + resp <- eventtypes.Result{Err: err} + } + return + } + + log.Info("delivery handler: execution=%s member=%s", payload.ExecutionID, payload.MemberID) + + content := payload.Content + prefs := payload.Preferences + if content == nil { + log.Warn("delivery handler: nil content for execution=%s", payload.ExecutionID) + if ev.IsCall { + resp <- eventtypes.Result{Data: "no content"} + } + return + } + if prefs == nil { + if ev.IsCall { + resp <- eventtypes.Result{Data: "no preferences, skipped"} + } + return + } + + deliveryCtx := &robottypes.DeliveryContext{ + MemberID: payload.MemberID, + ExecutionID: payload.ExecutionID, + TeamID: payload.TeamID, + } + + var results []robottypes.ChannelResult + var lastErr error + + if prefs.Email != nil && prefs.Email.Enabled { + for _, target := range prefs.Email.Targets { + r := h.sendEmail(ctx, content, target, deliveryCtx) + results = append(results, r) + if !r.Success && lastErr == nil { + lastErr = fmt.Errorf("email delivery failed: %s", r.Error) + } + } + } + + if prefs.Webhook != nil && prefs.Webhook.Enabled { + for _, target := range prefs.Webhook.Targets { + r := h.postWebhook(ctx, content, target, deliveryCtx) + results = append(results, r) + if !r.Success && lastErr == nil { + lastErr = fmt.Errorf("webhook delivery failed: %s", r.Error) + } + } + } + + if prefs.Process != nil && prefs.Process.Enabled { + for _, target := range prefs.Process.Targets { + r := h.callProcess(ctx, content, target, deliveryCtx) + results = append(results, r) + if !r.Success && lastErr == nil { + lastErr = fmt.Errorf("process delivery failed: %s", r.Error) + } + } + } + + if lastErr != nil { + log.Error("delivery handler: partial failure execution=%s: %v", payload.ExecutionID, lastErr) + } + + if ev.IsCall { + resp <- eventtypes.Result{ + Data: map[string]interface{}{ + "execution_id": payload.ExecutionID, + "results": results, + }, + Err: lastErr, + } + } +} + +// ============================================================================ +// Email +// ============================================================================ + +func (h *robotHandler) sendEmail( + ctx context.Context, + content *robottypes.DeliveryContent, + target robottypes.EmailTarget, + deliveryCtx *robottypes.DeliveryContext, +) robottypes.ChannelResult { + now := time.Now() + targetID := strings.Join(target.To, ",") + if targetID == "" { + targetID = "no-recipients" + } + + result := robottypes.ChannelResult{ + Type: robottypes.DeliveryEmail, + Target: targetID, + SentAt: &now, + } + + svc := messenger.Instance + if svc == nil { + result.Error = "messenger service not available" + return result + } + + htmlBody, plainBody := buildEmailBody(target.Template, content) + msg := &messengerTypes.Message{ + To: target.To, + Subject: buildEmailSubject(target.Subject, target.Template, content, deliveryCtx), + Body: plainBody, + HTML: htmlBody, + Type: messengerTypes.MessageTypeEmail, + } + + attachments := convertAttachments(ctx, content.Attachments) + if len(attachments) > 0 { + msg.Attachments = attachments + } + + channel := robottypes.DefaultEmailChannel() + if err := svc.Send(ctx, channel, msg); err != nil { + result.Error = err.Error() + return result + } + + result.Success = true + result.Recipients = target.To + return result +} + +// ============================================================================ +// Webhook +// ============================================================================ + +func (h *robotHandler) postWebhook( + ctx context.Context, + content *robottypes.DeliveryContent, + target robottypes.WebhookTarget, + deliveryCtx *robottypes.DeliveryContext, +) robottypes.ChannelResult { + now := time.Now() + result := robottypes.ChannelResult{ + Type: robottypes.DeliveryWebhook, + Target: target.URL, + SentAt: &now, + } + + payload := map[string]interface{}{ + "event": "robot.delivery", + "timestamp": now.Format(time.RFC3339), + "execution_id": deliveryCtx.ExecutionID, + "member_id": deliveryCtx.MemberID, + "team_id": deliveryCtx.TeamID, + "trigger_type": deliveryCtx.TriggerType, + "content": map[string]interface{}{ + "summary": content.Summary, + "body": content.Body, + }, + } + + if len(content.Attachments) > 0 { + info := make([]map[string]interface{}, 0, len(content.Attachments)) + for _, att := range content.Attachments { + info = append(info, map[string]interface{}{ + "title": att.Title, + "description": att.Description, + "task_id": att.TaskID, + "file": att.File, + }) + } + payload["attachments"] = info + } + + payloadBytes, err := json.Marshal(payload) + if err != nil { + result.Error = fmt.Sprintf("failed to marshal payload: %v", err) + return result + } + + method := target.Method + if method == "" { + method = "POST" + } + + req, err := http.NewRequestWithContext(ctx, method, target.URL, bytes.NewReader(payloadBytes)) + if err != nil { + result.Error = fmt.Sprintf("failed to create request: %v", err) + return result + } + + req.Header.Set("Content-Type", "application/json") + for key, value := range target.Headers { + req.Header.Set(key, value) + } + + if target.Secret != "" { + signature := ComputeHMACSignature(payloadBytes, target.Secret) + req.Header.Set("X-Yao-Signature", signature) + req.Header.Set("X-Yao-Signature-Algorithm", "HMAC-SHA256") + } + + resp, err := h.httpClient.Do(req) + if err != nil { + result.Error = fmt.Sprintf("request failed: %v", err) + return result + } + defer resp.Body.Close() + + body, _ := io.ReadAll(resp.Body) + + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + result.Error = fmt.Sprintf("webhook returned status %d: %s", resp.StatusCode, string(body)) + return result + } + + result.Success = true + result.Details = map[string]interface{}{ + "status_code": resp.StatusCode, + "response": string(body), + } + return result +} + +// ============================================================================ +// Process +// ============================================================================ + +func (h *robotHandler) callProcess( + ctx context.Context, + content *robottypes.DeliveryContent, + target robottypes.ProcessTarget, + deliveryCtx *robottypes.DeliveryContext, +) robottypes.ChannelResult { + now := time.Now() + result := robottypes.ChannelResult{ + Type: robottypes.DeliveryProcess, + Target: target.Process, + SentAt: &now, + } + + args := make([]interface{}, 0, 1+len(target.Args)) + args = append(args, map[string]interface{}{ + "content": map[string]interface{}{ + "summary": content.Summary, + "body": content.Body, + "attachments": content.Attachments, + }, + "context": map[string]interface{}{ + "execution_id": deliveryCtx.ExecutionID, + "member_id": deliveryCtx.MemberID, + "team_id": deliveryCtx.TeamID, + "trigger_type": deliveryCtx.TriggerType, + }, + }) + args = append(args, target.Args...) + + proc, err := process.Of(target.Process, args...) + if err != nil { + result.Error = fmt.Sprintf("failed to create process: %v", err) + return result + } + proc.Context = ctx + + if err = proc.Execute(); err != nil { + result.Error = err.Error() + return result + } + + result.Success = true + result.Details = toJSONSerializable(proc.Value) + return result +} + +// ============================================================================ +// Helpers +// ============================================================================ + +func toJSONSerializable(v interface{}) interface{} { + if v == nil { + return nil + } + if _, err := json.Marshal(v); err != nil { + return fmt.Sprintf("%v", v) + } + return v +} + +func buildEmailSubject(subject, template string, content *robottypes.DeliveryContent, ctx *robottypes.DeliveryContext) string { + if subject != "" { + return subject + } + if content.Summary != "" { + return content.Summary + } + return fmt.Sprintf("Execution %s Complete", ctx.ExecutionID) +} + +func buildEmailBody(template string, content *robottypes.DeliveryContent) (string, string) { + markdown := content.Body + if markdown == "" { + markdown = content.Summary + } + html, err := text.MarkdownToHTML(markdown) + if err != nil { + return markdown, markdown + } + return html, markdown +} + +func convertAttachments(ctx context.Context, attachments []robottypes.DeliveryAttachment) []messengerTypes.Attachment { + if len(attachments) == 0 { + return nil + } + + result := make([]messengerTypes.Attachment, 0, len(attachments)) + for _, att := range attachments { + uploader, fileID, isWrapper := attachment.Parse(att.File) + if !isWrapper { + log.Warn("convertAttachments: skipping non-wrapper file value=%q title=%q", att.File, att.Title) + continue + } + manager, ok := attachment.Managers[uploader] + if !ok { + log.Warn("convertAttachments: manager not found uploader=%q file=%q title=%q (available: %v)", + uploader, att.File, att.Title, attachmentManagerKeys()) + continue + } + info, err := manager.Info(ctx, fileID) + if err != nil { + log.Warn("convertAttachments: failed to get file info fileID=%q uploader=%q: %v", fileID, uploader, err) + continue + } + content, err := manager.Read(ctx, fileID) + if err != nil { + log.Warn("convertAttachments: failed to read file fileID=%q uploader=%q: %v", fileID, uploader, err) + continue + } + log.Info("convertAttachments: added attachment filename=%q contentType=%q size=%d", info.Filename, info.ContentType, len(content)) + result = append(result, messengerTypes.Attachment{ + Filename: info.Filename, + ContentType: info.ContentType, + Content: content, + }) + } + return result +} + +// attachmentManagerKeys returns registered attachment manager names for debug logging. +func attachmentManagerKeys() []string { + keys := make([]string, 0, len(attachment.Managers)) + for k := range attachment.Managers { + keys = append(keys, k) + } + return keys +} + +// ComputeHMACSignature computes HMAC-SHA256 signature for webhook payload. +func ComputeHMACSignature(payload []byte, secret string) string { + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(payload) + return hex.EncodeToString(mac.Sum(nil)) +} + +// VerifyHMACSignature verifies the HMAC-SHA256 signature of a webhook payload. +func VerifyHMACSignature(payload []byte, secret, signature string) bool { + expected := ComputeHMACSignature(payload, secret) + return hmac.Equal([]byte(expected), []byte(signature)) } diff --git a/agent/robot/events/handlers_test.go b/agent/robot/events/handlers_test.go index d173258e..3119ada1 100644 --- a/agent/robot/events/handlers_test.go +++ b/agent/robot/events/handlers_test.go @@ -2,62 +2,155 @@ package events import ( "context" + "encoding/json" + "net/http" + "net/http/httptest" "testing" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + robottypes "github.com/yaoapp/yao/agent/robot/types" eventtypes "github.com/yaoapp/yao/event/types" ) -func TestDeliveryHandler_Handle(t *testing.T) { - handler := NewDeliveryHandler() - - t.Run("processes valid delivery payload", func(t *testing.T) { - ev := &eventtypes.Event{ - Type: Delivery, - ID: "test-event-1", - Payload: DeliveryPayload{ - ExecutionID: "exec-1", - MemberID: "member-1", - TeamID: "team-1", - Content: map[string]interface{}{"summary": "test"}, - }, - } - resp := make(chan eventtypes.Result, 1) - handler.Handle(context.Background(), ev, resp) - // Fire-and-forget: no response expected for Push - }) - - t.Run("handles call mode with response", func(t *testing.T) { - ev := &eventtypes.Event{ - Type: Delivery, - ID: "test-event-2", - IsCall: true, - Payload: DeliveryPayload{ - ExecutionID: "exec-2", - MemberID: "member-2", - }, - } - resp := make(chan eventtypes.Result, 1) - handler.Handle(context.Background(), ev, resp) - - result := <-resp - require.NotNil(t, result.Data) - assert.Contains(t, result.Data.(string), "exec-2") - }) - - t.Run("handles invalid payload gracefully", func(t *testing.T) { - ev := &eventtypes.Event{ - Type: Delivery, - Payload: "invalid", - } - resp := make(chan eventtypes.Result, 1) - handler.Handle(context.Background(), ev, resp) - }) +func newTestHandler() *robotHandler { + return &robotHandler{ + httpClient: http.DefaultClient, + } } -func TestDeliveryHandler_Shutdown(t *testing.T) { - handler := NewDeliveryHandler() +func TestRobotHandler_DeliveryWebhook(t *testing.T) { + var received map[string]interface{} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + decoder := json.NewDecoder(r.Body) + _ = decoder.Decode(&received) + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"ok":true}`)) + })) + defer server.Close() + + handler := newTestHandler() + ev := &eventtypes.Event{ + Type: Delivery, + ID: "test-ev-1", + IsCall: true, + Payload: DeliveryPayload{ + ExecutionID: "exec-1", + MemberID: "member-1", + TeamID: "team-1", + Content: &robottypes.DeliveryContent{ + Summary: "test summary", + Body: "test body", + }, + Preferences: &robottypes.DeliveryPreferences{ + Webhook: &robottypes.WebhookPreference{ + Enabled: true, + Targets: []robottypes.WebhookTarget{ + {URL: server.URL}, + }, + }, + }, + }, + } + + resp := make(chan eventtypes.Result, 1) + handler.Handle(context.Background(), ev, resp) + + result := <-resp + require.NotNil(t, result.Data) + assert.NoError(t, result.Err) + + data, ok := result.Data.(map[string]interface{}) + require.True(t, ok) + assert.Equal(t, "exec-1", data["execution_id"]) + + require.NotNil(t, received) + assert.Equal(t, "robot.delivery", received["event"]) +} + +func TestRobotHandler_DeliveryNoContent(t *testing.T) { + handler := newTestHandler() + ev := &eventtypes.Event{ + Type: Delivery, + ID: "test-ev-2", + IsCall: true, + Payload: DeliveryPayload{ + ExecutionID: "exec-2", + MemberID: "member-2", + TeamID: "team-2", + }, + } + + resp := make(chan eventtypes.Result, 1) + handler.Handle(context.Background(), ev, resp) + + result := <-resp + assert.Equal(t, "no content", result.Data) +} + +func TestRobotHandler_DeliveryNoPreferences(t *testing.T) { + handler := newTestHandler() + ev := &eventtypes.Event{ + Type: Delivery, + ID: "test-ev-3", + IsCall: true, + Payload: DeliveryPayload{ + ExecutionID: "exec-3", + MemberID: "member-3", + TeamID: "team-3", + Content: &robottypes.DeliveryContent{ + Summary: "test", + Body: "body", + }, + }, + } + + resp := make(chan eventtypes.Result, 1) + handler.Handle(context.Background(), ev, resp) + + result := <-resp + assert.Equal(t, "no preferences, skipped", result.Data) +} + +func TestRobotHandler_InvalidPayload(t *testing.T) { + handler := newTestHandler() + ev := &eventtypes.Event{ + Type: Delivery, + ID: "test-ev-4", + IsCall: true, + Payload: "invalid", + } + + resp := make(chan eventtypes.Result, 1) + handler.Handle(context.Background(), ev, resp) + + result := <-resp + assert.Error(t, result.Err) +} + +func TestRobotHandler_UnhandledEvent(t *testing.T) { + handler := newTestHandler() + ev := &eventtypes.Event{ + Type: "robot.unknown", + ID: "test-ev-5", + } + + resp := make(chan eventtypes.Result, 1) + handler.Handle(context.Background(), ev, resp) + // Fire-and-forget, no response expected +} + +func TestRobotHandler_Shutdown(t *testing.T) { + handler := newTestHandler() err := handler.Shutdown(context.Background()) assert.NoError(t, err) } + +func TestVerifyHMACSignature(t *testing.T) { + payload := []byte(`{"event":"robot.delivery"}`) + secret := "test-secret" + + sig := ComputeHMACSignature(payload, secret) + assert.True(t, VerifyHMACSignature(payload, secret, sig)) + assert.False(t, VerifyHMACSignature(payload, "wrong-secret", sig)) +} diff --git a/agent/robot/executor/standard/delivery.go b/agent/robot/executor/standard/delivery.go index 83e00394..24a9851c 100644 --- a/agent/robot/executor/standard/delivery.go +++ b/agent/robot/executor/standard/delivery.go @@ -7,43 +7,32 @@ import ( "time" "github.com/yaoapp/gou/model" + "github.com/yaoapp/kun/log" robotevents "github.com/yaoapp/yao/agent/robot/events" robottypes "github.com/yaoapp/yao/agent/robot/types" "github.com/yaoapp/yao/event" ) // RunDelivery executes P4: Delivery phase -// Calls the Delivery Agent to generate content, then routes to Delivery Center -// -// Input: -// - Full execution context (P0-P3) -// - Robot config -// -// Output: -// - DeliveryResult with content and channel results // // Process: // 1. Call Delivery Agent with full execution context // 2. Agent generates DeliveryContent (summary, body, attachments) -// 3. Route content to Delivery Center for actual delivery +// 3. Push delivery event for asynchronous routing via handlers func (e *Executor) RunDelivery(ctx *robottypes.Context, exec *robottypes.Execution, _ interface{}) error { - // Get robot for identity and resources robot := exec.GetRobot() if robot == nil { return fmt.Errorf("robot not found in execution") } - // Update UI field with i18n locale := getEffectiveLocale(robot, exec.Input) e.updateUIFields(ctx, exec, "", getLocalizedMessage(locale, "generating_delivery")) - // Get agent ID for delivery phase - agentID := "__yao.delivery" // default + agentID := "__yao.delivery" if robot.Config != nil && robot.Config.Resources != nil { agentID = robot.Config.Resources.GetPhaseAgent(robottypes.PhaseDelivery) } - // Build input for Delivery Agent formatter := NewInputFormatter() userContent := formatter.FormatDeliveryInput(exec, robot) @@ -51,7 +40,6 @@ func (e *Executor) RunDelivery(ctx *robottypes.Context, exec *robottypes.Executi return fmt.Errorf("no content available for delivery generation") } - // Call Delivery Agent caller := NewAgentCaller() caller.Connector = robot.LanguageModel result, err := caller.CallWithMessages(ctx, agentID, userContent) @@ -59,11 +47,8 @@ func (e *Executor) RunDelivery(ctx *robottypes.Context, exec *robottypes.Executi return fmt.Errorf("delivery agent (%s) call failed: %w", agentID, err) } - // Parse response as JSON - // Delivery Agent returns: { "content": { "summary": "...", "body": "...", "attachments": [...] } } data, err := result.GetJSON() if err != nil { - // Fallback: if not JSON, create minimal content from raw text content := result.GetText() if content == "" { return fmt.Errorf("delivery agent returned empty response") @@ -79,20 +64,17 @@ func (e *Executor) RunDelivery(ctx *robottypes.Context, exec *robottypes.Executi return e.pushDeliveryEvent(ctx, exec, robot) } - // Build DeliveryContent from JSON content := parseDeliveryContent(data) if content == nil { return fmt.Errorf("delivery agent (%s) returned invalid content", agentID) } - // Build DeliveryResult exec.Delivery = &robottypes.DeliveryResult{ RequestID: generateRequestID(exec.ID), Content: content, Success: true, } - // Push delivery event for asynchronous routing via handlers return e.pushDeliveryEvent(ctx, exec, robot) } @@ -100,7 +82,7 @@ func (e *Executor) RunDelivery(ctx *robottypes.Context, exec *robottypes.Executi // Registered handlers (see events/handlers.go) route to email/webhook/process channels. func (e *Executor) pushDeliveryEvent(ctx *robottypes.Context, exec *robottypes.Execution, robot *robottypes.Robot) error { prefs := buildDeliveryPreferences(robot) - event.Push(ctx.Context, robotevents.Delivery, robotevents.DeliveryPayload{ + _, err := event.Push(ctx.Context, robotevents.Delivery, robotevents.DeliveryPayload{ ExecutionID: exec.ID, MemberID: exec.MemberID, TeamID: exec.TeamID, @@ -108,61 +90,9 @@ func (e *Executor) pushDeliveryEvent(ctx *robottypes.Context, exec *robottypes.E Content: exec.Delivery.Content, Preferences: prefs, }) - return nil -} - -// routeToDeliveryCenter sends content to the Delivery Center for actual delivery -// The Delivery Center decides which channels to use based on robot/user preferences -// -// Delivery logic: -// 1. Manager email: ALWAYS send to manager if manager_id is set (mandatory) -// 2. Additional targets: Append configured email/webhook/process targets -func (e *Executor) routeToDeliveryCenter(ctx *robottypes.Context, exec *robottypes.Execution, robot *robottypes.Robot) error { - if exec.Delivery == nil || exec.Delivery.Content == nil { - return fmt.Errorf("no delivery content to route") - } - - // Build final delivery preferences by merging manager email + configured targets - prefs := buildDeliveryPreferences(robot) - if prefs == nil || !hasActiveChannels(prefs) { - exec.Delivery.Success = true - return nil - } - - // Update UI field to show delivery is in progress - locale := getEffectiveLocale(robot, exec.Input) - e.updateUIFields(ctx, exec, "", getLocalizedMessage(locale, "sending_delivery")) - - // Create Delivery Center and execute - center := NewDeliveryCenter() - results, err := center.Deliver(ctx, exec.Delivery.Content, &robottypes.DeliveryContext{ - MemberID: exec.MemberID, - ExecutionID: exec.ID, - TriggerType: exec.TriggerType, - TeamID: exec.TeamID, - }, prefs, robot) - - // Update delivery result - exec.Delivery.Results = results - now := time.Now() - exec.Delivery.SentAt = &now - if err != nil { - exec.Delivery.Success = false - exec.Delivery.Error = err.Error() - return err + log.Error("delivery event push failed: execution=%s error=%v", exec.ID, err) } - - // Check if all channels succeeded - allSuccess := true - for _, r := range results { - if !r.Success { - allSuccess = false - break - } - } - exec.Delivery.Success = allSuccess - return nil } @@ -172,26 +102,20 @@ func parseDeliveryContent(data map[string]interface{}) *robottypes.DeliveryConte return nil } - // Try to get content object contentData, ok := data["content"].(map[string]interface{}) if !ok { - // Fallback: maybe the data itself is the content contentData = data } content := &robottypes.DeliveryContent{} - // Parse summary if summary, ok := contentData["summary"].(string); ok { content.Summary = summary } - - // Parse body if body, ok := contentData["body"].(string); ok { content.Body = body } - // Parse attachments if attachments, ok := contentData["attachments"].([]interface{}); ok { for _, att := range attachments { if attMap, ok := att.(map[string]interface{}); ok { @@ -203,7 +127,6 @@ func parseDeliveryContent(data map[string]interface{}) *robottypes.DeliveryConte } } - // Validate: at least summary or body should be present if content.Summary == "" && content.Body == "" { return nil } @@ -211,7 +134,6 @@ func parseDeliveryContent(data map[string]interface{}) *robottypes.DeliveryConte return content } -// parseDeliveryAttachment parses a single attachment from the agent response func parseDeliveryAttachment(data map[string]interface{}) *robottypes.DeliveryAttachment { if data == nil { return nil @@ -232,7 +154,6 @@ func parseDeliveryAttachment(data map[string]interface{}) *robottypes.DeliveryAt att.File = file } - // At minimum, need title and file if att.Title == "" || att.File == "" { return nil } @@ -240,21 +161,17 @@ func parseDeliveryAttachment(data map[string]interface{}) *robottypes.DeliveryAt return att } -// generateRequestID generates a unique request ID for delivery tracking func generateRequestID(execID string) string { return fmt.Sprintf("dlv-%s-%d", execID, time.Now().UnixNano()%1000000) } -// getTaskDescription extracts a description from task messages func getTaskDescription(task robottypes.Task) string { if len(task.Messages) == 0 { return task.GoalRef } - // Try to get text from first message for _, msg := range task.Messages { if content, ok := msg.Content.(string); ok && content != "" { - // Truncate if too long if len(content) > 100 { return content[:97] + "..." } @@ -262,7 +179,6 @@ func getTaskDescription(task robottypes.Task) string { } } - // Fallback to goal reference if task.GoalRef != "" { return task.GoalRef } @@ -270,12 +186,10 @@ func getTaskDescription(task robottypes.Task) string { return "Task " + task.ID } -// truncateSummary truncates text to maxLen characters func truncateSummary(text string, maxLen int) string { if len(text) <= maxLen { return text } - // Find last space before maxLen to avoid cutting words truncated := text[:maxLen] if idx := strings.LastIndex(truncated, " "); idx > maxLen/2 { return truncated[:idx] + "..." @@ -283,26 +197,6 @@ func truncateSummary(text string, maxLen int) string { return truncated + "..." } -// hasActiveChannels checks if any delivery channel is configured -func hasActiveChannels(prefs *robottypes.DeliveryPreferences) bool { - if prefs == nil { - return false - } - if prefs.Email != nil && prefs.Email.Enabled && len(prefs.Email.Targets) > 0 { - return true - } - if prefs.Webhook != nil && prefs.Webhook.Enabled && len(prefs.Webhook.Targets) > 0 { - return true - } - if prefs.Process != nil && prefs.Process.Enabled && len(prefs.Process.Targets) > 0 { - return true - } - return false -} - -// buildDeliveryPreferences builds the final delivery preferences by: -// 1. Always including manager email if manager_id is set (mandatory) -// 2. Appending all configured email/webhook/process targets func buildDeliveryPreferences(robot *robottypes.Robot) *robottypes.DeliveryPreferences { if robot == nil { return nil @@ -310,26 +204,22 @@ func buildDeliveryPreferences(robot *robottypes.Robot) *robottypes.DeliveryPrefe prefs := &robottypes.DeliveryPreferences{} - // Step 1: Get manager email (mandatory if manager_id is set) managerEmail := robot.ManagerEmail if managerEmail == "" && robot.ManagerID != "" { managerEmail = getManagerEmail(robot.ManagerID) if managerEmail != "" { - robot.ManagerEmail = managerEmail // Cache for future use + robot.ManagerEmail = managerEmail } } - // Step 2: Build email targets (manager first, then configured targets) var emailTargets []robottypes.EmailTarget - // Add manager email as first target (mandatory) if managerEmail != "" { emailTargets = append(emailTargets, robottypes.EmailTarget{ To: []string{managerEmail}, }) } - // Append configured email targets if robot.Config != nil && robot.Config.Delivery != nil && robot.Config.Delivery.Email != nil { for _, target := range robot.Config.Delivery.Email.Targets { if len(target.To) > 0 { @@ -338,7 +228,6 @@ func buildDeliveryPreferences(robot *robottypes.Robot) *robottypes.DeliveryPrefe } } - // Set email preference if we have any targets if len(emailTargets) > 0 { prefs.Email = &robottypes.EmailPreference{ Enabled: true, @@ -346,14 +235,12 @@ func buildDeliveryPreferences(robot *robottypes.Robot) *robottypes.DeliveryPrefe } } - // Step 3: Copy webhook preferences from config (if enabled) if robot.Config != nil && robot.Config.Delivery != nil && robot.Config.Delivery.Webhook != nil { if robot.Config.Delivery.Webhook.Enabled && len(robot.Config.Delivery.Webhook.Targets) > 0 { prefs.Webhook = robot.Config.Delivery.Webhook } } - // Step 4: Copy process preferences from config (if enabled) if robot.Config != nil && robot.Config.Delivery != nil && robot.Config.Delivery.Process != nil { if robot.Config.Delivery.Process.Enabled && len(robot.Config.Delivery.Process.Targets) > 0 { prefs.Process = robot.Config.Delivery.Process @@ -363,8 +250,6 @@ func buildDeliveryPreferences(robot *robottypes.Robot) *robottypes.DeliveryPrefe return prefs } -// getManagerEmail retrieves the manager's email from __yao.member table by member_id -// manager_id in Robot refers to a member_id in __yao.member table func getManagerEmail(managerID string) string { if managerID == "" { return "" @@ -400,7 +285,6 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * var sb strings.Builder - // Robot identity if robot != nil && robot.Config != nil && robot.Config.Identity != nil { sb.WriteString("## Robot Identity\n\n") sb.WriteString(fmt.Sprintf("- **Role**: %s\n", robot.Config.Identity.Role)) @@ -412,7 +296,6 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * sb.WriteString("\n") } - // Trigger type sb.WriteString("## Execution Context\n\n") sb.WriteString(fmt.Sprintf("- **Trigger**: %s\n", exec.TriggerType)) sb.WriteString(fmt.Sprintf("- **Status**: %s\n", exec.Status)) @@ -423,25 +306,21 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * } sb.WriteString("\n") - // Inspiration (P0) - for clock trigger if exec.Inspiration != nil && exec.Inspiration.Content != "" { sb.WriteString("## Inspiration (P0)\n\n") sb.WriteString(exec.Inspiration.Content) sb.WriteString("\n\n") } - // Goals (P1) if exec.Goals != nil && exec.Goals.Content != "" { sb.WriteString("## Goals (P1)\n\n") sb.WriteString(exec.Goals.Content) sb.WriteString("\n\n") } - // Tasks (P2) if len(exec.Tasks) > 0 { sb.WriteString("## Tasks (P2)\n\n") for i, task := range exec.Tasks { - // Extract task description from messages if available taskDesc := getTaskDescription(task) sb.WriteString(fmt.Sprintf("%d. **%s** - %s\n", i+1, task.ID, taskDesc)) sb.WriteString(fmt.Sprintf(" - Executor: %s (%s)\n", task.ExecutorID, task.ExecutorType)) @@ -453,7 +332,6 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * sb.WriteString("\n") } - // Results (P3) - detailed if len(exec.Results) > 0 { sb.WriteString("## Results (P3)\n\n") @@ -471,7 +349,6 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * sb.WriteString(fmt.Sprintf("- **Duration**: %dms\n", result.Duration)) - // Validation if result.Validation != nil { if result.Validation.Passed { sb.WriteString(fmt.Sprintf("- **Validation**: ✓ Passed (score: %.2f)\n", result.Validation.Score)) @@ -485,7 +362,6 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * } } - // Output if result.Output != nil { sb.WriteString("\n**Output**:\n") if output, err := json.MarshalIndent(result.Output, "", " "); err == nil { @@ -497,7 +373,6 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * } } - // Error if result.Error != "" { sb.WriteString(fmt.Sprintf("\n**Error**: %s\n", result.Error)) } @@ -505,7 +380,6 @@ func (f *InputFormatter) FormatDeliveryInput(exec *robottypes.Execution, robot * sb.WriteString("\n") } - // Summary sb.WriteString(fmt.Sprintf("### Summary\n\n- **Total Tasks**: %d\n- **Succeeded**: %d\n- **Failed**: %d\n\n", len(exec.Results), successCount, failCount)) } diff --git a/agent/robot/executor/standard/delivery_center.go b/agent/robot/executor/standard/delivery_center.go deleted file mode 100644 index 0df85eb0..00000000 --- a/agent/robot/executor/standard/delivery_center.go +++ /dev/null @@ -1,444 +0,0 @@ -package standard - -import ( - "bytes" - "context" - "crypto/hmac" - "crypto/sha256" - "encoding/hex" - "encoding/json" - "fmt" - "io" - "net/http" - "strings" - "time" - - "github.com/yaoapp/gou/process" - "github.com/yaoapp/gou/text" - robottypes "github.com/yaoapp/yao/agent/robot/types" - "github.com/yaoapp/yao/attachment" - "github.com/yaoapp/yao/messenger" - messengerTypes "github.com/yaoapp/yao/messenger/types" -) - -// DeliveryCenter handles routing delivery content to configured channels -// It decides which channels to use based on robot/user preferences and executes the delivery -type DeliveryCenter struct { - httpClient *http.Client -} - -// NewDeliveryCenter creates a new DeliveryCenter instance -func NewDeliveryCenter() *DeliveryCenter { - return &DeliveryCenter{ - httpClient: &http.Client{ - Timeout: 30 * time.Second, - }, - } -} - -// Deliver sends content to all configured channels based on preferences -// Returns results for each channel target and any error -func (dc *DeliveryCenter) Deliver( - ctx *robottypes.Context, - content *robottypes.DeliveryContent, - deliveryCtx *robottypes.DeliveryContext, - prefs *robottypes.DeliveryPreferences, - robotInstance *robottypes.Robot, -) ([]robottypes.ChannelResult, error) { - if content == nil { - return nil, fmt.Errorf("delivery content is nil") - } - if prefs == nil { - return nil, nil // No preferences = no delivery - } - - var results []robottypes.ChannelResult - var lastErr error - - // Process email targets - if prefs.Email != nil && prefs.Email.Enabled { - for _, target := range prefs.Email.Targets { - result := dc.sendEmail(ctx.Context, content, target, deliveryCtx, robotInstance) - results = append(results, result) - if !result.Success && lastErr == nil { - lastErr = fmt.Errorf("email delivery failed: %s", result.Error) - } - } - } - - // Process webhook targets - if prefs.Webhook != nil && prefs.Webhook.Enabled { - for _, target := range prefs.Webhook.Targets { - result := dc.postWebhook(ctx.Context, content, target, deliveryCtx) - results = append(results, result) - if !result.Success && lastErr == nil { - lastErr = fmt.Errorf("webhook delivery failed: %s", result.Error) - } - } - } - - // Process process targets - if prefs.Process != nil && prefs.Process.Enabled { - for _, target := range prefs.Process.Targets { - result := dc.callProcess(ctx.Context, content, target, deliveryCtx) - results = append(results, result) - if !result.Success && lastErr == nil { - lastErr = fmt.Errorf("process delivery failed: %s", result.Error) - } - } - } - - return results, lastErr -} - -// sendEmail sends delivery content to a single email target -func (dc *DeliveryCenter) sendEmail( - ctx context.Context, - content *robottypes.DeliveryContent, - target robottypes.EmailTarget, - deliveryCtx *robottypes.DeliveryContext, - robotInstance *robottypes.Robot, -) robottypes.ChannelResult { - now := time.Now() - - // Build target identifier from recipients - targetID := strings.Join(target.To, ",") - if targetID == "" { - targetID = "no-recipients" - } - - result := robottypes.ChannelResult{ - Type: robottypes.DeliveryEmail, - Target: targetID, - SentAt: &now, - } - - // Get messenger service - svc := messenger.Instance - if svc == nil { - result.Error = "messenger service not available" - return result - } - - // Build email message with HTML content - htmlBody, plainBody := buildEmailBody(target.Template, content) - msg := &messengerTypes.Message{ - To: target.To, - Subject: buildEmailSubject(target.Subject, target.Template, content, deliveryCtx, robotInstance), - Body: plainBody, // Plain text fallback - HTML: htmlBody, // HTML content for rich email display - Type: messengerTypes.MessageTypeEmail, - } - - // Set From address from Robot's email (if configured) - if robotInstance != nil && robotInstance.RobotEmail != "" { - msg.From = robotInstance.RobotEmail - } - - // Convert attachments - attachments := convertAttachments(ctx, content.Attachments) - if len(attachments) > 0 { - msg.Attachments = attachments - } - - // Send email using global default channel - channel := robottypes.DefaultEmailChannel() - err := svc.Send(ctx, channel, msg) - if err != nil { - result.Error = err.Error() - return result - } - - result.Success = true - result.Recipients = target.To - - return result -} - -// postWebhook posts delivery content to a single webhook target -func (dc *DeliveryCenter) postWebhook( - ctx context.Context, - content *robottypes.DeliveryContent, - target robottypes.WebhookTarget, - deliveryCtx *robottypes.DeliveryContext, -) robottypes.ChannelResult { - now := time.Now() - result := robottypes.ChannelResult{ - Type: robottypes.DeliveryWebhook, - Target: target.URL, - SentAt: &now, - } - - // Build webhook payload - payload := map[string]interface{}{ - "event": "robot.delivery", - "timestamp": now.Format(time.RFC3339), - "execution_id": deliveryCtx.ExecutionID, - "member_id": deliveryCtx.MemberID, - "team_id": deliveryCtx.TeamID, - "trigger_type": deliveryCtx.TriggerType, - "content": map[string]interface{}{ - "summary": content.Summary, - "body": content.Body, - }, - } - - // Add attachments info (not the actual files) - if len(content.Attachments) > 0 { - attachmentInfo := make([]map[string]interface{}, 0, len(content.Attachments)) - for _, att := range content.Attachments { - attachmentInfo = append(attachmentInfo, map[string]interface{}{ - "title": att.Title, - "description": att.Description, - "task_id": att.TaskID, - "file": att.File, - }) - } - payload["attachments"] = attachmentInfo - } - - // Marshal payload - payloadBytes, err := json.Marshal(payload) - if err != nil { - result.Error = fmt.Sprintf("failed to marshal payload: %v", err) - return result - } - - // Build request - method := target.Method - if method == "" { - method = "POST" - } - - req, err := http.NewRequestWithContext(ctx, method, target.URL, bytes.NewReader(payloadBytes)) - if err != nil { - result.Error = fmt.Sprintf("failed to create request: %v", err) - return result - } - - req.Header.Set("Content-Type", "application/json") - - // Add custom headers - for key, value := range target.Headers { - req.Header.Set(key, value) - } - - // Add HMAC signature if secret is configured - if target.Secret != "" { - signature := computeHMACSignature(payloadBytes, target.Secret) - req.Header.Set("X-Yao-Signature", signature) - req.Header.Set("X-Yao-Signature-Algorithm", "HMAC-SHA256") - } - - // Send request - resp, err := dc.httpClient.Do(req) - if err != nil { - result.Error = fmt.Sprintf("request failed: %v", err) - return result - } - defer resp.Body.Close() - - // Read response body - body, _ := io.ReadAll(resp.Body) - - // Check status code - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - result.Error = fmt.Sprintf("webhook returned status %d: %s", resp.StatusCode, string(body)) - return result - } - - result.Success = true - result.Details = map[string]interface{}{ - "status_code": resp.StatusCode, - "response": string(body), - } - - return result -} - -// callProcess calls a Yao Process with delivery content -func (dc *DeliveryCenter) callProcess( - ctx context.Context, - content *robottypes.DeliveryContent, - target robottypes.ProcessTarget, - deliveryCtx *robottypes.DeliveryContext, -) robottypes.ChannelResult { - now := time.Now() - result := robottypes.ChannelResult{ - Type: robottypes.DeliveryProcess, - Target: target.Process, - SentAt: &now, - } - - // Build args: DeliveryContent as first arg, then additional args - args := make([]interface{}, 0, 1+len(target.Args)) - args = append(args, map[string]interface{}{ - "content": map[string]interface{}{ - "summary": content.Summary, - "body": content.Body, - "attachments": content.Attachments, - }, - "context": map[string]interface{}{ - "execution_id": deliveryCtx.ExecutionID, - "member_id": deliveryCtx.MemberID, - "team_id": deliveryCtx.TeamID, - "trigger_type": deliveryCtx.TriggerType, - }, - }) - args = append(args, target.Args...) - - // Create and execute process - proc, err := process.Of(target.Process, args...) - if err != nil { - result.Error = fmt.Sprintf("failed to create process: %v", err) - return result - } - proc.Context = ctx - - err = proc.Execute() - if err != nil { - result.Error = err.Error() - return result - } - - result.Success = true - // Convert proc.Value to JSON-serializable format to avoid func type issues - result.Details = toJSONSerializable(proc.Value) - - return result -} - -// toJSONSerializable ensures the value can be JSON serialized -// Returns the original value if serializable, or a string fallback if not -func toJSONSerializable(v interface{}) interface{} { - if v == nil { - return nil - } - - // Try to marshal to check if it's JSON serializable - _, err := json.Marshal(v) - if err != nil { - // If it can't be serialized (e.g., contains func), return a string representation - return fmt.Sprintf("%v", v) - } - - // Return original value if it's serializable - return v -} - -// buildEmailSubject builds the email subject line -func buildEmailSubject(subject, template string, content *robottypes.DeliveryContent, ctx *robottypes.DeliveryContext, robot *robottypes.Robot) string { - // Use explicit subject if provided - if subject != "" { - return subject - } - - // Get robot display name for subject prefix - robotName := "Robot" - if robot != nil && robot.DisplayName != "" { - robotName = robot.DisplayName - } - - // Use template-based subject if template is specified - // TODO: Implement template rendering - if template != "" { - return fmt.Sprintf("[%s] %s", robotName, content.Summary) - } - - // Default: use summary - if content.Summary != "" { - return fmt.Sprintf("[%s] %s", robotName, content.Summary) - } - - return fmt.Sprintf("[%s] Execution %s Complete", robotName, ctx.ExecutionID) -} - -// buildEmailBody builds the email body content -// buildEmailBody returns HTML and plain text versions of the email body -// Returns: (htmlBody, plainBody) -func buildEmailBody(template string, content *robottypes.DeliveryContent) (string, string) { - // TODO: Implement template rendering - // Get markdown content (used as plain text fallback) - markdown := content.Body - if markdown == "" { - markdown = content.Summary - } - - // Convert Markdown to HTML for rich email display - html, err := text.MarkdownToHTML(markdown) - if err != nil { - // Fallback: use markdown as both HTML and plain text - return markdown, markdown - } - - return html, markdown -} - -// convertAttachments converts DeliveryAttachment to messenger Attachment format -func convertAttachments(ctx context.Context, attachments []robottypes.DeliveryAttachment) []messengerTypes.Attachment { - if len(attachments) == 0 { - return nil - } - - result := make([]messengerTypes.Attachment, 0, len(attachments)) - - for _, att := range attachments { - // Parse file wrapper: __:// - uploader, fileID, isWrapper := attachment.Parse(att.File) - if !isWrapper { - // Skip non-wrapper attachments - continue - } - - // Get file info from attachment manager - manager, ok := attachment.Managers[uploader] - if !ok { - continue - } - - info, err := manager.Info(ctx, fileID) - if err != nil { - continue - } - - // Read file content - content, err := manager.Read(ctx, fileID) - if err != nil { - continue - } - - // Build messenger attachment - msgAtt := messengerTypes.Attachment{ - Filename: info.Filename, - ContentType: info.ContentType, - Content: content, - } - - result = append(result, msgAtt) - } - - return result -} - -// ============================================================================ -// Webhook Signature -// ============================================================================ - -// computeHMACSignature computes HMAC-SHA256 signature for webhook payload -// Returns hex-encoded signature string -func computeHMACSignature(payload []byte, secret string) string { - mac := hmac.New(sha256.New, []byte(secret)) - mac.Write(payload) - return hex.EncodeToString(mac.Sum(nil)) -} - -// VerifyHMACSignature verifies the HMAC-SHA256 signature of a webhook payload -// Headers: -// - X-Yao-Signature: hex-encoded HMAC-SHA256 signature -// - X-Yao-Signature-Algorithm: "HMAC-SHA256" -// -// Returns true if the signature is valid -func VerifyHMACSignature(payload []byte, secret, signature string) bool { - expected := computeHMACSignature(payload, secret) - return hmac.Equal([]byte(expected), []byte(signature)) -} diff --git a/agent/robot/executor/standard/delivery_test.go b/agent/robot/executor/standard/delivery_test.go index b93cb7ac..df704b44 100644 --- a/agent/robot/executor/standard/delivery_test.go +++ b/agent/robot/executor/standard/delivery_test.go @@ -2,9 +2,6 @@ package standard_test import ( "context" - "encoding/json" - "net/http" - "net/http/httptest" "strings" "testing" "time" @@ -34,7 +31,6 @@ func TestRunDeliveryBasic(t *testing.T) { robot := createDeliveryTestRobot(t, "robot.delivery") exec := createDeliveryTestExecution(robot) - // Set up execution context with P0-P3 results exec.Inspiration = &types.InspirationReport{ Content: "Morning analysis suggests focus on Q4 review.", } @@ -50,7 +46,6 @@ func TestRunDeliveryBasic(t *testing.T) { {TaskID: "task-002", Success: true, Duration: 800, Output: "Q4 sales exceeded expectations by 15%."}, } - // Run delivery phase e := standard.New() err := e.RunDelivery(ctx, exec, nil) @@ -85,7 +80,6 @@ func TestRunDeliveryBasic(t *testing.T) { require.NotNil(t, exec.Delivery) require.NotNil(t, exec.Delivery.Content) - // Content should mention the failure body := strings.ToLower(exec.Delivery.Content.Body) hasFailureInfo := strings.Contains(body, "fail") || strings.Contains(body, "error") || @@ -111,7 +105,6 @@ func TestRunDeliveryErrorHandling(t *testing.T) { ID: "test-exec-1", TriggerType: types.TriggerClock, } - // Don't set robot e := standard.New() err := e.RunDelivery(ctx, exec, nil) @@ -147,314 +140,23 @@ func TestRunDeliveryErrorHandling(t *testing.T) { } // ============================================================================ -// Delivery Center Tests +// Email Channel Config Tests // ============================================================================ -func TestDeliveryCenterWebhook(t *testing.T) { - t.Run("posts to webhook successfully", func(t *testing.T) { - // Create mock webhook server - var receivedPayload map[string]interface{} - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - assert.Equal(t, "POST", r.Method) - assert.Equal(t, "application/json", r.Header.Get("Content-Type")) - - decoder := json.NewDecoder(r.Body) - err := decoder.Decode(&receivedPayload) - assert.NoError(t, err) - - w.WriteHeader(http.StatusOK) - w.Write([]byte(`{"status": "received"}`)) - })) - defer server.Close() - - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test delivery completed", - Body: "# Test Report\n\nThis is a test.", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - {URL: server.URL}, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - require.NoError(t, err) - require.Len(t, results, 1) - assert.True(t, results[0].Success) - assert.Equal(t, types.DeliveryWebhook, results[0].Type) - assert.Equal(t, server.URL, results[0].Target) - - // Verify payload structure - assert.Equal(t, "robot.delivery", receivedPayload["event"]) - assert.Equal(t, "exec-001", receivedPayload["execution_id"]) - assert.Equal(t, "member-001", receivedPayload["member_id"]) - contentMap := receivedPayload["content"].(map[string]interface{}) - assert.Equal(t, "Test delivery completed", contentMap["summary"]) - }) - - t.Run("handles webhook failure", func(t *testing.T) { - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusInternalServerError) - w.Write([]byte(`{"error": "internal error"}`)) - })) - defer server.Close() - - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test", - Body: "Test body", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - {URL: server.URL}, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - assert.Error(t, err) // Should return error for failed delivery - require.Len(t, results, 1) - assert.False(t, results[0].Success) - assert.Contains(t, results[0].Error, "500") - }) - - t.Run("supports multiple webhook targets", func(t *testing.T) { - callCount := 0 - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - callCount++ - w.WriteHeader(http.StatusOK) - w.Write([]byte(`{"ok": true}`)) - })) - defer server.Close() - - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test", - Body: "Test body", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - {URL: server.URL + "/hook1"}, - {URL: server.URL + "/hook2"}, - {URL: server.URL + "/hook3"}, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - require.NoError(t, err) - require.Len(t, results, 3) - assert.Equal(t, 3, callCount) - - for _, r := range results { - assert.True(t, r.Success) - } - }) - - t.Run("includes custom headers", func(t *testing.T) { - var receivedHeaders http.Header - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - receivedHeaders = r.Header - w.WriteHeader(http.StatusOK) - })) - defer server.Close() - - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test", - Body: "Test body", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - { - URL: server.URL, - Headers: map[string]string{ - "X-Custom-Header": "custom-value", - "Authorization": "Bearer test-token", - }, - }, - }, - }, - } - - results, _ := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - require.Len(t, results, 1) - assert.True(t, results[0].Success) - assert.Equal(t, "custom-value", receivedHeaders.Get("X-Custom-Header")) - assert.Equal(t, "Bearer test-token", receivedHeaders.Get("Authorization")) - }) -} - -func TestDeliveryCenterNoChannels(t *testing.T) { - t.Run("succeeds with no channels configured", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test", - Body: "Test body", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - // No preferences - results, err := center.Deliver(ctx, content, deliveryCtx, nil, nil) - - assert.NoError(t, err) - assert.Empty(t, results) - }) - - t.Run("succeeds with disabled channels", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test", - Body: "Test body", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: false, // Disabled - Targets: []types.WebhookTarget{ - {URL: "http://example.com"}, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - assert.NoError(t, err) - assert.Empty(t, results) - }) -} - -func TestDeliveryCenterMixedChannels(t *testing.T) { - t.Run("delivers to multiple channel types", func(t *testing.T) { - webhookCalled := false - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - webhookCalled = true - w.WriteHeader(http.StatusOK) - })) - defer server.Close() - - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test", - Body: "Test body", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - {URL: server.URL}, - }, - }, - // Email would fail without messenger setup, but webhook should succeed - } - - results, _ := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - assert.True(t, webhookCalled) - require.Len(t, results, 1) - assert.True(t, results[0].Success) - }) -} - func TestDefaultEmailChannel(t *testing.T) { t.Run("returns default email channel", func(t *testing.T) { - // Default should be "default" assert.Equal(t, "default", types.DefaultEmailChannel()) }) t.Run("can set custom email channel", func(t *testing.T) { - // Save original original := types.DefaultEmailChannel() defer types.SetDefaultEmailChannel(original) - // Set custom channel types.SetDefaultEmailChannel("custom-email") assert.Equal(t, "custom-email", types.DefaultEmailChannel()) }) t.Run("ignores empty channel", func(t *testing.T) { - // Save original and restore after test original := types.DefaultEmailChannel() defer types.SetDefaultEmailChannel(original) @@ -464,53 +166,6 @@ func TestDefaultEmailChannel(t *testing.T) { } func TestRobotEmailInDelivery(t *testing.T) { - t.Run("robot email is passed to delivery center", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test", - Body: "Test body", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - // Robot with email configured - robot := &types.Robot{ - MemberID: "robot-001", - RobotEmail: "robot@example.com", - } - - // Webhook to verify robot is passed (email would fail without messenger) - webhookCalled := false - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - webhookCalled = true - w.WriteHeader(http.StatusOK) - })) - defer server.Close() - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - {URL: server.URL}, - }, - }, - } - - // Deliver with robot - results, _ := center.Deliver(ctx, content, deliveryCtx, prefs, robot) - - assert.True(t, webhookCalled) - require.Len(t, results, 1) - assert.True(t, results[0].Success) - }) - t.Run("robot email field is loaded from map", func(t *testing.T) { data := map[string]interface{}{ "member_id": "robot-001", @@ -527,7 +182,6 @@ func TestRobotEmailInDelivery(t *testing.T) { data := map[string]interface{}{ "member_id": "robot-001", "team_id": "team-001", - // robot_email not set } robot, err := types.NewRobotFromMap(data) @@ -536,6 +190,10 @@ func TestRobotEmailInDelivery(t *testing.T) { }) } +// ============================================================================ +// FormatDeliveryInput Tests +// ============================================================================ + func TestFormatDeliveryInput(t *testing.T) { formatter := standard.NewInputFormatter() @@ -603,538 +261,6 @@ func TestFormatDeliveryInput(t *testing.T) { }) } -// ============================================================================ -// Email Delivery Tests (requires messenger setup) -// ============================================================================ - -func TestDeliveryCenterEmail(t *testing.T) { - if testing.Short() { - t.Skip("Skipping integration test") - } - - testutils.Prepare(t) - defer testutils.Clean(t) - - // Set robot channel as default for tests - original := types.DefaultEmailChannel() - types.SetDefaultEmailChannel("robot") - defer types.SetDefaultEmailChannel(original) - - t.Run("sends email to single target", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Test Delivery Report", - Body: "This is a test delivery from Robot Agent.\n\n## Results\n- Task 1: Completed\n- Task 2: Completed", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-001", - TriggerType: types.TriggerClock, - TeamID: "team-001", - } - - robot := &types.Robot{ - MemberID: "robot-001", - RobotEmail: "robot@example.com", - } - - prefs := &types.DeliveryPreferences{ - Email: &types.EmailPreference{ - Enabled: true, - Targets: []types.EmailTarget{ - { - To: []string{"test@example.com"}, - Subject: "Robot Delivery Test", - }, - }, - }, - } - - // Send email - ignore billing/API errors, just verify the call was made - results, _ := center.Deliver(ctx, content, deliveryCtx, prefs, robot) - - require.Len(t, results, 1) - assert.Equal(t, types.DeliveryEmail, results[0].Type) - assert.Equal(t, "test@example.com", results[0].Target) - // Note: Success depends on messenger configuration - }) - - t.Run("sends email to multiple targets", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Multi-target Test", - Body: "Test body for multiple recipients", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-002", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - robot := &types.Robot{ - MemberID: "robot-001", - RobotEmail: "robot@example.com", - } - - prefs := &types.DeliveryPreferences{ - Email: &types.EmailPreference{ - Enabled: true, - Targets: []types.EmailTarget{ - { - To: []string{"user1@example.com"}, - Subject: "Report for User 1", - }, - { - To: []string{"user2@example.com", "user3@example.com"}, - Subject: "Report for Team", - }, - }, - }, - } - - results, _ := center.Deliver(ctx, content, deliveryCtx, prefs, robot) - - require.Len(t, results, 2) - assert.Equal(t, types.DeliveryEmail, results[0].Type) - assert.Equal(t, types.DeliveryEmail, results[1].Type) - assert.Equal(t, "user1@example.com", results[0].Target) - assert.Equal(t, "user2@example.com,user3@example.com", results[1].Target) - }) - - t.Run("sends email with attachments", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Report with Attachments", - Body: "Please find the attached report.", - Attachments: []types.DeliveryAttachment{ - { - Title: "Q4 Report.pdf", - Description: "Quarterly sales report", - TaskID: "task-001", - File: "__local://reports/q4-2024.pdf", - }, - { - Title: "Data Export.csv", - Description: "Raw data export", - TaskID: "task-002", - File: "__local://exports/data.csv", - }, - }, - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-003", - TriggerType: types.TriggerClock, - TeamID: "team-001", - } - - robot := &types.Robot{ - MemberID: "robot-001", - RobotEmail: "robot@example.com", - } - - prefs := &types.DeliveryPreferences{ - Email: &types.EmailPreference{ - Enabled: true, - Targets: []types.EmailTarget{ - { - To: []string{"manager@example.com"}, - Subject: "Weekly Report with Attachments", - }, - }, - }, - } - - // Send - attachment conversion may fail if files don't exist, but structure is tested - results, _ := center.Deliver(ctx, content, deliveryCtx, prefs, robot) - - require.Len(t, results, 1) - assert.Equal(t, types.DeliveryEmail, results[0].Type) - }) -} - -// ============================================================================ -// Process Delivery Tests (requires Yao process setup) -// ============================================================================ - -func TestDeliveryCenterProcess(t *testing.T) { - if testing.Short() { - t.Skip("Skipping integration test") - } - - testutils.Prepare(t) - defer testutils.Clean(t) - - t.Run("calls process with delivery content", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Process Test Summary", - Body: "This is the body content for process testing.", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-process-001", - TriggerType: types.TriggerEvent, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Process: &types.ProcessPreference{ - Enabled: true, - Targets: []types.ProcessTarget{ - { - Process: "scripts.tests.delivery.Handle", - }, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - require.NoError(t, err) - require.Len(t, results, 1) - assert.Equal(t, types.DeliveryProcess, results[0].Type) - assert.Equal(t, "scripts.tests.delivery.Handle", results[0].Target) - assert.True(t, results[0].Success) - - // Verify process received correct data (Details structure depends on process return) - assert.NotNil(t, results[0].Details, "process should return details") - }) - - t.Run("calls process with additional args", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Args Test", - Body: "Testing additional arguments", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-process-002", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Process: &types.ProcessPreference{ - Enabled: true, - Targets: []types.ProcessTarget{ - { - Process: "scripts.tests.delivery.Handle", - Args: []interface{}{"custom-arg-1", "custom-arg-2"}, - }, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - require.NoError(t, err) - require.Len(t, results, 1) - assert.True(t, results[0].Success) - assert.NotNil(t, results[0].Details, "process should return details with args") - }) - - t.Run("calls multiple process targets", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Multi-process Test", - Body: "Testing multiple process targets", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-process-003", - TriggerType: types.TriggerClock, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Process: &types.ProcessPreference{ - Enabled: true, - Targets: []types.ProcessTarget{ - { - Process: "scripts.tests.delivery.Handle", - }, - { - Process: "scripts.tests.delivery.Notify", - Args: []interface{}{"user-123", "push"}, - }, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - require.NoError(t, err) - require.Len(t, results, 2) - - // First process: Handle - assert.Equal(t, "scripts.tests.delivery.Handle", results[0].Target) - assert.True(t, results[0].Success) - - // Second process: Notify - assert.Equal(t, "scripts.tests.delivery.Notify", results[1].Target) - assert.True(t, results[1].Success) - }) - - t.Run("handles process failure gracefully", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Failure Test", - Body: "Testing process failure handling", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-process-004", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Process: &types.ProcessPreference{ - Enabled: true, - Targets: []types.ProcessTarget{ - { - Process: "scripts.tests.delivery.HandleWithFailure", - Args: []interface{}{true}, // shouldFail = true - }, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - assert.Error(t, err) - require.Len(t, results, 1) - assert.False(t, results[0].Success) - assert.Contains(t, results[0].Error, "Simulated process failure") - }) - - t.Run("handles process with attachments", func(t *testing.T) { - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Attachment Process Test", - Body: "Testing process with attachments", - Attachments: []types.DeliveryAttachment{ - { - Title: "Report.pdf", - Description: "Test report", - TaskID: "task-001", - File: "__local://test/report.pdf", - }, - }, - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-process-005", - TriggerType: types.TriggerClock, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Process: &types.ProcessPreference{ - Enabled: true, - Targets: []types.ProcessTarget{ - { - Process: "scripts.tests.delivery.HandleAttachments", - }, - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - require.NoError(t, err) - require.Len(t, results, 1) - assert.True(t, results[0].Success) - assert.NotNil(t, results[0].Details, "process should return details with attachments info") - }) -} - -// ============================================================================ -// Mixed Channel Tests -// ============================================================================ - -func TestDeliveryCenterAllChannels(t *testing.T) { - if testing.Short() { - t.Skip("Skipping integration test") - } - - testutils.Prepare(t) - defer testutils.Clean(t) - - // Set robot channel as default - original := types.DefaultEmailChannel() - types.SetDefaultEmailChannel("robot") - defer types.SetDefaultEmailChannel(original) - - t.Run("delivers to email, webhook, and process simultaneously", func(t *testing.T) { - // Setup webhook server - webhookCalled := false - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - webhookCalled = true - w.WriteHeader(http.StatusOK) - w.Write([]byte(`{"ok": true}`)) - })) - defer server.Close() - - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Full Channel Test", - Body: "Testing all delivery channels together", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-full-001", - TriggerType: types.TriggerClock, - TeamID: "team-001", - } - - robot := &types.Robot{ - MemberID: "robot-001", - RobotEmail: "robot@example.com", - } - - prefs := &types.DeliveryPreferences{ - Email: &types.EmailPreference{ - Enabled: true, - Targets: []types.EmailTarget{ - {To: []string{"user@example.com"}, Subject: "Test"}, - }, - }, - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - {URL: server.URL}, - }, - }, - Process: &types.ProcessPreference{ - Enabled: true, - Targets: []types.ProcessTarget{ - {Process: "scripts.tests.delivery.Handle"}, - }, - }, - } - - results, _ := center.Deliver(ctx, content, deliveryCtx, prefs, robot) - - // Should have 3 results (1 email + 1 webhook + 1 process) - require.Len(t, results, 3) - - // Verify each channel type - var emailResult, webhookResult, processResult *types.ChannelResult - for i := range results { - switch results[i].Type { - case types.DeliveryEmail: - emailResult = &results[i] - case types.DeliveryWebhook: - webhookResult = &results[i] - case types.DeliveryProcess: - processResult = &results[i] - } - } - - assert.NotNil(t, emailResult, "should have email result") - assert.NotNil(t, webhookResult, "should have webhook result") - assert.NotNil(t, processResult, "should have process result") - - // Webhook and process should succeed - assert.True(t, webhookCalled, "webhook should be called") - assert.True(t, webhookResult.Success, "webhook should succeed") - assert.True(t, processResult.Success, "process should succeed") - }) - - t.Run("partial failure does not stop other channels", func(t *testing.T) { - // Webhook that fails - failServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusInternalServerError) - })) - defer failServer.Close() - - // Webhook that succeeds - successCalled := false - successServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - successCalled = true - w.WriteHeader(http.StatusOK) - })) - defer successServer.Close() - - center := standard.NewDeliveryCenter() - ctx := types.NewContext(context.Background(), nil) - - content := &types.DeliveryContent{ - Summary: "Partial Failure Test", - Body: "Testing partial failure handling", - } - - deliveryCtx := &types.DeliveryContext{ - MemberID: "member-001", - ExecutionID: "exec-partial-001", - TriggerType: types.TriggerHuman, - TeamID: "team-001", - } - - prefs := &types.DeliveryPreferences{ - Webhook: &types.WebhookPreference{ - Enabled: true, - Targets: []types.WebhookTarget{ - {URL: failServer.URL}, // This will fail - {URL: successServer.URL}, // This should still be called - }, - }, - Process: &types.ProcessPreference{ - Enabled: true, - Targets: []types.ProcessTarget{ - {Process: "scripts.tests.delivery.Handle"}, // This should succeed - }, - }, - } - - results, err := center.Deliver(ctx, content, deliveryCtx, prefs, nil) - - // Should have error (from first webhook failure) - assert.Error(t, err) - - // But all targets should be attempted - require.Len(t, results, 3) - - // First webhook failed - assert.False(t, results[0].Success) - - // Second webhook and process should succeed - assert.True(t, successCalled, "second webhook should be called despite first failure") - assert.True(t, results[1].Success, "second webhook should succeed") - assert.True(t, results[2].Success, "process should succeed") - }) -} - // ============================================================================ // Helper Functions // ============================================================================