Initial commit
Build & Push / build (push) Canceled after 0s

This commit is contained in:
2026-08-05 19:31:35 +03:00
commit a33fa24e99
272 changed files with 40566 additions and 0 deletions
@@ -0,0 +1,45 @@
package database
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgtype"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (d *Database) CreatePost(ctx context.Context, post domain.Post) error {
query := `INSERT INTO post (id, channel_id, message_id, text, published_at)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (channel_id, message_id) DO UPDATE
SET text = EXCLUDED.text,
published_at = COALESCE(post.published_at, EXCLUDED.published_at),
updated_at = CURRENT_TIMESTAMP;`
txOrPool := transaction.TryExtractTX(ctx)
dto := createPostDTO{
ID: pgtype.UUID{Bytes: post.ID, Valid: true},
ChannelID: pgtype.UUID{Bytes: post.ChannelID, Valid: true},
MessageID: pgtype.Int4{Int32: int32(post.MessageID), Valid: true},
Text: pgtype.Text{String: post.Text, Valid: true},
PublishedAt: pgtype.Timestamptz{Time: post.PublishedAt, Valid: !post.PublishedAt.IsZero()},
}
_, err := txOrPool.Exec(ctx, query, dto.ID, dto.ChannelID, dto.MessageID, dto.Text, dto.PublishedAt)
if err != nil {
return fmt.Errorf("txOrPool.Exec: %w", err)
}
return nil
}
type createPostDTO struct {
ID pgtype.UUID
ChannelID pgtype.UUID
MessageID pgtype.Int4
Text pgtype.Text
PublishedAt pgtype.Timestamptz
}
@@ -0,0 +1,40 @@
package database
import (
"context"
"fmt"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (d *Database) CreateViewsSnapshot(ctx context.Context, snapshot domain.ViewsSnapshot) error {
query := `INSERT INTO post_views_history (id, views_count, fetched_at, post_id)
VALUES ($1, $2, $3, $4);`
txOrPool := transaction.TryExtractTX(ctx)
dto := createViewsSnapshotDTO{
ID: pgtype.UUID{Bytes: uuid.New(), Valid: true},
ViewsCount: pgtype.Int4{Int32: int32(snapshot.ViewsCount), Valid: true},
FetchedAt: pgtype.Timestamptz{Time: snapshot.FetchedAt, Valid: true},
PostID: pgtype.UUID{Bytes: snapshot.PostID, Valid: true},
}
_, err := txOrPool.Exec(ctx, query, dto.ID, dto.ViewsCount, dto.FetchedAt, dto.PostID)
if err != nil {
return fmt.Errorf("txOrPool.Exec: %w", err)
}
return nil
}
type createViewsSnapshotDTO struct {
ID pgtype.UUID
ViewsCount pgtype.Int4
FetchedAt pgtype.Timestamptz
PostID pgtype.UUID
}
@@ -0,0 +1,7 @@
package database
type Database struct{}
func New() *Database {
return &Database{}
}
@@ -0,0 +1,37 @@
package database
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgtype"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (d *Database) DeletePost(ctx context.Context, p domain.Post) error {
query := `UPDATE post
SET deleted_from_channel_at = CURRENT_TIMESTAMP,
updated_at = CURRENT_TIMESTAMP
WHERE channel_id = $1 AND message_id = $2;`
txOrPool := transaction.TryExtractTX(ctx)
dto := deletePostDTO{
ChannelID: pgtype.UUID{Bytes: p.ChannelID, Valid: true},
MessageID: pgtype.Int4{Int32: int32(p.MessageID), Valid: true},
}
_, err := txOrPool.Exec(ctx, query, dto.ChannelID, dto.MessageID)
if err != nil {
return fmt.Errorf("txOrPool.Exec: %w", err)
}
return nil
}
type deletePostDTO struct {
ChannelID pgtype.UUID
MessageID pgtype.Int4
}
@@ -0,0 +1,88 @@
package database
import (
"context"
"github.com/jackc/pgx/v5/pgtype"
"github.com/rs/zerolog/log"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (d *Database) GetChannels(ctx context.Context) []domain.Channel {
query := `SELECT
id,
telegram_id,
username,
title,
access_hash,
pts,
invite_link,
is_accessible
FROM channel
WHERE deleted_at IS NULL
AND is_accessible = true;`
txOrPool := transaction.TryExtractTX(ctx)
rows, err := txOrPool.Query(ctx, query)
if err != nil {
log.Error().Err(err).Msg("txOrPool.Query failed")
return []domain.Channel{}
}
defer rows.Close()
channels := make([]domain.Channel, 0)
for rows.Next() {
var dto getChannelsDTO
err = rows.Scan(dto.destination()...)
if err != nil {
log.Error().Err(err).Msg("rows.Scan failed")
return []domain.Channel{}
}
channels = append(channels, dto.toDomain())
}
return channels
}
type getChannelsDTO struct {
ID pgtype.UUID
TelegramID pgtype.Int8
Username pgtype.Text
Title pgtype.Text
AccessHash pgtype.Int8
Pts pgtype.Int4
InviteLink pgtype.Text
IsAccessible pgtype.Bool
}
func (dto *getChannelsDTO) destination() []any {
return []any{
&dto.ID,
&dto.TelegramID,
&dto.Username,
&dto.Title,
&dto.AccessHash,
&dto.Pts,
&dto.InviteLink,
&dto.IsAccessible,
}
}
func (dto *getChannelsDTO) toDomain() domain.Channel {
return domain.Channel{
ID: dto.ID.Bytes,
TelegramID: domain.NormalizeChatID(dto.TelegramID.Int64),
Username: dto.Username.String,
Title: dto.Title.String,
AccessHash: dto.AccessHash.Int64,
Pts: int(dto.Pts.Int32),
InviteLink: dto.InviteLink.String,
IsAccessible: dto.IsAccessible.Bool,
}
}
@@ -0,0 +1,89 @@
package database
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgtype"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (d *Database) GetChannelsWithTrackedPosts(ctx context.Context) ([]domain.Channel, error) {
query := `SELECT DISTINCT
c.id,
c.telegram_id,
c.username,
c.title,
c.access_hash,
c.pts,
c.invite_link,
c.is_accessible
FROM channel c
INNER JOIN placement p ON p.channel_id = c.id
INNER JOIN placement_post pp ON pp.placement_id = p.id
INNER JOIN post po ON po.id = pp.post_id
WHERE po.deleted_from_channel_at IS NULL
AND c.is_accessible = true;`
txOrPool := transaction.TryExtractTX(ctx)
rows, err := txOrPool.Query(ctx, query)
if err != nil {
return nil, fmt.Errorf("txOrPool.Query: %w", err)
}
defer rows.Close()
channels := make([]domain.Channel, 0)
for rows.Next() {
var dto getChannelsWithTrackedPostsDTO
err = rows.Scan(dto.destination()...)
if err != nil {
return nil, fmt.Errorf("rows.Scan: %w", err)
}
channels = append(channels, dto.toDomain())
}
return channels, nil
}
type getChannelsWithTrackedPostsDTO struct {
ID pgtype.UUID
TelegramID pgtype.Int8
Username pgtype.Text
Title pgtype.Text
AccessHash pgtype.Int8
Pts pgtype.Int4
InviteLink pgtype.Text
IsAccessible pgtype.Bool
}
func (dto *getChannelsWithTrackedPostsDTO) destination() []any {
return []any{
&dto.ID,
&dto.TelegramID,
&dto.Username,
&dto.Title,
&dto.AccessHash,
&dto.Pts,
&dto.InviteLink,
&dto.IsAccessible,
}
}
func (dto *getChannelsWithTrackedPostsDTO) toDomain() domain.Channel {
return domain.Channel{
ID: dto.ID.Bytes,
TelegramID: domain.NormalizeChatID(dto.TelegramID.Int64),
Username: dto.Username.String,
Title: dto.Title.String,
AccessHash: dto.AccessHash.Int64,
Pts: int(dto.Pts.Int32),
InviteLink: dto.InviteLink.String,
IsAccessible: dto.IsAccessible.Bool,
}
}
@@ -0,0 +1,91 @@
package database
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgtype"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (d *Database) GetTrackedPosts(ctx context.Context, channel domain.Channel) ([]domain.Post, error) {
query := `SELECT DISTINCT
p.id,
p.channel_id,
p.message_id,
p.text,
CASE
WHEN c.username IS NOT NULL THEN CONCAT('https://t.me/', c.username, '/', p.message_id)
WHEN c.telegram_id IS NOT NULL THEN CONCAT('https://t.me/c/', (-c.telegram_id - 1000000000000), '/', p.message_id)
ELSE ''
END as link,
COALESCE(
(SELECT pvh.views_count
FROM post_views_history pvh
WHERE pvh.post_id = p.id
ORDER BY pvh.fetched_at DESC
LIMIT 1),
0
) as views
FROM post p
INNER JOIN channel c ON c.id = p.channel_id
INNER JOIN placement_post pp ON pp.post_id = p.id
WHERE p.channel_id = $1
AND p.deleted_from_channel_at IS NULL;`
txOrPool := transaction.TryExtractTX(ctx)
rows, err := txOrPool.Query(ctx, query, pgtype.UUID{Bytes: channel.ID, Valid: true})
if err != nil {
return nil, fmt.Errorf("txOrPool.Query: %w", err)
}
defer rows.Close()
posts := make([]domain.Post, 0)
for rows.Next() {
var dto getTrackedPostsDTO
err = rows.Scan(dto.destination()...)
if err != nil {
return nil, fmt.Errorf("rows.Scan: %w", err)
}
posts = append(posts, dto.toDomain())
}
return posts, nil
}
type getTrackedPostsDTO struct {
ID pgtype.UUID
ChannelID pgtype.UUID
MessageID pgtype.Int4
Text pgtype.Text
Link pgtype.Text
Views pgtype.Int4
}
func (dto *getTrackedPostsDTO) destination() []any {
return []any{
&dto.ID,
&dto.ChannelID,
&dto.MessageID,
&dto.Text,
&dto.Link,
&dto.Views,
}
}
func (dto *getTrackedPostsDTO) toDomain() domain.Post {
return domain.Post{
ID: dto.ID.Bytes,
ChannelID: dto.ChannelID.Bytes,
MessageID: int(dto.MessageID.Int32),
Text: dto.Text.String,
Link: dto.Link.String,
Views: int(dto.Views.Int32),
}
}
@@ -0,0 +1,124 @@
package database
import (
"context"
"fmt"
"strings"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
"github.com/jackc/pgx/v5/pgtype"
)
func (d *Database) UpdateChannelIfNotAccessible(ctx context.Context, channel domain.Channel) error {
query := `UPDATE channel
SET telegram_id = $2,
username = $3,
title = $4,
access_hash = $5,
pts = 0,
is_accessible = TRUE,
invite_link = $6,
updated_at = CURRENT_TIMESTAMP
WHERE telegram_id = $1
AND is_accessible = FALSE;`
txOrPool := transaction.TryExtractTX(ctx)
username := strings.TrimSpace(channel.Username)
var usernameDTO pgtype.Text
if username != "" {
usernameDTO = pgtype.Text{String: username, Valid: true}
} else {
usernameDTO = pgtype.Text{Valid: false}
}
inviteLink := strings.TrimSpace(channel.InviteLink)
var inviteLinkDTO pgtype.Text
if inviteLink != "" {
inviteLinkDTO = pgtype.Text{String: inviteLink, Valid: true}
} else {
inviteLinkDTO = pgtype.Text{Valid: false}
}
result, err := txOrPool.Exec(ctx, query,
channel.TelegramID,
pgtype.Int8{Int64: channel.TelegramID, Valid: true},
usernameDTO,
pgtype.Text{String: channel.Title, Valid: true},
pgtype.Int8{Int64: channel.AccessHash, Valid: true},
inviteLinkDTO,
)
if err != nil {
return fmt.Errorf("txOrPool.Exec: %w", err)
}
// Log if no rows were updated (channel either doesn't exist or is already accessible)
rowsAffected := result.RowsAffected()
if rowsAffected == 0 {
// Channel doesn't exist or is already accessible - this is fine
return nil
}
return nil
}
func (d *Database) UpdateChannel(ctx context.Context, channel domain.Channel) error {
query := `UPDATE channel
SET telegram_id = $2,
username = $3,
title = $4,
access_hash = $5,
pts = $6,
is_accessible = $7,
invite_link = $8,
updated_at = CURRENT_TIMESTAMP
WHERE id = $1;`
txOrPool := transaction.TryExtractTX(ctx)
username := strings.TrimSpace(channel.Username)
var usernameDTO pgtype.Text
if username != "" {
usernameDTO = pgtype.Text{String: username, Valid: true}
} else {
usernameDTO = pgtype.Text{Valid: false}
}
inviteLink := strings.TrimSpace(channel.InviteLink)
var inviteLinkDTO pgtype.Text
if inviteLink != "" {
inviteLinkDTO = pgtype.Text{String: inviteLink, Valid: true}
} else {
inviteLinkDTO = pgtype.Text{Valid: false}
}
dto := updateChannelDTO{
ID: pgtype.UUID{Bytes: channel.ID, Valid: true},
TelegramID: pgtype.Int8{Int64: channel.TelegramID, Valid: true},
Username: usernameDTO,
Title: pgtype.Text{String: channel.Title, Valid: true},
AccessHash: pgtype.Int8{Int64: channel.AccessHash, Valid: true},
Pts: pgtype.Int4{Int32: int32(channel.Pts), Valid: true},
IsAccessible: pgtype.Bool{Bool: channel.IsAccessible, Valid: true},
InviteLink: inviteLinkDTO,
}
_, err := txOrPool.Exec(ctx, query, dto.ID, dto.TelegramID, dto.Username, dto.Title, dto.AccessHash, dto.Pts, dto.IsAccessible, dto.InviteLink)
if err != nil {
return fmt.Errorf("txOrPool.Exec: %w", err)
}
return nil
}
type updateChannelDTO struct {
ID pgtype.UUID
TelegramID pgtype.Int8
Username pgtype.Text
Title pgtype.Text
AccessHash pgtype.Int8
Pts pgtype.Int4
IsAccessible pgtype.Bool
InviteLink pgtype.Text
}
@@ -0,0 +1,37 @@
package telegram
import (
"context"
"fmt"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
"github.com/gotd/td/tg"
)
func (t *Telegram) UpdatePostsViews(ctx context.Context, channel domain.Channel, posts []domain.Post) error {
ids := make([]int, len(posts))
for i, p := range posts {
ids[i] = p.MessageID
}
req := &tg.MessagesGetMessagesViewsRequest{
Peer: &tg.InputPeerChannel{
ChannelID: channel.ChannelID(),
AccessHash: channel.AccessHash,
},
ID: ids,
Increment: false,
}
resp, err := t.API().MessagesGetMessagesViews(ctx, req)
if err != nil {
return fmt.Errorf("get messages views: %w", err)
}
// Telegram гарантирует, что resp.Views соответствует порядку ids
for i, v := range resp.Views {
posts[i].Views = v.Views
}
return nil
}
@@ -0,0 +1,125 @@
package telegram
import (
"context"
"fmt"
"time"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
"github.com/gotd/td/tg"
)
func (t *Telegram) GetChannelDiff(ctx context.Context, channel domain.Channel, limit int) (domain.ChannelDiff, error) {
req := &tg.UpdatesGetChannelDifferenceRequest{
Channel: &tg.InputChannel{
ChannelID: channel.ChannelID(),
AccessHash: channel.AccessHash,
},
Filter: &tg.ChannelMessagesFilterEmpty{},
Pts: channel.Pts,
Limit: limit,
}
rawDiff, err := t.API().UpdatesGetChannelDifference(ctx, req)
if err != nil {
return domain.ChannelDiff{}, fmt.Errorf("get difference: %w", err)
}
result := domain.ChannelDiff{}
switch d := rawDiff.(type) {
case *tg.UpdatesChannelDifferenceEmpty:
result.NewPts = d.Pts
case *tg.UpdatesChannelDifferenceTooLong:
if dialog, ok := d.Dialog.(*tg.Dialog); ok {
result.NewPts = dialog.Pts
}
result.NewPosts = extractPosts(channel, d.Messages)
result.UpdatedChannel = extractChannelMeta(channel, d.Chats)
case *tg.UpdatesChannelDifference:
result.NewPts = d.Pts
result.NewPosts = extractPosts(channel, d.NewMessages)
result.DeletedPosts = extractDeletedPosts(channel, d.OtherUpdates)
result.UpdatedChannel = extractChannelMeta(channel, d.Chats)
default:
return domain.ChannelDiff{}, fmt.Errorf("unexpected rawDiff type: %T", rawDiff)
}
return result, nil
}
func extractPosts(channel domain.Channel, msgs []tg.MessageClass) []domain.Post {
posts := make([]domain.Post, 0, len(msgs))
for _, raw := range msgs {
m, ok := raw.(*tg.Message)
if !ok {
continue
}
text := messageToHTML(m.Message, m.Entities)
publishedAt := time.Unix(int64(m.Date), 0).UTC()
p := domain.NewPost(channel, m.ID, text, m.Views, publishedAt)
posts = append(posts, p)
}
return posts
}
func extractDeletedPosts(channel domain.Channel, updates []tg.UpdateClass) []domain.Post {
var deleted []domain.Post
for _, upd := range updates {
if u, ok := upd.(*tg.UpdateDeleteChannelMessages); ok {
for _, msgID := range u.Messages {
p := domain.NewPost(channel, msgID, "", 0, time.Time{})
deleted = append(deleted, p)
}
}
}
return deleted
}
func extractChannelMeta(currentChannel domain.Channel, chats []tg.ChatClass) *domain.Channel {
for _, chat := range chats {
ch, ok := chat.(*tg.Channel)
if !ok {
continue
}
if ch.ID != currentChannel.ChannelID() {
continue
}
// Preserve username if Telegram doesn't provide a new one (private channels may not have username)
username := currentChannel.Username
if usernameVal, ok := ch.GetUsername(); ok && usernameVal != "" {
username = usernameVal
}
// Keep current access hash if new one is not provided or is zero
accessHash := currentChannel.AccessHash
if accessHashVal, ok := ch.GetAccessHash(); ok && accessHashVal != 0 {
accessHash = accessHashVal
}
updated := domain.Channel{
ID: currentChannel.ID,
TelegramID: domain.ChatIDFromChannelID(ch.ID),
Username: username,
Title: ch.Title,
AccessHash: accessHash,
InviteLink: currentChannel.InviteLink, // Preserve invite link (never returned in channel diff)
Pts: currentChannel.Pts,
IsAccessible: currentChannel.IsAccessible,
}
return &updated
}
return nil
}
@@ -0,0 +1,49 @@
package telegram
import (
"context"
"fmt"
"time"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
"github.com/gotd/td/tg"
"github.com/rs/zerolog/log"
)
func (t *Telegram) GetChannelHistory(ctx context.Context, channel domain.Channel, limit int) ([]domain.Post, error) {
log.Info().Msg(channel.String())
req := &tg.MessagesGetHistoryRequest{
Peer: &tg.InputPeerChannel{
ChannelID: channel.ChannelID(),
AccessHash: channel.AccessHash,
},
Limit: limit,
}
resp, err := t.API().MessagesGetHistory(ctx, req)
if err != nil {
return nil, fmt.Errorf("resp: %w", err)
}
messages, ok := resp.(*tg.MessagesChannelMessages)
if !ok {
return nil, nil
}
posts := make([]domain.Post, 0, len(messages.Messages))
for _, raw := range messages.Messages {
m, ok := raw.(*tg.Message)
if !ok {
continue
}
text := messageToHTML(m.Message, m.Entities)
publishedAt := time.Unix(int64(m.Date), 0).UTC()
p := domain.NewPost(channel, m.ID, text, m.Views, publishedAt)
posts = append(posts, p)
}
return posts, nil
}
@@ -0,0 +1,197 @@
package telegram
import (
"html"
"sort"
"strconv"
"strings"
"unicode/utf8"
"github.com/gotd/td/tg"
)
type htmlEntity struct {
start int
end int
openTag string
closeTag string
length int
}
type htmlEvent struct {
pos int
tag string
isStart bool
length int
}
func messageToHTML(text string, entities []tg.MessageEntityClass) string {
if text == "" {
return ""
}
if len(entities) == 0 {
return html.EscapeString(text)
}
ranges := make([]htmlEntity, 0, len(entities))
needed := make([]int, 0, len(entities)*2)
for _, raw := range entities {
offset, length, openTag, closeTag, ok := htmlEntityMeta(raw)
if !ok {
continue
}
start := offset
end := offset + length
if length <= 0 {
continue
}
ranges = append(ranges, htmlEntity{
start: start,
end: end,
openTag: openTag,
closeTag: closeTag,
length: length,
})
needed = append(needed, start, end)
}
if len(ranges) == 0 {
return html.EscapeString(text)
}
positions := utf16PositionsToBytes(text, needed)
events := make([]htmlEvent, 0, len(ranges)*2)
for _, r := range ranges {
startByte, okStart := positions[r.start]
endByte, okEnd := positions[r.end]
if !okStart || !okEnd || startByte > endByte {
continue
}
events = append(events, htmlEvent{
pos: startByte,
tag: r.openTag,
isStart: true,
length: r.length,
})
events = append(events, htmlEvent{
pos: endByte,
tag: r.closeTag,
isStart: false,
length: r.length,
})
}
sort.SliceStable(events, func(i, j int) bool {
if events[i].pos != events[j].pos {
return events[i].pos < events[j].pos
}
if events[i].isStart != events[j].isStart {
return !events[i].isStart
}
if events[i].isStart {
return events[i].length > events[j].length
}
return events[i].length < events[j].length
})
var b strings.Builder
last := 0
for _, ev := range events {
if ev.pos > last {
b.WriteString(html.EscapeString(text[last:ev.pos]))
}
b.WriteString(ev.tag)
last = ev.pos
}
if last < len(text) {
b.WriteString(html.EscapeString(text[last:]))
}
return b.String()
}
func htmlEntityMeta(entity tg.MessageEntityClass) (offset int, length int, openTag string, closeTag string, ok bool) {
switch e := entity.(type) {
case *tg.MessageEntityBold:
return e.Offset, e.Length, "<b>", "</b>", true
case *tg.MessageEntityItalic:
return e.Offset, e.Length, "<i>", "</i>", true
case *tg.MessageEntityUnderline:
return e.Offset, e.Length, "<u>", "</u>", true
case *tg.MessageEntityStrike:
return e.Offset, e.Length, "<s>", "</s>", true
case *tg.MessageEntityCode:
return e.Offset, e.Length, "<code>", "</code>", true
case *tg.MessageEntityPre:
if e.Language != "" {
lang := html.EscapeString(e.Language)
return e.Offset, e.Length, `<pre><code class="language-` + lang + `">`, "</code></pre>", true
}
return e.Offset, e.Length, "<pre>", "</pre>", true
case *tg.MessageEntityTextURL:
url := html.EscapeString(e.URL)
return e.Offset, e.Length, `<a href="` + url + `">`, "</a>", true
case *tg.MessageEntityMentionName:
userID := html.EscapeString(strconv.FormatInt(e.UserID, 10))
return e.Offset, e.Length, `<a href="tg://user?id=` + userID + `">`, "</a>", true
case *tg.MessageEntitySpoiler:
return e.Offset, e.Length, `<span class="tg-spoiler">`, "</span>", true
case *tg.MessageEntityBlockquote:
if e.Collapsed {
return e.Offset, e.Length, `<blockquote expandable>`, "</blockquote>", true
}
return e.Offset, e.Length, "<blockquote>", "</blockquote>", true
case *tg.MessageEntityCustomEmoji:
id := html.EscapeString(strconv.FormatInt(e.DocumentID, 10))
return e.Offset, e.Length, `<tg-emoji emoji-id="` + id + `">`, "</tg-emoji>", true
default:
return 0, 0, "", "", false
}
}
func utf16PositionsToBytes(text string, needed []int) map[int]int {
result := make(map[int]int, len(needed))
needSet := make(map[int]struct{}, len(needed))
for _, n := range needed {
needSet[n] = struct{}{}
}
utf16Pos := 0
if _, ok := needSet[0]; ok {
result[0] = 0
}
for i, r := range text {
if _, ok := needSet[utf16Pos]; ok {
result[utf16Pos] = i
}
step := utf16RuneLen(r)
if step == 2 {
if _, ok := needSet[utf16Pos+1]; ok {
result[utf16Pos+1] = i
}
}
utf16Pos += step
}
if _, ok := needSet[utf16Pos]; ok {
result[utf16Pos] = len(text)
}
return result
}
func utf16RuneLen(r rune) int {
const surrSelf = 0x10000
if r >= surrSelf && r <= utf8.MaxRune {
return 2
}
return 1
}
@@ -0,0 +1,33 @@
package telegram
import (
"context"
"fmt"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
"github.com/gotd/td/tg"
)
func (t *Telegram) GetChannelPTS(ctx context.Context, channel domain.Channel) (int, error) {
req := []tg.InputDialogPeerClass{
&tg.InputDialogPeer{
Peer: &tg.InputPeerChannel{
ChannelID: channel.ChannelID(),
AccessHash: channel.AccessHash,
},
},
}
dialogs, err := t.API().MessagesGetPeerDialogs(ctx, req)
if err != nil {
return 0, fmt.Errorf("peer dialogs: %w", err)
}
if len(dialogs.Dialogs) > 0 {
if d, ok := dialogs.Dialogs[0].(*tg.Dialog); ok {
return d.Pts, nil
}
}
return 0, nil
}
@@ -0,0 +1,58 @@
package telegram
import (
"context"
"fmt"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
"github.com/gotd/td/tg"
)
func (t *Telegram) ParseChannelMeta(ctx context.Context, username string) (domain.Channel, error) {
resolved, err := t.API().ContactsResolveUsername(ctx, &tg.ContactsResolveUsernameRequest{
Username: username,
})
if err != nil {
return domain.Channel{}, fmt.Errorf("resolve username: %w", err)
}
if len(resolved.Chats) == 0 {
return domain.Channel{}, fmt.Errorf("no chats in resolve result")
}
ch, ok := resolved.Chats[0].(*tg.Channel)
if !ok {
return domain.Channel{}, fmt.Errorf("not a channel")
}
req := []tg.InputDialogPeerClass{
&tg.InputDialogPeer{
Peer: &tg.InputPeerChannel{
ChannelID: ch.ID,
AccessHash: ch.AccessHash,
},
},
}
dialogs, err := t.API().MessagesGetPeerDialogs(ctx, req)
if err != nil {
return domain.Channel{}, fmt.Errorf("peer dialogs: %w", err)
}
pts := 0
if len(dialogs.Dialogs) > 0 {
d, ok := dialogs.Dialogs[0].(*tg.Dialog)
if ok {
pts = d.Pts
}
}
return domain.Channel{
TelegramID: domain.ChatIDFromChannelID(ch.ID),
Username: username,
Title: ch.Title,
AccessHash: ch.AccessHash,
Pts: pts,
IsAccessible: true,
}, nil
}
@@ -0,0 +1,133 @@
package telegram
import (
"context"
"fmt"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
"github.com/gotd/td/telegram/deeplink"
"github.com/gotd/td/tg"
)
func (t *Telegram) ParseChannelMetaByInvite(ctx context.Context, inviteLink string) (domain.Channel, error) {
link, err := deeplink.Parse(inviteLink)
if err != nil {
return domain.Channel{}, fmt.Errorf("parse invite link: %w", err)
}
if link.Type != deeplink.Join {
return domain.Channel{}, fmt.Errorf("invite link is not a join link")
}
hash := link.Args.Get("invite")
if hash == "" {
return domain.Channel{}, fmt.Errorf("invite link missing hash")
}
info, err := t.API().MessagesCheckChatInvite(ctx, hash)
if err != nil {
return domain.Channel{}, fmt.Errorf("check invite: %w", err)
}
var channel *tg.Channel
switch v := info.(type) {
case *tg.ChatInviteAlready:
channel = extractChannelFromChat(v.Chat)
case *tg.ChatInvite:
updates, err := t.API().MessagesImportChatInvite(ctx, hash)
if err != nil {
return domain.Channel{}, fmt.Errorf("import invite: %w", err)
}
channel = extractChannelFromUpdates(updates)
case *tg.ChatInvitePeek:
// ChatInvitePeek means we can preview the channel without joining (public channels)
channel = extractChannelFromChat(v.Chat)
default:
return domain.Channel{}, fmt.Errorf("unexpected invite response: %T", v)
}
if channel == nil {
return domain.Channel{}, fmt.Errorf("no channel in invite response")
}
if channel.AccessHash == 0 {
return domain.Channel{}, fmt.Errorf("channel access hash missing")
}
pts, err := getChannelPTS(ctx, t.API(), channel)
if err != nil {
return domain.Channel{}, fmt.Errorf("get peer dialogs: %w", err)
}
username := ""
if usernameVal, ok := channel.GetUsername(); ok && usernameVal != "" {
username = usernameVal
}
return domain.Channel{
TelegramID: domain.ChatIDFromChannelID(channel.ID),
Username: username,
Title: channel.Title,
AccessHash: channel.AccessHash,
Pts: pts,
InviteLink: inviteLink,
IsAccessible: true,
}, nil
}
func extractChannelFromChat(chat tg.ChatClass) *tg.Channel {
switch v := chat.(type) {
case *tg.Channel:
return v
case *tg.ChannelForbidden:
return &tg.Channel{
ID: v.ID,
AccessHash: v.AccessHash,
Title: v.Title,
}
default:
return nil
}
}
func extractChannelFromUpdates(updates tg.UpdatesClass) *tg.Channel {
var chats []tg.ChatClass
switch v := updates.(type) {
case *tg.Updates:
chats = v.Chats
case *tg.UpdatesCombined:
chats = v.Chats
default:
return nil
}
for _, chat := range chats {
if channel := extractChannelFromChat(chat); channel != nil {
return channel
}
}
return nil
}
func getChannelPTS(ctx context.Context, api *tg.Client, channel *tg.Channel) (int, error) {
req := []tg.InputDialogPeerClass{
&tg.InputDialogPeer{
Peer: &tg.InputPeerChannel{
ChannelID: channel.ID,
AccessHash: channel.AccessHash,
},
},
}
dialogs, err := api.MessagesGetPeerDialogs(ctx, req)
if err != nil {
return 0, err
}
if len(dialogs.Dialogs) > 0 {
if d, ok := dialogs.Dialogs[0].(*tg.Dialog); ok {
return d.Pts, nil
}
}
return 0, nil
}
@@ -0,0 +1,20 @@
package telegram
import (
"github.com/TelegramExchange/pkg/telegram"
"github.com/gotd/td/tg"
)
type Telegram struct {
client *telegram.Client
}
func New(client *telegram.Client) *Telegram {
return &Telegram{
client: client,
}
}
func (t *Telegram) API() *tg.Client {
return t.client.API()
}
+91
View File
@@ -0,0 +1,91 @@
package app
import (
"context"
"errors"
"fmt"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/rs/zerolog/log"
"github.com/TelegramExchange/pkg/postgres"
"github.com/TelegramExchange/pkg/telegram"
"github.com/TelegramExchange/pkg/transaction"
"github.com/TelegramExchange/tgex-backend/tg_parser/config"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/adapter/database"
tgadapter "github.com/TelegramExchange/tgex-backend/tg_parser/internal/adapter/telegram"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/controller/httpserver"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/controller/worker"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/usecase"
)
func Run(ctx context.Context, c config.Config) error {
// Telegram client
tgClient, err := telegram.New(c.Telegram)
if err != nil {
return fmt.Errorf("telegram.New: %w", err)
}
// PostgreSQL
pgPool, err := postgres.New(ctx, c.Postgres)
if err != nil {
return fmt.Errorf("broker.New: %w", err)
}
defer pgPool.Close()
transaction.Init(pgPool)
// Adapters
tg := tgadapter.New(tgClient)
db := database.New()
// UseCase
uc := usecase.New(tg, db)
// Controllers
channelWorker := worker.NewChannelWorker(uc, c.ChannelWorker)
viewsWorker := worker.NewViewsWorker(uc, c.ViewsWorker)
httpServer := httpserver.New(uc, c.HTTP.Addr)
serverErr := make(chan error, 1)
go func() {
err := httpServer.ListenAndServe()
if err != nil && !errors.Is(err, http.ErrServerClosed) {
serverErr <- err
}
}()
log.Info().Str("addr", c.HTTP.Addr).Msg("HTTP server started")
log.Info().Msg("App started")
sig := make(chan os.Signal, 1)
signal.Notify(sig, os.Interrupt, syscall.SIGTERM)
select {
case <-sig:
case err := <-serverErr:
return fmt.Errorf("http server: %w", err)
}
log.Info().Msg("App got signal to stop")
// Controllers
viewsWorker.Stop()
channelWorker.Stop()
ctxShutdown, cancel := context.WithTimeout(ctx, 5*time.Second)
defer cancel()
if err := httpServer.Shutdown(ctxShutdown); err != nil {
log.Error().Err(err).Msg("HTTP server shutdown failed")
}
// Adapters
tgClient.Close()
log.Info().Msg("App stopped")
return nil
}
@@ -0,0 +1,152 @@
package httpserver
import (
"context"
"encoding/json"
"net/http"
"strings"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/usecase"
"github.com/gotd/td/tgerr"
"github.com/rs/zerolog/log"
)
type Server struct {
httpServer *http.Server
}
func New(uc *usecase.UseCase, addr string) *Server {
mux := http.NewServeMux()
mux.HandleFunc("/fetch-telegram-channel", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodGet {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
username := strings.TrimSpace(r.URL.Query().Get("username"))
username = strings.TrimPrefix(username, "@")
if username == "" {
http.Error(w, "missing username", http.StatusBadRequest)
return
}
channel, err := uc.FetchChannelMeta(r.Context(), username)
if err != nil {
if isNotFoundError(err) {
http.Error(w, "channel not found", http.StatusNotFound)
return
}
log.Error().Err(err).Str("username", username).Msg("fetch channel meta failed")
http.Error(w, "internal error", http.StatusInternalServerError)
return
}
resp := struct {
ID string `json:"id"`
TelegramID int64 `json:"telegram_id"`
Username string `json:"username"`
Title string `json:"title"`
AccessHash int64 `json:"access_hash"`
Pts int `json:"pts"`
}{
ID: channel.ID.String(),
TelegramID: channel.TelegramID,
Username: channel.Username,
Title: channel.Title,
AccessHash: channel.AccessHash,
Pts: channel.Pts,
}
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(resp); err != nil {
log.Error().Err(err).Msg("encode channel response")
}
})
mux.HandleFunc("/resolve-channel-by-invite", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
var payload struct {
InviteLink string `json:"invite_link"`
}
if err := json.NewDecoder(r.Body).Decode(&payload); err != nil {
http.Error(w, "invalid payload", http.StatusBadRequest)
return
}
inviteLink := strings.TrimSpace(payload.InviteLink)
if inviteLink == "" {
http.Error(w, "missing invite_link", http.StatusBadRequest)
return
}
channel, err := uc.FetchChannelMetaByInvite(r.Context(), inviteLink)
if err != nil {
if isNotFoundError(err) {
http.Error(w, "channel not found", http.StatusNotFound)
return
}
log.Error().Err(err).Msg("resolve channel by invite failed")
http.Error(w, "internal error", http.StatusInternalServerError)
return
}
resp := struct {
ID string `json:"id"`
TelegramID int64 `json:"telegram_id"`
Username string `json:"username"`
Title string `json:"title"`
AccessHash int64 `json:"access_hash"`
Pts int `json:"pts"`
}{
ID: channel.ID.String(),
TelegramID: channel.TelegramID,
Username: channel.Username,
Title: channel.Title,
AccessHash: channel.AccessHash,
Pts: channel.Pts,
}
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(resp); err != nil {
log.Error().Err(err).Msg("encode channel response")
}
})
return &Server{
httpServer: &http.Server{
Addr: addr,
Handler: mux,
},
}
}
func (s *Server) ListenAndServe() error {
return s.httpServer.ListenAndServe()
}
func (s *Server) Shutdown(ctx context.Context) error {
return s.httpServer.Shutdown(ctx)
}
func isNotFoundError(err error) bool {
if tgerr.Is(
err,
"USERNAME_NOT_OCCUPIED",
"USERNAME_INVALID",
"CHANNEL_INVALID",
"CHANNEL_PRIVATE",
"INVITE_HASH_INVALID",
"INVITE_HASH_EXPIRED",
"INVITE_HASH_EMPTY",
) {
return true
}
msg := err.Error()
return strings.Contains(msg, "no chats in resolve result") ||
strings.Contains(msg, "not a channel") ||
strings.Contains(msg, "invite link")
}
@@ -0,0 +1,85 @@
package worker
import (
"context"
"sync"
"time"
"github.com/rs/zerolog/log"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/usecase"
)
type ChannelConfig struct {
MaxWorkers int `envconfig:"WORKER__MAX_WORKERS" default:"1"`
ChannelDelay time.Duration `envconfig:"WORKER__CHANNEL_DELAY" default:"500ms"`
RequestTimeout time.Duration `envconfig:"WORKER__REQUEST_TIMEOUT" default:"10s"`
MessagesLimit int `envconfig:"WORKER__MESSAGES_LIMIT" default:"20"`
PollInterval time.Duration `envconfig:"WORKER__POLL_INTERVAL" default:"120s"`
}
type ChannelWorker struct {
usecase *usecase.UseCase
config ChannelConfig
stop chan struct{}
done chan struct{}
}
func NewChannelWorker(uc *usecase.UseCase, cfg ChannelConfig) *ChannelWorker {
w := &ChannelWorker{
usecase: uc,
config: cfg,
stop: make(chan struct{}),
done: make(chan struct{}),
}
go w.run()
return w
}
func (w *ChannelWorker) run() {
log.Info().Msg("channel worker: started")
var wg sync.WaitGroup
wg.Add(w.config.MaxWorkers)
for range w.config.MaxWorkers {
go w.spawnWorker(&wg)
}
wg.Wait()
log.Info().Msg("channel worker: stopped")
close(w.done)
}
func (w *ChannelWorker) spawnWorker(wg *sync.WaitGroup) {
defer wg.Done()
limiter := time.NewTicker(w.config.ChannelDelay)
defer limiter.Stop()
poll := time.NewTicker(w.config.PollInterval)
defer poll.Stop()
for {
select {
case <-w.stop:
return
case <-poll.C:
err := w.usecase.FetchChannels(context.Background())
if err != nil {
log.Error().Err(err).Msg("usecase.FetchChannels error")
}
}
}
}
func (w *ChannelWorker) Stop() {
close(w.stop)
<-w.done
}
@@ -0,0 +1,65 @@
package worker
import (
"context"
"time"
"github.com/rs/zerolog/log"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/usecase"
)
type ViewsConfig struct {
Interval time.Duration `envconfig:"WORKER__VIEWS_INTERVAL" default:"1800s"`
RequestTimeout time.Duration `envconfig:"WORKER__REQUEST_TIMEOUT" default:"10s"`
}
type ViewsWorker struct {
usecase *usecase.UseCase
config ViewsConfig
stop chan struct{}
done chan struct{}
}
func NewViewsWorker(uc *usecase.UseCase, cfg ViewsConfig) *ViewsWorker {
w := &ViewsWorker{
usecase: uc,
config: cfg,
stop: make(chan struct{}),
done: make(chan struct{}),
}
go w.run()
return w
}
func (w *ViewsWorker) run() {
log.Info().Msg("views worker: started")
for {
select {
case <-w.stop:
log.Info().Msg("views worker: stopped")
close(w.done)
return
case <-time.After(w.config.Interval):
log.Debug().Dur("interval", w.config.Interval).Msg("views worker: refresh tick")
ctx, cancel := context.WithTimeout(context.Background(), w.config.RequestTimeout)
err := w.usecase.FetchViews(ctx)
if err != nil {
log.Error().Err(err).Msg("views worker: FetchViews error")
} else {
log.Debug().Msg("views worker: refresh finished")
}
cancel()
}
}
}
func (w *ViewsWorker) Stop() {
close(w.stop)
<-w.done
}
+68
View File
@@ -0,0 +1,68 @@
package domain
import (
"fmt"
"github.com/google/uuid"
)
type Channel struct {
ID uuid.UUID
TelegramID int64
Username string
Title string
AccessHash int64
Pts int
InviteLink string
IsAccessible bool
}
const channelIDOffset int64 = 1000000000000
// ChannelID returns the positive channel identifier expected by Telegram API.
func (c Channel) ChannelID() int64 {
return ChannelIDFromChatID(c.TelegramID)
}
// ChannelIDFromChatID converts stored chat IDs (Bot API style) to channel IDs used by TDLib.
func ChannelIDFromChatID(chatID int64) int64 {
if chatID >= 0 {
return chatID
}
return -chatID - channelIDOffset
}
// ChatIDFromChannelID converts Telegram channel IDs to Bot API style chat IDs (-100...).
func ChatIDFromChannelID(channelID int64) int64 {
return -channelIDOffset - channelID
}
// NormalizeChatID ensures that TelegramID is stored in chat-id form (negative).
func NormalizeChatID(id int64) int64 {
if id < 0 {
return id
}
return ChatIDFromChannelID(id)
}
func (c Channel) String() string {
return fmt.Sprintf(
"Channel{id=%s telegram_id=%d username=%q title=%q access_hash=%d pts=%d invite_link=%t}",
c.ID,
c.TelegramID,
c.Username,
c.Title,
c.AccessHash,
c.Pts,
c.InviteLink != "",
)
}
type ChannelDiff struct {
NewPts int
NewPosts []Post
DeletedPosts []Post
UpdatedChannel *Channel
}
+36
View File
@@ -0,0 +1,36 @@
package domain
import (
"fmt"
"time"
"github.com/google/uuid"
)
type Post struct {
ID uuid.UUID
ChannelID uuid.UUID
MessageID int
Text string
Link string
Views int
PublishedAt time.Time
}
func NewPost(channel Channel, messageID int, text string, views int, publishedAt time.Time) Post {
link := ""
if channel.Username != "" {
link = fmt.Sprintf("https://t.me/%s/%d", channel.Username, messageID)
} else if channel.TelegramID != 0 {
link = fmt.Sprintf("https://t.me/c/%d/%d", ChannelIDFromChatID(channel.TelegramID), messageID)
}
return Post{
ID: uuid.New(),
ChannelID: channel.ID,
MessageID: messageID,
Text: text,
Link: link,
Views: views,
PublishedAt: publishedAt,
}
}
@@ -0,0 +1,14 @@
package domain
import (
"time"
"github.com/google/uuid"
)
type ViewsSnapshot struct {
ViewsCount int
FetchedAt time.Time
PostID uuid.UUID
}
@@ -0,0 +1,22 @@
package usecase
import (
"context"
"fmt"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (uc *UseCase) FetchChannelMeta(ctx context.Context, username string) (domain.Channel, error) {
channel, err := uc.telegram.ParseChannelMeta(ctx, username)
if err != nil {
return domain.Channel{}, fmt.Errorf("uc.telegram.ParseChannelMeta: %s", err)
}
err = uc.database.UpdateChannelIfNotAccessible(ctx, channel)
if err != nil {
return domain.Channel{}, fmt.Errorf("uc.database.UpdateChannelIfNotAccessible: %s", err)
}
return channel, nil
}
@@ -0,0 +1,22 @@
package usecase
import (
"context"
"fmt"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (uc *UseCase) FetchChannelMetaByInvite(ctx context.Context, inviteLink string) (domain.Channel, error) {
channel, err := uc.telegram.ParseChannelMetaByInvite(ctx, inviteLink)
if err != nil {
return domain.Channel{}, fmt.Errorf("uc.telegram.ParseChannelMetaByInvite: %s", err)
}
err = uc.database.UpdateChannelIfNotAccessible(ctx, channel)
if err != nil {
return domain.Channel{}, fmt.Errorf("uc.database.UpdateChannelIfNotAccessible: %s", err)
}
return channel, nil
}
@@ -0,0 +1,177 @@
package usecase
import (
"context"
"errors"
"fmt"
"github.com/TelegramExchange/pkg/transaction"
"github.com/gotd/td/tgerr"
"github.com/rs/zerolog/log"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (uc *UseCase) FetchChannels(ctx context.Context) error {
channels := uc.database.GetChannels(ctx)
log.Debug().Msg("start fetch channels")
for _, c := range channels {
if c.Pts == 0 {
go uc.initializeChannel(ctx, c)
continue
}
err := uc.processChannel(ctx, c)
if err != nil {
if isPrivateChannelError(err) {
log.Warn().Int64("telegram_id", c.TelegramID).Msg("Private channel, attempting rejoin")
go uc.rejoinChannel(ctx, c)
return nil
}
return fmt.Errorf("uc.processChannel: %w", err)
}
}
return nil
}
func isPrivateChannelError(err error) bool {
return tgerr.Is(err, "CHANNEL_PRIVATE", "CHANNEL_INVALID", "CHANNEL_FORBIDDEN")
}
func (uc *UseCase) processChannel(ctx context.Context, channel domain.Channel) error {
diff, err := uc.telegram.GetChannelDiff(ctx, channel, 20)
if err != nil {
return fmt.Errorf("telegram.GetChannelDiff: %w", err)
}
for _, p := range diff.NewPosts {
log.Info().Msgf("New post: %s - Views: %d, Text: %.10s...", p.Link, p.Views, p.Text)
err = uc.database.CreatePost(ctx, p)
if err != nil {
return fmt.Errorf("database.CreatePost: %w", err)
}
}
for _, p := range diff.DeletedPosts {
log.Info().Msgf("Deleted post: %s - Views: %d", p.Link, p.Views)
err = uc.database.DeletePost(ctx, p)
if err != nil {
return fmt.Errorf("database.DeletePost: %w", err)
}
}
// Update channel metadata if changed
if diff.UpdatedChannel != nil {
diff.UpdatedChannel.Pts = diff.NewPts
err = uc.database.UpdateChannel(ctx, *diff.UpdatedChannel)
if err != nil {
return fmt.Errorf("database.UpdateChannel: %w", err)
}
return nil
}
channel.Pts = diff.NewPts
err = uc.database.UpdateChannel(ctx, channel)
if err != nil {
return fmt.Errorf("database.UpdateChannel: %w", err)
}
return nil
}
func (uc *UseCase) rejoinChannel(ctx context.Context, channel domain.Channel) {
defer func() {
err := uc.database.UpdateChannel(ctx, channel)
if err != nil {
log.Err(err).Msg("database.UpdateChannel (rejoinChannel)")
}
}()
if channel.InviteLink == "" {
log.Warn().Int64("telegram_id", channel.TelegramID).Msg("No invite link stored, marking inaccessible")
channel.IsAccessible = false
return
}
updated, err := uc.telegram.ParseChannelMetaByInvite(ctx, channel.InviteLink)
if err != nil {
log.Warn().Err(err).Int64("telegram_id", channel.TelegramID).Msg("Invite rejoin failed, marking inaccessible")
channel.IsAccessible = false
return
}
channel.TelegramID = updated.TelegramID
channel.Title = updated.Title
channel.AccessHash = 0
channel.Pts = 0
channel.IsAccessible = true
}
func (uc *UseCase) initializeChannel(ctx context.Context, channel domain.Channel) {
log.Debug().Stringer("ch", channel).Msg("init channel")
var (
updated domain.Channel
err error
)
switch {
case channel.Username != "":
updated, err = uc.telegram.ParseChannelMeta(ctx, channel.Username)
case channel.InviteLink != "":
updated, err = uc.telegram.ParseChannelMetaByInvite(ctx, channel.InviteLink)
default:
err = errors.New("channel state invalid")
}
if err != nil {
if isPrivateChannelError(err) {
channel.IsAccessible = false
err = uc.database.UpdateChannel(ctx, channel)
if err != nil {
log.Err(err).Msg("initializeChannel.uc.database.UpdateChannel")
}
}
log.Error().Err(err).Msg("initializeChannel.ParseChannel")
return
}
channel.Title = updated.Title
channel.AccessHash = updated.AccessHash
channel.Pts = updated.Pts
const initPosts = 20
posts, err := uc.telegram.GetChannelHistory(ctx, channel, initPosts)
if err != nil {
log.Error().Err(err).Msg("telegram.GetChannelHistory")
return
}
err = transaction.Wrap(ctx, func(ctx context.Context) error {
err = uc.database.UpdateChannel(ctx, channel)
if err != nil {
return fmt.Errorf("database.UpdateChannel: %w", err)
}
for _, p := range posts {
log.Info().Msgf("New post (init): %s - Views: %d", p.Link, p.Views)
err = uc.database.CreatePost(ctx, p)
if err != nil {
return fmt.Errorf("database.CreatePost: %w", err)
}
}
return nil
})
if err != nil {
log.Error().Err(err).Msg("transaction.Wrap")
}
}
+54
View File
@@ -0,0 +1,54 @@
package usecase
import (
"context"
"time"
"github.com/rs/zerolog/log"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
func (uc *UseCase) FetchViews(ctx context.Context) error {
channels, err := uc.database.GetChannelsWithTrackedPosts(ctx)
if err != nil {
return err
}
for _, channel := range channels {
posts, err := uc.database.GetTrackedPosts(ctx, channel)
if err != nil {
return err
}
err = uc.telegram.UpdatePostsViews(ctx, channel, posts)
if err != nil {
if isPrivateChannelError(err) {
channel.IsAccessible = false
err = uc.database.UpdateChannel(ctx, channel)
if err != nil {
log.Err(err).Msg("initializeChannel.uc.database.UpdateChannel")
}
return nil
}
return err
}
for _, p := range posts {
v := domain.ViewsSnapshot{
ViewsCount: p.Views,
FetchedAt: time.Now().UTC(),
PostID: p.ID,
}
err = uc.database.CreateViewsSnapshot(ctx, v)
if err != nil {
return err
}
log.Info().Msgf("New ViewsSnapshot: %s - Post ID: %d, Views: %d", p.Link, p.MessageID, p.Views)
}
}
return nil
}
+41
View File
@@ -0,0 +1,41 @@
package usecase
import (
"context"
"github.com/TelegramExchange/tgex-backend/tg_parser/internal/domain"
)
type Telegram interface {
ParseChannelMeta(ctx context.Context, username string) (domain.Channel, error)
ParseChannelMetaByInvite(ctx context.Context, inviteLink string) (domain.Channel, error)
GetChannelPTS(ctx context.Context, channel domain.Channel) (int, error)
GetChannelHistory(ctx context.Context, channel domain.Channel, limit int) ([]domain.Post, error)
GetChannelDiff(ctx context.Context, channel domain.Channel, limit int) (domain.ChannelDiff, error)
UpdatePostsViews(ctx context.Context, channel domain.Channel, posts []domain.Post) error
}
type Database interface {
CreatePost(ctx context.Context, post domain.Post) error
CreateViewsSnapshot(ctx context.Context, snapshot domain.ViewsSnapshot) error
UpdateChannelIfNotAccessible(ctx context.Context, channel domain.Channel) error
GetChannels(ctx context.Context) []domain.Channel
UpdateChannel(ctx context.Context, channel domain.Channel) error
GetChannelsWithTrackedPosts(ctx context.Context) ([]domain.Channel, error)
GetTrackedPosts(ctx context.Context, channel domain.Channel) ([]domain.Post, error)
DeletePost(ctx context.Context, p domain.Post) error
}
type UseCase struct {
telegram Telegram
database Database
}
func New(telegram Telegram, database Database) *UseCase {
return &UseCase{
telegram: telegram,
database: database,
}
}