Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 0 additions & 13 deletions agent/ask.go
Original file line number Diff line number Diff line change
Expand Up @@ -469,19 +469,6 @@ func History(account, threadID string, max int) []QueryMessage {
return out
}

// Sent records the id a client's own protocol gave an answer, so a reply to it
// finds this conversation. Mail's Message-ID; nothing for a client without one.
func Sent(accountID, threadID, ref string) {
if threadID == "" || strings.TrimSpace(ref) == "" {
return
}
for _, m := range thread.Messages(accountID, threadID, 1) {
if m.Role == thread.RoleAgent {
thread.SetRef(accountID, m.ID, ref)
}
}
}

func threadID(th *thread.Thread) string {
if th == nil {
return ""
Expand Down
23 changes: 9 additions & 14 deletions agent/mail/mail.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ import (
"strings"
"time"

"github.com/google/uuid"

"mu/agent"
"mu/internal/app"
"mu/internal/event"
Expand Down Expand Up @@ -231,10 +233,9 @@ func answerMail(m mail.InboundMail) {
// the account has to be able to see that it happened.
via := agent.Via{From: m.From}

// The conversation the answer will be recorded on, filled in once the run
// has happened. deliver closes over it so it can note the id the reply went
// out under, which is what the next message in the thread looks for.
var threadRef string
// Reserve the mail identity before recording the answer. The independent
// Inbox consumer can then project delivery without adding another turn.
answerRef := "<" + uuid.NewString() + ".reply@" + domain + ">"

record := func(prompt, answer string, err error) string {
return agent.Record(agent.Recorded{
Expand Down Expand Up @@ -284,15 +285,9 @@ func answerMail(m mail.InboundMail) {
// every message from somebody writing to their own agent.
to, cc := replyTo(m)
plain = introduction(m.Owner, m, from) + plain
sent, err := mail.SendReplyAll(m.Owner, name, from, to, cc, subject,
plain, app.RenderString(plain), m.MessageID, m.References)
// The id the answer went out under, so the reply to *it* finds this
// turn. Recorded even when delivery failed below, because a message
// that reached the far side and then errored still gets answered.
// Against the conversation, which is what the next message looks
// in — the workflow record is how this answer was produced, not what
// was said.
agent.Sent(m.Owner, threadRef, sent)
_, err := mail.SendReplyAll(m.Owner, name, from, to, cc, subject,
plain, app.RenderString(plain), m.MessageID, m.References, answerRef)
Comment on lines +288 to +289

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve the fallback reply identity after model errors

When QueryWithOpts returns an error after producing partial output, Ask deliberately clears AnswerRef before recording that partial agent turn, but this call sends the replacement failure notice under the reserved answerRef. The Inbox arrival therefore cannot deduplicate against the recorded agent turn: the stale partial answer remains and the delivered notice is added as a separate RolePerson message from the agent address (and may open a separate projected thread). The removed agent.Sent path previously attached the actual sent ID to that partial turn, so the fallback needs to be recorded as the agent answer with the reserved identity before delivery.

AGENTS.md reference: AGENTS.md:L291-L300

Useful? React with 👍 / 👎.


if err != nil {
app.Log("mail", "agent %s could not reply to %s: %v", name, m.From, err)
// Recorded as an error against the run, because a reply that
Expand Down Expand Up @@ -371,12 +366,12 @@ func answerMail(m mail.InboundMail) {
As: from,
Ref: m.InReplyTo + " " + m.References,
MessageRef: m.MessageID,
AnswerRef: answerRef,
From: m.From,
FromName: m.FromName,
Via: via,
})
answer := res.Text
threadRef = res.Thread
if err != nil {
app.Log("mail", "agent %s failed on mail from %s: %v", name, m.From, err)
deliver(res.Flow, prompt, "I could not answer that one. Try again, or ask a different way.")
Expand Down
10 changes: 10 additions & 0 deletions inbox/imapbridge.go
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,16 @@ func asMessages(accountID string, t thread.Thread, domain string) []*mail.Messag
}
}

// Mail answers carry their sending address in From (AskRequest.As).
// They belong to the mail transport even before delivery is indexed;
// projecting them here exposes a second .conversation sender to IMAP.
if m.Role == thread.RoleAgent && m.Workflow != "" && strings.HasSuffix(strings.ToLower(m.From), "@"+strings.ToLower(domain)) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Suppress preflight mail replies before native delivery

When agent.Ask exits through safety.NeverAllowed or the insufficient-balance check, it records the reply with the reserved AnswerRef and local From, but with an empty Workflow. Because native delivery happens only after Ask returns, an IMAP fetch in that interval misses the lookup at lines 110–118 and also fails this new condition, assigning a UID to a synthetic .conversation message; after delivery, that message disappears and a different native-mail ID/UID appears, recreating the duplicate/expunge race this change is intended to eliminate. Identify reserved mail answers without requiring a workflow.

AGENTS.md reference: AGENTS.md:L306-L310

Useful? React with 👍 / 👎.

if strings.HasPrefix(m.Ref, "<") && strings.HasSuffix(m.Ref, ">") {
prev = m.Ref
}
continue
}

one := &mail.Message{
Bridged: true,
ID: id,
Expand Down
30 changes: 30 additions & 0 deletions inbox/imapbridge_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,3 +35,33 @@ func TestCheckinBridgeSkipsNativeMailAndKeepsReferences(t *testing.T) {
t.Fatalf("bad bridge: %+v", got)
}
}

func TestCheckinMailAnswerIsNotBridgedBeforeOrAfterDelivery(t *testing.T) {
const owner = "checkin-mail-answer"
defer thread.Forget(owner)
th := thread.Open(owner, thread.WebClient, "checkin:mail-answer")
from := "agent@" + mail.ConfiguredDomain()
ref := "<checkin-answer@test>"
answer := thread.Message{Account: owner, Thread: th.ID, Role: thread.RoleAgent, Text: "One answer", From: from, Ref: ref, Workflow: "mail-run"}
id := thread.Add(answer)
if got := Bridge(owner); len(got) != 0 {
t.Fatalf("mail answer bridged before delivery: %+v", got)
}
if err := mail.SendMessageTo(mail.Delivery{FromID: from, ToID: owner, Subject: "Re: Daily Checkin", Body: "One answer", MessageID: ref}); err != nil {
t.Fatal(err)
}
// The independent arrival consumer may run before the mail responder returns.
if added := thread.Add(thread.Message{Account: owner, Thread: th.ID, Role: thread.RolePerson, Text: "One answer", From: from, Ref: ref, Workflow: "mail-run"}); added != id {
t.Fatal("delivery recorded a second conversation turn")
}
if got := Bridge(owner); len(got) != 0 {
t.Fatalf("mail answer bridged after delivery: %+v", got)
}
// Old answers whose delivery raced the late reference update are suppressed too.
thread.Add(thread.Message{Account: owner, Thread: th.ID, Role: thread.RoleAgent, Text: "Older mail answer", From: from, Workflow: "older-mail-run"})
thread.Add(thread.Message{Account: owner, Thread: th.ID, Role: thread.RoleAgent, Text: "One answer"})
got := Bridge(owner)
if len(got) != 1 || got[0].Body != "One answer" || got[0].InReplyTo != ref {
t.Fatalf("lost genuine web answer or mail references: %+v", got)
}
}
19 changes: 19 additions & 0 deletions service/mail/checkin_imap_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,3 +46,22 @@ func TestLocalSubmissionKeepsClientMessageID(t *testing.T) {
t.Fatalf("client identity lost: %+v", stored)
}
}

func TestMailReplyUsesReservedIdentity(t *testing.T) {
const ref = "<reserved-reply@example.test>"
body, id := buildExternalTo("Micro", "agent@example.test", "", "person@outside.test", nil, "Re: Checkin", "Answer", "<p>Answer</p>", "<question@example.test>", "", ref)
if id != ref || !strings.Contains(string(body), "Message-ID: "+ref+"\r\n") {
t.Fatal("outbound reply replaced reserved identity")
}
const owner = "reserved-reply-owner"
auth.SetAccountForTest(&auth.Account{ID: owner, Approved: true})
defer auth.RemoveAccountForTest(owner)
got, err := SendReplyAll(owner, "Micro", "agent@"+ConfiguredDomain(), owner, nil, "Re: Checkin", "Answer", "<p>Answer</p>", "<question@example.test>", "", ref)
if err != nil {
t.Fatal(err)
}
stored := FindMessageByMessageID(ref)
if got != ref || stored == nil || stored.ToID != owner {
t.Fatalf("local reply replaced reserved identity: %q %+v", got, stored)
}
}
17 changes: 12 additions & 5 deletions service/mail/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,7 @@ func buildExternal(displayName, from, replyTo, to, subject, bodyPlain, bodyHTML
return buildExternalTo(displayName, from, replyTo, to, nil, subject, bodyPlain, bodyHTML, replyToMsgID, references)
}

func buildExternalTo(displayName, from, replyTo, to string, cc []string, subject, bodyPlain, bodyHTML string, replyToMsgID, references string) ([]byte, string) {
func buildExternalTo(displayName, from, replyTo, to string, cc []string, subject, bodyPlain, bodyHTML string, replyToMsgID, references string, identity ...string) ([]byte, string) {
// Extract username from email for Message-ID
username := from
if strings.Contains(from, "@") {
Expand All @@ -200,6 +200,9 @@ func buildExternalTo(displayName, from, replyTo, to string, cc []string, subject

// Generate unique Message-ID for threading
messageID := fmt.Sprintf("<%d.%s@%s>", time.Now().UnixNano(), username, ConfiguredDomain())
if len(identity) > 0 && identity[0] != "" && !strings.ContainsAny(identity[0], "\r\n") {
messageID = identity[0]
}

// Generate boundary for multipart
boundary := fmt.Sprintf("----=_Part_%d", time.Now().UnixNano())
Expand Down Expand Up @@ -499,7 +502,10 @@ func DKIMStatus() (enabled bool, domain, selector string) {
// and it asks per recipient: a thread with one local person and one outside is
// both, which is the case a single branch at the call site would get wrong.
func SendReplyAll(fromID, displayName, from, to string, cc []string, subject, bodyPlain, bodyHTML,
inReplyTo, references string) (string, error) {
inReplyTo, references, messageID string) (string, error) {
if strings.ContainsAny(messageID, "\r\n") {
return "", fmt.Errorf("invalid Message-ID")
}

var outside []string
var here []string
Expand All @@ -522,15 +528,16 @@ func SendReplyAll(fromID, displayName, from, to string, cc []string, subject, bo
// this instance are still owed their copy — the relay being down is not
// their problem, and returning early meant one bad address on a thread
// silenced the answer for everybody on it.
var messageID string
var relayErr error
if len(outside) > 0 {
id, err := queueReply(fromID, displayName, from, outside[0], outside[1:], subject,
bodyPlain, bodyHTML, inReplyTo, references)
bodyPlain, bodyHTML, inReplyTo, references, messageID)
if err != nil {
relayErr = err
}
messageID = id
if id != "" {
messageID = id
}
}

// And everybody here, delivered rather than relayed. Each on their own,
Expand Down
4 changes: 2 additions & 2 deletions service/mail/outbox.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,8 @@ type queuedMail struct {
LastError string `json:"last_error,omitempty"`
}

func queueReply(owner, display, from, to string, cc []string, subject, plain, html, parent, refs string) (string, error) {
message, id := buildExternalTo(display, from, "", to, cc, subject, plain, html, parent, refs)
func queueReply(owner, display, from, to string, cc []string, subject, plain, html, parent, refs string, identity ...string) (string, error) {
message, id := buildExternalTo(display, from, "", to, cc, subject, plain, html, parent, refs, identity...)
return enqueueMail(owner, queuedMail{From: from, Subject: subject, MessageID: id,
Message: signExternal(message), Recipients: append([]string{to}, cc...)})
}
Expand Down
Loading