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/action" "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/statemachine" ) 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(statemachine.OrchestratorWorkflow) w.RegisterWorkflow(statemachine.TaskUnitWorkflow) w.RegisterWorkflow(statemachine.TestWorkflow) w.RegisterWorkflow(statemachine.RoutingWorkflow) w.RegisterWorkflow(statemachine.WorkflowGraphQuery) // Register all activities w.RegisterActivity(action.CloneRepoActivity) w.RegisterActivity(action.GitWorktreeAddActivity) w.RegisterActivity(action.GitCommitActivity) w.RegisterActivity(action.GitPushActivity) w.RegisterActivity(action.GitSquashMergeActivity) w.RegisterActivity(action.GitDiffActivity) w.RegisterActivity(action.PrepareSkillsActivity) w.RegisterActivity(action.PlanningActivity) w.RegisterActivity(action.ImplementerActivity) w.RegisterActivity(action.JudgeActivity) // Integration and lessons activities - register when fully tested w.RegisterActivity(action.RunIntegrationTestActivity) // w.RegisterActivity(action.UpdateLessonsActivity) // w.RegisterActivity(action.ReadLessonsActivity) // Routing workflow activities w.RegisterActivity(action.LLMRouterActivity) w.RegisterActivity(action.ValidateWorkflowSpecActivity) w.RegisterActivity(action.ValidateCronWorkflowSpecActivity) // Analysis activities w.RegisterActivity(action.AnalyzeCodeActivity) w.RegisterActivity(action.SecurityScanActivity) w.RegisterActivity(action.GenerateReportActivity) // Notification and utility activities w.RegisterActivity(action.NotifyStatusActivity) w.RegisterActivity(action.ArchiveResultsActivity) w.RegisterActivity(action.DeploymentPreCheckActivity) w.RegisterActivity(action.ApproveWorkflowActivity) // Authentication activities w.RegisterActivity(action.AssumeRoleActivity) // Memory activities w.RegisterActivity(action.RetrieveMemoryActivity) // GraphRAG activities w.RegisterActivity(action.FetchCanvasRelationsActivity) w.RegisterActivity(action.QueryGraphRAGActivity) w.RegisterActivity(action.CanvasReasonerActivity) w.RegisterActivity(action.IndexGraphRAGActivity) w.RegisterActivity(action.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") } }