Files
kami 3588da9e28 Merge branch 'fix/g06' into fix/integrated
# Conflicts:
#	cmd/mavend/memoryeval.go
2026-08-01 14:20:04 +04:00

242 lines
9.8 KiB
Go

// 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"
"strings"
"time"
"unicode"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/email"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"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
// captureTimeout — how long the capture writes get, separately from the
// extraction budget. A candidate the model already produced must not be lost
// because the model was slow.
const captureTimeout = 30 * time.Second
// maxMailboxChars — a mailbox name is an IMAP folder, not free text. It ends up
// in the provenance string, which is a small controlled vocabulary.
const maxMailboxChars = 64
// validMailbox checks the name this method is willing to write provenance for.
// Empty is refused: "email:" is not a source. So is anything with a control
// character or a space-only value, so the source string stays greppable and
// stays one token.
func validMailbox(s string) (string, error) {
s = strings.TrimSpace(s)
if s == "" {
return "", fmt.Errorf("mail intake: mailbox is required")
}
if len([]rune(s)) > maxMailboxChars {
return "", fmt.Errorf("mail intake: mailbox name too long")
}
for _, r := range s {
if r < 0x20 || r == 0x7f || unicode.IsSpace(r) {
return "", fmt.Errorf("mail intake: mailbox name has whitespace or a control character")
}
}
return s, nil
}
// 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
// bus — the unified intake journal (Vikunja #283). This path captures
// through the store directly rather than through ipc.CoreAPI, so the
// decorator in intake.go does not see it and the publish is explicit here.
// nil is a working no-op.
bus *event.Bus
}
// 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, bus *event.Bus) *mailIntake {
if cfg.Email == nil {
return nil
}
lp, ok := phr.(*phraser.LLMPhraser)
if !ok {
// The phraser is not an *LLMPhraser. Today that means there is no
// llama-server; if anything ever WRAPS the phraser it will mean that
// instead, so the line names the assertion rather than guessing why.
log.Printf("mail intake: configured but the phraser is not an *phraser.LLMPhraser (%T) — mail ingestion disabled", phr)
return nil
}
timeout := time.Duration(cfg.Email.Timeout)
if timeout <= 0 {
timeout = config.DefaultEmailTimeout
}
// Background client: extraction is a job nobody is waiting on, and it shares
// one llama-server slot with the voice turn. Through the gate it yields to
// anything he is waiting for and only one extraction runs at a time, so a
// first poll of 25 unseen messages cannot queue 25 model calls in front of
// him. See llm.Gate.
ex := email.NewExtractor(llmBackgroundClientFor(lp, timeout), cfg.Email.MaxTasks, contextBlockFn(cfg, time.Now))
// The NORMALISED bound, not the configured one: with "email": {} in
// mavend.json the configured value is 0 and the daemon allows three.
log.Printf("mail intake: enabled (max %d candidates per message, timeout %s)", ex.Max(), timeout)
return &mailIntake{st: st, ex: ex, timeout: timeout, now: time.Now, bus: bus}
}
// 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, under an ExternalID naming the message and the
// span it was extracted from. That key is unique over every row whatever its
// status, so a mailbox re-read after a restart produces Created=0 — and, more
// to the point, a task he already marked done is not re-proposed the next time
// the same unread message is read again.
func (m *mailIntake) ingest(ctx context.Context, req ipc.IngestMailReq) (ipc.IngestMailResp, error) {
// The mailbox name becomes provenance ("email:INBOX"), and the source
// vocabulary is what the loop's rules trust. An empty name gave "email:" and
// an arbitrary string gave an arbitrary source under that namespace.
mailbox, err := validMailbox(req.Mailbox)
if err != nil {
return ipc.IngestMailResp{}, err
}
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
}
// The timeout scopes the EXTRACTION and nothing else. It used to wrap the
// capture writes too, so a model that answered at 119 seconds of a 120
// second budget left the first CaptureTask one second and the third none:
// the work was done, the answer was good, and it was dropped with a
// deadline error. Config calls this a per-message extraction budget, and now
// it is one.
exCtx, cancel := context.WithTimeout(ctx, m.timeout)
cands, err := m.ex.Extract(exCtx, msg)
cancel()
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
}
// A fresh budget for the writes, derived from the caller's context rather
// than from the extraction's. Encrypted-store writes are fast; what this
// bounds is a stuck store, not the model.
ctx, cancel = context.WithTimeout(ctx, captureTimeout)
defer cancel()
source := email.SourcePrefix + 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,
}
t.ExternalID = mailExternalID(source, req.UID, c.Text)
if due, ok := email.ParseDue(c.Due); ok {
t.Due = &due
}
res, err := m.st.CaptureTask(ctx, t)
if err != nil {
return resp, fmt.Errorf("mail intake: capture: %w", err)
}
resp.TaskIDs = append(resp.TaskIDs, res.ID)
if res.Created {
resp.Created++
// Only a row that was actually created. CaptureTask dedupes on
// normalised text among live rows, so a mailbox re-read after a
// restart must not refill the journal with tasks already in it.
m.bus.Publish(publishableTask(t, now), now)
}
}
// 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, bus *event.Bus) {
mi := newMailIntake(st, phr, cfg, bus)
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]) + "…"
}
// mailExternalID names the message and the span a candidate was extracted
// from. The mailbox and UID identify the message; the normalised text
// identifies which of the candidates in it this is, so a message yielding two
// tasks gets two keys and a re-read of it gets neither twice.
//
// UIDs are stable per mailbox, and a mailbox that renumbers (UIDVALIDITY
// changing) re-proposes its tasks once, which is the safe direction.
func mailExternalID(source string, uid uint32, text string) string {
return fmt.Sprintf("%s#%d:%s", source, uid, store.NormalizeTaskText(text))
}