diff --git a/internal/tasks/task_ticker.go b/internal/tasks/task_ticker.go index 13ccf9a3..fbe85509 100644 --- a/internal/tasks/task_ticker.go +++ b/internal/tasks/task_ticker.go @@ -18,11 +18,11 @@ import ( ) const ( - defaultRecoveryInterval = 5 * time.Minute - defaultStaleThreshold = 2 * time.Hour - defaultInReviewThreshold = 4 * time.Hour - followupCooldown = 5 * time.Minute - defaultFollowupInterval = 30 * time.Minute + defaultRecoveryInterval = 5 * time.Minute + defaultStaleThreshold = 2 * time.Hour + defaultInReviewThreshold = 4 * time.Hour + followupCooldown = 5 * time.Minute + defaultFollowupInterval = 30 * time.Minute ) // TaskTicker periodically recovers stale tasks and re-dispatches pending work. @@ -237,9 +237,13 @@ func (t *TaskTicker) notifyLeaders(ctx context.Context, tasks []store.RecoveredT // Resolve PeerKind from first task's metadata for correct session routing (#266). var peerKind string - if fullTask, err := t.teams.GetTask(ctx, scopeTasks[0].ID); err == nil && fullTask != nil && fullTask.Metadata != nil { - if pk, ok := fullTask.Metadata["peer_kind"].(string); ok { - peerKind = pk + var fullTask *store.TeamTaskData + if task, err := t.teams.GetTask(ctx, scopeTasks[0].ID); err == nil { + fullTask = task + if fullTask != nil && fullTask.Metadata != nil { + if pk, ok := fullTask.Metadata["peer_kind"].(string); ok { + peerKind = pk + } } } @@ -247,6 +251,7 @@ func (t *TaskTicker) notifyLeaders(ctx context.Context, tasks []store.RecoveredT Channel: channel, SenderID: "ticker:system", ChatID: chatID, + Metadata: tools.TaskLocalKeyMetadata(fullTask), AgentID: lead.AgentKey, UserID: team.CreatedBy, PeerKind: peerKind, @@ -334,11 +339,7 @@ func (t *TaskTicker) processTeamFollowups(ctx context.Context, tasks []store.Tea } content := fmt.Sprintf("Reminder (%s): %s", countLabel, task.FollowupMessage) - if !t.msgBus.TryPublishOutbound(bus.OutboundMessage{ - Channel: task.FollowupChannel, - ChatID: task.FollowupChatID, - Content: content, - }) { + if !t.msgBus.TryPublishOutbound(followupOutboundMessage(task, content)) { slog.Warn("task_ticker: outbound buffer full, skipping followup", "task_id", task.ID) continue } @@ -370,6 +371,16 @@ func (t *TaskTicker) processTeamFollowups(ctx context.Context, tasks []store.Tea } } +func followupOutboundMessage(task *store.TeamTaskData, content string) bus.OutboundMessage { + message := bus.OutboundMessage{ + Channel: task.FollowupChannel, + ChatID: task.FollowupChatID, + Content: content, + } + message.Metadata = tools.TaskLocalKeyMetadata(task) + return message +} + // followupInterval parses the team's followup_interval_minutes setting. func followupInterval(team store.TeamData) time.Duration { if team.Settings != nil { diff --git a/internal/tasks/task_ticker_test.go b/internal/tasks/task_ticker_test.go index 03a052c2..d93802f1 100644 --- a/internal/tasks/task_ticker_test.go +++ b/internal/tasks/task_ticker_test.go @@ -13,6 +13,7 @@ import ( "github.com/nextlevelbuilder/goclaw/internal/bus" "github.com/nextlevelbuilder/goclaw/internal/store" + "github.com/nextlevelbuilder/goclaw/internal/tools" ) // ─── minimal stub stores ─────────────────────────────────────────────────── @@ -460,3 +461,33 @@ func TestProcessFollowups_NoTasksIsNoop(t *testing.T) { t.Errorf("IncrementFollowupCount calls = %d, want 0", ts.incrementCalls.Load()) } } + +func TestFollowupOutboundMessage_UsesLocalKey(t *testing.T) { + task := &store.TeamTaskData{ + FollowupChannel: "telegram", + FollowupChatID: "-100123456", + Metadata: map[string]any{ + tools.TaskMetaLocalKey: "-100123456:topic:47", + }, + } + + got := followupOutboundMessage(task, "Reminder (1): ping") + if got.Metadata == nil { + t.Fatal("expected metadata to be populated") + } + if got.Metadata["local_key"] != "-100123456:topic:47" { + t.Fatalf("local_key = %q, want %q", got.Metadata["local_key"], "-100123456:topic:47") + } +} + +func TestFollowupOutboundMessage_OmitsLocalKeyWhenMissing(t *testing.T) { + task := &store.TeamTaskData{ + FollowupChannel: "telegram", + FollowupChatID: "-100123456", + } + + got := followupOutboundMessage(task, "Reminder (1): ping") + if got.Metadata != nil { + t.Fatalf("expected metadata to be nil, got %#v", got.Metadata) + } +} diff --git a/internal/tools/team_tasks_blocker.go b/internal/tools/team_tasks_blocker.go index a994fd77..6e486428 100644 --- a/internal/tools/team_tasks_blocker.go +++ b/internal/tools/team_tasks_blocker.go @@ -77,6 +77,7 @@ func (t *TeamTasksTool) handleBlockerComment( Channel: task.Channel, SenderID: "system:escalation", ChatID: task.ChatID, + Metadata: TaskLocalKeyMetadata(task), Content: escalationMsg, UserID: store.UserIDFromContext(ctx), PeerKind: blockerPeerKind, diff --git a/internal/tools/team_tool_helpers.go b/internal/tools/team_tool_helpers.go index 8d594c00..dc64d934 100644 --- a/internal/tools/team_tool_helpers.go +++ b/internal/tools/team_tool_helpers.go @@ -21,6 +21,27 @@ func (m *TeamToolManager) broadcastTeamEvent(ctx context.Context, name string, p bus.BroadcastForTenant(m.msgBus, name, store.TenantIDFromContext(ctx), payload) } +func reviewOutboundMessage(task *store.TeamTaskData, content string) bus.OutboundMessage { + message := bus.OutboundMessage{ + Channel: task.Channel, + ChatID: task.ChatID, + Content: content, + } + message.Metadata = TaskLocalKeyMetadata(task) + return message +} + +// TaskLocalKeyMetadata extracts local_key from task metadata for Telegram forum topic routing. +func TaskLocalKeyMetadata(task *store.TeamTaskData) map[string]string { + if task == nil || task.Metadata == nil { + return nil + } + if localKey, ok := task.Metadata[TaskMetaLocalKey].(string); ok && localKey != "" { + return map[string]string{TaskMetaLocalKey: localKey} + } + return nil +} + // resolveTeamRole returns the calling agent's role in the team. // Unlike requireLead(), this does NOT bypass for teammate channel — // workspace RBAC must respect actual roles even for teammate agents. @@ -204,9 +225,5 @@ func (m *TeamToolManager) notifyChannelReview(task *store.TeamTaskData) { return } content := fmt.Sprintf("🔔 Escalation: \"%s\" requires human review (task %s).", task.Subject, task.Identifier) - m.msgBus.PublishOutbound(bus.OutboundMessage{ - Channel: task.Channel, - ChatID: task.ChatID, - Content: content, - }) + m.msgBus.PublishOutbound(reviewOutboundMessage(task, content)) } diff --git a/internal/tools/team_tool_helpers_test.go b/internal/tools/team_tool_helpers_test.go new file mode 100644 index 00000000..49dc8d73 --- /dev/null +++ b/internal/tools/team_tool_helpers_test.go @@ -0,0 +1,61 @@ +package tools + +import ( + "testing" + + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +func TestReviewOutboundMessage_UsesLocalKey(t *testing.T) { + task := &store.TeamTaskData{ + Channel: "telegram", + ChatID: "-100123456", + Metadata: map[string]any{ + TaskMetaLocalKey: "-100123456:topic:47", + }, + } + + got := reviewOutboundMessage(task, "review needed") + if got.Metadata == nil { + t.Fatal("expected metadata to be populated") + } + if got.Metadata["local_key"] != "-100123456:topic:47" { + t.Fatalf("local_key = %q, want %q", got.Metadata["local_key"], "-100123456:topic:47") + } +} + +func TestReviewOutboundMessage_OmitsLocalKeyWhenMissing(t *testing.T) { + task := &store.TeamTaskData{ + Channel: "telegram", + ChatID: "-100123456", + } + + got := reviewOutboundMessage(task, "review needed") + if got.Metadata != nil { + t.Fatalf("expected metadata to be nil, got %#v", got.Metadata) + } +} + +func TestTaskLocalKeyMetadata(t *testing.T) { + t.Run("uses local key", func(t *testing.T) { + task := &store.TeamTaskData{ + Metadata: map[string]any{ + TaskMetaLocalKey: "-100123456:topic:47", + }, + } + + got := TaskLocalKeyMetadata(task) + if got == nil { + t.Fatal("expected metadata to be populated") + } + if got[TaskMetaLocalKey] != "-100123456:topic:47" { + t.Fatalf("local_key = %q, want %q", got[TaskMetaLocalKey], "-100123456:topic:47") + } + }) + + t.Run("omits local key when missing", func(t *testing.T) { + if got := TaskLocalKeyMetadata(&store.TeamTaskData{}); got != nil { + t.Fatalf("expected metadata to be nil, got %#v", got) + } + }) +}