Commit e1fca8b
Eric Bower
·
2025-08-07 22:21:35 -0400 EDT
parent 4ef30c4
feat(feeds): a wild `cron` property appears Now users can specify a digest interval for their rss-to-email posts using a cron format. The primary limitation is we run our feed processor every minute so we don't support the seconds specificity in cron. Reference: https://github.com/adhocore/gronx?tab=readme-ov-file#cron-expression
6 files changed,
+66,
-22
M
go.mod
+1,
-0
| ... | ... | @@ -23,6 +23,7 @@ toolchain go1.24.0 | |
| 23 | 23 | require ( | |
| 24 | 24 | git.sr.ht/~delthas/senpai v0.4.0 | |
| 25 | 25 | git.sr.ht/~rockorager/vaxis v0.14.1-0.20250527151737-5530f9f4bcf6 | |
| 26 | + | github.com/adhocore/gronx v1.19.6 | |
| 26 | 27 | github.com/alecthomas/chroma/v2 v2.15.0 | |
| 27 | 28 | github.com/antoniomika/syncmap v1.0.0 | |
| 28 | 29 | github.com/araddon/dateparse v0.0.0-20210429162001-6b43995a97de |
M
go.sum
+2,
-0
| ... | ... | @@ -54,6 +54,8 @@ github.com/PuerkitoBio/goquery v1.10.2/go.mod h1:0guWGjcLu9AYC7C1GHnpysHy056u9aE | |
| 54 | 54 | github.com/RoaringBitmap/roaring v1.2.1/go.mod h1:icnadbWcNyfEHlYdr+tDlOTih1Bf/h+rzPpv4sbomAA= | |
| 55 | 55 | github.com/RoaringBitmap/roaring v1.9.4 h1:yhEIoH4YezLYT04s1nHehNO64EKFTop/wBhxv2QzDdQ= | |
| 56 | 56 | github.com/RoaringBitmap/roaring v1.9.4/go.mod h1:6AXUsoIEzDTFFQCe1RbGA6uFONMhvejWj5rqITANK90= | |
| 57 | + | github.com/adhocore/gronx v1.19.6 h1:5KNVcoR9ACgL9HhEqCm5QXsab/gI4QDIybTAWcXDKDc= | |
| 58 | + | github.com/adhocore/gronx v1.19.6/go.mod h1:7oUY1WAU8rEJWmAxXR2DN0JaO4gi9khSgKjiRypqteg= | |
| 57 | 59 | github.com/alecthomas/assert/v2 v2.11.0 h1:2Q9r3ki8+JYXvGsDyBXwH3LcJ+WK5D0gc5E8vS6K3D0= | |
| 58 | 60 | github.com/alecthomas/assert/v2 v2.11.0/go.mod h1:Bze95FyfUr7x34QZrjL+XP+0qgp/zg8yS+TtBj1WA3k= | |
| 59 | 61 | github.com/alecthomas/chroma/v2 v2.2.0/go.mod h1:vf4zrexSH54oEjJ7EdB65tGNHmH3pGZmVkgTP5RHvAs= |
+10,
-3
| ... | ... | @@ -70,13 +70,20 @@ func Middleware(dbpool db.DB, cfg *shared.ConfigSite) pssh.SSHServerMiddleware { | |
| 70 | 70 | _, _ = fmt.Fprintln(writer, "Filename\tLast Digest\tNext Digest\tInterval\tFailed Attempts") | |
| 71 | 71 | for _, post := range posts.Data { | |
| 72 | 72 | parsed := shared.ListParseText(post.Text) | |
| 73 | - | digestOption := DigestOptionToTime(*post.Data.LastDigest, parsed.DigestInterval) | |
| 73 | + | ||
| 74 | + | nextDigest := "" | |
| 75 | + | if parsed.Cron != "" { | |
| 76 | + | nextDigest = parsed.Cron | |
| 77 | + | } else { | |
| 78 | + | digestOption := DigestOptionToTime(*post.Data.LastDigest, parsed.DigestInterval) | |
| 79 | + | nextDigest = digestOption.Format(time.RFC3339) | |
| 80 | + | } | |
| 74 | 81 | _, _ = fmt.Fprintf( | |
| 75 | 82 | writer, | |
| 76 | 83 | "%s\t%s\t%s\t%s\t%d/10\r\n", | |
| 77 | 84 | post.Filename, | |
| 78 | 85 | post.Data.LastDigest.Format(time.RFC3339), | |
| 79 | - | digestOption.Format(time.RFC3339), | |
| 86 | + | nextDigest, | |
| 80 | 87 | parsed.DigestInterval, | |
| 81 | 88 | post.Data.Attempts, | |
| 82 | 89 | ) |
| ... | ... | @@ -123,7 +130,7 @@ func Middleware(dbpool db.DB, cfg *shared.ConfigSite) pssh.SSHServerMiddleware { | |
| 123 | 130 | } | |
| 124 | 131 | _, _ = fmt.Fprintf(sesh, "running feed post: %s\r\n", filename) | |
| 125 | 132 | fetcher := NewFetcher(dbpool, cfg) | |
| 126 | - | err = fetcher.RunPost(logger, user, post, true) | |
| 133 | + | err = fetcher.RunPost(logger, user, post, true, time.Now().UTC()) | |
| 127 | 134 | if err != nil { | |
| 128 | 135 | _, _ = fmt.Fprintln(sesh.Stderr(), err) | |
| 129 | 136 | } |
+42,
-18
| ... | ... | @@ -15,6 +15,7 @@ import ( | |
| 15 | 15 | "text/template" | |
| 16 | 16 | "time" | |
| 17 | 17 | ||
| 18 | + | "github.com/adhocore/gronx" | |
| 18 | 19 | "github.com/emersion/go-sasl" | |
| 19 | 20 | "github.com/emersion/go-smtp" | |
| 20 | 21 | "github.com/mmcdole/gofeed" |
| ... | ... | @@ -129,26 +130,34 @@ type Fetcher struct { | |
| 129 | 130 | cfg *shared.ConfigSite | |
| 130 | 131 | db db.DB | |
| 131 | 132 | auth sasl.Client | |
| 133 | + | gron *gronx.Gronx | |
| 132 | 134 | } | |
| 133 | 135 | ||
| 134 | 136 | func NewFetcher(dbpool db.DB, cfg *shared.ConfigSite) *Fetcher { | |
| 135 | 137 | smtPass := os.Getenv("PICO_SMTP_PASS") | |
| 136 | 138 | emailLogin := os.Getenv("PICO_SMTP_USER") | |
| 137 | 139 | auth := sasl.NewPlainClient("", emailLogin, smtPass) | |
| 140 | + | gron := gronx.New() | |
| 138 | 141 | return &Fetcher{ | |
| 139 | 142 | db: dbpool, | |
| 140 | 143 | cfg: cfg, | |
| 141 | 144 | auth: auth, | |
| 145 | + | gron: gron, | |
| 142 | 146 | } | |
| 143 | 147 | } | |
| 144 | 148 | ||
| 145 | - | func (f *Fetcher) Validate(post *db.Post, parsed *shared.ListParsedText) error { | |
| 149 | + | func (f *Fetcher) Validate(post *db.Post, parsed *shared.ListParsedText, now time.Time) error { | |
| 146 | 150 | lastDigest := post.Data.LastDigest | |
| 147 | 151 | if lastDigest == nil { | |
| 148 | 152 | return nil | |
| 149 | 153 | } | |
| 150 | 154 | ||
| 151 | - | now := time.Now().UTC() | |
| 155 | + | toTheMin := time.Date( | |
| 156 | + | now.Year(), now.Month(), now.Day(), | |
| 157 | + | now.Hour(), now.Minute(), | |
| 158 | + | 0, 0, // zero out second and nano-second for cron | |
| 159 | + | now.Location(), | |
| 160 | + | ) | |
| 152 | 161 | ||
| 153 | 162 | expiresAt := post.ExpiresAt | |
| 154 | 163 | if expiresAt != nil { |
| ... | ... | @@ -157,14 +166,29 @@ func (f *Fetcher) Validate(post *db.Post, parsed *shared.ListParsedText) error { | |
| 157 | 166 | } | |
| 158 | 167 | } | |
| 159 | 168 | ||
| 160 | - | digestAt := DigestOptionToTime(*lastDigest, parsed.DigestInterval) | |
| 161 | - | if digestAt.After(now) { | |
| 162 | - | return fmt.Errorf("(%s) not time to digest, skipping", digestAt.Format(time.RFC3339)) | |
| 169 | + | if parsed.Cron != "" { | |
| 170 | + | isDue, err := f.gron.IsDue(parsed.Cron, toTheMin) | |
| 171 | + | if err != nil { | |
| 172 | + | return fmt.Errorf("cron error, skipping; err: %w", err) | |
| 173 | + | } | |
| 174 | + | if !isDue { | |
| 175 | + | nextTime, _ := gronx.NextTick(parsed.Cron, true) | |
| 176 | + | return fmt.Errorf( | |
| 177 | + | "cron not time to digest, skipping; cur run: %s, next run: %s", | |
| 178 | + | f.gron.C.GetRef(), | |
| 179 | + | nextTime, | |
| 180 | + | ) | |
| 181 | + | } | |
| 182 | + | } else if parsed.DigestInterval != "" { | |
| 183 | + | digestAt := DigestOptionToTime(*lastDigest, parsed.DigestInterval) | |
| 184 | + | if digestAt.After(now) { | |
| 185 | + | return fmt.Errorf("(%s) not time to digest, skipping", digestAt.Format(time.RFC3339)) | |
| 186 | + | } | |
| 163 | 187 | } | |
| 164 | 188 | return nil | |
| 165 | 189 | } | |
| 166 | 190 | ||
| 167 | - | func (f *Fetcher) RunPost(logger *slog.Logger, user *db.User, post *db.Post, skipValidation bool) error { | |
| 191 | + | func (f *Fetcher) RunPost(logger *slog.Logger, user *db.User, post *db.Post, skipValidation bool, now time.Time) error { | |
| 168 | 192 | logger = logger.With("filename", post.Filename) | |
| 169 | 193 | logger.Info("running feed post") | |
| 170 | 194 |
| ... | ... | @@ -178,8 +202,8 @@ func (f *Fetcher) RunPost(logger *slog.Logger, user *db.User, post *db.Post, ski | |
| 178 | 202 | } | |
| 179 | 203 | } | |
| 180 | 204 | ||
| 181 | - | logger.Info("last digest at", "lastDigest", post.Data.LastDigest.Format(time.RFC3339)) | |
| 182 | - | err := f.Validate(post, parsed) | |
| 205 | + | logger.Info("last digest", "timestamp", post.Data.LastDigest.Format(time.RFC3339)) | |
| 206 | + | err := f.Validate(post, parsed, now) | |
| 183 | 207 | if err != nil { | |
| 184 | 208 | logger.Info("validation failed", "err", err) | |
| 185 | 209 | if skipValidation { |
| ... | ... | @@ -212,9 +236,8 @@ func (f *Fetcher) RunPost(logger *slog.Logger, user *db.User, post *db.Post, ski | |
| 212 | 236 | urls = append(urls, u) | |
| 213 | 237 | } | |
| 214 | 238 | ||
| 215 | - | now := time.Now().UTC() | |
| 216 | 239 | if post.ExpiresAt == nil { | |
| 217 | - | expiresAt := time.Now().AddDate(0, 12, 0) | |
| 240 | + | expiresAt := now.AddDate(0, 12, 0) | |
| 218 | 241 | post.ExpiresAt = &expiresAt | |
| 219 | 242 | } | |
| 220 | 243 | _, err = f.db.UpdatePost(post) |
| ... | ... | @@ -291,8 +314,9 @@ Also, we have centralized logs in our pico.sh TUI that will display realtime fee | |
| 291 | 314 | return nil | |
| 292 | 315 | } | |
| 293 | 316 | ||
| 294 | - | func (f *Fetcher) RunUser(user *db.User) error { | |
| 317 | + | func (f *Fetcher) RunUser(user *db.User, now time.Time) error { | |
| 295 | 318 | logger := shared.LoggerWithUser(f.cfg.Logger, user) | |
| 319 | + | logger.Info("run user") | |
| 296 | 320 | posts, err := f.db.FindPostsForUser(&db.Pager{Num: 100}, user.ID, "feeds") | |
| 297 | 321 | if err != nil { | |
| 298 | 322 | return err |
| ... | ... | @@ -303,7 +327,7 @@ func (f *Fetcher) RunUser(user *db.User) error { | |
| 303 | 327 | } | |
| 304 | 328 | ||
| 305 | 329 | for _, post := range posts.Data { | |
| 306 | - | err = f.RunPost(logger, user, post, false) | |
| 330 | + | err = f.RunPost(logger, user, post, false, now) | |
| 307 | 331 | if err != nil { | |
| 308 | 332 | logger.Error("run post failed", "err", err) | |
| 309 | 333 | } |
| ... | ... | @@ -561,16 +585,16 @@ func (f *Fetcher) SendEmail(logger *slog.Logger, username, email, subject string | |
| 561 | 585 | return err | |
| 562 | 586 | } | |
| 563 | 587 | ||
| 564 | - | func (f *Fetcher) Run(logger *slog.Logger) error { | |
| 588 | + | func (f *Fetcher) Run(now time.Time) error { | |
| 565 | 589 | users, err := f.db.FindUsers() | |
| 566 | 590 | if err != nil { | |
| 567 | 591 | return err | |
| 568 | 592 | } | |
| 569 | 593 | ||
| 570 | 594 | for _, user := range users { | |
| 571 | - | err := f.RunUser(user) | |
| 595 | + | err := f.RunUser(user, now) | |
| 572 | 596 | if err != nil { | |
| 573 | - | logger.Error("run user failed", "err", err) | |
| 597 | + | f.cfg.Logger.Error("run user failed", "err", err) | |
| 574 | 598 | continue | |
| 575 | 599 | } | |
| 576 | 600 | } |
| ... | ... | @@ -583,12 +607,12 @@ func (f *Fetcher) Loop() { | |
| 583 | 607 | for { | |
| 584 | 608 | logger.Info("running digest emailer") | |
| 585 | 609 | ||
| 586 | - | err := f.Run(logger) | |
| 610 | + | err := f.Run(time.Now().UTC()) | |
| 587 | 611 | if err != nil { | |
| 588 | 612 | logger.Error("run failed", "err", err) | |
| 589 | 613 | } | |
| 590 | 614 | ||
| 591 | - | logger.Info("digest emailer finished, waiting 10 mins") | |
| 592 | - | time.Sleep(10 * time.Minute) | |
| 615 | + | logger.Info("digest emailer finished, waiting 1min ...") | |
| 616 | + | time.Sleep(1 * time.Minute) | |
| 593 | 617 | } | |
| 594 | 618 | } |
+4,
-1
| ... | ... | @@ -1561,7 +1561,10 @@ func (me *PsqlDB) InsertFeedItems(postID string, items []*db.FeedItem) error { | |
| 1561 | 1561 | item.Data, | |
| 1562 | 1562 | ) | |
| 1563 | 1563 | if err != nil { | |
| 1564 | - | return fmt.Errorf("post id:%s, guid:%s, err:%w", item.PostID, item.GUID, err) | |
| 1564 | + | return fmt.Errorf( | |
| 1565 | + | "post id:%s, link:%s, guid:%s, err:%w", | |
| 1566 | + | item.PostID, item.Data.Link, item.GUID, err, | |
| 1567 | + | ) | |
| 1565 | 1568 | } | |
| 1566 | 1569 | } | |
| 1567 | 1570 |
| ... | ... | @@ -11,6 +11,7 @@ import ( | |
| 11 | 11 | ||
| 12 | 12 | "slices" | |
| 13 | 13 | ||
| 14 | + | "github.com/adhocore/gronx" | |
| 14 | 15 | "github.com/araddon/dateparse" | |
| 15 | 16 | ) | |
| 16 | 17 |
| ... | ... | @@ -52,6 +53,7 @@ type ListMetaData struct { | |
| 52 | 53 | Tags []string | |
| 53 | 54 | ListType string // https://developer.mozilla.org/en-US/docs/Web/CSS/list-style-type | |
| 54 | 55 | DigestInterval string | |
| 56 | + | Cron string | |
| 55 | 57 | Email string | |
| 56 | 58 | InlineContent bool // allows content inlining to be disabled in feeds.pico.sh emails | |
| 57 | 59 | } |
| ... | ... | @@ -130,6 +132,11 @@ func TokenToMetaField(meta *ListMetaData, token *SplitToken) error { | |
| 130 | 132 | ) | |
| 131 | 133 | } | |
| 132 | 134 | meta.DigestInterval = token.Value | |
| 135 | + | case "cron": | |
| 136 | + | if !gronx.IsValid(token.Value) { | |
| 137 | + | return fmt.Errorf("(%s) is not in a valid cron format: https://github.com/adhocore/gronx?tab=readme-ov-file#cron-expression", token.Value) | |
| 138 | + | } | |
| 139 | + | meta.Cron = token.Value | |
| 133 | 140 | case "email": | |
| 134 | 141 | meta.Email = token.Value | |
| 135 | 142 | case "inline_content": |