command.go 2.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293
  1. package command
  2. //go:generate go run github.com/v2fly/v2ray-core/v4/common/errors/errorgen
  3. import (
  4. "context"
  5. "time"
  6. "google.golang.org/grpc"
  7. core "github.com/v2fly/v2ray-core/v4"
  8. "github.com/v2fly/v2ray-core/v4/common"
  9. "github.com/v2fly/v2ray-core/v4/features/routing"
  10. "github.com/v2fly/v2ray-core/v4/features/stats"
  11. )
  12. // routingServer is an implementation of RoutingService.
  13. type routingServer struct {
  14. router routing.Router
  15. routingStats stats.Channel
  16. }
  17. // NewRoutingServer creates a statistics service with statistics manager.
  18. func NewRoutingServer(router routing.Router, routingStats stats.Channel) RoutingServiceServer {
  19. return &routingServer{
  20. router: router,
  21. routingStats: routingStats,
  22. }
  23. }
  24. func (s *routingServer) TestRoute(ctx context.Context, request *TestRouteRequest) (*RoutingContext, error) {
  25. if request.RoutingContext == nil {
  26. return nil, newError("Invalid routing request.")
  27. }
  28. route, err := s.router.PickRoute(AsRoutingContext(request.RoutingContext))
  29. if err != nil {
  30. return nil, err
  31. }
  32. if request.PublishResult && s.routingStats != nil {
  33. ctx, _ := context.WithTimeout(context.Background(), 4*time.Second) // nolint: govet
  34. s.routingStats.Publish(ctx, route)
  35. }
  36. return AsProtobufMessage(request.FieldSelectors)(route), nil
  37. }
  38. func (s *routingServer) SubscribeRoutingStats(request *SubscribeRoutingStatsRequest, stream RoutingService_SubscribeRoutingStatsServer) error {
  39. if s.routingStats == nil {
  40. return newError("Routing statistics not enabled.")
  41. }
  42. genMessage := AsProtobufMessage(request.FieldSelectors)
  43. subscriber, err := stats.SubscribeRunnableChannel(s.routingStats)
  44. if err != nil {
  45. return err
  46. }
  47. defer stats.UnsubscribeClosableChannel(s.routingStats, subscriber)
  48. for {
  49. select {
  50. case value, ok := <-subscriber:
  51. if !ok {
  52. return newError("Upstream closed the subscriber channel.")
  53. }
  54. route, ok := value.(routing.Route)
  55. if !ok {
  56. return newError("Upstream sent malformed statistics.")
  57. }
  58. err := stream.Send(genMessage(route))
  59. if err != nil {
  60. return err
  61. }
  62. case <-stream.Context().Done():
  63. return stream.Context().Err()
  64. }
  65. }
  66. }
  67. func (s *routingServer) mustEmbedUnimplementedRoutingServiceServer() {}
  68. type service struct {
  69. v *core.Instance
  70. }
  71. func (s *service) Register(server *grpc.Server) {
  72. common.Must(s.v.RequireFeatures(func(router routing.Router, stats stats.Manager) {
  73. RegisterRoutingServiceServer(server, NewRoutingServer(router, nil))
  74. }))
  75. }
  76. func init() {
  77. common.Must(common.RegisterConfig((*Config)(nil), func(ctx context.Context, cfg interface{}) (interface{}, error) {
  78. s := core.MustFromContext(ctx)
  79. return &service{v: s}, nil
  80. }))
  81. }