mirror of
https://github.com/luxfi/netrunner.git
synced 2026-07-27 00:04:23 +00:00
feat(zapwire): PauseNode over ZAP
This commit is contained in:
@@ -174,6 +174,15 @@ func (c *Client) RestartNode(ctx context.Context, req *types.RestartNodeRequest)
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
// PauseNode pauses a single node (process suspended, peers see it as offline).
|
||||
func (c *Client) PauseNode(ctx context.Context, req *types.PauseNodeRequest) (*types.PauseNodeResponse, error) {
|
||||
resp := &types.PauseNodeResponse{}
|
||||
if err := c.callSub(ctx, OpStart, SubPauseNode, req, resp); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
// Stop terminates the cluster and returns the final snapshot.
|
||||
func (c *Client) Stop(ctx context.Context) (*types.StopResponse, error) {
|
||||
resp := &types.StopResponse{}
|
||||
|
||||
@@ -45,6 +45,7 @@ type stubBackend struct {
|
||||
lastAddNode *types.AddNodeRequest
|
||||
lastRemoveNode *types.RemoveNodeRequest
|
||||
lastRestartNode *types.RestartNodeRequest
|
||||
lastPauseNode *types.PauseNodeRequest
|
||||
}
|
||||
|
||||
func (b *stubBackend) Ping(ctx context.Context) (*types.PingResponse, error) {
|
||||
@@ -157,6 +158,22 @@ func (b *stubBackend) RestartNode(ctx context.Context, req *types.RestartNodeReq
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (b *stubBackend) PauseNode(ctx context.Context, req *types.PauseNodeRequest) (*types.PauseNodeResponse, error) {
|
||||
b.lastPauseNode = req
|
||||
return &types.PauseNodeResponse{
|
||||
ClusterInfo: &types.ClusterInfo{
|
||||
NodeNames: []string{req.Name},
|
||||
NodeInfos: map[string]*types.NodeInfo{
|
||||
req.Name: {Name: req.Name, Paused: true},
|
||||
},
|
||||
PID: 42,
|
||||
RootDataDir: "/tmp/netrunner",
|
||||
Healthy: true,
|
||||
NetworkID: 12345,
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (b *stubBackend) Stop(ctx context.Context) (*types.StopResponse, error) {
|
||||
return &types.StopResponse{
|
||||
ClusterInfo: &types.ClusterInfo{
|
||||
@@ -422,6 +439,31 @@ func TestE2ERestartNode(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestE2EPauseNode(t *testing.T) {
|
||||
be := &stubBackend{}
|
||||
srv, teardown := runServer(t, be)
|
||||
defer teardown()
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
client, err := Dial(ctx, srv.Addr())
|
||||
if err != nil {
|
||||
t.Fatalf("Dial: %v", err)
|
||||
}
|
||||
defer client.Close()
|
||||
|
||||
resp, err := client.PauseNode(ctx, &types.PauseNodeRequest{Name: "alpha"})
|
||||
if err != nil {
|
||||
t.Fatalf("PauseNode: %v", err)
|
||||
}
|
||||
if be.lastPauseNode == nil || be.lastPauseNode.Name != "alpha" {
|
||||
t.Fatalf("Name round-trip: %+v", be.lastPauseNode)
|
||||
}
|
||||
if !resp.ClusterInfo.NodeInfos["alpha"].Paused {
|
||||
t.Fatalf("Paused flag not propagated")
|
||||
}
|
||||
}
|
||||
|
||||
func TestE2EStop(t *testing.T) {
|
||||
srv, teardown := runServer(t, &stubBackend{})
|
||||
defer teardown()
|
||||
|
||||
@@ -26,6 +26,7 @@ type Backend interface {
|
||||
AddNode(ctx context.Context, req *types.AddNodeRequest) (*types.AddNodeResponse, error)
|
||||
RemoveNode(ctx context.Context, req *types.RemoveNodeRequest) (*types.RemoveNodeResponse, error)
|
||||
RestartNode(ctx context.Context, req *types.RestartNodeRequest) (*types.RestartNodeResponse, error)
|
||||
PauseNode(ctx context.Context, req *types.PauseNodeRequest) (*types.PauseNodeResponse, error)
|
||||
}
|
||||
|
||||
// Server hosts a netrunner control RPC over ZAP.
|
||||
@@ -172,6 +173,16 @@ func (s *Server) handleStartSub(
|
||||
return 0, nil, err
|
||||
}
|
||||
return msgType, encodeResp(resp), nil
|
||||
case SubPauseNode:
|
||||
req := &types.PauseNodeRequest{}
|
||||
if err := req.Decode(zap.NewReader(body)); err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
resp, err := s.be.PauseNode(ctx, req)
|
||||
if err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
return msgType, encodeResp(resp), nil
|
||||
default:
|
||||
return 0, nil, fmt.Errorf("zapwire: unknown SubOp 0x%02x for OpStart", uint8(sub))
|
||||
}
|
||||
|
||||
@@ -433,6 +433,48 @@ func (r *RestartNodeResponse) Decode(rd *zap.Reader) error {
|
||||
return r.ClusterInfo.Decode(rd)
|
||||
}
|
||||
|
||||
// PauseNodeRequest matches rpcpb.PauseNodeRequest.
|
||||
type PauseNodeRequest struct {
|
||||
Name string
|
||||
}
|
||||
|
||||
func (p *PauseNodeRequest) Encode(b *zap.Buffer) { b.WriteString(p.Name) }
|
||||
func (p *PauseNodeRequest) Decode(rd *zap.Reader) error {
|
||||
s, err := rd.ReadString()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
p.Name = s
|
||||
return nil
|
||||
}
|
||||
|
||||
// PauseNodeResponse carries the post-pause cluster snapshot.
|
||||
type PauseNodeResponse struct {
|
||||
ClusterInfo *ClusterInfo
|
||||
}
|
||||
|
||||
func (p *PauseNodeResponse) Encode(b *zap.Buffer) {
|
||||
if p.ClusterInfo == nil {
|
||||
b.WriteUint8(0)
|
||||
return
|
||||
}
|
||||
b.WriteUint8(1)
|
||||
p.ClusterInfo.Encode(b)
|
||||
}
|
||||
|
||||
func (p *PauseNodeResponse) Decode(rd *zap.Reader) error {
|
||||
present, err := rd.ReadUint8()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if present == 0 {
|
||||
p.ClusterInfo = nil
|
||||
return nil
|
||||
}
|
||||
p.ClusterInfo = &ClusterInfo{}
|
||||
return p.ClusterInfo.Decode(rd)
|
||||
}
|
||||
|
||||
// StopRequest is empty.
|
||||
type StopRequest struct{}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user