81 lines
1.8 KiB
Go
81 lines
1.8 KiB
Go
package scheduler
|
|
|
|
import (
|
|
"context"
|
|
"log/slog"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
"github.com/riverqueue/river"
|
|
"github.com/riverqueue/river/riverdriver/riverpgxv5"
|
|
"github.com/riverqueue/river/rivermigrate"
|
|
|
|
"github.com/nomadcode/nomadcode-core/internal/notification"
|
|
"github.com/nomadcode/nomadcode-core/internal/storage"
|
|
)
|
|
|
|
type Client struct {
|
|
river *river.Client[pgx.Tx]
|
|
driver *riverpgxv5.Driver
|
|
logger *slog.Logger
|
|
}
|
|
|
|
func New(pool *pgxpool.Pool, store *storage.Store, notifications *notification.Service, logger *slog.Logger) (*Client, error) {
|
|
workers := river.NewWorkers()
|
|
river.AddWorker(workers, &TaskWorker{
|
|
Store: store,
|
|
Notifications: notifications,
|
|
Logger: logger,
|
|
})
|
|
|
|
driver := riverpgxv5.New(pool)
|
|
riverClient, err := river.NewClient[pgx.Tx](driver, &river.Config{
|
|
Logger: logger,
|
|
Queues: map[string]river.QueueConfig{
|
|
river.QueueDefault: {MaxWorkers: 2},
|
|
},
|
|
Workers: workers,
|
|
MaxAttempts: 3,
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &Client{
|
|
river: riverClient,
|
|
driver: driver,
|
|
logger: logger,
|
|
}, nil
|
|
}
|
|
|
|
func (c *Client) Migrate(ctx context.Context) error {
|
|
migrator, err := rivermigrate.New[pgx.Tx](c.driver, &rivermigrate.Config{
|
|
Logger: c.logger,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
result, err := migrator.Migrate(ctx, rivermigrate.DirectionUp, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if c.logger != nil && len(result.Versions) > 0 {
|
|
c.logger.Info("river migrations applied", "count", len(result.Versions))
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (c *Client) Start(ctx context.Context) error {
|
|
return c.river.Start(ctx)
|
|
}
|
|
|
|
func (c *Client) Stop(ctx context.Context) error {
|
|
return c.river.Stop(ctx)
|
|
}
|
|
|
|
func (c *Client) EnqueueTask(ctx context.Context, taskID string) error {
|
|
_, err := c.river.Insert(ctx, TaskJobArgs{TaskID: taskID}, nil)
|
|
return err
|
|
}
|