mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
feat(rpc2): audit log on error handler
. #22617142 Signed-off-by: slasher <shenjie1@oppo.com>
This commit is contained in:
parent
0612bfc8ff
commit
fc32992b0b
@ -15,61 +15,94 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"os"
|
||||
"reflect"
|
||||
|
||||
"github.com/desertbit/grumble"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/api/shardnode"
|
||||
"github.com/cubefs/cubefs/blobstore/cli/common"
|
||||
"github.com/cubefs/cubefs/blobstore/cli/common/fmt"
|
||||
"github.com/cubefs/cubefs/blobstore/cli/config"
|
||||
"github.com/cubefs/cubefs/blobstore/common/rpc2"
|
||||
)
|
||||
|
||||
type (
|
||||
sCodec = rpc2.AnyCodec[struct{ val string }]
|
||||
bCodec = rpc2.AnyCodec[struct{ val []byte }]
|
||||
)
|
||||
var types = map[string]func() rpc2.Codec{
|
||||
"shardnode.GetBlobArgs": func() rpc2.Codec { return new(shardnode.GetBlobArgs) },
|
||||
"shardnode.GetBlobRet": func() rpc2.Codec { return new(shardnode.GetBlobRet) },
|
||||
|
||||
"nil": func() rpc2.Codec { return nil },
|
||||
"rpc2.NoParameter": func() rpc2.Codec { return rpc2.NoParameter },
|
||||
}
|
||||
|
||||
func helpType(t any) string {
|
||||
if t == nil {
|
||||
return "nil"
|
||||
}
|
||||
typ := reflect.TypeOf(t)
|
||||
if typ.Kind() == reflect.Ptr {
|
||||
typ = typ.Elem()
|
||||
}
|
||||
if typ.Kind() != reflect.Struct {
|
||||
panic(typ.Name())
|
||||
}
|
||||
return typ.PkgPath() + "." + typ.Name()
|
||||
}
|
||||
|
||||
func helps() (s string) {
|
||||
for n, ft := range types {
|
||||
s += "\n" + n + " -> " + helpType(ft())
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func cmdRpc2Request(c *grumble.Context) error {
|
||||
cli := config.Rpc2Client
|
||||
addr, path := c.Args.String("addr"), c.Args.String("path")
|
||||
paraVal, rstVal := c.Flags.String("parameter"), c.Flags.String("result")
|
||||
addr, path, request := c.Args.String("addr"), c.Args.String("path"), c.Args.String("request")
|
||||
paraType, rstType := c.Flags.String("parameter"), c.Flags.String("result")
|
||||
|
||||
if c.Flags.Bool("readable") {
|
||||
var para, rst sCodec
|
||||
para.Value.val = paraVal
|
||||
if err := cli.Request(common.CmdContext(), addr, path, ¶, &rst); err != nil {
|
||||
paraf, exist := types[paraType]
|
||||
if !exist {
|
||||
return fmt.Errorf("not found parameter(%s), types:%s", paraType, helps())
|
||||
}
|
||||
para := paraf()
|
||||
if para != nil && para != rpc2.NoParameter {
|
||||
if err := common.Unmarshal([]byte(request), para); err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Println("result:", rst.Value.val)
|
||||
return nil
|
||||
}
|
||||
if s, ok := para.(interface{ GoString() string }); ok {
|
||||
fmt.Println("parameter: " + s.GoString())
|
||||
} else if s, ok := para.(interface{ String() string }); ok {
|
||||
fmt.Println("parameter: " + s.String())
|
||||
} else {
|
||||
fmt.Println("parameter: " + common.Readable(para))
|
||||
}
|
||||
|
||||
val, err := os.ReadFile(paraVal)
|
||||
if err != nil {
|
||||
rstf, exist := types[rstType]
|
||||
if !exist {
|
||||
return fmt.Errorf("not found result(%s), types:%s", rstType, helps())
|
||||
}
|
||||
rst := rstf()
|
||||
if err := cli.Request(common.CmdContext(), addr, path, para, rst); err != nil {
|
||||
return err
|
||||
}
|
||||
var para, rst bCodec
|
||||
para.Value.val = val
|
||||
if err := cli.Request(common.CmdContext(), addr, path, ¶, &rst); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.WriteFile(rstVal, rst.Value.val, 0o644)
|
||||
fmt.Println("result: ", common.Readable(rst))
|
||||
return nil
|
||||
}
|
||||
|
||||
func registerRpc2(app *grumble.App) {
|
||||
rpc2Command := &grumble.Command{
|
||||
Name: "rpc2",
|
||||
Help: "simple client of rpc2",
|
||||
LongHelp: "rpc2 simple client struct request and response",
|
||||
LongHelp: "rpc2 simple client struct request and response, types:\n" + helps(),
|
||||
Args: func(a *grumble.Args) {
|
||||
a.String("addr", "request address")
|
||||
a.String("path", "request path")
|
||||
a.String("request", "request parameter json value")
|
||||
},
|
||||
Flags: func(f *grumble.Flags) {
|
||||
f.BoolL("readable", false, "para and args is readable")
|
||||
f.StringL("parameter", "", "request parameter")
|
||||
f.StringL("result", "", "response result")
|
||||
f.StringL("parameter", "", "request parameter name")
|
||||
f.StringL("result", "", "response result name")
|
||||
},
|
||||
Run: cmdRpc2Request,
|
||||
}
|
||||
|
||||
@ -28,24 +28,30 @@ func runClient() {
|
||||
},
|
||||
Timeout: util.Duration{Duration: time.Second},
|
||||
}
|
||||
ctx := context.Background()
|
||||
getAddr := func() string { return listenon[int(time.Now().UnixNano())%len(listenon)] }
|
||||
{
|
||||
var para paraCodec
|
||||
para.Value = message{I: 7, S: "ping string"}
|
||||
log.Infof("before request para : %+v", para)
|
||||
if err := client.Request(context.Background(),
|
||||
listenon[int(time.Now().UnixNano())%len(listenon)],
|
||||
"/kick", ¶, ¶); err != nil {
|
||||
err := client.Request(ctx, getAddr(), "/kick", ¶, nil)
|
||||
if err != nil {
|
||||
panic(rpc2.ErrorString(err))
|
||||
}
|
||||
err = client.Request(ctx, getAddr(), "/error", ¶, nil)
|
||||
if st, _, _ := rpc2.DetectError(err); st != 567 {
|
||||
panic(rpc2.ErrorString(err))
|
||||
}
|
||||
err = client.Request(ctx, getAddr(), "/panic", ¶, nil)
|
||||
if st, _, _ := rpc2.DetectError(err); st != rpc2.DefaultStatusPanic {
|
||||
panic(rpc2.ErrorString(err))
|
||||
}
|
||||
log.Infof("after request result : %+v", para)
|
||||
}
|
||||
|
||||
var para paraCodec
|
||||
para.Value = message{I: 7, S: "ping string"}
|
||||
buff := []byte("ping")
|
||||
req, _ := rpc2.NewRequest(context.Background(),
|
||||
listenon[int(time.Now().UnixNano())%len(listenon)],
|
||||
"/ping", ¶, bytes.NewReader(buff))
|
||||
req, _ := rpc2.NewRequest(ctx, getAddr(), "/ping", ¶, bytes.NewReader(buff))
|
||||
req.Trailer.SetLen("trailer-1", 1)
|
||||
req.AfterBody = func() error {
|
||||
req.Trailer.Set("trailer-1", "xX")
|
||||
|
||||
@ -58,8 +58,17 @@ func setUp2() (*rpc2.Router, []rpc2.Interceptor) {
|
||||
router.Middleware(handleMiddleware1, handleMiddleware2)
|
||||
router.Register("/ping", handlePing)
|
||||
router.Register("/kick", handleKick)
|
||||
router.Register("/error", handleError)
|
||||
router.Register("/panic", handlePanic)
|
||||
router.Register("/stream", handleStream)
|
||||
return router, nil
|
||||
return router, []rpc2.Interceptor{interceptor{"i1"}, interceptor{"i2"}}
|
||||
}
|
||||
|
||||
type interceptor struct{ id string }
|
||||
|
||||
func (i interceptor) Handle(w rpc2.ResponseWriter, req *rpc2.Request, h rpc2.Handle) error {
|
||||
log.Info("interceptor-" + i.id)
|
||||
return h(w, req)
|
||||
}
|
||||
|
||||
func handleMiddleware1(w rpc2.ResponseWriter, req *rpc2.Request) error {
|
||||
@ -72,6 +81,20 @@ func handleMiddleware2(w rpc2.ResponseWriter, req *rpc2.Request) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func handleKick(_ rpc2.ResponseWriter, req *rpc2.Request) error {
|
||||
var para paraCodec
|
||||
req.ParseParameter(¶)
|
||||
return nil
|
||||
}
|
||||
|
||||
func handleError(rpc2.ResponseWriter, *rpc2.Request) error {
|
||||
return rpc2.NewError(567, "", "")
|
||||
}
|
||||
|
||||
func handlePanic(rpc2.ResponseWriter, *rpc2.Request) error {
|
||||
panic("handle panic")
|
||||
}
|
||||
|
||||
func handlePing(w rpc2.ResponseWriter, req *rpc2.Request) error {
|
||||
log.Info(req.RequestHeader.String())
|
||||
var para paraCodec
|
||||
@ -105,13 +128,6 @@ func handlePing(w rpc2.ResponseWriter, req *rpc2.Request) error {
|
||||
return err
|
||||
}
|
||||
|
||||
func handleKick(w rpc2.ResponseWriter, req *rpc2.Request) error {
|
||||
var para paraCodec
|
||||
req.ParseParameter(¶)
|
||||
para.Value.S = "response -> " + para.Value.S
|
||||
return w.WriteOK(¶)
|
||||
}
|
||||
|
||||
func handleStream(_ rpc2.ResponseWriter, req *rpc2.Request) error {
|
||||
var para paraCodec
|
||||
req.ParseParameter(¶)
|
||||
|
||||
@ -32,6 +32,8 @@ type ResponseWriter interface {
|
||||
WriteHeader(status int, obj Marshaler) error
|
||||
// WriteOK object in body
|
||||
WriteOK(obj Marshaler) error
|
||||
// SetError fill error's reason to response header
|
||||
SetError(err error)
|
||||
Flush() error
|
||||
// io.Writer
|
||||
io.ReaderFrom
|
||||
@ -99,6 +101,12 @@ func (resp *response) Trailer() *FixedHeader {
|
||||
return &resp.hdr.Trailer
|
||||
}
|
||||
|
||||
func (resp *response) SetError(err error) {
|
||||
_, reason, detail := DetectError(err)
|
||||
resp.hdr.Reason = reason
|
||||
resp.hdr.Error = detail.Error()
|
||||
}
|
||||
|
||||
func (resp *response) WriteOK(obj Marshaler) error {
|
||||
if resp.hasWroteHeader {
|
||||
return nil
|
||||
|
||||
@ -135,5 +135,20 @@ func (r *Router) handle(w ResponseWriter, req *Request) (err error) {
|
||||
}
|
||||
}
|
||||
err = handle(w, req)
|
||||
if req.stream != nil { // stream
|
||||
return
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
status, _, _ := DetectError(err)
|
||||
w.SetError(err)
|
||||
w.WriteHeader(status, NoParameter)
|
||||
}
|
||||
if err = w.WriteOK(nil); err != nil {
|
||||
return
|
||||
}
|
||||
if err = w.Flush(); err != nil {
|
||||
return
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue
Block a user