Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ee7bec11e3 | |||
| f42d1594ef |
@@ -7,6 +7,7 @@
|
|||||||
/mavpoll
|
/mavpoll
|
||||||
/mavcaldav
|
/mavcaldav
|
||||||
/mavwaked
|
/mavwaked
|
||||||
|
/mavmaild
|
||||||
|
|
||||||
# Certs (private keys, don't commit)
|
# Certs (private keys, don't commit)
|
||||||
certs/
|
certs/
|
||||||
@@ -36,6 +37,8 @@ deploy/db_key.env
|
|||||||
deploy/telegram.env
|
deploy/telegram.env
|
||||||
# zenmoney API token, read by mavpoll (never in argv, never committed)
|
# zenmoney API token, read by mavpoll (never in argv, never committed)
|
||||||
deploy/zenmoney.token
|
deploy/zenmoney.token
|
||||||
|
# IMAP password, read by mavmaild (never in argv, never committed)
|
||||||
|
deploy/imap.password
|
||||||
|
|
||||||
# Temp files
|
# Temp files
|
||||||
/tmp/
|
/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`:
|
and libs wired through the Makefile — **do not** call `go build` on them bare, use `make`:
|
||||||
|
|
||||||
```sh
|
```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 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
|
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). |
|
| `mavenclient` | Voice loop client (mic → stt → core → tts). |
|
||||||
| `mavpoll` | Telegram long-poll reach. |
|
| `mavpoll` | Telegram long-poll reach. |
|
||||||
| `mavcaldav` | CalDAV calendar sync. |
|
| `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
|
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
|
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/mavttsd ./cmd/mavttsd && \
|
||||||
go build -o /out/mavweb ./cmd/mavweb && \
|
go build -o /out/mavweb ./cmd/mavweb && \
|
||||||
go build -o /out/mavpoll ./cmd/mavpoll && \
|
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
|
# 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
|
# 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
|
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:
|
build-stt:
|
||||||
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
|
CGO_CFLAGS="$(CGO_CFLAGS)" CGO_LDFLAGS="$(CGO_LDFLAGS)" LD_LIBRARY_PATH="$(shell pwd)/deps/lib" \
|
||||||
@@ -50,6 +50,9 @@ build-poll:
|
|||||||
build-caldav:
|
build-caldav:
|
||||||
$(GO) build $(GOFLAGS) -o mavcaldav ./cmd/mavcaldav/
|
$(GO) build $(GOFLAGS) -o mavcaldav ./cmd/mavcaldav/
|
||||||
|
|
||||||
|
build-mail:
|
||||||
|
$(GO) build $(GOFLAGS) -o mavmaild ./cmd/mavmaild/
|
||||||
|
|
||||||
run-web: build-web
|
run-web: build-web
|
||||||
./mavweb -addr :9200 -voice 127.0.0.1:9100
|
./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/'
|
@echo ' sudo cp onnxruntime-linux-x64-1.15.1/lib/libonnxruntime.so* /usr/local/lib/'
|
||||||
|
|
||||||
clean:
|
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,158 @@
|
|||||||
|
// mavend/mail.go — core's half of the email reader (Vikunja #246,
|
||||||
|
// docs/plans/01-email-reader.md).
|
||||||
|
//
|
||||||
|
// The split: cmd/mavmaild holds the IMAP credential, connects to the mailbox
|
||||||
|
// and converts messages to plaintext; it hands each message to core over
|
||||||
|
// ipc.MethodIngestMail. Core runs the extraction on the resident model —
|
||||||
|
// llama-server lives in this process, spawned by the phraser — and writes what
|
||||||
|
// comes back through the one task intake seam.
|
||||||
|
//
|
||||||
|
// What this file may produce is exactly one thing: rows in `tasks` with status
|
||||||
|
// "candidate". No fact, no reminder, no note, no nudge, no calendar event. A
|
||||||
|
// 1.7B misreading a mail can therefore put a wrong line on a review page and
|
||||||
|
// nothing else; it can never make Maven speak, and it can never make her
|
||||||
|
// recite something out of an advert as true.
|
||||||
|
//
|
||||||
|
// Off unless configured twice over: no `email` block in mavend.json ⇒ the IPC
|
||||||
|
// method does not exist; no llama-server phraser ⇒ same. A reader pointed at a
|
||||||
|
// core that is not set up for mail gets ErrUnknownMethod rather than silence.
|
||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"log"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/config"
|
||||||
|
"github.com/kami/maven/internal/email"
|
||||||
|
"github.com/kami/maven/internal/ipc"
|
||||||
|
"github.com/kami/maven/internal/llm"
|
||||||
|
"github.com/kami/maven/internal/phraser"
|
||||||
|
"github.com/kami/maven/internal/store"
|
||||||
|
)
|
||||||
|
|
||||||
|
// evidenceMaxChars — how much of the subject line is kept as a candidate's
|
||||||
|
// evidence. Enough to recognise the mail on /tasks, not enough to turn the task
|
||||||
|
// list into a copy of his mailbox.
|
||||||
|
const evidenceMaxChars = 160
|
||||||
|
|
||||||
|
// mailIntake — extraction + capture for one message at a time.
|
||||||
|
type mailIntake struct {
|
||||||
|
st *store.Store
|
||||||
|
ex *email.Extractor
|
||||||
|
timeout time.Duration
|
||||||
|
now func() time.Time
|
||||||
|
}
|
||||||
|
|
||||||
|
// newMailIntake returns nil when mail ingestion must not be available, which is
|
||||||
|
// the default. Both preconditions are real:
|
||||||
|
//
|
||||||
|
// - no cfg.Email ⇒ not configured, and a capability is off unless configured;
|
||||||
|
// - no llama-server phraser ⇒ nothing to extract with. There is deliberately
|
||||||
|
// no keyword fallback: "the subject line became a task" is not extraction,
|
||||||
|
// it is a mailbox rendered as a to-do list, and it would fill the review
|
||||||
|
// page faster than he could clear it.
|
||||||
|
func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config) *mailIntake {
|
||||||
|
if cfg.Email == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
lp, ok := phr.(*phraser.LLMPhraser)
|
||||||
|
if !ok {
|
||||||
|
log.Printf("mail intake: configured but no llama-server phraser — mail ingestion disabled")
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
timeout := time.Duration(cfg.Email.Timeout)
|
||||||
|
if timeout <= 0 {
|
||||||
|
timeout = config.DefaultEmailTimeout
|
||||||
|
}
|
||||||
|
ex := email.NewExtractor(llm.New(lp.BaseURL(), timeout), cfg.Email.MaxTasks, contextBlockFn(cfg, time.Now))
|
||||||
|
log.Printf("mail intake: enabled (max %d candidates per message, timeout %s)", cfg.Email.MaxTasks, timeout)
|
||||||
|
return &mailIntake{st: st, ex: ex, timeout: timeout, now: time.Now}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ingest handles one ipc.MethodIngestMail call.
|
||||||
|
//
|
||||||
|
// Junk and empty messages are answered Skipped without touching the model — the
|
||||||
|
// reader's header filter is what keeps the resident model off newsletters.
|
||||||
|
//
|
||||||
|
// Every candidate is captured with Status "candidate", Source "email:<mailbox>"
|
||||||
|
// and the subject as Evidence. CaptureTask dedupes on normalised text among
|
||||||
|
// live rows, so a mailbox re-read after a restart produces Created=0 rather
|
||||||
|
// than a second copy of every task.
|
||||||
|
func (m *mailIntake) ingest(ctx context.Context, req ipc.IngestMailReq) (ipc.IngestMailResp, error) {
|
||||||
|
msg := email.Message{
|
||||||
|
UID: req.UID,
|
||||||
|
From: req.From,
|
||||||
|
Subject: req.Subject,
|
||||||
|
Date: req.Date,
|
||||||
|
Body: req.Body,
|
||||||
|
Junk: req.Junk,
|
||||||
|
}
|
||||||
|
if msg.Junk || (msg.Subject == "" && msg.Body == "") {
|
||||||
|
return ipc.IngestMailResp{Skipped: true}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(ctx, m.timeout)
|
||||||
|
defer cancel()
|
||||||
|
cands, err := m.ex.Extract(ctx, msg)
|
||||||
|
if err != nil {
|
||||||
|
// The error from internal/email never carries mail text; keep it that way
|
||||||
|
// by not adding the subject here.
|
||||||
|
return ipc.IngestMailResp{}, fmt.Errorf("mail intake: uid %d: %w", req.UID, err)
|
||||||
|
}
|
||||||
|
if len(cands) == 0 {
|
||||||
|
return ipc.IngestMailResp{}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
source := email.SourcePrefix + req.Mailbox
|
||||||
|
evidence := truncateRunes(req.Subject, evidenceMaxChars)
|
||||||
|
now := m.now()
|
||||||
|
var resp ipc.IngestMailResp
|
||||||
|
for _, c := range cands {
|
||||||
|
t := store.Task{
|
||||||
|
CreatedTs: now,
|
||||||
|
Text: c.Text,
|
||||||
|
Source: source,
|
||||||
|
Evidence: evidence,
|
||||||
|
// The one status this path may ever write. Anything Maven derived from
|
||||||
|
// something she read is a suggestion until he confirms it on /tasks.
|
||||||
|
Status: store.TaskCandidate,
|
||||||
|
}
|
||||||
|
if due, ok := email.ParseDue(c.Due); ok {
|
||||||
|
t.Due = &due
|
||||||
|
}
|
||||||
|
id, created, err := m.st.CaptureTask(ctx, t)
|
||||||
|
if err != nil {
|
||||||
|
return resp, fmt.Errorf("mail intake: capture: %w", err)
|
||||||
|
}
|
||||||
|
resp.TaskIDs = append(resp.TaskIDs, id)
|
||||||
|
if created {
|
||||||
|
resp.Created++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// Counts only: the log line names the mailbox and the UID, never the subject,
|
||||||
|
// the sender or the task text. Reviewing a candidate is what /tasks is for.
|
||||||
|
log.Printf("mail intake: %s uid %d → %d candidate(s), %d new", source, req.UID, len(cands), resp.Created)
|
||||||
|
return resp, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// wireMailIntake installs the IPC hook, or leaves it nil so the method reports
|
||||||
|
// ErrUnknownMethod. Called on both startup paths (unlocked boot and passkey
|
||||||
|
// unlock) so mail behaves the same either way.
|
||||||
|
func wireMailIntake(srv *ipc.Server, st *store.Store, phr phraser.Phraser, cfg *config.Config) {
|
||||||
|
mi := newMailIntake(st, phr, cfg)
|
||||||
|
if mi == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
srv.IngestMailFn = mi.ingest
|
||||||
|
}
|
||||||
|
|
||||||
|
// truncateRunes cuts a string to n runes, marking the cut.
|
||||||
|
func truncateRunes(s string, n int) string {
|
||||||
|
r := []rune(s)
|
||||||
|
if len(r) <= n {
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
return string(r[:n]) + "…"
|
||||||
|
}
|
||||||
@@ -0,0 +1,182 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/config"
|
||||||
|
"github.com/kami/maven/internal/email"
|
||||||
|
"github.com/kami/maven/internal/ipc"
|
||||||
|
"github.com/kami/maven/internal/llm"
|
||||||
|
"github.com/kami/maven/internal/store"
|
||||||
|
)
|
||||||
|
|
||||||
|
// mailLLM — a canned extraction reply.
|
||||||
|
type mailLLM struct {
|
||||||
|
reply string
|
||||||
|
calls int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *mailLLM) Complete(_ context.Context, _ llm.Req) (string, error) {
|
||||||
|
m.calls++
|
||||||
|
return m.reply, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func newTestIntake(t *testing.T, reply string) (*mailIntake, *store.Store, *mailLLM) {
|
||||||
|
t.Helper()
|
||||||
|
st := newTestStore(t)
|
||||||
|
fake := &mailLLM{reply: reply}
|
||||||
|
return &mailIntake{
|
||||||
|
st: st,
|
||||||
|
ex: email.NewExtractor(fake, 0, nil),
|
||||||
|
timeout: 5 * time.Second,
|
||||||
|
now: func() time.Time { return time.Date(2026, 8, 1, 10, 0, 0, 0, time.UTC) },
|
||||||
|
}, st, fake
|
||||||
|
}
|
||||||
|
|
||||||
|
func ingestReq() ipc.IngestMailReq {
|
||||||
|
return ipc.IngestMailReq{
|
||||||
|
Mailbox: "INBOX", UID: 42,
|
||||||
|
From: "billing@isp.example",
|
||||||
|
Subject: "Счёт за интернет",
|
||||||
|
Body: "Оплатите счёт до 5 августа.",
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// The one property that matters: a mail-derived task is a candidate, attributed
|
||||||
|
// to the mailbox, with the subject as reviewable evidence — and nothing else is
|
||||||
|
// written.
|
||||||
|
func TestIngestCapturesCandidates(t *testing.T) {
|
||||||
|
mi, st, _ := newTestIntake(t, `[{"text":"оплатить счёт за интернет","due":"2026-08-05"}]`)
|
||||||
|
resp, err := mi.ingest(context.Background(), ingestReq())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ingest: %v", err)
|
||||||
|
}
|
||||||
|
if resp.Created != 1 || len(resp.TaskIDs) != 1 {
|
||||||
|
t.Fatalf("resp = %+v, want one created task", resp)
|
||||||
|
}
|
||||||
|
tasks, err := st.ListTasks(context.Background(), "")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list: %v", err)
|
||||||
|
}
|
||||||
|
if len(tasks) != 1 {
|
||||||
|
t.Fatalf("got %d tasks, want 1", len(tasks))
|
||||||
|
}
|
||||||
|
got := tasks[0]
|
||||||
|
if got.Status != store.TaskCandidate {
|
||||||
|
t.Errorf("status = %q, want %q — mail may only produce candidates", got.Status, store.TaskCandidate)
|
||||||
|
}
|
||||||
|
if got.Source != "email:INBOX" {
|
||||||
|
t.Errorf("source = %q, want email:INBOX", got.Source)
|
||||||
|
}
|
||||||
|
if got.Evidence != "Счёт за интернет" {
|
||||||
|
t.Errorf("evidence = %q, want the subject line", got.Evidence)
|
||||||
|
}
|
||||||
|
if got.Due == nil || got.Due.Format("2006-01-02") != "2026-08-05" {
|
||||||
|
t.Errorf("due = %v, want 2026-08-05", got.Due)
|
||||||
|
}
|
||||||
|
// Nothing else may have been written: no reminder, no fact.
|
||||||
|
rem, err := st.ListReminders(context.Background(), 10)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("list reminders: %v", err)
|
||||||
|
}
|
||||||
|
if len(rem) != 0 {
|
||||||
|
t.Errorf("mail created %d reminders; a misread mail must never be able to fire", len(rem))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Re-reading a mailbox must not grow the list — CaptureTask dedupes among live
|
||||||
|
// rows, and the intake relies on exactly that.
|
||||||
|
func TestIngestSameMailTwiceIsIdempotent(t *testing.T) {
|
||||||
|
mi, st, _ := newTestIntake(t, `[{"text":"оплатить счёт","due":""}]`)
|
||||||
|
if _, err := mi.ingest(context.Background(), ingestReq()); err != nil {
|
||||||
|
t.Fatalf("first ingest: %v", err)
|
||||||
|
}
|
||||||
|
resp, err := mi.ingest(context.Background(), ingestReq())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("second ingest: %v", err)
|
||||||
|
}
|
||||||
|
if resp.Created != 0 || len(resp.TaskIDs) != 1 {
|
||||||
|
t.Errorf("resp = %+v, want the existing row and Created=0", resp)
|
||||||
|
}
|
||||||
|
tasks, _ := st.ListTasks(context.Background(), "")
|
||||||
|
if len(tasks) != 1 {
|
||||||
|
t.Errorf("got %d tasks after two reads, want 1", len(tasks))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestIngestJunkSkipsTheModel(t *testing.T) {
|
||||||
|
mi, st, fake := newTestIntake(t, `[{"text":"купить со скидкой","due":""}]`)
|
||||||
|
req := ingestReq()
|
||||||
|
req.Junk = true
|
||||||
|
resp, err := mi.ingest(context.Background(), req)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ingest: %v", err)
|
||||||
|
}
|
||||||
|
if !resp.Skipped || resp.Created != 0 {
|
||||||
|
t.Errorf("resp = %+v, want skipped", resp)
|
||||||
|
}
|
||||||
|
if fake.calls != 0 {
|
||||||
|
t.Errorf("model called %d times for junk, want 0", fake.calls)
|
||||||
|
}
|
||||||
|
if tasks, _ := st.ListTasks(context.Background(), ""); len(tasks) != 0 {
|
||||||
|
t.Errorf("junk produced %d tasks, want 0", len(tasks))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestIngestEmptyMessageSkipped(t *testing.T) {
|
||||||
|
mi, _, fake := newTestIntake(t, "[]")
|
||||||
|
resp, err := mi.ingest(context.Background(), ipc.IngestMailReq{Mailbox: "INBOX", UID: 1})
|
||||||
|
if err != nil || !resp.Skipped {
|
||||||
|
t.Fatalf("resp = %+v, err = %v; want skipped", resp, err)
|
||||||
|
}
|
||||||
|
if fake.calls != 0 {
|
||||||
|
t.Errorf("model called %d times for an empty message, want 0", fake.calls)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestIngestNoTasksWritesNothing(t *testing.T) {
|
||||||
|
mi, st, _ := newTestIntake(t, "[]")
|
||||||
|
resp, err := mi.ingest(context.Background(), ingestReq())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("ingest: %v", err)
|
||||||
|
}
|
||||||
|
if resp.Created != 0 || len(resp.TaskIDs) != 0 || resp.Skipped {
|
||||||
|
t.Errorf("resp = %+v, want nothing captured and not skipped", resp)
|
||||||
|
}
|
||||||
|
if tasks, _ := st.ListTasks(context.Background(), ""); len(tasks) != 0 {
|
||||||
|
t.Errorf("got %d tasks, want 0", len(tasks))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestIngestTruncatesEvidence(t *testing.T) {
|
||||||
|
mi, st, _ := newTestIntake(t, `[{"text":"дело","due":""}]`)
|
||||||
|
req := ingestReq()
|
||||||
|
req.Subject = strings.Repeat("щ", 400)
|
||||||
|
if _, err := mi.ingest(context.Background(), req); err != nil {
|
||||||
|
t.Fatalf("ingest: %v", err)
|
||||||
|
}
|
||||||
|
tasks, _ := st.ListTasks(context.Background(), "")
|
||||||
|
if len(tasks) != 1 {
|
||||||
|
t.Fatalf("got %d tasks, want 1", len(tasks))
|
||||||
|
}
|
||||||
|
if n := len([]rune(tasks[0].Evidence)); n > evidenceMaxChars+1 {
|
||||||
|
t.Errorf("evidence kept %d runes, want ≤ %d", n, evidenceMaxChars)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Off unless configured: no email block ⇒ no intake, so the IPC method does not
|
||||||
|
// exist at all.
|
||||||
|
func TestNewMailIntakeOffWithoutConfig(t *testing.T) {
|
||||||
|
st := newTestStore(t)
|
||||||
|
if mi := newMailIntake(st, nil, &config.Config{}); mi != nil {
|
||||||
|
t.Error("no email block must mean no mail intake")
|
||||||
|
}
|
||||||
|
// Configured but with a non-LLM phraser: still off — there is no fallback
|
||||||
|
// extraction, by design.
|
||||||
|
if mi := newMailIntake(st, nil, &config.Config{Email: &config.EmailConfig{}}); mi != nil {
|
||||||
|
t.Error("without a llama-server phraser there is nothing to extract with")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -320,6 +320,13 @@ func run(args []string) error {
|
|||||||
|
|
||||||
srv.StepUp = func(ctx context.Context) error { return passkeySess.Assert(ctx, auth.Scope{}) }
|
srv.StepUp = func(ctx context.Context) error { return passkeySess.Assert(ctx, auth.Scope{}) }
|
||||||
|
|
||||||
|
// Mail ingestion (Vikunja #246): the hook stays nil unless an email block is
|
||||||
|
// configured and there is a llama-server to extract with, in which case
|
||||||
|
// ipc.MethodIngestMail reports ErrUnknownMethod.
|
||||||
|
if !locked {
|
||||||
|
wireMailIntake(srv, st, phr, cfg)
|
||||||
|
}
|
||||||
|
|
||||||
// WrapKeyFn — wraps the env key with a passkey credential public key and
|
// WrapKeyFn — wraps the env key with a passkey credential public key and
|
||||||
// persists the wrapped blob. Only wired when the daemon has the key in
|
// persists the wrapped blob. Only wired when the daemon has the key in
|
||||||
// memory (env key mode). Called by mavweb after passkey enrollment.
|
// memory (env key mode). Called by mavweb after passkey enrollment.
|
||||||
@@ -457,6 +464,7 @@ func run(args []string) error {
|
|||||||
}
|
}
|
||||||
srv.SetAPI(newAPI)
|
srv.SetAPI(newAPI)
|
||||||
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
|
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
|
||||||
|
wireMailIntake(srv, st, phr, cfg)
|
||||||
|
|
||||||
// Start voice server.
|
// Start voice server.
|
||||||
if voiceW != nil {
|
if voiceW != nil {
|
||||||
|
|||||||
@@ -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
|
- sockets:/run/maven
|
||||||
# - ./deploy/zenmoney.token:/run/secrets/zenmoney.token:ro
|
# - ./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:
|
volumes:
|
||||||
dbdata:
|
dbdata:
|
||||||
sockets:
|
sockets:
|
||||||
|
|||||||
@@ -74,7 +74,12 @@ func Requirement(m ipc.Method) Authority {
|
|||||||
// existing analogue.
|
// existing analogue.
|
||||||
ipc.MethodCaptureTask,
|
ipc.MethodCaptureTask,
|
||||||
ipc.MethodListTasks,
|
ipc.MethodListTasks,
|
||||||
ipc.MethodSetTaskStatus:
|
ipc.MethodSetTaskStatus,
|
||||||
|
// Mail ingestion (Vikunja #246). AuthRead because of what the method can
|
||||||
|
// produce: candidate tasks and nothing else. It cannot write a fact, set a
|
||||||
|
// reminder, or touch the tool allowlist, so a compromised mail reader can
|
||||||
|
// at worst put junk on a review page he clears in one click.
|
||||||
|
ipc.MethodIngestMail:
|
||||||
return AuthRead
|
return AuthRead
|
||||||
}
|
}
|
||||||
// Unknown method ⇒ AuthRead, but ipc.dispatch returns ErrUnknownMethod
|
// Unknown method ⇒ AuthRead, but ipc.dispatch returns ErrUnknownMethod
|
||||||
|
|||||||
@@ -150,6 +150,12 @@ type Config struct {
|
|||||||
// absent ⇒ no evaluation loop at all. See MemoryEvalConfig.
|
// absent ⇒ no evaluation loop at all. See MemoryEvalConfig.
|
||||||
MemoryEval *MemoryEvalConfig `json:"memory_eval,omitempty"`
|
MemoryEval *MemoryEvalConfig `json:"memory_eval,omitempty"`
|
||||||
|
|
||||||
|
// Email — mail ingestion (Vikunja #246). nil / absent ⇒ core refuses
|
||||||
|
// ipc.MethodIngestMail outright, so a mail reader cannot make Maven read a
|
||||||
|
// mailbox by merely existing. See EmailConfig; the IMAP host and credential
|
||||||
|
// live in the reader (cmd/mavmaild), never here.
|
||||||
|
Email *EmailConfig `json:"email,omitempty"`
|
||||||
|
|
||||||
// Praxis — the ecosystem attention-state service. When configured, maven
|
// Praxis — the ecosystem attention-state service. When configured, maven
|
||||||
// calls the Praxis HTTP tools API for attention listing and item lifecycle.
|
// calls the Praxis HTTP tools API for attention listing and item lifecycle.
|
||||||
// Maven never touches Praxis's database directly (ecosystem invariant: no
|
// Maven never touches Praxis's database directly (ecosystem invariant: no
|
||||||
@@ -418,6 +424,26 @@ type MemoryEvalConfig struct {
|
|||||||
MinConfidence float64 `json:"min_confidence,omitempty"`
|
MinConfidence float64 `json:"min_confidence,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// EmailConfig — core's half of the email reader: how many task candidates one
|
||||||
|
// message may produce, and how long the extraction call may take.
|
||||||
|
//
|
||||||
|
// There is deliberately nothing about a mailbox here. Core does not connect to
|
||||||
|
// IMAP, does not know an account exists, and holds no mail credential — the
|
||||||
|
// reader daemon does, the same split mavpoll uses for the zenmoney token. This
|
||||||
|
// block only says "extraction is allowed, with these bounds".
|
||||||
|
type EmailConfig struct {
|
||||||
|
// MaxTasks — candidates per message. 0 ⇒ email.MaxCandidates (3).
|
||||||
|
MaxTasks int `json:"max_tasks,omitempty"`
|
||||||
|
|
||||||
|
// Timeout — per-message extraction budget. 0 ⇒ DefaultEmailTimeout. This is
|
||||||
|
// a Thinking model reading a mail; nobody is waiting on the answer, but a
|
||||||
|
// hung llama-server must not pin the reader's connection forever.
|
||||||
|
Timeout Duration `json:"timeout,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// DefaultEmailTimeout — extraction budget per message.
|
||||||
|
const DefaultEmailTimeout = 2 * time.Minute
|
||||||
|
|
||||||
// PhraserConfig — the LLM-backed phraser seam. The daemon spawns llama-server
|
// PhraserConfig — the LLM-backed phraser seam. The daemon spawns llama-server
|
||||||
// as a managed subprocess and sends chat-completion requests to phrase nudge
|
// as a managed subprocess and sends chat-completion requests to phrase nudge
|
||||||
// and reminder messages. nil ⇒ the template-based Stub is used instead.
|
// and reminder messages. nil ⇒ the template-based Stub is used instead.
|
||||||
@@ -607,6 +633,12 @@ func (c *Config) applyDefaults() {
|
|||||||
c.MemoryEval.Interval = Duration(DefaultMemoryEvalInterval)
|
c.MemoryEval.Interval = Duration(DefaultMemoryEvalInterval)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Same rule again: absent stays nil (⇒ mail ingestion refused), present gets
|
||||||
|
// the timeout default so `{}` is a valid "on with the defaults".
|
||||||
|
if c.Email != nil && c.Email.Timeout <= 0 {
|
||||||
|
c.Email.Timeout = Duration(DefaultEmailTimeout)
|
||||||
|
}
|
||||||
|
|
||||||
if c.Voice != nil {
|
if c.Voice != nil {
|
||||||
if c.Voice.RouterThreshold <= 0 {
|
if c.Voice.RouterThreshold <= 0 {
|
||||||
c.Voice.RouterThreshold = DefaultRouterThreshold
|
c.Voice.RouterThreshold = DefaultRouterThreshold
|
||||||
|
|||||||
@@ -0,0 +1,226 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/llm"
|
||||||
|
"github.com/kami/maven/internal/persona"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Extraction — turning one mail into task CANDIDATES, and nothing else.
|
||||||
|
//
|
||||||
|
// The output of this file can only ever become rows in `tasks` with status
|
||||||
|
// "candidate" (store.TaskCandidate), written through the one intake seam
|
||||||
|
// (ipc.CaptureTaskReq, Vikunja #130). That bound is the whole design:
|
||||||
|
//
|
||||||
|
// - No reminder. A reminder FIRES; it speaks to him unprompted. A 1.7B that
|
||||||
|
// misreads "встреча была в четверг" as a future appointment would then wake
|
||||||
|
// him up about it. A candidate that is wrong is a line on a review page he
|
||||||
|
// dismisses in one click, which is the correct cost of a model being wrong
|
||||||
|
// about someone's mail.
|
||||||
|
// - No fact. A fact is a claim Maven will later recite as true. Nothing read
|
||||||
|
// out of a marketing mail deserves that standing.
|
||||||
|
// - No calendar event, no note, no action. Extraction writes candidates or
|
||||||
|
// writes nothing.
|
||||||
|
//
|
||||||
|
// The due date the model may return is stored on the candidate (tasks.due_ts),
|
||||||
|
// which no scheduler reads — it is there so the review page can sort by it.
|
||||||
|
//
|
||||||
|
// Privacy: the mail text goes to the resident model on this box and nowhere
|
||||||
|
// else. It is never search input (CLAUDE.md: "his notes and facts are never
|
||||||
|
// search input" — mail is the same class), and Evidence keeps only the subject
|
||||||
|
// line, so the review page shows him where a candidate came from without the
|
||||||
|
// store growing a copy of his mailbox.
|
||||||
|
|
||||||
|
// MaxCandidates — at most this many candidates per message, enforced by the
|
||||||
|
// grammar. A mail with four tasks in it is a mail he has to read himself; a
|
||||||
|
// model allowed ten will produce ten.
|
||||||
|
const MaxCandidates = 3
|
||||||
|
|
||||||
|
// SourcePrefix — provenance for everything this package captures. The mailbox
|
||||||
|
// name is appended: "email:INBOX". Same vocabulary as tap:voice / poll:netdata.
|
||||||
|
const SourcePrefix = "email:"
|
||||||
|
|
||||||
|
// Candidate — one piece of work the model thinks the mail is asking for.
|
||||||
|
type Candidate struct {
|
||||||
|
Text string `json:"text"`
|
||||||
|
// Due — "YYYY-MM-DD" or empty. A date the model read out of the text, not a
|
||||||
|
// date it computed: relative wording ("до пятницы") is left in Text, because
|
||||||
|
// a small model resolving "пятница" against today's date gets it wrong often
|
||||||
|
// enough that a stored wrong date is worse than no date.
|
||||||
|
Due string `json:"due"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// Completer — the llama-server seam, same shape memeval and the router use, so
|
||||||
|
// the one resident model serves this caller too.
|
||||||
|
type Completer interface {
|
||||||
|
Complete(ctx context.Context, r llm.Req) (string, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Extractor reads a message and returns candidates. It holds no store and no
|
||||||
|
// writer on purpose: this type cannot persist anything, so "extraction never
|
||||||
|
// acts" is a property of the code, not of a review.
|
||||||
|
type Extractor struct {
|
||||||
|
llm Completer
|
||||||
|
// MaxCandidates — 0 ⇒ MaxCandidates.
|
||||||
|
max int
|
||||||
|
// ContextBlock — the shared persona block, optional. Extraction output is
|
||||||
|
// not spoken, so the persona matters less here than in the phraser; it is
|
||||||
|
// wired anyway so a candidate reads in her voice on the review page.
|
||||||
|
contextBlock func() string
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewExtractor(c Completer, max int, contextBlock func() string) *Extractor {
|
||||||
|
if max <= 0 || max > MaxCandidates {
|
||||||
|
max = MaxCandidates
|
||||||
|
}
|
||||||
|
return &Extractor{llm: c, max: max, contextBlock: contextBlock}
|
||||||
|
}
|
||||||
|
|
||||||
|
// extractGrammar — GBNF pinning the answer to a bounded array of fixed-shape
|
||||||
|
// candidates. Same reasoning as memeval's evalGrammar and the router's
|
||||||
|
// routeGrammar: the shape and the length bound are what keep a small model from
|
||||||
|
// drifting into prose or spending the token budget repeating one field.
|
||||||
|
//
|
||||||
|
// The empty array is reachable, deliberately: most mail contains no task, and a
|
||||||
|
// model with no way to say "nothing" invents something.
|
||||||
|
const extractGrammar = `
|
||||||
|
root ::= "[" ws (item ("," ws item){0,2})? ws "]"
|
||||||
|
item ::= "{" ws "\"text\"" ws ":" ws text "," ws "\"due\"" ws ":" ws due ws "}"
|
||||||
|
text ::= "\"" ([^"\\] | "\\" .){1,120} "\""
|
||||||
|
due ::= "\"\"" | "\"" [0-9]{4} "-" [0-9]{2} "-" [0-9]{2} "\""
|
||||||
|
ws ::= [ \t\n]*
|
||||||
|
`
|
||||||
|
|
||||||
|
// extractSystem — the extraction prompt.
|
||||||
|
//
|
||||||
|
// Written around the two failure modes a small model has on this task: it
|
||||||
|
// summarises when asked to extract (turning a mail into "письмо от Антона"),
|
||||||
|
// and it invents an obligation from any polite closing sentence. Hence the
|
||||||
|
// insistence on a verb phrase, and the explicit permission to return [].
|
||||||
|
const extractSystem = `Ты читаешь одно письмо из его почты и достаёшь из него дела, которые письмо от него требует.
|
||||||
|
|
||||||
|
Правила:
|
||||||
|
- Отвечай ТОЛЬКО массивом JSON. Каждый элемент: {"text": "...", "due": "ГГГГ-ММ-ДД" или ""}.
|
||||||
|
- text — короткая формулировка дела по-русски, с глаголом: "оплатить счёт за интернет", "отправить акт". Не пересказывай письмо и не описывай его.
|
||||||
|
- Дело — это то, что должен сделать ОН. Рассылка, реклама, уведомление, отчёт, письмо «просто к сведению» — дел не содержат.
|
||||||
|
- Если письмо ничего от него не требует, верни пустой массив []. Это нормальный ответ, так бывает чаще всего.
|
||||||
|
- Ничего не придумывай. Если срока в письме нет — "".
|
||||||
|
- due заполняй только когда в письме стоит конкретная дата. Слова вроде «до пятницы» оставь в text, дату не вычисляй.
|
||||||
|
- Максимум три дела. Лучше одно точное, чем три общих.`
|
||||||
|
|
||||||
|
// Extract returns the candidates in one message.
|
||||||
|
//
|
||||||
|
// Junk is refused without an LLM call — cheapest possible defence, and the
|
||||||
|
// reason the header filter exists. An empty message (no subject, no body) is
|
||||||
|
// likewise not worth a round trip.
|
||||||
|
//
|
||||||
|
// A parse failure is an error the caller logs and moves past. It is never
|
||||||
|
// silently turned into zero candidates, because "the model went off the rails"
|
||||||
|
// and "the mail contains no task" want different reactions from a human reading
|
||||||
|
// the log.
|
||||||
|
func (e *Extractor) Extract(ctx context.Context, msg Message) ([]Candidate, error) {
|
||||||
|
if msg.Junk {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
user := renderForModel(msg)
|
||||||
|
if user == "" {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
raw, err := e.llm.Complete(ctx, llm.Req{
|
||||||
|
System: persona.Prepend(e.contextBlock, extractSystem),
|
||||||
|
User: user,
|
||||||
|
Grammar: extractGrammar,
|
||||||
|
MaxTokens: 512,
|
||||||
|
RepeatPenalty: 1.1,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("email: extract: %w", err)
|
||||||
|
}
|
||||||
|
items, err := parseCandidates(raw)
|
||||||
|
if err != nil {
|
||||||
|
// The raw reply is NOT in the error: it is a transformation of his mail,
|
||||||
|
// and this error reaches the daemon log.
|
||||||
|
return nil, fmt.Errorf("email: extract: unparsable reply (%d bytes)", len(raw))
|
||||||
|
}
|
||||||
|
out := make([]Candidate, 0, len(items))
|
||||||
|
seen := map[string]bool{}
|
||||||
|
for _, it := range items {
|
||||||
|
it.Text = strings.TrimSpace(it.Text)
|
||||||
|
if it.Text == "" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
key := strings.ToLower(strings.Join(strings.Fields(it.Text), " "))
|
||||||
|
if seen[key] {
|
||||||
|
continue // the model repeating itself is not two tasks
|
||||||
|
}
|
||||||
|
seen[key] = true
|
||||||
|
if _, ok := ParseDue(it.Due); !ok {
|
||||||
|
it.Due = "" // a date the grammar allowed but the calendar does not
|
||||||
|
}
|
||||||
|
out = append(out, it)
|
||||||
|
if len(out) >= e.max {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// renderForModel is the user turn: subject, sender and body, labelled. Only
|
||||||
|
// these three fields — no headers, no recipient list, no message-id, nothing
|
||||||
|
// that would let the model start reasoning about routing metadata.
|
||||||
|
func renderForModel(msg Message) string {
|
||||||
|
var b strings.Builder
|
||||||
|
if msg.From != "" {
|
||||||
|
fmt.Fprintf(&b, "От: %s\n", msg.From)
|
||||||
|
}
|
||||||
|
if msg.Subject != "" {
|
||||||
|
fmt.Fprintf(&b, "Тема: %s\n", msg.Subject)
|
||||||
|
}
|
||||||
|
if msg.Body != "" {
|
||||||
|
fmt.Fprintf(&b, "\n%s\n", msg.Body)
|
||||||
|
}
|
||||||
|
if msg.Subject == "" && msg.Body == "" {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
return b.String()
|
||||||
|
}
|
||||||
|
|
||||||
|
// parseCandidates decodes the grammar-constrained reply, tolerating the
|
||||||
|
// wrappers a Thinking model sometimes leaves around it (a fenced block, or
|
||||||
|
// leading reasoning before the array).
|
||||||
|
func parseCandidates(raw string) ([]Candidate, error) {
|
||||||
|
s := strings.TrimSpace(raw)
|
||||||
|
if i := strings.Index(s, "["); i > 0 {
|
||||||
|
s = s[i:]
|
||||||
|
}
|
||||||
|
if j := strings.LastIndex(s, "]"); j >= 0 {
|
||||||
|
s = s[:j+1]
|
||||||
|
}
|
||||||
|
var out []Candidate
|
||||||
|
if err := json.Unmarshal([]byte(s), &out); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return out, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ParseDue turns the model's "YYYY-MM-DD" into a time in UTC. Exported because
|
||||||
|
// the daemon-side intake stores it on the candidate.
|
||||||
|
//
|
||||||
|
// The zero-value/empty case returns ok=false rather than an error: no date is
|
||||||
|
// the common answer, not a failure.
|
||||||
|
func ParseDue(s string) (time.Time, bool) {
|
||||||
|
s = strings.TrimSpace(s)
|
||||||
|
if s == "" {
|
||||||
|
return time.Time{}, false
|
||||||
|
}
|
||||||
|
t, err := time.Parse("2006-01-02", s)
|
||||||
|
if err != nil {
|
||||||
|
return time.Time{}, false
|
||||||
|
}
|
||||||
|
return t, true
|
||||||
|
}
|
||||||
@@ -0,0 +1,143 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/kami/maven/internal/llm"
|
||||||
|
)
|
||||||
|
|
||||||
|
// fakeLLM returns a canned reply and records the request, so a test can assert
|
||||||
|
// on the grammar and on what of the mail was sent.
|
||||||
|
type fakeLLM struct {
|
||||||
|
reply string
|
||||||
|
err error
|
||||||
|
got llm.Req
|
||||||
|
calls int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *fakeLLM) Complete(_ context.Context, r llm.Req) (string, error) {
|
||||||
|
f.calls++
|
||||||
|
f.got = r
|
||||||
|
return f.reply, f.err
|
||||||
|
}
|
||||||
|
|
||||||
|
func msgFor(subject, body string) Message {
|
||||||
|
return Message{UID: 1, From: "anton@example.org", Subject: subject, Body: body}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestExtractCandidates(t *testing.T) {
|
||||||
|
f := &fakeLLM{reply: `[{"text":"отправить акт","due":""},{"text":"оплатить счёт","due":"2026-08-05"}]`}
|
||||||
|
e := NewExtractor(f, 0, nil)
|
||||||
|
got, err := e.Extract(context.Background(), msgFor("Акт и счёт", "Надо отправить акт и оплатить счёт до 5 августа."))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("extract: %v", err)
|
||||||
|
}
|
||||||
|
if len(got) != 2 {
|
||||||
|
t.Fatalf("got %d candidates, want 2: %+v", len(got), got)
|
||||||
|
}
|
||||||
|
if got[0].Text != "отправить акт" || got[1].Due != "2026-08-05" {
|
||||||
|
t.Errorf("candidates = %+v", got)
|
||||||
|
}
|
||||||
|
if f.got.Grammar == "" {
|
||||||
|
t.Error("extraction must be grammar-constrained")
|
||||||
|
}
|
||||||
|
// The subject and body go to the model; nothing else about the message does.
|
||||||
|
if !strings.Contains(f.got.User, "Акт и счёт") || !strings.Contains(f.got.User, "оплатить счёт") {
|
||||||
|
t.Errorf("user turn = %q", f.got.User)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestExtractEmptyArrayIsNotAnError(t *testing.T) {
|
||||||
|
f := &fakeLLM{reply: "[]"}
|
||||||
|
got, err := NewExtractor(f, 0, nil).Extract(context.Background(), msgFor("FYI", "Просто к сведению."))
|
||||||
|
if err != nil || len(got) != 0 {
|
||||||
|
t.Fatalf("got (%v, %v), want (empty, nil) — no task is the normal answer", got, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Junk must never reach the model: the header filter exists so the resident
|
||||||
|
// model is not spent on newsletters.
|
||||||
|
func TestExtractSkipsJunkWithoutCallingModel(t *testing.T) {
|
||||||
|
f := &fakeLLM{reply: `[{"text":"купить всё со скидкой","due":""}]`}
|
||||||
|
msg := msgFor("Скидки", "Sale!")
|
||||||
|
msg.Junk = true
|
||||||
|
got, err := NewExtractor(f, 0, nil).Extract(context.Background(), msg)
|
||||||
|
if err != nil || got != nil {
|
||||||
|
t.Fatalf("got (%v, %v), want (nil, nil)", got, err)
|
||||||
|
}
|
||||||
|
if f.calls != 0 {
|
||||||
|
t.Errorf("model called %d times for junk, want 0", f.calls)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestExtractEmptyMessageIsNotSent(t *testing.T) {
|
||||||
|
f := &fakeLLM{reply: "[]"}
|
||||||
|
if _, err := NewExtractor(f, 0, nil).Extract(context.Background(), Message{UID: 3}); err != nil {
|
||||||
|
t.Fatalf("extract: %v", err)
|
||||||
|
}
|
||||||
|
if f.calls != 0 {
|
||||||
|
t.Errorf("model called %d times for an empty message, want 0", f.calls)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestExtractCaps(t *testing.T) {
|
||||||
|
f := &fakeLLM{reply: `[{"text":"a","due":""},{"text":"b","due":""},{"text":"c","due":""}]`}
|
||||||
|
got, err := NewExtractor(f, 2, nil).Extract(context.Background(), msgFor("s", "b"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("extract: %v", err)
|
||||||
|
}
|
||||||
|
if len(got) != 2 {
|
||||||
|
t.Errorf("got %d, want the configured cap of 2", len(got))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestExtractDropsRepeatsAndBadDates(t *testing.T) {
|
||||||
|
f := &fakeLLM{reply: `[{"text":"Отправить акт","due":"2026-02-31"},{"text":"отправить акт","due":""},{"text":" ","due":""}]`}
|
||||||
|
got, err := NewExtractor(f, 0, nil).Extract(context.Background(), msgFor("s", "b"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("extract: %v", err)
|
||||||
|
}
|
||||||
|
if len(got) != 1 {
|
||||||
|
t.Fatalf("got %d candidates, want 1 (repeat and blank dropped): %+v", len(got), got)
|
||||||
|
}
|
||||||
|
if got[0].Due != "" {
|
||||||
|
t.Errorf("due = %q, want empty — 2026-02-31 is not a date", got[0].Due)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A Thinking model sometimes wraps the array; and when it emits something
|
||||||
|
// unparsable the caller must hear about it rather than see "no tasks".
|
||||||
|
func TestParseCandidatesTolerance(t *testing.T) {
|
||||||
|
got, err := parseCandidates("думаю... [{\"text\":\"x\",\"due\":\"\"}] всё")
|
||||||
|
if err != nil || len(got) != 1 || got[0].Text != "x" {
|
||||||
|
t.Fatalf("got (%+v, %v)", got, err)
|
||||||
|
}
|
||||||
|
if _, err := parseCandidates("нет никакого JSON"); err == nil {
|
||||||
|
t.Error("unparsable output must be an error")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestExtractParseErrorHidesMailText(t *testing.T) {
|
||||||
|
f := &fakeLLM{reply: "он просил отправить акт, вот такой ответ"}
|
||||||
|
_, err := NewExtractor(f, 0, nil).Extract(context.Background(), msgFor("Акт", "секретный текст"))
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("want an error")
|
||||||
|
}
|
||||||
|
if strings.Contains(err.Error(), "акт") || strings.Contains(err.Error(), "секретный") {
|
||||||
|
t.Errorf("error text leaks mail content: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseDue(t *testing.T) {
|
||||||
|
if _, ok := ParseDue(""); ok {
|
||||||
|
t.Error("empty due must be (zero, false)")
|
||||||
|
}
|
||||||
|
if got, ok := ParseDue("2026-08-05"); !ok || got.Year() != 2026 || got.Month() != 8 || got.Day() != 5 {
|
||||||
|
t.Errorf("ParseDue = (%v, %v)", got, ok)
|
||||||
|
}
|
||||||
|
if _, ok := ParseDue("05.08.2026"); ok {
|
||||||
|
t.Error("a non-ISO date must not parse")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -26,21 +26,23 @@ type FetchSince struct {
|
|||||||
Since time.Time
|
Since time.Time
|
||||||
Max int
|
Max int
|
||||||
Skip func(uid uint32) bool
|
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
|
// 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
|
// the configuration of a mailbox and the secret for it are never the same value
|
||||||
// sitting in the same place.
|
// sitting in the same place.
|
||||||
func (f FetchSince) Run(password string) ([]Message, error) {
|
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 == "" {
|
if f.Addr == "" || f.User == "" || f.Mailbox == "" {
|
||||||
return nil, fmt.Errorf("email: mailbox not configured (addr/user/mailbox)")
|
return nil, fmt.Errorf("email: mailbox not configured (addr/user/mailbox)")
|
||||||
}
|
}
|
||||||
dial := f.dial
|
|
||||||
if dial == nil {
|
if dial == nil {
|
||||||
dial = Dial
|
dial = Dial
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,13 +21,12 @@ func TestFetchSinceRun(t *testing.T) {
|
|||||||
Since: time.Date(2026, 7, 30, 0, 0, 0, 0, time.UTC),
|
Since: time.Date(2026, 7, 30, 0, 0, 0, 0, time.UTC),
|
||||||
Max: 2,
|
Max: 2,
|
||||||
Skip: func(uid uint32) bool { return uid == 3 },
|
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 {
|
if err != nil {
|
||||||
t.Fatalf("run: %v", err)
|
t.Fatalf("run: %v", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -138,6 +138,44 @@ type CaptureTaskResp struct {
|
|||||||
Created bool `json:"created"`
|
Created bool `json:"created"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// IngestMailReq — one message a mail reader has fetched, handed to core for
|
||||||
|
// extraction (Vikunja #246).
|
||||||
|
//
|
||||||
|
// The mail reader (cmd/mavmaild) holds the IMAP credential and core never sees
|
||||||
|
// it, the same split mavpoll uses for the zenmoney token. What crosses this
|
||||||
|
// boundary is only the message text, because extraction runs on the resident
|
||||||
|
// model and llama-server lives inside core's process.
|
||||||
|
//
|
||||||
|
// Body is already plaintext and truncated by internal/email; core does not
|
||||||
|
// re-parse MIME and never stores the body. Junk means the reader's header
|
||||||
|
// filter already classified the message as bulk — core is told rather than
|
||||||
|
// asked, so a junk message can be counted without a model call.
|
||||||
|
//
|
||||||
|
// This method is available only when core has an email block configured AND a
|
||||||
|
// llama-server phraser; otherwise it answers ErrUnknownMethod, which is what
|
||||||
|
// "off unless configured" looks like at the wire.
|
||||||
|
type IngestMailReq struct {
|
||||||
|
Mailbox string `json:"mailbox"`
|
||||||
|
UID uint32 `json:"uid"`
|
||||||
|
From string `json:"from,omitempty"`
|
||||||
|
Subject string `json:"subject,omitempty"`
|
||||||
|
Date string `json:"date,omitempty"`
|
||||||
|
Body string `json:"body,omitempty"`
|
||||||
|
Junk bool `json:"junk,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// IngestMailResp — what core did with the message. TaskIDs are the rows
|
||||||
|
// CaptureTask returned; Created counts the ones that were new (a re-read
|
||||||
|
// mailbox dedupes to Created=0). Skipped is set when nothing was asked of the
|
||||||
|
// model at all — junk, or an empty message.
|
||||||
|
//
|
||||||
|
// Nothing here echoes the mail back. The reader logs counts.
|
||||||
|
type IngestMailResp struct {
|
||||||
|
TaskIDs []int64 `json:"task_ids,omitempty"`
|
||||||
|
Created int `json:"created"`
|
||||||
|
Skipped bool `json:"skipped,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
type listTasksReq struct {
|
type listTasksReq struct {
|
||||||
Status string `json:"status"` // "" all | "live" | candidate|open|done|dropped
|
Status string `json:"status"` // "" all | "live" | candidate|open|done|dropped
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -448,6 +448,17 @@ func (c *Client) SetTaskStatus(ctx context.Context, id int64, status string, ts
|
|||||||
return c.call(ctx, MethodSetTaskStatus, setTaskStatusReq{ID: id, Status: status, Ts: ts}, nil)
|
return c.call(ctx, MethodSetTaskStatus, setTaskStatusReq{ID: id, Status: status, Ts: ts}, nil)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// IngestMail hands one fetched message to core for extraction. ErrUnknownMethod
|
||||||
|
// means core has no email block configured — the caller should stop asking, not
|
||||||
|
// retry.
|
||||||
|
func (c *Client) IngestMail(ctx context.Context, req IngestMailReq) (IngestMailResp, error) {
|
||||||
|
var r IngestMailResp
|
||||||
|
if err := c.call(ctx, MethodIngestMail, req, &r); err != nil {
|
||||||
|
return IngestMailResp{}, err
|
||||||
|
}
|
||||||
|
return r, nil
|
||||||
|
}
|
||||||
|
|
||||||
func (c *Client) DismissProposedRoutine(ctx context.Context, id int64) error {
|
func (c *Client) DismissProposedRoutine(ctx context.Context, id int64) error {
|
||||||
return c.call(ctx, MethodDismissProposedRoutine, dismissProposedRoutineReq{ID: id}, nil)
|
return c.call(ctx, MethodDismissProposedRoutine, dismissProposedRoutineReq{ID: id}, nil)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -565,3 +565,36 @@ func mustJSON(v any) []byte {
|
|||||||
}
|
}
|
||||||
return b
|
return b
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestIngestMail_OffUnlessConfigured — with no IngestMailFn set (the default,
|
||||||
|
// and what an unconfigured core looks like) the method does not exist. A mail
|
||||||
|
// reader gets a refusal it can act on rather than a silent success.
|
||||||
|
func TestIngestMail_OffUnlessConfigured(t *testing.T) {
|
||||||
|
_, _, cli, _ := newServerWithStore(t)
|
||||||
|
if _, err := cli.IngestMail(context.Background(), IngestMailReq{Mailbox: "INBOX", UID: 1}); !errors.Is(err, ErrUnknownMethod) {
|
||||||
|
t.Fatalf("IngestMail error = %v, want ErrUnknownMethod", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestIngestMail_Hook — when the daemon wires the hook, the message crosses the
|
||||||
|
// boundary intact and the response comes back.
|
||||||
|
func TestIngestMail_Hook(t *testing.T) {
|
||||||
|
_, srv, cli, _ := newServerWithStore(t)
|
||||||
|
var got IngestMailReq
|
||||||
|
srv.IngestMailFn = func(_ context.Context, req IngestMailReq) (IngestMailResp, error) {
|
||||||
|
got = req
|
||||||
|
return IngestMailResp{TaskIDs: []int64{7}, Created: 1}, nil
|
||||||
|
}
|
||||||
|
resp, err := cli.IngestMail(context.Background(), IngestMailReq{
|
||||||
|
Mailbox: "INBOX", UID: 12, Subject: "Счёт", Body: "Оплатить.", Junk: false,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("IngestMail: %v", err)
|
||||||
|
}
|
||||||
|
if resp.Created != 1 || len(resp.TaskIDs) != 1 || resp.TaskIDs[0] != 7 {
|
||||||
|
t.Errorf("resp = %+v", resp)
|
||||||
|
}
|
||||||
|
if got.UID != 12 || got.Subject != "Счёт" || got.Body != "Оплатить." {
|
||||||
|
t.Errorf("req across the wire = %+v", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+33
-4
@@ -421,6 +421,17 @@ type Server struct {
|
|||||||
// Set by the daemon; nil ⇒ MethodStoreEncryptionKey returns ErrUnknownMethod.
|
// Set by the daemon; nil ⇒ MethodStoreEncryptionKey returns ErrUnknownMethod.
|
||||||
WrapKeyFn WrapKeyFunc
|
WrapKeyFn WrapKeyFunc
|
||||||
|
|
||||||
|
// IngestMailFn — extracts task candidates from one fetched message. Set by
|
||||||
|
// the daemon only when an email block is configured AND there is a
|
||||||
|
// llama-server to extract with; nil ⇒ MethodIngestMail returns
|
||||||
|
// ErrUnknownMethod, so a mail reader pointed at a core that is not
|
||||||
|
// configured for mail is refused rather than silently ignored.
|
||||||
|
//
|
||||||
|
// Like StepUp/WrapKeyFn/UnlockFn this bypasses CoreAPI: it is not a store
|
||||||
|
// operation, it needs the resident model, and it must not become a method
|
||||||
|
// every CoreAPI implementation has to carry.
|
||||||
|
IngestMailFn IngestMailFunc
|
||||||
|
|
||||||
// UnlockFn — unwraps the store encryption key from the wrapped blob using
|
// UnlockFn — unwraps the store encryption key from the wrapped blob using
|
||||||
// the passkey credential public key, opens the encrypted store, and wires
|
// the passkey credential public key, opens the encrypted store, and wires
|
||||||
// the rest of the daemon (voice, loop, delivery). Set by the daemon when
|
// the rest of the daemon (voice, loop, delivery). Set by the daemon when
|
||||||
@@ -439,6 +450,9 @@ type WrapKeyFunc func(ctx context.Context, publicKey []byte) error
|
|||||||
// public key and completes daemon initialization.
|
// public key and completes daemon initialization.
|
||||||
type UnlockFunc func(ctx context.Context, publicKey []byte) error
|
type UnlockFunc func(ctx context.Context, publicKey []byte) error
|
||||||
|
|
||||||
|
// IngestMailFunc — core-side mail extraction. Returns what was captured.
|
||||||
|
type IngestMailFunc func(ctx context.Context, req IngestMailReq) (IngestMailResp, error)
|
||||||
|
|
||||||
// CheckFunc — the auth hook signature. Wired by the daemon (auth.Gate.Check
|
// CheckFunc — the auth hook signature. Wired by the daemon (auth.Gate.Check
|
||||||
// satisfies this); dispatch calls it once per request after param-unmarshal
|
// satisfies this); dispatch calls it once per request after param-unmarshal
|
||||||
// independence (it gets the raw params, may unmarshal what it needs — ipc
|
// independence (it gets the raw params, may unmarshal what it needs — ipc
|
||||||
@@ -606,9 +620,10 @@ func withoutParams[R any](fn func(ctx context.Context, api CoreAPI) (R, error))
|
|||||||
// existed) as an argument — so SetAPI's runtime swap (the unlock transition)
|
// existed) as an argument — so SetAPI's runtime swap (the unlock transition)
|
||||||
// is still honored on the very next request with no extra plumbing here.
|
// is still honored on the very next request with no extra plumbing here.
|
||||||
//
|
//
|
||||||
// MethodAssertStepUp, MethodStoreEncryptionKey and MethodUnlock are NOT in
|
// MethodAssertStepUp, MethodStoreEncryptionKey, MethodUnlock and
|
||||||
// this table: they bypass CoreAPI entirely (s.StepUp / s.WrapKeyFn /
|
// MethodIngestMail are NOT in this table: they bypass CoreAPI entirely
|
||||||
// s.UnlockFn), so dispatch special-cases them before consulting the table.
|
// (s.StepUp / s.WrapKeyFn / s.UnlockFn / s.IngestMailFn), so dispatch
|
||||||
|
// special-cases them before consulting the table.
|
||||||
var methodTable = map[Method]handlerFunc{
|
var methodTable = map[Method]handlerFunc{
|
||||||
MethodWriteFact: withParams(func(ctx context.Context, api CoreAPI, p WriteFactReq) (idResp, error) {
|
MethodWriteFact: withParams(func(ctx context.Context, api CoreAPI, p WriteFactReq) (idResp, error) {
|
||||||
id, err := api.WriteFact(ctx, p)
|
id, err := api.WriteFact(ctx, p)
|
||||||
@@ -816,7 +831,7 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// These three bypass CoreAPI entirely — they drive Server fields set
|
// These bypass CoreAPI entirely — they drive Server fields set
|
||||||
// directly by the daemon (StepUp / WrapKeyFn / UnlockFn), not store
|
// directly by the daemon (StepUp / WrapKeyFn / UnlockFn), not store
|
||||||
// state, so they can never be table entries keyed on a CoreAPI method.
|
// state, so they can never be table entries keyed on a CoreAPI method.
|
||||||
switch req.Method {
|
switch req.Method {
|
||||||
@@ -845,6 +860,20 @@ func (s *Server) dispatch(ctx context.Context, req Request) (json.RawMessage, er
|
|||||||
return marshalResult(nil), s.UnlockFn(ctx, p.PublicKey)
|
return marshalResult(nil), s.UnlockFn(ctx, p.PublicKey)
|
||||||
}
|
}
|
||||||
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
||||||
|
|
||||||
|
case MethodIngestMail:
|
||||||
|
if s.IngestMailFn != nil {
|
||||||
|
var p IngestMailReq
|
||||||
|
if err := unmarshalParams(req.Params, &p); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
resp, err := s.IngestMailFn(ctx, p)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return marshalResult(resp), nil
|
||||||
|
}
|
||||||
|
return nil, fmt.Errorf("%w: %s", ErrUnknownMethod, req.Method)
|
||||||
}
|
}
|
||||||
|
|
||||||
h, ok := methodTable[req.Method]
|
h, ok := methodTable[req.Method]
|
||||||
|
|||||||
@@ -50,6 +50,7 @@ const (
|
|||||||
MethodCaptureTask Method = "capture_task"
|
MethodCaptureTask Method = "capture_task"
|
||||||
MethodListTasks Method = "list_tasks"
|
MethodListTasks Method = "list_tasks"
|
||||||
MethodSetTaskStatus Method = "set_task_status"
|
MethodSetTaskStatus Method = "set_task_status"
|
||||||
|
MethodIngestMail Method = "ingest_mail"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Request — one frame from module to core. Params is the JSON-encoded argument
|
// Request — one frame from module to core. Params is the JSON-encoded argument
|
||||||
|
|||||||
Reference in New Issue
Block a user