From 6f2452d4deccb44ff4cf6f74218e603d116f64ad Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 5 Dec 2025 09:59:06 +0800 Subject: [PATCH] feat(service): vlogs service add database log query. (#6298) feat(service): vlogs service add database log query. (#6273) * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. * add database logs query. Co-authored-by: Bearix <2669184984@qq.com> --- service/pkg/api/req.go | 18 +- service/pkg/auth/authenticate.go | 40 +++++ service/vlogs/{server => config}/config.go | 2 +- service/vlogs/main.go | 5 +- .../vlogs/{server/query.go => query/app.go} | 20 +-- .../query_test.go => query/app_test.go} | 6 +- service/vlogs/query/database.go | 97 ++++++++++ service/vlogs/query/database_test.go | 170 ++++++++++++++++++ service/vlogs/server/handler.go | 164 +++++++++++++++++ service/vlogs/server/helper.go | 58 ++++++ service/vlogs/server/server.go | 135 +------------- 11 files changed, 567 insertions(+), 148 deletions(-) rename service/vlogs/{server => config}/config.go (97%) rename service/vlogs/{server/query.go => query/app.go} (88%) rename service/vlogs/{server/query_test.go => query/app_test.go} (96%) create mode 100644 service/vlogs/query/database.go create mode 100644 service/vlogs/query/database_test.go create mode 100644 service/vlogs/server/handler.go create mode 100644 service/vlogs/server/helper.go diff --git a/service/pkg/api/req.go b/service/pkg/api/req.go index 026f4b377..c8a8093da 100644 --- a/service/pkg/api/req.go +++ b/service/pkg/api/req.go @@ -78,7 +78,7 @@ type JSONQuery struct { Value string } -type VlogsRequest struct { +type VlogsLaunchpadRequest struct { Time string `json:"time"` Namespace string `json:"namespace"` App string `json:"app"` @@ -96,7 +96,21 @@ type VlogsRequest struct { EndTime string `json:"endTime,omitempty"` } -type VlogsResponse struct { +type VlogsDatabaseRequest struct { + Time string `json:"time"` + Namespace string `json:"namespace"` + Limit string `json:"limit,omitempty"` + NumberMode string `json:"numberMode,omitempty"` + NumberLevel string `json:"numberLevel,omitempty"` + Pvc []string `json:"pvc,omitempty"` + Type []string `json:"type,omitempty"` + Container []string `json:"container,omitempty"` + Keyword string `json:"keyword,omitempty"` + StartTime string `json:"startTime,omitempty"` + EndTime string `json:"endTime,omitempty"` +} + +type VlogsLaunchpadResponse struct { Time string `json:"_time"` Message string `json:"_msg"` Container string `json:"container"` diff --git a/service/pkg/auth/authenticate.go b/service/pkg/auth/authenticate.go index 721bd29c5..896061329 100644 --- a/service/pkg/auth/authenticate.go +++ b/service/pkg/auth/authenticate.go @@ -21,6 +21,8 @@ var ( ErrNilNs = errors.New("namespace not found") ErrNoAuth = errors.New("no permission for this namespace") ErrNoSealosHost = errors.New("unable to get the sealos host") + ErrPVCNotFound = errors.New("pvc not found in namespace") + ErrNilPVCUID = errors.New("pvc uid not provided") whiteListKubernetesHosts []string ) @@ -104,6 +106,44 @@ func Authenticate(ns, kc string) error { return nil } +func AuthenticatePVC(ns, kc string, pvcUIDs []string) error { + if ns == "" { + return ErrNilNs + } + if len(pvcUIDs) == 0 { + return ErrNilPVCUID + } + config, err := clientcmd.RESTConfigFromKubeConfig([]byte(kc)) + if err != nil { + return fmt.Errorf("kubeconfig failed: %w", err) + } + if !IsWhitelistKubernetesHost(config.Host) { + if k8shost := GetKubernetesHostFromEnv(); k8shost != "" { + config.Host = k8shost + } else { + return ErrNoSealosHost + } + } + client, err := kubernetes.NewForConfig(config) + if err != nil { + return fmt.Errorf("failed to new client: %w", err) + } + pvcList, err := client.CoreV1().PersistentVolumeClaims(ns).List(context.TODO(), metav1.ListOptions{}) + if err != nil { + return fmt.Errorf("failed to list pvcs: %w", err) + } + pvcUIDSet := make(map[string]struct{}, len(pvcList.Items)) + for _, pvc := range pvcList.Items { + pvcUIDSet["pvc-"+string(pvc.UID)] = struct{}{} + } + for _, uid := range pvcUIDs { + if _, ok := pvcUIDSet[uid]; !ok { + return fmt.Errorf("%w: %s", ErrPVCNotFound, uid) + } + } + return nil +} + func IsWhitelistKubernetesHost(host string) bool { for _, h := range whiteListKubernetesHosts { if h == host { diff --git a/service/vlogs/server/config.go b/service/vlogs/config/config.go similarity index 97% rename from service/vlogs/server/config.go rename to service/vlogs/config/config.go index 52c762127..e40928acc 100644 --- a/service/vlogs/server/config.go +++ b/service/vlogs/config/config.go @@ -1,4 +1,4 @@ -package server +package config import ( "fmt" diff --git a/service/vlogs/main.go b/service/vlogs/main.go index 8e66e5ec3..52179fed9 100644 --- a/service/vlogs/main.go +++ b/service/vlogs/main.go @@ -9,6 +9,7 @@ import ( "os" "time" + "github.com/labring/sealos/service/vlogs/config" vlogsServer "github.com/labring/sealos/service/vlogs/server" ) @@ -16,7 +17,7 @@ type RestartableServer struct { configFile string } -func (rs *RestartableServer) Serve(c *vlogsServer.Config) { +func (rs *RestartableServer) Serve(c *config.Config) { vs, err := vlogsServer.NewVLogsServer(c) if err != nil { fmt.Printf("Failed to create auth server: %s\n", err) @@ -53,7 +54,7 @@ func main() { return } - config, err := vlogsServer.InitConfig(cf) + config, err := config.InitConfig(cf) if err != nil { fmt.Println(err) return diff --git a/service/vlogs/server/query.go b/service/vlogs/query/app.go similarity index 88% rename from service/vlogs/server/query.go rename to service/vlogs/query/app.go index 11321ddab..b35a5b073 100644 --- a/service/vlogs/server/query.go +++ b/service/vlogs/query/app.go @@ -1,4 +1,4 @@ -package server +package query import ( "fmt" @@ -17,7 +17,7 @@ type VLogsQuery struct { query string } -func (v *VLogsQuery) getQuery(req *api.VlogsRequest) (string, error) { +func (v *VLogsQuery) GetQuery(req *api.VlogsLaunchpadRequest) (string, error) { if req.PodQuery == modeTrue { query := v.generatePodListQuery(req) return query, nil @@ -41,7 +41,7 @@ func EscapeSingleQuoted(s string) string { return s } -func (v *VLogsQuery) generatePodListQuery(req *api.VlogsRequest) string { +func (v *VLogsQuery) generatePodListQuery(req *api.VlogsLaunchpadRequest) string { var item string if len(req.Time) != 0 { item = fmt.Sprintf( @@ -57,7 +57,7 @@ func (v *VLogsQuery) generatePodListQuery(req *api.VlogsRequest) string { return v.query } -func (v *VLogsQuery) generateKeywordQuery(req *api.VlogsRequest) { +func (v *VLogsQuery) generateKeywordQuery(req *api.VlogsLaunchpadRequest) { if len(req.Keyword) == 0 { return } else { @@ -65,7 +65,7 @@ func (v *VLogsQuery) generateKeywordQuery(req *api.VlogsRequest) { } } -func (v *VLogsQuery) generateJSONQuery(req *api.VlogsRequest) error { +func (v *VLogsQuery) generateJSONQuery(req *api.VlogsLaunchpadRequest) error { if req.JSONMode != modeTrue { return nil } @@ -95,7 +95,7 @@ func (v *VLogsQuery) generateJSONQuery(req *api.VlogsRequest) error { return nil } -func (v *VLogsQuery) generateStreamQuery(req *api.VlogsRequest) { +func (v *VLogsQuery) generateStreamQuery(req *api.VlogsLaunchpadRequest) { namespace := EscapeSingleQuoted(req.Namespace) var builder strings.Builder switch { @@ -145,7 +145,7 @@ func (v *VLogsQuery) generateStreamQuery(req *api.VlogsRequest) { v.query += builder.String() } -func (v *VLogsQuery) generateCommonQuery(req *api.VlogsRequest) { +func (v *VLogsQuery) generateCommonQuery(req *api.VlogsLaunchpadRequest) { var builder strings.Builder var item string if len(req.Time) != 0 { @@ -166,7 +166,7 @@ func (v *VLogsQuery) generateCommonQuery(req *api.VlogsRequest) { // if query number,dont use limit param if req.NumberMode == modeFalse { var item string - if hasNonDigits(req.Limit) { + if req.Limit == "" || HasNonDigits(req.Limit) { item = ` | limit '100' ` } else { item = fmt.Sprintf(` | limit '%s' `, EscapeSingleQuoted(req.Limit)) @@ -188,7 +188,7 @@ var allowedNumberLevels = map[string]struct{}{ "s": {}, } -func hasNonDigits(s string) bool { +func HasNonDigits(s string) bool { _, err := strconv.Atoi(s) return err != nil } @@ -198,7 +198,7 @@ func isValidNumberLevel(level string) bool { return ok } -func (v *VLogsQuery) generateNumberQuery(req *api.VlogsRequest) { +func (v *VLogsQuery) generateNumberQuery(req *api.VlogsLaunchpadRequest) { if req.NumberMode == modeTrue { if isValidNumberLevel(req.NumberLevel) { item := fmt.Sprintf( diff --git a/service/vlogs/server/query_test.go b/service/vlogs/query/app_test.go similarity index 96% rename from service/vlogs/server/query_test.go rename to service/vlogs/query/app_test.go index e69c06273..84c03f366 100644 --- a/service/vlogs/server/query_test.go +++ b/service/vlogs/query/app_test.go @@ -1,6 +1,8 @@ -package server +package query -import "testing" +import ( + "testing" +) func TestEscapeSingleQuoted(t *testing.T) { type args struct { diff --git a/service/vlogs/query/database.go b/service/vlogs/query/database.go new file mode 100644 index 000000000..07b1ffc41 --- /dev/null +++ b/service/vlogs/query/database.go @@ -0,0 +1,97 @@ +package query + +import ( + "fmt" + "strings" + + "github.com/labring/sealos/service/pkg/api" +) + +type DBLogsQuery struct { + query string +} + +func (v *DBLogsQuery) GetDBQuery(req *api.VlogsDatabaseRequest) (string, error) { + v.query = "" + v.generateVolumeUIDQuery(req) + v.generateKeywordQuery(req) + v.generateContainerQuery(req) + v.generateTypeQuery(req) + v.generateCommonQuery(req) + v.generateNumberQuery(req) + v.generateSortQuery(req) + fmt.Printf("database query: %s\n", v.query) + return v.query, nil +} + +func (v *DBLogsQuery) generateVolumeUIDQuery(req *api.VlogsDatabaseRequest) { + if len(req.Pvc) == 0 { + return + } + pvcPattern := strings.Join(req.Pvc, "|") + v.query += fmt.Sprintf(`{volume_uid=~'%s'} `, EscapeSingleQuoted(pvcPattern)) +} + +func (v *DBLogsQuery) generateKeywordQuery(req *api.VlogsDatabaseRequest) { + if req.Keyword != "" { + v.query += fmt.Sprintf(`'%s' `, EscapeSingleQuoted(req.Keyword)) + } +} + +func (v *DBLogsQuery) generateContainerQuery(req *api.VlogsDatabaseRequest) { + if len(req.Container) == 0 { + return + } + escapedContainers := make([]string, 0, len(req.Container)) + for _, container := range req.Container { + escaped := EscapeSingleQuoted(container) + escapedContainers = append(escapedContainers, "'"+escaped+"'") + } + containerList := strings.Join(escapedContainers, ",") + v.query += fmt.Sprintf(`container:in(%s) `, containerList) +} + +func (v *DBLogsQuery) generateTypeQuery(req *api.VlogsDatabaseRequest) { + if len(req.Type) == 0 { + return + } + escapedTypes := make([]string, 0, len(req.Type)) + for _, logType := range req.Type { + escaped := EscapeSingleQuoted(logType) + escapedTypes = append(escapedTypes, "'"+escaped+"'") + } + typeList := strings.Join(escapedTypes, ",") + v.query += fmt.Sprintf(`log_type:in(%s) `, typeList) +} + +func (v *DBLogsQuery) generateCommonQuery(req *api.VlogsDatabaseRequest) { + var filters []string + if req.Time != "" { + filters = append(filters, fmt.Sprintf(`time'%s' `, EscapeSingleQuoted(req.Time))) + } + if len(filters) > 0 { + v.query += strings.Join(filters, " ") + } + if req.NumberMode != modeTrue { + limit := req.Limit + if limit == "" || HasNonDigits(req.Limit) { + limit = "100" + } + v.query += " | limit " + EscapeSingleQuoted(limit) + } +} + +func (v *DBLogsQuery) generateNumberQuery(req *api.VlogsDatabaseRequest) { + if req.NumberMode == modeTrue && isValidNumberLevel(req.NumberLevel) { + v.query += fmt.Sprintf( + ` | stats by (_time:1%s) count() logs_total`, + EscapeSingleQuoted(req.NumberLevel), + ) + } +} + +func (v *DBLogsQuery) generateSortQuery(req *api.VlogsDatabaseRequest) { + if req.NumberMode != modeTrue { + v.query += ` | sort by (_time) desc` + } +} diff --git a/service/vlogs/query/database_test.go b/service/vlogs/query/database_test.go new file mode 100644 index 000000000..ed42c58e1 --- /dev/null +++ b/service/vlogs/query/database_test.go @@ -0,0 +1,170 @@ +package query + +import ( + "testing" + + "github.com/labring/sealos/service/pkg/api" +) + +func TestDBLogsQuery_GetDBQuery(t *testing.T) { + tests := []struct { + name string + req *api.VlogsDatabaseRequest + want string + }{ + // Normal cases + { + name: "normal query with all filters", + req: &api.VlogsDatabaseRequest{ + Pvc: []string{"pvc-001"}, + Keyword: "error", + Container: []string{"nginx"}, + Type: []string{"stdout"}, + Time: "2024-01-01", + Limit: "50", + }, + want: "{volume_uid=~'pvc-001'} 'error' container:in('nginx') log_type:in('stdout') time'2024-01-01' | limit 50 | sort by (_time) desc", + }, + + // Empty value tests + { + name: "empty pvc list", + req: &api.VlogsDatabaseRequest{Pvc: []string{}}, + want: " | limit 100 | sort by (_time) desc", + }, + { + name: "empty keyword", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc-001"}, Keyword: ""}, + want: "{volume_uid=~'pvc-001'} | limit 100 | sort by (_time) desc", + }, + { + name: "empty container list", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc-001"}, Container: []string{}}, + want: "{volume_uid=~'pvc-001'} | limit 100 | sort by (_time) desc", + }, + { + name: "empty type list", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc-001"}, Type: []string{}}, + want: "{volume_uid=~'pvc-001'} | limit 100 | sort by (_time) desc", + }, + { + name: "empty time", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc-001"}, Time: ""}, + want: "{volume_uid=~'pvc-001'} | limit 100 | sort by (_time) desc", + }, + + // Invalid input tests + { + name: "invalid limit with letters", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc-001"}, Limit: "abc"}, + want: "{volume_uid=~'pvc-001'} | limit 100 | sort by (_time) desc", + }, + { + name: "invalid limit with special characters", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc-001"}, Limit: "50@#$"}, + want: "{volume_uid=~'pvc-001'} | limit 100 | sort by (_time) desc", + }, + { + name: "invalid number level", + req: &api.VlogsDatabaseRequest{ + Pvc: []string{"pvc-001"}, + NumberMode: modeTrue, + NumberLevel: "invalid", + }, + want: "{volume_uid=~'pvc-001'} ", + }, + + // Special character escaping tests + { + name: "pvc with single quote", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc'001"}}, + want: "{volume_uid=~'pvc\\'001'} | limit 100 | sort by (_time) desc", + }, + { + name: "keyword with single quote", + req: &api.VlogsDatabaseRequest{Keyword: "it's error"}, + want: "'it\\'s error' | limit 100 | sort by (_time) desc", + }, + { + name: "container with single quote", + req: &api.VlogsDatabaseRequest{Container: []string{"nginx'test"}}, + want: "container:in('nginx\\'test') | limit 100 | sort by (_time) desc", + }, + { + name: "pvc with backslash", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc\\001"}}, + want: "{volume_uid=~'pvc\\\\001'} | limit 100 | sort by (_time) desc", + }, + + // SQL injection tests + { + name: "keyword sql injection attempt", + req: &api.VlogsDatabaseRequest{Keyword: "'; DROP TABLE logs; --"}, + want: "'\\'; DROP TABLE logs; --' | limit 100 | sort by (_time) desc", + }, + { + name: "container sql injection attempt", + req: &api.VlogsDatabaseRequest{Container: []string{"nginx' OR '1'='1"}}, + want: "container:in('nginx\\' OR \\'1\\'=\\'1') | limit 100 | sort by (_time) desc", + }, + + // Number mode tests + { + name: "number mode with hour level", + req: &api.VlogsDatabaseRequest{ + Pvc: []string{"pvc-001"}, + NumberMode: modeTrue, + NumberLevel: "h", + }, + want: "{volume_uid=~'pvc-001'} | stats by (_time:1h) count() logs_total", + }, + { + name: "number mode with minute level", + req: &api.VlogsDatabaseRequest{ + Pvc: []string{"pvc-001"}, + NumberMode: modeTrue, + NumberLevel: "m", + }, + want: "{volume_uid=~'pvc-001'} | stats by (_time:1m) count() logs_total", + }, + { + name: "number mode disabled", + req: &api.VlogsDatabaseRequest{ + Pvc: []string{"pvc-001"}, + NumberMode: "false", + }, + want: "{volume_uid=~'pvc-001'} | limit 100 | sort by (_time) desc", + }, + + // Edge cases + { + name: "completely empty request", + req: &api.VlogsDatabaseRequest{}, + want: " | limit 100 | sort by (_time) desc", + }, + { + name: "very large limit", + req: &api.VlogsDatabaseRequest{Limit: "999999999999"}, + want: " | limit 999999999999 | sort by (_time) desc", + }, + { + name: "multiple pvcs", + req: &api.VlogsDatabaseRequest{Pvc: []string{"pvc-001", "pvc-002", "pvc-003"}}, + want: "{volume_uid=~'pvc-001|pvc-002|pvc-003'} | limit 100 | sort by (_time) desc", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + v := &DBLogsQuery{} + got, err := v.GetDBQuery(tt.req) + if err != nil { + t.Errorf("GetDBQuery() error = %v", err) + return + } + if got != tt.want { + t.Errorf("GetDBQuery()\ngot = %q\nwant = %q", got, tt.want) + } + }) + } +} diff --git a/service/vlogs/server/handler.go b/service/vlogs/server/handler.go new file mode 100644 index 000000000..021def131 --- /dev/null +++ b/service/vlogs/server/handler.go @@ -0,0 +1,164 @@ +package server + +import ( + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + + "github.com/labring/sealos/service/pkg/api" + "github.com/labring/sealos/service/pkg/auth" + "github.com/labring/sealos/service/vlogs/query" + "github.com/labring/sealos/service/vlogs/request" +) + +func (vl *VLogsServer) queryLogsByParams(rw http.ResponseWriter, req *http.Request) error { + resp, err := vl.executeQuery(req) + if err != nil { + return err + } + defer resp.Close() + _, err = io.Copy(rw, resp) + if err != nil { + return err + } + return nil +} + +func (vl *VLogsServer) queryPodList(rw http.ResponseWriter, req *http.Request) error { + resp, err := vl.executeQuery(req) + if err != nil { + return err + } + defer resp.Close() + body, err := io.ReadAll(resp) + if err != nil { + return fmt.Errorf("failed to read response body: %w", err) + } + if len(body) == 0 { + return errors.New("response body is empty") + } + podList, err := vl.extractUniquePods(body) + if err != nil { + return fmt.Errorf("failed to extract pod list: %w", err) + } + rw.Header().Set("Content-Type", "application/json") + if err := json.NewEncoder(rw).Encode(podList); err != nil { + return fmt.Errorf("failed to write response: %w", err) + } + return nil +} + +func (vl *VLogsServer) executeQuery(req *http.Request) (io.ReadCloser, error) { + vlogsReq, kubeConfig, err := vl.verifyParams(req) + if err != nil { + return nil, fmt.Errorf("bad request (%w)", err) + } + err = auth.Authenticate(vlogsReq.Namespace, kubeConfig) + if err != nil { + return nil, fmt.Errorf("authentication failed (%w)", err) + } + var vlogs query.VLogsQuery + query, err := vlogs.GetQuery(vlogsReq) + if err != nil { + return nil, fmt.Errorf("failed to parse request body: %w", err) + } + resp, err := request.QueryLogsByParams(&request.QueryParams{ + Path: vl.path, + Query: query, + Username: vl.username, + Password: vl.password, + StartTime: vlogsReq.StartTime, + EndTime: vlogsReq.EndTime, + }) + if err != nil { + return nil, fmt.Errorf("query failed (%w)", err) + } + return resp, nil +} + +func (vl *VLogsServer) verifyParams(req *http.Request) (*api.VlogsLaunchpadRequest, string, error) { + kubeConfig, err := vl.extractKubeConfig(req) + if err != nil { + return nil, "", err + } + vlogsReq := &api.VlogsLaunchpadRequest{} + err = json.NewDecoder(req.Body).Decode(&vlogsReq) + if err != nil { + return nil, "", fmt.Errorf("failed to parse request body: %w", err) + } + if vlogsReq.Namespace == "" { + return nil, "", errors.New("failed to get namespace") + } + if err := vl.validateTimeParams(vlogsReq.StartTime, vlogsReq.EndTime, vlogsReq.Time); err != nil { + return nil, "", err + } + return vlogsReq, kubeConfig, nil +} + +func (vl *VLogsServer) queryDBLogs(rw http.ResponseWriter, req *http.Request) error { + resp, err := vl.executeDBQuery(req) + if err != nil { + return err + } + defer resp.Close() + _, err = io.Copy(rw, resp) + if err != nil { + return err + } + return nil +} + +func (vl *VLogsServer) executeDBQuery(req *http.Request) (io.ReadCloser, error) { + vlogsReq, kubeConfig, err := vl.verifyDBParams(req) + if err != nil { + return nil, fmt.Errorf("bad request (%w)", err) + } + err = auth.Authenticate(vlogsReq.Namespace, kubeConfig) + if err != nil { + return nil, fmt.Errorf("authentication failed (%w)", err) + } + err = auth.AuthenticatePVC(vlogsReq.Namespace, kubeConfig, vlogsReq.Pvc) + if err != nil { + return nil, fmt.Errorf("authentication pvc failed (%w)", err) + } + var vlogs query.DBLogsQuery + query, err := vlogs.GetDBQuery(vlogsReq) + if err != nil { + return nil, fmt.Errorf("failed to parse request body: %w", err) + } + resp, err := request.QueryLogsByParams(&request.QueryParams{ + Path: vl.path, + Query: query, + Username: vl.username, + Password: vl.password, + StartTime: vlogsReq.StartTime, + EndTime: vlogsReq.EndTime, + }) + if err != nil { + return nil, fmt.Errorf("query failed (%w)", err) + } + return resp, nil +} + +func (vl *VLogsServer) verifyDBParams( + req *http.Request, +) (*api.VlogsDatabaseRequest, string, error) { + kubeConfig, err := vl.extractKubeConfig(req) + if err != nil { + return nil, "", err + } + vlogsReq := &api.VlogsDatabaseRequest{} + err = json.NewDecoder(req.Body).Decode(&vlogsReq) + if err != nil { + return nil, "", fmt.Errorf("failed to parse request body: %w", err) + } + if vlogsReq.Namespace == "" { + return nil, "", errors.New("failed to get namespace") + } + if err := vl.validateTimeParams(vlogsReq.StartTime, vlogsReq.EndTime, vlogsReq.Time); err != nil { + return nil, "", err + } + return vlogsReq, kubeConfig, nil +} diff --git a/service/vlogs/server/helper.go b/service/vlogs/server/helper.go new file mode 100644 index 000000000..30bf1152f --- /dev/null +++ b/service/vlogs/server/helper.go @@ -0,0 +1,58 @@ +package server + +import ( + "bufio" + "bytes" + "encoding/json" + "errors" + "fmt" + "net/http" + "net/url" + + "github.com/labring/sealos/service/pkg/api" +) + +func (vl *VLogsServer) extractKubeConfig(req *http.Request) (string, error) { + kubeConfig := req.Header.Get("Authorization") + if config, err := url.PathUnescape(kubeConfig); err == nil { + kubeConfig = config + } else { + return "", fmt.Errorf("failed to PathUnescape : %w", err) + } + return kubeConfig, nil +} + +func (vl *VLogsServer) validateTimeParams(startTime, endTime, time string) error { + if startTime == "" && endTime != "" { + return errors.New("failed to get start time") + } + if startTime != "" && endTime == "" { + return errors.New("failed to get end time") + } + if startTime != "" && endTime != "" && time != "" { + return errors.New("not to provide 3 time params") + } + return nil +} + +func (vl *VLogsServer) extractUniquePods(body []byte) ([]string, error) { + scanner := bufio.NewScanner(bytes.NewReader(body)) + uniquePods := make(map[string]struct{}) + for scanner.Scan() { + var entry api.VlogsLaunchpadResponse + lineBytes := scanner.Bytes() + err := json.Unmarshal(lineBytes, &entry) + if err != nil { + continue + } + uniquePods[entry.Pod] = struct{}{} + } + if err := scanner.Err(); err != nil { + return nil, fmt.Errorf("error reading response: %w", err) + } + podList := make([]string, 0, len(uniquePods)) + for pod := range uniquePods { + podList = append(podList, pod) + } + return podList, nil +} diff --git a/service/vlogs/server/server.go b/service/vlogs/server/server.go index d34bee50b..31808c33e 100644 --- a/service/vlogs/server/server.go +++ b/service/vlogs/server/server.go @@ -1,19 +1,12 @@ package server import ( - "bufio" - "bytes" - "encoding/json" "errors" "fmt" - "io" "log/slog" "net/http" - "net/url" - "github.com/labring/sealos/service/pkg/api" - "github.com/labring/sealos/service/pkg/auth" - "github.com/labring/sealos/service/vlogs/request" + "github.com/labring/sealos/service/vlogs/config" ) type VLogsServer struct { @@ -22,7 +15,7 @@ type VLogsServer struct { password string } -func NewVLogsServer(config *Config) (*VLogsServer, error) { +func NewVLogsServer(config *config.Config) (*VLogsServer, error) { vl := &VLogsServer{ path: config.Server.Path, username: config.Server.Username, @@ -61,129 +54,9 @@ func (vl *VLogsServer) queryConvert( return vl.queryLogsByParams, nil case "/queryPodList": return vl.queryPodList, nil + case "/queryLogsByPod": + return vl.queryDBLogs, nil default: return nil, errors.New("unknown url path") } } - -func (vl *VLogsServer) verifyParams(req *http.Request) (*api.VlogsRequest, string, error) { - kubeConfig := req.Header.Get("Authorization") - if config, err := url.PathUnescape(kubeConfig); err == nil { - kubeConfig = config - } else { - return nil, "", fmt.Errorf("failed to PathUnescape : %w", err) - } - vlogsReq := &api.VlogsRequest{} - err := json.NewDecoder(req.Body).Decode(&vlogsReq) - if err != nil { - return nil, "", fmt.Errorf("failed to parse request body: %w", err) - } - if vlogsReq.Namespace == "" { - return nil, "", errors.New("failed to get namespace") - } - if vlogsReq.StartTime == "" && vlogsReq.EndTime != "" { - return nil, "", errors.New("failed to get start time") - } - if vlogsReq.StartTime != "" && vlogsReq.EndTime == "" { - return nil, "", errors.New("failed to get end time") - } - if vlogsReq.StartTime != "" && vlogsReq.EndTime != "" && vlogsReq.Time != "" { - return nil, "", errors.New("not to provide 3 time params") - } - return vlogsReq, kubeConfig, nil -} - -func (vl *VLogsServer) authenticate(namespace, kubeConfig string) error { - err := auth.Authenticate(namespace, kubeConfig) - if err != nil { - return fmt.Errorf("authentication failed (%w)", err) - } - return nil -} - -func (vl *VLogsServer) executeQuery(req *http.Request) (io.ReadCloser, error) { - vlogsReq, kubeConfig, err := vl.verifyParams(req) - if err != nil { - return nil, fmt.Errorf("bad request (%w)", err) - } - err = vl.authenticate(vlogsReq.Namespace, kubeConfig) - if err != nil { - return nil, err - } - var vlogs VLogsQuery - query, err := vlogs.getQuery(vlogsReq) - if err != nil { - return nil, fmt.Errorf("failed to parse request body: %w", err) - } - resp, err := request.QueryLogsByParams(&request.QueryParams{ - Path: vl.path, - Query: query, - Username: vl.username, - Password: vl.password, - StartTime: vlogsReq.StartTime, - EndTime: vlogsReq.EndTime, - }) - if err != nil { - return nil, fmt.Errorf("query failed (%w)", err) - } - return resp, nil -} - -func (vl *VLogsServer) queryLogsByParams(rw http.ResponseWriter, req *http.Request) error { - resp, err := vl.executeQuery(req) - if err != nil { - return err - } - defer resp.Close() - _, err = io.Copy(rw, resp) - if err != nil { - return err - } - return nil -} - -func (vl *VLogsServer) queryPodList(rw http.ResponseWriter, req *http.Request) error { - resp, err := vl.executeQuery(req) - if err != nil { - return err - } - defer resp.Close() - body, err := io.ReadAll(resp) - if err != nil { - return fmt.Errorf("failed to read response body: %w", err) - } - if len(body) == 0 { - return errors.New("response body is empty") - } - podList, err := vl.extractUniquePods(body) - if err != nil { - return fmt.Errorf("failed to extract pod list: %w", err) - } - rw.Header().Set("Content-Type", "application/json") - if err := json.NewEncoder(rw).Encode(podList); err != nil { - return fmt.Errorf("failed to write response: %w", err) - } - return nil -} - -func (vl *VLogsServer) extractUniquePods(body []byte) ([]string, error) { - scanner := bufio.NewScanner(bytes.NewReader(body)) - uniquePods := make(map[string]struct{}) - for scanner.Scan() { - var entry api.VlogsResponse - lineBytes := scanner.Bytes() - err := json.Unmarshal(lineBytes, &entry) - if err != nil { - continue - } - uniquePods[entry.Pod] = struct{}{} - } - if err := scanner.Err(); err != nil { - return nil, fmt.Errorf("error reading response: %w", err) - } - podList := make([]string, 0, len(uniquePods)) - for pod := range uniquePods { - podList = append(podList, pod) - } - return podList, nil -}