mirror of
https://github.com/R0m1k3/FlowReader.git
synced 2026-10-11 17:28:05 +02:00
115 lines
2.8 KiB
Go
115 lines
2.8 KiB
Go
package service
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/michael/flowreader/internal/domain"
|
|
"github.com/michael/flowreader/internal/parser"
|
|
"github.com/michael/flowreader/internal/ws"
|
|
)
|
|
|
|
// FetchService handles fetching and parsing feeds.
|
|
type FetchService struct {
|
|
feedRepo domain.FeedRepository
|
|
articleRepo domain.ArticleRepository
|
|
parser *parser.FeedParser
|
|
hub *ws.Hub
|
|
}
|
|
|
|
// NewFetchService creates a new fetch service.
|
|
func NewFetchService(feedRepo domain.FeedRepository, articleRepo domain.ArticleRepository, hub *ws.Hub) *FetchService {
|
|
return &FetchService{
|
|
feedRepo: feedRepo,
|
|
articleRepo: articleRepo,
|
|
parser: parser.NewFeedParser(),
|
|
hub: hub,
|
|
}
|
|
}
|
|
|
|
// FetchFeed fetches a single feed and updates articles.
|
|
func (s *FetchService) FetchFeed(ctx context.Context, feedID uuid.UUID) error {
|
|
feed, err := s.feedRepo.GetByID(feedID)
|
|
if err != nil {
|
|
return fmt.Errorf("getting feed: %w", err)
|
|
}
|
|
if feed == nil {
|
|
return fmt.Errorf("feed not found: %s", feedID)
|
|
}
|
|
|
|
// Parse the feed
|
|
parsedFeed, err := s.parser.Parse(ctx, feed.URL, feed.ID)
|
|
if err != nil {
|
|
// Update feed with error
|
|
s.feedRepo.UpdateFetchStatus(feed.ID, time.Now(), err.Error())
|
|
return fmt.Errorf("parsing feed: %w", err)
|
|
}
|
|
|
|
// Update feed metadata
|
|
// Only update title if it's empty or looks like a URL (initial state)
|
|
if feed.Title == "" || feed.Title == feed.URL {
|
|
feed.Title = parsedFeed.Title
|
|
}
|
|
|
|
feed.Description = parsedFeed.Description
|
|
feed.SiteURL = parsedFeed.SiteURL
|
|
feed.ImageURL = parsedFeed.ImageURL
|
|
if err := s.feedRepo.Update(feed); err != nil {
|
|
log.Printf("Warning: failed to update feed metadata: %v", err)
|
|
}
|
|
|
|
// Ingest articles
|
|
if len(parsedFeed.Articles) > 0 {
|
|
if err := s.articleRepo.CreateBatch(parsedFeed.Articles); err != nil {
|
|
return fmt.Errorf("ingesting articles: %w", err)
|
|
}
|
|
|
|
// Broadcast update
|
|
if s.hub != nil {
|
|
s.hub.Broadcast("new_articles", map[string]interface{}{
|
|
"feed_id": feed.ID,
|
|
"feed_title": feed.Title,
|
|
"count": len(parsedFeed.Articles),
|
|
})
|
|
}
|
|
}
|
|
|
|
// Mark fetch as successful
|
|
s.feedRepo.UpdateFetchStatus(feed.ID, time.Now(), "")
|
|
|
|
return nil
|
|
}
|
|
|
|
// FetchAllPending fetches all feeds that need updating.
|
|
func (s *FetchService) FetchAllPending(ctx context.Context, concurrency int) (int, error) {
|
|
feeds, err := s.feedRepo.GetFeedsToFetch(100)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("getting feeds to fetch: %w", err)
|
|
}
|
|
|
|
if len(feeds) == 0 {
|
|
return 0, nil
|
|
}
|
|
|
|
// Simple sequential fetch for now (Story 2.5 will add worker pool)
|
|
fetchedCount := 0
|
|
for _, feed := range feeds {
|
|
select {
|
|
case <-ctx.Done():
|
|
return fetchedCount, ctx.Err()
|
|
default:
|
|
}
|
|
|
|
if err := s.FetchFeed(ctx, feed.ID); err != nil {
|
|
log.Printf("Error fetching feed %s: %v", feed.URL, err)
|
|
} else {
|
|
fetchedCount++
|
|
}
|
|
}
|
|
|
|
return fetchedCount, nil
|
|
}
|