From b2a2b6252d161afa1fd3c9cc3e2dc61eb920ed6d Mon Sep 17 00:00:00 2001 From: sawka Date: Fri, 19 Aug 2022 17:14:53 -0700 Subject: [PATCH] send cmd updates for donepk --- pkg/remote/remote.go | 24 +++++++++++++++--------- pkg/sstore/dbops.go | 19 ++++++++++++++++--- 2 files changed, 31 insertions(+), 12 deletions(-) diff --git a/pkg/remote/remote.go b/pkg/remote/remote.go index 3bd5ac21..410ba219 100644 --- a/pkg/remote/remote.go +++ b/pkg/remote/remote.go @@ -48,9 +48,8 @@ const ( var GlobalStore *Store type Store struct { - Lock *sync.Mutex - Map map[string]*MShellProc // key=remoteid - CmdStatusCallback func(ck base.CommandKey, status string) + Lock *sync.Mutex + Map map[string]*MShellProc // key=remoteid } type RemoteState struct { @@ -440,13 +439,17 @@ func makeDataAckPacket(ck base.CommandKey, fdNum int, ackLen int, err error) *pa } func (msh *MShellProc) handleCmdDonePacket(donePk *packet.CmdDonePacketType) { - err := sstore.UpdateCmdDonePk(context.Background(), donePk) + update, err := sstore.UpdateCmdDonePk(context.Background(), donePk) if err != nil { fmt.Printf("[error] updating cmddone: %v\n", err) return } - if GlobalStore.CmdStatusCallback != nil { - GlobalStore.CmdStatusCallback(donePk.CK, sstore.CmdStatusDone) + if update != nil { + // TODO fix timing issue (this update gets to the FE before run-command returns for short lived commands) + go func() { + time.Sleep(10 * time.Millisecond) + sstore.MainBus.SendUpdate(donePk.CK.GetSessionId(), update) + }() } return } @@ -461,10 +464,13 @@ func (msh *MShellProc) handleCmdErrorPacket(errPk *packet.CmdErrorPacketType) { } func (msh *MShellProc) notifyHangups_nolock() { - if GlobalStore.CmdStatusCallback != nil { - for _, ck := range msh.RunningCmds { - GlobalStore.CmdStatusCallback(ck, sstore.CmdStatusHangup) + for _, ck := range msh.RunningCmds { + cmd, err := sstore.GetCmdById(context.Background(), ck.GetSessionId(), ck.GetCmdId()) + if err != nil { + continue } + update := sstore.LineUpdate{Cmd: cmd} + sstore.MainBus.SendUpdate(ck.GetSessionId(), update) } msh.RunningCmds = nil } diff --git a/pkg/sstore/dbops.go b/pkg/sstore/dbops.go index 954d6d08..8df780cd 100644 --- a/pkg/sstore/dbops.go +++ b/pkg/sstore/dbops.go @@ -427,15 +427,28 @@ func GetCmdById(ctx context.Context, sessionId string, cmdId string) (*CmdType, return cmd, nil } -func UpdateCmdDonePk(ctx context.Context, donePk *packet.CmdDonePacketType) error { +func UpdateCmdDonePk(ctx context.Context, donePk *packet.CmdDonePacketType) (UpdatePacket, error) { if donePk == nil || donePk.CK.IsEmpty() { - return fmt.Errorf("invalid cmddone packet (no ck)") + return nil, fmt.Errorf("invalid cmddone packet (no ck)") } - return WithTx(ctx, func(tx *TxWrap) error { + var rtnCmd *CmdType + txErr := WithTx(ctx, func(tx *TxWrap) error { query := `UPDATE cmd SET status = ?, donepk = ? WHERE sessionid = ? AND cmdid = ?` tx.ExecWrap(query, CmdStatusDone, quickJson(donePk), donePk.CK.GetSessionId(), donePk.CK.GetCmdId()) + var err error + rtnCmd, err = GetCmdById(tx.Context(), donePk.CK.GetSessionId(), donePk.CK.GetCmdId()) + if err != nil { + return err + } return nil }) + if txErr != nil { + return nil, txErr + } + if rtnCmd == nil { + return nil, fmt.Errorf("cmd data not found for ck[%s]", donePk.CK) + } + return LineUpdate{Cmd: rtnCmd}, nil } func AppendCmdErrorPk(ctx context.Context, errPk *packet.CmdErrorPacketType) error {