summaryrefslogtreecommitdiff
path: root/fetch
diff options
context:
space:
mode:
authornytpu <alex@nytpu.com>2021-03-18 13:43:22 -0600
committernytpu <alex@nytpu.com>2021-03-18 13:43:22 -0600
commit326edc97c5b20ba0fed5a7d770d8fd61d0215ed3 (patch)
tree8c067206e8a5c97266bb3b4c7f4f00c17e396dd4 /fetch
parent364e8fbc210ba71c2cff5bb479860ac1130e178a (diff)
add RefreshAll with workers
Diffstat (limited to 'fetch')
-rw-r--r--fetch/fetch.go73
1 files changed, 73 insertions, 0 deletions
diff --git a/fetch/fetch.go b/fetch/fetch.go
index 5537339..428457f 100644
--- a/fetch/fetch.go
+++ b/fetch/fetch.go
@@ -9,6 +9,7 @@ package fetch
import (
"fmt"
"net/url"
+ "sync"
"golang.nytpu.com/comitium/core"
)
@@ -42,3 +43,75 @@ func Page(data *core.FullData, remote *url.URL, title string) error {
return fmt.Errorf("Unsupported protocol '%s'", remote.Scheme)
}
}
+
+type refreshJob struct {
+ typ string // "feed" or "page"
+ url url.URL
+}
+
+// RefreshAll will check all feeds and pages in a core.FullData for updates
+func RefreshAll(data *core.FullData, numWorkers int) error {
+ var wg sync.WaitGroup
+
+ data.RLock()
+ numJobs := len(data.Feeds) + len(data.Pages)
+ if numJobs == 0 {
+ data.RUnlock()
+ return nil
+ }
+
+ if numWorkers < 1 {
+ numWorkers = 1
+ }
+
+ jobs := make(chan refreshJob, numJobs)
+ returns := make(chan error, numJobs)
+
+ // start workers but jobs is blocking for right now
+ for w := 0; w < numWorkers; w++ {
+ wg.Add(1)
+ go refreshWorker(&wg, jobs, returns, data)
+ }
+
+ // get all keys in maps
+ feedKeys := make([]url.URL, 0, len(data.Feeds))
+ for k := range data.Feeds {
+ // we know that the keys are already validated uris, so no err
+ u, _ := url.ParseRequestURI(k)
+ feedKeys = append(feedKeys, *u)
+ }
+ pageKeys := make([]url.URL, 0, len(data.Pages))
+ for k := range data.Pages {
+ u, _ := url.ParseRequestURI(k)
+ pageKeys = append(pageKeys, *u)
+ }
+ data.RUnlock()
+
+ for _, v := range feedKeys {
+ jobs <- refreshJob{"feed", v}
+ }
+ for _, v := range pageKeys {
+ jobs <- refreshJob{"page", v}
+ }
+ close(jobs)
+
+ wg.Wait()
+
+ for i := 0; i < numJobs; i++ {
+ if err := <-returns; err != nil {
+ return err
+ }
+ }
+ return nil
+}
+
+func refreshWorker(wg *sync.WaitGroup, jobs <-chan refreshJob, returns chan<- error, data *core.FullData) {
+ defer wg.Done()
+ for j := range jobs {
+ if j.typ == "feed" {
+ returns <- Feed(data, &j.url, "")
+ } else if j.typ == "page" {
+ returns <- Page(data, &j.url, "")
+ }
+ }
+}