Files
FlowReader/internal/worker/fetcher.go
T
Antigravity AgentandClaude Opus 5.5 03e57e4308 fix(backend): harden security and speed up feeds and article API
Security:
- WebSocket events are routed to their owner only (no cross-user leak);
  hub close is idempotent (fixes double-close panic), adds ping/pong and
  write deadlines.
- Session tokens stored as SHA-256 (migration 008 keeps sessions valid);
  single-query auth middleware puts the user in the request context.
- Client IP only trusts X-Forwarded-For from TRUSTED_PROXIES; rate limiter
  map is bounded; per-user limit on AI summaries.
- Argon2id at OWASP minimum with a concurrency cap; constant-time login
  for unknown emails; atomic first-admin bootstrap; REGISTRATION_ENABLED.
- CSP/HSTS/COOP headers, same-origin guard on mutations, body size limits,
  wider SSRF denylist, bounded feed/page/AI response reads, generic errors.
- Upgrade chi, pgx, x/net, x/text, x/crypto (known CVEs); commit go.sum.

Performance:
- List endpoints return a plain-text excerpt and reading time instead of
  full HTML; content is sanitized once at ingest (legacy rows backfilled).
- Keyset pagination on (sort_at, id) with matching partial indexes;
  redundant indexes dropped (migration 007).
- Fetcher: bounded worker pool, conditional GET (ETag/Last-Modified),
  exponential backoff, dedupe before insert, column-safe truncation,
  retention-aware ingest, per-user refresh coalescing.
- Read/favorite/read-all are single ownership-scoped statements.
- gzip compression, immutable caching for hashed assets, path-safe SPA
  handler, server timeouts; expired sessions purged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-09 07:34:08 +02:00

80 lines
1.6 KiB
Go

// Package worker provides background job processing.
package worker
import (
"context"
"log"
"sync"
"time"
"github.com/michael/flowreader/internal/service"
)
// FeedFetcher handles periodic feed fetching.
type FeedFetcher struct {
fetchService *service.FetchService
interval time.Duration
concurrency int
stopCh chan struct{}
wg sync.WaitGroup
}
// NewFeedFetcher creates a new feed fetcher worker. The interval is how often
// due feeds are looked for; each feed carries its own next_fetch_at.
func NewFeedFetcher(fetchService *service.FetchService, interval time.Duration, concurrency int) *FeedFetcher {
return &FeedFetcher{
fetchService: fetchService,
interval: interval,
concurrency: concurrency,
stopCh: make(chan struct{}),
}
}
// Start begins the background fetch loop.
func (f *FeedFetcher) Start() {
f.wg.Add(1)
go f.run()
log.Printf("Feed fetcher started (interval: %s, concurrency: %d)", f.interval, f.concurrency)
}
// Stop gracefully stops the fetcher.
func (f *FeedFetcher) Stop() {
close(f.stopCh)
f.wg.Wait()
log.Println("Feed fetcher stopped")
}
func (f *FeedFetcher) run() {
defer f.wg.Done()
// Initial fetch on startup
f.fetch()
ticker := time.NewTicker(f.interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
f.fetch()
case <-f.stopCh:
return
}
}
}
func (f *FeedFetcher) fetch() {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
count, err := f.fetchService.FetchAllPending(ctx)
if err != nil {
log.Printf("Feed fetch error: %v", err)
return
}
if count > 0 {
log.Printf("Fetched %d feeds", count)
}
}