2026-09-05 00:43:25 -07:00
|
|
|
package main
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
|
|
|
|
"flag"
|
|
|
|
|
"log"
|
|
|
|
|
"os"
|
|
|
|
|
"os/signal"
|
|
|
|
|
"sync"
|
|
|
|
|
"syscall"
|
|
|
|
|
|
|
|
|
|
"go.temporal.io/sdk/client"
|
|
|
|
|
"go.temporal.io/sdk/worker"
|
|
|
|
|
|
|
|
|
|
"github.com/rockliang/poimen/workflows/action"
|
|
|
|
|
"github.com/rockliang/poimen/workflows/internal/api"
|
|
|
|
|
"github.com/rockliang/poimen/workflows/internal/config"
|
|
|
|
|
"github.com/rockliang/poimen/workflows/pkg/db"
|
|
|
|
|
"github.com/rockliang/poimen/workflows/statemachine"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
func main() {
|
|
|
|
|
var (
|
|
|
|
|
apiPort = flag.Int("port", 8080, "HTTP API port")
|
|
|
|
|
verbose = flag.Bool("verbose", false, "verbose logging")
|
|
|
|
|
)
|
|
|
|
|
flag.Parse()
|
|
|
|
|
|
|
|
|
|
logger := log.New(os.Stdout, "[poimen-server] ", log.LstdFlags|log.Lshortfile)
|
|
|
|
|
|
|
|
|
|
// Load configuration
|
|
|
|
|
cfg, err := config.LoadConfig()
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Fatalf("failed to load config: %v", err)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Connect to database (memory-db via K8s CNPG)
|
|
|
|
|
logger.Println("connecting to database...")
|
|
|
|
|
database, err := db.New(os.Getenv("DATABASE_URL"))
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Fatalf("failed to connect to database: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer database.Close()
|
|
|
|
|
logger.Println("✓ Connected to database")
|
|
|
|
|
|
|
|
|
|
// Connect to Temporal
|
|
|
|
|
logger.Printf("connecting to Temporal at %s", cfg.Temporal.HostPort)
|
|
|
|
|
c, err := client.Dial(client.Options{
|
|
|
|
|
HostPort: cfg.Temporal.HostPort,
|
|
|
|
|
Namespace: cfg.Temporal.Namespace,
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Fatalf("failed to connect to temporal: %v", err)
|
|
|
|
|
}
|
|
|
|
|
defer c.Close()
|
|
|
|
|
|
|
|
|
|
logger.Println("✓ Connected to Temporal")
|
|
|
|
|
|
|
|
|
|
// Create and start Temporal worker
|
|
|
|
|
w := worker.New(c, "default", worker.Options{})
|
|
|
|
|
|
|
|
|
|
// Register RoutingWorkflow
|
|
|
|
|
w.RegisterWorkflow(statemachine.RoutingWorkflow)
|
|
|
|
|
|
|
|
|
|
// Register activities
|
|
|
|
|
w.RegisterActivity(action.CloneRepoActivity)
|
|
|
|
|
w.RegisterActivity(action.AnalyzeCodeActivity)
|
|
|
|
|
w.RegisterActivity(action.SecurityScanActivity)
|
|
|
|
|
w.RegisterActivity(action.GenerateReportActivity)
|
|
|
|
|
w.RegisterActivity(action.DeploymentPreCheckActivity)
|
|
|
|
|
w.RegisterActivity(action.NotifyStatusActivity)
|
|
|
|
|
w.RegisterActivity(action.ApproveWorkflowActivity)
|
|
|
|
|
w.RegisterActivity(action.ArchiveResultsActivity)
|
|
|
|
|
w.RegisterActivity(action.RetrieveMemoryActivity)
|
|
|
|
|
w.RegisterActivity(action.AssumeRoleActivity)
|
|
|
|
|
w.RegisterActivity(action.LLMInferenceActivity)
|
|
|
|
|
w.RegisterActivity(action.LLMBatchInferenceActivity)
|
2026-09-05 00:53:56 -07:00
|
|
|
w.RegisterActivity(action.CanvasReasonerActivity)
|
2026-09-05 00:43:25 -07:00
|
|
|
|
|
|
|
|
var wg sync.WaitGroup
|
|
|
|
|
errChan := make(chan error, 2)
|
|
|
|
|
|
|
|
|
|
// Start Temporal worker
|
|
|
|
|
wg.Add(1)
|
|
|
|
|
go func() {
|
|
|
|
|
defer wg.Done()
|
|
|
|
|
logger.Println("starting Temporal worker...")
|
|
|
|
|
if err := w.Run(worker.InterruptCh()); err != nil {
|
|
|
|
|
errChan <- err
|
|
|
|
|
}
|
|
|
|
|
}()
|
|
|
|
|
|
|
|
|
|
// Start HTTP API server
|
|
|
|
|
wg.Add(1)
|
|
|
|
|
go func() {
|
|
|
|
|
defer wg.Done()
|
|
|
|
|
server := api.NewServer(database, c, logger)
|
|
|
|
|
logger.Printf("starting API server on port %d", *apiPort)
|
|
|
|
|
if err := server.Start(*apiPort); err != nil {
|
|
|
|
|
errChan <- err
|
|
|
|
|
}
|
|
|
|
|
}()
|
|
|
|
|
|
|
|
|
|
// Wait for interrupt signal
|
|
|
|
|
sigChan := make(chan os.Signal, 1)
|
|
|
|
|
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
|
|
|
|
|
|
|
|
|
|
go func() {
|
|
|
|
|
sig := <-sigChan
|
|
|
|
|
logger.Printf("received signal: %v", sig)
|
|
|
|
|
w.Stop()
|
|
|
|
|
}()
|
|
|
|
|
|
|
|
|
|
// Monitor for errors
|
|
|
|
|
go func() {
|
|
|
|
|
err := <-errChan
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Printf("error: %v", err)
|
|
|
|
|
w.Stop()
|
|
|
|
|
}
|
|
|
|
|
}()
|
|
|
|
|
|
|
|
|
|
wg.Wait()
|
|
|
|
|
logger.Println("✓ Server stopped gracefully")
|
|
|
|
|
}
|