diff --git a/_bmad-output/implementation-artifacts/sprint-status.yaml b/_bmad-output/implementation-artifacts/sprint-status.yaml index 3eb749d..a3c8f06 100644 --- a/_bmad-output/implementation-artifacts/sprint-status.yaml +++ b/_bmad-output/implementation-artifacts/sprint-status.yaml @@ -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 diff --git a/go.mod b/go.mod index 823ce78..867b851 100644 --- a/go.mod +++ b/go.mod @@ -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 ) diff --git a/internal/parser/parser.go b/internal/parser/parser.go new file mode 100644 index 0000000..09f079c --- /dev/null +++ b/internal/parser/parser.go @@ -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 +} diff --git a/internal/repository/article.go b/internal/repository/article.go new file mode 100644 index 0000000..6caebde --- /dev/null +++ b/internal/repository/article.go @@ -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 +} diff --git a/internal/service/fetch.go b/internal/service/fetch.go new file mode 100644 index 0000000..513bb92 --- /dev/null +++ b/internal/service/fetch.go @@ -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 +}