Skip to content

Commit 1565cda

Browse files
committed
logger/enable topic
1 parent 02d4674 commit 1565cda

3 files changed

Lines changed: 104 additions & 3 deletions

File tree

backend/cmd/main.go

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ import (
3232
"github.com/HyperloopUPV-H8/h9-backend/pkg/broker"
3333
connection_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/connection"
3434
data_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/data"
35+
logger_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/logger"
3536
order_topic "github.com/HyperloopUPV-H8/h9-backend/pkg/broker/topics/order"
3637
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport"
3738
"github.com/HyperloopUPV-H8/h9-backend/pkg/transport/network/sniffer"
@@ -155,10 +156,12 @@ func main() {
155156
defer dataTopic.Stop()
156157
connectionTopic := connection_topic.NewUpdateTopic()
157158
orderTopic := order_topic.NewSendTopic()
159+
loggerTopic := logger_topic.NewEnableTopic()
158160

159161
broker.AddTopic(data_topic.UpdateName, dataTopic)
160162
broker.AddTopic(connection_topic.UpdateName, connectionTopic)
161163
broker.AddTopic(order_topic.SendName, orderTopic)
164+
broker.AddTopic(logger_topic.EnableName, loggerTopic)
162165

163166
connections := make(chan *websocket.Client)
164167
upgrader := websocket.NewUpgrader(connections)
@@ -280,6 +283,26 @@ func main() {
280283
if err != nil {
281284
fmt.Println("Error pushing record to logger: ", err)
282285
}
286+
case logger_topic.EnableName:
287+
status, ok := push.(*logger_topic.Status)
288+
if !ok {
289+
trace.Error().Any("push", push).Msg("error casting push to enable")
290+
fmt.Printf("Push Type: %v\n", push)
291+
return
292+
}
293+
294+
var err error
295+
if status.Enable() {
296+
err = loggerHandler.Start()
297+
} else {
298+
err = loggerHandler.Stop()
299+
}
300+
301+
if err != nil {
302+
status.Fulfill(!status.Enable())
303+
} else {
304+
status.Fulfill(status.Enable())
305+
}
283306
default:
284307
fmt.Printf("unknow topic %s\n", push.Topic())
285308
}

backend/pkg/broker/topics/logger/enable.go

Lines changed: 80 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,29 @@
11
package data
22

33
import (
4+
"encoding/json"
5+
"fmt"
6+
"sync/atomic"
7+
48
"github.com/HyperloopUPV-H8/h9-backend/pkg/abstraction"
59
"github.com/HyperloopUPV-H8/h9-backend/pkg/websocket"
610
)
711

812
const EnableName abstraction.BrokerTopic = "logger/enable"
13+
const ResponseName abstraction.BrokerTopic = "logger/response"
914

1015
type Enable struct {
11-
pool *websocket.Pool
12-
api abstraction.BrokerAPI
16+
isRunning *atomic.Bool
17+
pool *websocket.Pool
18+
api abstraction.BrokerAPI
19+
}
20+
21+
func NewEnableTopic() *Enable {
22+
enable := &Enable{
23+
isRunning: &atomic.Bool{},
24+
}
25+
enable.isRunning.Store(false)
26+
return enable
1327
}
1428

1529
func (enable *Enable) Topic() abstraction.BrokerTopic {
@@ -24,10 +38,74 @@ func (enable *Enable) Pull(request abstraction.BrokerRequest) (abstraction.Broke
2438
return nil, nil
2539
}
2640

41+
func (enable *Enable) ClientMessage(id websocket.ClientId, message *websocket.Message) {
42+
switch message.Topic {
43+
case EnableName:
44+
err := enable.handleToggle(id, message)
45+
if err != nil {
46+
fmt.Printf("error handling logger: %v\n", err)
47+
}
48+
}
49+
}
50+
51+
func (enable *Enable) handleToggle(id websocket.ClientId, message *websocket.Message) error {
52+
var request bool
53+
err := json.Unmarshal(message.Payload, &request)
54+
if err != nil {
55+
return err
56+
}
57+
58+
status := newStatus(request)
59+
go enable.api.UserPush(status)
60+
61+
go func() {
62+
enable.isRunning.Store(<-status.response)
63+
enable.broadcastState()
64+
}()
65+
return nil
66+
}
67+
68+
func (enable *Enable) broadcastState() error {
69+
payload, err := json.Marshal(enable.isRunning.Load())
70+
if err != nil {
71+
return err
72+
}
73+
74+
enable.pool.Broadcast(websocket.Message{
75+
Topic: ResponseName,
76+
Payload: payload,
77+
})
78+
return nil
79+
}
80+
2781
func (enable *Enable) SetPool(pool *websocket.Pool) {
2882
enable.pool = pool
2983
}
3084

3185
func (enable *Enable) SetAPI(api abstraction.BrokerAPI) {
3286
enable.api = api
3387
}
88+
89+
type Status struct {
90+
request bool
91+
response chan bool
92+
}
93+
94+
func newStatus(request bool) *Status {
95+
return &Status{
96+
request: request,
97+
response: make(chan bool),
98+
}
99+
}
100+
101+
func (status *Status) Topic() abstraction.BrokerTopic {
102+
return EnableName
103+
}
104+
105+
func (status *Status) Fulfill(response bool) {
106+
status.response <- response
107+
}
108+
109+
func (status *Status) Enable() bool {
110+
return status.request
111+
}

backend/pkg/logger/logger.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,7 +111,7 @@ func (logger *Logger) Stop() error {
111111
defer logger.subloggersLock.Unlock()
112112

113113
if !logger.running.CompareAndSwap(true, false) {
114-
fmt.Printf("Logger already stopped")
114+
fmt.Println("Logger already stopped")
115115
return nil
116116
}
117117

0 commit comments

Comments
 (0)