fix(messaging): non-destructive read_inbox + stalemate UPDATE race guard

read_inbox now requires explicit MarkRead (default false). Worker-queue callers
opt in. Resolves bugs-synapbus #30674 where consecutive identical calls returned
0 the second time and produced inconsistent views with the claim/process/done
loop and StalemateWorker.

failTimedOutProcessing UPDATE now re-checks claimed_at < cutoff so a fresh
re-claim between SELECT and UPDATE can't be stomped to failed.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
Algis Dumbris
2026-05-08 11:39:40 +03:00
co-authored by Claude Opus 4.7
parent 121d05f875
commit 8814e87e16
7 changed files with 174 additions and 24 deletions
+7 -2
View File
@@ -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})`,
},
},
},
{
+4
View File
@@ -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)
+12
View File
@@ -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.
+29 -14
View File
@@ -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,
})
}
+57 -4
View File
@@ -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()
+21 -4
View File
@@ -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,
+44
View File
@@ -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()