From 8006e6d64c9eef3b3c3c8b68c8165428aa410915 Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 06:26:01 +0700 Subject: [PATCH 1/8] feat(rules): add Rule interface, Pipeline, and WhitelistRule Co-Authored-By: Claude Opus 4.7 --- internal/rules/pipeline.go | 30 +++++++++ internal/rules/pipeline_test.go | 106 ++++++++++++++++++++++++++++++++ internal/rules/types.go | 41 ++++++++++++ internal/rules/types_test.go | 42 +++++++++++++ internal/rules/whitelist.go | 30 +++++++++ 5 files changed, 249 insertions(+) create mode 100644 internal/rules/pipeline.go create mode 100644 internal/rules/pipeline_test.go create mode 100644 internal/rules/types.go create mode 100644 internal/rules/types_test.go create mode 100644 internal/rules/whitelist.go diff --git a/internal/rules/pipeline.go b/internal/rules/pipeline.go new file mode 100644 index 0000000..63f89bd --- /dev/null +++ b/internal/rules/pipeline.go @@ -0,0 +1,30 @@ +package rules + +// Pipeline evaluates a sequence of rules, stopping on first rejection. +type Pipeline struct { + rules []Rule +} + +// NewPipeline creates a rules pipeline from the given rules. +func NewPipeline(rules []Rule) *Pipeline { + return &Pipeline{rules: rules} +} + +// Evaluate runs all rules in order. Returns the first rejection or a merged accept. +func (p *Pipeline) Evaluate(ctx EmailContext) RuleResult { + merged := Accept() + for _, r := range p.rules { + result := r.Evaluate(ctx) + if !result.Accepted { + return result + } + for k, v := range result.Metadata { + if merged.Metadata == nil { + merged.Metadata = map[string]string{k: v} + } else { + merged.Metadata[k] = v + } + } + } + return merged +} \ No newline at end of file diff --git a/internal/rules/pipeline_test.go b/internal/rules/pipeline_test.go new file mode 100644 index 0000000..ac37a5f --- /dev/null +++ b/internal/rules/pipeline_test.go @@ -0,0 +1,106 @@ +package rules + +import ( + "testing" +) + +func TestPipeline_NoRules_Accepts(t *testing.T) { + p := NewPipeline(nil) + result := p.Evaluate(EmailContext{Sender: "anyone@test.com"}) + if !result.Accepted { + t.Fatal("expected accept with no rules") + } +} + +func TestPipeline_AllRulesAccept_Accepts(t *testing.T) { + allowAll := &stubRule{name: "allow-all", result: Accept()} + p := NewPipeline([]Rule{allowAll}) + result := p.Evaluate(EmailContext{Sender: "a@b.com"}) + if !result.Accepted { + t.Fatal("expected accept") + } +} + +func TestPipeline_OneRuleRejects_Rejects(t *testing.T) { + rejector := &stubRule{name: "rejector", result: Reject("nope")} + p := NewPipeline([]Rule{rejector}) + result := p.Evaluate(EmailContext{Sender: "a@b.com"}) + if result.Accepted { + t.Fatal("expected reject") + } + if result.Reason != "nope" { + t.Fatalf("expected reason 'nope', got %s", result.Reason) + } +} + +func TestPipeline_MergesMetadata(t *testing.T) { + r1 := &stubRule{name: "r1", result: Accept().WithMetadata("a", "1")} + r2 := &stubRule{name: "r2", result: Accept().WithMetadata("b", "2")} + p := NewPipeline([]Rule{r1, r2}) + result := p.Evaluate(EmailContext{}) + if result.Metadata["a"] != "1" { + t.Fatal("missing metadata a") + } + if result.Metadata["b"] != "2" { + t.Fatal("missing metadata b") + } +} + +func TestPipeline_StopsOnFirstReject(t *testing.T) { + rejector := &stubRule{name: "rejector", result: Reject("stop")} + neverCalled := &trackingRule{} + p := NewPipeline([]Rule{rejector, neverCalled}) + p.Evaluate(EmailContext{}) + if neverCalled.called { + t.Fatal("expected second rule not to be called after rejection") + } +} + +func TestWhitelistRule_EmptyList_Accepts(t *testing.T) { + r := NewWhitelistRule(nil) + result := r.Evaluate(EmailContext{Sender: "anyone@test.com"}) + if !result.Accepted { + t.Fatal("empty whitelist should accept all") + } +} + +func TestWhitelistRule_SenderInList_Accepts(t *testing.T) { + r := NewWhitelistRule([]string{"allowed@test.com"}) + result := r.Evaluate(EmailContext{Sender: "allowed@test.com"}) + if !result.Accepted { + t.Fatal("expected accept for whitelisted sender") + } +} + +func TestWhitelistRule_SenderNotInList_Rejects(t *testing.T) { + r := NewWhitelistRule([]string{"allowed@test.com"}) + result := r.Evaluate(EmailContext{Sender: "unknown@test.com"}) + if result.Accepted { + t.Fatal("expected reject for non-whitelisted sender") + } +} + +func TestWhitelistRule_CaseInsensitive(t *testing.T) { + r := NewWhitelistRule([]string{"Allowed@Test.com"}) + result := r.Evaluate(EmailContext{Sender: "allowed@test.com"}) + if !result.Accepted { + t.Fatal("whitelist should be case-insensitive") + } +} + +// stubs + +type stubRule struct { + name string + result RuleResult +} + +func (s *stubRule) Name() string { return s.name } +func (s *stubRule) Evaluate(_ EmailContext) RuleResult { return s.result } + +type trackingRule struct { + called bool +} + +func (t *trackingRule) Name() string { return "tracker" } +func (t *trackingRule) Evaluate(_ EmailContext) RuleResult { t.called = true; return Accept() } \ No newline at end of file diff --git a/internal/rules/types.go b/internal/rules/types.go new file mode 100644 index 0000000..8c21fe8 --- /dev/null +++ b/internal/rules/types.go @@ -0,0 +1,41 @@ +package rules + +// EmailContext contains the data available to rules for evaluation. +type EmailContext struct { + Sender string + Subject string + MessageID string +} + +// RuleResult is the output of evaluating a single rule. +type RuleResult struct { + Accepted bool + Reason string + Metadata map[string]string +} + +// Accept returns a passing RuleResult. +func Accept() RuleResult { + return RuleResult{Accepted: true} +} + +// Reject returns a failing RuleResult with a reason. +func Reject(reason string) RuleResult { + return RuleResult{Accepted: false, Reason: reason} +} + +// WithMetadata attaches a key-value pair to the result. +func (r RuleResult) WithMetadata(key, value string) RuleResult { + if r.Metadata == nil { + r.Metadata = map[string]string{key: value} + } else { + r.Metadata[key] = value + } + return r +} + +// Rule evaluates an email against operator-configurable policy. +type Rule interface { + Name() string + Evaluate(ctx EmailContext) RuleResult +} \ No newline at end of file diff --git a/internal/rules/types_test.go b/internal/rules/types_test.go new file mode 100644 index 0000000..2bc2be4 --- /dev/null +++ b/internal/rules/types_test.go @@ -0,0 +1,42 @@ +package rules + +import ( + "testing" +) + +func TestRuleResult_Accept(t *testing.T) { + r := Accept() + if !r.Accepted { + t.Fatal("expected Accepted=true") + } + if r.Reason != "" { + t.Fatal("expected no reason for accept") + } +} + +func TestRuleResult_Reject(t *testing.T) { + r := Reject("blocked sender") + if r.Accepted { + t.Fatal("expected Accepted=false") + } + if r.Reason != "blocked sender" { + t.Fatalf("expected reason 'blocked sender', got %s", r.Reason) + } +} + +func TestRuleResult_WithMetadata(t *testing.T) { + r := Accept().WithMetadata("routing", "high-priority") + if r.Metadata["routing"] != "high-priority" { + t.Fatalf("expected routing=high-priority, got %s", r.Metadata["routing"]) + } +} + +func TestEmailContext_Fields(t *testing.T) { + ctx := EmailContext{ + Sender: "user@example.com", + Subject: "Test", + } + if ctx.Sender != "user@example.com" { + t.Fatal("sender mismatch") + } +} \ No newline at end of file diff --git a/internal/rules/whitelist.go b/internal/rules/whitelist.go new file mode 100644 index 0000000..0eaca8a --- /dev/null +++ b/internal/rules/whitelist.go @@ -0,0 +1,30 @@ +package rules + +import "strings" + +// WhitelistRule accepts emails from configured senders. +// An empty list accepts all senders. +type WhitelistRule struct { + emails map[string]bool +} + +// NewWhitelistRule creates a whitelist rule from a list of email addresses. +func NewWhitelistRule(emails []string) *WhitelistRule { + m := make(map[string]bool, len(emails)) + for _, e := range emails { + m[strings.ToLower(e)] = true + } + return &WhitelistRule{emails: m} +} + +func (r *WhitelistRule) Name() string { return "whitelist" } + +func (r *WhitelistRule) Evaluate(ctx EmailContext) RuleResult { + if len(r.emails) == 0 { + return Accept() + } + if r.emails[strings.ToLower(ctx.Sender)] { + return Accept() + } + return Reject("not whitelisted") +} \ No newline at end of file -- 2.54.0 From ab992abc268252720acfaae518bf66dcad8c57e4 Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 06:26:56 +0700 Subject: [PATCH 2/8] feat(rules): add BlocklistRule Co-Authored-By: Claude Opus 4.7 --- internal/rules/blocklist.go | 26 ++++++++++++++++++++++++ internal/rules/blocklist_test.go | 35 ++++++++++++++++++++++++++++++++ 2 files changed, 61 insertions(+) create mode 100644 internal/rules/blocklist.go create mode 100644 internal/rules/blocklist_test.go diff --git a/internal/rules/blocklist.go b/internal/rules/blocklist.go new file mode 100644 index 0000000..643b6aa --- /dev/null +++ b/internal/rules/blocklist.go @@ -0,0 +1,26 @@ +package rules + +import "strings" + +// BlocklistRule rejects emails from configured senders. +type BlocklistRule struct { + emails map[string]bool +} + +// NewBlocklistRule creates a blocklist rule from a list of email addresses. +func NewBlocklistRule(emails []string) *BlocklistRule { + m := make(map[string]bool, len(emails)) + for _, e := range emails { + m[strings.ToLower(e)] = true + } + return &BlocklistRule{emails: m} +} + +func (r *BlocklistRule) Name() string { return "blocklist" } + +func (r *BlocklistRule) Evaluate(ctx EmailContext) RuleResult { + if r.emails[strings.ToLower(ctx.Sender)] { + return Reject("blocked sender") + } + return Accept() +} diff --git a/internal/rules/blocklist_test.go b/internal/rules/blocklist_test.go new file mode 100644 index 0000000..3081885 --- /dev/null +++ b/internal/rules/blocklist_test.go @@ -0,0 +1,35 @@ +package rules + +import "testing" + +func TestBlocklistRule_SenderBlocked_Rejects(t *testing.T) { + r := NewBlocklistRule([]string{"spam@test.com"}) + result := r.Evaluate(EmailContext{Sender: "spam@test.com"}) + if result.Accepted { + t.Fatal("expected reject for blocked sender") + } +} + +func TestBlocklistRule_SenderNotBlocked_Accepts(t *testing.T) { + r := NewBlocklistRule([]string{"spam@test.com"}) + result := r.Evaluate(EmailContext{Sender: "ok@test.com"}) + if !result.Accepted { + t.Fatal("expected accept for non-blocked sender") + } +} + +func TestBlocklistRule_EmptyList_Accepts(t *testing.T) { + r := NewBlocklistRule(nil) + result := r.Evaluate(EmailContext{Sender: "anyone@test.com"}) + if !result.Accepted { + t.Fatal("empty blocklist should accept all") + } +} + +func TestBlocklistRule_CaseInsensitive(t *testing.T) { + r := NewBlocklistRule([]string{"Spam@Test.com"}) + result := r.Evaluate(EmailContext{Sender: "spam@test.com"}) + if result.Accepted { + t.Fatal("blocklist should be case-insensitive") + } +} -- 2.54.0 From aa05e3f7c0c09b0e2494d06a722d430e37370fb5 Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 06:38:37 +0700 Subject: [PATCH 3/8] feat: wire external rules pipeline into ingress and dispatch - Add DispatchContext field to Task model for rule metadata storage - Add BlocklistEmails config (BLOCKLIST_EMAILS env var) - Wire rules.Pipeline into IMAPWatcher, replacing inline whitelist check - Pass rule metadata through DispatchContext to OpenClaw dispatch - Remove isWhitelisted method (now handled by WhitelistRule) Co-Authored-By: Claude Opus 4.7 --- internal/ai_client/client.go | 17 +++++++++- internal/config/config.go | 9 +++++ internal/database/models.go | 1 + internal/mail/imap.go | 59 ++++++++++++++++++++------------- internal/rules/pipeline.go | 2 +- internal/rules/pipeline_test.go | 6 ++-- internal/rules/types.go | 2 +- internal/rules/types_test.go | 2 +- internal/rules/whitelist.go | 2 +- 9 files changed, 69 insertions(+), 31 deletions(-) diff --git a/internal/ai_client/client.go b/internal/ai_client/client.go index c395f06..1a8b589 100644 --- a/internal/ai_client/client.go +++ b/internal/ai_client/client.go @@ -73,7 +73,7 @@ func (d *Dispatcher) Dispatch(task *database.Task) { SessionID: task.ThreadID, History: history, CallbackURL: callbackURL, - Metadata: map[string]string{"task_uuid": task.TaskUUID}, + Metadata: buildDispatchMetadata(task), } // Update status to AI_PROCESSING before first attempt. @@ -195,3 +195,18 @@ func (d *Dispatcher) buildHistory(threadID string) ([]historyEntry, error) { return history, nil } + +// buildDispatchMetadata constructs the metadata map for OpenClaw dispatch, +// merging the task UUID with any DispatchContext from the rules pipeline. +func buildDispatchMetadata(task *database.Task) map[string]string { + metadata := map[string]string{"task_uuid": task.TaskUUID} + if task.DispatchContext != "" { + var ctx map[string]string + if err := json.Unmarshal([]byte(task.DispatchContext), &ctx); err == nil { + for k, v := range ctx { + metadata[k] = v + } + } + } + return metadata +} diff --git a/internal/config/config.go b/internal/config/config.go index 13e8d1d..e0adcb0 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -31,6 +31,7 @@ type Config struct { BridgeCallbackToken string SystemEmail string WhitelistEmails []string + BlocklistEmails []string // Optional IMAPProxyURL string @@ -74,6 +75,14 @@ func Load() (*Config, error) { } } + if raw := os.Getenv("BLOCKLIST_EMAILS"); raw != "" { + for _, email := range strings.Split(raw, ",") { + if trimmed := strings.TrimSpace(email); trimmed != "" { + cfg.BlocklistEmails = append(cfg.BlocklistEmails, trimmed) + } + } + } + if cfg.ListenAddr == "" { cfg.ListenAddr = ":8080" } diff --git a/internal/database/models.go b/internal/database/models.go index f779396..a1cbac0 100644 --- a/internal/database/models.go +++ b/internal/database/models.go @@ -24,6 +24,7 @@ type Task struct { Sender string Subject string BodyPlain string + DispatchContext string // JSON metadata from external rules pipeline. AIResponse string Status string `gorm:"index;not null"` AttemptOpenClaw int `gorm:"default:0"` diff --git a/internal/mail/imap.go b/internal/mail/imap.go index 6fe6f32..da22273 100644 --- a/internal/mail/imap.go +++ b/internal/mail/imap.go @@ -3,6 +3,7 @@ package mail import ( "context" "crypto/tls" + "encoding/json" "fmt" "io" "log/slog" @@ -21,6 +22,7 @@ import ( "thuanle.me/claw-email-bridge/internal/config" "thuanle.me/claw-email-bridge/internal/database" "thuanle.me/claw-email-bridge/internal/logging" + "thuanle.me/claw-email-bridge/internal/rules" ) // IMAPWatcher monitors an IMAP mailbox via IDLE and processes new emails. @@ -32,6 +34,8 @@ type IMAPWatcher struct { dialIMAP func(addr string, options *imapclient.Options) (*imapclient.Client, error) dialIMAPViaProxy func(addr, proxyURL string, options *imapclient.Options) (*imapclient.Client, error) + rules *rules.Pipeline + // onReceived is called after a task is saved as RECEIVED. // This will be wired to the OpenClaw dispatch in a future step. OnReceived func(task *database.Task) @@ -57,11 +61,20 @@ var socks5DialerFactory = proxy.SOCKS5 // NewIMAPWatcher creates a new IMAP watcher. func NewIMAPWatcher(cfg *config.Config, db *gorm.DB) *IMAPWatcher { + var externalRules []rules.Rule + if len(cfg.WhitelistEmails) > 0 { + externalRules = append(externalRules, rules.NewWhitelistRule(cfg.WhitelistEmails)) + } + if len(cfg.BlocklistEmails) > 0 { + externalRules = append(externalRules, rules.NewBlocklistRule(cfg.BlocklistEmails)) + } + return &IMAPWatcher{ cfg: cfg, db: db, dialIMAP: imapclient.DialTLS, dialIMAPViaProxy: dialTLSViaSOCKS5, + rules: rules.NewPipeline(externalRules), } } @@ -343,10 +356,15 @@ func (w *IMAPWatcher) processMessage(c *imapclient.Client, buf *imapclient.Fetch return } - // 2. Whitelist: skip if sender not in allowed list. - if !w.isWhitelisted(sender) { - log.Info("imap: sender not whitelisted, ignoring", "sender", sender) - w.saveIgnored(messageID, sender, subject, "not whitelisted") + // 2. External rules: evaluate operator-configurable policy. + ruleResult := w.rules.Evaluate(rules.EmailContext{ + Sender: sender, + Subject: subject, + MessageID: messageID, + }) + if !ruleResult.Accepted { + log.Info("imap: rejected by external rules", "sender", sender, "reason", ruleResult.Reason) + w.saveIgnored(messageID, sender, subject, ruleResult.Reason) w.markSeen(c, buf.UID) return } @@ -404,13 +422,14 @@ func (w *IMAPWatcher) processMessage(c *imapclient.Client, buf *imapclient.Fetch // 5. Save task with status RECEIVED. taskUUID := uuid.New().String() task := &database.Task{ - TaskUUID: taskUUID, - ThreadID: threadID, - MessageID: messageID, - Sender: sender, - Subject: subject, - BodyPlain: bodyPlain, - Status: database.StatusReceived, + TaskUUID: taskUUID, + ThreadID: threadID, + MessageID: messageID, + Sender: sender, + Subject: subject, + BodyPlain: bodyPlain, + DispatchContext: marshalDispatchContext(ruleResult.Metadata), + Status: database.StatusReceived, } taskLog := logging.TaskLogger(taskUUID, threadID, messageID) @@ -453,19 +472,13 @@ func (w *IMAPWatcher) saveIgnored(messageID, sender, subject, reason string) { } } -// isWhitelisted checks if the sender is in the allowed list. -// Empty whitelist means all senders are allowed. -func (w *IMAPWatcher) isWhitelisted(sender string) bool { - if len(w.cfg.WhitelistEmails) == 0 { - return true +// marshalDispatchContext serializes rule metadata to JSON. +func marshalDispatchContext(metadata map[string]string) string { + if len(metadata) == 0 { + return "" } - lower := strings.ToLower(sender) - for _, email := range w.cfg.WhitelistEmails { - if strings.ToLower(email) == lower { - return true - } - } - return false + data, _ := json.Marshal(metadata) + return string(data) } // markSeen flags a message as \Seen in IMAP. diff --git a/internal/rules/pipeline.go b/internal/rules/pipeline.go index 63f89bd..d8debb0 100644 --- a/internal/rules/pipeline.go +++ b/internal/rules/pipeline.go @@ -27,4 +27,4 @@ func (p *Pipeline) Evaluate(ctx EmailContext) RuleResult { } } return merged -} \ No newline at end of file +} diff --git a/internal/rules/pipeline_test.go b/internal/rules/pipeline_test.go index ac37a5f..fe6f951 100644 --- a/internal/rules/pipeline_test.go +++ b/internal/rules/pipeline_test.go @@ -95,12 +95,12 @@ type stubRule struct { result RuleResult } -func (s *stubRule) Name() string { return s.name } +func (s *stubRule) Name() string { return s.name } func (s *stubRule) Evaluate(_ EmailContext) RuleResult { return s.result } type trackingRule struct { called bool } -func (t *trackingRule) Name() string { return "tracker" } -func (t *trackingRule) Evaluate(_ EmailContext) RuleResult { t.called = true; return Accept() } \ No newline at end of file +func (t *trackingRule) Name() string { return "tracker" } +func (t *trackingRule) Evaluate(_ EmailContext) RuleResult { t.called = true; return Accept() } diff --git a/internal/rules/types.go b/internal/rules/types.go index 8c21fe8..3afb3f7 100644 --- a/internal/rules/types.go +++ b/internal/rules/types.go @@ -38,4 +38,4 @@ func (r RuleResult) WithMetadata(key, value string) RuleResult { type Rule interface { Name() string Evaluate(ctx EmailContext) RuleResult -} \ No newline at end of file +} diff --git a/internal/rules/types_test.go b/internal/rules/types_test.go index 2bc2be4..1ac2d9a 100644 --- a/internal/rules/types_test.go +++ b/internal/rules/types_test.go @@ -39,4 +39,4 @@ func TestEmailContext_Fields(t *testing.T) { if ctx.Sender != "user@example.com" { t.Fatal("sender mismatch") } -} \ No newline at end of file +} diff --git a/internal/rules/whitelist.go b/internal/rules/whitelist.go index 0eaca8a..a3bead9 100644 --- a/internal/rules/whitelist.go +++ b/internal/rules/whitelist.go @@ -27,4 +27,4 @@ func (r *WhitelistRule) Evaluate(ctx EmailContext) RuleResult { return Accept() } return Reject("not whitelisted") -} \ No newline at end of file +} -- 2.54.0 From 92121278bb8e19f2c84b995424b6ecb26f5d7ea8 Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 06:39:50 +0700 Subject: [PATCH 4/8] docs: add ADR-0002, update flow and add BLOCKLIST_EMAILS MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - ADR-0002: external rules pipeline architecture - README: updated flow to show safety checks → rules → dispatch - .env.example: added BLOCKLIST_EMAILS Co-Authored-By: Claude Opus 4.7 --- .env.example | 1 + README.md | 3 ++- docs/adr/0002-external-rules-pipeline.md | 30 ++++++++++++++++++++++++ docs/adr/README.md | 1 + 4 files changed, 34 insertions(+), 1 deletion(-) create mode 100644 docs/adr/0002-external-rules-pipeline.md diff --git a/.env.example b/.env.example index 71cbc60..d7a2948 100644 --- a/.env.example +++ b/.env.example @@ -18,6 +18,7 @@ OPENCLAW_API_KEY= BRIDGE_CALLBACK_TOKEN=a-strong-random-token SYSTEM_EMAIL=bridge@example.com WHITELIST_EMAILS=user1@example.com,user2@example.com +BLOCKLIST_EMAILS= # === Optional === IMAP_PROXY_URL= diff --git a/README.md b/README.md index 014ad0c..69f3daf 100644 --- a/README.md +++ b/README.md @@ -3,7 +3,7 @@ Tự động xử lý email bằng OpenClaw AI theo luồng: ``` -Ingress (IMAP) → Dispatch (OpenClaw) → Callback → Egress (SMTP reply) +Ingress (IMAP) → Safety Checks → External Rules → Apply Context → Dispatch (OpenClaw) → Callback → Egress (SMTP reply) ``` ## Tính năng @@ -12,6 +12,7 @@ Ingress (IMAP) → Dispatch (OpenClaw) → Callback → Egress (SMTP reply) - **Threading** — giữ nguyên email thread qua `In-Reply-To` / `References` - **Idempotency** — chống duplicate theo `message_id` - **Anti-loop & Whitelist** — bỏ qua email từ chính hệ thống, chỉ xử lý sender được phép +- **External Rules** — pipeline configurable: whitelist, blocklist, và rule metadata cho dispatch - **Retry** — tối đa 3 lần với backoff 1s → 5s → 15s cho cả OpenClaw và SMTP - **Stateful** — SQLite lưu trạng thái pipeline, không mất dữ liệu khi restart diff --git a/docs/adr/0002-external-rules-pipeline.md b/docs/adr/0002-external-rules-pipeline.md new file mode 100644 index 0000000..f86859c --- /dev/null +++ b/docs/adr/0002-external-rules-pipeline.md @@ -0,0 +1,30 @@ +--- +type: ADR +id: "0002" +title: "Ingress uses a configurable external rules pipeline before dispatch" +status: active +date: 2026-04-28 +--- + +## Context + +The IMAP ingress flow previously hardcoded whitelist and idempotency checks directly in processMessage. There was no extension point for operator-defined policy such as sender blocking or context-based routing. The flow mixed internal safety rules (anti-loop) with configurable policy (whitelist). + +## Decision + +Separate internal safety checks (anti-loop, idempotency) from external configurable rules. External rules run in a pipeline after safety checks. Each rule implements a `Rule` interface that evaluates `EmailContext` and returns `RuleResult` (accept/reject + metadata). Rule metadata is persisted as `DispatchContext` on the Task and passed to OpenClaw dispatch. + +Pipeline order: anti-loop (internal) → idempotency (internal) → external rules → threading → save → dispatch. + +## Options considered + +- **Pipeline with Rule interface** (chosen): extensible, testable, each rule is isolated. +- **Keep everything in processMessage**: simpler but not configurable or testable. +- **Plugin-based rules via config files**: more flexible but over-engineered for current needs. + +## Consequences + +New rules can be added by implementing the Rule interface and registering in config. +Whitelist and blocklist are now env-configurable and independently testable. +Future rules (domain routing, context selection) fit the same pipeline. +Task model has a new DispatchContext column (auto-migrated). diff --git a/docs/adr/README.md b/docs/adr/README.md index 74d212a..18f416b 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -73,3 +73,4 @@ Optional. Capture important external input or review notes. | ID | Title | Status | |----|-------|--------| | 0001 | IMAP thread resolution prefers In-Reply-To and falls back to References | active | +| 0002 | Ingress uses a configurable external rules pipeline before dispatch | active | -- 2.54.0 From 562d3b01c393eb399d08efde9adeff002fcfc2ff Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 06:53:05 +0700 Subject: [PATCH 5/8] fix: reorder ingress checks and update requirements.md - Move idempotency before external rules per ADR-0002 - Update requirements.md: new pipeline flow, BLOCKLIST_EMAILS, dispatch_context metadata Co-Authored-By: Claude Opus 4.7 --- internal/mail/imap.go | 26 +++++++++++++------------- requirements.md | 8 +++++--- 2 files changed, 18 insertions(+), 16 deletions(-) diff --git a/internal/mail/imap.go b/internal/mail/imap.go index da22273..ab9bd12 100644 --- a/internal/mail/imap.go +++ b/internal/mail/imap.go @@ -356,7 +356,19 @@ func (w *IMAPWatcher) processMessage(c *imapclient.Client, buf *imapclient.Fetch return } - // 2. External rules: evaluate operator-configurable policy. + // 2. Idempotency: skip if message_id already exists. + existing, err := database.FindByMessageID(w.db, messageID) + if err != nil { + log.Error("imap: db lookup failed", "error", err) + return + } + if existing != nil { + log.Info("imap: duplicate message_id, skipping") + w.markSeen(c, buf.UID) + return + } + + // 3. External rules: evaluate operator-configurable policy. ruleResult := w.rules.Evaluate(rules.EmailContext{ Sender: sender, Subject: subject, @@ -369,18 +381,6 @@ func (w *IMAPWatcher) processMessage(c *imapclient.Client, buf *imapclient.Fetch return } - // 3. Idempotency: skip if message_id already exists. - existing, err := database.FindByMessageID(w.db, messageID) - if err != nil { - log.Error("imap: db lookup failed", "error", err) - return - } - if existing != nil { - log.Info("imap: duplicate message_id, skipping") - w.markSeen(c, buf.UID) - return - } - // 4. Threading: determine thread_id. threadID := messageID inReplyTo := "" diff --git a/requirements.md b/requirements.md index 9ea3f7c..26782fe 100644 --- a/requirements.md +++ b/requirements.md @@ -68,12 +68,12 @@ Danh sách trạng thái hợp lệ: - Nếu có `IMAP_PROXY_URL` thì kết nối IMAP qua proxy; nếu không thì direct. - Khi có email mới: 1. Anti-loop: nếu `sender` trùng email hệ thống -> `IGNORED`. - 2. Whitelist: nếu không nằm trong danh sách cho phép -> `IGNORED`. - 3. Idempotency: nếu `message_id` đã tồn tại -> bỏ qua. + 2. Idempotency: nếu `message_id` đã tồn tại -> bỏ qua. + 3. External rules: chạy pipeline các rule configurable (whitelist, blocklist, v.v.). Nếu reject -> `IGNORED`. Rule metadata được lưu vào `dispatch_context`. 4. Threading: - Nếu có `In-Reply-To`/`References`: tìm email cha trong DB, kế thừa `thread_id`. - Nếu không có: `thread_id = message_id`. - 5. Lưu DB với `status = RECEIVED`. + 5. Lưu DB với `status = RECEIVED`, bao gồm `dispatch_context` từ rules. 4.2 Dispatch (OpenClaw) - Lấy nội dung email hiện tại và history của cùng `thread_id` (chỉ các bản ghi `COMPLETED`). @@ -84,6 +84,7 @@ Danh sách trạng thái hợp lệ: - `history` (optional) - `callback_url` - `metadata.task_uuid` + - `metadata.*` từ `dispatch_context` (rule-derived context) - Cập nhật `status = AI_PROCESSING`. Retry OpenClaw: @@ -146,6 +147,7 @@ Required: - `OPENCLAW_URL`, `OPENCLAW_API_KEY` (nếu dùng) - `BRIDGE_CALLBACK_TOKEN` - `WHITELIST_EMAILS` (comma-separated) +- `BLOCKLIST_EMAILS` (comma-separated, optional) - `SYSTEM_EMAIL` - `IMAP_PROXY_URL` (optional) -- 2.54.0 From 5e7a479991218fe8849e4261997fe8d10f7a490e Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 07:01:20 +0700 Subject: [PATCH 6/8] fix: protect task_uuid from DispatchContext overwrite and update acceptance criteria - Skip task_uuid key when merging DispatchContext into dispatch metadata - Update requirements.md acceptance criteria and test matrix for external rules Co-Authored-By: Claude Opus 4.7 --- internal/ai_client/client.go | 4 +++- requirements.md | 7 +++++-- 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/internal/ai_client/client.go b/internal/ai_client/client.go index 1a8b589..fac4030 100644 --- a/internal/ai_client/client.go +++ b/internal/ai_client/client.go @@ -204,7 +204,9 @@ func buildDispatchMetadata(task *database.Task) map[string]string { var ctx map[string]string if err := json.Unmarshal([]byte(task.DispatchContext), &ctx); err == nil { for k, v := range ctx { - metadata[k] = v + if k != "task_uuid" { + metadata[k] = v + } } } } diff --git a/requirements.md b/requirements.md index 26782fe..5c91693 100644 --- a/requirements.md +++ b/requirements.md @@ -168,7 +168,7 @@ Ví dụ: 3. OpenClaw/SMTP retry đúng thứ tự 1s -> 5s -> 15s. 4. Hết retry thì `FAILED` và có `last_error`. 5. Callback sai token trả 401 và không đổi state task. -6. Anti-loop + whitelist + idempotency hoạt động đúng. +6. Anti-loop + idempotency + external rules hoạt động đúng thứ tự. 10) Test matrix tối thiểu @@ -177,4 +177,7 @@ Ví dụ: 3. SMTP fail hết retry -> `FAILED`, có `last_error`. 4. Callback sai token -> 401, state không đổi. 5. Duplicate `message_id` -> không tạo task mới. -6. Sender là chính hệ thống hoặc ngoài whitelist -> `IGNORED`. +6. Sender là chính hệ thống -> `IGNORED`. +7. Sender ngoài whitelist -> `IGNORED` bởi external rules. +8. Sender trong blocklist -> `IGNORED` bởi external rules. +9. Rule metadata được truyền qua dispatch context đến OpenClaw. -- 2.54.0 From c127018bc91202b3530f9282b144be8507ce70a1 Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 07:08:19 +0700 Subject: [PATCH 7/8] fix: add dispatch_context to schema docs and test buildDispatchMetadata - Add dispatch_context column to requirements.md task schema section - Export buildDispatchMetadata for test coverage - Add tests: context forwarding, task_uuid overwrite protection, empty context Co-Authored-By: Claude Opus 4.7 --- internal/ai_client/client.go | 5 ++++ internal/ai_client/client_test.go | 47 +++++++++++++++++++++++++++++++ requirements.md | 1 + 3 files changed, 53 insertions(+) diff --git a/internal/ai_client/client.go b/internal/ai_client/client.go index fac4030..f305c00 100644 --- a/internal/ai_client/client.go +++ b/internal/ai_client/client.go @@ -196,6 +196,11 @@ func (d *Dispatcher) buildHistory(threadID string) ([]historyEntry, error) { return history, nil } +// ExportBuildDispatchMetadata exports buildDispatchMetadata for testing. +func ExportBuildDispatchMetadata(task *database.Task) map[string]string { + return buildDispatchMetadata(task) +} + // buildDispatchMetadata constructs the metadata map for OpenClaw dispatch, // merging the task UUID with any DispatchContext from the rules pipeline. func buildDispatchMetadata(task *database.Task) map[string]string { diff --git a/internal/ai_client/client_test.go b/internal/ai_client/client_test.go index 6172e3e..08f12c7 100644 --- a/internal/ai_client/client_test.go +++ b/internal/ai_client/client_test.go @@ -12,6 +12,53 @@ import ( "thuanle.me/claw-email-bridge/internal/database" ) +func TestBuildDispatchMetadata_ForwardsContext(t *testing.T) { + task := &database.Task{ + TaskUUID: "uuid-ctx-test", + DispatchContext: `{"rule":"whitelist","source":"env"}`, + } + metadata := ai_client.ExportBuildDispatchMetadata(task) + + if metadata["task_uuid"] != "uuid-ctx-test" { + t.Errorf("expected task_uuid=uuid-ctx-test, got %q", metadata["task_uuid"]) + } + if metadata["rule"] != "whitelist" { + t.Errorf("expected rule=whitelist, got %q", metadata["rule"]) + } + if metadata["source"] != "env" { + t.Errorf("expected source=env, got %q", metadata["source"]) + } +} + +func TestBuildDispatchMetadata_ProtectsTaskUUID(t *testing.T) { + task := &database.Task{ + TaskUUID: "original-uuid", + DispatchContext: `{"task_uuid":"malicious-override","extra":"data"}`, + } + metadata := ai_client.ExportBuildDispatchMetadata(task) + + if metadata["task_uuid"] != "original-uuid" { + t.Errorf("task_uuid should not be overwritten, got %q", metadata["task_uuid"]) + } + if metadata["extra"] != "data" { + t.Errorf("expected extra=data, got %q", metadata["extra"]) + } +} + +func TestBuildDispatchMetadata_NoContext(t *testing.T) { + task := &database.Task{ + TaskUUID: "uuid-no-ctx", + } + metadata := ai_client.ExportBuildDispatchMetadata(task) + + if len(metadata) != 1 { + t.Errorf("expected 1 key, got %d", len(metadata)) + } + if metadata["task_uuid"] != "uuid-no-ctx" { + t.Errorf("expected task_uuid=uuid-no-ctx, got %q", metadata["task_uuid"]) + } +} + // Test matrix #2: OpenClaw fail 2 lần, lần 3 thành công → COMPLETED, attempt_openclaw = 3. func TestDispatch_RetryThenSuccess(t *testing.T) { tdb := database.NewTestDB(t) diff --git a/requirements.md b/requirements.md index 5c91693..7f8e58a 100644 --- a/requirements.md +++ b/requirements.md @@ -47,6 +47,7 @@ Bảng `tasks`: - `status` TEXT - `attempt_openclaw` INTEGER DEFAULT 0 - `attempt_smtp` INTEGER DEFAULT 0 +- `dispatch_context` TEXT (JSON metadata từ external rules pipeline, truyền qua OpenClaw dispatch) - `last_error` TEXT NULL - `next_retry_at` DATETIME NULL (chỉ cần nếu retry xử lý theo scheduler/worker) - `created_at` DATETIME -- 2.54.0 From 094bd35de44e838c248bde342b7e71465a91bffd Mon Sep 17 00:00:00 2001 From: thuanle Date: Tue, 28 Apr 2026 07:12:49 +0700 Subject: [PATCH 8/8] style: gofmt client_test.go Co-Authored-By: Claude Opus 4.7 --- internal/ai_client/client_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/internal/ai_client/client_test.go b/internal/ai_client/client_test.go index 08f12c7..a6bd822 100644 --- a/internal/ai_client/client_test.go +++ b/internal/ai_client/client_test.go @@ -14,8 +14,8 @@ import ( func TestBuildDispatchMetadata_ForwardsContext(t *testing.T) { task := &database.Task{ - TaskUUID: "uuid-ctx-test", - DispatchContext: `{"rule":"whitelist","source":"env"}`, + TaskUUID: "uuid-ctx-test", + DispatchContext: `{"rule":"whitelist","source":"env"}`, } metadata := ai_client.ExportBuildDispatchMetadata(task) -- 2.54.0