118 lines
3.2 KiB
Go
118 lines
3.2 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/mattermost"
|
|
modelopenai "github.com/nomadcode/nomadcode-core/internal/adapters/openai"
|
|
"github.com/nomadcode/nomadcode-core/internal/agent"
|
|
"github.com/nomadcode/nomadcode-core/internal/config"
|
|
"github.com/nomadcode/nomadcode-core/internal/db"
|
|
apphttp "github.com/nomadcode/nomadcode-core/internal/http"
|
|
"github.com/nomadcode/nomadcode-core/internal/notification"
|
|
"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,
|
|
}, logger)
|
|
notificationService := notification.NewService(mattermostClient, logger)
|
|
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)
|
|
}
|
|
|
|
taskScheduler, err := scheduler.New(pool, store, notificationService, agentClient, modelClient, 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 := workflow.NewService(store, taskScheduler, logger)
|
|
handler := apphttp.NewHandler(pool, workflowService, logger)
|
|
server := &http.Server{
|
|
Addr: cfg.HTTPAddr,
|
|
Handler: apphttp.NewRouter(handler, logger, apphttp.AuthConfig{Username: cfg.AuthUsername, Password: cfg.AuthPassword}),
|
|
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 := taskScheduler.Stop(shutdownCtx); err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|