mirror of
https://github.com/postmannen/ctrl.git
synced 2024-12-15 17:51:15 +00:00
fixed timeouts for copy messages
This commit is contained in:
parent
29a382838c
commit
3f5b58ffab
1 changed files with 16 additions and 0 deletions
|
@ -252,6 +252,10 @@ func (m methodREQCopySrc) handler(proc process, message Message, node string) ([
|
||||||
//msg.Method = REQToFile
|
//msg.Method = REQToFile
|
||||||
msg.Method = REQCopyDst
|
msg.Method = REQCopyDst
|
||||||
msg.Data = cb
|
msg.Data = cb
|
||||||
|
msg.ACKTimeout = message.ACKTimeout
|
||||||
|
msg.Retries = message.Retries
|
||||||
|
msg.ReplyACKTimeout = message.ReplyACKTimeout
|
||||||
|
msg.ReplyRetries = message.ReplyRetries
|
||||||
// msg.Directory = dstDir
|
// msg.Directory = dstDir
|
||||||
// msg.FileName = dstFile
|
// msg.FileName = dstFile
|
||||||
|
|
||||||
|
@ -527,6 +531,8 @@ func copySrcSubProcFunc(proc process, cia copyInitialData, cancel context.Cancel
|
||||||
IsSubPublishedMsg: true,
|
IsSubPublishedMsg: true,
|
||||||
ACKTimeout: initialMessage.ACKTimeout,
|
ACKTimeout: initialMessage.ACKTimeout,
|
||||||
Retries: initialMessage.Retries,
|
Retries: initialMessage.Retries,
|
||||||
|
ReplyACKTimeout: initialMessage.ReplyACKTimeout,
|
||||||
|
ReplyRetries: initialMessage.ReplyRetries,
|
||||||
}
|
}
|
||||||
|
|
||||||
fmt.Printf(" * DEBUG: ACKTimeout:%v, Retries: %v\n", initialMessage.ACKTimeout, initialMessage.Retries)
|
fmt.Printf(" * DEBUG: ACKTimeout:%v, Retries: %v\n", initialMessage.ACKTimeout, initialMessage.Retries)
|
||||||
|
@ -597,6 +603,8 @@ func copySrcSubProcFunc(proc process, cia copyInitialData, cancel context.Cancel
|
||||||
IsSubPublishedMsg: true,
|
IsSubPublishedMsg: true,
|
||||||
ACKTimeout: initialMessage.ACKTimeout,
|
ACKTimeout: initialMessage.ACKTimeout,
|
||||||
Retries: initialMessage.Retries,
|
Retries: initialMessage.Retries,
|
||||||
|
ReplyACKTimeout: initialMessage.ReplyACKTimeout,
|
||||||
|
ReplyRetries: initialMessage.ReplyRetries,
|
||||||
}
|
}
|
||||||
|
|
||||||
sam, err := newSubjectAndMessage(msg)
|
sam, err := newSubjectAndMessage(msg)
|
||||||
|
@ -658,6 +666,8 @@ func copyDstSubProcFunc(proc process, cia copyInitialData, message Message, canc
|
||||||
IsSubPublishedMsg: true,
|
IsSubPublishedMsg: true,
|
||||||
ACKTimeout: message.ACKTimeout,
|
ACKTimeout: message.ACKTimeout,
|
||||||
Retries: message.Retries,
|
Retries: message.Retries,
|
||||||
|
ReplyACKTimeout: message.ReplyACKTimeout,
|
||||||
|
ReplyRetries: message.ReplyRetries,
|
||||||
}
|
}
|
||||||
|
|
||||||
sam, err := newSubjectAndMessage(msg)
|
sam, err := newSubjectAndMessage(msg)
|
||||||
|
@ -758,6 +768,8 @@ func copyDstSubProcFunc(proc process, cia copyInitialData, message Message, canc
|
||||||
IsSubPublishedMsg: true,
|
IsSubPublishedMsg: true,
|
||||||
ACKTimeout: message.ACKTimeout,
|
ACKTimeout: message.ACKTimeout,
|
||||||
Retries: message.Retries,
|
Retries: message.Retries,
|
||||||
|
ReplyACKTimeout: message.ReplyACKTimeout,
|
||||||
|
ReplyRetries: message.ReplyRetries,
|
||||||
}
|
}
|
||||||
|
|
||||||
sam, err := newSubjectAndMessage(msg)
|
sam, err := newSubjectAndMessage(msg)
|
||||||
|
@ -789,6 +801,8 @@ func copyDstSubProcFunc(proc process, cia copyInitialData, message Message, canc
|
||||||
IsSubPublishedMsg: true,
|
IsSubPublishedMsg: true,
|
||||||
ACKTimeout: message.ACKTimeout,
|
ACKTimeout: message.ACKTimeout,
|
||||||
Retries: message.Retries,
|
Retries: message.Retries,
|
||||||
|
ReplyACKTimeout: message.ReplyACKTimeout,
|
||||||
|
ReplyRetries: message.ReplyRetries,
|
||||||
}
|
}
|
||||||
|
|
||||||
sam, err := newSubjectAndMessage(msg)
|
sam, err := newSubjectAndMessage(msg)
|
||||||
|
@ -941,6 +955,8 @@ func copyDstSubProcFunc(proc process, cia copyInitialData, message Message, canc
|
||||||
IsSubPublishedMsg: true,
|
IsSubPublishedMsg: true,
|
||||||
ACKTimeout: message.ACKTimeout,
|
ACKTimeout: message.ACKTimeout,
|
||||||
Retries: message.Retries,
|
Retries: message.Retries,
|
||||||
|
ReplyACKTimeout: message.ReplyACKTimeout,
|
||||||
|
ReplyRetries: message.ReplyRetries,
|
||||||
}
|
}
|
||||||
|
|
||||||
sam, err := newSubjectAndMessage(msg)
|
sam, err := newSubjectAndMessage(msg)
|
||||||
|
|
Loading…
Reference in a new issue