MetricsRestController.java 9.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207
  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.cluster.service.SchemaRegistryService;
  6. import com.provectus.kafka.ui.model.*;
  7. import lombok.RequiredArgsConstructor;
  8. import lombok.extern.log4j.Log4j2;
  9. import org.apache.commons.lang3.tuple.Pair;
  10. import org.springframework.http.HttpStatus;
  11. import org.springframework.http.ResponseEntity;
  12. import org.springframework.web.bind.annotation.RestController;
  13. import org.springframework.web.server.ServerWebExchange;
  14. import reactor.core.publisher.Flux;
  15. import reactor.core.publisher.Mono;
  16. import javax.validation.Valid;
  17. import java.util.Collections;
  18. import java.util.List;
  19. import java.util.function.Function;
  20. @RestController
  21. @RequiredArgsConstructor
  22. @Log4j2
  23. public class MetricsRestController implements ApiClustersApi {
  24. private final ClusterService clusterService;
  25. private final SchemaRegistryService schemaRegistryService;
  26. @Override
  27. public Mono<ResponseEntity<Flux<Cluster>>> getClusters(ServerWebExchange exchange) {
  28. return Mono.just(ResponseEntity.ok(Flux.fromIterable(clusterService.getClusters())));
  29. }
  30. @Override
  31. public Mono<ResponseEntity<BrokerMetrics>> getBrokersMetrics(String clusterName, Integer id, ServerWebExchange exchange) {
  32. return clusterService.getBrokerMetrics(clusterName, id)
  33. .map(ResponseEntity::ok)
  34. .onErrorReturn(ResponseEntity.notFound().build());
  35. }
  36. @Override
  37. public Mono<ResponseEntity<ClusterMetrics>> getClusterMetrics(String clusterName, ServerWebExchange exchange) {
  38. return clusterService.getClusterMetrics(clusterName)
  39. .map(ResponseEntity::ok)
  40. .onErrorReturn(ResponseEntity.notFound().build());
  41. }
  42. @Override
  43. public Mono<ResponseEntity<ClusterStats>> getClusterStats(String clusterName, ServerWebExchange exchange) {
  44. return clusterService.getClusterStats(clusterName)
  45. .map(ResponseEntity::ok)
  46. .onErrorReturn(ResponseEntity.notFound().build());
  47. }
  48. @Override
  49. public Mono<ResponseEntity<Flux<Topic>>> getTopics(String clusterName, ServerWebExchange exchange) {
  50. return Mono.just(ResponseEntity.ok(Flux.fromIterable(clusterService.getTopics(clusterName))));
  51. }
  52. @Override
  53. public Mono<ResponseEntity<TopicDetails>> getTopicDetails(String clusterName, String topicName, ServerWebExchange exchange) {
  54. return Mono.just(
  55. clusterService.getTopicDetails(clusterName, topicName)
  56. .map(ResponseEntity::ok)
  57. .orElse(ResponseEntity.notFound().build())
  58. );
  59. }
  60. @Override
  61. public Mono<ResponseEntity<Flux<TopicConfig>>> getTopicConfigs(String clusterName, String topicName, ServerWebExchange exchange) {
  62. return Mono.just(
  63. clusterService.getTopicConfigs(clusterName, topicName)
  64. .map(Flux::fromIterable)
  65. .map(ResponseEntity::ok)
  66. .orElse(ResponseEntity.notFound().build())
  67. );
  68. }
  69. @Override
  70. 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) {
  71. return parseConsumerPosition(seekType, seekTo)
  72. .map(consumerPosition -> ResponseEntity.ok(clusterService.getMessages(clusterName, topicName, consumerPosition, q, limit)));
  73. }
  74. @Override
  75. public Mono<ResponseEntity<Topic>> createTopic(String clusterName, @Valid Mono<TopicFormData> topicFormData, ServerWebExchange exchange) {
  76. return clusterService.createTopic(clusterName, topicFormData)
  77. .map(s -> new ResponseEntity<>(s, HttpStatus.OK))
  78. .switchIfEmpty(Mono.just(ResponseEntity.notFound().build()));
  79. }
  80. @Override
  81. public Mono<ResponseEntity<Flux<Broker>>> getBrokers(String clusterName, ServerWebExchange exchange) {
  82. return Mono.just(ResponseEntity.ok(clusterService.getBrokers(clusterName)));
  83. }
  84. @Override
  85. public Mono<ResponseEntity<Flux<ConsumerGroup>>> getConsumerGroups(String clusterName, ServerWebExchange exchange) {
  86. return clusterService.getConsumerGroups(clusterName)
  87. .map(Flux::fromIterable)
  88. .map(ResponseEntity::ok)
  89. .switchIfEmpty(Mono.just(ResponseEntity.notFound().build())); // TODO: check behaviour on cluster not found and empty groups list
  90. }
  91. @Override
  92. public Mono<ResponseEntity<SchemaSubject>> getLatestSchema(String clusterName, String subject, ServerWebExchange exchange) {
  93. return schemaRegistryService.getLatestSchemaVersionBySubject(clusterName, subject).map(ResponseEntity::ok);
  94. }
  95. @Override
  96. public Mono<ResponseEntity<SchemaSubject>> getSchemaByVersion(String clusterName, String subject, Integer version, ServerWebExchange exchange) {
  97. return schemaRegistryService.getSchemaSubjectByVersion(clusterName, subject, version).map(ResponseEntity::ok);
  98. }
  99. @Override
  100. public Mono<ResponseEntity<Flux<SchemaSubject>>> getSchemas(String clusterName, ServerWebExchange exchange) {
  101. Flux<SchemaSubject> subjects = schemaRegistryService.getAllLatestVersionSchemas(clusterName);
  102. return Mono.just(ResponseEntity.ok(subjects));
  103. }
  104. @Override
  105. public Mono<ResponseEntity<Flux<SchemaSubject>>> getAllVersionsBySubject(String clusterName, String subjectName, ServerWebExchange exchange) {
  106. Flux<SchemaSubject> schemas = schemaRegistryService.getAllVersionsBySubject(clusterName, subjectName);
  107. return Mono.just(ResponseEntity.ok(schemas));
  108. }
  109. @Override
  110. public Mono<ResponseEntity<Void>> deleteLatestSchema(String clusterName, String subject, ServerWebExchange exchange) {
  111. return schemaRegistryService.deleteLatestSchemaSubject(clusterName, subject);
  112. }
  113. @Override
  114. public Mono<ResponseEntity<Void>> deleteSchemaByVersion(String clusterName, String subjectName, Integer version, ServerWebExchange exchange) {
  115. return schemaRegistryService.deleteSchemaSubjectByVersion(clusterName, subjectName, version);
  116. }
  117. @Override
  118. public Mono<ResponseEntity<Void>> deleteSchema(String clusterName, String subjectName, ServerWebExchange exchange) {
  119. return schemaRegistryService.deleteSchemaSubjectEntirely(clusterName, subjectName);
  120. }
  121. @Override
  122. public Mono<ResponseEntity<SchemaSubject>> createNewSchema(String clusterName,
  123. @Valid Mono<NewSchemaSubject> newSchemaSubject,
  124. ServerWebExchange exchange) {
  125. return schemaRegistryService
  126. .registerNewSchema(clusterName, newSchemaSubject)
  127. .map(ResponseEntity::ok);
  128. }
  129. @Override
  130. public Mono<ResponseEntity<ConsumerGroupDetails>> getConsumerGroup(String clusterName, String consumerGroupId, ServerWebExchange exchange) {
  131. return clusterService.getConsumerGroupDetail(clusterName, consumerGroupId).map(ResponseEntity::ok);
  132. }
  133. @Override
  134. public Mono<ResponseEntity<Topic>> updateTopic(String clusterId, String topicName, @Valid Mono<TopicFormData> topicFormData, ServerWebExchange exchange) {
  135. return clusterService.updateTopic(clusterId, topicName, topicFormData).map(ResponseEntity::ok);
  136. }
  137. @Override
  138. public Mono<ResponseEntity<CompatibilityLevel>> getGlobalSchemaCompatibilityLevel(String clusterName, ServerWebExchange exchange) {
  139. return schemaRegistryService.getGlobalSchemaCompatibilityLevel(clusterName)
  140. .map(ResponseEntity::ok)
  141. .defaultIfEmpty(ResponseEntity.notFound().build());
  142. }
  143. @Override
  144. public Mono<ResponseEntity<Void>> updateGlobalSchemaCompatibilityLevel(String clusterName, @Valid Mono<CompatibilityLevel> compatibilityLevel, ServerWebExchange exchange) {
  145. log.info("Updating schema compatibility globally");
  146. return schemaRegistryService.updateSchemaCompatibility(clusterName, compatibilityLevel)
  147. .map(ResponseEntity::ok);
  148. }
  149. @Override
  150. public Mono<ResponseEntity<CompatibilityCheckResponse>> checkSchemaCompatibility(String clusterName, String subject,
  151. @Valid Mono<NewSchemaSubject> newSchemaSubject,
  152. ServerWebExchange exchange) {
  153. return schemaRegistryService.checksSchemaCompatibility(clusterName, subject, newSchemaSubject)
  154. .map(ResponseEntity::ok);
  155. }
  156. @Override
  157. public Mono<ResponseEntity<Void>> updateSchemaCompatibilityLevel(String clusterName, String subject, @Valid Mono<CompatibilityLevel> compatibilityLevel, ServerWebExchange exchange) {
  158. log.info("Updating schema compatibility for subject: {}", subject);
  159. return schemaRegistryService.updateSchemaCompatibility(clusterName, subject, compatibilityLevel)
  160. .map(ResponseEntity::ok);
  161. }
  162. private Mono<ConsumerPosition> parseConsumerPosition(SeekType seekType, List<String> seekTo) {
  163. return Mono.justOrEmpty(seekTo)
  164. .defaultIfEmpty(Collections.emptyList())
  165. .flatMapIterable(Function.identity())
  166. .map(p -> {
  167. String[] splited = p.split("::");
  168. if (splited.length != 2) {
  169. throw new IllegalArgumentException("Wrong seekTo argument format. See API docs for details");
  170. }
  171. return Pair.of(Integer.parseInt(splited[0]), Long.parseLong(splited[1]));
  172. })
  173. .collectMap(Pair::getKey, Pair::getValue)
  174. .map(positions -> new ConsumerPosition(seekType != null ? seekType : SeekType.BEGINNING, positions));
  175. }
  176. }