- Update milestone and SDD documents - Add webhook HTTP receiver with idempotency checks - Add config for webhook endpoints - Implement gito events processing - Add HTTP handlers and router updates - Archive completed task files
258 lines
8.7 KiB
Go
258 lines
8.7 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"log/slog"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"syscall"
|
|
"time"
|
|
|
|
agenta2a "github.com/nomadcode/nomadcode-core/internal/adapters/a2a"
|
|
"github.com/nomadcode/nomadcode-core/internal/adapters/jira"
|
|
"github.com/nomadcode/nomadcode-core/internal/adapters/mattermost"
|
|
modelopenai "github.com/nomadcode/nomadcode-core/internal/adapters/openai"
|
|
"github.com/nomadcode/nomadcode-core/internal/adapters/plane"
|
|
"github.com/nomadcode/nomadcode-core/internal/agent"
|
|
"github.com/nomadcode/nomadcode-core/internal/config"
|
|
"github.com/nomadcode/nomadcode-core/internal/db"
|
|
"github.com/nomadcode/nomadcode-core/internal/gitoevents"
|
|
"github.com/nomadcode/nomadcode-core/internal/gitosync"
|
|
apphttp "github.com/nomadcode/nomadcode-core/internal/http"
|
|
"github.com/nomadcode/nomadcode-core/internal/notification"
|
|
"github.com/nomadcode/nomadcode-core/internal/protosocket"
|
|
"github.com/nomadcode/nomadcode-core/internal/roadmapsyncpipeline"
|
|
"github.com/nomadcode/nomadcode-core/internal/scheduler"
|
|
"github.com/nomadcode/nomadcode-core/internal/storage"
|
|
"github.com/nomadcode/nomadcode-core/internal/workflow"
|
|
)
|
|
|
|
func main() {
|
|
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
|
|
if err := run(logger); err != nil {
|
|
logger.Error("server stopped", "error", err)
|
|
os.Exit(1)
|
|
}
|
|
}
|
|
|
|
func run(logger *slog.Logger) error {
|
|
cfg := config.Load()
|
|
|
|
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
|
defer stop()
|
|
|
|
pool, err := db.NewPool(ctx, cfg.DatabaseURL)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer pool.Close()
|
|
|
|
store := storage.NewStore(pool)
|
|
mattermostClient := mattermost.NewClient(mattermost.Config{
|
|
BaseURL: cfg.MattermostBaseURL,
|
|
Token: cfg.MattermostToken,
|
|
ChannelID: cfg.MattermostChannelID,
|
|
}, logger)
|
|
planeClient := plane.NewClient(plane.Config{
|
|
BaseURL: cfg.PlaneBaseURL,
|
|
Token: cfg.PlaneToken,
|
|
}, logger)
|
|
jiraClient := jira.NewClient(jira.Config{
|
|
BaseURL: cfg.JiraBaseURL,
|
|
Email: cfg.JiraEmail,
|
|
APIToken: cfg.JiraAPIToken,
|
|
}, logger)
|
|
if cfg.JiraBaseURL != "" {
|
|
logger.Info("jira adapter initialized", "base_url", cfg.JiraBaseURL)
|
|
}
|
|
protoSocketServer := protosocket.NewServer(protosocket.Config{
|
|
HeartbeatIntervalSec: cfg.ProtoSocketHeartbeatIntervalSec,
|
|
HeartbeatWaitSec: cfg.ProtoSocketHeartbeatWaitSec,
|
|
}, logger)
|
|
|
|
notificationService := notification.NewService(
|
|
logger,
|
|
mattermost.NewTaskNotificationSink(mattermostClient),
|
|
protosocket.NewTaskEventBroadcaster(protoSocketServer),
|
|
)
|
|
modelClient := modelopenai.NewClient(modelopenai.Config{
|
|
BaseURL: cfg.ModelBaseURL,
|
|
APIKey: cfg.ModelAPIKey,
|
|
Model: cfg.ModelName,
|
|
ContextSize: cfg.ModelContextSize,
|
|
TimeoutSec: cfg.ModelTimeoutSec,
|
|
}, logger)
|
|
|
|
var agentClient agent.Client
|
|
if cfg.A2AEdgeURL != "" {
|
|
agentClient = agenta2a.NewClient(agenta2a.Config{
|
|
URL: cfg.A2AEdgeURL,
|
|
Token: cfg.A2AToken,
|
|
TimeoutSec: cfg.A2ATimeoutSec,
|
|
}, logger)
|
|
logger.Info("a2a edge input enabled", "url", cfg.A2AEdgeURL)
|
|
}
|
|
|
|
lifecycle := workflow.NewLifecycle(store, logger)
|
|
|
|
// Plane-origin Milestone creation sync orchestrator. Jira is intentionally
|
|
// not wired here; this slice handles the Plane projection path only. The
|
|
// self-actor guard is active only when PLANE_SELF_ACTOR_ID is configured.
|
|
creationSync := roadmapsyncpipeline.NewService(store, planeClient, cfg.PlaneSelfActorID)
|
|
|
|
workflowService := workflow.NewService(store, nil, logger)
|
|
|
|
taskScheduler, err := scheduler.New(pool, lifecycle, notificationService, agentClient, modelClient, creationSync, workflowService, store, time.Duration(cfg.WorkflowTaskTimeoutSec)*time.Second, logger)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := taskScheduler.Migrate(ctx); err != nil {
|
|
return err
|
|
}
|
|
if err := taskScheduler.Start(ctx); err != nil {
|
|
return err
|
|
}
|
|
|
|
workflowService.SetEnqueuer(taskScheduler)
|
|
|
|
// Build a shared Gito bridge when either the proto-socket runner or the HTTP
|
|
// webhook consumer is enabled. Both share the same scanner/bridge assembly.
|
|
var gitoBridge *gitosync.Bridge
|
|
if cfg.GitoBranchEventsEnabled() || cfg.GitoWebhookConsumerEnabled() {
|
|
var err error
|
|
gitoBridge, err = newGitoBridge(cfg, planeClient, taskScheduler, logger)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
// Gito branch event consumer (proto-socket): drives Plane-origin Milestone
|
|
// creation sync from a develop push. Started only when the endpoint, repo,
|
|
// local develop checkout, and Todo state id are all configured; otherwise the
|
|
// Core server behaves exactly as before. Cancelled with the root context on shutdown.
|
|
if cfg.GitoBranchEventsEnabled() {
|
|
gitoRunner, err := newGitoRunner(cfg, gitoBridge, logger)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
reconnectingRunner := gitosync.NewReconnectingRunner(gitoRunner, logger)
|
|
go func() {
|
|
if err := reconnectingRunner.Run(ctx); err != nil {
|
|
logger.Error("gito branch event runner stopped", "error", err)
|
|
}
|
|
}()
|
|
logger.Info("gito branch event consumer enabled",
|
|
"repo_id", cfg.GitoRepoID, "branch", cfg.GitoBranch)
|
|
}
|
|
projectBinder := apphttp.NewProjectBinder(store)
|
|
handler := apphttp.NewHandler(
|
|
pool,
|
|
workflowService,
|
|
projectBinder,
|
|
logger,
|
|
apphttp.WorkItemProvider{ID: "plane", Reader: planeClient},
|
|
apphttp.WorkItemProvider{ID: "jira", Reader: jiraClient},
|
|
)
|
|
handler.SetPlaneWebhookSecret(cfg.PlaneWebhookSecret)
|
|
handler.SetPlaneWebhookDispatchConfig(apphttp.PlaneWebhookDispatchConfig{
|
|
WorkspaceID: cfg.PlaneWorkspaceID,
|
|
WorkspaceSlug: cfg.PlaneWorkspaceSlug,
|
|
BacklogStateID: cfg.PlaneCreationBacklogStateID,
|
|
AgentAssigneeID: cfg.PlaneAgentAssigneeID,
|
|
SelfActorID: cfg.PlaneSelfActorID,
|
|
})
|
|
if cfg.GitoWebhookConsumerEnabled() {
|
|
handler.SetGitoWebhookConfig(apphttp.GitoWebhookConfig{
|
|
Secret: cfg.GitoWebhookSecret,
|
|
RepoID: cfg.GitoRepoID,
|
|
Branch: cfg.GitoBranch,
|
|
})
|
|
handler.SetGitoBranchEventHandler(&gitoBridgeHTTPAdapter{bridge: gitoBridge})
|
|
logger.Info("gito http webhook consumer enabled",
|
|
"repo_id", cfg.GitoRepoID, "branch", cfg.GitoBranch)
|
|
}
|
|
|
|
protosocket.NewTaskChannels(workflowService).Register(protoSocketServer.Dispatcher())
|
|
|
|
server := &http.Server{
|
|
Addr: cfg.HTTPAddr,
|
|
Handler: apphttp.NewRouter(handler, logger, apphttp.AuthConfig{Username: cfg.AuthUsername, Password: cfg.AuthPassword}, protoSocketServer, cfg.ProtoSocketPath),
|
|
ReadHeaderTimeout: 5 * time.Second,
|
|
}
|
|
|
|
serverErr := make(chan error, 1)
|
|
go func() {
|
|
logger.Info("server listening", "addr", cfg.HTTPAddr, "env", cfg.AppEnv)
|
|
if err := server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
|
serverErr <- err
|
|
return
|
|
}
|
|
serverErr <- nil
|
|
}()
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
case err := <-serverErr:
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
|
|
if err := server.Shutdown(shutdownCtx); err != nil {
|
|
return err
|
|
}
|
|
if err := protoSocketServer.Close(); err != nil {
|
|
return err
|
|
}
|
|
if err := taskScheduler.Stop(shutdownCtx); err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// newGitoBridge assembles the shared Gito bridge: an exec-backed develop scanner,
|
|
// the Plane work item reader, the scheduler enqueuer, and an in-memory
|
|
// duplicate-revision guard. Both the proto-socket runner and the HTTP webhook
|
|
// consumer call this and share the resulting bridge.
|
|
func newGitoBridge(cfg config.Config, reader *plane.Client, enqueuer *scheduler.Client, logger *slog.Logger) (*gitosync.Bridge, error) {
|
|
scanner, err := gitosync.NewBranchRevisionScanner(gitosync.ExecCommandRunner{}, gitosync.ScannerConfig{
|
|
DevelopRepoPath: cfg.GitoDevelopRepoPath,
|
|
RemoteName: cfg.GitoRemoteName,
|
|
Branch: cfg.GitoBranch,
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return gitosync.NewBridge(scanner, reader, enqueuer, nil, gitosync.BridgeConfig{
|
|
RepoID: cfg.GitoRepoID,
|
|
Branch: cfg.GitoBranch,
|
|
TodoStateID: cfg.RoadmapCreationTodoStateID,
|
|
}, logger)
|
|
}
|
|
|
|
// newGitoRunner wires a pre-built bridge into a proto-socket runner.
|
|
func newGitoRunner(cfg config.Config, bridge *gitosync.Bridge, logger *slog.Logger) (*gitoevents.Client, error) {
|
|
return gitosync.NewRunner(cfg.GitoProtoSocketURL, cfg.GitoRepoID, cfg.GitoBranch, bridge, logger)
|
|
}
|
|
|
|
// gitoBridgeHTTPAdapter adapts *gitosync.Bridge to apphttp.GitoBranchEventHandler.
|
|
// The http package defines its own event struct to avoid an import cycle between
|
|
// internal/http → internal/gitoevents → internal/protosocket ← protosocket tests → internal/http.
|
|
type gitoBridgeHTTPAdapter struct {
|
|
bridge *gitosync.Bridge
|
|
}
|
|
|
|
func (a *gitoBridgeHTTPAdapter) Handle(ctx context.Context, ev apphttp.GitoBranchUpdatedEvent) error {
|
|
return a.bridge.Handle(ctx, gitoevents.BranchUpdatedEvent{
|
|
RepoID: ev.RepoID,
|
|
Branch: ev.Branch,
|
|
Before: ev.Before,
|
|
After: ev.After,
|
|
})
|
|
}
|