1
0
Fork 0
mirror of https://github.com/postmannen/ctrl.git synced 2024-12-14 12:37:31 +00:00

Added initital structure for methodREQToSocket

This commit is contained in:
postmannen 2021-07-01 10:05:34 +02:00
parent 8da6d85a85
commit 00b439d6d2
2 changed files with 39 additions and 1 deletions

View file

@ -79,6 +79,8 @@ func (p process) ProcessesStart() {
if p.configuration.StartSubREQTailFile.OK {
p.startup.subREQTailFile(p)
}
p.startup.subREQToSocket(p)
}
// ---------------------------------------------------------------------------------------
@ -249,3 +251,11 @@ func (s startup) subREQTailFile(p process) {
// fmt.Printf("*** %#v\n", proc)
go proc.spawnWorker(p.processes, p.natsConn)
}
func (s startup) subREQToSocket(p process) {
log.Printf("Starting write to socket subscriber: %#v\n", p.node)
sub := newSubject(REQToSocket, string(p.node))
proc := newProcess(p.natsConn, p.processes, p.toRingbufferCh, p.configuration, sub, p.errorCh, processKindSubscriber, []Node{"*"}, nil)
// fmt.Printf("*** %#v\n", proc)
go proc.spawnWorker(p.processes, p.natsConn)
}

View file

@ -121,6 +121,8 @@ const (
REQHttpGet Method = "REQHttpGet"
// Tail file
REQTailFile Method = "REQTailFile"
// Write to steward socket
REQToSocket Method = "REQToSocket"
)
// The mapping of all the method constants specified, what type
@ -179,6 +181,9 @@ func (m Method) GetMethodsAvailable() MethodsAvailable {
REQTailFile: methodREQTailFile{
commandOrEvent: EventACK,
},
REQToSocket: methodREQToSocket{
commandOrEvent: EventACK,
},
},
}
@ -189,7 +194,7 @@ func (m Method) GetMethodsAvailable() MethodsAvailable {
// the Stew client for knowing what of the req types are generally
// used as reply methods.
func (m Method) GetReplyMethods() []Method {
rm := []Method{REQToConsole, REQToFile, REQToFileAppend}
rm := []Method{REQToConsole, REQToFile, REQToFileAppend, REQToSocket}
return rm
}
@ -1096,3 +1101,26 @@ func (m methodREQnCliCommandCont) handler(proc process, message Message, node st
ackMsg := []byte("confirmed from: " + node + ": " + fmt.Sprint(message.ID))
return ackMsg, nil
}
// ---
type methodREQToSocket struct {
commandOrEvent CommandOrEvent
}
func (m methodREQToSocket) getKind() CommandOrEvent {
return m.commandOrEvent
}
func (m methodREQToSocket) handler(proc process, message Message, node string) ([]byte, error) {
for _, d := range message.Data {
// Write the data to the socket here.
fmt.Printf("Info: Data to write to socket: %v\n", d)
}
ackMsg := []byte("confirmed from: " + node + ": " + fmt.Sprint(message.ID))
return ackMsg, nil
}
// ----