validate input to agent worker register (#2231)

* validate input to agent worker register

* prevent multiple registrations

* cleanup
This commit is contained in:
Paul Wells
2023-11-09 00:58:59 -08:00
committed by GitHub
parent 229c80ac17
commit aa797c749c
+23
View File
@@ -198,10 +198,27 @@ func (s *AgentHandler) HandleConnection(conn *websocket.Conn) {
}
func (s *AgentHandler) handleRegister(worker *worker, msg *livekit.RegisterWorkerRequest) {
if err := s.doHandleRegister(worker, msg); err != nil {
logger.Errorw("failed to register worker", err, "workerID", msg.WorkerId, "jobType", msg.Type)
worker.conn.Close()
}
}
func (s *AgentHandler) doHandleRegister(worker *worker, msg *livekit.RegisterWorkerRequest) error {
if msg.WorkerId == "" {
return errors.New("invalid worker id")
}
s.mu.Lock()
if worker.id != "" {
s.mu.Unlock()
return errors.New("worker already registered")
}
switch msg.Type {
case livekit.JobType_JT_ROOM:
worker.id = msg.WorkerId
worker.jobType = msg.Type
delete(s.unregistered, worker.conn)
s.roomWorkers[worker.id] = worker
@@ -216,6 +233,7 @@ func (s *AgentHandler) handleRegister(worker *worker, msg *livekit.RegisterWorke
case livekit.JobType_JT_PUBLISHER:
worker.id = msg.WorkerId
worker.jobType = msg.Type
delete(s.unregistered, worker.conn)
s.publisherWorkers[worker.id] = worker
@@ -227,6 +245,9 @@ func (s *AgentHandler) handleRegister(worker *worker, msg *livekit.RegisterWorke
s.publisherRegistered = true
}
}
default:
s.mu.Unlock()
return errors.New("invalid job type")
}
s.mu.Unlock()
@@ -241,6 +262,8 @@ func (s *AgentHandler) handleRegister(worker *worker, msg *livekit.RegisterWorke
if err != nil {
logger.Errorw("failed to write server message", err)
}
return nil
}
func (s *AgentHandler) handleAvailability(w *worker, msg *livekit.AvailabilityResponse) {