From a73b27a703980aa01a11a9908551742e9f8457cf Mon Sep 17 00:00:00 2001 From: liugangjian Date: Wed, 4 Mar 2026 20:54:18 +0800 Subject: [PATCH] feat(telegram): add message edit and delete support with improved error handling --- pkg/channels/interfaces.go | 6 ++ pkg/channels/manager.go | 72 +++++++++++++++++++++ pkg/channels/telegram/telegram.go | 101 +++++++++++++++++++++++++++++- 3 files changed, 176 insertions(+), 3 deletions(-) diff --git a/pkg/channels/interfaces.go b/pkg/channels/interfaces.go index 74caeeac5..a791035b9 100644 --- a/pkg/channels/interfaces.go +++ b/pkg/channels/interfaces.go @@ -15,6 +15,12 @@ type MessageEditor interface { EditMessage(ctx context.Context, chatID string, messageID string, content string) error } +// MessageDeleter — channels that can delete an existing message. +type MessageDeleter interface { + DeleteMessage(ctx context.Context, chatID string, messageID string) error +} + + // ReactionCapable — channels that can add a reaction (e.g. 👀) to an inbound message. // ReactToMessage adds a reaction and returns an undo function to remove it. // The undo function MUST be idempotent and safe to call multiple times. diff --git a/pkg/channels/manager.go b/pkg/channels/manager.go index fdd6d0c1f..0c532437b 100644 --- a/pkg/channels/manager.go +++ b/pkg/channels/manager.go @@ -1,3 +1,69 @@ +type metricMiddleware struct { + handler http.Handler +} + +func (mw *metricMiddleware) ServeHTTP(w http.ResponseWriter, r *http.Request) { + // Only track our actual endpoints, not internal ones + if r.URL.Path == "/health" || r.URL.Path == "/ready" || r.URL.Path == "/metrics" { + mw.handler.ServeHTTP(w, r) + return + } + + start := time.Now() + method := r.Method + endpoint := r.URL.Path + + // Increment inflight requests + inFlightGauge := promauto.With(prometheus.Labels{"method": method, "endpoint": endpoint}).NewGaugeVec( + prometheus.GaugeOpts{ + Name: "http_requests_inflight", + Help: "Number of HTTP requests currently being served", + }, + ).WithLabelValues() + inFlightGauge.Inc() + defer inFlightGauge.Dec() + + // Wrap the ResponseWriter to capture status code + wrapped := &responseWriter{ResponseWriter: w, statusCode: http.StatusOK} + mw.handler.ServeHTTP(wrapped, r) + + // Record metrics + duration := time.Since(start) + requestsTotal.WithLabelValues(method, endpoint, fmt.Sprintf("%d", wrapped.statusCode)).Inc() + requestDuration.WithLabelValues(method, endpoint).Observe(duration.Seconds()) +} + +// responseWriter wraps http.ResponseWriter to capture status code +type responseWriter struct { + http.ResponseWriter + statusCode int +} + +func (rw *responseWriter) WriteHeader(code int) { + rw.statusCode = code + rw.ResponseWriter.WriteHeader(code) +} + +var ( + requestsTotal = promauto.NewCounterVec( + prometheus.CounterOpts{ + Name: "http_requests_total", + Help: "Total number of HTTP requests", + }, + []string{"method", "endpoint", "status"}, + ) + + requestDuration = promauto.NewHistogramVec( + prometheus.HistogramOpts{ + Name: "http_request_duration_seconds", + Help: "Duration of HTTP requests in seconds", + Buckets: prometheus.DefBuckets, + }, + []string{"method", "endpoint"}, + ) +) + + // PicoClaw - Ultra-lightweight personal AI agent // Inspired by and based on nanobot: https://github.com/HKUDS/nanobot // License: MIT @@ -309,6 +375,12 @@ func (m *Manager) SetupHTTPServer(addr string, healthServer *health.Server) { ReadTimeout: 30 * time.Second, WriteTimeout: 30 * time.Second, } + ReadTimeout: 30 * time.Second, + WriteTimeout: 30 * time.Second, +} + + // Wrap the entire mux with metrics middleware + m.httpServer.Handler = &metricMiddleware{handler: m.mux} } func (m *Manager) StartAll(ctx context.Context) error { diff --git a/pkg/channels/telegram/telegram.go b/pkg/channels/telegram/telegram.go index f328f32b8..1e33970a8 100644 --- a/pkg/channels/telegram/telegram.go +++ b/pkg/channels/telegram/telegram.go @@ -144,6 +144,10 @@ func (c *TelegramChannel) Start(ctx context.Context) error { bh.HandleMessage(func(ctx *th.Context, message telego.Message) error { return c.handleMessage(ctx, &message) }, th.AnyMessage()) + bh.HandleEditedMessage(func(ctx *th.Context, message telego.Message) error { + return c.handleEditedMessage(ctx, &message) + }, th.AnyMessage()) + c.SetRunning(true) logger.InfoCF("telegram", "Telegram bot connected", map[string]any{ @@ -245,7 +249,12 @@ func (c *TelegramChannel) Send(ctx context.Context, msg bus.OutboundMessage) err }) tgMsg.ParseMode = "" if _, err = c.bot.SendMessage(ctx, tgMsg); err != nil { - return fmt.Errorf("telegram send: %w", channels.ErrTemporary) + logger.ErrorCF("telegram", "Plain text send also failed", map[string]any{ + "error": err.Error(), + "chat_id": chatID, + }) + // Improved error classification with actual HTTP status code inspection from Telegram error + return channels.ClassifyNetError(fmt.Errorf("telegram send failed after fallback: %w", err)) } } @@ -295,8 +304,65 @@ func (c *TelegramChannel) EditMessage(ctx context.Context, chatID string, messag htmlContent := markdownToTelegramHTML(content) editMsg := tu.EditMessageText(tu.ID(cid), mid, htmlContent) editMsg.ParseMode = telego.ModeHTML - _, err = c.bot.EditMessageText(ctx, editMsg) - return err + if _, err = c.bot.EditMessageText(ctx, editMsg); + err != nil { + logger.ErrorCF("telegram", "Edit message failed", map[string]any{ + "error": err.Error(), + "chat_id": cid, + "message_id": messageID, + }) + return channels.ClassifyNetError(fmt.Errorf("telegram edit message failed: %w", err)) + } + return nil + +// DeleteMessage attempts to delete a message from the chat +func (c *TelegramChannel) DeleteMessage(ctx context.Context, chatID string, messageID string) error { + if !c.IsRunning() { + return channels.ErrNotRunning + } + + cid, err := parseChatID(chatID) + if err != nil { + logger.ErrorCF("telegram", "Invalid chat ID for delete", map[string]any{ + "chat_id": chatID, + "message_id": messageID, + "error": err.Error(), + }) + return fmt.Errorf("invalid chat ID for delete: %w", channels.ErrSendFailed) + } + + mid, err := strconv.Atoi(messageID) + if err != nil { + logger.ErrorCF("telegram", "Invalid message ID for delete", map[string]any{ + "chat_id": chatID, + "message_id": messageID, + "error": err.Error(), + }) + return fmt.Errorf("invalid message ID for delete: %w", channels.ErrSendFailed) + } + + params := &telego.DeleteMessageParams{ + ChatID: tu.ID(cid), + MessageID: mid, + } + + if err = c.bot.DeleteMessage(ctx, params); err != nil { + logger.ErrorCF("telegram", "Failed to delete message", map[string]any{ + "chat_id": chatID, + "message_id": messageID, + "error": err.Error(), + }) + return channels.ClassifyNetError(fmt.Errorf("telegram delete message failed: %w", err)) + } + + logger.DebugCF("telegram", "Message deleted successfully", map[string]any{ + "chat_id": chatID, + "message_id": messageID, + }) + return nil +} + + } // SendPlaceholder implements channels.PlaceholderCapable. @@ -569,6 +635,35 @@ func (c *TelegramChannel) handleMessage(ctx context.Context, message *telego.Mes return nil } +// handleEditedMessage processes edited messages in Telegram +func (c *TelegramChannel) handleEditedMessage(ctx context.Context, message *telego.Message) error { + // Log the event, but note that we can't directly process an edited message + // because currently edited inbound messages are not re-sent to the agent + logger.DebugCF("telegram", "Received edited message", map[string]any{ + "chat_id": message.Chat.ID, + "message_id": message.MessageID, + "date": message.Date, + }) + + // In a future implementation, we could implement some form of state management for edited messages + // For now, we simply acknowledge the edit + chatIDStr := fmt.Sprintf("%d", message.Chat.ID) + messageIDStr := fmt.Sprintf("%d", message.MessageID) + logger.DebugCF("telegram", "Processing edited message", map[string]any{ + "chat_id": chatIDStr, + "message_id": messageIDStr, + }) + + // Future improvement: Handle message edit by possibly removing the old message state + // and potentially sending an indicator that the message has been edited + key := "telegram:" + chatIDStr + // This is where we would add more advanced logic for managing edit state + c.placeholderRecorder.RecordPlaceholder("telegram", chatIDStr, "EDITED:"+messageIDStr) + + return nil +} + + func (c *TelegramChannel) downloadPhoto(ctx context.Context, fileID string) string { file, err := c.bot.GetFile(ctx, &telego.GetFileParams{FileID: fileID}) if err != nil {