nomadcode/services/core/internal/scheduler/river.go

87 lines
2.1 KiB
Go

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
}
func New(pool *pgxpool.Pool, lifecycle TaskLifecycle, notifications *notification.Service, agentClient agent.Client, modelClient model.Client, runTimeout time.Duration, logger *slog.Logger) (*Client, error) {
workers := river.NewWorkers()
river.AddWorker(workers, &TaskWorker{
Lifecycle: lifecycle,
Notifications: notifications,
Agent: agentClient,
Model: modelClient,
RunTimeout: runTimeout,
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
}