feat(vlogs service): fix stderr handling and improve number logic validation. (#5859)

* fix stderr and number logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix(vlogs service): change time query logic.

fix ci lint.

bug fix.

bug fix.

bug fix.

* Update service/vlogs/server/query.go

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>

* bug fix.

* bug fix.

* bug fix.

* bug fix.

---------

Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
Bearix
2025-09-16 16:33:14 +08:00
committed by GitHub
co-authored by Copilot
parent c0fcb760a3
commit 837dc6ff50
7 changed files with 336 additions and 243 deletions
+2
View File
@@ -91,6 +91,8 @@ type VlogsRequest struct {
Keyword string `json:"keyword,omitempty"`
JSONQuery []JSONQuery `json:"jsonQuery,omitempty"`
PodQuery string `json:"podQuery,omitempty"`
StartTime string `json:"startTime,omitempty"`
EndTime string `json:"endTime,omitempty"`
}
type VlogsResponse struct {
+1
View File
@@ -12,6 +12,7 @@ PLATFORMS ?= linux_arm64 linux_amd64
GOOS=linux
CGO_ENABLED=0
GOARCH=$(shell go env GOARCH)
TARGETARCH ?= $(GOARCH)
GO_BUILD_FLAGS=-trimpath -ldflags "-s -w"
+5 -3
View File
@@ -7,6 +7,7 @@ import (
"net"
"net/http"
"os"
"time"
vlogsServer "github.com/labring/sealos/service/vlogs/server"
)
@@ -16,15 +17,16 @@ type RestartableServer struct {
}
func (rs *RestartableServer) Serve(c *vlogsServer.Config) {
var vs, err = vlogsServer.NewVLogsServer(c)
vs, err := vlogsServer.NewVLogsServer(c)
if err != nil {
fmt.Printf("Failed to create auth server: %s\n", err)
return
}
hs := &http.Server{
Addr: c.Server.ListenAddress,
Handler: vs,
Addr: c.Server.ListenAddress,
Handler: vs,
ReadHeaderTimeout: 30 * time.Second,
}
var listener net.Listener
+49 -29
View File
@@ -1,48 +1,68 @@
package request
import (
"crypto/tls"
"fmt"
"io"
"log"
"net/http"
"net/url"
)
func generateReq(path string, username string, password string, query string) (*http.Request, error) {
baseURL, err := url.Parse(path + "/select/logsql/query")
if err != nil {
return nil, fmt.Errorf("can not parser API URL: %v", err)
}
params := url.Values{}
params.Add("query", query)
baseURL.RawQuery = params.Encode()
req, err := http.NewRequest("GET", baseURL.String(), nil)
if err != nil {
return nil, fmt.Errorf("create HTTP req error: %v", err)
}
req.SetBasicAuth(username, password)
return req, nil
type QueryParams struct {
Path string
Username string
Password string
Query string
StartTime string
EndTime string
}
func QueryLogsByParams(path string, username string, password string, query string) (*http.Response, error) {
httpClient := &http.Client{
Transport: &http.Transport{
// nosemgrep
TLSClientConfig: &tls.Config{
InsecureSkipVerify: true,
},
},
}
req, err := generateReq(path, username, password, query)
func QueryLogsByParams(query *QueryParams) (io.ReadCloser, error) {
httpClient := http.DefaultClient
req, err := generateReq(query)
if err != nil {
return nil, err
}
resp, err := httpClient.Do(req)
if err != nil {
return nil, fmt.Errorf("HTTP req error: %v", err)
return nil, fmt.Errorf("HTTP req error: %w", err)
}
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("res error,err info: %+v", resp)
defer resp.Body.Close()
body, readErr := io.ReadAll(resp.Body)
if readErr != nil {
return nil, fmt.Errorf(
"HTTP %d error, unable to read error details: %w",
resp.StatusCode,
readErr,
)
}
log.Printf("=== Victoria Logs Query Failed ===")
log.Printf("Status Code: %d", resp.StatusCode)
log.Printf("Server: %s", resp.Header.Get("X-Server-Hostname"))
log.Printf("Request URL: %s", req.URL.String())
log.Printf("Error Content: %s", string(body))
log.Printf("===================================")
return nil, fmt.Errorf("victoria Logs query failed [%d]: %s", resp.StatusCode, string(body))
}
return resp, nil
return resp.Body, nil
}
func generateReq(query *QueryParams) (*http.Request, error) {
parsedURL, err := url.Parse(query.Path)
if err != nil {
return nil, fmt.Errorf("can not parser API URL: %w", err)
}
parsedURL = parsedURL.JoinPath("select/logsql/query")
params := url.Values{}
params.Add("query", query.Query)
params.Add("start", query.StartTime)
params.Add("end", query.EndTime)
parsedURL.RawQuery = params.Encode()
req, err := http.NewRequest(http.MethodGet, parsedURL.String(), nil)
if err != nil {
return nil, fmt.Errorf("create HTTP req error: %w", err)
}
req.SetBasicAuth(query.Username, query.Password)
return req, nil
}
+2 -2
View File
@@ -21,11 +21,11 @@ type ServeConfig struct {
func InitConfig(configPath string) (*Config, error) {
configData, err := os.ReadFile(configPath)
if err != nil {
return nil, fmt.Errorf("could not read %s: %s", configPath, err)
return nil, fmt.Errorf("could not read %s: %w", configPath, err)
}
c := &Config{}
if err := yaml.Unmarshal(configData, c); err != nil {
return nil, fmt.Errorf("could not parse config: %s", err)
return nil, fmt.Errorf("could not parse config: %w", err)
}
return c, nil
}
+176
View File
@@ -0,0 +1,176 @@
package server
import (
"fmt"
"strings"
"github.com/labring/sealos/service/pkg/api"
)
const (
modeTrue = "true"
modeFalse = "false"
)
type VLogsQuery struct {
query string
}
func (v *VLogsQuery) getQuery(req *api.VlogsRequest) (string, error) {
if req.PodQuery == modeTrue {
query := v.generatePodListQuery(req)
return query, nil
}
v.generateKeywordQuery(req)
v.generateStreamQuery(req)
v.generateCommonQuery(req)
err := v.generateJSONQuery(req)
if err != nil {
return "", err
}
v.generateDropQuery()
v.generateNumberQuery(req)
return v.query, nil
}
func (v *VLogsQuery) generatePodListQuery(req *api.VlogsRequest) string {
var item string
if len(req.Time) != 0 {
item = fmt.Sprintf(
`{namespace="%s"} _time:%s app:="%s" | Drop _stream_id,_stream,app,job,namespace,node`,
req.Namespace,
req.Time,
req.App,
)
} else {
item = fmt.Sprintf(`{namespace="%s"} app:="%s" | Drop _stream_id,_stream,app,job,namespace,node`, req.Namespace, req.App)
}
v.query += item
return v.query
}
func (v *VLogsQuery) generateKeywordQuery(req *api.VlogsRequest) {
v.query += fmt.Sprintf("%s ", req.Keyword)
}
func (v *VLogsQuery) generateJSONQuery(req *api.VlogsRequest) error {
if req.JSONMode != modeTrue {
return nil
}
var builder strings.Builder
builder.WriteString(" | unpack_json")
if len(req.JSONQuery) > 0 {
for _, jsonQuery := range req.JSONQuery {
var item string
switch jsonQuery.Mode {
case "=":
item = fmt.Sprintf("| %s:=%s ", jsonQuery.Key, jsonQuery.Value)
case "!=":
item = fmt.Sprintf("| %s:(!=%s) ", jsonQuery.Key, jsonQuery.Value)
case "~":
item = fmt.Sprintf("| %s:%s ", jsonQuery.Key, jsonQuery.Value)
case "!~":
item = fmt.Sprintf("| %s:(!~%s) ", jsonQuery.Key, jsonQuery.Value)
default:
return fmt.Errorf("invalid JSON query mode: %s", jsonQuery.Mode)
}
builder.WriteString(item)
}
}
v.query += builder.String()
return nil
}
func (v *VLogsQuery) generateStreamQuery(req *api.VlogsRequest) {
var builder strings.Builder
switch {
case len(req.Pod) == 0 && len(req.Container) == 0:
// Generate query based only on namespace
builder.WriteString(fmt.Sprintf(`{namespace="%s"}`, req.Namespace))
case len(req.Pod) == 0:
// Generate query based on container
for i, container := range req.Container {
builder.WriteString(
fmt.Sprintf(`{container="%s",namespace="%s"}`, container, req.Namespace),
)
if i != len(req.Container)-1 {
builder.WriteString(" OR ")
}
}
case len(req.Container) == 0:
// Generate query based on pod
for i, pod := range req.Pod {
builder.WriteString(fmt.Sprintf(`{pod="%s",namespace="%s"}`, pod, req.Namespace))
if i != len(req.Pod)-1 {
builder.WriteString(" OR ")
}
}
default:
// Generate query based on both pod and container
for i, container := range req.Container {
for j, pod := range req.Pod {
builder.WriteString(
fmt.Sprintf(
`{container="%s",namespace="%s",pod="%s"}`,
container,
req.Namespace,
pod,
),
)
if i != len(req.Container)-1 || j != len(req.Pod)-1 {
builder.WriteString(" OR ")
}
}
}
}
v.query += builder.String()
}
func (v *VLogsQuery) generateCommonQuery(req *api.VlogsRequest) {
var builder strings.Builder
var item string
if len(req.Time) != 0 {
item = fmt.Sprintf(`_time:%s app:="%s" `, req.Time, req.App)
} else {
item = fmt.Sprintf(`app:="%s" `, req.App)
}
builder.WriteString(item)
// if query stderr and number,using stderr first.
if req.StderrMode == modeTrue {
item := `| stream:="stderr" `
builder.WriteString(item)
}
// if query number,dont use limit param
if req.NumberMode == modeFalse {
item := fmt.Sprintf(` | limit %s `, req.Limit)
builder.WriteString(item)
}
v.query += builder.String()
}
func (v *VLogsQuery) generateDropQuery() {
v.query += "| Drop _stream_id,_stream,app,job,namespace,node"
}
// allowedNumberLevels defines the set of valid NumberLevel values.
var allowedNumberLevels = map[string]struct{}{
"m": {},
"h": {},
"d": {},
"s": {},
}
func isValidNumberLevel(level string) bool {
_, ok := allowedNumberLevels[level]
return ok
}
func (v *VLogsQuery) generateNumberQuery(req *api.VlogsRequest) {
if req.NumberMode == modeTrue {
if isValidNumberLevel(req.NumberLevel) {
item := fmt.Sprintf(" | stats by (_time:1%s) count() logs_total ", req.NumberLevel)
v.query += item
}
// else: invalid NumberLevel, do not add to query
}
}
+101 -209
View File
@@ -2,13 +2,14 @@ package server
import (
"bufio"
"bytes"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net/http"
"net/url"
"strings"
"github.com/labring/sealos/service/pkg/api"
"github.com/labring/sealos/service/pkg/auth"
@@ -21,9 +22,6 @@ type VLogsServer struct {
password string
}
const modeTrue = "true"
const modeFalse = "false"
func NewVLogsServer(config *Config) (*VLogsServer, error) {
vl := &VLogsServer{
path: config.Server.Path,
@@ -36,52 +34,108 @@ func NewVLogsServer(config *Config) (*VLogsServer, error) {
func (vl *VLogsServer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
query, err := vl.queryConvert(req)
if err != nil {
http.Error(rw, fmt.Sprintf("query %s error: %s", req.URL.Path, err), http.StatusInternalServerError)
http.Error(
rw,
fmt.Sprintf("query %s error: %s", req.URL.Path, err),
http.StatusInternalServerError,
)
return
}
err = query(rw, req)
if err != nil {
http.Error(rw, fmt.Sprintf("query %s error: %s", req.URL.Path, err), http.StatusInternalServerError)
http.Error(
rw,
fmt.Sprintf("query %s error: %s", req.URL.Path, err),
http.StatusInternalServerError,
)
slog.Error("%s error: %s", req.URL.Path, err)
return
}
}
func (vl *VLogsServer) queryConvert(req *http.Request) (func(rw http.ResponseWriter, req *http.Request) error, error) {
func (vl *VLogsServer) queryConvert(
req *http.Request,
) (func(rw http.ResponseWriter, req *http.Request) error, error) {
switch req.URL.Path {
case "/queryLogsByParams":
return vl.queryLogsByParams, nil
case "/queryPodList":
return vl.queryPodList, nil
default:
return nil, fmt.Errorf("unknown url path")
return nil, errors.New("unknown url path")
}
}
func (vl *VLogsServer) authenticate(req *http.Request) (string, error) {
kubeConfig, namespace, query, err := vl.generateParamsRequest(req)
if err != nil {
return "", fmt.Errorf("bad request (%s)", err)
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
}
err = auth.Authenticate(namespace, kubeConfig)
func (vl *VLogsServer) authenticate(namespace, kubeConfig string) error {
err := auth.Authenticate(namespace, kubeConfig)
if err != nil {
return "", fmt.Errorf("authentication failed (%s)", err)
return fmt.Errorf("authentication failed (%w)", err)
}
return query, nil
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 {
query, err := vl.authenticate(req)
resp, err := vl.executeQuery(req)
if err != nil {
return err
}
resp, err := request.QueryLogsByParams(vl.path, vl.username, vl.password, query)
if err != nil {
return fmt.Errorf("query failed (%s)", err)
}
defer resp.Body.Close()
_, err = io.Copy(rw, resp.Body)
defer resp.Close()
_, err = io.Copy(rw, resp)
if err != nil {
return err
}
@@ -89,209 +143,47 @@ func (vl *VLogsServer) queryLogsByParams(rw http.ResponseWriter, req *http.Reque
}
func (vl *VLogsServer) queryPodList(rw http.ResponseWriter, req *http.Request) error {
query, err := vl.authenticate(req)
resp, err := vl.executeQuery(req)
if err != nil {
return err
}
resp, err := request.QueryLogsByParams(vl.path, vl.username, vl.password, query)
defer resp.Close()
body, err := io.ReadAll(resp)
if err != nil {
return fmt.Errorf("query failed (%s)", err)
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return fmt.Errorf("failed to read response body: %v", err)
return fmt.Errorf("failed to read response body: %w", err)
}
if len(body) == 0 {
return fmt.Errorf("response body is empty")
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
}
scanner := bufio.NewScanner(strings.NewReader(string(body)))
var logs []api.VlogsResponse
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
line := scanner.Text()
err := json.Unmarshal([]byte(line), &entry)
lineBytes := scanner.Bytes()
err := json.Unmarshal(lineBytes, &entry)
if err != nil {
continue
}
logs = append(logs, entry)
uniquePods[entry.Pod] = struct{}{}
}
uniquePods := make(map[string]struct{})
for _, log := range logs {
uniquePods[log.Pod] = struct{}{}
if err := scanner.Err(); err != nil {
return nil, fmt.Errorf("error reading response: %w", err)
}
var podList []string
podList := make([]string, 0, len(uniquePods))
for pod := range uniquePods {
podList = append(podList, pod)
}
rw.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(rw).Encode(podList); err != nil {
return fmt.Errorf("failed to write response: %v", err)
}
return nil
}
func (vl *VLogsServer) generateParamsRequest(req *http.Request) (string, string, string, error) {
kubeConfig := req.Header.Get("Authorization")
if config, err := url.PathUnescape(kubeConfig); err == nil {
kubeConfig = config
} else {
return "", "", "", fmt.Errorf("failed to PathUnescape : %s", err)
}
var query string
vlogsReq := &api.VlogsRequest{}
err := json.NewDecoder(req.Body).Decode(&vlogsReq)
if err != nil {
return "", "", "", fmt.Errorf("failed to parse request body: %s", err)
}
if vlogsReq.Namespace == "" {
return "", "", "", fmt.Errorf("failed to get namespace")
}
var vlogs VLogsQuery
query, err = vlogs.getQuery(vlogsReq)
if err != nil {
return "", "", "", fmt.Errorf("failed to parse request body: %s", err)
}
return kubeConfig, vlogsReq.Namespace, query, nil
}
type VLogsQuery struct {
query string
}
func (v *VLogsQuery) getQuery(req *api.VlogsRequest) (string, error) {
if req.PodQuery == modeTrue {
query := v.generatePodListQuery(req)
return query, nil
}
v.generateKeywordQuery(req)
v.generateStreamQuery(req)
v.generateCommonQuery(req)
err := v.generateJSONQuery(req)
if err != nil {
return "", err
}
v.generateStdQuery(req)
v.generateDropQuery()
v.generateNumberQuery(req)
return v.query, nil
}
func (v *VLogsQuery) generatePodListQuery(req *api.VlogsRequest) string {
var builder strings.Builder
item := fmt.Sprintf(`{namespace="%s"} _time:%s app:="%s" | Drop _stream_id,_stream,app,job,namespace,node`, req.Namespace, req.Time, req.App)
builder.WriteString(item)
v.query += builder.String()
return v.query
}
func (v *VLogsQuery) generateKeywordQuery(req *api.VlogsRequest) {
var builder strings.Builder
builder.WriteString(req.Keyword)
builder.WriteString(" ")
v.query += builder.String()
}
func (v *VLogsQuery) generateJSONQuery(req *api.VlogsRequest) error {
if req.JSONMode != modeTrue {
return nil
}
var builder strings.Builder
builder.WriteString(" | unpack_json")
if len(req.JSONQuery) > 0 {
for _, jsonQuery := range req.JSONQuery {
var item string
switch jsonQuery.Mode {
case "=":
item = fmt.Sprintf("| %s:=%s ", jsonQuery.Key, jsonQuery.Value)
case "!=":
item = fmt.Sprintf("| %s:(!=%s) ", jsonQuery.Key, jsonQuery.Value)
case "~":
item = fmt.Sprintf("| %s:%s ", jsonQuery.Key, jsonQuery.Value)
case "!~":
item = fmt.Sprintf("| %s:(!~%s) ", jsonQuery.Key, jsonQuery.Value)
default:
return fmt.Errorf("invalid JSON query mode: %s", jsonQuery.Mode)
}
builder.WriteString(item)
}
}
v.query += builder.String()
return nil
}
func (v *VLogsQuery) generateStreamQuery(req *api.VlogsRequest) {
var builder strings.Builder
if len(req.Pod) == 0 && len(req.Container) == 0 {
// Generate query based only on namespace
builder.WriteString(fmt.Sprintf(`{namespace="%s"}`, req.Namespace))
} else if len(req.Pod) == 0 {
// Generate query based on container
for i, container := range req.Container {
builder.WriteString(fmt.Sprintf(`{container="%s",namespace="%s"}`, container, req.Namespace))
if i != len(req.Container)-1 {
builder.WriteString(" OR ")
}
}
} else if len(req.Container) == 0 {
// Generate query based on pod
for i, pod := range req.Pod {
builder.WriteString(fmt.Sprintf(`{pod="%s",namespace="%s"}`, pod, req.Namespace))
if i != len(req.Pod)-1 {
builder.WriteString(" OR ")
}
}
} else {
// Generate query based on both pod and container
for i, container := range req.Container {
for j, pod := range req.Pod {
builder.WriteString(fmt.Sprintf(`{container="%s",namespace="%s",pod="%s"}`, container, req.Namespace, pod))
if i != len(req.Container)-1 || j != len(req.Pod)-1 {
builder.WriteString(" OR ")
}
}
}
}
v.query += builder.String()
}
func (v *VLogsQuery) generateStdQuery(req *api.VlogsRequest) {
var builder strings.Builder
if req.StderrMode == modeTrue {
item := `| stream:="stderr" `
builder.WriteString(item)
}
v.query += builder.String()
}
func (v *VLogsQuery) generateCommonQuery(req *api.VlogsRequest) {
var builder strings.Builder
item := fmt.Sprintf(`_time:%s app:="%s" `, req.Time, req.App)
builder.WriteString(item)
// if query number,dont use limit param
if req.NumberMode == modeFalse {
item := fmt.Sprintf(` | limit %s `, req.Limit)
builder.WriteString(item)
}
v.query += builder.String()
}
func (v *VLogsQuery) generateDropQuery() {
var builder strings.Builder
builder.WriteString("| Drop _stream_id,_stream,app,job,namespace,node")
v.query += builder.String()
}
func (v *VLogsQuery) generateNumberQuery(req *api.VlogsRequest) {
var builder strings.Builder
if req.NumberMode == modeTrue {
item := fmt.Sprintf(" | stats by (_time:1%s) count() logs_total ", req.NumberLevel)
builder.WriteString(item)
v.query += builder.String()
}
return podList, nil
}