77031ef6024d541cd8ba966084d103a9e4b4be57
- Author
- Pavle Portic <git@theedgeofrage.com>
- Committer
- Pavle Portic <git@theedgeofrage.com>
- Date
Message
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 {