Rssmon2/exchange/exchange.go

197 lines
5.1 KiB
Go

package exchange
import (
"errors"
"local1/logger"
"local3/rssmon2/monitor"
"local3/rssmon2/rss"
"local3/rssmon2/server"
"local3/rssmon2/store"
"sort"
"strings"
"time"
)
const nsForFeeds = "FEEDS"
type Exchange struct {
Mon *monitor.Monitor
SClient store.Client
Srv *server.Server
allFeeds map[string]*rss.Feed
}
func New(mon *monitor.Monitor, sclient store.Client, srv *server.Server) *Exchange {
return &Exchange{
Mon: mon,
SClient: sclient,
Srv: srv,
allFeeds: make(map[string]*rss.Feed),
}
}
func (ex *Exchange) LoadDB() error {
oldFeeds, err := ex.SClient.List(nsForFeeds, "", true, -1)
if err != nil {
return err
}
for _, feedID := range oldFeeds {
b, err := ex.SClient.Get(nsForFeeds, feedID)
if err != nil {
return err
}
feed, err := rss.Deserialize(b)
if err != nil {
return err
}
ex.allFeeds[feed.Link] = feed
if err := ex.Mon.Submit(feed.Link, feed.Interval); err != nil {
return err
}
}
return nil
}
func (ex *Exchange) NewFeed(url, itemFilter, contentFilter string, tags []string, interval time.Duration) {
feed, err := rss.New(url, itemFilter, contentFilter, tags, interval)
if err != nil {
logger.Logf("can't create new RSS %q: %v", url, err)
return
}
ex.allFeeds[url] = feed
if err := ex.Mon.Submit(url, feed.Interval); err != nil {
logger.Logf("Cannot accept new feed %q: %v", url, err)
}
}
func (ex *Exchange) GetFeedRSS(url string, n int) (string, error) {
feed, ok := ex.allFeeds[url]
if !ok {
return "", errors.New("unknown feed " + url)
}
itemKeys, err := ex.SClient.List(feed.ID(), "", false, n)
if err != nil {
return "", err
}
items := make([]*rss.Item, len(itemKeys))
for i := range itemKeys {
b, err := ex.SClient.Get(feed.ID(), itemKeys[i])
if err != nil {
return "", errors.New("cannot get feed item " + itemKeys[i])
}
items[i], err = rss.DeserializeItem(b)
if err != nil {
return "", errors.New("cannot deserialize feed item" + itemKeys[i])
}
}
return rss.ToRSS(feed, items)
}
func (ex *Exchange) GetFeedItem(ID string) (string, error) {
b, err := ex.SClient.Get(strings.Split(ID, ".")[0], strings.Join(strings.Split(ID, ".")[1:], "."))
if err != nil {
return "", errors.New("cannot get feed item " + ID)
}
item, err := rss.DeserializeItem(b)
if err != nil {
return "", errors.New("cannot deserialize feed item" + ID)
}
return item.Content, nil
}
func (ex *Exchange) GetFeedTagRSS(tag string) (string, error) {
feedNames, err := ex.SClient.List(nsForFeeds, "", true, -1)
if err != nil {
return "", err
}
combinedItems := []*rss.Item{}
for _, feedName := range feedNames {
b, err := ex.SClient.Get(nsForFeeds, feedName)
if err != nil {
return "", err
}
feed, err := rss.Deserialize(b)
if err != nil {
return "", err
}
for i := range feed.Tags {
if feed.Tags[i] == tag {
itemKeys, err := ex.SClient.List(feed.ID(), "", false, 20) // TODO variable, or pick most recent n or something
if err != nil {
return "", err
}
for i := range itemKeys {
b, err := ex.SClient.Get(feed.ID(), itemKeys[i])
if err != nil {
return "", errors.New("cannot get feed item " + itemKeys[i])
}
item, err := rss.DeserializeItem(b)
if err != nil {
return "", errors.New("cannot deserialize feed item" + itemKeys[i])
}
combinedItems = append(combinedItems, item)
}
break
}
}
}
combinedFeed, err := rss.New(tag, "", "", nil, time.Minute)
if err != nil {
return "", err
}
sort.Slice(combinedItems, func(i, j int) bool {
return !combinedItems[i].TS.Before(combinedItems[j].TS)
})
return rss.ToRSS(combinedFeed, combinedItems)
}
func (ex *Exchange) UpdateFeed(url string) {
feed, ok := ex.allFeeds[url]
if !ok {
f, err := rss.New(url, "", "", nil, time.Minute)
if err != nil {
logger.Logf("cannot identify unknown feed triggered in monitor: %q: %v", url, err)
return
}
b, err := ex.SClient.Get(nsForFeeds, f.ID())
if err != nil {
logger.Logf("cannot get unknown feed triggered in monitor: %q: %v", url, err)
return
}
feed, err = rss.Deserialize(b)
if err != nil {
logger.Logf("cannot deserialize feed triggered in monitor: %q: %v", url, err)
return
}
}
items, err := ex.allFeeds[url].Update()
if err != nil {
logger.Logf("can't update old RSS %q: %v", url, err)
return
}
b, err := feed.Serialize()
if err != nil {
logger.Logf("can't serialize to save RSS %q: %v", url, err)
return
}
if err := ex.SClient.Set(nsForFeeds, feed.ID(), b); err != nil {
logger.Logf("can't save RSS %q.%q: %v", nsForFeeds, feed.ID(), err)
return
}
logger.Log("Saved feed", feed)
for i := range items {
b, err := items[i].Serialize()
if err != nil {
logger.Logf("can't save rss item %q.%q: %v", url, items[i].Link, err)
return
}
if err := ex.SClient.Set(feed.ID(), items[i].ID(), b); err != nil {
logger.Logf("can't save rss item %q.%q: %v", feed.ID(), items[i].ID(), err)
return
}
//logger.Log("Saved feed item", feed.ID(), items[i].ID(), items[i])
}
logger.Logf("Saved %d feed items for %s", len(items), feed.Title)
}