fix pipe closed issue, and a lot of other stuff
This commit is contained in:
132
server/server.go
132
server/server.go
@@ -8,10 +8,10 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
|
||||
_ "embed"
|
||||
|
||||
@@ -19,6 +19,8 @@ import (
|
||||
"github.com/docker/docker/client"
|
||||
"github.com/juls0730/flux/pkg"
|
||||
_ "github.com/mattn/go-sqlite3"
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -31,7 +33,8 @@ var (
|
||||
Level: 0,
|
||||
},
|
||||
}
|
||||
Flux *FluxServer
|
||||
Flux *FluxServer
|
||||
logger *zap.SugaredLogger
|
||||
)
|
||||
|
||||
type FluxServerConfig struct {
|
||||
@@ -46,92 +49,125 @@ type FluxServer struct {
|
||||
rootDir string
|
||||
appManager *AppManager
|
||||
dockerClient *client.Client
|
||||
Logger *zap.SugaredLogger
|
||||
}
|
||||
|
||||
func NewFluxServer() *FluxServer {
|
||||
dockerClient, err := client.NewClientWithOpts(client.FromEnv)
|
||||
if err != nil {
|
||||
logger.Fatalw("Failed to create docker client", zap.Error(err))
|
||||
}
|
||||
|
||||
rootDir := os.Getenv("FLUXD_ROOT_DIR")
|
||||
if rootDir == "" {
|
||||
rootDir = "/var/fluxd"
|
||||
}
|
||||
|
||||
db, err := sql.Open("sqlite3", filepath.Join(rootDir, "fluxd.db"))
|
||||
if err != nil {
|
||||
logger.Fatalw("Failed to open database", zap.Error(err))
|
||||
}
|
||||
|
||||
_, err = db.Exec(string(schemaBytes))
|
||||
if err != nil {
|
||||
logger.Fatalw("Failed to create database schema", zap.Error(err))
|
||||
}
|
||||
|
||||
return &FluxServer{
|
||||
db: db,
|
||||
proxy: &Proxy{},
|
||||
appManager: new(AppManager),
|
||||
rootDir: rootDir,
|
||||
dockerClient: dockerClient,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *FluxServer) Stop() {
|
||||
s.Logger.Sync()
|
||||
}
|
||||
|
||||
func NewServer() *FluxServer {
|
||||
Flux = new(FluxServer)
|
||||
verbosity, err := strconv.Atoi(os.Getenv("FLUXD_VERBOSITY"))
|
||||
if err != nil {
|
||||
verbosity = 0
|
||||
}
|
||||
|
||||
config := zap.NewProductionConfig()
|
||||
|
||||
if os.Getenv("DEBUG") == "true" {
|
||||
config = zap.NewDevelopmentConfig()
|
||||
verbosity = -1
|
||||
}
|
||||
|
||||
config.Level = zap.NewAtomicLevelAt(zapcore.Level(verbosity))
|
||||
|
||||
lameLogger, err := config.Build()
|
||||
logger = lameLogger.Sugar()
|
||||
|
||||
if err != nil {
|
||||
logger.Fatalw("Failed to create logger", zap.Error(err))
|
||||
}
|
||||
|
||||
Flux = NewFluxServer()
|
||||
Flux.Logger = logger
|
||||
|
||||
var serverConfig FluxServerConfig
|
||||
|
||||
Flux.rootDir = os.Getenv("FLUXD_ROOT_DIR")
|
||||
if Flux.rootDir == "" {
|
||||
Flux.rootDir = "/var/fluxd"
|
||||
}
|
||||
|
||||
// parse config, if it doesnt exist, create it and use the default config
|
||||
configPath := filepath.Join(Flux.rootDir, "config.json")
|
||||
if _, err := os.Stat(configPath); err != nil {
|
||||
if err := os.MkdirAll(Flux.rootDir, 0755); err != nil {
|
||||
log.Fatalf("Failed to create fluxd directory: %v\n", err)
|
||||
logger.Fatalw("Failed to create fluxd directory", zap.Error(err))
|
||||
}
|
||||
|
||||
configBytes, err := json.Marshal(DefaultConfig)
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to marshal default config: %v\n", err)
|
||||
logger.Fatalw("Failed to marshal default config", zap.Error(err))
|
||||
}
|
||||
|
||||
log.Printf("Config file not found, creating default config file at %s\n", configPath)
|
||||
logger.Debugw("Config file not found creating default config file at", zap.String("path", configPath))
|
||||
if err := os.WriteFile(configPath, configBytes, 0644); err != nil {
|
||||
log.Fatalf("Failed to write config file: %v\n", err)
|
||||
logger.Fatalw("Failed to write config file", zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
configFile, err := os.ReadFile(configPath)
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to read config file: %v\n", err)
|
||||
logger.Fatalw("Failed to read config file", zap.Error(err))
|
||||
}
|
||||
|
||||
if err := json.Unmarshal(configFile, &serverConfig); err != nil {
|
||||
log.Fatalf("Failed to parse config file: %v\n", err)
|
||||
logger.Fatalw("Failed to parse config file", zap.Error(err))
|
||||
}
|
||||
|
||||
Flux.config = serverConfig
|
||||
|
||||
Flux.dockerClient, err = client.NewClientWithOpts(client.FromEnv)
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to create docker client: %v\n", err)
|
||||
}
|
||||
|
||||
log.Printf("Pulling builder image %s, this may take a while...\n", serverConfig.Builder)
|
||||
|
||||
logger.Infof("Pulling builder image %s this may take a while...", serverConfig.Builder)
|
||||
events, err := Flux.dockerClient.ImagePull(context.Background(), fmt.Sprintf("%s:latest", serverConfig.Builder), image.PullOptions{})
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to pull builder image: %v\n", err)
|
||||
logger.Fatalw("Failed to pull builder image", zap.Error(err))
|
||||
}
|
||||
|
||||
// wait for the iamge to be pulled
|
||||
// blocking wait for the iamge to be pulled
|
||||
io.Copy(io.Discard, events)
|
||||
|
||||
log.Printf("Successfully pulled builder image %s\n", serverConfig.Builder)
|
||||
logger.Infow("Successfully pulled builder image", zap.String("image", serverConfig.Builder))
|
||||
|
||||
if err := os.MkdirAll(filepath.Join(Flux.rootDir, "apps"), 0755); err != nil {
|
||||
log.Fatalf("Failed to create apps directory: %v\n", err)
|
||||
logger.Fatalw("Failed to create apps directory", zap.Error(err))
|
||||
}
|
||||
|
||||
Flux.db, err = sql.Open("sqlite3", filepath.Join(Flux.rootDir, "fluxd.db"))
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to open database: %v\n", err)
|
||||
}
|
||||
|
||||
_, err = Flux.db.Exec(string(schemaBytes))
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to create database schema: %v\n", err)
|
||||
}
|
||||
|
||||
Flux.appManager = new(AppManager)
|
||||
Flux.appManager.Init()
|
||||
|
||||
Flux.proxy = &Proxy{}
|
||||
|
||||
port := os.Getenv("FLUXD_PROXY_PORT")
|
||||
if port == "" {
|
||||
port = "7465"
|
||||
}
|
||||
|
||||
go func() {
|
||||
log.Printf("Proxy server starting on http://127.0.0.1:%s\n", port)
|
||||
logger.Infof("Proxy server starting on http://127.0.0.1:%s", port)
|
||||
if err := http.ListenAndServe(fmt.Sprintf(":%s", port), Flux.proxy); err != nil && err != http.ErrServerClosed {
|
||||
log.Fatalf("Proxy server error: %v", err)
|
||||
logger.Fatalw("Proxy server error", zap.Error(err))
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -142,7 +178,7 @@ func (s *FluxServer) UploadAppCode(code io.Reader, projectConfig pkg.ProjectConf
|
||||
var err error
|
||||
projectPath := filepath.Join(s.rootDir, "apps", projectConfig.Name)
|
||||
if err = os.MkdirAll(projectPath, 0755); err != nil {
|
||||
log.Printf("Failed to create project directory: %v\n", err)
|
||||
logger.Errorw("Failed to create project directory", zap.Error(err))
|
||||
return "", err
|
||||
}
|
||||
|
||||
@@ -156,7 +192,7 @@ func (s *FluxServer) UploadAppCode(code io.Reader, projectConfig pkg.ProjectConf
|
||||
if s.config.Compression.Enabled {
|
||||
gzReader, err = gzip.NewReader(code)
|
||||
if err != nil {
|
||||
log.Printf("Failed to create gzip reader: %v\n", err)
|
||||
logger.Infow("Failed to create gzip reader", zap.Error(err))
|
||||
return "", err
|
||||
}
|
||||
}
|
||||
@@ -168,14 +204,14 @@ func (s *FluxServer) UploadAppCode(code io.Reader, projectConfig pkg.ProjectConf
|
||||
tarReader = tar.NewReader(code)
|
||||
}
|
||||
|
||||
log.Printf("Extracting files for %s...\n", projectPath)
|
||||
logger.Infow("Extracting files for project", zap.String("project", projectPath))
|
||||
for {
|
||||
header, err := tarReader.Next()
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
log.Printf("Failed to read tar header: %v\n", err)
|
||||
logger.Debugw("Failed to read tar header", zap.Error(err))
|
||||
return "", err
|
||||
}
|
||||
|
||||
@@ -186,24 +222,24 @@ func (s *FluxServer) UploadAppCode(code io.Reader, projectConfig pkg.ProjectConf
|
||||
switch header.Typeflag {
|
||||
case tar.TypeDir:
|
||||
if err = os.MkdirAll(path, 0755); err != nil {
|
||||
log.Printf("Failed to extract directory: %v\n", err)
|
||||
logger.Debugw("Failed to extract directory", zap.Error(err))
|
||||
return "", err
|
||||
}
|
||||
case tar.TypeReg:
|
||||
if err = os.MkdirAll(filepath.Dir(path), 0755); err != nil {
|
||||
log.Printf("Failed to extract directory: %v\n", err)
|
||||
logger.Debugw("Failed to extract directory", zap.Error(err))
|
||||
return "", err
|
||||
}
|
||||
|
||||
outFile, err := os.Create(path)
|
||||
if err != nil {
|
||||
log.Printf("Failed to extract file: %v\n", err)
|
||||
logger.Debugw("Failed to extract file", zap.Error(err))
|
||||
return "", err
|
||||
}
|
||||
defer outFile.Close()
|
||||
|
||||
if _, err = io.Copy(outFile, tarReader); err != nil {
|
||||
log.Printf("Failed to copy file during extraction: %v\n", err)
|
||||
logger.Debugw("Failed to copy file during extraction", zap.Error(err))
|
||||
return "", err
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user