feat(story-2.2): add feed parser with gofeed and article repository

This commit is contained in:
Michael committed 2026-02-04 14:54:54 +01:00
1 parent 1249e1bd42
commit 12f0c6879f
5 files changed
+646 -2

No files matched your search

@@ -45,8 +45,8 @@ development_status:
# Epic 2: Feed Core Engine
epic-2: in-progress
2-1-feed-model-add-feed-by-url: review
2-2-feed-metadata-auto-discovery: backlog
2-1-feed-model-add-feed-by-url: done
2-2-feed-metadata-auto-discovery: review
2-3-opml-import: backlog
2-4-feed-management-crud: backlog
2-5-background-feed-fetcher: backlog
+8
View File
@@ -8,11 +8,19 @@ require (
)
require (
github.com/PuerkitoBio/goquery v1.8.0 // indirect
github.com/andybalholm/cascadia v1.3.1 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20221227161230-091c0ba34f0a // indirect
github.com/jackc/puddle/v2 v2.2.1 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/mmcdole/gofeed v1.3.0 // indirect
github.com/mmcdole/goxpp v1.1.1-0.20240225020742-a0c311522b23 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
golang.org/x/crypto v0.17.0 // indirect
golang.org/x/net v0.10.0 // indirect
golang.org/x/sync v0.1.0 // indirect
golang.org/x/text v0.14.0 // indirect
)
+136
View File
@@ -0,0 +1,136 @@
// Package parser provides RSS/Atom feed parsing functionality.
package parser
import (
"context"
"fmt"
"net/http"
"time"
"github.com/google/uuid"
"github.com/michael/flowreader/internal/domain"
"github.com/mmcdole/gofeed"
)
// FeedParser handles RSS/Atom feed parsing.
type FeedParser struct {
client *http.Client
parser *gofeed.Parser
}
// NewFeedParser creates a new feed parser.
func NewFeedParser() *FeedParser {
client := &http.Client{
Timeout: 30 * time.Second,
}
return &FeedParser{
client: client,
parser: gofeed.NewParser(),
}
}
// ParsedFeed contains the parsed feed data.
type ParsedFeed struct {
Title string
Description string
SiteURL string
ImageURL string
Articles []*domain.Article
}
// Parse fetches and parses a feed URL.
func (p *FeedParser) Parse(ctx context.Context, feedURL string, feedID uuid.UUID) (*ParsedFeed, error) {
// Create request with context
req, err := http.NewRequestWithContext(ctx, http.MethodGet, feedURL, nil)
if err != nil {
return nil, fmt.Errorf("creating request: %w", err)
}
req.Header.Set("User-Agent", "FlowReader/1.0 (RSS Reader)")
// Fetch the feed
resp, err := p.client.Do(req)
if err != nil {
return nil, fmt.Errorf("fetching feed: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("unexpected status code: %d", resp.StatusCode)
}
// Parse the feed
feed, err := p.parser.Parse(resp.Body)
if err != nil {
return nil, fmt.Errorf("parsing feed: %w", err)
}
// Extract metadata
parsed := &ParsedFeed{
Title: feed.Title,
Description: feed.Description,
}
if feed.Link != "" {
parsed.SiteURL = feed.Link
}
if feed.Image != nil && feed.Image.URL != "" {
parsed.ImageURL = feed.Image.URL
}
// Convert items to articles
for _, item := range feed.Items {
article := &domain.Article{
ID: uuid.New(),
FeedID: feedID,
GUID: getGUID(item),
Title: item.Title,
}
if item.Link != "" {
article.URL = item.Link
}
if item.Content != "" {
article.Content = item.Content
}
if item.Description != "" {
article.Summary = item.Description
}
if item.Author != nil {
article.Author = item.Author.Name
} else if len(item.Authors) > 0 {
article.Author = item.Authors[0].Name
}
if item.Image != nil && item.Image.URL != "" {
article.ImageURL = item.Image.URL
}
if item.PublishedParsed != nil {
article.PublishedAt = item.PublishedParsed
} else if item.UpdatedParsed != nil {
article.PublishedAt = item.UpdatedParsed
}
article.CreatedAt = time.Now()
parsed.Articles = append(parsed.Articles, article)
}
return parsed, nil
}
// getGUID returns a unique identifier for the feed item.
func getGUID(item *gofeed.Item) string {
if item.GUID != "" {
return item.GUID
}
if item.Link != "" {
return item.Link
}
return item.Title // Last resort fallback
}
+402
View File
@@ -0,0 +1,402 @@
package repository
import (
"context"
"errors"
"fmt"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/michael/flowreader/internal/domain"
)
// ArticleRepository implements domain.ArticleRepository using PostgreSQL.
type ArticleRepository struct {
pool *pgxpool.Pool
}
// NewArticleRepository creates a new article repository.
func NewArticleRepository(pool *pgxpool.Pool) *ArticleRepository {
return &ArticleRepository{pool: pool}
}
// Create inserts a new article into the database.
func (r *ArticleRepository) Create(article *domain.Article) error {
ctx := context.Background()
query := `
INSERT INTO articles (id, feed_id, guid, title, url, content, summary, author, image_url, published_at, created_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
ON CONFLICT (feed_id, guid) DO NOTHING
`
_, err := r.pool.Exec(ctx, query,
article.ID,
article.FeedID,
article.GUID,
article.Title,
nullString(article.URL),
nullString(article.Content),
nullString(article.Summary),
nullString(article.Author),
nullString(article.ImageURL),
article.PublishedAt,
article.CreatedAt,
)
if err != nil {
return fmt.Errorf("creating article: %w", err)
}
return nil
}
// CreateBatch inserts multiple articles into the database.
func (r *ArticleRepository) CreateBatch(articles []*domain.Article) error {
ctx := context.Background()
batch := &pgx.Batch{}
query := `
INSERT INTO articles (id, feed_id, guid, title, url, content, summary, author, image_url, published_at, created_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
ON CONFLICT (feed_id, guid) DO NOTHING
`
for _, article := range articles {
batch.Queue(query,
article.ID,
article.FeedID,
article.GUID,
article.Title,
nullString(article.URL),
nullString(article.Content),
nullString(article.Summary),
nullString(article.Author),
nullString(article.ImageURL),
article.PublishedAt,
article.CreatedAt,
)
}
results := r.pool.SendBatch(ctx, batch)
defer results.Close()
for range articles {
if _, err := results.Exec(); err != nil {
return fmt.Errorf("batch insert: %w", err)
}
}
return nil
}
// GetByID retrieves an article by its ID.
func (r *ArticleRepository) GetByID(id uuid.UUID) (*domain.Article, error) {
ctx := context.Background()
query := `
SELECT a.id, a.feed_id, a.guid, a.title, a.url, a.content, a.summary, a.author,
a.image_url, a.published_at, a.is_read, a.is_favorite, a.read_at, a.created_at,
f.title as feed_title
FROM articles a
JOIN feeds f ON f.id = a.feed_id
WHERE a.id = $1
`
article, err := r.scanArticle(r.pool.QueryRow(ctx, query, id))
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
return nil, fmt.Errorf("getting article by ID: %w", err)
}
return article, nil
}
// GetByFeedID retrieves articles for a specific feed.
func (r *ArticleRepository) GetByFeedID(feedID uuid.UUID, limit, offset int) ([]*domain.Article, error) {
ctx := context.Background()
query := `
SELECT a.id, a.feed_id, a.guid, a.title, a.url, a.content, a.summary, a.author,
a.image_url, a.published_at, a.is_read, a.is_favorite, a.read_at, a.created_at,
f.title as feed_title
FROM articles a
JOIN feeds f ON f.id = a.feed_id
WHERE a.feed_id = $1
ORDER BY a.published_at DESC NULLS LAST, a.created_at DESC
LIMIT $2 OFFSET $3
`
rows, err := r.pool.Query(ctx, query, feedID, limit, offset)
if err != nil {
return nil, fmt.Errorf("querying articles: %w", err)
}
defer rows.Close()
return r.scanArticles(rows)
}
// GetByUserID retrieves articles for all user's feeds.
func (r *ArticleRepository) GetByUserID(userID uuid.UUID, limit, offset int, unreadOnly bool) ([]*domain.Article, error) {
ctx := context.Background()
var query string
if unreadOnly {
query = `
SELECT a.id, a.feed_id, a.guid, a.title, a.url, a.content, a.summary, a.author,
a.image_url, a.published_at, a.is_read, a.is_favorite, a.read_at, a.created_at,
f.title as feed_title
FROM articles a
JOIN feeds f ON f.id = a.feed_id
WHERE f.user_id = $1 AND a.is_read = false
ORDER BY a.published_at DESC NULLS LAST, a.created_at DESC
LIMIT $2 OFFSET $3
`
} else {
query = `
SELECT a.id, a.feed_id, a.guid, a.title, a.url, a.content, a.summary, a.author,
a.image_url, a.published_at, a.is_read, a.is_favorite, a.read_at, a.created_at,
f.title as feed_title
FROM articles a
JOIN feeds f ON f.id = a.feed_id
WHERE f.user_id = $1
ORDER BY a.published_at DESC NULLS LAST, a.created_at DESC
LIMIT $2 OFFSET $3
`
}
rows, err := r.pool.Query(ctx, query, userID, limit, offset)
if err != nil {
return nil, fmt.Errorf("querying articles: %w", err)
}
defer rows.Close()
return r.scanArticles(rows)
}
// GetByGUID retrieves an article by its GUID within a feed.
func (r *ArticleRepository) GetByGUID(feedID uuid.UUID, guid string) (*domain.Article, error) {
ctx := context.Background()
query := `
SELECT a.id, a.feed_id, a.guid, a.title, a.url, a.content, a.summary, a.author,
a.image_url, a.published_at, a.is_read, a.is_favorite, a.read_at, a.created_at,
f.title as feed_title
FROM articles a
JOIN feeds f ON f.id = a.feed_id
WHERE a.feed_id = $1 AND a.guid = $2
`
article, err := r.scanArticle(r.pool.QueryRow(ctx, query, feedID, guid))
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
return nil, fmt.Errorf("getting article by GUID: %w", err)
}
return article, nil
}
// MarkAsRead marks an article as read.
func (r *ArticleRepository) MarkAsRead(id uuid.UUID) error {
ctx := context.Background()
query := `UPDATE articles SET is_read = true, read_at = $2 WHERE id = $1`
_, err := r.pool.Exec(ctx, query, id, time.Now())
if err != nil {
return fmt.Errorf("marking as read: %w", err)
}
return nil
}
// MarkAsUnread marks an article as unread.
func (r *ArticleRepository) MarkAsUnread(id uuid.UUID) error {
ctx := context.Background()
query := `UPDATE articles SET is_read = false, read_at = NULL WHERE id = $1`
_, err := r.pool.Exec(ctx, query, id)
if err != nil {
return fmt.Errorf("marking as unread: %w", err)
}
return nil
}
// MarkAllAsRead marks all articles in a feed as read.
func (r *ArticleRepository) MarkAllAsRead(feedID uuid.UUID) error {
ctx := context.Background()
query := `UPDATE articles SET is_read = true, read_at = $2 WHERE feed_id = $1 AND is_read = false`
_, err := r.pool.Exec(ctx, query, feedID, time.Now())
if err != nil {
return fmt.Errorf("marking all as read: %w", err)
}
return nil
}
// ToggleFavorite toggles the favorite status of an article.
func (r *ArticleRepository) ToggleFavorite(id uuid.UUID) error {
ctx := context.Background()
query := `UPDATE articles SET is_favorite = NOT is_favorite WHERE id = $1`
_, err := r.pool.Exec(ctx, query, id)
if err != nil {
return fmt.Errorf("toggling favorite: %w", err)
}
return nil
}
// GetFavorites retrieves favorited articles for a user.
func (r *ArticleRepository) GetFavorites(userID uuid.UUID, limit, offset int) ([]*domain.Article, error) {
ctx := context.Background()
query := `
SELECT a.id, a.feed_id, a.guid, a.title, a.url, a.content, a.summary, a.author,
a.image_url, a.published_at, a.is_read, a.is_favorite, a.read_at, a.created_at,
f.title as feed_title
FROM articles a
JOIN feeds f ON f.id = a.feed_id
WHERE f.user_id = $1 AND a.is_favorite = true
ORDER BY a.published_at DESC NULLS LAST, a.created_at DESC
LIMIT $2 OFFSET $3
`
rows, err := r.pool.Query(ctx, query, userID, limit, offset)
if err != nil {
return nil, fmt.Errorf("querying favorites: %w", err)
}
defer rows.Close()
return r.scanArticles(rows)
}
// CountUnread counts unread articles for a feed.
func (r *ArticleRepository) CountUnread(feedID uuid.UUID) (int, error) {
ctx := context.Background()
query := `SELECT COUNT(*) FROM articles WHERE feed_id = $1 AND is_read = false`
var count int
err := r.pool.QueryRow(ctx, query, feedID).Scan(&count)
if err != nil {
return 0, fmt.Errorf("counting unread: %w", err)
}
return count, nil
}
// scanArticle scans a single article row.
func (r *ArticleRepository) scanArticle(row pgx.Row) (*domain.Article, error) {
var article domain.Article
var url, content, summary, author, imageURL, feedTitle *string
var publishedAt, readAt *time.Time
err := row.Scan(
&article.ID,
&article.FeedID,
&article.GUID,
&article.Title,
&url,
&content,
&summary,
&author,
&imageURL,
&publishedAt,
&article.IsRead,
&article.IsFavorite,
&readAt,
&article.CreatedAt,
&feedTitle,
)
if err != nil {
return nil, err
}
if url != nil {
article.URL = *url
}
if content != nil {
article.Content = *content
}
if summary != nil {
article.Summary = *summary
}
if author != nil {
article.Author = *author
}
if imageURL != nil {
article.ImageURL = *imageURL
}
if feedTitle != nil {
article.FeedTitle = *feedTitle
}
article.PublishedAt = publishedAt
article.ReadAt = readAt
return &article, nil
}
// scanArticles scans multiple article rows.
func (r *ArticleRepository) scanArticles(rows pgx.Rows) ([]*domain.Article, error) {
var articles []*domain.Article
for rows.Next() {
var article domain.Article
var url, content, summary, author, imageURL, feedTitle *string
var publishedAt, readAt *time.Time
err := rows.Scan(
&article.ID,
&article.FeedID,
&article.GUID,
&article.Title,
&url,
&content,
&summary,
&author,
&imageURL,
&publishedAt,
&article.IsRead,
&article.IsFavorite,
&readAt,
&article.CreatedAt,
&feedTitle,
)
if err != nil {
return nil, fmt.Errorf("scanning article: %w", err)
}
if url != nil {
article.URL = *url
}
if content != nil {
article.Content = *content
}
if summary != nil {
article.Summary = *summary
}
if author != nil {
article.Author = *author
}
if imageURL != nil {
article.ImageURL = *imageURL
}
if feedTitle != nil {
article.FeedTitle = *feedTitle
}
article.PublishedAt = publishedAt
article.ReadAt = readAt
articles = append(articles, &article)
}
return articles, nil
}
// nullString returns nil if string is empty.
func nullString(s string) *string {
if s == "" {
return nil
}
return &s
}
+98
View File
@@ -0,0 +1,98 @@
package service
import (
"context"
"fmt"
"log"
"time"
"github.com/google/uuid"
"github.com/michael/flowreader/internal/domain"
"github.com/michael/flowreader/internal/parser"
)
// FetchService handles feed fetching and article ingestion.
type FetchService struct {
feedRepo domain.FeedRepository
articleRepo domain.ArticleRepository
parser *parser.FeedParser
}
// NewFetchService creates a new fetch service.
func NewFetchService(feedRepo domain.FeedRepository, articleRepo domain.ArticleRepository) *FetchService {
return &FetchService{
feedRepo: feedRepo,
articleRepo: articleRepo,
parser: parser.NewFeedParser(),
}
}
// 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
parsed, 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
feed.Title = parsed.Title
feed.Description = parsed.Description
feed.SiteURL = parsed.SiteURL
feed.ImageURL = parsed.ImageURL
if err := s.feedRepo.Update(feed); err != nil {
log.Printf("Warning: failed to update feed metadata: %v", err)
}
// Insert new articles (ON CONFLICT DO NOTHING handles duplicates)
if len(parsed.Articles) > 0 {
if err := s.articleRepo.CreateBatch(parsed.Articles); err != nil {
return fmt.Errorf("creating articles: %w", err)
}
}
// 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
}