mirror of
https://github.com/wavetermdev/backup.git
synced 2026-08-05 13:57:07 -07:00
client support for rpc cancel (#249)
This commit is contained in:
@@ -57,27 +57,30 @@ func sendRpcRequestResponseStreamHelper[T any](w *wshutil.WshRpc, command string
|
||||
if err != nil {
|
||||
rtnErr(respChan, err)
|
||||
return respChan
|
||||
} else {
|
||||
go func() {
|
||||
defer close(respChan)
|
||||
for {
|
||||
if reqHandler.ResponseDone() {
|
||||
break
|
||||
}
|
||||
resp, err := reqHandler.NextResponse()
|
||||
if err != nil {
|
||||
respChan <- wshrpc.RespOrErrorUnion[T]{Error: err}
|
||||
break
|
||||
}
|
||||
var respData T
|
||||
err = utilfn.ReUnmarshal(&respData, resp)
|
||||
if err != nil {
|
||||
respChan <- wshrpc.RespOrErrorUnion[T]{Error: err}
|
||||
break
|
||||
}
|
||||
respChan <- wshrpc.RespOrErrorUnion[T]{Response: respData}
|
||||
}
|
||||
}()
|
||||
}
|
||||
opts.StreamCancelFn = func() {
|
||||
// TODO coordinate the cancel with the for loop below
|
||||
reqHandler.SendCancel()
|
||||
}
|
||||
go func() {
|
||||
defer close(respChan)
|
||||
for {
|
||||
if reqHandler.ResponseDone() {
|
||||
break
|
||||
}
|
||||
resp, err := reqHandler.NextResponse()
|
||||
if err != nil {
|
||||
respChan <- wshrpc.RespOrErrorUnion[T]{Error: err}
|
||||
break
|
||||
}
|
||||
var respData T
|
||||
err = utilfn.ReUnmarshal(&respData, resp)
|
||||
if err != nil {
|
||||
respChan <- wshrpc.RespOrErrorUnion[T]{Error: err}
|
||||
break
|
||||
}
|
||||
respChan <- wshrpc.RespOrErrorUnion[T]{Response: respData}
|
||||
}
|
||||
}()
|
||||
return respChan
|
||||
}
|
||||
|
||||
@@ -105,6 +105,8 @@ type RpcOpts struct {
|
||||
Timeout int `json:"timeout,omitempty"`
|
||||
NoResponse bool `json:"noresponse,omitempty"`
|
||||
Route string `json:"route,omitempty"`
|
||||
|
||||
StreamCancelFn func() `json:"-"` // this is an *output* parameter, set by the handler
|
||||
}
|
||||
|
||||
type RpcContext struct {
|
||||
|
||||
Reference in New Issue
Block a user