diff --git a/internal/actions/registry.go b/internal/actions/registry.go index f24d3a3..9e95441 100644 --- a/internal/actions/registry.go +++ b/internal/actions/registry.go @@ -85,24 +85,29 @@ func allActions() []Action { { Name: "read_inbox", Category: "messaging", - Description: "Check your message inbox for pending messages. Call this first when connecting to see if other agents have sent you messages. Returns unread/pending direct messages addressed to you.", + Description: "Peek at your message inbox. Idempotent and side-effect free by default — does not mark messages as read or change inbox state. Pass mark_read: true to advance the read pointer past the returned messages (legacy worker-queue behavior). To process messages with a lock, use claim_messages + mark_done instead.", Params: []Param{ {Name: "limit", Type: "number", Description: "Maximum number of messages to return (default 50)", Default: "50"}, {Name: "status_filter", Type: "string", Description: "Filter by message status: pending, processing, done, failed"}, {Name: "include_read", Type: "boolean", Description: "Include previously read messages (default false)", Default: "false"}, + {Name: "mark_read", Type: "boolean", Description: "Advance the read pointer past returned messages (default false; pure peek)", Default: "false"}, {Name: "min_priority", Type: "number", Description: "Minimum priority filter (1-10)"}, {Name: "from_agent", Type: "string", Description: "Filter by sender agent name"}, }, Returns: "JSON with messages array and count", Examples: []Example{ { - Description: "Check for new messages", + Description: "Peek at unread messages without consuming them", Code: `call("read_inbox", {})`, }, { Description: "Read high-priority messages from a specific agent", Code: `call("read_inbox", {"min_priority": 8, "from_agent": "coordinator"})`, }, + { + Description: "Legacy worker-queue: fetch unread and mark them read", + Code: `call("read_inbox", {"mark_read": true})`, + }, }, }, { diff --git a/internal/mcp/bridge.go b/internal/mcp/bridge.go index 7efa2cd..32d70ef 100644 --- a/internal/mcp/bridge.go +++ b/internal/mcp/bridge.go @@ -241,12 +241,16 @@ func (b *ServiceBridge) callSendMessage(ctx context.Context, args map[string]any } func (b *ServiceBridge) callReadInbox(ctx context.Context, args map[string]any) (any, error) { + // read_inbox is a pure peek by default. Callers that want the legacy + // worker-queue behavior (fetch unread + advance the read pointer) must + // pass mark_read: true explicitly. See bug 30674. opts := messaging.ReadOptions{ Limit: getInt(args, "limit", 50), Status: getString(args, "status_filter", ""), MinPriority: getInt(args, "min_priority", 0), FromAgent: getString(args, "from_agent", ""), IncludeRead: getBool(args, "include_read", false), + MarkRead: getBool(args, "mark_read", false), } page, err := b.msgService.ReadInbox(ctx, b.agentName, opts) diff --git a/internal/messaging/options.go b/internal/messaging/options.go index 07cacd3..b9ac49d 100644 --- a/internal/messaging/options.go +++ b/internal/messaging/options.go @@ -12,6 +12,14 @@ type SendOptions struct { } // ReadOptions configures inbox reading behavior. +// +// Read state is now an explicit, opt-in side effect. Callers that want the +// historical "fetch unread + mark them read" worker-queue behavior must set +// both IncludeRead=false (filter to unread) and MarkRead=true (advance the +// per-conversation read pointer). The default behavior is a pure peek that +// never mutates inbox state — this keeps read_inbox idempotent under retries +// and prevents it from racing with the claim/process/done loop and the +// StalemateWorker (see #bugs-synapbus message 30674). type ReadOptions struct { Status string `json:"status,omitempty"` FromAgent string `json:"from_agent,omitempty"` @@ -22,6 +30,10 @@ type ReadOptions struct { After string `json:"after,omitempty"` Before string `json:"before,omitempty"` IncludeRead bool `json:"include_read,omitempty"` + // MarkRead, when true, advances the per-conversation read pointer for the + // returned messages. When false (the default) ReadInbox is a pure peek and + // does not mutate any state. + MarkRead bool `json:"mark_read,omitempty"` } // SearchOptions configures message search behavior. diff --git a/internal/messaging/service.go b/internal/messaging/service.go index 6add5ec..161632f 100644 --- a/internal/messaging/service.go +++ b/internal/messaging/service.go @@ -247,7 +247,17 @@ func (s *MessagingService) SendMessage(ctx context.Context, from, to, body strin return msg, nil } -// ReadInbox returns messages for an agent and advances the read position. +// ReadInbox returns messages for an agent. By default this is a pure peek and +// does not mutate any state. When opts.MarkRead is true the per-conversation +// read pointer is advanced past the returned messages — this is the explicit +// opt-in for the historical "fetch unread + mark read" worker-queue flow. +// +// Splitting the read-state side effect from the fetch fixes a race against the +// claim/process/done loop and the StalemateWorker: a destructive default meant +// that `read_inbox` calls between `claim_messages` and `mark_done` could move +// the read pointer past messages that the StalemateWorker still needed to +// reason about, producing inconsistent views across `read_inbox`, `my_status`, +// and the worker's reminder/escalation queries (see bug 30674). func (s *MessagingService) ReadInbox(ctx context.Context, agentName string, opts ReadOptions) (*PaginatedMessages, error) { messages, err := s.store.GetInboxMessages(ctx, agentName, opts) if err != nil { @@ -264,33 +274,38 @@ func (s *MessagingService) ReadInbox(ctx context.Context, agentName string, opts limit = 50 } - // Advance inbox state for each conversation - conversationMaxID := make(map[int64]int64) - for _, msg := range messages { - if msg.ID > conversationMaxID[msg.ConversationID] { - conversationMaxID[msg.ConversationID] = msg.ID + if opts.MarkRead { + // Advance inbox state for each conversation, taking the max returned + // message ID per conversation as the new read watermark. + conversationMaxID := make(map[int64]int64) + for _, msg := range messages { + if msg.ID > conversationMaxID[msg.ConversationID] { + conversationMaxID[msg.ConversationID] = msg.ID + } } - } - for convID, maxMsgID := range conversationMaxID { - if err := s.store.UpdateInboxState(ctx, agentName, convID, maxMsgID); err != nil { - s.logger.Error("failed to update inbox state", - "agent", agentName, - "conversation_id", convID, - "error", err, - ) + for convID, maxMsgID := range conversationMaxID { + if err := s.store.UpdateInboxState(ctx, agentName, convID, maxMsgID); err != nil { + s.logger.Error("failed to update inbox state", + "agent", agentName, + "conversation_id", convID, + "error", err, + ) + } } } s.logger.Info("inbox read", "agent", agentName, "message_count", len(messages), + "mark_read", opts.MarkRead, ) if s.tracer != nil { s.tracer.Record(ctx, agentName, "read_inbox", map[string]any{ "message_count": len(messages), "include_read": opts.IncludeRead, + "mark_read": opts.MarkRead, }) } diff --git a/internal/messaging/service_test.go b/internal/messaging/service_test.go index fa801a2..fc366a6 100644 --- a/internal/messaging/service_test.go +++ b/internal/messaging/service_test.go @@ -222,8 +222,8 @@ func TestMessagingService_ReadInbox_ReadUnread(t *testing.T) { t.Fatalf("SendMessage: %v", err) } - // First read - result, err := svc.ReadInbox(ctx, "receiver", ReadOptions{}) + // Worker-queue path: explicit MarkRead advances the read pointer. + result, err := svc.ReadInbox(ctx, "receiver", ReadOptions{MarkRead: true}) if err != nil { t.Fatalf("ReadInbox: %v", err) } @@ -231,7 +231,7 @@ func TestMessagingService_ReadInbox_ReadUnread(t *testing.T) { t.Fatal("expected messages on first read") } - // Second read without include_read + // Second read without include_read — read pointer was advanced, so 0. result, err = svc.ReadInbox(ctx, "receiver", ReadOptions{}) if err != nil { t.Fatalf("ReadInbox: %v", err) @@ -240,7 +240,7 @@ func TestMessagingService_ReadInbox_ReadUnread(t *testing.T) { t.Errorf("got %d messages on second read (no include_read), want 0", len(result.Messages)) } - // With include_read + // With include_read — sees the marked-read messages too. result, err = svc.ReadInbox(ctx, "receiver", ReadOptions{IncludeRead: true}) if err != nil { t.Fatalf("ReadInbox: %v", err) @@ -250,6 +250,59 @@ func TestMessagingService_ReadInbox_ReadUnread(t *testing.T) { } } +// TestMessagingService_ReadInbox_NonDestructiveDefault verifies that the +// default ReadInbox call is a pure peek and does not mutate inbox read state. +// This is the bug-30674 regression test: agents calling read_inbox repeatedly +// must see the same unread messages until they explicitly opt in to MarkRead +// or use the claim/process/done loop. The destructive default raced against +// the StalemateWorker by silently advancing the per-conversation read pointer +// while a message was still being claimed and processed. +func TestMessagingService_ReadInbox_NonDestructiveDefault(t *testing.T) { + svc, _ := newTestService(t) + ctx := context.Background() + + if _, err := svc.SendMessage(ctx, "sender", "receiver", "msg 1", SendOptions{Priority: 3}); err != nil { + t.Fatalf("SendMessage: %v", err) + } + if _, err := svc.SendMessage(ctx, "sender", "receiver", "msg 2", SendOptions{Priority: 8}); err != nil { + t.Fatalf("SendMessage: %v", err) + } + + cases := []struct { + name string + opts ReadOptions + }{ + {"default options", ReadOptions{}}, + {"with limit", ReadOptions{Limit: 10}}, + {"with from_agent filter", ReadOptions{FromAgent: "sender"}}, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + first, err := svc.ReadInbox(ctx, "receiver", tc.opts) + if err != nil { + t.Fatalf("first ReadInbox: %v", err) + } + if len(first.Messages) != 2 { + t.Fatalf("first read returned %d messages, want 2", len(first.Messages)) + } + + second, err := svc.ReadInbox(ctx, "receiver", tc.opts) + if err != nil { + t.Fatalf("second ReadInbox: %v", err) + } + if len(second.Messages) != len(first.Messages) { + t.Errorf("second read returned %d messages, want %d (read_inbox must be idempotent by default)", + len(second.Messages), len(first.Messages)) + } + if second.Total != first.Total { + t.Errorf("second read total = %d, want %d (totals must match across idempotent calls)", + second.Total, first.Total) + } + }) + } +} + func TestMessagingService_ReadInbox_Filters(t *testing.T) { svc, _ := newTestService(t) ctx := context.Background() diff --git a/internal/messaging/stalemate.go b/internal/messaging/stalemate.go index 39752ad..5fc5fdb 100644 --- a/internal/messaging/stalemate.go +++ b/internal/messaging/stalemate.go @@ -221,11 +221,17 @@ func (w *StalemateWorker) failTimedOutProcessing(ctx context.Context) int64 { metadata := map[string]any{"error": "claim timeout exceeded"} metaBytes, _ := json.Marshal(metadata) - // Update directly via DB since the store's UpdateMessageStatus requires the claiming agent - _, err := w.db.ExecContext(ctx, + // Update directly via DB since the store's UpdateMessageStatus requires + // the claiming agent. The UPDATE re-checks both status='processing' AND + // claimed_at < cutoff to close a TOCTOU window between the SELECT above + // and this UPDATE: if a legitimate claimer ran mark_done (status flip) + // or a fresh claim_messages refreshed claimed_at between the two, this + // no-ops instead of stomping live work. RowsAffected = 0 is the signal + // the row escaped the stale window before the worker reached it. + res, err := w.db.ExecContext(ctx, `UPDATE messages SET status = ?, metadata = ?, updated_at = CURRENT_TIMESTAMP - WHERE id = ? AND status = 'processing'`, - StatusFailed, string(metaBytes), dm.ID, + WHERE id = ? AND status = 'processing' AND claimed_at < ?`, + StatusFailed, string(metaBytes), dm.ID, cutoff, ) if err != nil { w.logger.Error("auto-fail message failed", @@ -234,6 +240,17 @@ func (w *StalemateWorker) failTimedOutProcessing(ctx context.Context) int64 { ) continue } + affected, _ := res.RowsAffected() + if affected == 0 { + // Message was reclaimed, completed, or otherwise moved out of the + // stale window between SELECT and UPDATE. Skip silently — the + // next worker tick will re-evaluate. + w.logger.Debug("stale processing message no longer stale; skipped", + "message_id", dm.ID, + "claimed_by", dm.ClaimedBy, + ) + continue + } w.logger.Info("auto-failed stale processing message", "message_id", dm.ID, "from_agent", dm.FromAgent, diff --git a/internal/messaging/stalemate_test.go b/internal/messaging/stalemate_test.go index 9ca1778..c248757 100644 --- a/internal/messaging/stalemate_test.go +++ b/internal/messaging/stalemate_test.go @@ -136,6 +136,50 @@ func TestStalemateWorker_ProcessingTimeout_NotExpired(t *testing.T) { } } +// TestStalemateWorker_ProcessingTimeout_RaceGuard verifies that the +// auto-fail UPDATE re-checks claimed_at < cutoff and won't stomp a row that +// was legitimately re-claimed (claimed_at refreshed) between the worker's +// SELECT scan and its row-by-row UPDATE. This guards the TOCTOU window +// the stale-worker race depends on. +func TestStalemateWorker_ProcessingTimeout_RaceGuard(t *testing.T) { + svc, db := newStalemateTestService(t) + ctx := context.Background() + + // Insert a message that *was* stale at SELECT time. + oldClaimedAt := time.Now().Add(-25 * time.Hour) + msgID := insertStaleMessage(t, db, "sender", "receiver", "racing task", + StatusProcessing, time.Now().Add(-26*time.Hour), &oldClaimedAt, "receiver") + + config := DefaultStalemateConfig() + config.ProcessingTimeout = 24 * time.Hour + lookup := &stubChannelLookup{channelID: 0, err: fmt.Errorf("no channel")} + worker := NewStalemateWorker(db, svc, lookup, config) + + // Simulate the race: between the worker's SELECT (which would have picked + // this row) and its UPDATE, the legitimate claimer refreshes claimed_at to + // "now". With the cutoff predicate in place the UPDATE no-ops instead of + // silently failing live work. + freshClaimedAt := time.Now() + if _, err := db.ExecContext(ctx, + `UPDATE messages SET claimed_at = ? WHERE id = ?`, + freshClaimedAt, msgID, + ); err != nil { + t.Fatalf("refresh claimed_at: %v", err) + } + + worker.checkStaleMessages(ctx) + + var status string + if err := db.QueryRowContext(ctx, + `SELECT status FROM messages WHERE id = ?`, msgID, + ).Scan(&status); err != nil { + t.Fatalf("query message: %v", err) + } + if status != StatusProcessing { + t.Errorf("status = %q, want %q — re-claimed message must not be auto-failed", status, StatusProcessing) + } +} + func TestStalemateWorker_PendingReminder(t *testing.T) { svc, db := newStalemateTestService(t) ctx := context.Background()