chore: sync local updates
This commit is contained in:
82
mengyamonitor-backend-server/internal/ingest/server.go
Normal file
82
mengyamonitor-backend-server/internal/ingest/server.go
Normal file
@@ -0,0 +1,82 @@
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"errors"
|
||||
"io"
|
||||
"log"
|
||||
|
||||
"mengyamonitor-backend-server/internal/gen/mengyav1"
|
||||
"mengyamonitor-backend-server/internal/hub"
|
||||
"mengyamonitor-backend-server/internal/store"
|
||||
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
type Server struct {
|
||||
mengyav1.UnimplementedAgentIngestServer
|
||||
Store *store.Store
|
||||
Hub *hub.Hub
|
||||
}
|
||||
|
||||
func (s *Server) Connect(stream grpc.BidiStreamingServer[mengyav1.AgentMessage, mengyav1.ControlMessage]) error {
|
||||
var authedID string
|
||||
for {
|
||||
msg, err := stream.Recv()
|
||||
if err == io.EOF {
|
||||
return nil
|
||||
}
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
ctx := stream.Context()
|
||||
|
||||
switch {
|
||||
case msg.GetHello() != nil:
|
||||
h := msg.GetHello()
|
||||
srv, err := s.Store.VerifyAgent(ctx, h.GetServerId(), h.GetAgentKey())
|
||||
if err != nil {
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
log.Printf("ingest: auth failed server_id=%q: no such id (SERVER_ID must be the UUID from admin table, not the display name)", h.GetServerId())
|
||||
} else {
|
||||
log.Printf("ingest: auth failed server_id=%s: %v", h.GetServerId(), err)
|
||||
}
|
||||
return status.Errorf(codes.Unauthenticated, "invalid credentials")
|
||||
}
|
||||
authedID = srv.ID
|
||||
_ = stream.Send(&mengyav1.ControlMessage{
|
||||
Payload: &mengyav1.ControlMessage_Ack{
|
||||
Ack: &mengyav1.Ack{Message: "hello_ok"},
|
||||
},
|
||||
})
|
||||
|
||||
case msg.GetHeartbeat() != nil:
|
||||
if authedID == "" {
|
||||
continue
|
||||
}
|
||||
_ = stream.Send(&mengyav1.ControlMessage{
|
||||
Payload: &mengyav1.ControlMessage_Ack{
|
||||
Ack: &mengyav1.Ack{Message: "pong"},
|
||||
},
|
||||
})
|
||||
|
||||
case msg.GetMetrics() != nil:
|
||||
if authedID == "" {
|
||||
return status.Errorf(codes.FailedPrecondition, "send hello first")
|
||||
}
|
||||
m := msg.GetMetrics()
|
||||
if len(m.GetPayloadJson()) == 0 {
|
||||
continue
|
||||
}
|
||||
s.Hub.NotifyMetricsUpdate(authedID, []byte(m.GetPayloadJson()))
|
||||
_ = stream.Send(&mengyav1.ControlMessage{
|
||||
Payload: &mengyav1.ControlMessage_Ack{
|
||||
Ack: &mengyav1.Ack{Message: "metrics_ok"},
|
||||
},
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user