switch to go-pkgz/notify package: telegram

This commit is contained in:
Dmitry Verkhoturov
2022-04-29 13:32:15 -05:00
committed by Umputun
parent ad0ac693de
commit 59fb68ab2d
4 changed files with 71 additions and 901 deletions
+37 -419
View File
@@ -1,25 +1,13 @@
package notify
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
neturl "net/url"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
log "github.com/go-pkgz/lgr"
"github.com/go-pkgz/repeater"
ntf "github.com/go-pkgz/notify"
"github.com/hashicorp/go-multierror"
"github.com/microcosm-cc/bluemonday"
"github.com/pkg/errors"
"golang.org/x/net/html"
)
// TelegramParams contain settings for telegram notifications
@@ -29,78 +17,29 @@ type TelegramParams struct {
Timeout time.Duration // http client timeout
UserNotifications bool // flag which enables user notifications
ErrorMsg, SuccessMsg string // messages for successful and unsuccessful subscription requests to bot
apiPrefix string // changed only in tests
}
// Telegram implements notify.Destination for telegram
type Telegram struct {
TelegramParams
*ntf.Telegram
// Identifier of the first update to be requested.
// Should be equal to LastSeenUpdateID + 1
// See https://core.telegram.org/bots/api#getupdates
updateOffset int
apiPollInterval time.Duration // interval to check updates from Telegram API and answer to users
expiredCleanupInterval time.Duration // interval to check and clean up expired notification requests
username string // bot username
run int32 // non-zero if Run goroutine has started
requests struct {
sync.RWMutex
data map[string]tgAuthRequest
}
AdminChannelID string // unique identifier for the target chat or username of the target channel (in the format @channelusername)
UserNotifications bool // flag which enables user notifications
}
// telegramMsg is used to send message trough Telegram bot API
type telegramMsg struct {
Text string `json:"text"`
ParseMode string `json:"parse_mode,omitempty"`
}
type tgAuthRequest struct {
confirmed bool // whether login request has been confirmed and user info set
expires time.Time
telegramID string
user string
site string
}
// TelegramBotInfo structure contains information about telegram bot, which is used from whole telegram API response
type TelegramBotInfo struct {
Username string `json:"username"`
}
const telegramTimeOut = 5000 * time.Millisecond
const telegramAPIPrefix = "https://api.telegram.org/bot"
const tgPollInterval = time.Second * 5
const tgCleanupInterval = time.Minute * 5
// NewTelegram makes telegram bot for notifications
func NewTelegram(params TelegramParams) (*Telegram, error) {
res := Telegram{TelegramParams: params}
if res.apiPrefix == "" {
res.apiPrefix = telegramAPIPrefix
}
if res.Timeout == 0 {
res.Timeout = telegramTimeOut
}
res.apiPollInterval = tgPollInterval
res.expiredCleanupInterval = tgCleanupInterval
log.Printf("[DEBUG] create new telegram notifier for api=%s, timeout=%s", res.apiPrefix, res.Timeout)
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
botInfo, err := res.botInfo(ctx)
client, err := ntf.NewTelegram(ntf.TelegramParams{
Token: params.Token,
Timeout: params.Timeout,
ErrorMsg: params.ErrorMsg,
SuccessMsg: params.SuccessMsg,
})
if err != nil {
return nil, errors.Wrapf(err, "can't retrieve bot info from Telegram API")
return nil, err
}
res.username = botInfo.Username
res.requests.data = make(map[string]tgAuthRequest)
return &res, nil
return &Telegram{Telegram: client, AdminChannelID: params.AdminChannelID, UserNotifications: params.UserNotifications}, nil
}
// Send to telegram recipients
@@ -108,121 +47,55 @@ func (t *Telegram) Send(ctx context.Context, req Request) error {
log.Printf("[DEBUG] send telegram notification for comment ID %s", req.Comment.ID)
result := new(multierror.Error)
msg, err := buildMessage(req)
if err != nil {
return fmt.Errorf("failed to make telegram message body for comment ID %s: %w", req.Comment.ID, err)
}
msg := t.buildMessage(req)
if t.AdminChannelID != "" {
err := t.sendMessage(ctx, msg, t.AdminChannelID)
result = multierror.Append(errors.Wrapf(err,
"problem sending admin telegram notification about comment ID %s to %s", req.Comment.ID, t.AdminChannelID),
)
err := t.Telegram.Send(ctx, fmt.Sprintf("telegram:%s?parseMode=HTML", t.AdminChannelID), msg)
if err != nil {
result = multierror.Append(result,
fmt.Errorf("problem sending admin telegram notification about comment ID %s to %s: %w",
req.Comment.ID, t.AdminChannelID, err,
),
)
}
}
if t.UserNotifications {
for _, user := range req.Telegrams {
err := t.sendMessage(ctx, msg, user)
result = multierror.Append(errors.Wrapf(err,
"problem sending user telegram notification about comment ID %s to %q", req.Comment.ID, user),
)
err := t.Telegram.Send(ctx, fmt.Sprintf("telegram:%s?parseMode=HTML", user), msg)
if err != nil {
result = multierror.Append(result,
fmt.Errorf("problem sending user telegram notification about comment ID %s to %q: %w",
req.Comment.ID, user, err,
),
)
}
}
}
return result.ErrorOrNil()
}
func (t *Telegram) sendMessage(ctx context.Context, b []byte, chatID string) error {
if _, err := strconv.ParseInt(chatID, 10, 64); err != nil {
chatID = "@" + chatID // if chatID not a number enforce @ prefix
}
url := fmt.Sprintf("sendMessage?chat_id=%s&disable_web_page_preview=true", chatID)
return t.Request(ctx, url, b, &struct{}{})
}
// buildMessage generates message for generic notification about new comment
func buildMessage(req Request) ([]byte, error) {
func (t *Telegram) buildMessage(req Request) string {
commentURLPrefix := req.Comment.Locator.URL + uiNav
msg := fmt.Sprintf(`<a href=%q>%s</a>`, commentURLPrefix+req.Comment.ID, escapeTelegramText(req.Comment.User.Name))
msg := fmt.Sprintf(`<a href=%q>%s</a>`, commentURLPrefix+req.Comment.ID, ntf.EscapeTelegramText(req.Comment.User.Name))
if req.Comment.ParentID != "" {
msg += fmt.Sprintf(" -> <a href=%q>%s</a>", commentURLPrefix+req.parent.ID, escapeTelegramText(req.parent.User.Name))
msg += fmt.Sprintf(" -> <a href=%q>%s</a>", commentURLPrefix+req.parent.ID, ntf.EscapeTelegramText(req.parent.User.Name))
}
msg += fmt.Sprintf("\n\n%s", telegramSupportedHTML(req.Comment.Text))
msg += fmt.Sprintf("\n\n%s", ntf.TelegramSupportedHTML(req.Comment.Text))
if req.Comment.ParentID != "" {
msg += fmt.Sprintf("\n\n\"<i>%s</i>\"", telegramSupportedHTML(req.parent.Text))
msg += fmt.Sprintf("\n\n\"<i>%s</i>\"", ntf.TelegramSupportedHTML(req.parent.Text))
}
if req.Comment.PostTitle != "" {
msg += fmt.Sprintf("\n\n↦ <a href=%q>%s</a>", req.Comment.Locator.URL, escapeTelegramText(req.Comment.PostTitle))
msg += fmt.Sprintf("\n\n↦ <a href=%q>%s</a>", req.Comment.Locator.URL, ntf.EscapeTelegramText(req.Comment.PostTitle))
}
body := telegramMsg{Text: msg, ParseMode: "HTML"}
b, err := json.Marshal(body)
if err != nil {
return nil, err
}
return b, nil
}
// returns HTML with only tags allowed in Telegram HTML message payload, also trims ending newlines
// https://core.telegram.org/bots/api#html-style
func telegramSupportedHTML(htmlText string) string {
adjustedHTMLText := adjustHTMLTags(htmlText)
p := bluemonday.NewPolicy()
p.AllowElements("b", "strong", "i", "em", "u", "ins", "s", "strike", "del", "a", "code", "pre")
p.AllowAttrs("href").OnElements("a")
p.AllowAttrs("class").OnElements("code")
return strings.TrimRight(p.Sanitize(adjustedHTMLText), "\n")
}
// returns text sanitized of symbols not allowed inside other HTML tags in Telegram HTML message payload
// https://core.telegram.org/bots/api#html-style
func escapeTelegramText(text string) string {
// order is important
text = strings.ReplaceAll(text, "&", "&amp;")
text = strings.ReplaceAll(text, "<", "&lt;")
text = strings.ReplaceAll(text, ">", "&gt;")
return text
}
// telegram not allow h1-h6 tags
// replace these tags with a combination of <b> and <i> for visual distinction
func adjustHTMLTags(htmlText string) string {
buff := strings.Builder{}
tokenizer := html.NewTokenizer(strings.NewReader(htmlText))
for {
if tokenizer.Next() == html.ErrorToken {
return buff.String()
}
token := tokenizer.Token()
switch token.Type {
case html.StartTagToken, html.EndTagToken:
switch token.Data {
case "h1", "h2", "h3":
if token.Type == html.StartTagToken {
buff.WriteString("<b>")
}
if token.Type == html.EndTagToken {
buff.WriteString("</b>")
}
case "h4", "h5", "h6":
if token.Type == html.StartTagToken {
buff.WriteString("<i><b>")
}
if token.Type == html.EndTagToken {
buff.WriteString("</b></i>")
}
default:
buff.WriteString(token.String())
}
default:
buff.WriteString(token.String())
}
}
return msg
}
// SendVerification is not needed for telegram
@@ -230,132 +103,8 @@ func (t *Telegram) SendVerification(_ context.Context, _ VerificationRequest) er
return nil
}
// TelegramUpdate contains update information, which is used from whole telegram API response
type TelegramUpdate struct {
Result []struct {
UpdateID int `json:"update_id"`
Message struct {
Chat struct {
ID int `json:"id"`
Name string `json:"first_name"`
Type string `json:"type"`
} `json:"chat"`
Text string `json:"text"`
} `json:"message"`
} `json:"result"`
}
// GetBotUsername returns bot username
func (t *Telegram) GetBotUsername() string {
return t.username
}
// AddToken adds token
func (t *Telegram) AddToken(token, user, site string, expires time.Time) {
t.requests.Lock()
t.requests.data[token] = tgAuthRequest{
expires: expires,
user: user,
site: site,
}
t.requests.Unlock()
}
// CheckToken verifies incoming token, returns the user address if it's confirmed and empty string otherwise
func (t *Telegram) CheckToken(token, user string) (telegram, site string, err error) {
t.requests.RLock()
authRequest, ok := t.requests.data[token]
t.requests.RUnlock()
if !ok {
return "", "", fmt.Errorf("request is not found")
}
if time.Now().After(authRequest.expires) {
t.requests.Lock()
delete(t.requests.data, token)
t.requests.Unlock()
return "", "", fmt.Errorf("request expired")
}
if !authRequest.confirmed {
return "", "", fmt.Errorf("request is not verified yet")
}
if authRequest.user != user {
return "", "", fmt.Errorf("user does not match original requester")
}
// Delete request
t.requests.Lock()
delete(t.requests.data, token)
t.requests.Unlock()
return authRequest.telegramID, authRequest.site, nil
}
// Run starts processing login requests sent in Telegram, required for user notifications to work
// Blocks caller
func (t *Telegram) Run(ctx context.Context) {
atomic.AddInt32(&t.run, 1)
processUpdatedTicker := time.NewTicker(t.apiPollInterval)
cleanupTicker := time.NewTicker(t.expiredCleanupInterval)
for {
select {
case <-ctx.Done():
processUpdatedTicker.Stop()
cleanupTicker.Stop()
atomic.AddInt32(&t.run, -1)
return
case <-processUpdatedTicker.C:
updates, err := t.getUpdates(ctx)
if err != nil {
log.Printf("[WARN] Error while getting telegram updates: %v", err)
continue
}
t.processUpdates(ctx, updates)
case <-cleanupTicker.C:
now := time.Now()
t.requests.Lock()
for key, req := range t.requests.data {
if now.After(req.expires) {
delete(t.requests.data, key)
}
}
t.requests.Unlock()
}
}
}
// ProcessUpdate is alternative to Run, it processes provided plain text update from Telegram
// so that caller could get updates and send it not only there but to multiple sources
func (t *Telegram) ProcessUpdate(ctx context.Context, textUpdate string) error {
if atomic.LoadInt32(&t.run) != 0 {
return fmt.Errorf("the Run goroutine should not be used with ProcessUpdate")
}
defer func() {
// as Run goroutine is not running, clean up old requests on each update
// even if we hit json decode error
now := time.Now()
t.requests.Lock()
for key, req := range t.requests.data {
if now.After(req.expires) {
delete(t.requests.data, key)
}
}
t.requests.Unlock()
}()
var updates TelegramUpdate
if err := json.Unmarshal([]byte(textUpdate), &updates); err != nil {
return fmt.Errorf("failed to decode provided telegram update: %w", err)
}
t.processUpdates(ctx, &updates)
return nil
}
func (t *Telegram) String() string {
result := "telegram"
result := t.Telegram.String()
if t.AdminChannelID != "" {
result += " with admin notifications to " + t.AdminChannelID
}
@@ -364,134 +113,3 @@ func (t *Telegram) String() string {
}
return result
}
// getUpdates fetches incoming updates
func (t *Telegram) getUpdates(ctx context.Context) (*TelegramUpdate, error) {
url := `getUpdates?allowed_updates=["message"]`
if t.updateOffset != 0 {
url += fmt.Sprintf("&offset=%d", t.updateOffset)
}
var result TelegramUpdate
err := t.Request(ctx, url, nil, &result)
if err != nil {
return nil, fmt.Errorf("failed to fetch updates: %w", err)
}
for _, u := range result.Result {
if u.UpdateID >= t.updateOffset {
t.updateOffset = u.UpdateID + 1
}
}
return &result, nil
}
// processUpdates processes a batch of updates from telegram servers
func (t *Telegram) processUpdates(ctx context.Context, updates *TelegramUpdate) {
for _, update := range updates.Result {
if update.Message.Chat.Type != "private" {
continue
}
if !strings.HasPrefix(update.Message.Text, "/start ") {
continue
}
token := strings.TrimPrefix(update.Message.Text, "/start ")
t.requests.RLock()
authRequest, ok := t.requests.data[token]
if !ok { // No such token
t.requests.RUnlock()
if t.ErrorMsg != "" {
if err := t.sendText(ctx, update.Message.Chat.ID, t.ErrorMsg); err != nil {
log.Printf("[WARN] failed to notify telegram peer: %v", err)
}
}
continue
}
t.requests.RUnlock()
authRequest.confirmed = true
authRequest.telegramID = strconv.Itoa(update.Message.Chat.ID)
t.requests.Lock()
t.requests.data[token] = authRequest
t.requests.Unlock()
if err := t.sendText(ctx, update.Message.Chat.ID, t.SuccessMsg); err != nil {
log.Printf("[ERROR] failed to notify telegram peer: %v", err)
}
}
}
// sendText sends a plain text message to telegram peer
func (t *Telegram) sendText(ctx context.Context, recipientID int, msg string) error {
url := fmt.Sprintf("sendMessage?chat_id=%d&text=%s", recipientID, neturl.PathEscape(msg))
return t.Request(ctx, url, nil, &struct{}{})
}
// botInfo returns info about configured bot
func (t *Telegram) botInfo(ctx context.Context) (*TelegramBotInfo, error) {
var resp = struct {
Result *TelegramBotInfo `json:"result"`
}{}
err := t.Request(ctx, "getMe", nil, &resp)
if err != nil {
return nil, err
}
if resp.Result == nil {
return nil, fmt.Errorf("received empty result")
}
return resp.Result, nil
}
// Request makes a request to the Telegram API and return the result
func (t *Telegram) Request(ctx context.Context, method string, b []byte, data interface{}) error {
return repeater.NewDefault(3, time.Millisecond*250).Do(ctx, func() error {
url := fmt.Sprintf("%s%s/%s", t.apiPrefix, t.Token, method)
var req *http.Request
var err error
if b == nil {
req, err = http.NewRequestWithContext(ctx, "GET", url, http.NoBody)
} else {
req, err = http.NewRequestWithContext(ctx, "POST", url, bytes.NewReader(b))
req.Header.Set("Content-Type", "application/json; charset=utf-8")
}
if err != nil {
return fmt.Errorf("failed to create request: %w", err)
}
client := http.Client{Timeout: t.Timeout}
resp, err := client.Do(req)
if err != nil {
return fmt.Errorf("failed to send request: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return t.parseError(resp.Body, resp.StatusCode)
}
if err = json.NewDecoder(resp.Body).Decode(data); err != nil {
return fmt.Errorf("failed to decode json response: %w", err)
}
return nil
})
}
func (t *Telegram) parseError(r io.Reader, statusCode int) error {
tgErr := struct {
Description string `json:"description"`
}{}
if err := json.NewDecoder(r).Decode(&tgErr); err != nil {
return fmt.Errorf("unexpected telegram API status code %d", statusCode)
}
return fmt.Errorf("unexpected telegram API status code %d, error: %q", statusCode, tgErr.Description)
}
+30 -476
View File
@@ -2,507 +2,61 @@ package notify
import (
"context"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/go-chi/chi/v5"
ntf "github.com/go-pkgz/notify"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/umputun/remark42/backend/app/store"
)
func TestTelegram_New(t *testing.T) {
ts := mockTelegramServer(nil)
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "good-token",
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
assert.NotNil(t, tb)
assert.Equal(t, tb.Timeout, time.Second*5)
assert.Equal(t, "remark_test", tb.AdminChannelID, "@ added")
_, err = NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "empty-json",
apiPrefix: ts.URL + "/",
})
assert.EqualError(t, err, "can't retrieve bot info from Telegram API: received empty result")
st := time.Now()
_, err = NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "non-json-resp",
Timeout: 2 * time.Second,
apiPrefix: ts.URL + "/",
})
func TestTelegram_NewError(t *testing.T) {
tb, err := NewTelegram(TelegramParams{})
assert.Error(t, err)
assert.Contains(t, err.Error(), "failed to decode json response:")
assert.True(t, time.Since(st) >= 250*3*time.Millisecond)
_, err = NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "404",
Timeout: 2 * time.Second,
apiPrefix: ts.URL + "/",
})
assert.EqualError(t, err, "can't retrieve bot info from Telegram API: unexpected telegram API status code 404")
_, err = NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "no-such-thing",
apiPrefix: "http://127.0.0.1:4321/",
})
require.Error(t, err)
assert.Contains(t, err.Error(), "can't retrieve bot info from Telegram API")
assert.Contains(t, err.Error(), "dial tcp 127.0.0.1:4321: connect: connection refused")
_, err = NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "no-such-thing",
apiPrefix: "",
})
assert.Error(t, err, "empty api url not allowed")
_, err = NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "good-token",
Timeout: 2 * time.Second,
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err, "0 timeout allowed as default")
tb, err = NewTelegram(TelegramParams{
AdminChannelID: "1234567890",
Token: "good-token",
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
assert.NotNil(t, tb)
assert.Equal(t, "1234567890", tb.AdminChannelID, "no @ prefix")
}
func TestTelegram_GetBotUsername(t *testing.T) {
ts := mockTelegramServer(nil)
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "good-token",
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
assert.NotNil(t, tb)
assert.Equal(t, "remark42_test_bot", tb.GetBotUsername())
assert.Nil(t, tb)
}
func TestTelegram_Send(t *testing.T) {
ts := mockTelegramServer(nil)
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
tb := Telegram{
AdminChannelID: "remark_test",
Token: "good-token",
UserNotifications: true,
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
assert.NotNil(t, tb)
c := store.Comment{Text: "some text", ParentID: "1", ID: "999", Locator: store.Locator{URL: "http://example.org/"}}
Telegram: &ntf.Telegram{}, // broken sender due to unset API
}
assert.Equal(t, "telegram notifications destination with admin notifications to remark_test with user notifications enabled", tb.String())
c := store.Comment{Text: "some text", ParentID: "1", ID: "999", PostTitle: "[test title]", Locator: store.Locator{URL: "http://example.org/"}}
c.User.Name = "from"
cp := store.Comment{Text: `<p>some parent text with a <a href="http://example.org">link</a> and special text:<br>& < > &</p>`}
cp.User.Name = "to"
err = tb.Send(context.TODO(), Request{Comment: c, parent: cp, Telegrams: []string{"test_user_channel"}})
assert.NoError(t, err)
c.PostTitle = "test title"
err = tb.Send(context.TODO(), Request{Comment: c, parent: cp})
assert.NoError(t, err)
err = tb.Send(context.TODO(), Request{Comment: c, parent: cp})
assert.NoError(t, err)
c.PostTitle = "[test title]"
err = tb.Send(context.TODO(), Request{Comment: c, parent: cp})
assert.NoError(t, err)
tb, err = NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "non-json-resp",
UserNotifications: true,
apiPrefix: ts.URL + "/",
})
assert.Nil(t, tb)
assert.Error(t, err, "should fail")
tb = &Telegram{
TelegramParams: TelegramParams{
AdminChannelID: "remark_test",
Token: "non-json-resp",
UserNotifications: true,
apiPrefix: ts.URL + "/",
}}
err = tb.Send(context.TODO(), Request{Comment: c, parent: cp, Telegrams: []string{"test_user_channel"}})
require.Error(t, err)
assert.Contains(t, err.Error(), "unexpected telegram API status code 404", "send on broken tg")
assert.Equal(t, "telegram with admin notifications to remark_test with user notifications enabled", tb.String())
// bad API URL
tb.apiPrefix = "http://non-existent"
err = tb.Send(context.TODO(), Request{Comment: c, parent: cp, Telegrams: []string{"test_user_channel"}})
err := tb.Send(context.Background(), Request{Comment: c, parent: cp, Telegrams: []string{"test_user_channel"}})
assert.Error(t, err)
assert.Contains(t, err.Error(), "2 errors occurred")
assert.Contains(t, err.Error(), "problem sending user telegram notification about comment ID 999 to \"test_user_channel\"")
assert.Contains(t, err.Error(), "problem sending admin telegram notification about comment ID 999 to remark_test")
// test buildMessage separately for message text
res, err := buildMessage(Request{Comment: c, parent: cp})
assert.NoError(t, err)
assert.Equal(t, `{"text":"\u003ca href=\"http://example.org/#remark42__comment-999\"\u003efrom\u003c/a\u003e -\u003e \u003ca href=\"http://example.org/#remark42__comment-\"\u003eto\u003c/a\u003e\n\n`+
`some text\n\n`+
`\"\u003ci\u003esome parent text with a \u003ca href=\"http://example.org\"\u003elink\u003c/a\u003e and special text:\u0026amp; \u0026lt; \u0026gt; \u0026amp;\u003c/i\u003e\"\n\n`+
`↦ \u003ca href=\"http://example.org/\"\u003e[test title]\u003c/a\u003e","parse_mode":"HTML"}`,
string(res))
res := tb.buildMessage(Request{Comment: c, parent: cp})
assert.Equal(t, `<a href="http://example.org/#remark42__comment-999">from</a> -> <a href="http://example.org/#remark42__comment-">to</a>
some text
"<i>some parent text with a <a href="http://example.org">link</a> and special text:&amp; &lt; &gt; &amp;</i>"
↦ <a href="http://example.org/">[test title]</a>`,
res)
// special case for text with h1-h6 header
ch := store.Comment{Text: "<h1>Hello</h1><h6>World</h6>", ID: "555", Locator: store.Locator{URL: "http://example.org/"}}
ch.User.Name = "from"
res, err = buildMessage(Request{Comment: ch})
assert.NoError(t, err)
assert.Equal(t, `{"text":"\u003ca href=\"http://example.org/#remark42__comment-555\"\u003efrom\u003c/a\u003e\n\n`+
`\u003cb\u003eHello\u003c/b\u003e\u003ci\u003e\u003cb\u003eWorld\u003c/b\u003e\u003c/i\u003e`+
`","parse_mode":"HTML"}`,
string(res))
res = tb.buildMessage(Request{Comment: ch})
assert.Equal(t, `<a href="http://example.org/#remark42__comment-555">from</a>
<b>Hello</b><i><b>World</b></i>`,
res)
}
func TestTelegram_SendVerification(t *testing.T) {
ts := mockTelegramServer(nil)
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "good-token",
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
assert.NotNil(t, tb)
// empty VerificationRequest should return no error and no nothing, as well as any other
assert.NoError(t, tb.SendVerification(context.TODO(), VerificationRequest{}))
}
const getUpdatesResp = `{
"ok": true,
"result": [
{
"update_id": 998,
"message": {
"chat": {
"type": "group"
}
}
},
{
"update_id": 999,
"message": {
"text": "not starting with /start",
"chat": {
"type": "private"
}
}
},
{
"update_id": 1000,
"message": {
"message_id": 4,
"from": {
"id": 313131313,
"is_bot": false,
"first_name": "Joe",
"username": "joe123",
"language_code": "en"
},
"chat": {
"id": 313131313,
"first_name": "Joe",
"username": "joe123",
"type": "private"
},
"date": 1601665548,
"text": "/start token",
"entities": [
{
"offset": 0,
"length": 6,
"type": "bot_command"
}
]
}
}
]
}`
func TestTelegram_GetUpdatesFlow(t *testing.T) {
first := true
ts := mockTelegramServer(func(w http.ResponseWriter, r *http.Request) {
if strings.Contains(r.URL.String(), "sendMessage") {
// respond normally to processUpdates attempt to send message back to user
_, _ = w.Write([]byte("{}"))
return
}
// responses to get updates calls to API
if first {
assert.Equal(t, "", r.URL.Query().Get("offset"))
first = false
} else {
assert.Equal(t, "1001", r.URL.Query().Get("offset"))
}
_, _ = w.Write([]byte(getUpdatesResp))
})
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "xxxsupersecretxxx",
UserNotifications: true,
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
// send request with no offset
upd, err := tb.getUpdates(context.Background())
assert.NoError(t, err)
assert.Len(t, upd.Result, 3)
assert.Equal(t, 1001, tb.updateOffset)
assert.Equal(t, "/start token", upd.Result[len(upd.Result)-1].Message.Text)
tb.AddToken("token", "user", "site", time.Now().Add(time.Minute))
_, _, err = tb.CheckToken("token", "user")
assert.Error(t, err)
tb.processUpdates(context.Background(), upd)
tgID, site, err := tb.CheckToken("token", "user")
assert.NoError(t, err)
assert.Equal(t, "313131313", tgID)
assert.Equal(t, "site", site)
// send request with offset
_, err = tb.getUpdates(context.Background())
assert.NoError(t, err)
}
func TestTelegram_ProcessUpdateFlow(t *testing.T) {
ts := mockTelegramServer(func(w http.ResponseWriter, r *http.Request) {
// respond normally to processUpdates attempt to send message back to user
_, _ = w.Write([]byte("{}"))
})
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "xxxsupersecretxxx",
UserNotifications: true,
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
tb.AddToken("token", "user", "site", time.Now().Add(time.Minute))
tb.AddToken("expired token", "user", "site", time.Now().Add(-time.Minute))
assert.Len(t, tb.requests.data, 2)
_, _, err = tb.CheckToken("token", "user")
assert.Error(t, err)
assert.NoError(t, tb.ProcessUpdate(context.Background(), getUpdatesResp))
assert.Len(t, tb.requests.data, 1, "expired token was cleaned up")
tgID, site, err := tb.CheckToken("token", "user")
assert.NoError(t, err)
assert.Len(t, tb.requests.data, 0, "token is deleted after successful check")
assert.Equal(t, "313131313", tgID)
assert.Equal(t, "site", site)
tb.AddToken("expired token", "user", "site", time.Now().Add(-time.Minute))
assert.Len(t, tb.requests.data, 1)
assert.EqualError(t, tb.ProcessUpdate(context.Background(), ""), "failed to decode provided telegram update: unexpected end of JSON input")
assert.Len(t, tb.requests.data, 0, "expired token should be cleaned up despite the error")
}
const sendMessageResp = `{
"ok": true,
"result": {
"message_id": 100,
"from": {
"id": 666666666,
"is_bot": true,
"first_name": "Test auth bot",
"username": "TestAuthBot"
},
"chat": {
"id": 313131313,
"first_name": "Joe",
"username": "joe123",
"type": "private"
},
"date": 1602430546,
"text": "123"
}
}`
func TestTelegram_SendText(t *testing.T) {
ts := mockTelegramServer(func(w http.ResponseWriter, r *http.Request) {
assert.Equal(t, "123", r.URL.Query().Get("chat_id"))
assert.Equal(t, "hello there", r.URL.Query().Get("text"))
_, _ = w.Write([]byte(sendMessageResp))
})
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "xxxsupersecretxxx",
UserNotifications: true,
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
err = tb.sendText(context.Background(), 123, "hello there")
assert.NoError(t, err)
}
const errorResp = `{"ok":false,"error_code":400,"description":"Very bad request"}`
func TestTelegram_Error(t *testing.T) {
ts := mockTelegramServer(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusBadRequest)
_, _ = w.Write([]byte(errorResp))
})
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "xxxsupersecretxxx",
UserNotifications: true,
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
_, err = tb.getUpdates(context.Background())
assert.EqualError(t, err, "failed to fetch updates: unexpected telegram API status code 400, error: \"Very bad request\"")
}
func TestTelegram_TokenVerification(t *testing.T) {
ts := mockTelegramServer(func(w http.ResponseWriter, r *http.Request) {
if strings.Contains(r.URL.String(), "sendMessage") {
// respond normally to processUpdates attempt to send message back to user
_, _ = w.Write([]byte("{}"))
return
}
// responses to get updates calls to API
_, _ = w.Write([]byte(getUpdatesResp))
})
defer ts.Close()
tb, err := NewTelegram(TelegramParams{
AdminChannelID: "remark_test",
Token: "good-token",
apiPrefix: ts.URL + "/",
})
assert.NoError(t, err)
assert.NotNil(t, tb)
tb.AddToken("token", "user", "site", time.Now().Add(time.Minute))
assert.Len(t, tb.requests.data, 1)
// wrong token
tgID, site, err := tb.CheckToken("unknown token", "user")
assert.Empty(t, tgID)
assert.Empty(t, site)
assert.EqualError(t, err, "request is not found")
// right token and user, not verified yet
tgID, site, err = tb.CheckToken("token", "user")
assert.Empty(t, tgID)
assert.Empty(t, site)
assert.EqualError(t, err, "request is not verified yet")
// confirm request
authRequest, ok := tb.requests.data["token"]
assert.True(t, ok)
authRequest.confirmed = true
authRequest.telegramID = "telegramID"
tb.requests.data["token"] = authRequest
// wrong user
tgID, site, err = tb.CheckToken("token", "wrong user")
assert.Empty(t, tgID)
assert.Empty(t, site)
assert.EqualError(t, err, "user does not match original requester")
// successful check
tgID, site, err = tb.CheckToken("token", "user")
assert.NoError(t, err)
assert.Equal(t, "telegramID", tgID)
assert.Equal(t, "site", site)
// expired token
tb.AddToken("expired token", "user", "site", time.Now().Add(-time.Minute))
tgID, site, err = tb.CheckToken("expired token", "user")
assert.Empty(t, tgID)
assert.Empty(t, site)
assert.EqualError(t, err, "request expired")
assert.Len(t, tb.requests.data, 0)
// expired token, cleaned up by the cleanup
tb.apiPollInterval = time.Millisecond * 15
tb.expiredCleanupInterval = time.Millisecond * 10
ctx, cancel := context.WithCancel(context.Background())
go tb.Run(ctx)
assert.Eventually(t, func() bool {
return tb.ProcessUpdate(ctx, "").Error() == "the Run goroutine should not be used with ProcessUpdate"
}, time.Millisecond*100, time.Millisecond*10, "ProcessUpdate should not work same time as Run")
tb.AddToken("expired token", "user", "site", time.Now().Add(-time.Minute))
tb.requests.RLock()
assert.Len(t, tb.requests.data, 1)
tb.requests.RUnlock()
time.Sleep(tb.expiredCleanupInterval * 2)
tb.requests.RLock()
assert.Len(t, tb.requests.data, 0)
tb.requests.RUnlock()
cancel()
// give enough time for Run() to finish
time.Sleep(tb.expiredCleanupInterval)
}
const getMeResp = `{"ok": true,
"result": {
"first_name": "comments_test",
"id": 707381019,
"is_bot": true,
"username": "remark42_test_bot"
}}`
func mockTelegramServer(h http.HandlerFunc) *httptest.Server {
if h != nil {
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if strings.Contains(r.URL.String(), "getMe") {
_, _ = w.Write([]byte(getMeResp))
return
}
h(w, r)
}))
}
router := chi.NewRouter()
router.Get("/good-token/getMe", func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(getMeResp))
})
router.Get("/empty-json/getMe", func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(`{}`))
})
router.Get("/non-json-resp/getMe", func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(`not-a-json`))
})
router.Get("/404/getMe", func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(404)
})
router.Post("/good-token/sendMessage", func(w http.ResponseWriter, r *http.Request) {
_, _ = w.Write([]byte(`{"ok": true}`))
})
return httptest.NewServer(router)
tb := Telegram{}
// empty VerificationRequest should return no error and do nothing, as well as any other
assert.NoError(t, tb.SendVerification(context.Background(), VerificationRequest{}))
}
+2 -3
View File
@@ -11,8 +11,7 @@ import (
"time"
log "github.com/go-pkgz/lgr"
"github.com/umputun/remark42/backend/app/notify"
ntf "github.com/go-pkgz/notify"
)
type tgRequester interface {
@@ -45,7 +44,7 @@ func DispatchTelegramUpdates(ctx context.Context, requester tgRequester, receive
url += fmt.Sprintf("&offset=%d", updateOffset)
}
var update notify.TelegramUpdate
var update ntf.TelegramUpdate
err := requester.Request(ctx, url, nil, &update)
if err != nil {
+2 -3
View File
@@ -7,9 +7,8 @@ import (
"testing"
"time"
ntf "github.com/go-pkgz/notify"
"github.com/stretchr/testify/assert"
"github.com/umputun/remark42/backend/app/notify"
)
func TestDispatchTelegramUpdates(t *testing.T) {
@@ -59,7 +58,7 @@ func (m *mockTGUpdatesReceiver) String() string {
}
func (m *mockTGUpdatesReceiver) ProcessUpdate(_ context.Context, textUpdate string) error {
var result notify.TelegramUpdate
var result ntf.TelegramUpdate
err := json.Unmarshal([]byte(textUpdate), &result)
assert.NoError(m.t, err)
if m.hit < 2 {