Add mavmaild, the read-only IMAP poller that feeds mail intake (#246)
The extraction seam landed on the previous branch but nothing fed it. This adds the daemon that does: every interval it opens one mailbox read-only (EXAMINE + BODY.PEEK, so reading leaves no \Seen behind), fetches the UIDs it has not handed over yet, and posts each message to core over ingest_mail. Core runs the model and writes task candidates; this daemon writes nothing and cannot create a reminder. It is a separate daemon because of the credential. mavpoll set the precedent with the zenmoney token (#125): the module talking to the third party holds the secret, reads it from a file so it never lands in argv, in docker-compose.yml or in shell history, and core never sees it. There is deliberately no -password flag, and a test asserts that. Off unless configured at both ends: without -password-file the daemon refuses to start, and if core has no email block the first ingest returns ErrUnknownMethod, which disables the reader instead of hammering a socket that will keep refusing. A seen-UID state file (0600, atomic write) keeps a restart from re-extracting the whole lookback window; correctness does not depend on it, since capture dedupes on normalised text. Logs are counts and UIDs — no subject, sender or body. Verified with an in-process IMAP server and a fake core: bulk mail is filtered before core is asked, seen UIDs are not re-fetched, a failed ingest is retried next poll, ErrUnknownMethod stops at the first message, and state survives a restart. The live half is untested by design — no IMAP credential exists on this box; setup is written up as QA steps. Vikunja #246
This commit is contained in:
@@ -7,6 +7,7 @@
|
||||
/mavpoll
|
||||
/mavcaldav
|
||||
/mavwaked
|
||||
/mavmaild
|
||||
|
||||
# Certs (private keys, don't commit)
|
||||
certs/
|
||||
@@ -36,6 +37,8 @@ deploy/db_key.env
|
||||
deploy/telegram.env
|
||||
# zenmoney API token, read by mavpoll (never in argv, never committed)
|
||||
deploy/zenmoney.token
|
||||
# IMAP password, read by mavmaild (never in argv, never committed)
|
||||
deploy/imap.password
|
||||
|
||||
# Temp files
|
||||
/tmp/
|
||||
|
||||
@@ -34,7 +34,7 @@ CGO daemons (`mavend`, `mavsttd`, `mavttsd`, `mavenclient`) need the vendored to
|
||||
and libs wired through the Makefile — **do not** call `go build` on them bare, use `make`:
|
||||
|
||||
```sh
|
||||
make build # all 8 binaries
|
||||
make build # all 9 binaries
|
||||
make build-web # single daemon (pure-Go ones: web/waked/poll/caldav build without CGO)
|
||||
make test # go test -race across ./internal/... ./cmd/... with CGO env set
|
||||
```
|
||||
@@ -62,6 +62,7 @@ Pure-Go packages (`router`, `memory`, `mavweb`, …) run under a plain `go test
|
||||
| `mavenclient` | Voice loop client (mic → stt → core → tts). |
|
||||
| `mavpoll` | Telegram long-poll reach. |
|
||||
| `mavcaldav` | CalDAV calendar sync. |
|
||||
| `mavmaild` | Mail reader (IMAP, read-only). Holds the IMAP password; core never sees it. |
|
||||
|
||||
Daemons are wired socket-to-socket, not linked. `internal/ipc` is the client/server wire
|
||||
protocol; the config in `deploy/mavend.json` (with `${VAR}` env expansion from gitignored
|
||||
|
||||
+2
-1
@@ -51,7 +51,8 @@ RUN go build -o /out/mavend ./cmd/mavend && \
|
||||
go build -o /out/mavttsd ./cmd/mavttsd && \
|
||||
go build -o /out/mavweb ./cmd/mavweb && \
|
||||
go build -o /out/mavpoll ./cmd/mavpoll && \
|
||||
go build -o /out/mavcaldav ./cmd/mavcaldav
|
||||
go build -o /out/mavcaldav ./cmd/mavcaldav && \
|
||||
go build -o /out/mavmaild ./cmd/mavmaild
|
||||
|
||||
# llama.cpp Vulkan build — the phraser/router LFM engine (llama-server). Built
|
||||
# from source (not a prebuilt vendored blob) so the binary's glibc/GLIBCXX match
|
||||
|
||||
@@ -20,7 +20,7 @@ PIPER_ESPEAK := $(shell pwd)/deps/piper/espeak-ng-data
|
||||
|
||||
all: build
|
||||
|
||||
build: build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav
|
||||
build: build-stt build-tts build-daemon build-client build-waked build-web build-poll build-caldav build-mail
|
||||
|
||||
build-stt:
|
||||
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
|
||||
@@ -50,6 +50,9 @@ build-poll:
|
||||
build-caldav:
|
||||
$(GO) build $(GOFLAGS) -o mavcaldav ./cmd/mavcaldav/
|
||||
|
||||
build-mail:
|
||||
$(GO) build $(GOFLAGS) -o mavmaild ./cmd/mavmaild/
|
||||
|
||||
run-web: build-web
|
||||
./mavweb -addr :9200 -voice 127.0.0.1:9100
|
||||
|
||||
@@ -185,4 +188,4 @@ download-embedder:
|
||||
@echo ' sudo cp onnxruntime-linux-x64-1.15.1/lib/libonnxruntime.so* /usr/local/lib/'
|
||||
|
||||
clean:
|
||||
rm -f mavend mavenclient mavsttd mavttsd mavweb mavpoll mavcaldav mavwaked
|
||||
rm -f mavend mavenclient mavsttd mavttsd mavweb mavpoll mavcaldav mavwaked mavmaild
|
||||
|
||||
@@ -0,0 +1,332 @@
|
||||
// mavmaild — the mail reader module (Vikunja #246,
|
||||
// docs/plans/01-email-reader.md).
|
||||
//
|
||||
// Every so often it opens one IMAP mailbox read-only, fetches the messages it
|
||||
// has not read yet, and hands each one to core over ipc.MethodIngestMail. Core
|
||||
// runs the extraction on the resident model and writes what comes back as task
|
||||
// CANDIDATES he reviews on /tasks. Nothing here writes to the store, nothing
|
||||
// here can create a reminder, and nothing here speaks.
|
||||
//
|
||||
// Why a separate daemon rather than a loop inside mavend, when extraction has
|
||||
// to happen in mavend anyway: the credential. mavpoll set the precedent with the
|
||||
// zenmoney token (#125) — the module that talks to a third party holds the
|
||||
// secret, reads it from a FILE so it never appears in `ps`, in
|
||||
// docker-compose.yml or in shell history, and core never sees it. Core learns
|
||||
// that mail exists only as message text on one IPC method; it cannot connect to
|
||||
// the mailbox even if it wanted to, and a compromised core yields no mail
|
||||
// password.
|
||||
//
|
||||
// Off unless configured: without -password-file there is nothing to run, and
|
||||
// the daemon says so and exits. If core has no `email` block the very first
|
||||
// ingest comes back ErrUnknownMethod and this daemon stops polling instead of
|
||||
// hammering a socket that will keep refusing.
|
||||
//
|
||||
// Mail is personal, so the log is counts and UIDs: how many messages were
|
||||
// fetched, how many were bulk, how many candidates came back. No subject, no
|
||||
// sender, no body, ever — reviewing a candidate is what /tasks is for.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/email"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
func main() {
|
||||
if err := run(os.Args[1:]); err != nil {
|
||||
fmt.Fprintln(os.Stderr, "mavmaild:", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
func run(args []string) error {
|
||||
fs := flag.NewFlagSet("mavmaild", flag.ContinueOnError)
|
||||
socket := fs.String("socket", "", "core IPC socket path (required)")
|
||||
server := fs.String("imap", "", "IMAP server, host or host:993 (required)")
|
||||
user := fs.String("user", "", "IMAP username (required)")
|
||||
passFile := fs.String("password-file", "", "file holding the IMAP password (required — never passed as a flag value)")
|
||||
mailbox := fs.String("mailbox", "INBOX", "mailbox to read, read-only")
|
||||
interval := fs.Duration("interval", 15*time.Minute, "how often to read the mailbox")
|
||||
lookback := fs.Duration("lookback", 72*time.Hour, "how far back to search on each poll")
|
||||
max := fs.Int("max", 25, "most messages to fetch in one poll")
|
||||
timeout := fs.Duration("timeout", 30*time.Second, "IMAP network timeout")
|
||||
statePath := fs.String("state", "", "file remembering which UIDs were read (default: none — every poll re-reads the window)")
|
||||
if err := fs.Parse(args); err != nil {
|
||||
return err
|
||||
}
|
||||
if *socket == "" {
|
||||
return fmt.Errorf("-socket is required")
|
||||
}
|
||||
if *server == "" || *user == "" || *passFile == "" {
|
||||
return fmt.Errorf("mail reading is off unless configured: set -imap, -user and -password-file")
|
||||
}
|
||||
|
||||
// The password is read from a file, never taken as a flag value: an argv
|
||||
// secret is visible in `ps` to every user on the box and lands in the compose
|
||||
// file and the shell history. Read once at start — a rotated password means a
|
||||
// restart, which is cheaper than re-reading his credential every quarter hour.
|
||||
raw, err := os.ReadFile(*passFile)
|
||||
if err != nil {
|
||||
return fmt.Errorf("read password file: %w", err)
|
||||
}
|
||||
password := strings.TrimSpace(string(raw))
|
||||
if password == "" {
|
||||
return fmt.Errorf("password file %s is empty", *passFile)
|
||||
}
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
defer stop()
|
||||
|
||||
core, err := ipc.DialWait(*socket, 60*time.Second)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer core.Close()
|
||||
|
||||
r := &reader{
|
||||
core: core,
|
||||
addr: *server,
|
||||
user: *user,
|
||||
mailbox: *mailbox,
|
||||
lookback: *lookback,
|
||||
max: *max,
|
||||
timeout: *timeout,
|
||||
state: newSeenState(*statePath),
|
||||
}
|
||||
if err := r.state.load(); err != nil {
|
||||
// A missing or corrupt state file must not stop mail from being read: the
|
||||
// worst case is re-reading the window, and capture dedupes on text.
|
||||
log.Printf("mavmaild: state: %v (starting from an empty seen-set)", err)
|
||||
}
|
||||
|
||||
// The password is never logged, not even its length.
|
||||
log.Printf("mavmaild: reading %s on %s every %s (lookback %s, max %d/poll)",
|
||||
*mailbox, *server, *interval, *lookback, *max)
|
||||
|
||||
r.pollOnce(ctx, password) // don't idle a full interval on start
|
||||
t := time.NewTicker(*interval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Printf("mavmaild: bye")
|
||||
return nil
|
||||
case <-t.C:
|
||||
if r.disabled {
|
||||
// Core told us mail ingestion is not configured. Nothing will change
|
||||
// without a core restart, and a restart restarts us too.
|
||||
log.Printf("mavmaild: core does not accept mail — idling")
|
||||
return nil
|
||||
}
|
||||
r.pollOnce(ctx, password)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// mailIngester — the slice of core this daemon uses. One method: hand over a
|
||||
// message. It cannot write a fact, create a reminder or read the store, and the
|
||||
// interface says so.
|
||||
type mailIngester interface {
|
||||
IngestMail(ctx context.Context, req ipc.IngestMailReq) (ipc.IngestMailResp, error)
|
||||
}
|
||||
|
||||
type reader struct {
|
||||
core mailIngester
|
||||
addr string
|
||||
user string
|
||||
mailbox string
|
||||
lookback time.Duration
|
||||
max int
|
||||
timeout time.Duration
|
||||
state *seenState
|
||||
|
||||
// dial — connection seam for the tests; nil ⇒ implicit TLS.
|
||||
dial func(addr string, timeout time.Duration) (*email.Conn, error)
|
||||
|
||||
// disabled — core answered ErrUnknownMethod, i.e. it has no email block.
|
||||
disabled bool
|
||||
}
|
||||
|
||||
// pollOnce — one read of the mailbox, then one ingest per message.
|
||||
//
|
||||
// A fetch error aborts this poll and nothing else; the next tick tries again.
|
||||
// An ingest error for one message does not skip the rest — one mail the model
|
||||
// choked on should not hide the four behind it.
|
||||
func (r *reader) pollOnce(ctx context.Context, password string) {
|
||||
msgs, err := r.fetch(password)
|
||||
if err != nil {
|
||||
// The error may name a UID; it never names a subject or a sender.
|
||||
log.Printf("mavmaild: fetch: %v", err)
|
||||
if len(msgs) == 0 {
|
||||
return
|
||||
}
|
||||
}
|
||||
var junk, candidates, created int
|
||||
for _, m := range msgs {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
if m.Junk {
|
||||
junk++
|
||||
// Marked seen without a model call: the header filter already decided,
|
||||
// and re-classifying it every quarter hour would be pure waste.
|
||||
r.state.mark(m.UID)
|
||||
continue
|
||||
}
|
||||
resp, err := r.core.IngestMail(ctx, ipc.IngestMailReq{
|
||||
Mailbox: r.mailbox,
|
||||
UID: m.UID,
|
||||
From: m.From,
|
||||
Subject: m.Subject,
|
||||
Date: m.Date,
|
||||
Body: m.Body,
|
||||
})
|
||||
if errors.Is(err, ipc.ErrUnknownMethod) {
|
||||
log.Printf("mavmaild: core has no email block configured — mail ingestion is off; stopping")
|
||||
r.disabled = true
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
// Not marked seen: an ingest that failed should be retried next poll.
|
||||
log.Printf("mavmaild: ingest uid %d: %v", m.UID, err)
|
||||
continue
|
||||
}
|
||||
r.state.mark(m.UID)
|
||||
candidates += len(resp.TaskIDs)
|
||||
created += resp.Created
|
||||
}
|
||||
if err := r.state.save(); err != nil {
|
||||
log.Printf("mavmaild: state: %v", err)
|
||||
}
|
||||
log.Printf("mavmaild: %s: %d read, %d bulk, %d candidate(s), %d new", r.mailbox, len(msgs), junk, candidates, created)
|
||||
}
|
||||
|
||||
// fetch reads the mailbox. Messages already in the seen-set are not fetched at
|
||||
// all, so a steady mailbox costs one SEARCH per poll and nothing else.
|
||||
func (r *reader) fetch(password string) ([]email.Message, error) {
|
||||
f := email.FetchSince{
|
||||
Addr: r.addr,
|
||||
User: r.user,
|
||||
Mailbox: r.mailbox,
|
||||
Timeout: r.timeout,
|
||||
Since: time.Now().Add(-r.lookback),
|
||||
Max: r.max,
|
||||
Skip: r.state.seen,
|
||||
}
|
||||
return f.RunWith(password, r.dial)
|
||||
}
|
||||
|
||||
// ---- seen state ------------------------------------------------------------
|
||||
|
||||
// seenState — the UIDs already handed to core, persisted so a restart does not
|
||||
// re-read (and re-extract, at multi-second LLM cost) the whole lookback window.
|
||||
//
|
||||
// Correctness does not depend on it: ipc.CaptureTask dedupes on normalised text
|
||||
// among live tasks, so a re-read produces no duplicate rows. This exists to save
|
||||
// the model's time, which is why a broken state file is a log line rather than a
|
||||
// failure.
|
||||
//
|
||||
// UIDs are per-mailbox and monotonic, so the set is kept as a high-water mark
|
||||
// plus the stragglers above it. If the server ever changes UIDVALIDITY, UIDs
|
||||
// reset and the window is simply re-read once — dedupe absorbs it.
|
||||
type seenState struct {
|
||||
path string
|
||||
high uint32
|
||||
set map[uint32]bool
|
||||
dirty bool
|
||||
}
|
||||
|
||||
func newSeenState(path string) *seenState {
|
||||
return &seenState{path: path, set: map[uint32]bool{}}
|
||||
}
|
||||
|
||||
type seenFile struct {
|
||||
High uint32 `json:"high"`
|
||||
UIDs []uint32 `json:"uids,omitempty"`
|
||||
}
|
||||
|
||||
func (s *seenState) seen(uid uint32) bool {
|
||||
return uid <= s.high || s.set[uid]
|
||||
}
|
||||
|
||||
func (s *seenState) mark(uid uint32) {
|
||||
if s.seen(uid) {
|
||||
return
|
||||
}
|
||||
s.set[uid] = true
|
||||
s.dirty = true
|
||||
// Advance the high-water mark through any contiguous run, so the explicit set
|
||||
// stays small on a mailbox read in order.
|
||||
for {
|
||||
next := s.high + 1
|
||||
if !s.set[next] {
|
||||
break
|
||||
}
|
||||
delete(s.set, next)
|
||||
s.high = next
|
||||
}
|
||||
}
|
||||
|
||||
func (s *seenState) load() error {
|
||||
if s.path == "" {
|
||||
return nil
|
||||
}
|
||||
b, err := os.ReadFile(s.path)
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
return nil // first run
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var f seenFile
|
||||
if err := json.Unmarshal(b, &f); err != nil {
|
||||
return fmt.Errorf("parse %s: %w", s.path, err)
|
||||
}
|
||||
s.high = f.High
|
||||
for _, u := range f.UIDs {
|
||||
s.set[u] = true
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// save writes the state atomically (temp file + rename), 0600: it is a list of
|
||||
// message ids from his mailbox, which is metadata about his mail.
|
||||
func (s *seenState) save() error {
|
||||
if s.path == "" || !s.dirty {
|
||||
return nil
|
||||
}
|
||||
uids := make([]uint32, 0, len(s.set))
|
||||
for u := range s.set {
|
||||
uids = append(uids, u)
|
||||
}
|
||||
sort.Slice(uids, func(i, j int) bool { return uids[i] < uids[j] })
|
||||
b, err := json.Marshal(seenFile{High: s.high, UIDs: uids})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
tmp := s.path + ".tmp"
|
||||
if err := os.MkdirAll(filepath.Dir(s.path), 0o700); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := os.WriteFile(tmp, b, 0o600); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := os.Rename(tmp, s.path); err != nil {
|
||||
return err
|
||||
}
|
||||
s.dirty = false
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,254 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/kami/maven/internal/email"
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// ---- a scripted IMAP server, same shape internal/email's tests use ---------
|
||||
|
||||
type fakeIMAP struct {
|
||||
msgs map[uint32]string
|
||||
uids []uint32
|
||||
cmds []string
|
||||
}
|
||||
|
||||
func (f *fakeIMAP) serve(c net.Conn) {
|
||||
defer c.Close()
|
||||
fmt.Fprint(c, "* OK fake ready\r\n")
|
||||
r := bufio.NewReader(c)
|
||||
for {
|
||||
line, err := r.ReadString('\n')
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
parts := strings.SplitN(strings.TrimRight(line, "\r\n"), " ", 2)
|
||||
if len(parts) != 2 {
|
||||
return
|
||||
}
|
||||
tag, cmd := parts[0], parts[1]
|
||||
f.cmds = append(f.cmds, cmd)
|
||||
upper := strings.ToUpper(cmd)
|
||||
switch {
|
||||
case strings.HasPrefix(upper, "LOGIN"), strings.HasPrefix(upper, "EXAMINE"):
|
||||
fmt.Fprintf(c, "%s OK\r\n", tag)
|
||||
case strings.HasPrefix(upper, "UID SEARCH"):
|
||||
var ids []string
|
||||
for _, u := range f.uids {
|
||||
ids = append(ids, strconv.FormatUint(uint64(u), 10))
|
||||
}
|
||||
fmt.Fprintf(c, "* SEARCH %s\r\n%s OK\r\n", strings.Join(ids, " "), tag)
|
||||
case strings.HasPrefix(upper, "UID FETCH"):
|
||||
uid, _ := strconv.ParseUint(strings.Fields(cmd)[2], 10, 32)
|
||||
if raw, ok := f.msgs[uint32(uid)]; ok {
|
||||
fmt.Fprintf(c, "* 1 FETCH (UID %d BODY[] {%d}\r\n%s)\r\n", uid, len(raw), raw)
|
||||
}
|
||||
fmt.Fprintf(c, "%s OK\r\n", tag)
|
||||
case strings.HasPrefix(upper, "LOGOUT"):
|
||||
fmt.Fprintf(c, "* BYE\r\n%s OK\r\n", tag)
|
||||
return
|
||||
default:
|
||||
fmt.Fprintf(c, "%s BAD\r\n", tag)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (f *fakeIMAP) dial(_ string, timeout time.Duration) (*email.Conn, error) {
|
||||
cli, srv := net.Pipe()
|
||||
go f.serve(srv)
|
||||
return email.NewConn(cli, timeout)
|
||||
}
|
||||
|
||||
// ---- a fake core -----------------------------------------------------------
|
||||
|
||||
type fakeCore struct {
|
||||
got []ipc.IngestMailReq
|
||||
resp ipc.IngestMailResp
|
||||
err error
|
||||
}
|
||||
|
||||
func (c *fakeCore) IngestMail(_ context.Context, req ipc.IngestMailReq) (ipc.IngestMailResp, error) {
|
||||
c.got = append(c.got, req)
|
||||
if c.err != nil {
|
||||
return ipc.IngestMailResp{}, c.err
|
||||
}
|
||||
return c.resp, nil
|
||||
}
|
||||
|
||||
func mail(subject, body string, extraHeaders ...string) string {
|
||||
h := "Subject: " + subject + "\r\nFrom: a@b.c\r\nContent-Type: text/plain; charset=utf-8\r\n"
|
||||
for _, e := range extraHeaders {
|
||||
h += e + "\r\n"
|
||||
}
|
||||
return h + "\r\n" + body + "\r\n"
|
||||
}
|
||||
|
||||
func newTestReader(t *testing.T, f *fakeIMAP, core *fakeCore, statePath string) *reader {
|
||||
t.Helper()
|
||||
return &reader{
|
||||
core: core, addr: "mail.example:993", user: "kami", mailbox: "INBOX",
|
||||
lookback: 72 * time.Hour, max: 25, timeout: 5 * time.Second,
|
||||
state: newSeenState(statePath),
|
||||
dial: f.dial,
|
||||
}
|
||||
}
|
||||
|
||||
func TestPollHandsMessagesToCore(t *testing.T) {
|
||||
f := &fakeIMAP{
|
||||
uids: []uint32{1, 2},
|
||||
msgs: map[uint32]string{
|
||||
1: mail("Счёт", "Оплатить до 5 августа."),
|
||||
2: mail("Скидки", "Sale!", "List-Unsubscribe: <mailto:u@x>"),
|
||||
},
|
||||
}
|
||||
core := &fakeCore{resp: ipc.IngestMailResp{TaskIDs: []int64{1}, Created: 1}}
|
||||
r := newTestReader(t, f, core, "")
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
|
||||
// The newsletter is filtered before core is asked: only the real mail crosses.
|
||||
if len(core.got) != 1 {
|
||||
t.Fatalf("core saw %d messages, want 1 (the bulk one must not cross): %+v", len(core.got), core.got)
|
||||
}
|
||||
got := core.got[0]
|
||||
if got.UID != 1 || got.Mailbox != "INBOX" || got.Subject != "Счёт" {
|
||||
t.Errorf("ingest req = %+v", got)
|
||||
}
|
||||
if !strings.Contains(got.Body, "Оплатить") {
|
||||
t.Errorf("body = %q", got.Body)
|
||||
}
|
||||
}
|
||||
|
||||
// A second poll must not re-send what core already saw — extraction is a
|
||||
// multi-second LLM call per message.
|
||||
func TestPollSkipsSeenUIDs(t *testing.T) {
|
||||
f := &fakeIMAP{uids: []uint32{5}, msgs: map[uint32]string{5: mail("Счёт", "текст")}}
|
||||
core := &fakeCore{}
|
||||
r := newTestReader(t, f, core, "")
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
if len(core.got) != 1 {
|
||||
t.Errorf("core saw %d messages over two polls, want 1", len(core.got))
|
||||
}
|
||||
}
|
||||
|
||||
// An ingest that failed is NOT marked seen: the next poll retries it.
|
||||
func TestPollRetriesFailedIngest(t *testing.T) {
|
||||
f := &fakeIMAP{uids: []uint32{5}, msgs: map[uint32]string{5: mail("Счёт", "текст")}}
|
||||
core := &fakeCore{err: fmt.Errorf("llama-server is warming up")}
|
||||
r := newTestReader(t, f, core, "")
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
core.err = nil
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
if len(core.got) != 2 {
|
||||
t.Errorf("core saw %d attempts, want 2 (a failed ingest is retried)", len(core.got))
|
||||
}
|
||||
}
|
||||
|
||||
// Core without an email block ⇒ stop, don't hammer the socket.
|
||||
func TestPollStopsWhenCoreRefusesMail(t *testing.T) {
|
||||
f := &fakeIMAP{uids: []uint32{1, 2}, msgs: map[uint32]string{1: mail("a", "b"), 2: mail("c", "d")}}
|
||||
core := &fakeCore{err: fmt.Errorf("call: %w", ipc.ErrUnknownMethod)}
|
||||
r := newTestReader(t, f, core, "")
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
if !r.disabled {
|
||||
t.Error("ErrUnknownMethod must disable the reader")
|
||||
}
|
||||
if len(core.got) != 1 {
|
||||
t.Errorf("core saw %d messages, want 1 — stop at the first refusal", len(core.got))
|
||||
}
|
||||
}
|
||||
|
||||
func TestSeenStatePersists(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "state", "seen.json")
|
||||
f := &fakeIMAP{uids: []uint32{9}, msgs: map[uint32]string{9: mail("Счёт", "текст")}}
|
||||
core := &fakeCore{}
|
||||
r := newTestReader(t, f, core, path)
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
|
||||
fi, err := os.Stat(path)
|
||||
if err != nil {
|
||||
t.Fatalf("state file: %v", err)
|
||||
}
|
||||
// A list of message ids from his mailbox is metadata about his mail.
|
||||
if perm := fi.Mode().Perm(); perm != 0o600 {
|
||||
t.Errorf("state file mode = %v, want 0600", perm)
|
||||
}
|
||||
|
||||
// A fresh reader with the same state file must not re-read the message.
|
||||
core2 := &fakeCore{}
|
||||
r2 := newTestReader(t, f, core2, path)
|
||||
if err := r2.state.load(); err != nil {
|
||||
t.Fatalf("load: %v", err)
|
||||
}
|
||||
r2.pollOnce(context.Background(), "secret")
|
||||
if len(core2.got) != 0 {
|
||||
t.Errorf("after a restart core saw %d messages, want 0", len(core2.got))
|
||||
}
|
||||
}
|
||||
|
||||
func TestSeenStateHighWaterMark(t *testing.T) {
|
||||
s := newSeenState("")
|
||||
s.mark(1)
|
||||
s.mark(3)
|
||||
s.mark(2)
|
||||
if s.high != 3 {
|
||||
t.Errorf("high = %d, want 3 (contiguous run collapses)", s.high)
|
||||
}
|
||||
if len(s.set) != 0 {
|
||||
t.Errorf("explicit set = %v, want empty", s.set)
|
||||
}
|
||||
if !s.seen(2) || s.seen(4) {
|
||||
t.Errorf("seen(2)=%v seen(4)=%v", s.seen(2), s.seen(4))
|
||||
}
|
||||
}
|
||||
|
||||
func TestSeenStateCorruptFileIsNotFatal(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "seen.json")
|
||||
if err := os.WriteFile(path, []byte("{not json"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := newSeenState(path)
|
||||
if err := s.load(); err == nil {
|
||||
t.Error("a corrupt state file should report an error the caller logs")
|
||||
}
|
||||
if s.seen(1) {
|
||||
t.Error("a corrupt state file must leave an empty seen-set, not a poisoned one")
|
||||
}
|
||||
}
|
||||
|
||||
// Off unless configured, and the credential is never a flag value.
|
||||
func TestRunRequiresConfig(t *testing.T) {
|
||||
if err := run([]string{}); err == nil {
|
||||
t.Error("no -socket must be an error")
|
||||
}
|
||||
if err := run([]string{"-socket", "/tmp/nope.sock"}); err == nil {
|
||||
t.Error("no mailbox configuration must be an error, not a default mailbox")
|
||||
}
|
||||
// There is no -password flag at all: only -password-file.
|
||||
if err := run([]string{"-socket", "/x", "-imap", "h", "-user", "u", "-password", "p"}); err == nil ||
|
||||
!strings.Contains(err.Error(), "flag provided but not defined") {
|
||||
t.Errorf("a -password flag must not exist; err = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunRejectsEmptyPasswordFile(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "pass")
|
||||
if err := os.WriteFile(path, []byte(" \n"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err := run([]string{"-socket", "/x/y.sock", "-imap", "h", "-user", "u", "-password-file", path})
|
||||
if err == nil || !strings.Contains(err.Error(), "empty") {
|
||||
t.Errorf("an empty password file must be refused before dialling; err = %v", err)
|
||||
}
|
||||
}
|
||||
@@ -113,6 +113,31 @@ services:
|
||||
- sockets:/run/maven
|
||||
# - ./deploy/zenmoney.token:/run/secrets/zenmoney.token:ro
|
||||
|
||||
# The mail reader (Vikunja #246) is OFF and commented out: it needs an IMAP
|
||||
# account, and there is none on this box. mavmaild reads the password from a
|
||||
# FILE so it never appears in `ps`, in this file, or in shell history — the
|
||||
# same rule mavpoll follows for the zenmoney token. Core never sees the
|
||||
# password: the reader hands core message text on one IPC method, and core
|
||||
# writes what the model extracts as task CANDIDATES he reviews on /tasks.
|
||||
# Nothing here can create a reminder, so a misread mail cannot fire.
|
||||
#
|
||||
# To enable: write the password to deploy/imap.password (0600, gitignored),
|
||||
# add an "email": {} block to deploy/mavend.json, and uncomment this service.
|
||||
# mavmaild:
|
||||
# <<: *image
|
||||
# command: ["mavmaild", "-socket", "/run/maven/mavend.sock",
|
||||
# "-imap", "imap.example.org:993",
|
||||
# "-user", "kami@example.org",
|
||||
# "-password-file", "/run/secrets/imap.password",
|
||||
# "-mailbox", "INBOX",
|
||||
# "-interval", "15m",
|
||||
# "-state", "/var/lib/maven/mail-seen.json"]
|
||||
# depends_on: [mavend]
|
||||
# volumes:
|
||||
# - sockets:/run/maven
|
||||
# - dbdata:/var/lib/maven
|
||||
# - ./deploy/imap.password:/run/secrets/imap.password:ro
|
||||
|
||||
volumes:
|
||||
dbdata:
|
||||
sockets:
|
||||
|
||||
@@ -26,21 +26,23 @@ type FetchSince struct {
|
||||
Since time.Time
|
||||
Max int
|
||||
Skip func(uid uint32) bool
|
||||
|
||||
// dial is the connection seam. nil means Dial (implicit TLS); the tests set
|
||||
// it to an in-process fake. Unexported so no configuration path can point
|
||||
// the reader at a non-TLS transport.
|
||||
dial func(addr string, timeout time.Duration) (*Conn, error)
|
||||
}
|
||||
|
||||
// Run performs one read. password is passed here, not stored in the struct, so
|
||||
// the configuration of a mailbox and the secret for it are never the same value
|
||||
// sitting in the same place.
|
||||
func (f FetchSince) Run(password string) ([]Message, error) {
|
||||
return f.RunWith(password, nil)
|
||||
}
|
||||
|
||||
// RunWith is Run with an explicit connection function, which is how the reader
|
||||
// daemon and the tests substitute an in-process server. nil ⇒ Dial, i.e.
|
||||
// implicit TLS with certificate verification; there is no configuration path
|
||||
// that reaches this, so no deployment can end up talking cleartext IMAP.
|
||||
func (f FetchSince) RunWith(password string, dial func(addr string, timeout time.Duration) (*Conn, error)) ([]Message, error) {
|
||||
if f.Addr == "" || f.User == "" || f.Mailbox == "" {
|
||||
return nil, fmt.Errorf("email: mailbox not configured (addr/user/mailbox)")
|
||||
}
|
||||
dial := f.dial
|
||||
if dial == nil {
|
||||
dial = Dial
|
||||
}
|
||||
|
||||
@@ -21,13 +21,12 @@ func TestFetchSinceRun(t *testing.T) {
|
||||
Since: time.Date(2026, 7, 30, 0, 0, 0, 0, time.UTC),
|
||||
Max: 2,
|
||||
Skip: func(uid uint32) bool { return uid == 3 },
|
||||
dial: func(addr string, timeout time.Duration) (*Conn, error) {
|
||||
cli, srv := net.Pipe()
|
||||
go f.serve(t, srv)
|
||||
return NewConn(cli, timeout)
|
||||
},
|
||||
}
|
||||
msgs, err := fs.Run("secret")
|
||||
msgs, err := fs.RunWith("secret", func(addr string, timeout time.Duration) (*Conn, error) {
|
||||
cli, srv := net.Pipe()
|
||||
go f.serve(t, srv)
|
||||
return NewConn(cli, timeout)
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("run: %v", err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user