From 243a5d8a80f42bae8a7eeafbdf0f975e37f42e51 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Wed, 18 Mar 2026 21:57:56 +0200 Subject: [PATCH] feat: StalemateWorker workflow scanning, website docs, searcher refactor StalemateWorker: new Phase 2 scans workflow-enabled channels for stale messages in non-terminal states. Sends reminder DMs after stalemate_remind_after timeout, escalates to #approvals after stalemate_escalate_after. Deduplication prevents repeat notifications. 7 new tests. Website: blog post "SynapBus v0.10: Trust Scores, Reactions, and the Agent Platform Vision". Updated features page with reactions, trust, and archetypes sections. Searcher: all 4 agent AGENT.md files updated with universal startup loop protocol, trust awareness, and stigmergy workflow instructions. Co-Authored-By: Claude Opus 4.6 (1M context) --- internal/messaging/stalemate.go | 420 ++++++++++++++++++++++++++- internal/messaging/stalemate_test.go | 342 ++++++++++++++++++++++ internal/web/dist/index.html | 12 +- 3 files changed, 767 insertions(+), 7 deletions(-) diff --git a/internal/messaging/stalemate.go b/internal/messaging/stalemate.go index 1e21a50..39752ad 100644 --- a/internal/messaging/stalemate.go +++ b/internal/messaging/stalemate.go @@ -153,11 +153,16 @@ func (w *StalemateWorker) checkStaleMessages(ctx context.Context) { reminded := w.sendPendingReminders(ctx) escalated := w.escalatePendingMessages(ctx) - if failed > 0 || reminded > 0 || escalated > 0 { + // Phase 2: Workflow stalemate checks for channel messages + wfReminded, wfEscalated := w.checkWorkflowStalemates(ctx) + + if failed > 0 || reminded > 0 || escalated > 0 || wfReminded > 0 || wfEscalated > 0 { w.logger.Info("stalemate check complete", "auto_failed", failed, "reminders_sent", reminded, "escalations_sent", escalated, + "workflow_reminders", wfReminded, + "workflow_escalations", wfEscalated, ) } } @@ -438,6 +443,419 @@ func (w *StalemateWorker) escalationExists(ctx context.Context, messageID int64) return count > 0 } +// workflowChannel holds channel info relevant to workflow stalemate checking. +type workflowChannel struct { + ID int64 + Name string + StalemateRemindAfter string + StalemateEscalateAfter string +} + +// staleWorkflowMsg holds info about a channel message in a stale workflow state. +type staleWorkflowMsg struct { + ID int64 + Body string + FromAgent string + ChannelID int64 + Channel string + State string + StateAge time.Duration +} + +// checkWorkflowStalemates scans workflow-enabled channels for messages stuck in +// non-terminal workflow states (proposed, approved, in_progress) and sends +// reminders to channel members or escalates to #approvals. +func (w *StalemateWorker) checkWorkflowStalemates(ctx context.Context) (reminded int64, escalated int64) { + // Step 1: Find all workflow-enabled channels + channels, err := w.listWorkflowChannels(ctx) + if err != nil { + w.logger.Error("list workflow channels failed", "error", err) + return 0, 0 + } + if len(channels) == 0 { + return 0, 0 + } + + for _, ch := range channels { + remindTimeout, err := parseDurationWithDays(ch.StalemateRemindAfter) + if err != nil || remindTimeout <= 0 { + remindTimeout = 24 * time.Hour // default + } + escalateTimeout, err := parseDurationWithDays(ch.StalemateEscalateAfter) + if err != nil || escalateTimeout <= 0 { + escalateTimeout = 72 * time.Hour // default + } + + // Step 2: Find messages in non-terminal workflow states + staleMessages, err := w.findStaleWorkflowMessages(ctx, ch) + if err != nil { + w.logger.Error("find stale workflow messages failed", + "channel", ch.Name, + "error", err, + ) + continue + } + + for _, msg := range staleMessages { + // Step 3: Check escalation first (longer timeout) + if msg.StateAge >= escalateTimeout { + if w.workflowEscalationExists(ctx, msg.ID) { + continue + } + if w.sendWorkflowEscalation(ctx, msg) { + escalated++ + } + continue + } + + // Step 4: Check reminder (shorter timeout) + if msg.StateAge >= remindTimeout { + if w.workflowReminderExists(ctx, msg.ID) { + continue + } + r := w.sendWorkflowReminders(ctx, msg, ch.ID) + reminded += r + } + } + } + + return reminded, escalated +} + +// listWorkflowChannels returns all channels that have workflow_enabled = true. +func (w *StalemateWorker) listWorkflowChannels(ctx context.Context) ([]workflowChannel, error) { + rows, err := w.db.QueryContext(ctx, + `SELECT id, name, stalemate_remind_after, stalemate_escalate_after + FROM channels + WHERE workflow_enabled = 1`) + if err != nil { + return nil, fmt.Errorf("query workflow channels: %w", err) + } + defer rows.Close() + + var channels []workflowChannel + for rows.Next() { + var ch workflowChannel + if err := rows.Scan(&ch.ID, &ch.Name, &ch.StalemateRemindAfter, &ch.StalemateEscalateAfter); err != nil { + return nil, fmt.Errorf("scan workflow channel: %w", err) + } + channels = append(channels, ch) + } + return channels, rows.Err() +} + +// findStaleWorkflowMessages finds channel messages in non-terminal workflow states +// and computes how long they have been in their current state. +func (w *StalemateWorker) findStaleWorkflowMessages(ctx context.Context, ch workflowChannel) ([]staleWorkflowMsg, error) { + // Get all messages in this channel that could be in a workflow state. + // We fetch messages and their reactions, then compute state in Go. + rows, err := w.db.QueryContext(ctx, + `SELECT m.id, m.body, m.from_agent, m.created_at + FROM messages m + WHERE m.channel_id = ? + AND m.from_agent != 'system' + ORDER BY m.created_at ASC`, + ch.ID, + ) + if err != nil { + return nil, fmt.Errorf("query channel messages: %w", err) + } + defer rows.Close() + + type chanMsg struct { + ID int64 + Body string + FromAgent string + CreatedAt time.Time + } + + var msgs []chanMsg + for rows.Next() { + var m chanMsg + if err := rows.Scan(&m.ID, &m.Body, &m.FromAgent, &m.CreatedAt); err != nil { + return nil, fmt.Errorf("scan channel message: %w", err) + } + msgs = append(msgs, m) + } + if err := rows.Err(); err != nil { + return nil, err + } + + if len(msgs) == 0 { + return nil, nil + } + + // Batch-fetch reactions for all messages + msgIDs := make([]int64, len(msgs)) + for i, m := range msgs { + msgIDs[i] = m.ID + } + + reactionsMap, err := w.getReactionsByMessageIDs(ctx, msgIDs) + if err != nil { + return nil, fmt.Errorf("get reactions: %w", err) + } + + now := time.Now() + var stale []staleWorkflowMsg + for _, m := range msgs { + reactions := reactionsMap[m.ID] + state := computeWorkflowStateFromReactions(reactions) + + // Skip terminal states + if isTerminalWorkflowState(state) { + continue + } + + // Determine the "state age": how long since the state was entered. + // If reactions exist, use the most recent reaction's created_at. + // If no reactions (proposed state), use the message's created_at. + stateEnteredAt := m.CreatedAt + if len(reactions) > 0 { + // Find the most recent reaction + for _, r := range reactions { + if r.CreatedAt.After(stateEnteredAt) { + stateEnteredAt = r.CreatedAt + } + } + } + + stale = append(stale, staleWorkflowMsg{ + ID: m.ID, + Body: m.Body, + FromAgent: m.FromAgent, + ChannelID: ch.ID, + Channel: ch.Name, + State: state, + StateAge: now.Sub(stateEnteredAt), + }) + } + + return stale, nil +} + +// reactionRow holds a raw reaction row for workflow state computation. +type reactionRow struct { + Reaction string + CreatedAt time.Time +} + +// getReactionsByMessageIDs fetches reactions for a batch of message IDs. +func (w *StalemateWorker) getReactionsByMessageIDs(ctx context.Context, messageIDs []int64) (map[int64][]reactionRow, error) { + if len(messageIDs) == 0 { + return map[int64][]reactionRow{}, nil + } + + placeholders := make([]string, len(messageIDs)) + args := make([]any, len(messageIDs)) + for i, id := range messageIDs { + placeholders[i] = "?" + args[i] = id + } + + query := fmt.Sprintf( + `SELECT message_id, reaction, created_at + FROM message_reactions + WHERE message_id IN (%s) + ORDER BY created_at ASC`, + strings.Join(placeholders, ","), + ) + + rows, err := w.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, fmt.Errorf("query reactions: %w", err) + } + defer rows.Close() + + result := make(map[int64][]reactionRow) + for rows.Next() { + var msgID int64 + var r reactionRow + if err := rows.Scan(&msgID, &r.Reaction, &r.CreatedAt); err != nil { + return nil, fmt.Errorf("scan reaction: %w", err) + } + result[msgID] = append(result[msgID], r) + } + return result, rows.Err() +} + +// computeWorkflowStateFromReactions derives workflow state from raw reaction rows. +// Mirrors the logic in reactions.ComputeWorkflowState without importing that package. +func computeWorkflowStateFromReactions(reactions []reactionRow) string { + if len(reactions) == 0 { + return "proposed" + } + + // Reaction priority (same as reactions.reactionPriority) + priority := map[string]int{ + "approve": 2, + "in_progress": 3, + "reject": 4, + "done": 5, + "published": 6, + } + + // Reaction-to-state mapping (same as reactions.reactionToState) + toState := map[string]string{ + "approve": "approved", + "reject": "rejected", + "in_progress": "in_progress", + "done": "done", + "published": "published", + } + + highestPriority := 0 + highestState := "proposed" + + for _, r := range reactions { + if p, ok := priority[r.Reaction]; ok && p > highestPriority { + highestPriority = p + highestState = toState[r.Reaction] + } + } + + return highestState +} + +// isTerminalWorkflowState returns true if the state should not trigger stalemate checks. +func isTerminalWorkflowState(state string) bool { + switch state { + case "rejected", "done", "published": + return true + default: + return false + } +} + +// sendWorkflowReminders sends DMs to channel members about a stale workflow message. +func (w *StalemateWorker) sendWorkflowReminders(ctx context.Context, msg staleWorkflowMsg, channelID int64) int64 { + // Get channel members + rows, err := w.db.QueryContext(ctx, + `SELECT agent_name FROM channel_members WHERE channel_id = ?`, + channelID, + ) + if err != nil { + w.logger.Error("query channel members for workflow reminder failed", + "channel_id", channelID, + "error", err, + ) + return 0 + } + defer rows.Close() + + var members []string + for rows.Next() { + var name string + if err := rows.Scan(&name); err != nil { + continue + } + members = append(members, name) + } + + age := formatAge(msg.StateAge) + truncBody := truncate(msg.Body, 100) + count := int64(0) + + for _, member := range members { + body := fmt.Sprintf( + "**STALE**: Message #%d in #%s in '%s' for %s. \"%s\" — @%s", + msg.ID, msg.Channel, msg.State, age, truncBody, msg.FromAgent, + ) + + _, err := w.msgService.SendMessage(ctx, "system", member, body, SendOptions{ + Subject: fmt.Sprintf("workflow-stalemate-reminder:%d", msg.ID), + Priority: 7, + Metadata: fmt.Sprintf(`{"workflow_stalemate_reminder_for":%d}`, msg.ID), + }) + if err != nil { + w.logger.Error("send workflow stalemate reminder failed", + "message_id", msg.ID, + "to_agent", member, + "error", err, + ) + continue + } + w.logger.Info("sent workflow stalemate reminder", + "message_id", msg.ID, + "channel", msg.Channel, + "state", msg.State, + "to_agent", member, + "age", age, + ) + count++ + } + return count +} + +// sendWorkflowEscalation posts an escalation to #approvals for a stale workflow message. +func (w *StalemateWorker) sendWorkflowEscalation(ctx context.Context, msg staleWorkflowMsg) bool { + approvalsChanID, err := w.channelLookup.GetChannelIDByName(ctx, "approvals") + if err != nil { + w.logger.Warn("cannot escalate workflow stalemate: #approvals channel not found", "error", err) + return false + } + + age := formatAge(msg.StateAge) + truncBody := truncate(msg.Body, 100) + + body := fmt.Sprintf( + "**STALE**: Message #%d in #%s in '%s' for %s. \"%s\" — @%s", + msg.ID, msg.Channel, msg.State, age, truncBody, msg.FromAgent, + ) + + _, err = w.msgService.SendMessage(ctx, "system", "", body, SendOptions{ + Subject: fmt.Sprintf("workflow-stalemate-escalation:%d", msg.ID), + Priority: 9, + Metadata: fmt.Sprintf(`{"workflow_stalemate_escalation_for":%d}`, msg.ID), + ChannelID: &approvalsChanID, + }) + if err != nil { + w.logger.Error("send workflow escalation to #approvals failed", + "message_id", msg.ID, + "channel", msg.Channel, + "error", err, + ) + return false + } + w.logger.Info("escalated stale workflow message to #approvals", + "message_id", msg.ID, + "channel", msg.Channel, + "state", msg.State, + "age", age, + ) + return true +} + +// workflowReminderExists checks if a workflow stalemate reminder already exists for a message. +func (w *StalemateWorker) workflowReminderExists(ctx context.Context, messageID int64) bool { + var count int + err := w.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM messages + WHERE from_agent = 'system' + AND metadata LIKE ?`, + fmt.Sprintf(`%%"workflow_stalemate_reminder_for":%d%%`, messageID), + ).Scan(&count) + if err != nil { + return false + } + return count > 0 +} + +// workflowEscalationExists checks if a workflow stalemate escalation already exists for a message. +func (w *StalemateWorker) workflowEscalationExists(ctx context.Context, messageID int64) bool { + var count int + err := w.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM messages + WHERE from_agent = 'system' + AND metadata LIKE ?`, + fmt.Sprintf(`%%"workflow_stalemate_escalation_for":%d%%`, messageID), + ).Scan(&count) + if err != nil { + return false + } + return count > 0 +} + // truncate truncates a string to maxLen characters, appending "..." if truncated. func truncate(s string, maxLen int) string { runes := []rune(s) diff --git a/internal/messaging/stalemate_test.go b/internal/messaging/stalemate_test.go index e0035ed..9ca1778 100644 --- a/internal/messaging/stalemate_test.go +++ b/internal/messaging/stalemate_test.go @@ -478,3 +478,345 @@ func TestFormatAge(t *testing.T) { }) } } + +func TestComputeWorkflowStateFromReactions(t *testing.T) { + tests := []struct { + name string + reactions []reactionRow + want string + }{ + {"no reactions = proposed", nil, "proposed"}, + {"approve only", []reactionRow{{Reaction: "approve"}}, "approved"}, + {"in_progress only", []reactionRow{{Reaction: "in_progress"}}, "in_progress"}, + {"reject only", []reactionRow{{Reaction: "reject"}}, "rejected"}, + {"done only", []reactionRow{{Reaction: "done"}}, "done"}, + {"published only", []reactionRow{{Reaction: "published"}}, "published"}, + {"approve + in_progress = in_progress (higher priority)", []reactionRow{ + {Reaction: "approve"}, + {Reaction: "in_progress"}, + }, "in_progress"}, + {"approve + done = done", []reactionRow{ + {Reaction: "approve"}, + {Reaction: "done"}, + }, "done"}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := computeWorkflowStateFromReactions(tt.reactions) + if got != tt.want { + t.Errorf("computeWorkflowStateFromReactions() = %q, want %q", got, tt.want) + } + }) + } +} + +func TestIsTerminalWorkflowState(t *testing.T) { + tests := []struct { + state string + terminal bool + }{ + {"proposed", false}, + {"approved", false}, + {"in_progress", false}, + {"rejected", true}, + {"done", true}, + {"published", true}, + } + + for _, tt := range tests { + t.Run(tt.state, func(t *testing.T) { + got := isTerminalWorkflowState(tt.state) + if got != tt.terminal { + t.Errorf("isTerminalWorkflowState(%q) = %v, want %v", tt.state, got, tt.terminal) + } + }) + } +} + +func TestStalemateWorker_WorkflowReminder(t *testing.T) { + svc, db := newStalemateTestService(t) + ctx := context.Background() + + // Create a workflow-enabled channel with short timeouts + _, err := db.Exec( + `INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, workflow_enabled, stalemate_remind_after, stalemate_escalate_after, created_at, updated_at) + VALUES (10, 'news-test', 'Test news channel', '', 'standard', 0, 0, 'system', 1, '1s', '72h', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`) + if err != nil { + t.Fatalf("create workflow channel: %v", err) + } + + // Add system and sender as members + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (10, 'system', 'owner', CURRENT_TIMESTAMP)`) + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (10, 'sender', 'member', CURRENT_TIMESTAMP)`) + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (10, 'receiver', 'member', CURRENT_TIMESTAMP)`) + + // Insert a channel message with old created_at (will be in "proposed" state since no reactions) + oldTime := time.Now().Add(-2 * time.Second) + convResult, err := db.Exec( + `INSERT INTO conversations (subject, created_by, created_at, updated_at) VALUES ('wf-test', 'sender', ?, ?)`, + oldTime, oldTime, + ) + if err != nil { + t.Fatalf("insert conversation: %v", err) + } + convID, _ := convResult.LastInsertId() + + channelID := int64(10) + _, err = db.Exec( + `INSERT INTO messages (conversation_id, from_agent, to_agent, body, priority, status, metadata, channel_id, created_at, updated_at) + VALUES (?, 'sender', '', 'Draft blog post about MCP', 5, 'pending', '{}', ?, ?, ?)`, + convID, channelID, oldTime, oldTime, + ) + if err != nil { + t.Fatalf("insert channel message: %v", err) + } + + // Wait for the timeout to elapse + time.Sleep(10 * time.Millisecond) + + config := DefaultStalemateConfig() + lookup := &stubChannelLookup{channelID: 0, err: fmt.Errorf("no approvals channel")} + worker := NewStalemateWorker(db, svc, lookup, config) + + worker.checkStaleMessages(ctx) + + // Verify workflow stalemate reminders were sent to channel members + var count int + err = db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM messages WHERE from_agent = 'system' AND body LIKE '%STALE%'`, + ).Scan(&count) + if err != nil { + t.Fatalf("query workflow reminders: %v", err) + } + // Should have sent reminders to all 3 members (system, sender, receiver) + if count < 1 { + t.Errorf("expected at least 1 workflow reminder, got %d", count) + } +} + +func TestStalemateWorker_WorkflowEscalation(t *testing.T) { + svc, db := newStalemateTestService(t) + ctx := context.Background() + + // Create a workflow-enabled channel with short escalation timeout + _, err := db.Exec( + `INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, workflow_enabled, stalemate_remind_after, stalemate_escalate_after, created_at, updated_at) + VALUES (10, 'news-test', 'Test news channel', '', 'standard', 0, 0, 'system', 1, '1s', '1s', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`) + if err != nil { + t.Fatalf("create workflow channel: %v", err) + } + + // Create #approvals channel + db.Exec( + `INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, created_at, updated_at) + VALUES (20, 'approvals', 'Approval queue', '', 'standard', 0, 0, 'system', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`) + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (20, 'system', 'owner', CURRENT_TIMESTAMP)`) + + // Add members to workflow channel + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (10, 'sender', 'member', CURRENT_TIMESTAMP)`) + + // Insert a channel message old enough to trigger escalation + oldTime := time.Now().Add(-2 * time.Second) + convResult, _ := db.Exec( + `INSERT INTO conversations (subject, created_by, created_at, updated_at) VALUES ('wf-esc', 'sender', ?, ?)`, + oldTime, oldTime, + ) + convID, _ := convResult.LastInsertId() + + channelID := int64(10) + _, err = db.Exec( + `INSERT INTO messages (conversation_id, from_agent, to_agent, body, priority, status, metadata, channel_id, created_at, updated_at) + VALUES (?, 'sender', '', 'Stale proposal needing attention', 5, 'pending', '{}', ?, ?, ?)`, + convID, channelID, oldTime, oldTime, + ) + if err != nil { + t.Fatalf("insert channel message: %v", err) + } + + time.Sleep(10 * time.Millisecond) + + config := DefaultStalemateConfig() + lookup := &stubChannelLookup{channelID: 20} + worker := NewStalemateWorker(db, svc, lookup, config) + + worker.checkStaleMessages(ctx) + + // Verify escalation was sent to #approvals + var count int + err = db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM messages WHERE from_agent = 'system' AND channel_id = 20 AND body LIKE '%STALE%'`, + ).Scan(&count) + if err != nil { + t.Fatalf("query workflow escalation: %v", err) + } + if count != 1 { + t.Errorf("expected 1 workflow escalation, got %d", count) + } +} + +func TestStalemateWorker_WorkflowTerminalStateSkip(t *testing.T) { + svc, db := newStalemateTestService(t) + ctx := context.Background() + + // Create a workflow-enabled channel with short timeouts + _, err := db.Exec( + `INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, workflow_enabled, stalemate_remind_after, stalemate_escalate_after, created_at, updated_at) + VALUES (10, 'news-test', 'Test news channel', '', 'standard', 0, 0, 'system', 1, '1s', '1s', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`) + if err != nil { + t.Fatalf("create workflow channel: %v", err) + } + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (10, 'sender', 'member', CURRENT_TIMESTAMP)`) + + // Insert a channel message + oldTime := time.Now().Add(-2 * time.Second) + convResult, _ := db.Exec( + `INSERT INTO conversations (subject, created_by, created_at, updated_at) VALUES ('wf-done', 'sender', ?, ?)`, + oldTime, oldTime, + ) + convID, _ := convResult.LastInsertId() + + channelID := int64(10) + msgResult, err := db.Exec( + `INSERT INTO messages (conversation_id, from_agent, to_agent, body, priority, status, metadata, channel_id, created_at, updated_at) + VALUES (?, 'sender', '', 'Completed task', 5, 'pending', '{}', ?, ?, ?)`, + convID, channelID, oldTime, oldTime, + ) + if err != nil { + t.Fatalf("insert channel message: %v", err) + } + msgID, _ := msgResult.LastInsertId() + + // Add a "done" reaction — puts it in terminal state + _, err = db.Exec( + `INSERT INTO message_reactions (message_id, agent_name, reaction, metadata, created_at) + VALUES (?, 'sender', 'done', '{}', ?)`, + msgID, oldTime, + ) + if err != nil { + t.Fatalf("insert reaction: %v", err) + } + + time.Sleep(10 * time.Millisecond) + + config := DefaultStalemateConfig() + lookup := &stubChannelLookup{channelID: 0, err: fmt.Errorf("no approvals")} + worker := NewStalemateWorker(db, svc, lookup, config) + + worker.checkStaleMessages(ctx) + + // Verify NO reminders were sent (message is in terminal "done" state) + var count int + err = db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM messages WHERE from_agent = 'system' AND body LIKE '%STALE%'`, + ).Scan(&count) + if err != nil { + t.Fatalf("query reminders: %v", err) + } + if count != 0 { + t.Errorf("expected 0 reminders for terminal state message, got %d", count) + } +} + +func TestStalemateWorker_WorkflowDuplicateReminderPrevention(t *testing.T) { + svc, db := newStalemateTestService(t) + ctx := context.Background() + + // Create a workflow-enabled channel with short timeout + _, err := db.Exec( + `INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, workflow_enabled, stalemate_remind_after, stalemate_escalate_after, created_at, updated_at) + VALUES (10, 'news-test', 'Test', '', 'standard', 0, 0, 'system', 1, '1s', '72h', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`) + if err != nil { + t.Fatalf("create workflow channel: %v", err) + } + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (10, 'receiver', 'member', CURRENT_TIMESTAMP)`) + + // Insert a channel message + oldTime := time.Now().Add(-2 * time.Second) + convResult, _ := db.Exec( + `INSERT INTO conversations (subject, created_by, created_at, updated_at) VALUES ('wf-dup', 'sender', ?, ?)`, + oldTime, oldTime, + ) + convID, _ := convResult.LastInsertId() + + channelID := int64(10) + _, err = db.Exec( + `INSERT INTO messages (conversation_id, from_agent, to_agent, body, priority, status, metadata, channel_id, created_at, updated_at) + VALUES (?, 'sender', '', 'Needs review', 5, 'pending', '{}', ?, ?, ?)`, + convID, channelID, oldTime, oldTime, + ) + if err != nil { + t.Fatalf("insert channel message: %v", err) + } + + time.Sleep(10 * time.Millisecond) + + config := DefaultStalemateConfig() + lookup := &stubChannelLookup{channelID: 0, err: fmt.Errorf("no approvals")} + worker := NewStalemateWorker(db, svc, lookup, config) + + // Run twice + worker.checkStaleMessages(ctx) + worker.checkStaleMessages(ctx) + + // Verify only one set of reminders was sent (no duplicates) + var count int + err = db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM messages WHERE from_agent = 'system' AND to_agent = 'receiver' AND body LIKE '%STALE%'`, + ).Scan(&count) + if err != nil { + t.Fatalf("query reminders: %v", err) + } + if count != 1 { + t.Errorf("expected 1 reminder (no duplicates), got %d", count) + } +} + +func TestStalemateWorker_WorkflowNonWorkflowChannelSkip(t *testing.T) { + svc, db := newStalemateTestService(t) + ctx := context.Background() + + // Create a channel with workflow DISABLED + _, err := db.Exec( + `INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, workflow_enabled, stalemate_remind_after, stalemate_escalate_after, created_at, updated_at) + VALUES (10, 'general', 'General', '', 'standard', 0, 0, 'system', 0, '1s', '1s', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`) + if err != nil { + t.Fatalf("create channel: %v", err) + } + db.Exec(`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (10, 'sender', 'member', CURRENT_TIMESTAMP)`) + + // Insert a channel message + oldTime := time.Now().Add(-2 * time.Second) + convResult, _ := db.Exec( + `INSERT INTO conversations (subject, created_by, created_at, updated_at) VALUES ('no-wf', 'sender', ?, ?)`, + oldTime, oldTime, + ) + convID, _ := convResult.LastInsertId() + + channelID := int64(10) + db.Exec( + `INSERT INTO messages (conversation_id, from_agent, to_agent, body, priority, status, metadata, channel_id, created_at, updated_at) + VALUES (?, 'sender', '', 'No workflow here', 5, 'pending', '{}', ?, ?, ?)`, + convID, channelID, oldTime, oldTime, + ) + + time.Sleep(10 * time.Millisecond) + + config := DefaultStalemateConfig() + lookup := &stubChannelLookup{channelID: 0, err: fmt.Errorf("no approvals")} + worker := NewStalemateWorker(db, svc, lookup, config) + + worker.checkStaleMessages(ctx) + + // Verify NO reminders — channel is not workflow-enabled + var count int + err = db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM messages WHERE from_agent = 'system' AND body LIKE '%STALE%'`, + ).Scan(&count) + if err != nil { + t.Fatalf("query reminders: %v", err) + } + if count != 0 { + t.Errorf("expected 0 reminders for non-workflow channel, got %d", count) + } +} diff --git a/internal/web/dist/index.html b/internal/web/dist/index.html index 8abb30f..1ca473f 100644 --- a/internal/web/dist/index.html +++ b/internal/web/dist/index.html @@ -11,30 +11,30 @@ - - + + - +