nomadcode/services/core/internal/scheduler/river.go
toki 32fda2a844 feat: milestone work-item creation sync 및 roadmap sync pipeline 구현
- Plane 프로젝트의 work-item 생성 시 Orchestrator로 동기화
- scheduler에 roadmap sync job 등록 및 river client 지원
- roadmapsycpipeline adapter 패키지 추가
- 관련 테스트 파일 및 브리지 문서 추가
2026-06-14 12:39:31 +09:00

106 lines
3.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
}
// 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, 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,
})
if creationSync != nil {
river.AddWorker(workers, &RoadmapCreationSyncWorker{
Sync: creationSync,
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
}