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>
This commit is contained in:
github-actions[bot]
2025-12-05 09:59:06 +08:00
committed by GitHub
co-authored by Bearix
parent f4ce21521c
commit 6f2452d4de
11 changed files with 567 additions and 148 deletions
+16 -2
View File
@@ -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"`
+40
View File
@@ -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 {
@@ -1,4 +1,4 @@
package server
package config
import (
"fmt"
+3 -2
View File
@@ -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
@@ -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(
@@ -1,6 +1,8 @@
package server
package query
import "testing"
import (
"testing"
)
func TestEscapeSingleQuoted(t *testing.T) {
type args struct {
+97
View File
@@ -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`
}
}
+170
View File
@@ -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)
}
})
}
}
+164
View File
@@ -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
}
+58
View File
@@ -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
}
+4 -131
View File
@@ -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
}