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/agent" "github.com/nomadcode/nomadcode-core/internal/model" "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, agentClient agent.Client, modelClient model.Client, logger *slog.Logger) (*Client, error) { workers := river.NewWorkers() river.AddWorker(workers, &TaskWorker{ Store: store, Notifications: notifications, Agent: agentClient, Model: modelClient, 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 }