128 lines
4.3 KiB
Go
128 lines
4.3 KiB
Go
package main
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"fmt"
|
|
"log"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"strconv"
|
|
"strings"
|
|
"syscall"
|
|
"time"
|
|
|
|
"neuralhunt/internal/controllerui"
|
|
"neuralhunt/internal/customer"
|
|
"neuralhunt/internal/servicecontroller"
|
|
)
|
|
|
|
func loadDotEnv(path string) {
|
|
f, err := os.Open(path)
|
|
if err != nil {
|
|
return
|
|
}
|
|
defer f.Close()
|
|
s := bufio.NewScanner(f)
|
|
for s.Scan() {
|
|
line := strings.TrimSpace(s.Text())
|
|
if line == "" || strings.HasPrefix(line, "#") {
|
|
continue
|
|
}
|
|
if strings.HasPrefix(line, "export ") {
|
|
line = strings.TrimSpace(strings.TrimPrefix(line, "export "))
|
|
}
|
|
k, v, ok := strings.Cut(line, "=")
|
|
if !ok {
|
|
continue
|
|
}
|
|
k, v = strings.TrimSpace(k), strings.TrimSpace(v)
|
|
if k == "" {
|
|
continue
|
|
}
|
|
if _, exists := os.LookupEnv(k); exists {
|
|
continue
|
|
}
|
|
if len(v) >= 2 && ((v[0] == '"' && v[len(v)-1] == '"') || (v[0] == '\'' && v[len(v)-1] == '\'')) {
|
|
v = v[1 : len(v)-1]
|
|
}
|
|
_ = os.Setenv(k, v)
|
|
}
|
|
}
|
|
func env(k, d string) string {
|
|
if v := strings.TrimSpace(os.Getenv(k)); v != "" {
|
|
return v
|
|
}
|
|
return d
|
|
}
|
|
func intEnv(k string, d int) int {
|
|
v, e := strconv.Atoi(strings.TrimSpace(os.Getenv(k)))
|
|
if e != nil {
|
|
return d
|
|
}
|
|
return v
|
|
}
|
|
func durationEnv(k string, d time.Duration) time.Duration {
|
|
v := strings.TrimSpace(os.Getenv(k))
|
|
if v == "" {
|
|
return d
|
|
}
|
|
x, e := time.ParseDuration(v)
|
|
if e != nil {
|
|
return d
|
|
}
|
|
return x
|
|
}
|
|
|
|
func main() {
|
|
loadDotEnv(".env")
|
|
ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
|
|
defer cancel()
|
|
cfg := servicecontroller.Config{
|
|
ID: env("SC_ID", ""), Name: env("SC_NAME", env("SC_ID", "Service Controller")),
|
|
MasterURL: strings.TrimRight(env("SC_MASTER_URL", ""), "/"), AdvertiseURL: strings.TrimRight(env("SC_ADVERTISE_URL", ""), "/"), SharedSecret: env("SERVICE_CONTROLLER_SHARED_SECRET", ""),
|
|
ControlAddr: env("SC_CONTROL_ADDR", ":8102"), AdminAddr: env("SC_ADMIN_ADDR", ":8101"), AdminUser: env("SC_ADMIN_USER", "admin"), AdminPassword: env("SC_ADMIN_PASSWORD", ""), AdminCookieSecureMode: env("SC_ADMIN_COOKIE_SECURE", "auto"),
|
|
MaxWorkers: intEnv("SC_MAX_WORKERS", 500), MaxRunning: intEnv("SC_MAX_RUNNING_WORKERS", 100), HeartbeatInterval: durationEnv("SC_HEARTBEAT_INTERVAL", 15*time.Second),
|
|
DefaultWorkerImage: env("SC_WORKER_IMAGE", env("CS_WORKER_IMAGE", "neuralhunt-worker:local")), WorkerNetworkOverride: env("SC_WORKER_NETWORK", ""), WorkerRegisterURLOverride: env("SC_WORKER_REGISTER_URL", ""), GameURLOverride: env("SC_GAME_PUBLIC_URL", ""), RegistryUsername: env("SC_WORKER_REGISTRY_USERNAME", env("CS_WORKER_REGISTRY_USERNAME", "")), RegistryPassword: env("SC_WORKER_REGISTRY_PASSWORD", env("CS_WORKER_REGISTRY_PASSWORD", "")), RegistryServer: env("SC_WORKER_REGISTRY_SERVER", env("CS_WORKER_REGISTRY_SERVER", "")),
|
|
}
|
|
if err := servicecontroller.ValidateConfig(cfg); err != nil {
|
|
log.Fatal(err)
|
|
}
|
|
docker, err := customer.NewDockerClient(env("DOCKER_HOST", "unix:///var/run/docker.sock"))
|
|
if err != nil {
|
|
log.Fatal(err)
|
|
}
|
|
pingCtx, pingCancel := context.WithTimeout(ctx, 5*time.Second)
|
|
err = docker.Ping(pingCtx)
|
|
pingCancel()
|
|
if err != nil {
|
|
log.Fatalf("Docker Engine API unavailable: %v", err)
|
|
}
|
|
svc := servicecontroller.New(docker, cfg)
|
|
servers := []*http.Server{
|
|
{Addr: cfg.ControlAddr, Handler: svc.ControlRoutes(), ReadHeaderTimeout: 5 * time.Second, IdleTimeout: 45 * time.Second},
|
|
{Addr: cfg.AdminAddr, Handler: svc.AdminRoutes(controllerui.Handler()), ReadHeaderTimeout: 5 * time.Second, IdleTimeout: 60 * time.Second},
|
|
}
|
|
names := []string{"service-controller control/private", "service-controller admin/private"}
|
|
for i, s := range servers {
|
|
go func(n string, hs *http.Server) {
|
|
log.Printf("%s listener on %s", n, hs.Addr)
|
|
if err := hs.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
log.Fatal(err)
|
|
}
|
|
}(names[i], s)
|
|
}
|
|
// Start registration only after the control listener is live enough for the
|
|
// Master's bidirectional health check.
|
|
go func() { time.Sleep(250 * time.Millisecond); svc.RunRegistration(ctx) }()
|
|
log.Printf("service-controller %s (%s) registering at %s and advertising %s", cfg.ID, cfg.Name, cfg.MasterURL, cfg.AdvertiseURL)
|
|
<-ctx.Done()
|
|
shutdown, done := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer done()
|
|
for _, s := range servers {
|
|
_ = s.Shutdown(shutdown)
|
|
}
|
|
fmt.Println("service-controller stopped")
|
|
}
|