Separate database svc and minio svc (#4447)

* separate database svc and minio svc

* delete OBJECT_STORAGE_INSTANCE env

* fix deploy.yaml
This commit is contained in:
xuziyi
2024-01-02 14:53:11 +08:00
committed by GitHub
parent d73becbf40
commit a35524dfd3
14 changed files with 465 additions and 228 deletions
-2
View File
@@ -1,2 +0,0 @@
server:
addr: ":9090"
@@ -37,7 +37,6 @@ spec:
env:
- name: PROMETHEUS_SERVICE_HOST
value: http://kb-addon-prometheus-server.kb-system.svc.cluster.local
- name: OBJECT_STORAGE_INSTANCE
image: ghcr.io/labring/sealos-database-service:latest
imagePullPolicy: Always
name: database-monitor
@@ -45,9 +44,12 @@ spec:
- containerPort: 9090
protocol: TCP
resources:
limits:
cpu: 500m
memory: 1024Mi
requests:
cpu: 1m
memory: 500M
cpu: 5m
memory: 64Mi
securityContext:
allowPrivilegeEscalation: false
capabilities:
+8 -39
View File
@@ -4,53 +4,20 @@ import (
"flag"
"fmt"
"log"
"net"
"net/http"
"os"
"github.com/labring/sealos/service/database/server"
dbserver "github.com/labring/sealos/service/database/server"
"github.com/labring/sealos/service/pkg/server"
)
type RestartableServer struct {
configFile string
// server *server.PromServer
// hs *http.Server
}
func (rs *RestartableServer) Serve(c *server.Config) {
var ps, err = server.NewPromServer(c)
if err != nil {
fmt.Printf("Failed to create auth server: %s\n", err)
return
}
hs := &http.Server{
Addr: c.Server.ListenAddress,
Handler: ps,
}
var listener net.Listener
listener, err = net.Listen("tcp", c.Server.ListenAddress)
if err != nil {
fmt.Println(err)
return
}
fmt.Printf("Serve on %s\n", c.Server.ListenAddress)
if err := hs.Serve(listener); err != nil {
fmt.Println(err)
return
}
}
func main() {
log.SetOutput(os.Stdout) // 将日志输出定向到标准输出(stdout)
log.SetOutput(os.Stdout)
log.SetFlags(log.LstdFlags | log.Lshortfile)
flag.Parse()
cf := flag.Arg(0)
if cf == "" {
fmt.Println("Config file not sepcified")
fmt.Println("The config file is not specified")
return
}
@@ -59,8 +26,10 @@ func main() {
fmt.Println(err)
return
}
rs := RestartableServer{
configFile: cf,
rs := dbserver.DatabaseServer{
ConfigFile: cf,
}
rs.Serve(config)
}
+19 -183
View File
@@ -1,202 +1,38 @@
package server
import (
"encoding/json"
"fmt"
"log"
"net"
"net/http"
"net/url"
"text/template"
"github.com/labring/sealos/service/pkg/auth"
"github.com/labring/sealos/service/database/api"
"github.com/labring/sealos/service/database/request"
"github.com/labring/sealos/service/pkg/server"
)
type PromServer struct {
Config *Config
type DatabaseServer struct {
ConfigFile string
}
func NewPromServer(c *Config) (*PromServer, error) {
ps := &PromServer{
Config: c,
}
return ps, nil
}
func (ps *PromServer) Authenticate(pr *api.PromRequest) error {
return auth.Authenticate(pr.NS, pr.Pwd)
}
func (ps *PromServer) Request(pr *api.PromRequest) (*api.QueryResult, error) {
body, err := request.PrometheusPre(pr)
func (rs *DatabaseServer) Serve(c *server.Config) {
ps, err := server.NewPromServer(c)
if err != nil {
return nil, err
fmt.Printf("Failed to create auth server: %s\n", err)
return
}
var result *api.QueryResult
if err := json.Unmarshal(body, &result); err != nil {
return nil, err
}
return result, err
}
func (ps *PromServer) DBReq(pr *api.PromRequest) (*api.QueryResult, error) {
body, err := request.PrometheusNew(pr)
hs := &http.Server{
Addr: c.Server.ListenAddress,
Handler: ps,
}
listener, err := net.Listen("tcp", c.Server.ListenAddress)
if err != nil {
return nil, err
fmt.Println(err)
return
}
var result *api.QueryResult
if err := json.Unmarshal(body, &result); err != nil {
return nil, err
}
return result, nil
}
fmt.Printf("Serve on %s\n", c.Server.ListenAddress)
func (ps *PromServer) ParseRequest(req *http.Request) (*api.PromRequest, error) {
pr := &api.PromRequest{}
auth := req.Header.Get("Authorization")
if pwd, err := url.PathUnescape(auth); err == nil {
pr.Pwd = pwd
} else {
return nil, err
}
if pr.Pwd == "" {
return nil, api.ErrEmptyKubeconfig
}
if err := req.ParseForm(); err != nil {
return nil, err
}
for key, val := range req.Form {
switch key {
case "query":
pr.Query = val[0]
case "step":
pr.Range.Step = val[0]
case "start":
pr.Range.Start = val[0]
case "end":
pr.Range.End = val[0]
case "time":
pr.Range.Time = val[0]
case "namespace":
pr.NS = val[0]
case "type":
pr.Type = val[0]
case "app":
pr.Cluster = val[0]
}
}
if pr.NS == "" || pr.Query == "" {
return nil, api.ErrUncompleteParam
}
return pr, nil
}
// 获取客户端请求的信息
func (ps *PromServer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
pathPrefix := ""
switch {
case req.URL.Path == pathPrefix+"/query": // 将废弃
ps.doReqPre(rw, req)
case req.URL.Path == pathPrefix+"/q":
ps.doReqNew(rw, req)
default:
http.Error(rw, "Not found", http.StatusNotFound)
return
}
}
func (ps *PromServer) doReqNew(rw http.ResponseWriter, req *http.Request) {
pr, err := ps.ParseRequest(req)
if err != nil {
http.Error(rw, fmt.Sprintf("Bad request (%s)", err), http.StatusBadRequest)
log.Printf("Bad request (%s)\n", err)
return
}
if err := ps.Authenticate(pr); err != nil {
http.Error(rw, fmt.Sprintf("Authentication failed (%s)", err), http.StatusInternalServerError)
log.Printf("Authentication failed (%s)\n", err)
return
}
res, err := ps.DBReq(pr)
if err != nil {
http.Error(rw, fmt.Sprintf("Query failed (%s)", err), http.StatusInternalServerError)
log.Printf("Query failed (%s)\n", err)
return
}
result, err := json.Marshal(res)
if err != nil {
http.Error(rw, "Result failed (invalid query expression)", http.StatusInternalServerError)
log.Printf("Reulst failed (%s)\n", err)
return
}
rw.Header().Set("Content-Type", "application/json")
tmpl := template.New("responseTemplate").Delims("{{", "}}")
tmpl, err = tmpl.Parse(`{{.}}`)
if err != nil {
log.Printf("template failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
return
}
if err = tmpl.Execute(rw, string(result)); err != nil {
log.Printf("Reulst failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
return
}
}
func (ps *PromServer) doReqPre(rw http.ResponseWriter, req *http.Request) {
pr, err := ps.ParseRequest(req)
if err != nil {
http.Error(rw, fmt.Sprintf("Bad request (%s)", err), http.StatusBadRequest)
log.Printf("Bad request (%s)\n", err)
return
}
if err := ps.Authenticate(pr); err != nil {
http.Error(rw, fmt.Sprintf("Authentication failed (%s)", err), http.StatusInternalServerError)
log.Printf("Authentication failed (%s)\n", err)
log.Printf("Kubeconfig (%s)\n", pr.Pwd)
}
res, err := ps.Request(pr)
if err != nil {
http.Error(rw, fmt.Sprintf("Query failed (%s)", err), http.StatusInternalServerError)
log.Printf("Query failed (%s)\n", err)
return
}
result, err := json.Marshal(res)
if err != nil {
http.Error(rw, "Result failed (invalid query expression)", http.StatusInternalServerError)
log.Printf("Reulst failed (%s)\n", err)
return
}
rw.Header().Set("Content-Type", "application/json")
tmpl := template.New("responseTemplate").Delims("{{", "}}")
tmpl, err = tmpl.Parse(`{{.}}`)
if err != nil {
log.Printf("template failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
return
}
if err = tmpl.Execute(rw, string(result)); err != nil {
log.Printf("Reulst failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
if err := hs.Serve(listener); err != nil {
fmt.Println(err)
return
}
}
+9
View File
@@ -0,0 +1,9 @@
# FROM scratch
FROM gcr.io/distroless/static:nonroot
# FROM gengweifeng/gcr-io-distroless-static-nonroot
ARG TARGETARCH
COPY bin/service-minio-$TARGETARCH /manager
EXPOSE 9090
USER 65532:65532
ENTRYPOINT ["/manager"]
+55
View File
@@ -0,0 +1,55 @@
IMG ?= ghcr.io/labring/sealos-minio-service:latest
# Get the currently used golang install path (in GOPATH/bin, unless GOBIN is set)
ifeq (,$(shell go env GOBIN))
GOBIN=$(shell go env GOPATH)/bin
else
GOBIN=$(shell go env GOBIN)
endif
# only support linux, non cgo
PLATFORMS ?= linux_arm64 linux_amd64
GOOS=linux
CGO_ENABLED=0
GOARCH=$(shell go env GOARCH)
GO_BUILD_FLAGS=-trimpath -ldflags "-s -w"
.PHONY: all
all: build
##@ General
# The help target prints out all targets with their descriptions organized
# beneath their categories. The categories are represented by '##@' and the
# target descriptions by '##'. The awk commands is responsible for reading the
# entire set of makefiles included in this invocation, looking for lines of the
# file as xyz: ## something, and then pretty-format the target and help. Then,
# if there's a line with ##@ something, that gets pretty-printed as a category.
# More info on the usage of ANSI control characters for terminal formatting:
# https://en.wikipedia.org/wiki/ANSI_escape_code#SGR_parameters
# More info on the awk command:
# http://linuxcommand.org/lc3_adv_awk.php
.PHONY: help
help: ## Display this help.
@awk 'BEGIN {FS = ":.*##"; printf "\nUsage:\n make \033[36m<target>\033[0m\n"} /^[a-zA-Z_0-9-]+:.*?##/ { printf " \033[36m%-15s\033[0m %s\n", $$1, $$2 } /^##@/ { printf "\n\033[1m%s\033[0m\n", substr($$0, 5) } ' $(MAKEFILE_LIST)
##@ Build
.PHONY: clean
clean:
rm -f $(SERVICE_NAME)
.PHONY: build
build: clean ## Build service-hub binary.
CGO_ENABLED=$(CGO_ENABLED) GOOS=$(GOOS) go build $(GO_BUILD_FLAGS) -o bin/manager main.go
.PHONY: docker-build
docker-build: build
mv bin/manager bin/service-minio-${TARGETARCH}
docker build -t $(IMG) .
.PHONY: docker-push
docker-push:
docker push $(IMG)
+5
View File
@@ -0,0 +1,5 @@
FROM scratch
COPY registry registry
COPY manifests manifests
CMD ["kubectl apply -f manifests/deploy.yaml"]
@@ -0,0 +1,88 @@
apiVersion: v1
kind: ConfigMap
metadata:
labels:
app: object-storage-monitor
name: object-storage-monitor-config
namespace: objectstorage-system
data:
config.yml: |
server:
addr: ":9090"
---
apiVersion: apps/v1
kind: Deployment
metadata:
labels:
app: object-storage-monitor
name: object-storage-monitor
namespace: sealos
spec:
replicas: 1
selector:
matchLabels:
app: object-storage-monitor
strategy:
type: Recreate
template:
metadata:
labels:
app: object-storage-monitor
spec:
containers:
- args:
- /config/config.yml
command:
- /manager
env:
- name: PROMETHEUS_SERVICE_HOST
value: http://prometheus-object-storage.objectstorage-system.svc.cluster.local:9090
- name: OBJECT_STORAGE_INSTANCE
value: object-storage.objectstorage-system.svc.cluster.local:80
image: ghcr.io/labring/sealos-minio-service:latest
imagePullPolicy: Always
name: object-storage-monitor
ports:
- containerPort: 9090
protocol: TCP
resources:
limits:
cpu: 500m
memory: 1024Mi
requests:
cpu: 5m
memory: 64Mi
securityContext:
allowPrivilegeEscalation: false
capabilities:
drop:
- ALL
runAsNonRoot: true
terminationMessagePath: /dev/termination-log
terminationMessagePolicy: File
volumeMounts:
- mountPath: /config
name: config-vol
dnsPolicy: ClusterFirst
restartPolicy: Always
volumes:
- configMap:
defaultMode: 420
name: object-storage-monitor-config
name: config-vol
---
apiVersion: v1
kind: Service
metadata:
labels:
app: object-storage-monitor
name: object-storage-monitor
namespace: objectstorage-system
spec:
ports:
- name: http
port: 9090
protocol: TCP
targetPort: 9090
selector:
app: object-storage-monitor
+35
View File
@@ -0,0 +1,35 @@
package main
import (
"flag"
"fmt"
"log"
"os"
minioserver "github.com/labring/sealos/service/minio/server"
"github.com/labring/sealos/service/pkg/server"
)
func main() {
log.SetOutput(os.Stdout)
log.SetFlags(log.LstdFlags | log.Lshortfile)
flag.Parse()
cf := flag.Arg(0)
if cf == "" {
fmt.Println("The config file is not specified")
return
}
config, err := server.InitConfig(cf)
if err != nil {
fmt.Println(err)
return
}
rs := minioserver.MinioServer{
ConfigFile: cf,
}
rs.Serve(config)
}
+38
View File
@@ -0,0 +1,38 @@
package server
import (
"fmt"
"net"
"net/http"
"github.com/labring/sealos/service/pkg/server"
)
type MinioServer struct {
ConfigFile string
}
func (rs *MinioServer) Serve(c *server.Config) {
ps, err := server.NewPromServer(c)
if err != nil {
fmt.Printf("Failed to create auth server: %s\n", err)
return
}
hs := &http.Server{
Addr: c.Server.ListenAddress,
Handler: ps,
}
listener, err := net.Listen("tcp", c.Server.ListenAddress)
if err != nil {
fmt.Println(err)
return
}
fmt.Printf("Serve on %s\n", c.Server.ListenAddress)
if err := hs.Serve(listener); err != nil {
fmt.Println(err)
return
}
}
@@ -10,7 +10,7 @@ import (
"os"
"strings"
"github.com/labring/sealos/service/database/api"
"github.com/labring/sealos/service/pkg/api"
)
func Request(addr string, params *bytes.Buffer) ([]byte, error) {
+202
View File
@@ -0,0 +1,202 @@
package server
import (
"encoding/json"
"fmt"
"log"
"net/http"
"net/url"
"text/template"
"github.com/labring/sealos/service/pkg/auth"
"github.com/labring/sealos/service/pkg/api"
"github.com/labring/sealos/service/pkg/request"
)
type PromServer struct {
Config *Config
}
func NewPromServer(c *Config) (*PromServer, error) {
ps := &PromServer{
Config: c,
}
return ps, nil
}
func (ps *PromServer) Authenticate(pr *api.PromRequest) error {
return auth.Authenticate(pr.NS, pr.Pwd)
}
func (ps *PromServer) Request(pr *api.PromRequest) (*api.QueryResult, error) {
body, err := request.PrometheusPre(pr)
if err != nil {
return nil, err
}
var result *api.QueryResult
if err := json.Unmarshal(body, &result); err != nil {
return nil, err
}
return result, err
}
func (ps *PromServer) DBReq(pr *api.PromRequest) (*api.QueryResult, error) {
body, err := request.PrometheusNew(pr)
if err != nil {
return nil, err
}
var result *api.QueryResult
if err := json.Unmarshal(body, &result); err != nil {
return nil, err
}
return result, nil
}
func (ps *PromServer) ParseRequest(req *http.Request) (*api.PromRequest, error) {
pr := &api.PromRequest{}
auth := req.Header.Get("Authorization")
if pwd, err := url.PathUnescape(auth); err == nil {
pr.Pwd = pwd
} else {
return nil, err
}
if pr.Pwd == "" {
return nil, api.ErrEmptyKubeconfig
}
if err := req.ParseForm(); err != nil {
return nil, err
}
for key, val := range req.Form {
switch key {
case "query":
pr.Query = val[0]
case "step":
pr.Range.Step = val[0]
case "start":
pr.Range.Start = val[0]
case "end":
pr.Range.End = val[0]
case "time":
pr.Range.Time = val[0]
case "namespace":
pr.NS = val[0]
case "type":
pr.Type = val[0]
case "app":
pr.Cluster = val[0]
}
}
if pr.NS == "" || pr.Query == "" {
return nil, api.ErrUncompleteParam
}
return pr, nil
}
// 获取客户端请求的信息
func (ps *PromServer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
pathPrefix := ""
switch {
case req.URL.Path == pathPrefix+"/query": // 将废弃
ps.doReqPre(rw, req)
case req.URL.Path == pathPrefix+"/q":
ps.doReqNew(rw, req)
default:
http.Error(rw, "Not found", http.StatusNotFound)
return
}
}
func (ps *PromServer) doReqNew(rw http.ResponseWriter, req *http.Request) {
pr, err := ps.ParseRequest(req)
if err != nil {
http.Error(rw, fmt.Sprintf("Bad request (%s)", err), http.StatusBadRequest)
log.Printf("Bad request (%s)\n", err)
return
}
if err := ps.Authenticate(pr); err != nil {
http.Error(rw, fmt.Sprintf("Authentication failed (%s)", err), http.StatusInternalServerError)
log.Printf("Authentication failed (%s)\n", err)
return
}
res, err := ps.DBReq(pr)
if err != nil {
http.Error(rw, fmt.Sprintf("Query failed (%s)", err), http.StatusInternalServerError)
log.Printf("Query failed (%s)\n", err)
return
}
result, err := json.Marshal(res)
if err != nil {
http.Error(rw, "Result failed (invalid query expression)", http.StatusInternalServerError)
log.Printf("Reulst failed (%s)\n", err)
return
}
rw.Header().Set("Content-Type", "application/json")
tmpl := template.New("responseTemplate").Delims("{{", "}}")
tmpl, err = tmpl.Parse(`{{.}}`)
if err != nil {
log.Printf("template failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
return
}
if err = tmpl.Execute(rw, string(result)); err != nil {
log.Printf("Reulst failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
return
}
}
func (ps *PromServer) doReqPre(rw http.ResponseWriter, req *http.Request) {
pr, err := ps.ParseRequest(req)
if err != nil {
http.Error(rw, fmt.Sprintf("Bad request (%s)", err), http.StatusBadRequest)
log.Printf("Bad request (%s)\n", err)
return
}
if err := ps.Authenticate(pr); err != nil {
http.Error(rw, fmt.Sprintf("Authentication failed (%s)", err), http.StatusInternalServerError)
log.Printf("Authentication failed (%s)\n", err)
log.Printf("Kubeconfig (%s)\n", pr.Pwd)
}
res, err := ps.Request(pr)
if err != nil {
http.Error(rw, fmt.Sprintf("Query failed (%s)", err), http.StatusInternalServerError)
log.Printf("Query failed (%s)\n", err)
return
}
result, err := json.Marshal(res)
if err != nil {
http.Error(rw, "Result failed (invalid query expression)", http.StatusInternalServerError)
log.Printf("Reulst failed (%s)\n", err)
return
}
rw.Header().Set("Content-Type", "application/json")
tmpl := template.New("responseTemplate").Delims("{{", "}}")
tmpl, err = tmpl.Parse(`{{.}}`)
if err != nil {
log.Printf("template failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
return
}
if err = tmpl.Execute(rw, string(result)); err != nil {
log.Printf("Reulst failed: %s\n", err)
http.Error(rw, "Internal Server Error", http.StatusInternalServerError)
return
}
}