mirror of
https://github.com/wavetermdev/backup.git
synced 2026-08-05 13:57:07 -07:00
zsh reinit fixes (#477)
* reset command now initiates and completes async so there is feedback that something is happening when it takes a long time * switch from standard rpc to rpciter * checkpoint on reinit -- stream output, stats packet, logging to cmd pty, new endBytes for EOF * make generic versions of endbytes scanner and channel output funcs * update bash to use more modern state parsing (tricks learned from zsh) * verbose mode, fix stats output message * add a diff when verbose mode is on
This commit is contained in:
@@ -196,7 +196,7 @@ func (msh *MShellProc) EnsureShellType(ctx context.Context, shellType string) er
|
||||
return nil
|
||||
}
|
||||
// try to reinit the shell
|
||||
_, err := msh.ReInit(ctx, shellType)
|
||||
_, err := msh.ReInit(ctx, shellType, nil, false)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error trying to initialize shell %q: %v", shellType, err)
|
||||
}
|
||||
@@ -1401,33 +1401,60 @@ func makeReinitErrorUpdate(shellType string) sstore.ActivityUpdate {
|
||||
return rtn
|
||||
}
|
||||
|
||||
func (msh *MShellProc) ReInit(ctx context.Context, shellType string) (*packet.ShellStatePacketType, error) {
|
||||
func (msh *MShellProc) ReInit(ctx context.Context, shellType string, dataFn func([]byte), verbose bool) (rtnPk *packet.ShellStatePacketType, rtnErr error) {
|
||||
if !msh.IsConnected() {
|
||||
return nil, fmt.Errorf("cannot reinit, remote is not connected")
|
||||
}
|
||||
if shellType != packet.ShellType_bash && shellType != packet.ShellType_zsh {
|
||||
return nil, fmt.Errorf("invalid shell type %q", shellType)
|
||||
}
|
||||
if dataFn == nil {
|
||||
dataFn = func([]byte) {}
|
||||
}
|
||||
defer func() {
|
||||
if rtnErr != nil {
|
||||
sstore.UpdateActivityWrap(ctx, makeReinitErrorUpdate(shellType), "reiniterror")
|
||||
}
|
||||
}()
|
||||
startTs := time.Now()
|
||||
reinitPk := packet.MakeReInitPacket()
|
||||
reinitPk.ReqId = uuid.New().String()
|
||||
reinitPk.ShellType = shellType
|
||||
resp, err := msh.PacketRpcRaw(ctx, reinitPk)
|
||||
rpcIter, err := msh.PacketRpcIter(ctx, reinitPk)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if resp == nil {
|
||||
return nil, fmt.Errorf("no response")
|
||||
}
|
||||
ssPk, ok := resp.(*packet.ShellStatePacketType)
|
||||
if !ok {
|
||||
sstore.UpdateActivityWrap(ctx, makeReinitErrorUpdate(shellType), "reiniterror")
|
||||
if respPk, ok := resp.(*packet.ResponsePacketType); ok && respPk.Error != "" {
|
||||
return nil, fmt.Errorf("error reinitializing remote: %s", respPk.Error)
|
||||
defer rpcIter.Close()
|
||||
var ssPk *packet.ShellStatePacketType
|
||||
for {
|
||||
resp, err := rpcIter.Next(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return nil, fmt.Errorf("invalid reinit response (not an shellstate packet): %T", resp)
|
||||
if resp == nil {
|
||||
return nil, fmt.Errorf("channel closed with no response")
|
||||
}
|
||||
var ok bool
|
||||
ssPk, ok = resp.(*packet.ShellStatePacketType)
|
||||
if ok {
|
||||
break
|
||||
}
|
||||
respPk, ok := resp.(*packet.ResponsePacketType)
|
||||
if ok {
|
||||
if respPk.Error != "" {
|
||||
return nil, fmt.Errorf("error reinitializing remote: %s", respPk.Error)
|
||||
}
|
||||
return nil, fmt.Errorf("invalid response from waveshell")
|
||||
}
|
||||
dataPk, ok := resp.(*packet.FileDataPacketType)
|
||||
if ok {
|
||||
dataFn(dataPk.Data)
|
||||
continue
|
||||
}
|
||||
invalidPkStr := fmt.Sprintf("\r\ninvalid packettype from waveshell: %s\r\n", resp.GetType())
|
||||
dataFn([]byte(invalidPkStr))
|
||||
}
|
||||
if ssPk.State == nil {
|
||||
sstore.UpdateActivityWrap(ctx, makeReinitErrorUpdate(shellType), "reiniterror")
|
||||
if ssPk == nil || ssPk.State == nil {
|
||||
return nil, fmt.Errorf("invalid reinit response shellstate packet does not contain remote state")
|
||||
}
|
||||
// TODO: maybe we don't need to save statebase here. should be possible to save it on demand
|
||||
@@ -1438,10 +1465,29 @@ func (msh *MShellProc) ReInit(ctx context.Context, shellType string) (*packet.Sh
|
||||
return nil, fmt.Errorf("error storing remote state: %w", err)
|
||||
}
|
||||
msh.StateMap.SetCurrentState(ssPk.State.GetShellType(), ssPk.State)
|
||||
msh.WriteToPtyBuffer("initialized shell:%s state:%s\n", shellType, ssPk.State.GetHashVal(false))
|
||||
timeDur := time.Since(startTs)
|
||||
dataFn([]byte(makeShellInitOutputMsg(verbose, ssPk.State, ssPk.Stats, timeDur, false)))
|
||||
msh.WriteToPtyBuffer("%s", makeShellInitOutputMsg(false, ssPk.State, ssPk.Stats, timeDur, true))
|
||||
return ssPk, nil
|
||||
}
|
||||
|
||||
func makeShellInitOutputMsg(verbose bool, state *packet.ShellState, stats *packet.ShellStateStats, dur time.Duration, ptyMsg bool) string {
|
||||
if !verbose || ptyMsg {
|
||||
if ptyMsg {
|
||||
return fmt.Sprintf("initialized state shell:%s statehash:%s %dms\n", state.GetShellType(), state.GetHashVal(false), dur.Milliseconds())
|
||||
} else {
|
||||
return fmt.Sprintf("initialized connection state (shell:%s)\r\n", state.GetShellType())
|
||||
}
|
||||
}
|
||||
var buf bytes.Buffer
|
||||
buf.WriteString("-----\r\n")
|
||||
buf.WriteString(fmt.Sprintf("initialized connection shell:%s statehash:%s %dms\r\n", state.GetShellType(), state.GetHashVal(false), dur.Milliseconds()))
|
||||
if stats != nil {
|
||||
buf.WriteString(fmt.Sprintf(" outsize:%s size:%s env:%d, vars:%d, aliases:%d, funcs:%d\r\n", scbase.NumFormatDec(stats.OutputSize), scbase.NumFormatDec(stats.StateSize), stats.EnvCount, stats.VarCount, stats.AliasCount, stats.FuncCount))
|
||||
}
|
||||
return buf.String()
|
||||
}
|
||||
|
||||
func (msh *MShellProc) WriteFile(ctx context.Context, writePk *packet.WriteFilePacketType) (*packet.RpcResponseIter, error) {
|
||||
return msh.PacketRpcIter(ctx, writePk)
|
||||
}
|
||||
@@ -1690,7 +1736,7 @@ func (msh *MShellProc) initActiveShells() {
|
||||
return
|
||||
}
|
||||
for _, shellType := range activeShells {
|
||||
_, err = msh.ReInit(ctx, shellType)
|
||||
_, err = msh.ReInit(ctx, shellType, nil, false)
|
||||
if err != nil {
|
||||
msh.WriteToPtyBuffer("*error reiniting shell %q: %v\n", shellType, err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user