feat(telegram): add message edit and delete support with improved error handling
This commit is contained in:
parent
028605cfd0
commit
692e002adf
3 changed files with 176 additions and 3 deletions
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue