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"}, }, }) } } }