mirror of
https://github.com/wavetermdev/backup.git
synced 2026-08-05 13:57:07 -07:00
update tailer to use filenamegenerator. allow ptyonly option for getcmd
This commit is contained in:
+38
-29
@@ -39,6 +39,11 @@ type CmdWatchEntry struct {
|
|||||||
Tails []TailPos
|
Tails []TailPos
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type FileNameGenerator interface {
|
||||||
|
PtyOutFile(ck base.CommandKey) string
|
||||||
|
RunOutFile(ck base.CommandKey) string
|
||||||
|
}
|
||||||
|
|
||||||
func (w CmdWatchEntry) getTailPos(reqId string) (TailPos, bool) {
|
func (w CmdWatchEntry) getTailPos(reqId string) (TailPos, bool) {
|
||||||
for _, pos := range w.Tails {
|
for _, pos := range w.Tails {
|
||||||
if pos.ReqId == reqId {
|
if pos.ReqId == reqId {
|
||||||
@@ -76,9 +81,9 @@ func (pos TailPos) IsCurrent(entry CmdWatchEntry) bool {
|
|||||||
type Tailer struct {
|
type Tailer struct {
|
||||||
Lock *sync.Mutex
|
Lock *sync.Mutex
|
||||||
WatchList map[base.CommandKey]CmdWatchEntry
|
WatchList map[base.CommandKey]CmdWatchEntry
|
||||||
MHomeDir string
|
|
||||||
Watcher *fsnotify.Watcher
|
Watcher *fsnotify.Watcher
|
||||||
Sender *packet.PacketSender
|
Sender *packet.PacketSender
|
||||||
|
Gen FileNameGenerator
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *Tailer) updateTailPos_nolock(cmdKey base.CommandKey, reqId string, pos TailPos) {
|
func (t *Tailer) updateTailPos_nolock(cmdKey base.CommandKey, reqId string, pos TailPos) {
|
||||||
@@ -108,10 +113,9 @@ func (t *Tailer) removeTailPos_nolock(cmdKey base.CommandKey, reqId string) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// delete from watchlist, remove watches
|
// delete from watchlist, remove watches
|
||||||
fileNames := base.MakeCommandFileNamesWithHome(t.MHomeDir, cmdKey)
|
|
||||||
delete(t.WatchList, cmdKey)
|
delete(t.WatchList, cmdKey)
|
||||||
t.Watcher.Remove(fileNames.PtyOutFile)
|
t.Watcher.Remove(t.Gen.PtyOutFile(cmdKey))
|
||||||
t.Watcher.Remove(fileNames.RunnerOutFile)
|
t.Watcher.Remove(t.Gen.RunOutFile(cmdKey))
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *Tailer) getEntryAndPos_nolock(cmdKey base.CommandKey, reqId string) (CmdWatchEntry, TailPos, bool) {
|
func (t *Tailer) getEntryAndPos_nolock(cmdKey base.CommandKey, reqId string) (CmdWatchEntry, TailPos, bool) {
|
||||||
@@ -126,13 +130,12 @@ func (t *Tailer) getEntryAndPos_nolock(cmdKey base.CommandKey, reqId string) (Cm
|
|||||||
return entry, pos, true
|
return entry, pos, true
|
||||||
}
|
}
|
||||||
|
|
||||||
func MakeTailer(sender *packet.PacketSender) (*Tailer, error) {
|
func MakeTailer(sender *packet.PacketSender, gen FileNameGenerator) (*Tailer, error) {
|
||||||
mhomeDir := base.GetMShellHomeDir()
|
|
||||||
rtn := &Tailer{
|
rtn := &Tailer{
|
||||||
Lock: &sync.Mutex{},
|
Lock: &sync.Mutex{},
|
||||||
WatchList: make(map[base.CommandKey]CmdWatchEntry),
|
WatchList: make(map[base.CommandKey]CmdWatchEntry),
|
||||||
MHomeDir: mhomeDir,
|
|
||||||
Sender: sender,
|
Sender: sender,
|
||||||
|
Gen: gen,
|
||||||
}
|
}
|
||||||
var err error
|
var err error
|
||||||
rtn.Watcher, err = fsnotify.NewWatcher()
|
rtn.Watcher, err = fsnotify.NewWatcher()
|
||||||
@@ -156,13 +159,13 @@ func (t *Tailer) readDataFromFile(fileName string, pos int64, maxBytes int) ([]b
|
|||||||
return buf[0:nr], nil
|
return buf[0:nr], nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *Tailer) makeCmdDataPacket(fileNames *base.CommandFileNames, entry CmdWatchEntry, pos TailPos) (*packet.CmdDataPacketType, error) {
|
func (t *Tailer) makeCmdDataPacket(entry CmdWatchEntry, pos TailPos) (*packet.CmdDataPacketType, error) {
|
||||||
dataPacket := packet.MakeCmdDataPacket(pos.ReqId)
|
dataPacket := packet.MakeCmdDataPacket(pos.ReqId)
|
||||||
dataPacket.CK = entry.CmdKey
|
dataPacket.CK = entry.CmdKey
|
||||||
dataPacket.PtyPos = pos.TailPtyPos
|
dataPacket.PtyPos = pos.TailPtyPos
|
||||||
dataPacket.RunPos = pos.TailRunPos
|
dataPacket.RunPos = pos.TailRunPos
|
||||||
if entry.FilePtyLen > pos.TailPtyPos {
|
if entry.FilePtyLen > pos.TailPtyPos {
|
||||||
ptyData, err := t.readDataFromFile(fileNames.PtyOutFile, pos.TailPtyPos, MaxDataBytes)
|
ptyData, err := t.readDataFromFile(t.Gen.PtyOutFile(entry.CmdKey), pos.TailPtyPos, MaxDataBytes)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -170,7 +173,7 @@ func (t *Tailer) makeCmdDataPacket(fileNames *base.CommandFileNames, entry CmdWa
|
|||||||
dataPacket.PtyDataLen = len(ptyData)
|
dataPacket.PtyDataLen = len(ptyData)
|
||||||
}
|
}
|
||||||
if entry.FileRunLen > pos.TailRunPos {
|
if entry.FileRunLen > pos.TailRunPos {
|
||||||
runData, err := t.readDataFromFile(fileNames.RunnerOutFile, pos.TailRunPos, MaxDataBytes)
|
runData, err := t.readDataFromFile(t.Gen.RunOutFile(entry.CmdKey), pos.TailRunPos, MaxDataBytes)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
@@ -188,8 +191,7 @@ func (t *Tailer) runSingleDataTransfer(key base.CommandKey, reqId string) (*pack
|
|||||||
if !foundPos {
|
if !foundPos {
|
||||||
return nil, false, nil
|
return nil, false, nil
|
||||||
}
|
}
|
||||||
fileNames := base.MakeCommandFileNamesWithHome(t.MHomeDir, key)
|
dataPacket, dataErr := t.makeCmdDataPacket(entry, pos)
|
||||||
dataPacket, dataErr := t.makeCmdDataPacket(fileNames, entry, pos)
|
|
||||||
|
|
||||||
t.Lock.Lock()
|
t.Lock.Lock()
|
||||||
defer t.Lock.Unlock()
|
defer t.Lock.Unlock()
|
||||||
@@ -330,13 +332,12 @@ func max(v1 int64, v2 int64) int64 {
|
|||||||
return v2
|
return v2
|
||||||
}
|
}
|
||||||
|
|
||||||
func (entry *CmdWatchEntry) fillFilePos(scHomeDir string) {
|
func (entry *CmdWatchEntry) fillFilePos(gen FileNameGenerator) {
|
||||||
fileNames := base.MakeCommandFileNamesWithHome(scHomeDir, entry.CmdKey)
|
ptyInfo, _ := os.Stat(gen.PtyOutFile(entry.CmdKey))
|
||||||
ptyInfo, _ := os.Stat(fileNames.PtyOutFile)
|
|
||||||
if ptyInfo != nil {
|
if ptyInfo != nil {
|
||||||
entry.FilePtyLen = ptyInfo.Size()
|
entry.FilePtyLen = ptyInfo.Size()
|
||||||
}
|
}
|
||||||
runoutInfo, _ := os.Stat(fileNames.RunnerOutFile)
|
runoutInfo, _ := os.Stat(gen.RunOutFile(entry.CmdKey))
|
||||||
if runoutInfo != nil {
|
if runoutInfo != nil {
|
||||||
entry.FileRunLen = runoutInfo.Size()
|
entry.FileRunLen = runoutInfo.Size()
|
||||||
}
|
}
|
||||||
@@ -348,27 +349,33 @@ func (t *Tailer) RemoveWatch(pk *packet.UntailCmdPacketType) {
|
|||||||
t.removeTailPos_nolock(pk.CK, pk.ReqId)
|
t.removeTailPos_nolock(pk.CK, pk.ReqId)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *Tailer) AddFileWatches_nolock(fileNames *base.CommandFileNames) error {
|
func (t *Tailer) AddFileWatches_nolock(key base.CommandKey, ptyOnly bool) error {
|
||||||
err := t.Watcher.Add(fileNames.PtyOutFile)
|
ptyName := t.Gen.PtyOutFile(key)
|
||||||
|
runName := t.Gen.RunOutFile(key)
|
||||||
|
fmt.Printf("WATCH> add %s\n", ptyName)
|
||||||
|
err := t.Watcher.Add(ptyName)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
err = t.Watcher.Add(fileNames.RunnerOutFile)
|
if ptyOnly {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
err = t.Watcher.Add(runName)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Watcher.Remove(fileNames.PtyOutFile) // best effort clean up
|
t.Watcher.Remove(ptyName) // best effort clean up
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *Tailer) AddWatch(getPacket *packet.GetCmdPacketType) error {
|
// returns (up-to-date/done, error)
|
||||||
|
func (t *Tailer) AddWatch(getPacket *packet.GetCmdPacketType) (bool, error) {
|
||||||
if err := getPacket.CK.Validate("getcmd"); err != nil {
|
if err := getPacket.CK.Validate("getcmd"); err != nil {
|
||||||
return err
|
return false, err
|
||||||
}
|
}
|
||||||
if getPacket.ReqId == "" {
|
if getPacket.ReqId == "" {
|
||||||
return fmt.Errorf("getcmd, no reqid specified")
|
return false, fmt.Errorf("getcmd, no reqid specified")
|
||||||
}
|
}
|
||||||
fileNames := base.MakeCommandFileNamesWithHome(t.MHomeDir, getPacket.CK)
|
|
||||||
t.Lock.Lock()
|
t.Lock.Lock()
|
||||||
defer t.Lock.Unlock()
|
defer t.Lock.Unlock()
|
||||||
key := getPacket.CK
|
key := getPacket.CK
|
||||||
@@ -376,7 +383,7 @@ func (t *Tailer) AddWatch(getPacket *packet.GetCmdPacketType) error {
|
|||||||
if !foundEntry {
|
if !foundEntry {
|
||||||
// initialize entry, add watches
|
// initialize entry, add watches
|
||||||
entry = CmdWatchEntry{CmdKey: key}
|
entry = CmdWatchEntry{CmdKey: key}
|
||||||
entry.fillFilePos(t.MHomeDir)
|
entry.fillFilePos(t.Gen)
|
||||||
}
|
}
|
||||||
pos, foundPos := entry.getTailPos(getPacket.ReqId)
|
pos, foundPos := entry.getTailPos(getPacket.ReqId)
|
||||||
if !foundPos {
|
if !foundPos {
|
||||||
@@ -397,13 +404,15 @@ func (t *Tailer) AddWatch(getPacket *packet.GetCmdPacketType) error {
|
|||||||
entry.updateTailPos(pos.ReqId, pos)
|
entry.updateTailPos(pos.ReqId, pos)
|
||||||
if !pos.Follow && pos.IsCurrent(entry) {
|
if !pos.Follow && pos.IsCurrent(entry) {
|
||||||
// don't add to t.WatchList, don't t.AddFileWatches_nolock, send rpc response
|
// don't add to t.WatchList, don't t.AddFileWatches_nolock, send rpc response
|
||||||
go func() { t.Sender.SendResponse(getPacket.ReqId, true) }()
|
return true, nil
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
if !foundEntry {
|
if !foundEntry {
|
||||||
t.AddFileWatches_nolock(fileNames)
|
err := t.AddFileWatches_nolock(key, getPacket.PtyOnly)
|
||||||
|
if err != nil {
|
||||||
|
return false, err
|
||||||
|
}
|
||||||
}
|
}
|
||||||
t.WatchList[key] = entry
|
t.WatchList[key] = entry
|
||||||
t.tryStartRun_nolock(entry, pos)
|
t.tryStartRun_nolock(entry, pos)
|
||||||
return nil
|
return false, nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -39,7 +39,7 @@ const (
|
|||||||
DataEndPacketStr = "dataend"
|
DataEndPacketStr = "dataend"
|
||||||
ResponsePacketStr = "resp" // rpc-response
|
ResponsePacketStr = "resp" // rpc-response
|
||||||
DonePacketStr = "done"
|
DonePacketStr = "done"
|
||||||
CmdErrorPacketStr = "cmderror"
|
CmdErrorPacketStr = "cmderror" // command
|
||||||
MessagePacketStr = "message"
|
MessagePacketStr = "message"
|
||||||
GetCmdPacketStr = "getcmd" // rpc
|
GetCmdPacketStr = "getcmd" // rpc
|
||||||
UntailCmdPacketStr = "untailcmd" // rpc
|
UntailCmdPacketStr = "untailcmd" // rpc
|
||||||
@@ -275,12 +275,13 @@ func MakeUntailCmdPacket() *UntailCmdPacketType {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type GetCmdPacketType struct {
|
type GetCmdPacketType struct {
|
||||||
Type string `json:"type"`
|
Type string `json:"type"`
|
||||||
ReqId string `json:"reqid"`
|
ReqId string `json:"reqid"`
|
||||||
CK base.CommandKey `json:"ck"`
|
CK base.CommandKey `json:"ck"`
|
||||||
PtyPos int64 `json:"ptypos"`
|
PtyPos int64 `json:"ptypos"`
|
||||||
RunPos int64 `json:"runpos"`
|
RunPos int64 `json:"runpos"`
|
||||||
Tail bool `json:"tail,omitempty"`
|
Tail bool `json:"tail,omitempty"`
|
||||||
|
PtyOnly bool `json:"ptyonly,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
func (*GetCmdPacketType) GetType() string {
|
func (*GetCmdPacketType) GetType() string {
|
||||||
|
|||||||
Reference in New Issue
Block a user