Refactor zero-scale gateway into modular Go microservice packages
This commit is contained in:
@@ -1,464 +1,69 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"crypto/tls"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/http/httputil"
|
||||
"net/url"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/banner"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/config"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/proxy"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/reaper"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/sync"
|
||||
)
|
||||
|
||||
var (
|
||||
appsDir = "/opt/apps"
|
||||
configPath = "/etc/zero-scale/artifacts.json"
|
||||
allocatedPorts = make(map[string]int)
|
||||
lastActive = make(map[string]time.Time)
|
||||
mu sync.Mutex
|
||||
|
||||
idleTimeout = 15 * time.Minute
|
||||
checkPeriod = 30 * time.Second
|
||||
)
|
||||
|
||||
type ArtifactItem struct {
|
||||
URL string `json:"url"`
|
||||
Type string `json:"type"`
|
||||
}
|
||||
|
||||
type ConfigMapData struct {
|
||||
Artifacts map[string]ArtifactItem `json:"artifacts"`
|
||||
}
|
||||
|
||||
func main() {
|
||||
if dir := os.Getenv("APPS_DIR"); dir != "" {
|
||||
appsDir = dir
|
||||
}
|
||||
os.MkdirAll(appsDir, 0755)
|
||||
|
||||
go startReaper()
|
||||
go startArtifactSync()
|
||||
|
||||
http.HandleFunc("/", handleRequest)
|
||||
|
||||
log.Println("Starting Zero-Scale Gateway on :80...")
|
||||
if err := http.ListenAndServe(":80", nil); err != nil {
|
||||
log.Fatalf("Server error: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func handleRequest(w http.ResponseWriter, r *http.Request) {
|
||||
host := r.Host
|
||||
if strings.Contains(host, ":") {
|
||||
host = strings.Split(host, ":")[0]
|
||||
}
|
||||
log.Printf("Received %s request for host: %s (path: %s)", r.Method, host, r.URL.Path)
|
||||
|
||||
binaryPath := filepath.Join(appsDir, host)
|
||||
if !fileExists(binaryPath) {
|
||||
// Fallback lookup: match shortName (e.g. bdash2 -> bdash2.kube.fairfaxmedia.net)
|
||||
shortName := strings.Split(host, ".")[0]
|
||||
entries, err := os.ReadDir(appsDir)
|
||||
if err == nil {
|
||||
for _, entry := range entries {
|
||||
if !entry.IsDir() && (strings.HasPrefix(entry.Name(), shortName+".") || entry.Name() == shortName) {
|
||||
binaryPath = filepath.Join(appsDir, entry.Name())
|
||||
log.Printf("Resolved host %s to binary %s", host, binaryPath)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !fileExists(binaryPath) {
|
||||
http.Error(w, fmt.Sprintf("Application %s not found", host), http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
|
||||
mu.Lock()
|
||||
lastActive[host] = time.Now()
|
||||
|
||||
active, err := isServiceActive(unitName)
|
||||
// 1. Load configuration and CLI flags
|
||||
cfg, err := config.Load()
|
||||
if err != nil {
|
||||
mu.Unlock()
|
||||
log.Printf("Error checking status of unit %s: %v", unitName, err)
|
||||
http.Error(w, "Failed to check service status", http.StatusInternalServerError)
|
||||
return
|
||||
log.Fatalf("Configuration error: %v", err)
|
||||
}
|
||||
|
||||
port, exists := allocatedPorts[host]
|
||||
if active && !exists {
|
||||
// Discover active port from systemd unit Environment if gateway restarted
|
||||
if activePort := getActiveServicePort(unitName); activePort > 0 {
|
||||
allocatedPorts[host] = activePort
|
||||
port = activePort
|
||||
exists = true
|
||||
log.Printf("Discovered active service %s running on port %d", host, port)
|
||||
}
|
||||
// Ensure apps directory exists
|
||||
if err := os.MkdirAll(cfg.AppsDir, 0755); err != nil {
|
||||
log.Fatalf("Failed to create apps directory %s: %v", cfg.AppsDir, err)
|
||||
}
|
||||
|
||||
if !active || !exists {
|
||||
// Service is stopped or unallocated; allocate fresh free port and start
|
||||
if active {
|
||||
stopServiceLocked(host)
|
||||
}
|
||||
freePort, err := getFreePort()
|
||||
if err != nil {
|
||||
mu.Unlock()
|
||||
log.Printf("Failed to get free port for %s: %v", host, err)
|
||||
http.Error(w, "Internal Server Error: No available ports", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
allocatedPorts[host] = freePort
|
||||
port = freePort
|
||||
// 2. Instantiate gateway proxy
|
||||
gw := proxy.New(cfg)
|
||||
|
||||
log.Printf("Host %s matches binary %s. Starting sandboxed systemd unit on port %d...", host, binaryPath, port)
|
||||
err = startTransientService(host, port, binaryPath)
|
||||
if err != nil {
|
||||
delete(allocatedPorts, host)
|
||||
mu.Unlock()
|
||||
log.Printf("Failed to start transient systemd unit for %s: %v", host, err)
|
||||
http.Error(w, "Failed to start application process", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
mu.Unlock()
|
||||
// 3. Start background workers
|
||||
rp := reaper.New(cfg, gw.GetLastActiveMap, gw.StopService)
|
||||
go rp.Start()
|
||||
go sync.StartArtifactSync(cfg)
|
||||
|
||||
// Wait for the app socket to bind and respond
|
||||
targetAddr := fmt.Sprintf("127.0.0.1:%d", port)
|
||||
if err := waitPortReady(targetAddr, 30*time.Second); err != nil {
|
||||
log.Printf("Service %s on port %d did not bind in time: %v", host, port, err)
|
||||
stopService(host)
|
||||
http.Error(w, "Application startup timeout", http.StatusGatewayTimeout)
|
||||
return
|
||||
}
|
||||
log.Printf("Service %s successfully scaled up on port %d!", host, port)
|
||||
} else {
|
||||
mu.Unlock()
|
||||
// 4. Print startup banner with direct URLs and journalctl CLI commands
|
||||
banner.PrintStartup(cfg)
|
||||
|
||||
// 5. Configure HTTP server
|
||||
server := &http.Server{
|
||||
Addr: cfg.ListenAddr,
|
||||
Handler: gw,
|
||||
}
|
||||
|
||||
// Reverse proxy request
|
||||
targetURL, _ := url.Parse(fmt.Sprintf("http://127.0.0.1:%d", port))
|
||||
proxy := httputil.NewSingleHostReverseProxy(targetURL)
|
||||
proxy.ErrorHandler = func(w http.ResponseWriter, r *http.Request, err error) {
|
||||
log.Printf("Proxy error for %s on port %d: %v. Cleaning up service...", host, port, err)
|
||||
stopService(host)
|
||||
http.Error(w, "Bad Gateway", http.StatusBadGateway)
|
||||
// 6. Graceful shutdown handler
|
||||
stop := make(chan os.Signal, 1)
|
||||
signal.Notify(stop, os.Interrupt, syscall.SIGTERM)
|
||||
|
||||
go func() {
|
||||
if err := server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
|
||||
log.Fatalf("Gateway server error: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
<-stop
|
||||
log.Println("Shutting down Zero-Scale Gateway...")
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
if err := server.Shutdown(ctx); err != nil {
|
||||
log.Printf("Server shutdown error: %v", err)
|
||||
}
|
||||
proxy.ServeHTTP(w, r)
|
||||
}
|
||||
|
||||
// getFreePort queries kernel for free ephemeral port
|
||||
func getFreePort() (int, error) {
|
||||
addr, err := net.ResolveTCPAddr("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
l, err := net.ListenTCP("tcp", addr)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
defer l.Close()
|
||||
return l.Addr().(*net.TCPAddr).Port, nil
|
||||
}
|
||||
|
||||
// isServiceActive checks if the dynamic systemd unit is running
|
||||
func isServiceActive(unit string) (bool, error) {
|
||||
cmd := exec.Command("systemctl", "is-active", unit)
|
||||
err := cmd.Run()
|
||||
if err != nil {
|
||||
if _, ok := err.(*exec.ExitError); ok {
|
||||
return false, nil
|
||||
}
|
||||
return false, err
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// getActiveServicePort reads PORT from running systemd service environment
|
||||
func getActiveServicePort(unit string) int {
|
||||
out, err := exec.Command("systemctl", "show", "-p", "Environment", unit).Output()
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
str := string(out)
|
||||
for _, env := range strings.Fields(str) {
|
||||
if strings.HasPrefix(env, "PORT=") || strings.HasPrefix(env, "Environment=PORT=") {
|
||||
parts := strings.Split(env, "=")
|
||||
if len(parts) >= 2 {
|
||||
val := parts[len(parts)-1]
|
||||
if p, err := strconv.Atoi(val); err == nil && p > 0 {
|
||||
return p
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
// fetchK8sResource dynamically retrieves a ConfigMap or Secret from the in-cluster Kubernetes API
|
||||
func fetchK8sResource(resourceType, name string) (map[string]string, error) {
|
||||
tokenBytes, err := os.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/token")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
token := strings.TrimSpace(string(tokenBytes))
|
||||
|
||||
k8sHost := os.Getenv("KUBERNETES_SERVICE_HOST")
|
||||
if k8sHost == "" {
|
||||
k8sHost = "kubernetes.default.svc"
|
||||
}
|
||||
k8sPort := os.Getenv("KUBERNETES_SERVICE_PORT")
|
||||
if k8sPort == "" {
|
||||
k8sPort = "443"
|
||||
}
|
||||
|
||||
reqURL := fmt.Sprintf("https://%s:%s/api/v1/namespaces/default/%s/%s", k8sHost, k8sPort, resourceType, name)
|
||||
req, err := http.NewRequest("GET", reqURL, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
req.Header.Set("Authorization", "Bearer "+token)
|
||||
|
||||
tr := &http.Transport{
|
||||
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
||||
}
|
||||
client := &http.Client{Transport: tr, Timeout: 5 * time.Second}
|
||||
resp, err := client.Do(req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, fmt.Errorf("k8s API returned status %d for %s/%s", resp.StatusCode, resourceType, name)
|
||||
}
|
||||
|
||||
var res struct {
|
||||
Data map[string]string `json:"data"`
|
||||
}
|
||||
if err := json.NewDecoder(resp.Body).Decode(&res); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if resourceType == "secrets" {
|
||||
decoded := make(map[string]string)
|
||||
for k, v := range res.Data {
|
||||
b, err := base64.StdEncoding.DecodeString(v)
|
||||
if err == nil {
|
||||
decoded[k] = string(b)
|
||||
} else {
|
||||
decoded[k] = v
|
||||
}
|
||||
}
|
||||
return decoded, nil
|
||||
}
|
||||
|
||||
return res.Data, nil
|
||||
}
|
||||
|
||||
// startTransientService executes systemd-run with dynamic environment injection
|
||||
func startTransientService(host string, port int, binaryPath string) error {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
typePath := binaryPath + ".type"
|
||||
|
||||
runType := "elf"
|
||||
if fileExists(typePath) {
|
||||
tBytes, err := os.ReadFile(typePath)
|
||||
if err == nil {
|
||||
runType = strings.TrimSpace(string(tBytes))
|
||||
}
|
||||
}
|
||||
|
||||
baseArgs := []string{
|
||||
fmt.Sprintf("--unit=%s", unitName),
|
||||
"-p", "DynamicUser=yes",
|
||||
"-p", "PrivateTmp=yes",
|
||||
"-p", "ProtectSystem=strict",
|
||||
fmt.Sprintf("--setenv=PORT=%d", port),
|
||||
}
|
||||
|
||||
if _, err := os.Stat("/app/uploads"); err == nil {
|
||||
baseArgs = append(baseArgs, "-p", "BindPaths=/app/uploads:/app/uploads")
|
||||
}
|
||||
|
||||
envKeys := []string{
|
||||
"DATABASE_URL", "DATABASE_HOST", "DATABASE_PORT", "DATABASE_USER", "DATABASE_PASSWORD", "DATABASE_NAME",
|
||||
"PGHOST", "PGPORT", "PGUSER", "PGPASSWORD", "PGDATABASE",
|
||||
}
|
||||
for _, k := range envKeys {
|
||||
if v := os.Getenv(k); v != "" {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
|
||||
shortName := strings.Split(host, ".")[0]
|
||||
appCMName := "zero-scale-app-" + shortName
|
||||
if appData, err := fetchK8sResource("configmaps", appCMName); err == nil {
|
||||
for k, v := range appData {
|
||||
if k != "host" && k != "url" && k != "type" && k != "mounts" && k != "scale-to-zero" && k != "database_configmap" && k != "database_secret" {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
if dbCM := appData["database_configmap"]; dbCM != "" {
|
||||
if cmData, err := fetchK8sResource("configmaps", dbCM); err == nil {
|
||||
for k, v := range cmData {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
}
|
||||
if dbSec := appData["database_secret"]; dbSec != "" {
|
||||
if secData, err := fetchK8sResource("secrets", dbSec); err == nil {
|
||||
for k, v := range secData {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
var cmd *exec.Cmd
|
||||
if runType == "wasm" {
|
||||
args := append(baseArgs,
|
||||
"/usr/local/bin/wasmtime", "run",
|
||||
"--tcplisten", fmt.Sprintf("127.0.0.1:%d", port),
|
||||
binaryPath,
|
||||
)
|
||||
cmd = exec.Command("systemd-run", args...)
|
||||
} else {
|
||||
args := append(baseArgs, binaryPath, "-port", strconv.Itoa(port))
|
||||
cmd = exec.Command("systemd-run", args...)
|
||||
}
|
||||
|
||||
log.Printf("Executing: %s", cmd.String())
|
||||
return cmd.Run()
|
||||
}
|
||||
|
||||
// stopServiceLocked shuts down dynamic systemd service unit and clears port state (assumes mu is held)
|
||||
func stopServiceLocked(host string) error {
|
||||
unit := fmt.Sprintf("app-%s", host)
|
||||
cmd := exec.Command("systemctl", "stop", unit)
|
||||
_ = cmd.Run()
|
||||
cmdReset := exec.Command("systemctl", "reset-failed", unit)
|
||||
_ = cmdReset.Run()
|
||||
|
||||
delete(allocatedPorts, host)
|
||||
return nil
|
||||
}
|
||||
|
||||
// stopService safely locks mu and calls stopServiceLocked
|
||||
func stopService(host string) error {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
return stopServiceLocked(host)
|
||||
}
|
||||
|
||||
// waitPortReady checks port every 200ms
|
||||
func waitPortReady(addr string, timeout time.Duration) error {
|
||||
deadline := time.Now().Add(timeout)
|
||||
for time.Now().Before(deadline) {
|
||||
conn, err := net.DialTimeout("tcp", addr, 100*time.Millisecond)
|
||||
if err == nil {
|
||||
conn.Close()
|
||||
return nil
|
||||
}
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
}
|
||||
return fmt.Errorf("timeout waiting for %s", addr)
|
||||
}
|
||||
|
||||
// startReaper checks active services and stops idle ones
|
||||
func startReaper() {
|
||||
ticker := time.NewTicker(checkPeriod)
|
||||
for range ticker.C {
|
||||
mu.Lock()
|
||||
now := time.Now()
|
||||
for host, lastTime := range lastActive {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
active, err := isServiceActive(unitName)
|
||||
if err == nil && active && now.Sub(lastTime) > idleTimeout {
|
||||
log.Printf("Service %s has been idle for %v. Scaling to zero...", host, idleTimeout)
|
||||
if err := stopServiceLocked(host); err != nil {
|
||||
log.Printf("Failed to stop idle unit %s: %v", unitName, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
mu.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
// startArtifactSync polls artifacts config map and downloads binaries dynamically
|
||||
func startArtifactSync() {
|
||||
ticker := time.NewTicker(15 * time.Second)
|
||||
for range ticker.C {
|
||||
if !fileExists(configPath) {
|
||||
continue
|
||||
}
|
||||
data, err := os.ReadFile(configPath)
|
||||
if err != nil {
|
||||
log.Printf("Failed to read artifacts config: %v", err)
|
||||
continue
|
||||
}
|
||||
var cfg ConfigMapData
|
||||
if err := json.Unmarshal(data, &cfg); err != nil {
|
||||
log.Printf("Failed to parse artifacts config JSON: %v", err)
|
||||
continue
|
||||
}
|
||||
|
||||
for host, item := range cfg.Artifacts {
|
||||
targetPath := filepath.Join(appsDir, host)
|
||||
typePath := targetPath + ".type"
|
||||
|
||||
_ = os.WriteFile(typePath, []byte(item.Type), 0644)
|
||||
|
||||
if !fileExists(targetPath) {
|
||||
log.Printf("Downloading artifact for %s from %s...", host, item.URL)
|
||||
err := downloadFile(item.URL, targetPath)
|
||||
if err != nil {
|
||||
log.Printf("Failed to download %s: %v", item.URL, err)
|
||||
continue
|
||||
}
|
||||
os.Chmod(targetPath, 0755)
|
||||
log.Printf("Successfully deployed %s to %s", host, targetPath)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func fileExists(path string) bool {
|
||||
info, err := os.Stat(path)
|
||||
if os.IsNotExist(err) {
|
||||
return false
|
||||
}
|
||||
return err == nil && !info.IsDir()
|
||||
}
|
||||
|
||||
func downloadFile(url, targetPath string) error {
|
||||
resp, err := http.Get(url)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("bad status: %s", resp.Status)
|
||||
}
|
||||
|
||||
out, err := os.Create(targetPath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer out.Close()
|
||||
|
||||
_, err = io.Copy(out, resp.Body)
|
||||
return err
|
||||
log.Println("Zero-Scale Gateway stopped.")
|
||||
}
|
||||
|
||||
57
internal/banner/banner.go
Normal file
57
internal/banner/banner.go
Normal file
@@ -0,0 +1,57 @@
|
||||
package banner
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/config"
|
||||
)
|
||||
|
||||
// PrintStartup prints the initial gateway startup banner to console/journal
|
||||
func PrintStartup(cfg *config.Config) {
|
||||
log.Println("================================================================================")
|
||||
log.Printf("🚀 Zero-Scale Platform Gateway Ready on %s", cfg.ListenAddr)
|
||||
log.Printf("📁 Applications directory: %s", cfg.AppsDir)
|
||||
log.Printf("☸️ Kubernetes Pod: %s", cfg.PodName)
|
||||
log.Printf("☸️ Pod Log Viewer: %s", cfg.KubeLogURL())
|
||||
log.Printf("☸️ Pod Web Shell: %s", cfg.KubeShellURL())
|
||||
log.Println("--------------------------------------------------------------------------------")
|
||||
if entries, err := os.ReadDir(cfg.AppsDir); err == nil {
|
||||
for _, e := range entries {
|
||||
if !e.IsDir() && !strings.HasSuffix(e.Name(), ".type") && e.Name() != "lost+found" {
|
||||
PrintRegistration(e.Name())
|
||||
}
|
||||
}
|
||||
}
|
||||
log.Println("================================================================================")
|
||||
}
|
||||
|
||||
// PrintRegistration prints service registration details with Web Logs link and CLI command
|
||||
func PrintRegistration(host string) {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
log.Printf(" ⭐ Service Registered: https://%s/", host)
|
||||
log.Printf(" 🖥️ Web Logs: https://%s/_logs", host)
|
||||
log.Printf(" 📋 CLI Logs: kubectl exec -n default deployment/zero-scale-gateway -- journalctl -u %s -f", unitName)
|
||||
}
|
||||
|
||||
// PrintScaleUp prints the scale-up lifecycle banner
|
||||
func PrintScaleUp(host string, port int) {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
log.Println("--------------------------------------------------------------------------------")
|
||||
log.Printf("⚡ [SCALE-UP] %s (port %d)", host, port)
|
||||
log.Printf(" 🌐 App URL: https://%s/", host)
|
||||
log.Printf(" 🖥️ Web Logs: https://%s/_logs", host)
|
||||
log.Printf(" 📋 CLI Logs: kubectl exec -n default deployment/zero-scale-gateway -- journalctl -u %s -f", unitName)
|
||||
log.Println("--------------------------------------------------------------------------------")
|
||||
}
|
||||
|
||||
// PrintScaleToZero prints the scale-to-zero lifecycle banner
|
||||
func PrintScaleToZero(host string, idleTimeout time.Duration) {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
log.Printf("💤 [SCALE-TO-ZERO] %s idle for %v (deactivated unit %s)", host, idleTimeout, unitName)
|
||||
log.Printf(" 🖥️ Web Logs: https://%s/_logs", host)
|
||||
log.Printf(" 📋 CLI Logs: kubectl exec -n default deployment/zero-scale-gateway -- journalctl -u %s -n 50", unitName)
|
||||
}
|
||||
91
internal/config/config.go
Normal file
91
internal/config/config.go
Normal file
@@ -0,0 +1,91 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Config holds all configuration options for the zero-scale gateway
|
||||
type Config struct {
|
||||
ListenAddr string
|
||||
AppsDir string
|
||||
ConfigPath string
|
||||
IdleTimeout time.Duration
|
||||
CheckPeriod time.Duration
|
||||
PodName string
|
||||
KubeNamespace string
|
||||
}
|
||||
|
||||
// Load parses command line flags, environment variables, and defaults.
|
||||
// Precedence: CLI Flags > Environment Variables > Defaults
|
||||
func Load() (*Config, error) {
|
||||
defaultAddr := getEnv("LISTEN_ADDR", ":80")
|
||||
if port := os.Getenv("PORT"); port != "" && defaultAddr == ":80" {
|
||||
defaultAddr = ":" + port
|
||||
}
|
||||
defaultAppsDir := getEnv("APPS_DIR", "/opt/apps")
|
||||
defaultConfigPath := getEnv("CONFIG_PATH", "/etc/zero-scale/artifacts.json")
|
||||
defaultIdleTimeout := getEnvDuration("IDLE_TIMEOUT", 15*time.Minute)
|
||||
defaultCheckPeriod := getEnvDuration("CHECK_PERIOD", 30*time.Second)
|
||||
defaultNamespace := getEnv("POD_NAMESPACE", "default")
|
||||
|
||||
podName := os.Getenv("POD_NAME")
|
||||
if podName == "" {
|
||||
if h, err := os.Hostname(); err == nil && h != "" {
|
||||
podName = h
|
||||
} else {
|
||||
podName = "zero-scale-gateway"
|
||||
}
|
||||
}
|
||||
|
||||
addrFlag := flag.String("addr", defaultAddr, "Address to listen on (e.g. :80)")
|
||||
appsDirFlag := flag.String("apps-dir", defaultAppsDir, "Directory containing application binaries")
|
||||
configPathFlag := flag.String("config", defaultConfigPath, "Path to artifacts configuration JSON")
|
||||
idleTimeoutFlag := flag.Duration("idle-timeout", defaultIdleTimeout, "Duration of inactivity before scaling an app to zero")
|
||||
checkPeriodFlag := flag.Duration("check-period", defaultCheckPeriod, "Interval for idle check reaper")
|
||||
namespaceFlag := flag.String("namespace", defaultNamespace, "Kubernetes namespace")
|
||||
|
||||
flag.Parse()
|
||||
|
||||
cfg := &Config{
|
||||
ListenAddr: *addrFlag,
|
||||
AppsDir: *appsDirFlag,
|
||||
ConfigPath: *configPathFlag,
|
||||
IdleTimeout: *idleTimeoutFlag,
|
||||
CheckPeriod: *checkPeriodFlag,
|
||||
PodName: podName,
|
||||
KubeNamespace: *namespaceFlag,
|
||||
}
|
||||
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
// KubeLogURL returns the direct Kubernetes Dashboard web log URL for this pod
|
||||
func (c *Config) KubeLogURL() string {
|
||||
return fmt.Sprintf("https://kube.fairfaxmedia.net/#/log/%s/%s/pod?namespace=%s&container=systemd-runner",
|
||||
c.KubeNamespace, c.PodName, c.KubeNamespace)
|
||||
}
|
||||
|
||||
// KubeShellURL returns the direct Kubernetes Dashboard web shell URL for this pod
|
||||
func (c *Config) KubeShellURL() string {
|
||||
return fmt.Sprintf("https://kube.fairfaxmedia.net/#/shell/%s/%s/systemd-runner?namespace=%s",
|
||||
c.KubeNamespace, c.PodName, c.KubeNamespace)
|
||||
}
|
||||
|
||||
func getEnv(key, fallback string) string {
|
||||
if val := os.Getenv(key); val != "" {
|
||||
return val
|
||||
}
|
||||
return fallback
|
||||
}
|
||||
|
||||
func getEnvDuration(key string, fallback time.Duration) time.Duration {
|
||||
if val := os.Getenv(key); val != "" {
|
||||
if d, err := time.ParseDuration(val); err == nil {
|
||||
return d
|
||||
}
|
||||
}
|
||||
return fallback
|
||||
}
|
||||
347
internal/handlers/logs.go
Normal file
347
internal/handlers/logs.go
Normal file
@@ -0,0 +1,347 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"fmt"
|
||||
"html"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/config"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/runner"
|
||||
)
|
||||
|
||||
// HandleLogsView handles the web log viewer, SSE stream, and raw CLI text endpoints
|
||||
func HandleLogsView(w http.ResponseWriter, r *http.Request, host string, cfg *config.Config, getPortFunc func(string) int) {
|
||||
if appParam := r.URL.Query().Get("app"); appParam != "" {
|
||||
host = appParam
|
||||
}
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
cliCmd := fmt.Sprintf("kubectl exec -n default deployment/zero-scale-gateway -- journalctl -u %s -f", unitName)
|
||||
|
||||
// Raw output for curl users
|
||||
if r.URL.Query().Get("raw") == "1" {
|
||||
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
|
||||
out, err := exec.Command("journalctl", "-u", unitName, "-n", "250", "--no-pager").CombinedOutput()
|
||||
if err != nil {
|
||||
fmt.Fprintf(w, "Error querying journalctl for %s: %v\n", unitName, err)
|
||||
return
|
||||
}
|
||||
w.Write(out)
|
||||
return
|
||||
}
|
||||
|
||||
// SSE streaming endpoint for live logs
|
||||
if r.URL.Query().Get("stream") == "1" || r.URL.Path == "/_logs/stream" {
|
||||
w.Header().Set("Content-Type", "text/event-stream")
|
||||
w.Header().Set("Cache-Control", "no-cache")
|
||||
w.Header().Set("Connection", "keep-alive")
|
||||
w.Header().Set("Access-Control-Allow-Origin", "*")
|
||||
|
||||
flusher, ok := w.(http.Flusher)
|
||||
if !ok {
|
||||
http.Error(w, "Streaming unsupported", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
ctx := r.Context()
|
||||
cmd := exec.CommandContext(ctx, "journalctl", "-u", unitName, "-f", "-n", "50", "--no-pager")
|
||||
stdout, err := cmd.StdoutPipe()
|
||||
if err != nil {
|
||||
fmt.Fprintf(w, "data: [error opening journalctl pipe: %v]\n\n", err)
|
||||
flusher.Flush()
|
||||
return
|
||||
}
|
||||
if err := cmd.Start(); err != nil {
|
||||
fmt.Fprintf(w, "data: [error starting journalctl: %v]\n\n", err)
|
||||
flusher.Flush()
|
||||
return
|
||||
}
|
||||
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
for scanner.Scan() {
|
||||
line := scanner.Text()
|
||||
fmt.Fprintf(w, "data: %s\n\n", line)
|
||||
flusher.Flush()
|
||||
}
|
||||
_ = cmd.Wait()
|
||||
return
|
||||
}
|
||||
|
||||
// HTML Web Log Viewer UI
|
||||
active, _ := runner.IsServiceActive(unitName)
|
||||
statusBadge := `<span class="badge idle">💤 SCALED TO ZERO (IDLE)</span>`
|
||||
if active {
|
||||
port := getPortFunc(host)
|
||||
if port == 0 {
|
||||
port = runner.GetActiveServicePort(unitName)
|
||||
}
|
||||
statusBadge = fmt.Sprintf(`<span class="badge active">🟢 ACTIVE (Port %d)</span>`, port)
|
||||
}
|
||||
|
||||
initialLogs, _ := exec.Command("journalctl", "-u", unitName, "-n", "100", "--no-pager").CombinedOutput()
|
||||
if len(initialLogs) == 0 {
|
||||
initialLogs = []byte("-- No log entries recorded yet for " + unitName + " --\n")
|
||||
}
|
||||
|
||||
var appChips strings.Builder
|
||||
if entries, err := os.ReadDir(cfg.AppsDir); err == nil {
|
||||
for _, e := range entries {
|
||||
if !e.IsDir() && !strings.HasSuffix(e.Name(), ".type") && e.Name() != "lost+found" {
|
||||
currHost := e.Name()
|
||||
activeCls := ""
|
||||
if currHost == host {
|
||||
activeCls = " active"
|
||||
}
|
||||
appChips.WriteString(fmt.Sprintf(`<a href="/_logs?app=%s" class="chip%s">%s</a>`, currHost, activeCls, html.EscapeString(currHost)))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
w.Header().Set("Content-Type", "text/html; charset=utf-8")
|
||||
pageHTML := fmt.Sprintf(`<!DOCTYPE html>
|
||||
<html lang="en">
|
||||
<head>
|
||||
<meta charset="utf-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>%s - Zero-Scale Service Logs</title>
|
||||
<style>
|
||||
:root {
|
||||
--bg: #0d1117;
|
||||
--card: #161b22;
|
||||
--border: #30363d;
|
||||
--text: #c9d1d9;
|
||||
--muted: #8b949e;
|
||||
--green: #2ea043;
|
||||
--amber: #d29922;
|
||||
--blue: #58a6ff;
|
||||
--term-bg: #06090f;
|
||||
}
|
||||
* { box-sizing: border-box; margin: 0; padding: 0; }
|
||||
body {
|
||||
background-color: var(--bg);
|
||||
color: var(--text);
|
||||
font-family: -apple-system, BlinkMacSystemFont, "Segoe UI", Helvetica, Arial, sans-serif;
|
||||
padding: 24px;
|
||||
line-height: 1.5;
|
||||
}
|
||||
.container { max-width: 1200px; margin: 0 auto; }
|
||||
header {
|
||||
display: flex;
|
||||
justify-content: space-between;
|
||||
align-items: center;
|
||||
padding-bottom: 16px;
|
||||
border-bottom: 1px solid var(--border);
|
||||
margin-bottom: 20px;
|
||||
flex-wrap: wrap;
|
||||
gap: 12px;
|
||||
}
|
||||
h1 { font-size: 20px; font-weight: 600; display: flex; align-items: center; gap: 10px; }
|
||||
.badge {
|
||||
font-size: 12px;
|
||||
padding: 4px 10px;
|
||||
border-radius: 20px;
|
||||
font-weight: 600;
|
||||
}
|
||||
.badge.active { background: rgba(46, 160, 67, 0.15); color: #3fb950; border: 1px solid rgba(46, 160, 67, 0.4); }
|
||||
.badge.idle { background: rgba(210, 153, 34, 0.15); color: #e3b341; border: 1px solid rgba(210, 153, 34, 0.4); }
|
||||
.links { display: flex; gap: 10px; flex-wrap: wrap; }
|
||||
.btn {
|
||||
display: inline-flex;
|
||||
align-items: center;
|
||||
gap: 6px;
|
||||
background: var(--card);
|
||||
color: var(--text);
|
||||
text-decoration: none;
|
||||
padding: 6px 12px;
|
||||
border-radius: 6px;
|
||||
border: 1px solid var(--border);
|
||||
font-size: 13px;
|
||||
cursor: pointer;
|
||||
transition: all 0.15s ease;
|
||||
}
|
||||
.btn:hover { background: #21262d; border-color: #8b949e; color: #fff; }
|
||||
.btn-primary { background: #238636; border-color: rgba(240, 246, 252, 0.1); color: #fff; }
|
||||
.btn-primary:hover { background: #2ea043; border-color: rgba(240, 246, 252, 0.2); }
|
||||
.chips { margin-bottom: 16px; display: flex; gap: 8px; flex-wrap: wrap; align-items: center; }
|
||||
.chip-label { font-size: 12px; color: var(--muted); text-transform: uppercase; font-weight: 600; margin-right: 4px; }
|
||||
.chip {
|
||||
padding: 4px 10px;
|
||||
border-radius: 12px;
|
||||
background: var(--card);
|
||||
color: var(--muted);
|
||||
border: 1px solid var(--border);
|
||||
font-size: 12px;
|
||||
text-decoration: none;
|
||||
}
|
||||
.chip:hover { color: var(--text); border-color: var(--muted); }
|
||||
.chip.active { background: #1f6feb22; border-color: var(--blue); color: var(--blue); font-weight: 600; }
|
||||
.cli-box {
|
||||
background: var(--card);
|
||||
border: 1px solid var(--border);
|
||||
border-radius: 6px;
|
||||
padding: 10px 14px;
|
||||
margin-bottom: 16px;
|
||||
display: flex;
|
||||
justify-content: space-between;
|
||||
align-items: center;
|
||||
gap: 12px;
|
||||
font-family: ui-monospace, SFMono-Regular, "SF Mono", Menlo, Consolas, monospace;
|
||||
font-size: 13px;
|
||||
}
|
||||
.cli-code { color: #79c0ff; overflow-x: auto; white-space: nowrap; }
|
||||
.terminal-card {
|
||||
background: var(--term-bg);
|
||||
border: 1px solid var(--border);
|
||||
border-radius: 8px;
|
||||
overflow: hidden;
|
||||
box-shadow: 0 10px 24px rgba(0,0,0,0.5);
|
||||
}
|
||||
.terminal-header {
|
||||
background: #161b22;
|
||||
padding: 8px 14px;
|
||||
border-bottom: 1px solid var(--border);
|
||||
display: flex;
|
||||
justify-content: space-between;
|
||||
align-items: center;
|
||||
font-size: 12px;
|
||||
color: var(--muted);
|
||||
}
|
||||
.term-controls { display: flex; gap: 8px; align-items: center; }
|
||||
.term-dot { width: 10px; height: 10px; border-radius: 50%%; display: inline-block; }
|
||||
.dot-green { background: #3fb950; box-shadow: 0 0 6px #3fb950; }
|
||||
pre#logs {
|
||||
padding: 16px;
|
||||
font-family: ui-monospace, SFMono-Regular, "SF Mono", Menlo, Consolas, monospace;
|
||||
font-size: 12px;
|
||||
line-height: 1.6;
|
||||
color: #e6edf3;
|
||||
overflow-y: auto;
|
||||
max-height: 65vh;
|
||||
white-space: pre-wrap;
|
||||
word-break: break-all;
|
||||
}
|
||||
</style>
|
||||
</head>
|
||||
<body>
|
||||
<div class="container">
|
||||
<header>
|
||||
<div>
|
||||
<h1>⚡ %s %s</h1>
|
||||
<div style="font-size: 12px; color: var(--muted); margin-top: 4px;">Zero-Scale Isolated Journald Logs</div>
|
||||
</div>
|
||||
<div class="links">
|
||||
<a href="https://%s/" target="_blank" class="btn btn-primary">🌐 Open App</a>
|
||||
<a href="%s" target="_blank" class="btn">☸️ Kube Logs</a>
|
||||
<a href="%s" target="_blank" class="btn">☸️ Pod Shell</a>
|
||||
<a href="/_logs?app=%s&raw=1" target="_blank" class="btn">📄 Raw Text</a>
|
||||
</div>
|
||||
</header>
|
||||
|
||||
<div class="chips">
|
||||
<span class="chip-label">Services:</span>
|
||||
%s
|
||||
</div>
|
||||
|
||||
<div class="cli-box">
|
||||
<div class="cli-code"><span style="color:var(--muted)">$</span> %s</div>
|
||||
<button class="btn" id="copy-btn">Copy CLI</button>
|
||||
</div>
|
||||
|
||||
<div class="terminal-card">
|
||||
<div class="terminal-header">
|
||||
<div style="display:flex; align-items:center; gap:8px;">
|
||||
<span class="term-dot dot-green" id="stream-dot"></span>
|
||||
<span id="stream-status">Live Stream Active</span>
|
||||
</div>
|
||||
<div class="term-controls">
|
||||
<button class="btn" id="stream-toggle">Pause Stream</button>
|
||||
<button class="btn" id="clear-btn">Clear</button>
|
||||
</div>
|
||||
</div>
|
||||
<pre id="logs">%s</pre>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<script>
|
||||
const host = "%s";
|
||||
const cliCmd = "%s";
|
||||
const logElem = document.getElementById("logs");
|
||||
const dotElem = document.getElementById("stream-dot");
|
||||
const statusElem = document.getElementById("stream-status");
|
||||
const toggleBtn = document.getElementById("stream-toggle");
|
||||
const clearBtn = document.getElementById("clear-btn");
|
||||
const copyBtn = document.getElementById("copy-btn");
|
||||
|
||||
logElem.scrollTop = logElem.scrollHeight;
|
||||
|
||||
copyBtn.addEventListener("click", function() {
|
||||
navigator.clipboard.writeText(cliCmd);
|
||||
copyBtn.innerText = "Copied!";
|
||||
setTimeout(() => copyBtn.innerText = "Copy CLI", 2000);
|
||||
});
|
||||
|
||||
clearBtn.addEventListener("click", function() {
|
||||
logElem.textContent = "";
|
||||
});
|
||||
|
||||
let evtSource = null;
|
||||
let isStreaming = false;
|
||||
|
||||
function startStream() {
|
||||
evtSource = new EventSource("/_logs?app=" + encodeURIComponent(host) + "&stream=1");
|
||||
evtSource.onmessage = function(e) {
|
||||
logElem.textContent += e.data + "\n";
|
||||
logElem.scrollTop = logElem.scrollHeight;
|
||||
};
|
||||
evtSource.onerror = function() {
|
||||
statusElem.innerText = "Reconnecting...";
|
||||
dotElem.style.background = "#d29922";
|
||||
};
|
||||
evtSource.onopen = function() {
|
||||
statusElem.innerText = "Live Stream Active";
|
||||
dotElem.style.background = "#3fb950";
|
||||
};
|
||||
isStreaming = true;
|
||||
toggleBtn.innerText = "Pause Stream";
|
||||
}
|
||||
|
||||
function stopStream() {
|
||||
if (evtSource) {
|
||||
evtSource.close();
|
||||
evtSource = null;
|
||||
}
|
||||
isStreaming = false;
|
||||
statusElem.innerText = "Stream Paused";
|
||||
dotElem.style.background = "#8b949e";
|
||||
toggleBtn.innerText = "Resume Stream";
|
||||
}
|
||||
|
||||
toggleBtn.addEventListener("click", function() {
|
||||
if (isStreaming) stopStream();
|
||||
else startStream();
|
||||
});
|
||||
|
||||
startStream();
|
||||
</script>
|
||||
</body>
|
||||
</html>`,
|
||||
html.EscapeString(host),
|
||||
html.EscapeString(host),
|
||||
statusBadge,
|
||||
html.EscapeString(host),
|
||||
cfg.KubeLogURL(),
|
||||
cfg.KubeShellURL(),
|
||||
url.QueryEscape(host),
|
||||
appChips.String(),
|
||||
html.EscapeString(cliCmd),
|
||||
html.EscapeString(string(initialLogs)),
|
||||
html.EscapeString(host),
|
||||
cliCmd,
|
||||
)
|
||||
|
||||
w.Write([]byte(pageHTML))
|
||||
}
|
||||
181
internal/proxy/proxy.go
Normal file
181
internal/proxy/proxy.go
Normal file
@@ -0,0 +1,181 @@
|
||||
package proxy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
"net/http"
|
||||
"net/http/httputil"
|
||||
"net/url"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/banner"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/config"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/handlers"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/runner"
|
||||
)
|
||||
|
||||
// Gateway manages dynamic on-demand reverse proxying and zero-scale scaling
|
||||
type Gateway struct {
|
||||
cfg *config.Config
|
||||
allocatedPorts map[string]int
|
||||
lastActive map[string]time.Time
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
// New creates a new Gateway instance
|
||||
func New(cfg *config.Config) *Gateway {
|
||||
return &Gateway{
|
||||
cfg: cfg,
|
||||
allocatedPorts: make(map[string]int),
|
||||
lastActive: make(map[string]time.Time),
|
||||
}
|
||||
}
|
||||
|
||||
// GetPort returns the currently allocated port for host, or 0 if not allocated
|
||||
func (g *Gateway) GetPort(host string) int {
|
||||
g.mu.Lock()
|
||||
defer g.mu.Unlock()
|
||||
return g.allocatedPorts[host]
|
||||
}
|
||||
|
||||
// GetLastActiveMap returns a snapshot copy of the last active timestamps
|
||||
func (g *Gateway) GetLastActiveMap() map[string]time.Time {
|
||||
g.mu.Lock()
|
||||
defer g.mu.Unlock()
|
||||
cp := make(map[string]time.Time, len(g.lastActive))
|
||||
for k, v := range g.lastActive {
|
||||
cp[k] = v
|
||||
}
|
||||
return cp
|
||||
}
|
||||
|
||||
// StopService deallocates port and stops the dynamic unit
|
||||
func (g *Gateway) StopService(host string) error {
|
||||
g.mu.Lock()
|
||||
defer g.mu.Unlock()
|
||||
delete(g.allocatedPorts, host)
|
||||
delete(g.lastActive, host)
|
||||
return runner.StopService(host)
|
||||
}
|
||||
|
||||
// ServeHTTP handles incoming HTTP requests, intercepts logs, and reverse proxies
|
||||
func (g *Gateway) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
host := r.Host
|
||||
if strings.Contains(host, ":") {
|
||||
host = strings.Split(host, ":")[0]
|
||||
}
|
||||
|
||||
// Filter out noisy static asset logging
|
||||
isStatic := strings.HasPrefix(r.URL.Path, "/static/") ||
|
||||
strings.HasPrefix(r.URL.Path, "/_logs") ||
|
||||
r.URL.Path == "/favicon.ico" ||
|
||||
strings.HasSuffix(r.URL.Path, ".css") ||
|
||||
strings.HasSuffix(r.URL.Path, ".js") ||
|
||||
strings.HasSuffix(r.URL.Path, ".png") ||
|
||||
strings.HasSuffix(r.URL.Path, ".jpg") ||
|
||||
strings.HasSuffix(r.URL.Path, ".ico")
|
||||
if !isStatic {
|
||||
log.Printf("📥 [%s] %s %s", host, r.Method, r.URL.Path)
|
||||
}
|
||||
|
||||
// Intercept log viewer requests
|
||||
if r.URL.Path == "/_logs" || strings.HasPrefix(r.URL.Path, "/_logs/") {
|
||||
handlers.HandleLogsView(w, r, host, g.cfg, g.GetPort)
|
||||
return
|
||||
}
|
||||
|
||||
binaryPath := filepath.Join(g.cfg.AppsDir, host)
|
||||
if !runner.FileExists(binaryPath) {
|
||||
// Fallback lookup: match shortName prefix (e.g. bdash2 -> bdash2.kube.fairfaxmedia.net)
|
||||
shortName := strings.Split(host, ".")[0]
|
||||
entries, err := os.ReadDir(g.cfg.AppsDir)
|
||||
if err == nil {
|
||||
for _, entry := range entries {
|
||||
if !entry.IsDir() && (strings.HasPrefix(entry.Name(), shortName+".") || entry.Name() == shortName) {
|
||||
binaryPath = filepath.Join(g.cfg.AppsDir, entry.Name())
|
||||
log.Printf("Resolved host %s to binary %s", host, binaryPath)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !runner.FileExists(binaryPath) {
|
||||
http.Error(w, fmt.Sprintf("Application %s not found", host), http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
|
||||
g.mu.Lock()
|
||||
g.lastActive[host] = time.Now()
|
||||
|
||||
active, err := runner.IsServiceActive(unitName)
|
||||
if err != nil {
|
||||
g.mu.Unlock()
|
||||
log.Printf("Error checking status of unit %s: %v", unitName, err)
|
||||
http.Error(w, "Failed to check service status", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
port, exists := g.allocatedPorts[host]
|
||||
if active && !exists {
|
||||
if activePort := runner.GetActiveServicePort(unitName); activePort > 0 {
|
||||
g.allocatedPorts[host] = activePort
|
||||
port = activePort
|
||||
exists = true
|
||||
log.Printf("Discovered active service %s running on port %d", host, port)
|
||||
}
|
||||
}
|
||||
|
||||
if !active || !exists {
|
||||
// Stop any stale unit without holding lock
|
||||
if active {
|
||||
delete(g.allocatedPorts, host)
|
||||
delete(g.lastActive, host)
|
||||
_ = runner.StopService(host)
|
||||
}
|
||||
freePort, err := runner.GetFreePort()
|
||||
if err != nil {
|
||||
g.mu.Unlock()
|
||||
log.Printf("Failed to get free port for %s: %v", host, err)
|
||||
http.Error(w, "Internal Server Error: No available ports", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
g.allocatedPorts[host] = freePort
|
||||
port = freePort
|
||||
|
||||
log.Printf("Host %s matches binary %s. Starting sandboxed systemd unit on port %d...", host, binaryPath, port)
|
||||
if err := runner.StartTransientService(host, port, binaryPath, g.cfg); err != nil {
|
||||
delete(g.allocatedPorts, host)
|
||||
g.mu.Unlock()
|
||||
log.Printf("Failed to start service %s: %v", host, err)
|
||||
http.Error(w, fmt.Sprintf("Failed to start service %s", host), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
g.mu.Unlock()
|
||||
|
||||
targetAddr := fmt.Sprintf("127.0.0.1:%d", port)
|
||||
if err := runner.WaitPortReady(targetAddr, 30*time.Second); err != nil {
|
||||
log.Printf("Service %s on port %d did not bind in time: %v", host, port, err)
|
||||
_ = g.StopService(host)
|
||||
http.Error(w, "Application startup timeout", http.StatusGatewayTimeout)
|
||||
return
|
||||
}
|
||||
banner.PrintScaleUp(host, port)
|
||||
} else {
|
||||
g.mu.Unlock()
|
||||
}
|
||||
|
||||
targetURL, _ := url.Parse(fmt.Sprintf("http://127.0.0.1:%d", port))
|
||||
proxy := httputil.NewSingleHostReverseProxy(targetURL)
|
||||
proxy.ErrorHandler = func(w http.ResponseWriter, r *http.Request, err error) {
|
||||
log.Printf("Proxy error for %s on port %d: %v. Cleaning up service...", host, port, err)
|
||||
_ = g.StopService(host)
|
||||
http.Error(w, "Bad Gateway", http.StatusBadGateway)
|
||||
}
|
||||
proxy.ServeHTTP(w, r)
|
||||
}
|
||||
44
internal/reaper/reaper.go
Normal file
44
internal/reaper/reaper.go
Normal file
@@ -0,0 +1,44 @@
|
||||
package reaper
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/banner"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/config"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/runner"
|
||||
)
|
||||
|
||||
type Reaper struct {
|
||||
cfg *config.Config
|
||||
stopFunc func(string) error
|
||||
getActive func() map[string]time.Time
|
||||
}
|
||||
|
||||
func New(cfg *config.Config, getActive func() map[string]time.Time, stopFunc func(string) error) *Reaper {
|
||||
return &Reaper{
|
||||
cfg: cfg,
|
||||
getActive: getActive,
|
||||
stopFunc: stopFunc,
|
||||
}
|
||||
}
|
||||
|
||||
// Start begins the background reaper ticker loop
|
||||
func (r *Reaper) Start() {
|
||||
ticker := time.NewTicker(r.cfg.CheckPeriod)
|
||||
for range ticker.C {
|
||||
now := time.Now()
|
||||
activeMap := r.getActive()
|
||||
for host, lastTime := range activeMap {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
active, err := runner.IsServiceActive(unitName)
|
||||
if err == nil && active && now.Sub(lastTime) > r.cfg.IdleTimeout {
|
||||
banner.PrintScaleToZero(host, r.cfg.IdleTimeout)
|
||||
if err := r.stopFunc(host); err != nil {
|
||||
log.Printf("Failed to stop idle unit %s: %v", unitName, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
76
internal/runner/k8s.go
Normal file
76
internal/runner/k8s.go
Normal file
@@ -0,0 +1,76 @@
|
||||
package runner
|
||||
|
||||
import (
|
||||
"crypto/tls"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// FetchK8sResource dynamically retrieves a ConfigMap or Secret from the in-cluster Kubernetes API
|
||||
func FetchK8sResource(resourceType, name, namespace string) (map[string]string, error) {
|
||||
tokenBytes, err := os.ReadFile("/var/run/secrets/kubernetes.io/serviceaccount/token")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
token := strings.TrimSpace(string(tokenBytes))
|
||||
|
||||
k8sHost := os.Getenv("KUBERNETES_SERVICE_HOST")
|
||||
if k8sHost == "" {
|
||||
k8sHost = "kubernetes.default.svc"
|
||||
}
|
||||
k8sPort := os.Getenv("KUBERNETES_SERVICE_PORT")
|
||||
if k8sPort == "" {
|
||||
k8sPort = "443"
|
||||
}
|
||||
if namespace == "" {
|
||||
namespace = "default"
|
||||
}
|
||||
|
||||
reqURL := fmt.Sprintf("https://%s:%s/api/v1/namespaces/%s/%s/%s", k8sHost, k8sPort, namespace, resourceType, name)
|
||||
req, err := http.NewRequest("GET", reqURL, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
req.Header.Set("Authorization", "Bearer "+token)
|
||||
|
||||
tr := &http.Transport{
|
||||
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
|
||||
}
|
||||
client := &http.Client{Transport: tr, Timeout: 5 * time.Second}
|
||||
resp, err := client.Do(req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, fmt.Errorf("k8s API returned status %d for %s/%s", resp.StatusCode, resourceType, name)
|
||||
}
|
||||
|
||||
var res struct {
|
||||
Data map[string]string `json:"data"`
|
||||
}
|
||||
if err := json.NewDecoder(resp.Body).Decode(&res); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if resourceType == "secrets" {
|
||||
decoded := make(map[string]string)
|
||||
for k, v := range res.Data {
|
||||
b, err := base64.StdEncoding.DecodeString(v)
|
||||
if err == nil {
|
||||
decoded[k] = string(b)
|
||||
} else {
|
||||
decoded[k] = v
|
||||
}
|
||||
}
|
||||
return decoded, nil
|
||||
}
|
||||
|
||||
return res.Data, nil
|
||||
}
|
||||
139
internal/runner/systemd.go
Normal file
139
internal/runner/systemd.go
Normal file
@@ -0,0 +1,139 @@
|
||||
package runner
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/config"
|
||||
)
|
||||
|
||||
// IsServiceActive checks if the dynamic systemd unit is running
|
||||
func IsServiceActive(unit string) (bool, error) {
|
||||
cmd := exec.Command("systemctl", "is-active", unit)
|
||||
err := cmd.Run()
|
||||
if err != nil {
|
||||
if _, ok := err.(*exec.ExitError); ok {
|
||||
return false, nil
|
||||
}
|
||||
return false, err
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// GetActiveServicePort reads PORT from running systemd service environment
|
||||
func GetActiveServicePort(unit string) int {
|
||||
out, err := exec.Command("systemctl", "show", "-p", "Environment", unit).Output()
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
str := string(out)
|
||||
for _, env := range strings.Fields(str) {
|
||||
if strings.HasPrefix(env, "PORT=") || strings.HasPrefix(env, "Environment=PORT=") {
|
||||
parts := strings.Split(env, "=")
|
||||
if len(parts) >= 2 {
|
||||
val := parts[len(parts)-1]
|
||||
if p, err := strconv.Atoi(val); err == nil && p > 0 {
|
||||
return p
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
// StartTransientService executes systemd-run with dynamic environment injection
|
||||
func StartTransientService(host string, port int, binaryPath string, cfg *config.Config) error {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
typePath := binaryPath + ".type"
|
||||
|
||||
runType := "elf"
|
||||
if FileExists(typePath) {
|
||||
tBytes, err := os.ReadFile(typePath)
|
||||
if err == nil {
|
||||
runType = strings.TrimSpace(string(tBytes))
|
||||
}
|
||||
}
|
||||
|
||||
baseArgs := []string{
|
||||
fmt.Sprintf("--unit=%s", unitName),
|
||||
"-p", "DynamicUser=yes",
|
||||
"-p", "PrivateTmp=yes",
|
||||
"-p", "ProtectSystem=strict",
|
||||
fmt.Sprintf("--setenv=PORT=%d", port),
|
||||
}
|
||||
|
||||
// Mount persistent volume if present
|
||||
if _, err := os.Stat("/app/uploads"); err == nil {
|
||||
baseArgs = append(baseArgs, "-p", "BindPaths=/app/uploads:/app/uploads")
|
||||
}
|
||||
|
||||
// Forward standard database environment variables if present on host
|
||||
envKeys := []string{
|
||||
"DATABASE_URL", "DATABASE_HOST", "DATABASE_PORT", "DATABASE_USER", "DATABASE_PASSWORD", "DATABASE_NAME",
|
||||
"PGHOST", "PGPORT", "PGUSER", "PGPASSWORD", "PGDATABASE",
|
||||
}
|
||||
for _, k := range envKeys {
|
||||
if v := os.Getenv(k); v != "" {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch dynamic app ConfigMap and Secret from Kubernetes API
|
||||
shortName := strings.Split(host, ".")[0]
|
||||
appCMName := "zero-scale-app-" + shortName
|
||||
if appData, err := FetchK8sResource("configmaps", appCMName, cfg.KubeNamespace); err == nil {
|
||||
for k, v := range appData {
|
||||
if k != "host" && k != "url" && k != "type" && k != "mounts" && k != "scale-to-zero" && k != "database_configmap" && k != "database_secret" {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
if dbCM := appData["database_configmap"]; dbCM != "" {
|
||||
if cmData, err := FetchK8sResource("configmaps", dbCM, cfg.KubeNamespace); err == nil {
|
||||
for k, v := range cmData {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
}
|
||||
if dbSec := appData["database_secret"]; dbSec != "" {
|
||||
if secData, err := FetchK8sResource("secrets", dbSec, cfg.KubeNamespace); err == nil {
|
||||
for k, v := range secData {
|
||||
baseArgs = append(baseArgs, fmt.Sprintf("--setenv=%s=%s", k, v))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
var cmd *exec.Cmd
|
||||
if runType == "wasm" {
|
||||
args := append(baseArgs,
|
||||
"--setenv=WASM_PORT="+strconv.Itoa(port),
|
||||
"/usr/local/bin/wasmtime", "run",
|
||||
"--env", fmt.Sprintf("PORT=%d", port),
|
||||
"--tcplisten", fmt.Sprintf("127.0.0.1:%d", port),
|
||||
binaryPath,
|
||||
)
|
||||
log.Printf("Executing WASM runner: /usr/bin/systemd-run %s", strings.Join(args, " "))
|
||||
cmd = exec.Command("systemd-run", args...)
|
||||
} else {
|
||||
args := append(baseArgs, binaryPath, "-port", strconv.Itoa(port))
|
||||
log.Printf("Executing: /usr/bin/systemd-run %s", strings.Join(args, " "))
|
||||
cmd = exec.Command("systemd-run", args...)
|
||||
}
|
||||
|
||||
out, err := cmd.CombinedOutput()
|
||||
if err != nil {
|
||||
return fmt.Errorf("systemd-run failed: %v, output: %s", err, string(out))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// StopService stops the dynamic systemd unit
|
||||
func StopService(host string) error {
|
||||
unitName := fmt.Sprintf("app-%s", host)
|
||||
cmd := exec.Command("systemctl", "stop", unitName)
|
||||
return cmd.Run()
|
||||
}
|
||||
45
internal/runner/utils.go
Normal file
45
internal/runner/utils.go
Normal file
@@ -0,0 +1,45 @@
|
||||
package runner
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"time"
|
||||
)
|
||||
|
||||
// GetFreePort queries the kernel for a free ephemeral port
|
||||
func GetFreePort() (int, error) {
|
||||
addr, err := net.ResolveTCPAddr("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
l, err := net.ListenTCP("tcp", addr)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
defer l.Close()
|
||||
return l.Addr().(*net.TCPAddr).Port, nil
|
||||
}
|
||||
|
||||
// WaitPortReady polls targetAddr until it accepts TCP connections or timeout occurs
|
||||
func WaitPortReady(addr string, timeout time.Duration) error {
|
||||
deadline := time.Now().Add(timeout)
|
||||
for time.Now().Before(deadline) {
|
||||
conn, err := net.DialTimeout("tcp", addr, 100*time.Millisecond)
|
||||
if err == nil {
|
||||
conn.Close()
|
||||
return nil
|
||||
}
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
}
|
||||
return fmt.Errorf("timeout waiting for %s to bind", addr)
|
||||
}
|
||||
|
||||
// FileExists checks whether path exists and is not a directory
|
||||
func FileExists(path string) bool {
|
||||
info, err := os.Stat(path)
|
||||
if os.IsNotExist(err) {
|
||||
return false
|
||||
}
|
||||
return err == nil && !info.IsDir()
|
||||
}
|
||||
83
internal/sync/artifacts.go
Normal file
83
internal/sync/artifacts.go
Normal file
@@ -0,0 +1,83 @@
|
||||
package sync
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"time"
|
||||
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/config"
|
||||
"git.fairfaxmedia.net/open-cream-cheese/zero-scale-platform/internal/runner"
|
||||
)
|
||||
|
||||
type ArtifactItem struct {
|
||||
URL string `json:"url"`
|
||||
Type string `json:"type"`
|
||||
}
|
||||
|
||||
type ConfigMapData struct {
|
||||
Artifacts map[string]ArtifactItem `json:"artifacts"`
|
||||
}
|
||||
|
||||
// StartArtifactSync periodically checks the artifacts config and downloads new binaries
|
||||
func StartArtifactSync(cfg *config.Config) {
|
||||
ticker := time.NewTicker(15 * time.Second)
|
||||
for range ticker.C {
|
||||
if !runner.FileExists(cfg.ConfigPath) {
|
||||
continue
|
||||
}
|
||||
data, err := os.ReadFile(cfg.ConfigPath)
|
||||
if err != nil {
|
||||
log.Printf("Failed to read artifacts config: %v", err)
|
||||
continue
|
||||
}
|
||||
var cmCfg ConfigMapData
|
||||
if err := json.Unmarshal(data, &cmCfg); err != nil {
|
||||
log.Printf("Failed to parse artifacts config JSON: %v", err)
|
||||
continue
|
||||
}
|
||||
|
||||
for host, item := range cmCfg.Artifacts {
|
||||
targetPath := filepath.Join(cfg.AppsDir, host)
|
||||
typePath := targetPath + ".type"
|
||||
|
||||
_ = os.WriteFile(typePath, []byte(item.Type), 0644)
|
||||
|
||||
if !runner.FileExists(targetPath) {
|
||||
log.Printf("Downloading artifact for %s from %s...", host, item.URL)
|
||||
err := downloadFile(item.URL, targetPath)
|
||||
if err != nil {
|
||||
log.Printf("Failed to download %s: %v", item.URL, err)
|
||||
continue
|
||||
}
|
||||
os.Chmod(targetPath, 0755)
|
||||
log.Printf("Successfully deployed %s to %s", host, targetPath)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func downloadFile(url, targetPath string) error {
|
||||
resp, err := http.Get(url)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("bad status: %s", resp.Status)
|
||||
}
|
||||
|
||||
out, err := os.Create(targetPath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer out.Close()
|
||||
|
||||
_, err = io.Copy(out, resp.Body)
|
||||
return err
|
||||
}
|
||||
@@ -6,6 +6,8 @@ After=network.target
|
||||
Type=simple
|
||||
ExecStart=/usr/local/bin/gateway
|
||||
Restart=always
|
||||
StandardOutput=journal+console
|
||||
StandardError=journal+console
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
|
||||
Reference in New Issue
Block a user