mirror of
https://github.com/R0m1k3/FlowReader.git
synced 2026-10-11 17:28:05 +02:00
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>
136 lines
3.2 KiB
Go
136 lines
3.2 KiB
Go
package worker
|
|
|
|
import (
|
|
"context"
|
|
"log"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/michael/flowreader/internal/repository"
|
|
"github.com/michael/flowreader/internal/service"
|
|
"github.com/michael/flowreader/internal/utils"
|
|
)
|
|
|
|
// Cleaner handles periodic database maintenance.
|
|
type Cleaner struct {
|
|
repo *repository.ArticleRepository
|
|
authService *service.AuthService
|
|
interval time.Duration
|
|
stopCh chan struct{}
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
// NewCleaner creates a new database cleaner worker.
|
|
func NewCleaner(repo *repository.ArticleRepository, authService *service.AuthService, interval time.Duration) *Cleaner {
|
|
return &Cleaner{
|
|
repo: repo,
|
|
authService: authService,
|
|
interval: interval,
|
|
stopCh: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
// Start begins the background maintenance loop.
|
|
func (c *Cleaner) Start() {
|
|
c.wg.Add(1)
|
|
go c.run()
|
|
log.Printf("Maintenance worker started (interval: %s)", c.interval)
|
|
}
|
|
|
|
// Stop gracefully stops the cleaner.
|
|
func (c *Cleaner) Stop() {
|
|
close(c.stopCh)
|
|
c.wg.Wait()
|
|
log.Println("Maintenance worker stopped")
|
|
}
|
|
|
|
func (c *Cleaner) run() {
|
|
defer c.wg.Done()
|
|
|
|
// One-off: sanitize legacy articles and compute excerpts/word counts.
|
|
c.backfill()
|
|
|
|
// Initial cleanup on startup
|
|
c.cleanup()
|
|
|
|
ticker := time.NewTicker(c.interval)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ticker.C:
|
|
c.cleanup()
|
|
case <-c.stopCh:
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func (c *Cleaner) cleanup() {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
|
|
defer cancel()
|
|
|
|
count, err := c.repo.DeleteOldArticles(ctx, service.ArticleRetention)
|
|
if err != nil {
|
|
log.Printf("Maintenance cleanup error: %v", err)
|
|
} else if count > 0 {
|
|
log.Printf("Maintenance: cleaned up %d old articles", count)
|
|
}
|
|
|
|
if c.authService != nil {
|
|
if n, err := c.authService.PurgeExpiredSessions(); err != nil {
|
|
log.Printf("Maintenance session purge error: %v", err)
|
|
} else if n > 0 {
|
|
log.Printf("Maintenance: purged %d expired sessions", n)
|
|
}
|
|
}
|
|
}
|
|
|
|
// backfill processes articles stored before migration 007: their HTML is
|
|
// sanitized once and stored, and the list excerpt / word count is derived,
|
|
// so list endpoints never need to ship or sanitize full content again.
|
|
func (c *Cleaner) backfill() {
|
|
sanitizer := utils.NewContentSanitizer()
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
|
|
defer cancel()
|
|
|
|
total := 0
|
|
for {
|
|
select {
|
|
case <-c.stopCh:
|
|
return
|
|
default:
|
|
}
|
|
|
|
rows, err := c.repo.PendingBackfill(ctx, 200)
|
|
if err != nil {
|
|
log.Printf("Backfill error: %v", err)
|
|
return
|
|
}
|
|
if len(rows) == 0 {
|
|
break
|
|
}
|
|
for _, row := range rows {
|
|
content := sanitizer.Sanitize(row.Content)
|
|
summary := sanitizer.Sanitize(row.Summary)
|
|
plain := utils.PlainText(content)
|
|
if plain == "" {
|
|
plain = utils.PlainText(summary)
|
|
}
|
|
excerptSrc := utils.PlainText(summary)
|
|
if excerptSrc == "" {
|
|
excerptSrc = plain
|
|
}
|
|
if err := c.repo.UpdateDerived(ctx, row.ID, content, summary,
|
|
utils.Excerpt(excerptSrc, 320), utils.WordCount(plain)); err != nil {
|
|
log.Printf("Backfill error on %s: %v", row.ID, err)
|
|
return
|
|
}
|
|
}
|
|
total += len(rows)
|
|
}
|
|
if total > 0 {
|
|
log.Printf("Backfill: processed %d legacy articles", total)
|
|
}
|
|
}
|