Wire shared EventBus from main into PollDownloadWorker and server

Create EventBus in serveGrpc() and pass to both setupRiver() and
NewMusicAgregatorServer() so PollDownloadWorker and the gRPC server
share the same event bus instance.

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/claude-agent)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
Alexander
2026-05-20 14:57:02 +02:00
parent 908c37f73f
commit b5849d6ad1
2 changed files with 9 additions and 5 deletions
+8 -3
View File
@@ -29,6 +29,7 @@ import (
"homelab.lan/music-agregator/internal/analysis"
"homelab.lan/music-agregator/internal/config"
"homelab.lan/music-agregator/internal/database"
"homelab.lan/music-agregator/internal/eventbus"
"homelab.lan/music-agregator/internal/hello"
"homelab.lan/music-agregator/internal/indexer"
"homelab.lan/music-agregator/internal/metadata"
@@ -84,7 +85,7 @@ type riverSetup struct {
cacheRefreshWorker *indexer.CacheRefreshWorker
}
func setupRiver(ctx context.Context, cfg config.Config, db *database.DB, torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, musicfsClient *musicfs.Client, metadataService *metadata.MetadataService) *riverSetup {
func setupRiver(ctx context.Context, cfg config.Config, db *database.DB, torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, musicfsClient *musicfs.Client, metadataService *metadata.MetadataService, bus *eventbus.EventBus) *riverSetup {
cacheWorker := &indexer.CacheRefreshWorker{}
pollWorker := &workers.PollDownloadWorker{
Downloads: database.NewDownloadRepository(db.Pool),
@@ -98,6 +99,8 @@ func setupRiver(ctx context.Context, cfg config.Config, db *database.DB, torrent
MusicFSOriginID: cfg.MusicFS.OriginID,
MusicFSOriginRoot: cfg.MusicFS.OriginRoot,
MetadataService: metadataService,
EventBus: bus,
AlbumEvents: database.NewAlbumEventRepository(db.Pool),
}
riverWorkers := river.NewWorkers()
@@ -194,9 +197,11 @@ func serveGrpc(config config.Config) {
defer metadataConn.Close()
metadataService := metadata.NewMetadataService(metadataClient, db)
rs := setupRiver(ctx, config, db, torrentClient, pathMapper, musicfsClient, metadataService)
bus := eventbus.New()
musiscAgregatorSeerver, err := internal.NewMusicAgregatorServer(config, rs.client, torrentClient, pathMapper, db)
rs := setupRiver(ctx, config, db, torrentClient, pathMapper, musicfsClient, metadataService, bus)
musiscAgregatorSeerver, err := internal.NewMusicAgregatorServer(config, rs.client, torrentClient, pathMapper, db, bus)
if err != nil {
log.Fatal().Err(err).Msg("failed to create MusicAgregatorServer")
}
+1 -2
View File
@@ -26,13 +26,12 @@ type MusicAgregatorServer struct {
pb.UnimplementedMusicAgregatorServiceServer
}
func NewMusicAgregatorServer(cfg config.Config, riverClient *river.Client[pgx.Tx], torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, db *database.DB) (*MusicAgregatorServer, error) {
func NewMusicAgregatorServer(cfg config.Config, riverClient *river.Client[pgx.Tx], torrentClient torrent.TorrentClient, pathMapper *torrent.PathMapper, db *database.DB, bus *eventbus.EventBus) (*MusicAgregatorServer, error) {
service, err := NewMusicAgregatorService(cfg, riverClient, torrentClient, pathMapper, db)
if err != nil {
log.Err(err).Msg("failed to create MusicAgregatorService")
return nil, err
}
bus := eventbus.New()
return &MusicAgregatorServer{
service: service,
bus: bus,