Author SHA1 Message Date
thuanleandClaude Opus 4.7 f978ed44b8 feat: support http/https proxy for IMAP via HTTP CONNECT tunnel
CI / fmt (pull_request) Successful in 4m45s
CI / test (pull_request) Successful in 10m47s
Extend IMAP_PROXY_URL to accept http:// and https:// schemes in addition
to the existing socks5:// support. HTTP CONNECT tunneling is used for
both new schemes, with TLS to the proxy for https://. Proxy auth via
user:pass@ in the URL is supported for all schemes.

Refactors dialTLSViaSOCKS5 into a generic dialTLSViaProxy dispatcher.

Fixes #9

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-04-28 10:56:57 +07:00
thuanle b5941e5aa1 update agent instruction 2026-04-28 09:58:47 +07:00
thuanle c112ac37c9 docs: add design spec for issue 18 CI checks
Record the approved design to add go vet and pinned staticcheck jobs to PR CI with minimal workflow changes.
2026-04-28 07:25:52 +07:00
thuanleandClaude Opus 4.7 094bd35de4 style: gofmt client_test.go
CI / fmt (pull_request) Successful in 4m46s
CI / test (pull_request) Successful in 12m58s
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-04-28 07:12:49 +07:00
thuanleandClaude Opus 4.7 c127018bc9 fix: add dispatch_context to schema docs and test buildDispatchMetadata
CI / fmt (pull_request) Failing after 4m45s
CI / test (pull_request) Successful in 10m47s
- 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 <noreply@anthropic.com>
2026-04-28 07:08:19 +07:00
thuanleandClaude Opus 4.7 5e7a479991 fix: protect task_uuid from DispatchContext overwrite and update acceptance criteria
CI / fmt (pull_request) Successful in 4m44s
CI / test (pull_request) Successful in 12m48s
- 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 <noreply@anthropic.com>
2026-04-28 07:01:20 +07:00
thuanleandClaude Opus 4.7 562d3b01c3 fix: reorder ingress checks and update requirements.md
CI / fmt (pull_request) Successful in 4m43s
CI / test (pull_request) Failing after 21m37s
- 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 <noreply@anthropic.com>
2026-04-28 06:53:05 +07:00
thuanleandClaude Opus 4.7 92121278bb docs: add ADR-0002, update flow and add BLOCKLIST_EMAILS
CI / fmt (pull_request) Successful in 4m47s
CI / test (pull_request) Successful in 12m54s
- 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 <noreply@anthropic.com>
2026-04-28 06:39:50 +07:00
thuanleandClaude Opus 4.7 aa05e3f7c0 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 <noreply@anthropic.com>
2026-04-28 06:38:37 +07:00
thuanleandClaude Opus 4.7 ab992abc26 feat(rules): add BlocklistRule
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-04-28 06:26:56 +07:00
thuanleandClaude Opus 4.7 8006e6d64c feat(rules): add Rule interface, Pipeline, and WhitelistRule
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-04-28 06:26:01 +07:00
thuanle 2233fd232a update workflow - auto assign yourself as assignee/reviewer 2026-04-28 06:23:28 +07:00
thuanle 3985b9dae3 Merge pull request 'Require PR discussion comments in reviewer workflow' (#17) from fix/reviewer-pr-discussion-comments into main
Reviewed-on: #17
2026-04-28 06:13:27 +07:00
thuanle 28224ccb9f docs: require PR discussion comments in reviews
CI / fmt (pull_request) Successful in 4m46s
CI / test (pull_request) Successful in 10m45s
2026-04-28 06:10:55 +07:00
thuanle 668b64d594 Merge pull request 'Add ADR workflow, update skills, rename do-task to implement-task' (#15) from feat/adr-workflow into main
Reviewed-on: #15
2026-04-28 05:56:21 +07:00
24 changed files with 887 additions and 255 deletions
@@ -1,53 +0,0 @@
---
name: do-task
description: Use when handling a Gitea issue/task end-to-end, including validation, planning, implementation, PR creation, and review-feedback iteration.
---
# Do Task
## Overview
Process a single Gitea issue from analysis to merge-ready PR with explicit decision gates. Prefer Gitea MCP tools for all issue/PR actions.
## Workflow
1. **Load issue**
- Read issue by ID from Gitea.
- Extract: problem, expected behavior, impact, acceptance checks.
2. **Validate technical correctness**
- Verify against current code/docs/tests.
- Decide:
- **Valid issue** -> continue to planning.
- **Not valid / out of scope** -> comment rationale on issue and stop.
3. **If valid: plan and comment**
- Create implementation plan using Superpowers planning flow.
- Post summary plan to the issue before coding.
4. **Implement**
- Use Superpowers execution flow (TDD + verification before completion).
- Keep changes scoped strictly to issue requirements.
5. **Create PR**
- Open PR from feature branch.
- PR description must include:
- Summary of changes
- Test evidence (exact commands)
- `Fixes #<issue-id>` (or `Closes #<issue-id>`)
- Add reviewer: `codex`.
6. **Review feedback loop**
- Read all feedback.
- For each item:
- If technically valid -> implement + re-test + reply.
- If not valid -> reply with concise technical reasoning.
- Do not blindly accept external feedback without verification.
## Required Rules
- Use Gitea MCP tools first for issue/PR/review operations.
- Do not use `tea` or other CLI Gitea clients unless MCP is unavailable.
- Never claim completion without fresh verification output.
## Quick Command Pattern
- Invoke as: `/do-task <issue-id>`
- Example: `/do-task 12`
+3
View File
@@ -23,6 +23,8 @@ description: Phase 2 skill. Use when issue scope is locked by Product Owner (Pha
- Read issue: `mcp__gitea__.issue_read` (`method: "get"`).
- Read comments, find the **Final Requirements** comment: `mcp__gitea__.issue_read` (`method: "get_comments"`).
- If Final Requirements not found → stop, tell user to validate issue first.
- Read current Gitea user: `mcp__gitea__.get_me`.
- If current user is not already in the issue assignee list, update the issue with `mcp__gitea__.issue_write` (`method: "update"`) and set `assignees` to the existing assignees plus the current username.
2. **Read codebase**
- Use `mcp__gitea__.get_file_contents`, `mcp__gitea__.get_repository_tree` to understand current state.
@@ -72,6 +74,7 @@ description: Phase 2 skill. Use when issue scope is locked by Product Owner (Pha
## Rules
- Use Gitea MCP tools for all issue/PR operations.
- Preserve existing assignees when adding yourself unless the user explicitly asks to replace them.
- Never edit remote files directly — always work on local branch.
- Never claim completion without fresh verification output.
- Keep changes strictly scoped to the approved plan.
-133
View File
@@ -1,133 +0,0 @@
---
name: review-pr
description: Review pull requests in Gitea repositories using Gitea MCP for server-side PR metadata, comments, reviews, diffs, workflow state, and PR comments. Use local CLI tools such as git and test runners only for local repository inspection and validation. Use when the user asks to "Review PR", "review pull request", "check PR", "recheck PR after updates", or "post review comment". Inspect PR metadata and discussion, compare the PR head against its base branch, produce a clear PASS/FAIL verdict, and optionally post the result back to the PR.
---
# Review PR
## Local Preference
For repositories on `git.thuanle.me`, prefer the `mcp__gitea__` tools for Gitea-hosted pull request metadata, comments, reviews, server-side diffs, repository contents, branches, commits, workflow runs, and PR comments. Do not use `tea` or other command-line Gitea API clients unless the user explicitly asks for them or the MCP tools are unavailable for the required Gitea operation. Other local CLI tools such as `git`, `go`, `rg`, test runners, and formatters are OK for local repository inspection and validation.
## Workflow
Follow this sequence:
1. Identify the PR and current repo context.
- Use the current repository when the user is already inside it.
- If the user gives only a PR number and the repo context is ambiguous, ask a short clarifying question.
- Do not ask follow-up questions when the repository, PR, and requested reviewer action are already clear.
- If the user has already given a standing instruction earlier in the same thread such as "if you're confident, post it", carry that instruction forward for later re-reviews until the user changes it.
- Read the PR first with `mcp__gitea__.pull_request_read` using `method: "get"`, and read issue comments with `mcp__gitea__.issue_read` using `method: "get_comments"`.
- Read reviews with `mcp__gitea__.pull_request_read` using `method: "get_reviews"` when review state matters.
- Note the PR status exactly. If it is already merged or closed, say so explicitly before continuing.
2. Read the review history before judging the latest update.
- Find the latest blocking review comments first.
- If the author says they addressed feedback in a follow-up commit, focus on the new commits after that discussion.
- Still sanity-check the current full PR diff before concluding `PASS`.
3. Compare against the real base branch.
- Use PR metadata from MCP as the source of truth for base/head refs and SHAs.
- Prefer the PR base branch from MCP metadata; use `main` only when the base branch is not obvious from PR metadata.
4. Inspect the code changes.
- Use `mcp__gitea__.pull_request_read` with `method: "get_diff"` to read the PR diff.
- Use `mcp__gitea__.get_file_contents`, `mcp__gitea__.get_repository_tree`, and commit/branch MCP reads for surrounding context.
- If the user asks for a re-review after updates, optionally diff the latest fix commit range as a helper view, but do not skip the full PR sanity check.
5. Validate behavior where it matters.
- Prefer an isolated worktree when validation should run on the PR snapshot instead of the currently checked out branch.
- Use a repository-local worktree under `.worktrees/` by default, for example `git worktree add .worktrees/pr-<pr> origin/<head>`. Before creating it, run `git worktree list` and confirm the repo-local `.worktrees/` path.
- For `/Users/tm/working/thuanle/crypto/crypto-price-bot`, PR validation worktrees must be under `.worktrees/pr-<pr>` from the repository root. Do not use `/tmp` for this repo; if `.worktrees/` is unavailable, stop and ask the user before using any fallback path.
- Run validation commands inside that worktree.
- Remove it afterward with `git worktree remove .worktrees/pr-<pr>` when it is no longer needed.
- For other repositories, if `git worktree add` fails because the environment blocks writes to `.git/worktrees`, fall back to exporting a temporary snapshot for validation and say that you used the fallback.
- Run repo-appropriate checks when feasible.
- Prefer targeted checks first, then broader validation such as `go test ./...` and `go vet ./...` in Go repositories.
- Say explicitly when validation could not be run.
6. Write the verdict for the user.
- Start with `PASS` or `FAIL`.
- For `FAIL`, list blocking findings first, ordered by severity.
- Include file and line references whenever possible.
- Explain the concrete behavior impact, not style preferences.
- For `PASS`, say that no blocking issues were found and then add any non-blocking notes.
7. Post a review result proactively by default.
- When acting as a reviewer, prefer an official PR review over a plain issue comment if the result should affect PR state.
- Default stance: if the verdict is clear and you are confident in it, post the review without waiting for another prompt from the user.
- If the user has already asked you to post the result, or has given a standing instruction such as "if you're confident, post it", treat that as continuing authority for later re-reviews in the same thread unless the user revokes or narrows it.
- Ask the user only when there is a real ambiguity you cannot safely resolve, such as unclear repo/PR context, unclear whether they explicitly want a local-only verdict instead of a posted review, or a materially uncertain finding.
- If the user explicitly asks for a local verdict only, do not post.
- Use `mcp__gitea__.pull_request_review_write` for review posting.
- For normal review results, prefer a single `create` call with the final `state` and `body`. This is the default path for `PASS`, `FAIL`, and neutral review comments.
- Use `submit` only when you intentionally need a multi-step pending review flow, such as staging a pending review first and finalizing it later.
- Never send the same review narrative through both `create` and `submit`; that is a process bug and can surface as duplicate review content in the UI.
- Map the verdict to review state:
- `PASS` -> `APPROVED`
- `FAIL` -> `REQUEST_CHANGES`
- neutral/non-blocking note only -> `COMMENT`
- Use `mcp__gitea__.issue_write` with `method: "add_comment"` only for plain discussion comments that should not change PR review status.
- Do not add bracketed author/tool tags such as `[codex]`.
- Start directly with the review result or author response, depending on your role in the thread.
- Keep the review/comment concise and direct.
- Mention the validation commands you actually ran when posting review results.
## Command Patterns
Use these MCP operations directly:
- `mcp__gitea__.pull_request_read` with `method: "get"`
- `mcp__gitea__.pull_request_read` with `method: "get_diff"`
- `mcp__gitea__.pull_request_read` with `method: "get_reviews"`
- `mcp__gitea__.pull_request_read` with `method: "get_review_comments"` when a specific review thread is needed
- `mcp__gitea__.issue_read` with `method: "get_comments"`
- `mcp__gitea__.get_file_contents`
- `mcp__gitea__.get_repository_tree`
- `mcp__gitea__.get_commit`
- `mcp__gitea__.actions_run_read` for workflow status and logs
- `mcp__gitea__.pull_request_review_write` with `method: "create"` for the normal one-shot final review path
- Set the final `state` on `create`
- Put the final review text on `create`
- `mcp__gitea__.pull_request_review_write` with `method: "submit"` only when finishing an intentionally pending review
- `mcp__gitea__.issue_write` with `method: "add_comment"` when posting a plain PR comment
Do not use `tea` for Gitea PR review unless the user explicitly requests it or MCP cannot provide the needed data. Local CLI usage such as `git diff`, `git worktree`, and test commands remains acceptable when it materially improves local validation.
## Comment Template
Use this structure when posting a review result:
```text
PASS
No blocking issues found in the latest update.
Validation:
- <command 1>
- <command 2>
Non-blocking note:
- <optional note>
```
```text
FAIL
1. [High] <blocking issue with path:line and impact>
2. [Medium] <blocking issue with path:line and impact>
Validation:
- <command 1>
- <command 2>
```
When acting as a reviewer and the user asks to approve/reject or otherwise post the result, also submit the matching official review state:
- `PASS` -> `APPROVED`
- `FAIL` -> `REQUEST_CHANGES`
- neutral note only -> `COMMENT`
When responding as the PR author rather than reviewer, do not use `PASS`/`FAIL`; write a concise author response that explains what feedback was addressed and which validation ran.
Do not add filler. Keep the review specific, technical, and actionable.
+8 -2
View File
@@ -22,13 +22,17 @@ description: Phase 2 skill. Use when a PR is created or updated. Review PR again
1. **Load PR and issue**
- Read PR: `mcp__gitea__.pull_request_read` (`method: "get"`).
- If PR is already merged or closed, say so and stop.
- Read current Gitea user: `mcp__gitea__.get_me`.
- Read existing reviews: `mcp__gitea__.pull_request_read` (`method: "get_reviews"`).
- If the current user is not already requested as a reviewer and has not already reviewed this PR, add them with `mcp__gitea__.pull_request_write` (`method: "add_reviewers"`).
- Identify the linked issue from PR body (`Fixes #<id>` or `Closes #<id>`).
- Read issue comments, find the **Final Requirements** comment: `mcp__gitea__.issue_read` (`method: "get_comments"`).
- Read existing reviews: `mcp__gitea__.pull_request_read` (`method: "get_reviews"`).
- If PR is already merged or closed, say so and stop.
- Read PR discussion comments: `mcp__gitea__.issue_read` (`method: "get_comments"`), using the PR number as the index.
2. **Read review history** (for re-reviews)
- Get review comments: `mcp__gitea__.pull_request_read` (`method: "get_review_comments"`).
- Treat reviews, inline review comments, and PR discussion comments as separate context sources.
- If author addressed feedback in a follow-up commit, focus on new commits but still sanity-check the full diff.
3. **Get the diff**
@@ -99,9 +103,11 @@ Validation:
## Rules
- Use Gitea MCP tools for all PR/review operations.
- Do not add yourself as reviewer twice; check requested reviewers and existing reviews first.
- Do NOT edit code on the PR.
- FAIL if PR contains unrelated changes or scope creep.
- Post review proactively if verdict is clear.
- Do not post a verdict until PR discussion comments have been checked.
- Mention validation commands you actually ran.
- No bracketed author/tool tags.
- Accept `docs/adr/**` when it records a direct, durable decision tied to the issue scope.
+1
View File
@@ -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=
+5 -3
View File
@@ -3,16 +3,19 @@
## 1. Coding Principles
### 1.1 Simplicity First
- No out-of-scope features.
- No abstractions for single-use code.
- If 50 lines work instead of 200, write 50.
### 1.2 Surgical Changes
- Only change lines directly related to the request.
- Never refactor adjacent code unless asked.
- Remove orphaned imports/variables/functions your changes created.
### 1.3 Goal-Driven Execution
- Define success criteria before coding.
- Verify each step before moving to the next.
@@ -70,7 +73,7 @@ All work starts from a Gitea issue. Two phases, each with a dedicated role:
**Constraint:** Reviewer must NOT edit code on the PR.
1. Re-read the original issue — verify PR addresses the locked requirement.
2. Read PR diff and discussion history.
2. Read PR diff, review history, and PR discussion comments.
3. Validate in worktree if needed (`.worktrees/pr-<number>`, run `go test`, `go vet`, `gofmt`).
4. Check for scope creep — FAIL if PR contains unrelated changes or refactoring.
5. Write verdict — start with **PASS** or **FAIL**.
@@ -83,10 +86,9 @@ All work starts from a Gitea issue. Two phases, each with a dedicated role:
## 5. Gitea Tools
- If git remote contains `git.thuanle.me`, ALWAYS use Gitea MCP tools.
- Alway use Gitea MCP tools in this project.
- Scope: read PR/issue, list comments, post replies, create/edit PRs, reviews.
- Use `tea` CLI only as fallback when MCP is unavailable.
- Do not use `gh` for Gitea repositories.
- Local CLI tools (`git`, `go`, `rg`) are fine for local validation.
---
+2 -1
View File
@@ -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
+30
View File
@@ -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).
+1
View File
@@ -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 |
@@ -0,0 +1,76 @@
# Design: Issue #18 — Add `go vet` and `staticcheck` to CI
## Context
Issue #18 requires expanding PR CI validation for this Go repository beyond existing `gofmt` and `go test` checks.
Current state:
- CI is defined in `.gitea/workflows/ci.yml`.
- Existing jobs: `fmt` (`gofmt -l .`) and `test` (`go test ./...`).
- No `go vet` or `staticcheck` step exists.
## Goal
Add practical baseline static analysis to pull request CI by:
- Running `go vet ./...`.
- Running `staticcheck ./...`.
- Pinning `staticcheck` version in CI for deterministic behavior.
## Non-goals
- Adding `golangci-lint`.
- Refactoring application code unrelated to CI wiring.
- Changing workflow triggers.
## Approved approach
Use separate CI jobs for `vet` and `staticcheck` while keeping existing `fmt` and `test` jobs unchanged.
### Why this approach
- Keeps changes surgical and easy to review.
- Makes failures explicit per check (better debugging signal than a combined lint step).
- Preserves current CI behavior while extending validation coverage.
## Implementation design
Modify `.gitea/workflows/ci.yml` only.
### Job: `vet`
- `runs-on: linux`
- Steps:
1. `actions/checkout@v4`
2. `actions/setup-go@v5` with `go-version-file: go.mod`
3. `run: go vet ./...`
### Job: `staticcheck`
- `runs-on: linux`
- Steps:
1. `actions/checkout@v4`
2. `actions/setup-go@v5` with `go-version-file: go.mod`
3. Install pinned staticcheck version:
- `go install honnef.co/go/tools/cmd/staticcheck@2025.1.1`
4. Run staticcheck:
- `$(go env GOPATH)/bin/staticcheck ./...`
## Failure behavior
- CI must fail if `go vet ./...` exits non-zero.
- CI must fail if `staticcheck ./...` exits non-zero.
- Existing `fmt` and `test` failure behavior remains unchanged.
## Verification plan
Local verification before PR creation:
- `gofmt -l .`
- `go test ./...`
- `go vet ./...`
- `staticcheck ./...` (if tool is available locally)
CI verification through PR run:
- Confirm `fmt`, `test`, `vet`, and `staticcheck` jobs execute on pull requests.
- Confirm non-zero from `vet`/`staticcheck` marks workflow failed.
## Risks and mitigations
- Risk: pinned staticcheck version not present in runner cache.
- Mitigation: install in-job via `go install` each run.
- Risk: staticcheck findings on current code block CI.
- Mitigation: if encountered, address only findings required to satisfy issue scope.
## Acceptance mapping
- CI fail on vet errors → `vet` job runs `go vet ./...`.
- CI fail on staticcheck errors → `staticcheck` job runs pinned staticcheck binary.
- Existing checks preserved → `fmt` and `test` jobs remain in workflow.
- No `golangci-lint` → not introduced in workflow.
+23 -1
View File
@@ -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,25 @@ 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 {
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 {
if k != "task_uuid" {
metadata[k] = v
}
}
}
}
return metadata
}
+47
View File
@@ -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)
+9
View File
@@ -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"
}
+1
View File
@@ -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"`
+147 -51
View File
@@ -1,12 +1,16 @@
package mail
import (
"bufio"
"context"
"crypto/tls"
"encoding/base64"
"encoding/json"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"net/url"
"strings"
"time"
@@ -21,6 +25,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 +37,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)
@@ -55,13 +62,26 @@ func (d timeoutContextDialer) DialContext(ctx context.Context, network, address
var socks5DialerFactory = proxy.SOCKS5
var proxyTLSConfigForTest = func(host string) *tls.Config {
return &tls.Config{ServerName: host}
}
// 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,
dialIMAPViaProxy: dialTLSViaProxy,
rules: rules.NewPipeline(externalRules),
}
}
@@ -180,34 +200,20 @@ func (w *IMAPWatcher) connectAndWatch(ctx context.Context) error {
}
}
func dialTLSViaSOCKS5(addr, proxyURL string, options *imapclient.Options) (*imapclient.Client, error) {
func dialTLSViaProxy(addr, proxyURL string, options *imapclient.Options) (*imapclient.Client, error) {
u, err := url.Parse(proxyURL)
if err != nil {
return nil, fmt.Errorf("parse proxy url: %w", err)
}
if u.Scheme != "socks5" {
return nil, fmt.Errorf("unsupported proxy scheme: %s", u.Scheme)
}
var auth *proxy.Auth
if u.User != nil {
pw, _ := u.User.Password()
auth = &proxy.Auth{User: u.User.Username(), Password: pw}
}
dialer, err := socks5DialerFactory("tcp", u.Host, auth, timeoutContextDialer{timeout: imapDialTimeout})
if err != nil {
return nil, fmt.Errorf("create socks5 dialer: %w", err)
}
dialCtx, dialCancel := context.WithTimeout(context.Background(), imapDialTimeout)
defer dialCancel()
var conn net.Conn
if cd, ok := dialer.(proxy.ContextDialer); ok {
conn, err = cd.DialContext(dialCtx, "tcp", addr)
} else {
conn, err = dialer.Dial("tcp", addr)
switch u.Scheme {
case "socks5":
conn, err = dialViaSOCKS5(u, addr)
case "http", "https":
conn, err = dialViaCONNECT(u, addr)
default:
return nil, fmt.Errorf("unsupported proxy scheme: %s", u.Scheme)
}
if err != nil {
return nil, fmt.Errorf("proxy dial: %w", err)
@@ -230,6 +236,96 @@ func dialTLSViaSOCKS5(addr, proxyURL string, options *imapclient.Options) (*imap
return imapclient.New(tlsConn, options), nil
}
func dialViaSOCKS5(u *url.URL, targetAddr string) (net.Conn, error) {
var auth *proxy.Auth
if u.User != nil {
pw, _ := u.User.Password()
auth = &proxy.Auth{User: u.User.Username(), Password: pw}
}
dialer, err := socks5DialerFactory("tcp", u.Host, auth, timeoutContextDialer{timeout: imapDialTimeout})
if err != nil {
return nil, fmt.Errorf("create socks5 dialer: %w", err)
}
dialCtx, dialCancel := context.WithTimeout(context.Background(), imapDialTimeout)
defer dialCancel()
if cd, ok := dialer.(proxy.ContextDialer); ok {
return cd.DialContext(dialCtx, "tcp", targetAddr)
}
return dialer.Dial("tcp", targetAddr)
}
// proxyConn wraps a net.Conn, draining buffered data from a bufio.Reader
// before delegating reads to the underlying connection.
type proxyConn struct {
net.Conn
reader io.Reader
}
func (c *proxyConn) Read(b []byte) (int, error) {
return c.reader.Read(b)
}
// dialToProxy is the function used to establish a TCP connection to the proxy.
var dialToProxy = func(ctx context.Context, network, addr string) (net.Conn, error) {
return (&net.Dialer{Timeout: imapDialTimeout}).DialContext(ctx, network, addr)
}
func dialViaCONNECT(u *url.URL, targetAddr string) (net.Conn, error) {
ctx, cancel := context.WithTimeout(context.Background(), imapDialTimeout)
defer cancel()
conn, err := dialToProxy(ctx, "tcp", u.Host)
if err != nil {
return nil, fmt.Errorf("dial proxy %s: %w", u.Host, err)
}
if u.Scheme == "https" {
tlsConn := tls.Client(conn, proxyTLSConfigForTest(u.Hostname()))
if err := tlsConn.HandshakeContext(ctx); err != nil {
_ = tlsConn.Close()
return nil, fmt.Errorf("tls handshake to proxy: %w", err)
}
conn = tlsConn
}
_ = conn.SetDeadline(time.Now().Add(imapDialTimeout))
connectReq := fmt.Sprintf("CONNECT %s HTTP/1.1\r\nHost: %s\r\n", targetAddr, targetAddr)
if u.User != nil {
username := u.User.Username()
password, _ := u.User.Password()
creds := base64.StdEncoding.EncodeToString([]byte(username + ":" + password))
connectReq += fmt.Sprintf("Proxy-Authorization: Basic %s\r\n", creds)
}
connectReq += "\r\n"
if _, err := fmt.Fprint(conn, connectReq); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("send connect: %w", err)
}
br := bufio.NewReader(conn)
resp, err := http.ReadResponse(br, nil)
if err != nil {
_ = conn.Close()
return nil, fmt.Errorf("read connect response: %w", err)
}
if resp.StatusCode != http.StatusOK {
_ = conn.Close()
return nil, fmt.Errorf("proxy connect failed: %s", resp.Status)
}
_ = conn.SetDeadline(time.Time{})
if br.Buffered() > 0 {
return &proxyConn{Conn: conn, reader: io.MultiReader(br, conn)}, nil
}
return conn, nil
}
// fetchUnseen searches for UNSEEN messages and processes each one.
func (w *IMAPWatcher) fetchUnseen(c *imapclient.Client) error {
criteria := &imap.SearchCriteria{
@@ -343,15 +439,7 @@ 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")
w.markSeen(c, buf.UID)
return
}
// 3. Idempotency: skip if message_id already exists.
// 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)
@@ -363,6 +451,19 @@ func (w *IMAPWatcher) processMessage(c *imapclient.Client, buf *imapclient.Fetch
return
}
// 3. 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
}
// 4. Threading: determine thread_id.
threadID := messageID
inReplyTo := ""
@@ -404,13 +505,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 +555,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.
+213 -6
View File
@@ -1,10 +1,20 @@
package mail
import (
"bufio"
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/tls"
"crypto/x509"
"encoding/base64"
"errors"
"fmt"
"io"
"math/big"
"net"
"net/http"
"strings"
"testing"
"time"
@@ -73,13 +83,22 @@ func TestConnectAndWatch_UsesProxyDial_WhenProxySet(t *testing.T) {
}
}
func TestDialTLSViaSOCKS5_RejectsNonSocks5Scheme(t *testing.T) {
_, err := dialTLSViaSOCKS5("imap.example.com:993", "http://proxy:8080", nil)
func TestDialTLSViaProxy_RejectsUnsupportedScheme(t *testing.T) {
_, err := dialTLSViaProxy("imap.example.com:993", "ftp://proxy:21", nil)
if err == nil || !strings.Contains(err.Error(), "unsupported proxy scheme") {
t.Fatalf("expected unsupported scheme error, got %v", err)
}
}
func TestDialTLSViaProxy_RejectsMalformedURL(t *testing.T) {
_, err := dialTLSViaProxy("imap.example.com:993", "://bad", nil)
if err == nil {
t.Fatal("expected error for malformed URL")
}
}
// --- SOCKS5 path ---
type fakeContextDialer struct{}
func (fakeContextDialer) Dial(network, address string) (net.Conn, error) {
@@ -93,7 +112,7 @@ func (fakeContextDialer) DialContext(ctx context.Context, network, address strin
return nil, errors.New("used dialcontext")
}
func TestDialTLSViaSOCKS5_UsesDialContextWithTimeout(t *testing.T) {
func TestDialTLSViaProxy_SOCKS5_UsesDialContextWithTimeout(t *testing.T) {
orig := socks5DialerFactory
t.Cleanup(func() { socks5DialerFactory = orig })
@@ -104,7 +123,7 @@ func TestDialTLSViaSOCKS5_UsesDialContextWithTimeout(t *testing.T) {
return fakeContextDialer{}, nil
}
_, err := dialTLSViaSOCKS5("imap.example.com:993", "socks5://127.0.0.1:1080", nil)
_, err := dialTLSViaProxy("imap.example.com:993", "socks5://127.0.0.1:1080", nil)
if err == nil || !strings.Contains(err.Error(), "proxy dial: used dialcontext") {
t.Fatalf("expected DialContext path, got %v", err)
}
@@ -143,7 +162,7 @@ func (d contextConnDialer) DialContext(ctx context.Context, network, address str
return d.conn, nil
}
func TestDialTLSViaSOCKS5_SetsDeadlineForTLSHandshake(t *testing.T) {
func TestDialTLSViaProxy_SOCKS5_SetsDeadlineForTLSHandshake(t *testing.T) {
orig := socks5DialerFactory
t.Cleanup(func() { socks5DialerFactory = orig })
@@ -152,7 +171,7 @@ func TestDialTLSViaSOCKS5_SetsDeadlineForTLSHandshake(t *testing.T) {
return contextConnDialer{conn: conn}, nil
}
_, err := dialTLSViaSOCKS5("imap.example.com:993", "socks5://127.0.0.1:1080", nil)
_, err := dialTLSViaProxy("imap.example.com:993", "socks5://127.0.0.1:1080", nil)
if err == nil || !strings.Contains(err.Error(), "tls handshake") {
t.Fatalf("expected tls handshake error, got %v", err)
}
@@ -160,3 +179,191 @@ func TestDialTLSViaSOCKS5_SetsDeadlineForTLSHandshake(t *testing.T) {
t.Fatal("expected TLS handshake deadline to be set")
}
}
// --- HTTP CONNECT path ---
func startFakeProxy(t *testing.T, handler func(host string) bool) (addr string, cleanup func()) {
t.Helper()
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
done := make(chan struct{})
go func() {
defer close(done)
for {
conn, err := ln.Accept()
if err != nil {
return
}
go func(c net.Conn) {
defer c.Close()
br := bufio.NewReader(c)
req, err := http.ReadRequest(br)
if err != nil {
return
}
if handler(req.URL.Host) {
fmt.Fprintf(c, "HTTP/1.1 200 OK\r\n\r\n")
} else {
fmt.Fprintf(c, "HTTP/1.1 403 Forbidden\r\n\r\n")
}
}(conn)
}
}()
return ln.Addr().String(), func() {
ln.Close()
<-done
}
}
func TestDialTLSViaProxy_HTTPConnect_ReachesTLSHandshake(t *testing.T) {
addr, cleanup := startFakeProxy(t, func(host string) bool { return true })
defer cleanup()
_, err := dialTLSViaProxy("imap.example.com:993", fmt.Sprintf("http://%s", addr), nil)
if err == nil || !strings.Contains(err.Error(), "tls handshake") {
t.Fatalf("expected tls handshake error (CONNECT succeeded), got %v", err)
}
}
func TestDialTLSViaProxy_HTTPConnect_RejectsOnProxyFailure(t *testing.T) {
addr, cleanup := startFakeProxy(t, func(host string) bool { return false })
defer cleanup()
_, err := dialTLSViaProxy("imap.example.com:993", fmt.Sprintf("http://%s", addr), nil)
if err == nil || !strings.Contains(err.Error(), "proxy connect failed") {
t.Fatalf("expected proxy connect failed error, got %v", err)
}
}
func TestDialTLSViaProxy_HTTPConnect_SendsProxyAuth(t *testing.T) {
var gotAuth string
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
go func() {
defer ln.Close()
conn, err := ln.Accept()
if err != nil {
return
}
defer conn.Close()
br := bufio.NewReader(conn)
req, err := http.ReadRequest(br)
if err != nil {
return
}
gotAuth = req.Header.Get("Proxy-Authorization")
fmt.Fprintf(conn, "HTTP/1.1 200 OK\r\n\r\n")
}()
_, _ = dialTLSViaProxy("imap.example.com:993",
fmt.Sprintf("http://testuser:testpass@%s", ln.Addr().String()), nil)
expected := "Basic " + base64.StdEncoding.EncodeToString([]byte("testuser:testpass"))
if gotAuth != expected {
t.Fatalf("expected auth %q, got %q", expected, gotAuth)
}
}
func TestDialTLSViaProxy_HTTPConnect_DialError(t *testing.T) {
origDial := dialToProxy
t.Cleanup(func() { dialToProxy = origDial })
dialToProxy = func(ctx context.Context, network, addr string) (net.Conn, error) {
return nil, errors.New("dial refused")
}
_, err := dialTLSViaProxy("imap.example.com:993", "http://127.0.0.1:1", nil)
if err == nil || !strings.Contains(err.Error(), "dial proxy") {
t.Fatalf("expected dial proxy error, got %v", err)
}
}
// --- HTTPS CONNECT path ---
func generateTestCert(t *testing.T) tls.Certificate {
t.Helper()
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
t.Fatalf("generate key: %v", err)
}
tmpl := &x509.Certificate{
SerialNumber: big.NewInt(1),
NotBefore: time.Now(),
NotAfter: time.Now().Add(time.Hour),
KeyUsage: x509.KeyUsageDigitalSignature,
DNSNames: []string{"127.0.0.1"},
IPAddresses: []net.IP{net.ParseIP("127.0.0.1")},
}
certDER, err := x509.CreateCertificate(rand.Reader, tmpl, tmpl, &key.PublicKey, key)
if err != nil {
t.Fatalf("create cert: %v", err)
}
cert, err := x509.ParseCertificate(certDER)
if err != nil {
t.Fatalf("parse cert: %v", err)
}
return tls.Certificate{
Certificate: [][]byte{certDER},
PrivateKey: key,
Leaf: cert,
}
}
func TestDialTLSViaProxy_HTTPSConnect_ReachesTLSHandshake(t *testing.T) {
cert := generateTestCert(t)
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
tlsLn := tls.NewListener(ln, &tls.Config{
Certificates: []tls.Certificate{cert},
})
done := make(chan struct{})
go func() {
defer close(done)
for {
conn, err := tlsLn.Accept()
if err != nil {
return
}
go func(c net.Conn) {
defer c.Close()
br := bufio.NewReader(c)
req, err := http.ReadRequest(br)
if err != nil {
return
}
if req.URL.Host != "" {
fmt.Fprintf(c, "HTTP/1.1 200 OK\r\n\r\n")
}
}(conn)
}
}()
defer func() {
tlsLn.Close()
<-done
}()
// Override dialToProxy to return raw TCP, but inject InsecureSkipVerify
// so the TLS handshake to proxy succeeds with the test cert.
origDial := dialToProxy
t.Cleanup(func() { dialToProxy = origDial })
origProxyTLS := proxyTLSConfigForTest
t.Cleanup(func() { proxyTLSConfigForTest = origProxyTLS })
proxyTLSConfigForTest = func(host string) *tls.Config {
return &tls.Config{ServerName: host, InsecureSkipVerify: true}
}
_, err = dialTLSViaProxy("imap.example.com:993",
fmt.Sprintf("https://%s", ln.Addr().String()), nil)
if err == nil || !strings.Contains(err.Error(), "tls handshake") {
t.Fatalf("expected tls handshake error (HTTPS CONNECT succeeded), got %v", err)
}
}
+26
View File
@@ -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()
}
+35
View File
@@ -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")
}
}
+30
View File
@@ -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
}
+106
View File
@@ -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() }
+41
View File
@@ -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
}
+42
View File
@@ -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")
}
}
+30
View File
@@ -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")
}
+11 -5
View File
@@ -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
@@ -68,12 +69,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 +85,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 +148,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)
@@ -166,7 +169,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
@@ -175,4 +178,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.