- action/ → activity/ (Temporal activities) - statemachine/ → workflow/ (Temporal workflows) - Removed internal/api/ and cmd/server/ (api-gw handles HTTP, Temporal is the API) - Created pkg/types/types.go as single source of truth for all shared types - Extracted CallRoleLLM helper (DRY: implementer/planner/judge shared pattern) - Fixed circular import: workflow_graph_query uses string activity names - Fixed logger.logf → logger.Info/Warn (method didn't exist) - Fixed routing types: added Branches, Activity, BackoffSeconds, TaskActivity - Fixed db.Canvas.Name, db.Client→DB, GetWorkflow→FetchWorkflow - Removed unused imports - All tests pass, build clean, vet clean
153 lines
4.6 KiB
Go
153 lines
4.6 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"log"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"syscall"
|
|
"time"
|
|
|
|
"go.temporal.io/sdk/client"
|
|
"go.temporal.io/sdk/worker"
|
|
"github.com/rockliang/poimen/workflows/activity"
|
|
"github.com/rockliang/poimen/workflows/internal/config"
|
|
"github.com/rockliang/poimen/workflows/internal/health"
|
|
"github.com/rockliang/poimen/workflows/internal/logging"
|
|
"github.com/rockliang/poimen/workflows/workflow"
|
|
)
|
|
|
|
func main() {
|
|
// Initialize structured logging
|
|
if err := logging.InitLogger(); err != nil {
|
|
log.Fatalf("failed to initialize logger: %v", err)
|
|
}
|
|
defer logging.Sync()
|
|
|
|
// Load configuration
|
|
cfg, err := config.LoadConfig()
|
|
if err != nil {
|
|
logging.Fatal("failed to load config", logging.Err(err))
|
|
}
|
|
|
|
// Connect to Temporal
|
|
c, err := client.Dial(client.Options{
|
|
HostPort: cfg.Temporal.HostPort,
|
|
Namespace: cfg.Temporal.Namespace,
|
|
})
|
|
if err != nil {
|
|
logging.Fatal("failed to connect to temporal", logging.Err(err))
|
|
}
|
|
defer c.Close()
|
|
|
|
// Create worker
|
|
w := worker.New(c, "poimen-taskqueue", worker.Options{})
|
|
if w == nil {
|
|
logging.Fatal("failed to create worker")
|
|
}
|
|
|
|
// Register all workflows
|
|
w.RegisterWorkflow(workflow.OrchestratorWorkflow)
|
|
w.RegisterWorkflow(workflow.TaskUnitWorkflow)
|
|
w.RegisterWorkflow(workflow.TestWorkflow)
|
|
w.RegisterWorkflow(workflow.RoutingWorkflow)
|
|
w.RegisterWorkflow(workflow.WorkflowGraphQuery)
|
|
|
|
// Register all activities
|
|
w.RegisterActivity(activity.CloneRepoActivity)
|
|
w.RegisterActivity(activity.GitWorktreeAddActivity)
|
|
w.RegisterActivity(activity.GitCommitActivity)
|
|
w.RegisterActivity(activity.GitPushActivity)
|
|
w.RegisterActivity(activity.GitSquashMergeActivity)
|
|
w.RegisterActivity(activity.GitDiffActivity)
|
|
w.RegisterActivity(activity.PrepareSkillsActivity)
|
|
w.RegisterActivity(activity.PlanningActivity)
|
|
w.RegisterActivity(activity.ImplementerActivity)
|
|
w.RegisterActivity(activity.JudgeActivity)
|
|
// Integration and lessons activities - register when fully tested
|
|
w.RegisterActivity(activity.RunIntegrationTestActivity)
|
|
// w.RegisterActivity(activity.UpdateLessonsActivity)
|
|
// w.RegisterActivity(activity.ReadLessonsActivity)
|
|
|
|
// Routing workflow activities
|
|
w.RegisterActivity(activity.LLMRouterActivity)
|
|
w.RegisterActivity(activity.ValidateWorkflowSpecActivity)
|
|
w.RegisterActivity(activity.ValidateCronWorkflowSpecActivity)
|
|
|
|
// Analysis activities
|
|
w.RegisterActivity(activity.AnalyzeCodeActivity)
|
|
w.RegisterActivity(activity.SecurityScanActivity)
|
|
w.RegisterActivity(activity.GenerateReportActivity)
|
|
|
|
// Notification and utility activities
|
|
w.RegisterActivity(activity.NotifyStatusActivity)
|
|
w.RegisterActivity(activity.ArchiveResultsActivity)
|
|
w.RegisterActivity(activity.DeploymentPreCheckActivity)
|
|
w.RegisterActivity(activity.ApproveWorkflowActivity)
|
|
|
|
// Authentication activities
|
|
w.RegisterActivity(activity.AssumeRoleActivity)
|
|
|
|
// Memory activities
|
|
w.RegisterActivity(activity.RetrieveMemoryActivity)
|
|
|
|
// GraphRAG activities
|
|
w.RegisterActivity(activity.FetchCanvasRelationsActivity)
|
|
w.RegisterActivity(activity.QueryGraphRAGActivity)
|
|
w.RegisterActivity(activity.CanvasReasonerActivity)
|
|
w.RegisterActivity(activity.IndexGraphRAGActivity)
|
|
w.RegisterActivity(activity.CanvasCompatibilityActivity)
|
|
|
|
// Initialize health checker
|
|
healthChecker := health.NewChecker(c)
|
|
healthHandler := health.NewHandler(healthChecker)
|
|
|
|
// Set up HTTP server for health checks
|
|
mux := http.NewServeMux()
|
|
healthHandler.RegisterRoutes(mux)
|
|
|
|
healthServer := &http.Server{
|
|
Addr: ":8081",
|
|
Handler: mux,
|
|
}
|
|
|
|
// Start health check server in a goroutine
|
|
go func() {
|
|
log.Printf("Health check server listening on %s", healthServer.Addr)
|
|
if err := healthServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
log.Printf("health check server error: %v", err)
|
|
}
|
|
}()
|
|
|
|
// Set up signal handling for graceful shutdown
|
|
sigChan := make(chan os.Signal, 1)
|
|
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
|
|
|
|
// Run worker in a goroutine
|
|
workerErrChan := make(chan error, 1)
|
|
go func() {
|
|
logging.Info("starting worker", logging.String("queue", "poimen"))
|
|
if err := w.Run(worker.InterruptCh()); err != nil {
|
|
workerErrChan <- err
|
|
}
|
|
}()
|
|
|
|
// Wait for either worker error or signal
|
|
select {
|
|
case err := <-workerErrChan:
|
|
logging.Fatal("worker failed", logging.Err(err))
|
|
case sig := <-sigChan:
|
|
logging.Info("received signal", logging.String("signal", sig.String()))
|
|
w.Stop()
|
|
|
|
// Shutdown health check server
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if err := healthServer.Shutdown(ctx); err != nil {
|
|
logging.Warn("health check server shutdown error", logging.Err(err))
|
|
}
|
|
logging.Info("worker shutdown complete")
|
|
}
|
|
}
|