package scheduler import ( "context" "log/slog" "time" "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/workflow" ) type Client struct { river *river.Client[pgx.Tx] driver *riverpgxv5.Driver logger *slog.Logger } // New builds the scheduler client. creationSync is the optional Plane-origin // Milestone creation sync orchestrator: when non-nil its worker is registered so // EnqueueRoadmapCreationSync has a runtime worker; when nil only the task worker // is registered, preserving the default task scheduler behavior. func New(pool *pgxpool.Pool, lifecycle TaskLifecycle, notifications *notification.Service, agentClient agent.Client, modelClient model.Client, creationSync CreationSyncRunner, finalizer TaskFinalizer, slotUpdater WorkspaceSlotStateUpdater, runTimeout time.Duration, logger *slog.Logger) (*Client, error) { workers := river.NewWorkers() river.AddWorker(workers, &TaskWorker{ Lifecycle: lifecycle, Notifications: notifications, Agent: agentClient, Model: modelClient, SlotUpdater: slotUpdater, RunTimeout: runTimeout, Logger: logger, }) if creationSync != nil { river.AddWorker(workers, &RoadmapCreationSyncWorker{ Sync: creationSync, TaskFinalizer: finalizer, SlotUpdater: slotUpdater, 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: workflow.DefaultTaskMaxAttempts, }) 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 } // EnqueueRoadmapCreationSync enqueues a Plane-origin Milestone creation sync // job. It is the internal seam the develop-scan handoff calls to hand a matched // scan result to the orchestrator worker; it is intentionally separate from // EnqueueTask and is not exposed through a public HTTP route. func (c *Client) EnqueueRoadmapCreationSync(ctx context.Context, args RoadmapCreationSyncJobArgs) error { _, err := c.river.Insert(ctx, args, nil) return err }