77031ef6024d541cd8ba966084d103a9e4b4be57

Author
Pavle Portic <git@theedgeofrage.com>
Committer
Pavle Portic <git@theedgeofrage.com>
Date

Message

Add background fetch routine and handle fetch operations in parallel

Diff

This diff is truncated to protect this page.

  1diff --git a/Makefile b/Makefile
  2index 874daa5cf374fa594074ad795cfe1503595a1a9d..30d04d518182ddefd32f5f2d00dc5bea4d21a5b3 100644
  3--- a/Makefile
  4+++ b/Makefile
  5@@ -17,6 +17,7 @@ bin/moq:
  6 
  7 gen-mocks: bin/moq
  8 	./bin/moq -pkg db_mock -out ./mocks/db/db.go ./db DB
  9+	./bin/moq -pkg parser_mock -out ./mocks/feedparser/feedparser.go ./feedparser Parser
 10 	go fmt ./...
 11 
 12 bin/golangci-lint:
 13diff --git a/cmd/main.go b/cmd/main.go
 14index 795318be34faf7c868a7f7de4e3d15fc3af0f104..37e62ab29f57a1854e74865fd53c32af9bc1e9e9 100644
 15--- a/cmd/main.go
 16+++ b/cmd/main.go
 17@@ -13,6 +13,7 @@ import (
 18 
 19 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/config"
 20 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/db"
 21+	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/feedparser"
 22 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/handler"
 23 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/httpserver/auth"
 24 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/httpserver/ytrssil"
 25@@ -24,6 +25,16 @@ func init() {
 26 	time.Local = time.UTC
 27 }
 28 
 29+func fetcherRoutine(l log.Logger, h handler.Handler) {
 30+	for {
 31+		err := h.FetchVideos(context.Background())
 32+		if err != nil {
 33+			l.Log("level", "ERROR", "function", "main.fetcherRoutine", "call", "handler.FetchVideos", "err", err)
 34+		}
 35+		time.Sleep(5 * time.Minute)
 36+	}
 37+}
 38+
 39 func main() {
 40 	log := log.NewLogger()
 41 
 42@@ -32,14 +43,13 @@ func main() {
 43 		log.Log("level", "FATAL", "call", "config.Parse", "error", err)
 44 		return
 45 	}
 46-
 47 	db, err := db.NewPostgresDB(log, config.DB)
 48 	if err != nil {
 49 		log.Log("level", "FATAL", "call", "db.NewPostgresDB", "error", err)
 50 		return
 51 	}
 52-
 53-	handler := handler.New(log, db)
 54+	parser := feedparser.NewParser(log)
 55+	handler := handler.New(log, db, parser)
 56 	gin.SetMode(gin.ReleaseMode)
 57 	router, err := ytrssil.SetupGinRouter(
 58 		log,
 59@@ -79,6 +89,9 @@ func main() {
 60 		}
 61 	}()
 62 
 63+	// start periodic fetch videos routine
 64+	go fetcherRoutine(log, handler)
 65+
 66 	log.Log(
 67 		"level", "INFO",
 68 		"msg", "ytrssil API is starting up",
 69diff --git a/cmd/main_test.go b/cmd/main_test.go
 70index e42d98501cda6543b27d8918f00a032bb34cdd3c..9493c424cb82a4baa40dfb774e4148f6cf7b62b9 100644
 71--- a/cmd/main_test.go
 72+++ b/cmd/main_test.go
 73@@ -14,6 +14,7 @@ import (
 74 
 75 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/config"
 76 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/db"
 77+	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/feedparser"
 78 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/handler"
 79 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/httpserver/auth"
 80 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/httpserver/ytrssil"
 81@@ -34,8 +35,8 @@ func setupTestServer(t *testing.T, authEnabled bool) (*http.Server, db.DB) {
 82 	if !assert.NoError(t, err) {
 83 		return nil, nil
 84 	}
 85-
 86-	handler := handler.New(l, db)
 87+	parser := feedparser.NewParser(l)
 88+	handler := handler.New(l, db, parser)
 89 	gin.SetMode(gin.TestMode)
 90 	router, err := ytrssil.SetupGinRouter(
 91 		l,
 92diff --git a/docker-compose.yml b/docker-compose.yml
 93index 4e8863d63fe21899b0dcc28cd4456e3c6d7a25f0..2f58916e97f92363ffd48aca6565110bebd832f6 100644
 94--- a/docker-compose.yml
 95+++ b/docker-compose.yml
 96@@ -11,7 +11,7 @@ services:
 97     volumes:
 98       - postgres-data:/var/lib/postgresql/data
 99     healthcheck:
100-      test: [ "CMD-SHELL", "pg_isready -d $${POSTGRES_DB} -U $${POSTGRES_USER}" ]
101+      test: ["CMD-SHELL", "pg_isready -d $${POSTGRES_DB} -U $${POSTGRES_USER}"]
102       interval: 5s
103       timeout: 2s
104       retries: 5
105diff --git a/feedparser/feedparser.go b/feedparser/feedparser.go
106index afebb54e5432cf6e0c813e68bdc849829edfe4ed..03b913eb28f1ca9218f54fd34f182a1b4b5ad399 100644
107--- a/feedparser/feedparser.go
108+++ b/feedparser/feedparser.go
109@@ -6,6 +6,7 @@ import (
110 	"fmt"
111 	"io"
112 	"net/http"
113+	"sync"
114 
115 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/lib/log"
116 	"github.com/paulrosania/go-charset/charset"
117@@ -18,29 +19,31 @@ var (
118 
119 var urlFormat = "https://www.youtube.com/feeds/videos.xml?channel_id=%s"
120 
121-// Video struct for each video in the feed
122-type Video struct {
123-	ID        string `xml:"id"`
124-	Title     string `xml:"title"`
125-	Published Date   `xml:"published"`
126+type Parser interface {
127+	Parse(channelID string) (*Channel, error)
128+	ParseThreadSafe(channelID string, channelChan chan *Channel, errChan chan error, mu *sync.Mutex, wg *sync.WaitGroup)
129 }
130 
131-// Channel struct for RSS
132-type Channel struct {
133-	Name   string  `xml:"title"`
134-	Videos []Video `xml:"entry"`
135+type parser struct {
136+	log log.Logger
137 }
138 
139-func read(l log.Logger, url string) (io.ReadCloser, error) {
140+func NewParser(l log.Logger) *parser {
141+	return &parser{
142+		log: l,
143+	}
144+}
145+
146+func (p *parser) fetch(url string) (io.ReadCloser, error) {
147 	req, err := http.NewRequest("GET", url, nil)
148 	if err != nil {
149-		l.Log("level", "ERROR", "function", "feedparser.read", "call", "http.NewRequest", "error", err)
150+		p.log.Log("level", "ERROR", "function", "feedparser.fetch", "call", "http.NewRequest", "error", err)
151 		return nil, err
152 	}
153 
154 	response, err := http.DefaultClient.Do(req)
155 	if err != nil {
156-		l.Log("level", "ERROR", "function", "feedparser.read", "call", "http.Do", "error", err)
157+		p.log.Log("level", "ERROR", "function", "feedparser.fetch", "call", "http.Do", "error", err)
158 		return nil, err
159 	}
160 
161@@ -54,9 +57,9 @@ func read(l log.Logger, url string) (io.ReadCloser, error) {
162 }
163 
164 // Parse parses a YouTube channel XML feed from a channel ID
165-func Parse(l log.Logger, channelID string) (*Channel, error) {
166+func (p *parser) Parse(channelID string) (*Channel, error) {
167 	url := fmt.Sprintf(urlFormat, channelID)
168-	reader, err := read(l, url)
169+	reader, err := p.fetch(url)
170 	if err != nil {
171 		return nil, err
172 	}
173@@ -67,8 +70,23 @@ func Parse(l log.Logger, channelID string) (*Channel, error) {
174 
175 	var channel Channel
176 	if err := xmlDecoder.Decode(&channel); err != nil {
177-		l.Log("level", "ERROR", "function", "feedparser.read", "call", "xml.Decode", "error", err)
178+		p.log.Log("level", "ERROR", "function", "feedparser.Parse", "call", "xml.Decode", "error", err)
179 		return nil, fmt.Errorf("%w: %s", ErrParseFailed, err.Error())
180 	}
181+	channel.ID = channelID
182 	return &channel, nil
183 }
184+
185+// ParseThreadSafe calls Parse, but additionally accepts an out parameter to store the result,
186+// as well as a mutex and wait group to run multiple fetches in parallel
187+func (p *parser) ParseThreadSafe(
188+	channelID string, channelChan chan *Channel, errChan chan error, mu *sync.Mutex, wg *sync.WaitGroup,
189+) {
190+	channel, err := p.Parse(channelID)
191+
192+	mu.Lock()
193+	channelChan <- channel
194+	errChan <- err
195+	mu.Unlock()
196+	wg.Done()
197+}
198diff --git a/feedparser/models.go b/feedparser/models.go
199new file mode 100644
200index 0000000000000000000000000000000000000000..54a1c45532e340e2720068e4b5693c7560b5a14e
201--- /dev/null
202+++ b/feedparser/models.go
203@@ -0,0 +1,15 @@
204+package feedparser
205+
206+// Video struct for each video in the feed
207+type Video struct {
208+	ID        string `xml:"id"`
209+	Title     string `xml:"title"`
210+	Published Date   `xml:"published"`
211+}
212+
213+// Channel struct for RSS
214+type Channel struct {
215+	ID     string
216+	Name   string  `xml:"title"`
217+	Videos []Video `xml:"entry"`
218+}
219diff --git a/handler/channels.go b/handler/channels.go
220index 0b2b787e0146ea5988f25ba2ee7bbe3e70d710cd..afe5928ee56f5502faa66761cdba5013505ebc92 100644
221--- a/handler/channels.go
222+++ b/handler/channels.go
223@@ -5,12 +5,11 @@ import (
224 	"errors"
225 
226 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/db"
227-	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/feedparser"
228 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/models"
229 )
230 
231 func (h *handler) SubscribeToChannel(ctx context.Context, username string, channelID string) error {
232-	parsedChannel, err := feedparser.Parse(h.log, channelID)
233+	parsedChannel, err := h.parser.Parse(channelID)
234 	if err != nil {
235 		return err
236 	}
237diff --git a/handler/handler.go b/handler/handler.go
238index 8c95bf75ce8fb719695a057b0b3474103c9fb629..7a0a0413d302e4cbab2647e23e95d2868aabde9f 100644
239--- a/handler/handler.go
240+++ b/handler/handler.go
241@@ -4,6 +4,7 @@ import (
242 	"context"
243 
244 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/db"
245+	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/feedparser"
246 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/lib/log"
247 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/models"
248 )
249@@ -20,10 +21,15 @@ type Handler interface {
250 }
251 
252 type handler struct {
253-	log log.Logger
254-	db  db.DB
255+	log    log.Logger
256+	db     db.DB
257+	parser feedparser.Parser
258 }
259 
260-func New(log log.Logger, db db.DB) *handler {
261-	return &handler{log: log, db: db}
262+func New(log log.Logger, db db.DB, parser feedparser.Parser) *handler {
263+	return &handler{
264+		log:    log,
265+		db:     db,
266+		parser: parser,
267+	}
268 }
269diff --git a/handler/handler_test.go b/handler/handler_test.go
270index c9bd1472a1370ad04f50e6ec0ad1a9cbca60489d..53ef3613525e53821ce95ec95b742cac22c5f99f 100644
271--- a/handler/handler_test.go
272+++ b/handler/handler_test.go
273@@ -10,6 +10,7 @@ import (
274 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/config"
275 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/lib/log"
276 	db_mock "gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/mocks/db"
277+	parser_mock "gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/mocks/feedparser"
278 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/models"
279 )
280 
281@@ -35,7 +36,7 @@ func TestGetNewVideos(t *testing.T) {
282 				},
283 			}, nil
284 		},
285-	})
286+	}, &parser_mock.ParserMock{})
287 
288 	// Act
289 	resp, err := handler.GetNewVideos(context.TODO(), "username")
290diff --git a/handler/videos.go b/handler/videos.go
291index 142e5976fec02c728bc939bab23fa26c04fdb899..6edde23e8883e0b48d5b576e5fa7321118c133ee 100644
292--- a/handler/videos.go
293+++ b/handler/videos.go
294@@ -4,6 +4,7 @@ import (
295 	"context"
296 	"errors"
297 	"strings"
298+	"sync"
299 	"time"
300 
301 	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/db"
302@@ -37,7 +38,7 @@ func (h *handler) addVideoToAllSubscribers(ctx context.Context, channelID string
303 	return nil
304 }
305 
306-func (h *handler) fetchVideosForChannel(ctx context.Context, channelID string, parsedChannel *feedparser.Channel) {
307+func (h *handler) addVideosForChannel(ctx context.Context, parsedChannel *feedparser.Channel) {
308 	for _, parsedVideo := range parsedChannel.Videos {
309 		date, err := parsedVideo.Published.Parse()
310 		if err != nil {
311@@ -45,20 +46,20 @@ func (h *handler) fetchVideosForChannel(ctx context.Context, channelID string, p
312 			continue
313 		}
314 
315-		id := strings.Split(parsedVideo.ID, ":")[2]
316+		videoID := strings.Split(parsedVideo.ID, ":")[2]
317 		video := models.Video{
318-			ID:            id,
319+			ID:            videoID,
320 			Title:         parsedVideo.Title,
321 			PublishedTime: date,
322 		}
323-		err = h.db.AddVideo(ctx, video, channelID)
324+		err = h.db.AddVideo(ctx, video, parsedChannel.ID)
325 		if err != nil {
326 			if !errors.Is(err, db.ErrVideoExists) {
327 				h.log.Log("level", "WARNING", "call", "db.AddVideo", "err", err)
328 			}
329 			continue
330 		}
331-		err = h.addVideoToAllSubscribers(ctx, channelID, id)
332+		err = h.addVideoToAllSubscribers(ctx, parsedChannel.ID, videoID)
333 		if err != nil {
334 			continue
335 		}
336@@ -66,18 +67,29 @@ func (h *handler) fetchVideosForChannel(ctx context.Context, channelID string, p
337 }
338 
339 func (h *handler) FetchVideos(ctx context.Context) error {
340+	h.log.Log("level", "INFO", "msg", "fetching new videos for all channels")
341+
342 	channels, err := h.db.ListChannels(ctx)
343 	if err != nil {
344 		return err
345 	}
346-
347+	var parsedChannels = make(chan *feedparser.Channel, len(channels))
348+	var errors = make(chan error, len(channels))
349+	var wg sync.WaitGroup
350+	var mu sync.Mutex
351 	for _, channel := range channels {
352-		parsedChannel, err := feedparser.Parse(h.log, channel.ID)
353+		wg.Add(1)
354+		go h.parser.ParseThreadSafe(channel.ID, parsedChannels, errors, &mu, &wg)
355+	}
356+	wg.Wait()
357+
358+	for range channels {
359+		parsedChannel := <-parsedChannels
360+		err = <-errors
361 		if err != nil {
362 			continue
363 		}
364-
365-		h.fetchVideosForChannel(ctx, channel.ID, parsedChannel)
366+		h.addVideosForChannel(ctx, parsedChannel)
367 	}
368 
369 	return nil
370diff --git a/httpserver/ytrssil/api_setup_test.go b/httpserver/ytrssil/api_setup_test.go
371index 7284a41415d8101e9b73b62f97f1aec5f87e14a4..925a7e51865e5b4cb77ead943cabbd79ed077a3d 100644
372--- a/httpserver/ytrssil/api_setup_test.go
373+++ b/httpserver/ytrssil/api_setup_test.go
374@@ -28,7 +28,7 @@ func init() {
375 func setupTestServer(t *testing.T) *http.Server {
376 	l := log.NewNopLogger()
377 
378-	handler := handler.New(l, nil)
379+	handler := handler.New(l, nil, nil)
380 
381 	gin.SetMode(gin.TestMode)
382 	router, err := ytrssil.SetupGinRouter(
383diff --git a/mocks/feedparser/feedparser.go b/mocks/feedparser/feedparser.go
384new file mode 100644
385index 0000000000000000000000000000000000000000..64502bf73e3f70fcbe6e78c762d2907d4a2defc7
386--- /dev/null
387+++ b/mocks/feedparser/feedparser.go
388@@ -0,0 +1,143 @@
389+// Code generated by moq; DO NOT EDIT.
390+// github.com/matryer/moq
391+
392+package parser_mock
393+
394+import (
395+	"gitea.theedgeofrage.com/TheEdgeOfRage/ytrssil-api/feedparser"
396+	"sync"
397+)
398+
399+// Ensure, that ParserMock does implement feedparser.Parser.
400+// If this is not the case, regenerate this file with moq.
401+var _ feedparser.Parser = &ParserMock{}
402+
403+// ParserMock is a mock implementation of feedparser.Parser.
404+//
405+//	func TestSomethingThatUsesParser(t *testing.T) {
406+//
407+//		// make and configure a mocked feedparser.Parser
408+//		mockedParser := &ParserMock{
409+//			ParseFunc: func(channelID string) (*feedparser.Channel, error) {
410+//				panic("mock out the Parse method")
411+//			},
412+//			ParseThreadSafeFunc: func(channelID string, channelChan chan *feedparser.Channel, errChan chan error, mu *sync.Mutex, wg *sync.WaitGroup)  {
413+//				panic("mock out the ParseThreadSafe method")
414+//			},
415+//		}
416+//
417+//		// use mockedParser in code that requires feedparser.Parser
418+//		// and then make assertions.
419+//
420+//	}
421+type ParserMock struct {
422+	// ParseFunc mocks the Parse method.
423+	ParseFunc func(channelID string) (*feedparser.Channel, error)
424+
425+	// ParseThreadSafeFunc mocks the ParseThreadSafe method.
426+	ParseThreadSafeFunc func(channelID string, channelChan chan *feedparser.Channel, errChan chan error, mu *sync.Mutex, wg *sync.WaitGroup)
427+
428+	// calls tracks calls to the methods.
429+	calls struct {
430+		// Parse holds details about calls to the Parse method.
431+		Parse []struct {
432+			// ChannelID is the channelID argument value.
433+			ChannelID string
434+		}
435+		// ParseThreadSafe holds details about calls to the ParseThreadSafe method.
436+		ParseThreadSafe []struct {
437+			// ChannelID is the channelID argument value.
438+			ChannelID string
439+			// ChannelChan is the channelChan argument value.
440+			ChannelChan chan *feedparser.Channel
441+			// ErrChan is the errChan argument value.
442+			ErrChan chan error
443+			// Mu is the mu argument value.
444+			Mu *sync.Mutex
445+			// Wg is the wg argument value.
446+			Wg *sync.WaitGroup
447+		}
448+	}
449+	lockParse           sync.RWMutex
450+	lockParseThreadSafe sync.RWMutex
451+}
452+
453+// Parse calls ParseFunc.
454+func (mock *ParserMock) Parse(channelID string) (*feedparser.Channel, error) {
455+	if mock.ParseFunc == nil {
456+		panic("ParserMock.ParseFunc: method is nil but Parser.Parse was just called")
457+	}
458+	callInfo := struct {
459+		ChannelID string
460+	}{
461+		ChannelID: channelID,
462+	}
463+	mock.lockParse.Lock()
464+	mock.calls.Parse = append(mock.calls.Parse, callInfo)
465+	mock.lockParse.Unlock()
466+	return mock.ParseFunc(channelID)
467+}
468+
469+// ParseCalls gets all the calls that were made to Parse.
470+// Check the length with:
471+//
472+//	len(mockedParser.ParseCalls())
473+func (mock *ParserMock) ParseCalls() []struct {
474+	ChannelID string
475+} {
476+	var calls []struct {
477+		ChannelID string
478+	}
479+	mock.lockParse.RLock()
480+	calls = mock.calls.Parse
481+	mock.lockParse.RUnlock()
482+	return calls
483+}
484+
485+// ParseThreadSafe calls ParseThreadSafeFunc.
486+func (mock *ParserMock) ParseThreadSafe(channelID string, channelChan chan *feedparser.Channel, errChan chan error, mu *sync.Mutex, wg *sync.WaitGroup) {
487+	if mock.ParseThreadSafeFunc == nil {