mirror of
https://github.com/mmahdium/TGSS.git
synced 2026-08-12 08:32:48 +03:30
feat: add feed cache implementation and integrate with channel handler
This commit is contained in:
+3
-1
@@ -23,10 +23,12 @@ func main() {
|
||||
handlers.NewAuthHandler,
|
||||
|
||||
rss.NewRSSGenerator,
|
||||
rss.NewFeedCache,
|
||||
),
|
||||
telegram.Module,
|
||||
fx.Invoke(server.RegisterRoutes),
|
||||
fx.Invoke((server.RegisterCleanup)),
|
||||
fx.Invoke((server.RegisterRateLimiterCleanup)),
|
||||
fx.Invoke(rss.RegisterFeedCacheCleanup),
|
||||
fx.Logger(infra.NewFXLogger()),
|
||||
)
|
||||
|
||||
|
||||
@@ -13,18 +13,23 @@ import (
|
||||
)
|
||||
|
||||
type ChannelHandler struct {
|
||||
logger *zap.Logger
|
||||
tgService *telegram.Service
|
||||
rss *rss.RSSGenerator
|
||||
logger *zap.Logger
|
||||
tgService *telegram.Service
|
||||
rssGenerator *rss.RSSGenerator
|
||||
rssFeedCache *rss.FeedCache
|
||||
}
|
||||
|
||||
func NewChannelHandler(
|
||||
tgService *telegram.Service,
|
||||
logger *zap.Logger,
|
||||
rssGenerator *rss.RSSGenerator,
|
||||
rssFeedCache *rss.FeedCache,
|
||||
) *ChannelHandler {
|
||||
return &ChannelHandler{
|
||||
logger: logger,
|
||||
tgService: tgService,
|
||||
logger: logger,
|
||||
tgService: tgService,
|
||||
rssGenerator: rssGenerator,
|
||||
rssFeedCache: rssFeedCache,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -80,16 +85,24 @@ func (ch *ChannelHandler) GetMessagesRSS(c *gin.Context) {
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), config.MessageOperationTimeout)
|
||||
defer cancel()
|
||||
var rssFeed *rss.RSSFeed
|
||||
rssFeed = ch.rssFeedCache.GetFeedFromFeedCache(channelId, limit)
|
||||
|
||||
messages, err := ch.tgService.LastMessages(ctx, channelId, limit)
|
||||
if err != nil {
|
||||
ch.logger.Error("Error getting channel messages:", zap.Error(err))
|
||||
c.JSON(500, gin.H{"error": "Unable to get channel messages"})
|
||||
return
|
||||
if rssFeed == nil {
|
||||
messages, err := ch.tgService.LastMessages(ctx, channelId, limit)
|
||||
if err != nil {
|
||||
ch.logger.Error("Error getting channel messages:", zap.Error(err))
|
||||
c.JSON(500, gin.H{"error": "Unable to get channel messages"})
|
||||
return
|
||||
}
|
||||
|
||||
rssFeed = ch.rssGenerator.GenerateFeed(messages, channelId)
|
||||
ch.rssFeedCache.SetFeedToFeedCache(channelId, rssFeed)
|
||||
c.Header("X-Cache-Status", "MISS")
|
||||
} else {
|
||||
c.Header("X-Cache-Status", "HIT")
|
||||
}
|
||||
|
||||
rssFeed := ch.rss.GenerateFeed(messages, channelId)
|
||||
|
||||
c.Header("Content-Type", "application/rss+xml; charset=utf-8")
|
||||
c.Header("X-Content-Type-Options", "nosniff")
|
||||
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package rss
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"go.uber.org/fx"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
type FeedCache struct {
|
||||
logger *zap.Logger
|
||||
feedsCache map[string]*RSSFeed
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
func NewFeedCache(logger *zap.Logger) *FeedCache {
|
||||
return &FeedCache{logger: logger, feedsCache: map[string]*RSSFeed{}, mu: sync.Mutex{}}
|
||||
}
|
||||
|
||||
func (c *FeedCache) SetFeedToFeedCache(channelId string, feed *RSSFeed) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
c.feedsCache[channelId] = feed
|
||||
}
|
||||
|
||||
func (c *FeedCache) GetFeedFromFeedCache(channelId string, limit int) *RSSFeed {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
if feed, ok := c.feedsCache[channelId]; ok {
|
||||
feedBuildTime, err := time.Parse("Mon, 02 Jan 2006 15:04 -0700", feed.Channel.LastBuildDate)
|
||||
if err != nil || time.Since(feedBuildTime) > time.Minute*15 {
|
||||
delete(c.feedsCache, channelId)
|
||||
return nil
|
||||
}
|
||||
|
||||
if len(feed.Channel.Items) >= limit {
|
||||
feedCopy := *feed
|
||||
feedCopy.Channel.Items = feed.Channel.Items[len(feed.Channel.Items)-limit:]
|
||||
c.logger.Debug("Cache hit on channel ID", zap.String("channelId", channelId))
|
||||
return &feedCopy
|
||||
}
|
||||
return nil
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *FeedCache) CleanupFeedCache() {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
for channelId, feed := range c.feedsCache {
|
||||
feedBuildTime, _ := time.Parse("Mon, 02 Jan 2006 15:04 MST", feed.Channel.LastBuildDate)
|
||||
if time.Since(feedBuildTime) > time.Minute*15 {
|
||||
delete(c.feedsCache, channelId)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func RegisterFeedCacheCleanup(lc fx.Lifecycle, logger *zap.Logger, f *FeedCache) {
|
||||
var ticker *time.Ticker
|
||||
var stopChan chan struct{}
|
||||
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(ctx context.Context) error {
|
||||
logger.Info("Starting feed cache cleanup task", zap.Duration("interval", 30*time.Minute))
|
||||
ticker = time.NewTicker(20 * time.Minute)
|
||||
stopChan = make(chan struct{})
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
logger.Info("Running feed cache cleanup")
|
||||
f.CleanupFeedCache()
|
||||
logger.Info("Feed cache cleanup completed")
|
||||
case <-stopChan:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
},
|
||||
OnStop: func(ctx context.Context) error {
|
||||
logger.Info("Stopping feed cache cleanup task")
|
||||
if ticker != nil {
|
||||
ticker.Stop()
|
||||
}
|
||||
if stopChan != nil {
|
||||
close(stopChan)
|
||||
}
|
||||
return nil
|
||||
},
|
||||
})
|
||||
}
|
||||
@@ -52,7 +52,6 @@ func RegisterRoutes(params RouterParams) {
|
||||
|
||||
// TODO: improve accurecy in rss channel fields
|
||||
// TODO: add ui with templates under /setup with fetch
|
||||
// TODO: add gin level cache (look for higher limit)
|
||||
// TODO: image endpoint + hash and expiry
|
||||
authStatCtx, cancel := context.WithTimeout(context.Background(), config.AuthStatusTimeout)
|
||||
authStat, err := params.TgService.AuthStatus(authStatCtx)
|
||||
|
||||
@@ -54,7 +54,7 @@ func (r *RateLimiter) CleanupRateLimiter() {
|
||||
|
||||
}
|
||||
|
||||
func RegisterCleanup(lc fx.Lifecycle, logger *zap.Logger, r *RateLimiter) {
|
||||
func RegisterRateLimiterCleanup(lc fx.Lifecycle, logger *zap.Logger, r *RateLimiter) {
|
||||
var ticker *time.Ticker
|
||||
var stopChan chan struct{}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user