mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-11 03:13:24 +00:00
fix(telegram): preserve Telegram topic routing for delayed notifications (#850)
* Fix delayed Telegram topic routing * refactor: export TaskLocalKeyMetadata and fix formatting - Export TaskLocalKeyMetadata from tools package for reuse - Remove duplicate taskLocalKeyMetadata from tasks package - Fix gofmt indentation issue in notifyLeaders - Add trailing newline to team_tool_helpers_test.go --------- Co-authored-by: viettranx <viettranx@gmail.com>
This commit is contained in:
1 parent
7e798a5d13
commit
6608e3dafe
5 files changed
+139
-18
No files matched your search
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
Reference in new issue
Block a user