MetricsRestController.java 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135
  1. package com.provectus.kafka.ui.rest;
  2. import com.provectus.kafka.ui.api.ApiClustersApi;
  3. import com.provectus.kafka.ui.cluster.model.ConsumerPosition;
  4. import com.provectus.kafka.ui.cluster.service.ClusterService;
  5. import com.provectus.kafka.ui.model.*;
  6. import lombok.RequiredArgsConstructor;
  7. import org.apache.commons.lang3.tuple.Pair;
  8. import org.springframework.http.HttpStatus;
  9. import org.springframework.http.ResponseEntity;
  10. import org.springframework.web.bind.annotation.RestController;
  11. import org.springframework.web.server.ServerWebExchange;
  12. import reactor.core.publisher.Flux;
  13. import reactor.core.publisher.Mono;
  14. import javax.validation.Valid;
  15. import java.util.Collections;
  16. import java.util.List;
  17. import java.util.function.Function;
  18. @RestController
  19. @RequiredArgsConstructor
  20. public class MetricsRestController implements ApiClustersApi {
  21. private final ClusterService clusterService;
  22. @Override
  23. public Mono<ResponseEntity<Flux<Cluster>>> getClusters(ServerWebExchange exchange) {
  24. return Mono.just(ResponseEntity.ok(Flux.fromIterable(clusterService.getClusters())));
  25. }
  26. @Override
  27. public Mono<ResponseEntity<BrokerMetrics>> getBrokersMetrics(String clusterName, Integer id, ServerWebExchange exchange) {
  28. return clusterService.getBrokerMetrics(clusterName, id)
  29. .map(ResponseEntity::ok)
  30. .onErrorReturn(ResponseEntity.notFound().build());
  31. }
  32. @Override
  33. public Mono<ResponseEntity<ClusterMetrics>> getClusterMetrics(String clusterName, ServerWebExchange exchange) {
  34. return clusterService.getClusterMetrics(clusterName)
  35. .map(ResponseEntity::ok)
  36. .onErrorReturn(ResponseEntity.notFound().build());
  37. }
  38. @Override
  39. public Mono<ResponseEntity<ClusterStats>> getClusterStats(String clusterName, ServerWebExchange exchange) {
  40. return clusterService.getClusterStats(clusterName)
  41. .map(ResponseEntity::ok)
  42. .onErrorReturn(ResponseEntity.notFound().build());
  43. }
  44. @Override
  45. public Mono<ResponseEntity<Flux<Topic>>> getTopics(String clusterName, ServerWebExchange exchange) {
  46. return Mono.just(ResponseEntity.ok(Flux.fromIterable(clusterService.getTopics(clusterName))));
  47. }
  48. @Override
  49. public Mono<ResponseEntity<TopicDetails>> getTopicDetails(String clusterName, String topicName, ServerWebExchange exchange) {
  50. return Mono.just(
  51. clusterService.getTopicDetails(clusterName, topicName)
  52. .map(ResponseEntity::ok)
  53. .orElse(ResponseEntity.notFound().build())
  54. );
  55. }
  56. @Override
  57. public Mono<ResponseEntity<Flux<TopicConfig>>> getTopicConfigs(String clusterName, String topicName, ServerWebExchange exchange) {
  58. return Mono.just(
  59. clusterService.getTopicConfigs(clusterName, topicName)
  60. .map(Flux::fromIterable)
  61. .map(ResponseEntity::ok)
  62. .orElse(ResponseEntity.notFound().build())
  63. );
  64. }
  65. @Override
  66. public Mono<ResponseEntity<Flux<TopicMessage>>> getTopicMessages(String clusterName, String topicName, @Valid SeekType seekType, @Valid List<String> seekTo, @Valid Integer limit, @Valid String q, ServerWebExchange exchange) {
  67. return parseConsumerPosition(seekType, seekTo)
  68. .map(consumerPosition -> ResponseEntity.ok(clusterService.getMessages(clusterName, topicName, consumerPosition, q, limit)));
  69. }
  70. @Override
  71. public Mono<ResponseEntity<Topic>> createTopic(String clusterName, @Valid Mono<TopicFormData> topicFormData, ServerWebExchange exchange) {
  72. return clusterService.createTopic(clusterName, topicFormData)
  73. .map(s -> new ResponseEntity<>(s, HttpStatus.OK))
  74. .switchIfEmpty(Mono.just(ResponseEntity.notFound().build()));
  75. }
  76. @Override
  77. public Mono<ResponseEntity<Flux<Broker>>> getBrokers(String clusterName, ServerWebExchange exchange) {
  78. return Mono.just(ResponseEntity.ok(clusterService.getBrokers(clusterName)));
  79. }
  80. @Override
  81. public Mono<ResponseEntity<Flux<ConsumerGroup>>> getConsumerGroups(String clusterName, ServerWebExchange exchange) {
  82. return clusterService.getConsumerGroups(clusterName)
  83. .map(Flux::fromIterable)
  84. .map(ResponseEntity::ok)
  85. .switchIfEmpty(Mono.just(ResponseEntity.notFound().build())); // TODO: check behaviour on cluster not found and empty groups list
  86. }
  87. @Override
  88. public Mono<ResponseEntity<Flux<String>>> getSchemaSubjects(String clusterName, ServerWebExchange exchange) {
  89. Flux<String> subjects = clusterService.getSchemaSubjects(clusterName);
  90. return Mono.just(ResponseEntity.ok(subjects));
  91. }
  92. @Override
  93. public Mono<ResponseEntity<ConsumerGroupDetails>> getConsumerGroup(String clusterName, String consumerGroupId, ServerWebExchange exchange) {
  94. return clusterService.getConsumerGroupDetail(clusterName, consumerGroupId).map(ResponseEntity::ok);
  95. }
  96. @Override
  97. public Mono<ResponseEntity<Topic>> updateTopic(String clusterId, String topicName, @Valid Mono<TopicFormData> topicFormData, ServerWebExchange exchange) {
  98. return clusterService.updateTopic(clusterId, topicName, topicFormData).map(ResponseEntity::ok);
  99. }
  100. private Mono<ConsumerPosition> parseConsumerPosition(SeekType seekType, List<String> seekTo) {
  101. return Mono.justOrEmpty(seekTo)
  102. .defaultIfEmpty(Collections.emptyList())
  103. .flatMapIterable(Function.identity())
  104. .map(p -> {
  105. String[] splited = p.split("::");
  106. if (splited.length != 2) {
  107. throw new IllegalArgumentException("Wrong seekTo argument format. See API docs for details");
  108. }
  109. return Pair.of(Integer.parseInt(splited[0]), Long.parseLong(splited[1]));
  110. })
  111. .collectMap(Pair::getKey, Pair::getValue)
  112. .map(positions -> new ConsumerPosition(seekType != null ? seekType : SeekType.BEGINNING, positions));
  113. }
  114. }