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.
This commit is contained in:
parent
0c4e15463d
commit
5f5bcba2b2
8 changed files with 586 additions and 1566 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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"`
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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>
|
||||
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))
|
||||
}
|
||||
|
|
@ -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
|
||||
// ============================================================================
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue