2017-02-10 01:58:20 +00:00
|
|
|
package agent
|
|
|
|
|
|
|
|
import (
|
2020-05-21 14:05:04 +00:00
|
|
|
"context"
|
|
|
|
"io"
|
|
|
|
"net"
|
2017-02-10 01:58:20 +00:00
|
|
|
"net/http"
|
|
|
|
"strings"
|
|
|
|
|
2017-12-18 21:16:23 +00:00
|
|
|
"fmt"
|
|
|
|
"strconv"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
"github.com/hashicorp/consul/agent/consul/autopilot"
|
2020-05-21 14:05:04 +00:00
|
|
|
"github.com/hashicorp/go-msgpack/codec"
|
2017-12-18 21:16:23 +00:00
|
|
|
"github.com/hashicorp/nomad/api"
|
2020-06-04 21:14:43 +00:00
|
|
|
cstructs "github.com/hashicorp/nomad/client/structs"
|
2017-02-10 01:58:20 +00:00
|
|
|
"github.com/hashicorp/nomad/nomad/structs"
|
|
|
|
"github.com/hashicorp/raft"
|
|
|
|
)
|
|
|
|
|
2020-04-03 23:03:14 +00:00
|
|
|
// OperatorRequest is used route operator/raft API requests to the implementing
|
|
|
|
// functions.
|
2017-02-10 01:58:20 +00:00
|
|
|
func (s *HTTPServer) OperatorRequest(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
path := strings.TrimPrefix(req.URL.Path, "/v1/operator/raft/")
|
|
|
|
switch {
|
|
|
|
case strings.HasPrefix(path, "configuration"):
|
|
|
|
return s.OperatorRaftConfiguration(resp, req)
|
|
|
|
case strings.HasPrefix(path, "peer"):
|
|
|
|
return s.OperatorRaftPeer(resp, req)
|
|
|
|
default:
|
|
|
|
return nil, CodedError(404, ErrInvalidMethod)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// OperatorRaftConfiguration is used to inspect the current Raft configuration.
|
|
|
|
// This supports the stale query mode in case the cluster doesn't have a leader.
|
|
|
|
func (s *HTTPServer) OperatorRaftConfiguration(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
if req.Method != "GET" {
|
|
|
|
resp.WriteHeader(http.StatusMethodNotAllowed)
|
|
|
|
return nil, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
var args structs.GenericRequest
|
|
|
|
if done := s.parse(resp, req, &args.Region, &args.QueryOptions); done {
|
|
|
|
return nil, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
var reply structs.RaftConfigurationResponse
|
|
|
|
if err := s.agent.RPC("Operator.RaftGetConfiguration", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
return reply, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// OperatorRaftPeer supports actions on Raft peers. Currently we only support
|
|
|
|
// removing peers by address.
|
|
|
|
func (s *HTTPServer) OperatorRaftPeer(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
if req.Method != "DELETE" {
|
2018-02-08 00:47:44 +00:00
|
|
|
return nil, CodedError(404, ErrInvalidMethod)
|
2017-02-10 01:58:20 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
params := req.URL.Query()
|
2018-01-16 21:35:32 +00:00
|
|
|
_, hasID := params["id"]
|
|
|
|
_, hasAddress := params["address"]
|
|
|
|
|
|
|
|
if !hasID && !hasAddress {
|
2018-02-08 00:47:44 +00:00
|
|
|
return nil, CodedError(http.StatusBadRequest, "Must specify either ?id with the server's ID or ?address with IP:port of peer to remove")
|
2018-01-16 21:35:32 +00:00
|
|
|
}
|
|
|
|
if hasID && hasAddress {
|
2018-02-08 00:47:44 +00:00
|
|
|
return nil, CodedError(http.StatusBadRequest, "Must specify only one of ?id or ?address")
|
2017-02-10 01:58:20 +00:00
|
|
|
}
|
|
|
|
|
2018-01-16 21:35:32 +00:00
|
|
|
if hasID {
|
|
|
|
var args structs.RaftPeerByIDRequest
|
|
|
|
s.parseWriteRequest(req, &args.WriteRequest)
|
|
|
|
|
|
|
|
var reply struct{}
|
|
|
|
args.ID = raft.ServerID(params.Get("id"))
|
|
|
|
if err := s.agent.RPC("Operator.RaftRemovePeerByID", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
var args structs.RaftPeerByAddressRequest
|
|
|
|
s.parseWriteRequest(req, &args.WriteRequest)
|
|
|
|
|
|
|
|
var reply struct{}
|
|
|
|
args.Address = raft.ServerAddress(params.Get("address"))
|
|
|
|
if err := s.agent.RPC("Operator.RaftRemovePeerByAddress", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2017-02-10 01:58:20 +00:00
|
|
|
}
|
2018-01-16 21:35:32 +00:00
|
|
|
|
2017-02-10 01:58:20 +00:00
|
|
|
return nil, nil
|
|
|
|
}
|
2017-12-18 21:16:23 +00:00
|
|
|
|
|
|
|
// OperatorAutopilotConfiguration is used to inspect the current Autopilot configuration.
|
|
|
|
// This supports the stale query mode in case the cluster doesn't have a leader.
|
|
|
|
func (s *HTTPServer) OperatorAutopilotConfiguration(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
// Switch on the method
|
|
|
|
switch req.Method {
|
|
|
|
case "GET":
|
|
|
|
var args structs.GenericRequest
|
|
|
|
if done := s.parse(resp, req, &args.Region, &args.QueryOptions); done {
|
|
|
|
return nil, nil
|
|
|
|
}
|
|
|
|
|
2018-01-30 03:53:34 +00:00
|
|
|
var reply structs.AutopilotConfig
|
2017-12-18 21:16:23 +00:00
|
|
|
if err := s.agent.RPC("Operator.AutopilotGetConfiguration", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
out := api.AutopilotConfiguration{
|
|
|
|
CleanupDeadServers: reply.CleanupDeadServers,
|
2018-01-30 03:53:34 +00:00
|
|
|
LastContactThreshold: reply.LastContactThreshold,
|
2017-12-18 21:16:23 +00:00
|
|
|
MaxTrailingLogs: reply.MaxTrailingLogs,
|
2020-02-16 21:23:20 +00:00
|
|
|
MinQuorum: reply.MinQuorum,
|
2018-01-30 03:53:34 +00:00
|
|
|
ServerStabilizationTime: reply.ServerStabilizationTime,
|
|
|
|
EnableRedundancyZones: reply.EnableRedundancyZones,
|
2017-12-18 21:16:23 +00:00
|
|
|
DisableUpgradeMigration: reply.DisableUpgradeMigration,
|
2018-01-30 03:53:34 +00:00
|
|
|
EnableCustomUpgrades: reply.EnableCustomUpgrades,
|
2017-12-18 21:16:23 +00:00
|
|
|
CreateIndex: reply.CreateIndex,
|
|
|
|
ModifyIndex: reply.ModifyIndex,
|
|
|
|
}
|
|
|
|
|
|
|
|
return out, nil
|
|
|
|
|
|
|
|
case "PUT":
|
|
|
|
var args structs.AutopilotSetConfigRequest
|
2018-02-08 00:47:44 +00:00
|
|
|
s.parseWriteRequest(req, &args.WriteRequest)
|
2017-12-18 21:16:23 +00:00
|
|
|
|
|
|
|
var conf api.AutopilotConfiguration
|
2018-01-30 03:53:34 +00:00
|
|
|
if err := decodeBody(req, &conf); err != nil {
|
2018-02-08 00:47:44 +00:00
|
|
|
return nil, CodedError(http.StatusBadRequest, fmt.Sprintf("Error parsing autopilot config: %v", err))
|
2017-12-18 21:16:23 +00:00
|
|
|
}
|
|
|
|
|
2018-01-30 03:53:34 +00:00
|
|
|
args.Config = structs.AutopilotConfig{
|
2017-12-18 21:16:23 +00:00
|
|
|
CleanupDeadServers: conf.CleanupDeadServers,
|
2018-01-30 03:53:34 +00:00
|
|
|
LastContactThreshold: conf.LastContactThreshold,
|
2017-12-18 21:16:23 +00:00
|
|
|
MaxTrailingLogs: conf.MaxTrailingLogs,
|
2020-02-16 21:23:20 +00:00
|
|
|
MinQuorum: conf.MinQuorum,
|
2018-01-30 03:53:34 +00:00
|
|
|
ServerStabilizationTime: conf.ServerStabilizationTime,
|
|
|
|
EnableRedundancyZones: conf.EnableRedundancyZones,
|
2017-12-18 21:16:23 +00:00
|
|
|
DisableUpgradeMigration: conf.DisableUpgradeMigration,
|
2018-01-30 03:53:34 +00:00
|
|
|
EnableCustomUpgrades: conf.EnableCustomUpgrades,
|
2017-12-18 21:16:23 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// Check for cas value
|
|
|
|
params := req.URL.Query()
|
|
|
|
if _, ok := params["cas"]; ok {
|
|
|
|
casVal, err := strconv.ParseUint(params.Get("cas"), 10, 64)
|
|
|
|
if err != nil {
|
2018-02-08 00:47:44 +00:00
|
|
|
return nil, CodedError(http.StatusBadRequest, fmt.Sprintf("Error parsing cas value: %v", err))
|
2017-12-18 21:16:23 +00:00
|
|
|
}
|
|
|
|
args.Config.ModifyIndex = casVal
|
|
|
|
args.CAS = true
|
|
|
|
}
|
|
|
|
|
|
|
|
var reply bool
|
|
|
|
if err := s.agent.RPC("Operator.AutopilotSetConfiguration", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
// Only use the out value if this was a CAS
|
|
|
|
if !args.CAS {
|
|
|
|
return true, nil
|
|
|
|
}
|
|
|
|
return reply, nil
|
|
|
|
|
|
|
|
default:
|
2018-02-08 00:47:44 +00:00
|
|
|
return nil, CodedError(404, ErrInvalidMethod)
|
2017-12-18 21:16:23 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2018-01-16 21:35:32 +00:00
|
|
|
// OperatorServerHealth is used to get the health of the servers in the given Region.
|
2017-12-18 21:16:23 +00:00
|
|
|
func (s *HTTPServer) OperatorServerHealth(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
if req.Method != "GET" {
|
2018-02-08 00:47:44 +00:00
|
|
|
return nil, CodedError(404, ErrInvalidMethod)
|
2017-12-18 21:16:23 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
var args structs.GenericRequest
|
|
|
|
if done := s.parse(resp, req, &args.Region, &args.QueryOptions); done {
|
|
|
|
return nil, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
var reply autopilot.OperatorHealthReply
|
|
|
|
if err := s.agent.RPC("Operator.ServerHealth", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
// Reply with status 429 if something is unhealthy
|
|
|
|
if !reply.Healthy {
|
|
|
|
resp.WriteHeader(http.StatusTooManyRequests)
|
|
|
|
}
|
|
|
|
|
|
|
|
out := &api.OperatorHealthReply{
|
|
|
|
Healthy: reply.Healthy,
|
|
|
|
FailureTolerance: reply.FailureTolerance,
|
|
|
|
}
|
|
|
|
for _, server := range reply.Servers {
|
|
|
|
out.Servers = append(out.Servers, api.ServerHealth{
|
|
|
|
ID: server.ID,
|
|
|
|
Name: server.Name,
|
|
|
|
Address: server.Address,
|
|
|
|
Version: server.Version,
|
|
|
|
Leader: server.Leader,
|
|
|
|
SerfStatus: server.SerfStatus.String(),
|
2018-01-30 03:53:34 +00:00
|
|
|
LastContact: server.LastContact,
|
2017-12-18 21:16:23 +00:00
|
|
|
LastTerm: server.LastTerm,
|
|
|
|
LastIndex: server.LastIndex,
|
|
|
|
Healthy: server.Healthy,
|
|
|
|
Voter: server.Voter,
|
|
|
|
StableSince: server.StableSince.Round(time.Second).UTC(),
|
|
|
|
})
|
|
|
|
}
|
|
|
|
|
|
|
|
return out, nil
|
|
|
|
}
|
2018-09-28 04:27:38 +00:00
|
|
|
|
|
|
|
// OperatorSchedulerConfiguration is used to inspect the current Scheduler configuration.
|
|
|
|
// This supports the stale query mode in case the cluster doesn't have a leader.
|
|
|
|
func (s *HTTPServer) OperatorSchedulerConfiguration(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
// Switch on the method
|
|
|
|
switch req.Method {
|
|
|
|
case "GET":
|
2018-11-10 23:37:33 +00:00
|
|
|
return s.schedulerGetConfig(resp, req)
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
case "PUT", "POST":
|
|
|
|
return s.schedulerUpdateConfig(resp, req)
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
default:
|
|
|
|
return nil, CodedError(405, ErrInvalidMethod)
|
|
|
|
}
|
|
|
|
}
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
func (s *HTTPServer) schedulerGetConfig(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
var args structs.GenericRequest
|
|
|
|
if done := s.parse(resp, req, &args.Region, &args.QueryOptions); done {
|
|
|
|
return nil, nil
|
|
|
|
}
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
var reply structs.SchedulerConfigurationResponse
|
|
|
|
if err := s.agent.RPC("Operator.SchedulerGetConfiguration", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
setMeta(resp, &reply.QueryMeta)
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
return reply, nil
|
|
|
|
}
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
func (s *HTTPServer) schedulerUpdateConfig(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
var args structs.SchedulerSetConfigRequest
|
|
|
|
s.parseWriteRequest(req, &args.WriteRequest)
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
var conf api.SchedulerConfiguration
|
|
|
|
if err := decodeBody(req, &conf); err != nil {
|
|
|
|
return nil, CodedError(http.StatusBadRequest, fmt.Sprintf("Error parsing scheduler config: %v", err))
|
|
|
|
}
|
|
|
|
|
|
|
|
args.Config = structs.SchedulerConfiguration{
|
2020-04-24 14:47:43 +00:00
|
|
|
SchedulerAlgorithm: structs.SchedulerAlgorithm(conf.SchedulerAlgorithm),
|
2019-04-29 23:48:07 +00:00
|
|
|
PreemptionConfig: structs.PreemptionConfig{
|
2019-05-03 19:06:12 +00:00
|
|
|
SystemSchedulerEnabled: conf.PreemptionConfig.SystemSchedulerEnabled,
|
|
|
|
BatchSchedulerEnabled: conf.PreemptionConfig.BatchSchedulerEnabled,
|
|
|
|
ServiceSchedulerEnabled: conf.PreemptionConfig.ServiceSchedulerEnabled},
|
2018-11-10 23:37:33 +00:00
|
|
|
}
|
|
|
|
|
2020-04-24 14:47:43 +00:00
|
|
|
if err := args.Config.Validate(); err != nil {
|
|
|
|
return nil, CodedError(http.StatusBadRequest, err.Error())
|
|
|
|
}
|
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
// Check for cas value
|
|
|
|
params := req.URL.Query()
|
|
|
|
if _, ok := params["cas"]; ok {
|
|
|
|
casVal, err := strconv.ParseUint(params.Get("cas"), 10, 64)
|
|
|
|
if err != nil {
|
|
|
|
return nil, CodedError(http.StatusBadRequest, fmt.Sprintf("Error parsing cas value: %v", err))
|
2018-09-28 04:27:38 +00:00
|
|
|
}
|
2018-11-10 23:37:33 +00:00
|
|
|
args.Config.ModifyIndex = casVal
|
|
|
|
args.CAS = true
|
|
|
|
}
|
2018-09-28 04:27:38 +00:00
|
|
|
|
2018-11-10 23:37:33 +00:00
|
|
|
var reply structs.SchedulerSetConfigurationResponse
|
|
|
|
if err := s.agent.RPC("Operator.SchedulerSetConfiguration", &args, &reply); err != nil {
|
|
|
|
return nil, err
|
2018-09-28 04:27:38 +00:00
|
|
|
}
|
2018-11-10 23:37:33 +00:00
|
|
|
setIndex(resp, reply.Index)
|
|
|
|
return reply, nil
|
2018-09-28 04:27:38 +00:00
|
|
|
}
|
2020-05-21 14:05:04 +00:00
|
|
|
|
|
|
|
func (s *HTTPServer) SnapshotRequest(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
switch req.Method {
|
|
|
|
case "GET":
|
|
|
|
return s.snapshotSaveRequest(resp, req)
|
2020-06-04 21:14:43 +00:00
|
|
|
case "PUT", "POST":
|
|
|
|
return s.snapshotRestoreRequest(resp, req)
|
2020-05-21 14:05:04 +00:00
|
|
|
default:
|
|
|
|
return nil, CodedError(405, ErrInvalidMethod)
|
|
|
|
}
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *HTTPServer) snapshotSaveRequest(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
args := &structs.SnapshotSaveRequest{}
|
|
|
|
if s.parse(resp, req, &args.Region, &args.QueryOptions) {
|
|
|
|
return nil, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
var handler structs.StreamingRpcHandler
|
|
|
|
var handlerErr error
|
|
|
|
|
|
|
|
if server := s.agent.Server(); server != nil {
|
|
|
|
handler, handlerErr = server.StreamingRpcHandler("Operator.SnapshotSave")
|
|
|
|
} else if client := s.agent.Client(); client != nil {
|
|
|
|
handler, handlerErr = client.RemoteStreamingRpcHandler("Operator.SnapshotSave")
|
|
|
|
} else {
|
|
|
|
handlerErr = fmt.Errorf("misconfigured connection")
|
|
|
|
}
|
|
|
|
|
|
|
|
if handlerErr != nil {
|
|
|
|
return nil, CodedError(500, handlerErr.Error())
|
|
|
|
}
|
|
|
|
|
|
|
|
httpPipe, handlerPipe := net.Pipe()
|
|
|
|
decoder := codec.NewDecoder(httpPipe, structs.MsgpackHandle)
|
|
|
|
encoder := codec.NewEncoder(httpPipe, structs.MsgpackHandle)
|
|
|
|
|
|
|
|
// Create a goroutine that closes the pipe if the connection closes.
|
|
|
|
ctx, cancel := context.WithCancel(req.Context())
|
|
|
|
defer cancel()
|
|
|
|
go func() {
|
|
|
|
<-ctx.Done()
|
|
|
|
httpPipe.Close()
|
|
|
|
}()
|
|
|
|
|
2020-06-04 21:14:43 +00:00
|
|
|
errCh := make(chan HTTPCodedError, 2)
|
2020-05-21 14:05:04 +00:00
|
|
|
go func() {
|
|
|
|
defer cancel()
|
|
|
|
|
|
|
|
// Send the request
|
|
|
|
if err := encoder.Encode(args); err != nil {
|
|
|
|
errCh <- CodedError(500, err.Error())
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
var res structs.SnapshotSaveResponse
|
|
|
|
if err := decoder.Decode(&res); err != nil {
|
|
|
|
errCh <- CodedError(500, err.Error())
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
if res.ErrorMsg != "" {
|
|
|
|
errCh <- CodedError(res.ErrorCode, res.ErrorMsg)
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
resp.Header().Add("Digest", res.SnapshotChecksum)
|
|
|
|
|
|
|
|
_, err := io.Copy(resp, httpPipe)
|
|
|
|
if err != nil &&
|
|
|
|
err != io.EOF &&
|
|
|
|
!strings.Contains(err.Error(), "closed") &&
|
|
|
|
!strings.Contains(err.Error(), "EOF") {
|
|
|
|
errCh <- CodedError(500, err.Error())
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
errCh <- nil
|
|
|
|
}()
|
|
|
|
|
|
|
|
handler(handlerPipe)
|
|
|
|
cancel()
|
|
|
|
codedErr := <-errCh
|
|
|
|
|
|
|
|
return nil, codedErr
|
|
|
|
}
|
2020-06-04 21:14:43 +00:00
|
|
|
|
|
|
|
func (s *HTTPServer) snapshotRestoreRequest(resp http.ResponseWriter, req *http.Request) (interface{}, error) {
|
|
|
|
args := &structs.SnapshotRestoreRequest{}
|
|
|
|
s.parseWriteRequest(req, &args.WriteRequest)
|
|
|
|
|
|
|
|
var handler structs.StreamingRpcHandler
|
|
|
|
var handlerErr error
|
|
|
|
|
|
|
|
if server := s.agent.Server(); server != nil {
|
|
|
|
handler, handlerErr = server.StreamingRpcHandler("Operator.SnapshotRestore")
|
|
|
|
} else if client := s.agent.Client(); client != nil {
|
|
|
|
handler, handlerErr = client.RemoteStreamingRpcHandler("Operator.SnapshotRestore")
|
|
|
|
} else {
|
|
|
|
handlerErr = fmt.Errorf("misconfigured connection")
|
|
|
|
}
|
|
|
|
|
|
|
|
if handlerErr != nil {
|
|
|
|
return nil, CodedError(500, handlerErr.Error())
|
|
|
|
}
|
|
|
|
|
|
|
|
httpPipe, handlerPipe := net.Pipe()
|
|
|
|
decoder := codec.NewDecoder(httpPipe, structs.MsgpackHandle)
|
|
|
|
encoder := codec.NewEncoder(httpPipe, structs.MsgpackHandle)
|
|
|
|
|
|
|
|
// Create a goroutine that closes the pipe if the connection closes.
|
|
|
|
ctx, cancel := context.WithCancel(req.Context())
|
|
|
|
defer cancel()
|
|
|
|
go func() {
|
|
|
|
<-ctx.Done()
|
|
|
|
httpPipe.Close()
|
|
|
|
}()
|
|
|
|
|
|
|
|
errCh := make(chan HTTPCodedError, 2)
|
|
|
|
go func() {
|
|
|
|
defer cancel()
|
|
|
|
|
|
|
|
// Send the request
|
|
|
|
if err := encoder.Encode(args); err != nil {
|
|
|
|
errCh <- CodedError(500, err.Error())
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
go func() {
|
|
|
|
var wrapper cstructs.StreamErrWrapper
|
|
|
|
bytes := make([]byte, 1024)
|
|
|
|
|
|
|
|
for {
|
|
|
|
n, err := req.Body.Read(bytes)
|
|
|
|
if n > 0 {
|
|
|
|
wrapper.Payload = bytes[:n]
|
|
|
|
err := encoder.Encode(wrapper)
|
|
|
|
if err != nil {
|
|
|
|
errCh <- CodedError(500, err.Error())
|
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
if err != nil {
|
|
|
|
wrapper.Payload = nil
|
|
|
|
wrapper.Error = &cstructs.RpcError{Message: err.Error()}
|
|
|
|
err := encoder.Encode(wrapper)
|
|
|
|
if err != nil {
|
|
|
|
errCh <- CodedError(500, err.Error())
|
|
|
|
}
|
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
var res structs.SnapshotRestoreResponse
|
|
|
|
if err := decoder.Decode(&res); err != nil {
|
|
|
|
errCh <- CodedError(500, err.Error())
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
if res.ErrorMsg != "" {
|
|
|
|
errCh <- CodedError(res.ErrorCode, res.ErrorMsg)
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
errCh <- nil
|
|
|
|
}()
|
|
|
|
|
|
|
|
handler(handlerPipe)
|
|
|
|
cancel()
|
|
|
|
codedErr := <-errCh
|
|
|
|
|
|
|
|
return nil, codedErr
|
|
|
|
}
|