Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
93 changes: 89 additions & 4 deletions cmd/kleffd/main.go
Original file line number Diff line number Diff line change
@@ -1,13 +1,24 @@
package main

import (
"context"
"fmt"
"log"
"os"
"os/signal"
"syscall"

"github.com/kleffio/kleff-daemon/internal/adapters/out/db"
"github.com/kleffio/kleff-daemon/internal/adapters/out/observability/logging"
queueadapter "github.com/kleffio/kleff-daemon/internal/adapters/out/queue"
memrepo "github.com/kleffio/kleff-daemon/internal/adapters/out/repository/memory"
dockeradapter "github.com/kleffio/kleff-daemon/internal/adapters/out/runtime/docker"
k8sadapter "github.com/kleffio/kleff-daemon/internal/adapters/out/runtime/kubernetes"
"github.com/kleffio/kleff-daemon/internal/app/config"
"github.com/kleffio/kleff-daemon/internal/application/ports"
"github.com/kleffio/kleff-daemon/internal/workers"
"github.com/kleffio/kleff-daemon/internal/workers/jobs"
"k8s.io/client-go/rest"
)

func main() {
Expand All @@ -17,17 +28,91 @@ func main() {
}

baseLogger := logging.NewSlogAdapter()

daemonLog := baseLogger.With(ports.LogKeyNodeID, cfg.NodeID)

daemonLog.Info("Daemon starting", "runtime_mode", cfg.RuntimeMode)
// --- Runtime adapter: auto-detect Docker vs Kubernetes ---
runtime, err := detectRuntime(cfg, daemonLog)
if err != nil {
daemonLog.Error("Failed to initialize runtime adapter", err)
os.Exit(1)
}

// --- Database ---
sqliteDB, err := db.InitDB(cfg.DatabasePath)
if err != nil {
daemonLog.Error("Failed to initialize database", err, "path", cfg.DatabasePath)
os.Exit(1)
}
defer sqliteDB.Close()
daemonLog.Info("Database initialized successfully", "path", cfg.DatabasePath)

daemonLog.Info("Database initialized", "path", cfg.DatabasePath)

// --- Queue ---
var q ports.Queue
switch cfg.QueueBackend {
case config.QueueBackendRedis:
rq, err := queueadapter.NewRedisQueue(cfg.RedisURL, cfg.RedisPassword, cfg.RedisTLS)
if err != nil {
daemonLog.Error("Failed to initialize Redis queue", err)
os.Exit(1)
}
q = rq
daemonLog.Info("Queue backend: Redis", "url", cfg.RedisURL)
default:
q = queueadapter.NewMemoryQueue()
daemonLog.Info("Queue backend: in-memory")
}

// --- Repository ---
repo := memrepo.NewServerRepository()

// --- Dispatcher + workers ---
dispatcher := workers.NewDispatcher(q, 4)
dispatcher.Register(jobs.JobTypeServerProvision, workers.NewProvisionWorker(runtime, repo, daemonLog).Handle)
dispatcher.Register(jobs.JobTypeServerStart, workers.NewStartWorker(runtime, repo, daemonLog).Handle)
dispatcher.Register(jobs.JobTypeServerStop, workers.NewStopWorker(runtime, repo, daemonLog).Handle)
dispatcher.Register(jobs.JobTypeServerDelete, workers.NewDeleteWorker(runtime, repo, daemonLog).Handle)
dispatcher.Register(jobs.JobTypeServerRestart, workers.NewRestartWorker(runtime, repo, daemonLog).Handle)

daemonLog.Info("Daemon started", "node_id", cfg.NodeID, "grpc_port", cfg.GRPCPort)

ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()

dispatcher.Run(ctx)
daemonLog.Info("Daemon shutdown complete")
}

func detectRuntime(cfg *config.Config, logger ports.Logger) (ports.RuntimeAdapter, error) {
ctx := context.Background()

// Explicit kubeconfig → always Kubernetes, fail hard if unreachable.
if cfg.Kubeconfig != "" {
adapter, err := k8sadapter.New(cfg.Kubeconfig, cfg.KubeNamespace, cfg.NodeID)
if err != nil {
return nil, fmt.Errorf("kubernetes runtime unavailable: %w", err)
}
logger.Info("Runtime: Kubernetes", "kubeconfig", cfg.Kubeconfig)
return adapter, nil
}

// In-cluster environment → Kubernetes, fail hard if unreachable.
if _, err := rest.InClusterConfig(); err == nil {
adapter, err := k8sadapter.New("", cfg.KubeNamespace, cfg.NodeID)
if err != nil {
return nil, fmt.Errorf("kubernetes in-cluster runtime unavailable: %w", err)
}
logger.Info("Runtime: Kubernetes (in-cluster)")
return adapter, nil
}

// No Kubernetes — check if Docker is actually reachable before using it.
dockerAdapter, err := dockeradapter.New(cfg.NodeID)
if err != nil {
return nil, fmt.Errorf("no runtime available: kubernetes not detected, docker client failed: %w", err)
}
if err := dockerAdapter.Ping(ctx); err != nil {
return nil, fmt.Errorf("no runtime available: kubernetes not detected, docker unreachable: %w", err)
}
logger.Info("Runtime: Docker")
return dockerAdapter, nil
}
169 changes: 169 additions & 0 deletions cmd/testprovision/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,169 @@
package main

import (
"context"
"fmt"
"log"
"os"
"time"

dockeradapter "github.com/kleffio/kleff-daemon/internal/adapters/out/runtime/docker"
k8sadapter "github.com/kleffio/kleff-daemon/internal/adapters/out/runtime/kubernetes"
"github.com/kleffio/kleff-daemon/internal/application/ports"
)

// testprovision provisions a Minecraft server via either the Kubernetes or Docker adapter.
//
// Kubernetes usage (on the cluster node):
//
// kubectl proxy --port=8888 &
// ./testprovision k8s
// ./testprovision k8s cleanup
//
// Docker usage (local):
//
// ./testprovision docker
// ./testprovision docker cleanup
func main() {
const (
nodeID = "test-node"
serverID = "kleff-test-minecraft"
ownerID = "test-owner"
blueprintID = "minecraft-vanilla"
)

mode := "k8s"
if len(os.Args) > 1 {
mode = os.Args[1]
}
cleanup := len(os.Args) > 2 && os.Args[2] == "cleanup"

spec := ports.WorkloadSpec{
OwnerID: ownerID,
ServerID: serverID,
BlueprintID: blueprintID,
Image: "itzg/minecraft-server",
MemoryBytes: 2048 * 1024 * 1024,
CPUMillicores: 500,
EnvOverrides: map[string]string{
"EULA": "TRUE",
"TYPE": "VANILLA",
"VERSION": "LATEST",
"DIFFICULTY": "normal",
"MAX_PLAYERS": "20",
"MEMORY": "2G",
},
PortRequirements: []ports.PortRequirement{
{TargetPort: 25565, Protocol: "tcp"},
{TargetPort: 25565, Protocol: "udp"},
{TargetPort: 25575, Protocol: "tcp"},
},
RuntimeHints: ports.RuntimeHints{
PersistentStorage: true,
StoragePath: "/data",
},
}

switch mode {
case "docker":
runDocker(cleanup, serverID, spec)
default:
spec.RuntimeHints.KubernetesStrategy = "agones"
spec.RuntimeHints.ExposeUDP = true
runK8s(cleanup, serverID, spec)
}
}

func runDocker(cleanup bool, serverID string, spec ports.WorkloadSpec) {
adapter, err := dockeradapter.New("test-node")
if err != nil {
log.Fatalf("failed to create docker adapter: %v", err)
}

ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()

if cleanup {
fmt.Printf("Removing %s...\n", serverID)
if err := adapter.Remove(ctx, serverID); err != nil {
log.Fatalf("remove failed: %v", err)
}
fmt.Println("Removed.")
return
}

fmt.Printf("Provisioning %s (%s) via Docker...\n", serverID, spec.Image)
fmt.Printf(" Memory: %d MB CPU: %dm\n", spec.MemoryBytes/1024/1024, spec.CPUMillicores)
fmt.Printf(" Ports: 25565/tcp, 25565/udp, 25575/tcp\n")
fmt.Printf(" Storage: %s\n", spec.RuntimeHints.StoragePath)
fmt.Println(" Pulling image and starting container...")

start := time.Now()
server, err := adapter.Deploy(ctx, spec)
if err != nil {
log.Fatalf("deploy failed: %v", err)
}

fmt.Printf("\nServer is running! (took %s)\n", time.Since(start).Round(time.Second))
fmt.Printf(" RuntimeRef : %s\n", server.RuntimeRef)
fmt.Printf(" State : %s\n", server.State)

endpoint, err := adapter.Endpoint(ctx, serverID)
if err != nil {
fmt.Printf(" Endpoint : (could not resolve: %v)\n", err)
} else {
fmt.Printf(" Endpoint : %s\n", endpoint)
}

fmt.Printf("\nTo clean up: ./testprovision docker cleanup\n")
}

func runK8s(cleanup bool, serverID string, spec ports.WorkloadSpec) {
const (
proxyURL = "http://localhost:8888"
namespace = "default"
)

adapter, err := k8sadapter.New(proxyURL, namespace, "test-node")
if err != nil {
log.Fatalf("failed to create kubernetes adapter: %v", err)
}

ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
defer cancel()

if cleanup {
fmt.Printf("Removing %s...\n", serverID)
if err := adapter.Remove(ctx, serverID); err != nil {
log.Fatalf("remove failed: %v", err)
}
fmt.Println("Removed.")
return
}

fmt.Printf("Provisioning %s (%s) via Agones...\n", serverID, spec.Image)
fmt.Printf(" Memory: %d MB CPU: %dm\n", spec.MemoryBytes/1024/1024, spec.CPUMillicores)
fmt.Printf(" Ports: 25565/tcp, 25565/udp, 25575/tcp\n")
fmt.Println(" Waiting for GameServer to become Ready (up to 10 min)...")

start := time.Now()
server, err := adapter.Deploy(ctx, spec)
if err != nil {
log.Fatalf("deploy failed: %v", err)
}

fmt.Printf("\nServer is ready! (took %s)\n", time.Since(start).Round(time.Second))
fmt.Printf(" RuntimeRef : %s\n", server.RuntimeRef)
fmt.Printf(" State : %s\n", server.State)
fmt.Printf(" NodeID : %s\n", server.Labels.NodeID)
fmt.Printf(" ServerID : %s\n", server.Labels.ServerID)

endpoint, err := adapter.Endpoint(ctx, serverID)
if err != nil {
fmt.Printf(" Endpoint : (could not resolve: %v)\n", err)
} else {
fmt.Printf(" Endpoint : %s\n", endpoint)
}

fmt.Printf("\nTo clean up: ./testprovision k8s cleanup\n")
}
48 changes: 41 additions & 7 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -4,54 +4,88 @@ go 1.25.5

require (
github.com/alicebob/miniredis/v2 v2.37.0
github.com/docker/docker v28.5.2+incompatible
github.com/docker/go-connections v0.6.0
github.com/redis/go-redis/v9 v9.18.0
github.com/spf13/pflag v1.0.10
github.com/spf13/viper v1.21.0
k8s.io/api v0.29.0
k8s.io/apimachinery v0.29.0
k8s.io/client-go v0.29.0
modernc.org/sqlite v1.46.1
)

require (
github.com/Microsoft/go-winio v0.4.21 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/containerd/errdefs v1.0.0 // indirect
github.com/containerd/errdefs/pkg v0.3.0 // indirect
github.com/containerd/log v0.1.0 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect
github.com/distribution/reference v0.6.0 // indirect
github.com/docker/go-units v0.5.0 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/emicklei/go-restful/v3 v3.11.0 // indirect
github.com/felixge/httpsnoop v1.0.4 // indirect
github.com/fsnotify/fsnotify v1.9.0 // indirect
github.com/go-logr/logr v1.4.2 // indirect
github.com/go-logr/logr v1.4.3 // indirect
github.com/go-logr/stdr v1.2.2 // indirect
github.com/go-openapi/jsonpointer v0.19.6 // indirect
github.com/go-openapi/jsonreference v0.20.2 // indirect
github.com/go-openapi/swag v0.22.3 // indirect
github.com/go-viper/mapstructure/v2 v2.4.0 // indirect
github.com/gogo/protobuf v1.3.2 // indirect
github.com/golang/protobuf v1.5.4 // indirect
github.com/google/gnostic-models v0.6.8 // indirect
github.com/google/gofuzz v1.2.0 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/imdario/mergo v0.3.6 // indirect
github.com/josharian/intern v1.0.0 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/mailru/easyjson v0.7.7 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect
github.com/moby/docker-image-spec v1.3.1 // indirect
github.com/moby/sys/atomicwriter v0.1.0 // indirect
github.com/moby/term v0.5.2 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
github.com/morikuni/aec v1.1.0 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/ncruces/go-strftime v1.0.0 // indirect
github.com/opencontainers/go-digest v1.0.0 // indirect
github.com/opencontainers/image-spec v1.1.1 // indirect
github.com/pelletier/go-toml/v2 v2.2.4 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/sagikazarmark/locafero v0.11.0 // indirect
github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 // indirect
github.com/spf13/afero v1.15.0 // indirect
github.com/spf13/cast v1.10.0 // indirect
github.com/subosito/gotenv v1.6.0 // indirect
github.com/yuin/gopher-lua v1.1.1 // indirect
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.68.0 // indirect
go.opentelemetry.io/otel v1.43.0 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.43.0 // indirect
go.opentelemetry.io/otel/metric v1.43.0 // indirect
go.opentelemetry.io/otel/trace v1.43.0 // indirect
go.uber.org/atomic v1.11.0 // indirect
go.yaml.in/yaml/v3 v3.0.4 // indirect
golang.org/x/exp v0.0.0-20251023183803-a4bb9ffd2546 // indirect
golang.org/x/net v0.46.0 // indirect
golang.org/x/net v0.52.0 // indirect
golang.org/x/oauth2 v0.20.0 // indirect
golang.org/x/sys v0.37.0 // indirect
golang.org/x/term v0.36.0 // indirect
golang.org/x/text v0.30.0 // indirect
golang.org/x/sys v0.42.0 // indirect
golang.org/x/term v0.41.0 // indirect
golang.org/x/text v0.35.0 // indirect
golang.org/x/time v0.3.0 // indirect
google.golang.org/protobuf v1.34.2 // indirect
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c // indirect
google.golang.org/protobuf v1.36.11 // indirect
gopkg.in/inf.v0 v0.9.1 // indirect
gopkg.in/yaml.v2 v2.4.0 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
gotest.tools/v3 v3.5.2 // indirect
k8s.io/klog/v2 v2.130.1 // indirect
k8s.io/kube-openapi v0.0.0-20231010175941-2dd684a91f00 // indirect
k8s.io/utils v0.0.0-20230726121419-3b25d923346b // indirect
modernc.org/libc v1.67.6 // indirect
modernc.org/mathutil v1.7.1 // indirect
Expand Down
Loading
Loading