From a4f259f3200300bede94ac4e2b27a078f44c73bf Mon Sep 17 00:00:00 2001 From: Freedom <459102951@qq.com> Date: Tue, 14 Jul 2026 10:18:50 +0800 Subject: [PATCH 1/6] feat(mq): improve Kafka peek defaults and localize admin console Make Kafka message peek read all partitions from earliest when partition/offset are omitted, and wire MQ admin UI strings through i18n for all supported locales. Co-authored-by: Cursor --- .../java/com/dbx/agent/kafka/KafkaAgent.java | 170 ++++-- .../com/dbx/agent/kafka/KafkaAgentTest.java | 29 + .../connection/ConnectionDialog.vue | 88 +-- .../desktop/src/components/mq/BrokerPanel.vue | 54 +- .../src/components/mq/MonitoringPanel.vue | 220 ++++---- .../src/components/mq/MqAdminConsole.vue | 32 +- .../src/components/mq/NamespacesPanel.vue | 47 +- .../src/components/mq/PermissionsPanel.vue | 114 ++-- .../src/components/mq/PoliciesPanel.vue | 80 +-- .../components/mq/ProducerConsumerPanel.vue | 84 +-- .../desktop/src/components/mq/RawApiPanel.vue | 76 +-- .../src/components/mq/SendMessagePanel.vue | 131 +++-- .../src/components/mq/SubscriptionsPanel.vue | 124 ++--- .../src/components/mq/TenantsPanel.vue | 88 +-- apps/desktop/src/i18n/locales/en.ts | 507 ++++++++++++++++++ apps/desktop/src/i18n/locales/es.ts | 507 ++++++++++++++++++ apps/desktop/src/i18n/locales/it.ts | 507 ++++++++++++++++++ apps/desktop/src/i18n/locales/ja.ts | 507 ++++++++++++++++++ apps/desktop/src/i18n/locales/pt-BR.ts | 507 ++++++++++++++++++ apps/desktop/src/i18n/locales/zh-CN.ts | 507 ++++++++++++++++++ apps/desktop/src/i18n/locales/zh-TW.ts | 507 ++++++++++++++++++ .../src/lib/__tests__/mq/mqTenantForm.spec.ts | 4 +- .../lib/__tests__/mq/mqTokenErrors.spec.ts | 15 +- apps/desktop/src/lib/mq/mqTenantForm.ts | 4 +- apps/desktop/src/lib/mq/mqTokenErrors.ts | 17 +- crates/dbx-core/src/mq/adapters/kafka.rs | 64 ++- crates/dbx-core/src/mq/types.rs | 4 +- 27 files changed, 4372 insertions(+), 622 deletions(-) diff --git a/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java b/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java index ed732e0caa..6e03d95020 100644 --- a/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java +++ b/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java @@ -793,9 +793,9 @@ static OffsetSpec offsetSpecForPosition(String position, Long timestampMs) { private static Object peekMessages(JsonObject params) throws Exception { String topic = stringOrEmpty(params, "topic"); - int partition = intOrDefault(params, "partition", 0); - long offset = longOrDefault(params, "offset", 0); - int count = intOrDefault(params, "count", 10); + Integer partition = integerOrNull(params, "partition"); + Long offset = longOrNull(params, "offset"); + int count = Math.max(1, intOrDefault(params, "count", 10)); // Build a temporary consumer for peeking (no commit) Properties props = new Properties(); @@ -818,52 +818,134 @@ private static Object peekMessages(JsonObject params) throws Exception { props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none"); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, count); - TopicPartition tp = new TopicPartition(topic, partition); try (KafkaConsumer consumer = new KafkaConsumer<>(props)) { - consumer.assign(Collections.singletonList(tp)); + List candidatePartitions = resolvePeekPartitions(consumer, topic, partition); + if (candidatePartitions.isEmpty()) { + return Collections.singletonMap("messages", Collections.emptyList()); + } + Map beginningOffsets = - consumer.beginningOffsets(Collections.singletonList(tp), Duration.ofSeconds(5)); + consumer.beginningOffsets(candidatePartitions, Duration.ofSeconds(5)); Map endOffsets = - consumer.endOffsets(Collections.singletonList(tp), Duration.ofSeconds(5)); - long beginningOffset = beginningOffsets.getOrDefault(tp, 0L); - long endOffset = endOffsets.getOrDefault(tp, beginningOffset); - Long seekOffset = normalizePeekOffset(offset, beginningOffset, endOffset); - if (seekOffset == null) { + consumer.endOffsets(candidatePartitions, Duration.ofSeconds(5)); + + List readablePartitions = new ArrayList<>(); + Map seekOffsets = new LinkedHashMap<>(); + for (TopicPartition tp : candidatePartitions) { + long beginningOffset = beginningOffsets.getOrDefault(tp, 0L); + long endOffset = endOffsets.getOrDefault(tp, beginningOffset); + long requestedOffset = offset != null ? offset : beginningOffset; + Long seekOffset = normalizePeekOffset(requestedOffset, beginningOffset, endOffset); + if (seekOffset == null) { + continue; + } + readablePartitions.add(tp); + seekOffsets.put(tp, seekOffset); + } + if (readablePartitions.isEmpty()) { return Collections.singletonMap("messages", Collections.emptyList()); } - consumer.seek(tp, seekOffset); + + consumer.assign(readablePartitions); + for (Map.Entry entry : seekOffsets.entrySet()) { + consumer.seek(entry.getKey(), entry.getValue()); + } List> messages = new ArrayList<>(); - ConsumerRecords records = consumer.poll(Duration.ofSeconds(5)); - for (ConsumerRecord record : records) { - if (messages.size() >= count) break; - Map msg = new LinkedHashMap<>(); - msg.put("topic", record.topic()); - msg.put("partition", record.partition()); - msg.put("offset", record.offset()); - msg.put("timestamp", record.timestamp()); - msg.put("key", record.key()); - // Headers - Map headers = new LinkedHashMap<>(); - record.headers().forEach(h -> - headers.put(h.key(), new String(h.value(), StandardCharsets.UTF_8))); - msg.put("headers", headers); - // Payload - if (record.value() != null) { - msg.put("payloadBase64", Base64.getEncoder().encodeToString(record.value())); - String text = tryDecodeUtf8(record.value()); - if (text != null) { - msg.put("payloadText", text); + long deadlineNs = System.nanoTime() + Duration.ofSeconds(5).toNanos(); + while (messages.size() < count && System.nanoTime() < deadlineNs) { + ConsumerRecords records = consumer.poll(Duration.ofMillis(500)); + if (records.isEmpty()) { + break; + } + for (ConsumerRecord record : records) { + messages.add(peekedMessageFromRecord(record)); + if (messages.size() >= count) { + break; } - } else { - msg.put("payloadBase64", ""); } - messages.add(msg); + } + sortPeekedMessages(messages); + if (messages.size() > count) { + messages = new ArrayList<>(messages.subList(0, count)); } return Collections.singletonMap("messages", messages); } } + /** When partition is null, peek across every partition of the topic. */ + static List resolvePeekPartitions( + KafkaConsumer consumer, + String topic, + Integer partition + ) { + if (partition != null) { + return resolvePeekPartitions(topic, partition, Collections.emptyList()); + } + List infos = consumer.partitionsFor(topic, Duration.ofSeconds(5)); + if (infos == null || infos.isEmpty()) { + return Collections.emptyList(); + } + List available = infos.stream().map(PartitionInfo::partition).collect(Collectors.toList()); + return resolvePeekPartitions(topic, null, available); + } + + static List resolvePeekPartitions(String topic, Integer partition, List availablePartitions) { + if (partition != null) { + return Collections.singletonList(new TopicPartition(topic, partition)); + } + if (availablePartitions == null || availablePartitions.isEmpty()) { + return Collections.emptyList(); + } + return availablePartitions.stream() + .sorted() + .map(id -> new TopicPartition(topic, id)) + .collect(Collectors.toList()); + } + + static void sortPeekedMessages(List> messages) { + messages.sort((left, right) -> { + long leftTs = ((Number) left.getOrDefault("timestamp", 0L)).longValue(); + long rightTs = ((Number) right.getOrDefault("timestamp", 0L)).longValue(); + int byTs = Long.compare(leftTs, rightTs); + if (byTs != 0) { + return byTs; + } + int leftPartition = ((Number) left.getOrDefault("partition", 0)).intValue(); + int rightPartition = ((Number) right.getOrDefault("partition", 0)).intValue(); + int byPartition = Integer.compare(leftPartition, rightPartition); + if (byPartition != 0) { + return byPartition; + } + long leftOffset = ((Number) left.getOrDefault("offset", 0L)).longValue(); + long rightOffset = ((Number) right.getOrDefault("offset", 0L)).longValue(); + return Long.compare(leftOffset, rightOffset); + }); + } + + private static Map peekedMessageFromRecord(ConsumerRecord record) { + Map msg = new LinkedHashMap<>(); + msg.put("topic", record.topic()); + msg.put("partition", record.partition()); + msg.put("offset", record.offset()); + msg.put("timestamp", record.timestamp()); + msg.put("key", record.key()); + Map headers = new LinkedHashMap<>(); + record.headers().forEach(h -> + headers.put(h.key(), new String(h.value(), StandardCharsets.UTF_8))); + msg.put("headers", headers); + if (record.value() != null) { + msg.put("payloadBase64", Base64.getEncoder().encodeToString(record.value())); + String text = tryDecodeUtf8(record.value()); + if (text != null) { + msg.put("payloadText", text); + } + } else { + msg.put("payloadBase64", ""); + } + return msg; + } + private static Object sendMessage(JsonObject params) throws Exception { if (producer == null) { throw new IllegalStateException("Producer is not initialized. Call connect first."); @@ -1229,14 +1311,24 @@ private static String stringOrDefault(JsonObject object, String key, String fall return value == null ? fallback : value; } - private static int intOrDefault(JsonObject object, String key, int fallback) { + private static Integer integerOrNull(JsonObject object, String key) { JsonElement element = object.get(key); - return element == null || element.isJsonNull() ? fallback : element.getAsInt(); + return element == null || element.isJsonNull() ? null : element.getAsInt(); } - private static long longOrDefault(JsonObject object, String key, long fallback) { + private static Long longOrNull(JsonObject object, String key) { JsonElement element = object.get(key); - return element == null || element.isJsonNull() ? fallback : element.getAsLong(); + return element == null || element.isJsonNull() ? null : element.getAsLong(); + } + + private static int intOrDefault(JsonObject object, String key, int fallback) { + Integer value = integerOrNull(object, key); + return value == null ? fallback : value; + } + + private static long longOrDefault(JsonObject object, String key, long fallback) { + Long value = longOrNull(object, key); + return value == null ? fallback : value; } private static boolean boolOrDefault(JsonObject object, String key, boolean fallback) { diff --git a/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java b/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java index 1b4ad46db1..7003d1f574 100644 --- a/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java +++ b/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java @@ -4,6 +4,7 @@ import static org.junit.jupiter.api.Assertions.assertNull; import com.google.gson.JsonParser; +import java.util.List; import java.util.Map; import java.util.Properties; import org.junit.jupiter.api.Test; @@ -34,6 +35,34 @@ void returnsNoSeekOffsetWhenTopicHasNoReadableMessages() { assertNull(KafkaAgent.normalizePeekOffset(0, 5, 5)); } + @Test + void resolvePeekPartitionsUsesSinglePartitionWhenSpecified() { + var partitions = KafkaAgent.resolvePeekPartitions("events", 2, List.of(0, 1, 2)); + assertEquals(1, partitions.size()); + assertEquals(2, partitions.get(0).partition()); + assertEquals("events", partitions.get(0).topic()); + } + + @Test + void resolvePeekPartitionsUsesAllPartitionsWhenUnspecified() { + var partitions = KafkaAgent.resolvePeekPartitions("events", null, List.of(2, 0, 1)); + assertEquals(List.of(0, 1, 2), partitions.stream().map(org.apache.kafka.common.TopicPartition::partition).toList()); + } + + @Test + void sortPeekedMessagesOrdersByTimestampThenPartitionThenOffset() { + var messages = new java.util.ArrayList>(); + messages.add(Map.of("timestamp", 20L, "partition", 1, "offset", 1L)); + messages.add(Map.of("timestamp", 10L, "partition", 0, "offset", 5L)); + messages.add(Map.of("timestamp", 10L, "partition", 0, "offset", 2L)); + messages.add(Map.of("timestamp", 10L, "partition", 1, "offset", 0L)); + KafkaAgent.sortPeekedMessages(messages); + assertEquals(2L, messages.get(0).get("offset")); + assertEquals(5L, messages.get(1).get("offset")); + assertEquals(1, messages.get(2).get("partition")); + assertEquals(20L, messages.get(3).get("timestamp")); + } + @Test void appliesKerberosKafkaProperties() { Properties props = new Properties(); diff --git a/apps/desktop/src/components/connection/ConnectionDialog.vue b/apps/desktop/src/components/connection/ConnectionDialog.vue index 595670dbe8..5c3e8cdce8 100644 --- a/apps/desktop/src/components/connection/ConnectionDialog.vue +++ b/apps/desktop/src/components/connection/ConnectionDialog.vue @@ -465,17 +465,17 @@ const mqTlsSkipVerify = ref(false); const mqPinnedVersion = ref(pinnedVersionToSelection(undefined)); const mqTokenSigningMode = ref("none"); const mqTokenSigningKey = ref(""); -const mqSystemOptions: Array<{ value: MqSystemKind; label: string }> = [ - { value: "pulsar", label: "Apache Pulsar" }, - { value: "kafka", label: "Apache Kafka" }, -]; -const mqKafkaSecurityProtocolOptions = [ - { value: MQ_KAFKA_SECURITY_PROTOCOL_AUTO, label: "Auto" }, +const mqSystemOptions = computed(() => [ + { value: "pulsar" as const, label: t("connection.mqSystemPulsar") }, + { value: "kafka" as const, label: t("connection.mqSystemKafka") }, +]); +const mqKafkaSecurityProtocolOptions = computed(() => [ + { value: MQ_KAFKA_SECURITY_PROTOCOL_AUTO, label: t("connection.mqSecurityAuto") }, { value: "PLAINTEXT", label: "PLAINTEXT" }, { value: "SSL", label: "SSL" }, { value: "SASL_PLAINTEXT", label: "SASL_PLAINTEXT" }, { value: "SASL_SSL", label: "SASL_SSL" }, -]; +]); const mqKafkaSaslMechanismOptions = [ { value: "PLAIN", label: "PLAIN" }, { value: "SCRAM-SHA-256", label: "SCRAM-SHA-256" }, @@ -932,9 +932,9 @@ function buildMqAuth(): MqAuth { case "oauth2": return { kind: "oauth2", - issuerUrl: requireMqField(mqOauthIssuerUrl.value, "OAuth2 auth requires an issuer URL"), - clientId: requireMqField(mqOauthClientId.value, "OAuth2 auth requires a client ID"), - clientSecret: requireMqField(mqOauthClientSecret.value, "OAuth2 auth requires a client secret"), + issuerUrl: requireMqField(mqOauthIssuerUrl.value, t("connection.mqOauthIssuerRequired")), + clientId: requireMqField(mqOauthClientId.value, t("connection.mqOauthClientIdRequired")), + clientSecret: requireMqField(mqOauthClientSecret.value, t("connection.mqOauthClientSecretRequired")), audience: mqOauthAudience.value.trim() || undefined, scope: mqOauthScope.value.trim() || undefined, }; @@ -953,7 +953,7 @@ function buildMqTokenSigning() { if (mqTokenSigningMode.value === "none") return undefined; return { algorithm: mqTokenSigningMode.value, - key: requireMqField(mqTokenSigningKey.value, "Broker token signing key is required"), + key: requireMqField(mqTokenSigningKey.value, t("connection.mqTokenSigningKeyRequired")), }; } @@ -987,7 +987,7 @@ function buildMqAdminConfig(): MqAdminConfig { return { systemKind: mqSystemKind.value, - adminUrl: requireMqField(mqAdminUrl.value, "MQ Admin URL is required"), + adminUrl: requireMqField(mqAdminUrl.value, t("connection.mqAdminUrlRequired")), auth: buildMqAuth(), tlsSkipVerify: mqTlsSkipVerify.value || undefined, pinnedVersion: selectionToPinnedVersion(mqPinnedVersion.value), @@ -1231,7 +1231,7 @@ function applyMqAdminUrl(config: LegacyConnectionConfig, adminUrl: string) { try { parsed = new URL(adminUrl); } catch { - throw new Error("MQ Admin URL is invalid"); + throw new Error(t("connection.mqAdminUrlInvalid")); } const port = Number(parsed.port) || (parsed.protocol === "https:" ? 443 : 8080); config.host = parsed.hostname; @@ -1241,12 +1241,12 @@ function applyMqAdminUrl(config: LegacyConnectionConfig, adminUrl: string) { function applyMqKafkaBootstrapServers(config: LegacyConnectionConfig, bootstrapServers: string, securityProtocol?: string) { const first = normalizeKafkaBootstrapServers(bootstrapServers).split(",")[0]; - if (!first) throw new Error("Kafka bootstrap servers are required"); + if (!first) throw new Error(t("connection.mqBootstrapServersRequired")); let parsed: URL; try { parsed = new URL(`kafka://${first}`); } catch { - throw new Error("Kafka bootstrap servers are invalid"); + throw new Error(t("connection.mqBootstrapServersInvalid")); } config.host = parsed.hostname; config.port = Number(parsed.port) || 9092; @@ -4238,7 +4238,7 @@ function openExternalUrl(url: string) {
- +
- +
- +
- + + {{ t("mqRaw.queryParams") }} +
-

响应

+

{{ t("mqRaw.response") }}

HTTP {{ response.status }} - 文本响应 + {{ t("mqRaw.textResponse") }}
{{ response.text || formattedBody }}
-
尚未发送请求
+
{{ t("mqRaw.noRequestYet") }}
diff --git a/apps/desktop/src/components/mq/SendMessagePanel.vue b/apps/desktop/src/components/mq/SendMessagePanel.vue index 2fae77e4ee..42450b381a 100644 --- a/apps/desktop/src/components/mq/SendMessagePanel.vue +++ b/apps/desktop/src/components/mq/SendMessagePanel.vue @@ -1,5 +1,6 @@ + + diff --git a/apps/desktop/src/i18n/locales/en.ts b/apps/desktop/src/i18n/locales/en.ts index a83f700a05..e2f05d26f4 100644 --- a/apps/desktop/src/i18n/locales/en.ts +++ b/apps/desktop/src/i18n/locales/en.ts @@ -3625,6 +3625,23 @@ export default { project: "Project", openSource: "Open-source repository", officialDocs: "Official docs", + changelogTitle: "Changelog", + changelogDescription: "Browse what's new in each release.", + changelogLoading: "Loading changelog...", + changelogLoadFailed: "Failed to load changelog.", + changelogEmpty: "No changelog entries yet.", + changelogPublishedOn: "Published on {date}", + changelogNew: "NEW", + changelogLoadMore: "Load more", + changelogOpenWebsite: "View on website", + changelogOpenGitHub: "GitHub Release", + changelogNoDetails: "See GitHub Release for details.", + changelogRetry: "Retry", + changelogSectionAdded: "New Features", + changelogSectionImproved: "Improvements", + changelogSectionFixed: "Bug Fixes", + changelogSectionChanged: "Changes", + changelogSectionRemoved: "Removed", shortcutUppercaseSelection: "Convert selection to uppercase", shortcutLowercaseSelection: "Convert selection to lowercase", shortcutExPasteSqlInCondition: "ExPaste: paste as IN condition", diff --git a/apps/desktop/src/i18n/locales/es.ts b/apps/desktop/src/i18n/locales/es.ts index b1373e64c9..9d53c2fe8c 100644 --- a/apps/desktop/src/i18n/locales/es.ts +++ b/apps/desktop/src/i18n/locales/es.ts @@ -3415,6 +3415,23 @@ export default withEnglishFallback({ project: "Proyecto", openSource: "Repositorio de código abierto", officialDocs: "Documentación oficial", + changelogTitle: "Registro de cambios", + changelogDescription: "Consulta las novedades de cada versión.", + changelogLoading: "Cargando el registro de cambios...", + changelogLoadFailed: "No se pudo cargar el registro de cambios.", + changelogEmpty: "Aún no hay entradas en el registro de cambios.", + changelogPublishedOn: "Publicado el {date}", + changelogNew: "NEW", + changelogLoadMore: "Cargar más", + changelogOpenWebsite: "Ver en el sitio web", + changelogOpenGitHub: "GitHub Release", + changelogNoDetails: "Consulta GitHub Release para más detalles.", + changelogRetry: "Reintentar", + changelogSectionAdded: "Novedades", + changelogSectionImproved: "Mejoras", + changelogSectionFixed: "Correcciones", + changelogSectionChanged: "Cambios", + changelogSectionRemoved: "Eliminado", shortcutUppercaseSelection: "Convertir selección a mayúsculas", shortcutLowercaseSelection: "Convertir selección a minúsculas", shortcutExPasteSqlInCondition: "ExPaste: pegar como condición IN", diff --git a/apps/desktop/src/i18n/locales/it.ts b/apps/desktop/src/i18n/locales/it.ts index c8c503b5a0..cebf11465e 100644 --- a/apps/desktop/src/i18n/locales/it.ts +++ b/apps/desktop/src/i18n/locales/it.ts @@ -3413,6 +3413,23 @@ export default withEnglishFallback({ project: "Progetto", openSource: "Repository open-source", officialDocs: "Documenti ufficiali", + changelogTitle: "Registro modifiche", + changelogDescription: "Consulta le novità di ogni versione.", + changelogLoading: "Caricamento del registro modifiche...", + changelogLoadFailed: "Impossibile caricare il registro modifiche.", + changelogEmpty: "Nessuna voce nel registro modifiche.", + changelogPublishedOn: "Pubblicato il {date}", + changelogNew: "NEW", + changelogLoadMore: "Carica altro", + changelogOpenWebsite: "Apri sul sito", + changelogOpenGitHub: "GitHub Release", + changelogNoDetails: "Vedi GitHub Release per i dettagli.", + changelogRetry: "Riprova", + changelogSectionAdded: "Nuove funzionalità", + changelogSectionImproved: "Miglioramenti", + changelogSectionFixed: "Correzioni", + changelogSectionChanged: "Modifiche", + changelogSectionRemoved: "Rimossi", shortcutUppercaseSelection: "Converti selezione in maiuscolo", shortcutLowercaseSelection: "Converti selezione in minuscolo", shortcutExPasteSqlInCondition: "ExPaste: incolla come condizione IN", diff --git a/apps/desktop/src/i18n/locales/ja.ts b/apps/desktop/src/i18n/locales/ja.ts index 52bb66b6aa..e99f885303 100644 --- a/apps/desktop/src/i18n/locales/ja.ts +++ b/apps/desktop/src/i18n/locales/ja.ts @@ -3414,6 +3414,23 @@ export default withEnglishFallback({ project: "プロジェクト", openSource: "オープンソースリポジトリ", officialDocs: "公式ドキュメント", + changelogTitle: "更新履歴", + changelogDescription: "各バージョンの更新内容を確認できます。", + changelogLoading: "更新履歴を読み込み中...", + changelogLoadFailed: "更新履歴の読み込みに失敗しました。", + changelogEmpty: "更新履歴はまだありません。", + changelogPublishedOn: "公開日 {date}", + changelogNew: "NEW", + changelogLoadMore: "さらに読み込む", + changelogOpenWebsite: "公式サイトで見る", + changelogOpenGitHub: "GitHub Release", + changelogNoDetails: "詳細は GitHub Release を参照してください。", + changelogRetry: "再試行", + changelogSectionAdded: "新機能", + changelogSectionImproved: "改善", + changelogSectionFixed: "不具合修正", + changelogSectionChanged: "変更", + changelogSectionRemoved: "削除", shortcutUppercaseSelection: "選択範囲を大文字に変換", shortcutLowercaseSelection: "選択範囲を小文字に変換", shortcutExPasteSqlInCondition: "ExPaste: IN条件として貼り付け", diff --git a/apps/desktop/src/i18n/locales/pt-BR.ts b/apps/desktop/src/i18n/locales/pt-BR.ts index ecb162d147..5d31264bb2 100644 --- a/apps/desktop/src/i18n/locales/pt-BR.ts +++ b/apps/desktop/src/i18n/locales/pt-BR.ts @@ -3415,6 +3415,23 @@ export default withEnglishFallback({ project: "Projeto", openSource: "Repositório de código aberto", officialDocs: "Documentação oficial", + changelogTitle: "Registro de alterações", + changelogDescription: "Veja o que há de novo em cada versão.", + changelogLoading: "Carregando o registro de alterações...", + changelogLoadFailed: "Falha ao carregar o registro de alterações.", + changelogEmpty: "Ainda não há entradas no registro de alterações.", + changelogPublishedOn: "Publicado em {date}", + changelogNew: "NEW", + changelogLoadMore: "Carregar mais", + changelogOpenWebsite: "Ver no site", + changelogOpenGitHub: "GitHub Release", + changelogNoDetails: "Veja o GitHub Release para detalhes.", + changelogRetry: "Tentar novamente", + changelogSectionAdded: "Novidades", + changelogSectionImproved: "Melhorias", + changelogSectionFixed: "Correções", + changelogSectionChanged: "Alterações", + changelogSectionRemoved: "Removido", shortcutUppercaseSelection: "Converter seleção em maiúsculas", shortcutLowercaseSelection: "Converter seleção em minúsculas", shortcutExPasteSqlInCondition: "ExPaste: colar como condição IN", diff --git a/apps/desktop/src/i18n/locales/zh-CN.ts b/apps/desktop/src/i18n/locales/zh-CN.ts index 1c243db319..141c41cab7 100644 --- a/apps/desktop/src/i18n/locales/zh-CN.ts +++ b/apps/desktop/src/i18n/locales/zh-CN.ts @@ -3627,6 +3627,23 @@ export default withEnglishFallback({ project: "项目", openSource: "开源仓库", officialDocs: "官方文档", + changelogTitle: "更新日志", + changelogDescription: "查看各版本的更新内容。", + changelogLoading: "正在加载更新日志...", + changelogLoadFailed: "加载更新日志失败。", + changelogEmpty: "暂无更新日志。", + changelogPublishedOn: "发布于 {date}", + changelogNew: "NEW", + changelogLoadMore: "加载更多", + changelogOpenWebsite: "在官网查看", + changelogOpenGitHub: "GitHub Release", + changelogNoDetails: "详见 GitHub Release。", + changelogRetry: "重试", + changelogSectionAdded: "新功能", + changelogSectionImproved: "改进", + changelogSectionFixed: "问题修复", + changelogSectionChanged: "变更", + changelogSectionRemoved: "移除", }, driverStore: { progressJreExtract: "解压 JRE...", diff --git a/apps/desktop/src/i18n/locales/zh-TW.ts b/apps/desktop/src/i18n/locales/zh-TW.ts index ead172c169..feef91c539 100644 --- a/apps/desktop/src/i18n/locales/zh-TW.ts +++ b/apps/desktop/src/i18n/locales/zh-TW.ts @@ -3214,6 +3214,23 @@ export default withEnglishFallback({ project: "專案", openSource: "開源倉庫", officialDocs: "官方文件", + changelogTitle: "更新日誌", + changelogDescription: "查看各版本的更新內容。", + changelogLoading: "正在載入更新日誌...", + changelogLoadFailed: "載入更新日誌失敗。", + changelogEmpty: "暫無更新日誌。", + changelogPublishedOn: "發佈於 {date}", + changelogNew: "NEW", + changelogLoadMore: "載入更多", + changelogOpenWebsite: "在官網查看", + changelogOpenGitHub: "GitHub Release", + changelogNoDetails: "詳見 GitHub Release。", + changelogRetry: "重試", + changelogSectionAdded: "新功能", + changelogSectionImproved: "改進", + changelogSectionFixed: "問題修復", + changelogSectionChanged: "變更", + changelogSectionRemoved: "移除", mcpTab: "MCP", mcpTitle: "MCP Server", mcpDescription: "檢查 Claude Code、Cursor 等程式助理使用的 DBX MCP Server 安裝與版本狀態。", diff --git a/apps/desktop/src/lib/app/changelog.ts b/apps/desktop/src/lib/app/changelog.ts new file mode 100644 index 0000000000..17591d1a8a --- /dev/null +++ b/apps/desktop/src/lib/app/changelog.ts @@ -0,0 +1,68 @@ +export type ChangelogItem = { + title: string; + desc: string; +}; + +export type ChangelogSection = { + type: string; + title: string; + items: ChangelogItem[]; +}; + +export type ChangelogRelease = { + tag: string; + name: string; + date: string; + sections: ChangelogSection[]; +}; + +export type ChangelogData = { + updatedAt: string; + releases: ChangelogRelease[]; +}; + +export type ChangelogLang = "en" | "cn"; + +const cache = new Map>(); + +export function changelogLangFromLocale(locale: string): ChangelogLang { + return locale === "zh-CN" || locale === "zh-TW" ? "cn" : "en"; +} + +export function changelogWebsiteUrl(lang: ChangelogLang): string { + return `https://dbxio.com/${lang}/changelog`; +} + +export function changelogReleaseUrl(tag: string): string { + return `https://github.com/t8y2/dbx/releases/tag/${encodeURIComponent(tag)}`; +} + +export async function fetchChangelog(lang: ChangelogLang, options: { force?: boolean } = {}): Promise { + if (options.force) { + cache.delete(lang); + } + + let pending = cache.get(lang); + if (!pending) { + pending = loadChangelog(lang).catch((error) => { + cache.delete(lang); + throw error; + }); + cache.set(lang, pending); + } + + return pending; +} + +async function loadChangelog(lang: ChangelogLang): Promise { + const { fetchChangelog: fetchChangelogViaBackend } = await import("@/lib/backend/api"); + const data = await fetchChangelogViaBackend(lang); + if (!data || !Array.isArray(data.releases)) { + throw new Error("Invalid changelog payload"); + } + + return { + updatedAt: typeof data.updatedAt === "string" ? data.updatedAt : "", + releases: data.releases.filter((release) => release && typeof release.tag === "string"), + }; +} diff --git a/apps/desktop/src/lib/backend/api.ts b/apps/desktop/src/lib/backend/api.ts index a33b09a0ba..e3f3dfeb81 100644 --- a/apps/desktop/src/lib/backend/api.ts +++ b/apps/desktop/src/lib/backend/api.ts @@ -466,6 +466,7 @@ export const deleteHistoryEntry = forward("deleteHistoryEntry"); export const checkMcpServerStatus = forward("checkMcpServerStatus"); export const installMcpServer = forward("installMcpServer"); export const checkForUpdates = forward("checkForUpdates"); +export const fetchChangelog = forward("fetchChangelog"); export const getSystemProxyUrl = forward("getSystemProxyUrl"); export const downloadAndInstallUpdate = forward("downloadAndInstallUpdate"); export const getAppVersion = forward("getAppVersion"); diff --git a/apps/desktop/src/lib/backend/http.ts b/apps/desktop/src/lib/backend/http.ts index 6f9aeed05a..f485d9014d 100644 --- a/apps/desktop/src/lib/backend/http.ts +++ b/apps/desktop/src/lib/backend/http.ts @@ -2165,6 +2165,11 @@ export async function checkForUpdates(locale?: string): Promise { return get(`/api/update/check${query}`); } +export async function fetchChangelog(lang?: string): Promise { + const query = lang ? `?lang=${encodeURIComponent(lang)}` : ""; + return get(`/api/changelog${query}`); +} + export async function checkMcpServerStatus(): Promise { return { installed: false, diff --git a/apps/desktop/src/lib/backend/tauri.ts b/apps/desktop/src/lib/backend/tauri.ts index 5243ee56c0..93cb99a3b3 100644 --- a/apps/desktop/src/lib/backend/tauri.ts +++ b/apps/desktop/src/lib/backend/tauri.ts @@ -1391,6 +1391,10 @@ export async function checkForUpdates(locale?: string): Promise { return invoke("check_for_updates", { locale }); } +export async function fetchChangelog(lang?: string): Promise { + return invoke("fetch_changelog", { lang }); +} + export async function getSystemProxyUrl(): Promise { return invoke("get_system_proxy_url"); } diff --git a/crates/dbx-core/src/changelog.rs b/crates/dbx-core/src/changelog.rs new file mode 100644 index 0000000000..6fc108367a --- /dev/null +++ b/crates/dbx-core/src/changelog.rs @@ -0,0 +1,99 @@ +use serde::{Deserialize, Serialize}; + +const CHANGELOG_R2_PREFIX: &str = "changelog/releases-"; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct ChangelogItem { + pub title: String, + #[serde(default)] + pub desc: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ChangelogSection { + #[serde(default)] + pub r#type: String, + #[serde(default)] + pub title: String, + #[serde(default)] + pub items: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ChangelogRelease { + pub tag: String, + #[serde(default)] + pub name: String, + #[serde(default)] + pub date: String, + #[serde(default)] + pub sections: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct ChangelogData { + #[serde(default)] + pub updated_at: String, + #[serde(default)] + pub releases: Vec, +} + +/// Normalize UI locale / short lang into `cn` or `en` (matches R2 file naming). +pub fn normalize_changelog_lang(lang: &str) -> &'static str { + let trimmed = lang.trim(); + if trimmed.eq_ignore_ascii_case("cn") + || trimmed.eq_ignore_ascii_case("zh") + || trimmed.eq_ignore_ascii_case("zh-CN") + || trimmed.eq_ignore_ascii_case("zh-TW") + { + "cn" + } else { + "en" + } +} + +pub async fn fetch_changelog(lang: &str) -> Result { + let lang = normalize_changelog_lang(lang); + let client = build_changelog_http_client()?; + let url = format!("{}{CHANGELOG_R2_PREFIX}{lang}.json", crate::R2_CDN_BASE); + + let resp = client + .get(&url) + .header(reqwest::header::USER_AGENT, "dbx-changelog") + .send() + .await + .and_then(|r| r.error_for_status()) + .map_err(|e| format!("Failed to fetch changelog: {e}"))?; + + let mut data: ChangelogData = resp.json().await.map_err(|e| format!("Failed to parse changelog: {e}"))?; + data.releases.retain(|release| !release.tag.trim().is_empty()); + Ok(data) +} + +fn build_changelog_http_client() -> Result { + let mut builder = + reqwest::Client::builder().timeout(std::time::Duration::from_secs(15)).user_agent("dbx-changelog"); + + if let Some(proxy_url) = crate::update::system_proxy_url() { + let proxy = reqwest::Proxy::all(&proxy_url).map_err(|e| format!("Invalid system proxy URL: {e}"))?; + builder = builder.proxy(proxy); + } + + builder.build().map_err(|e| format!("Failed to create HTTP client: {e}")) +} + +#[cfg(test)] +mod tests { + use super::normalize_changelog_lang; + + #[test] + fn normalizes_changelog_lang() { + assert_eq!(normalize_changelog_lang("cn"), "cn"); + assert_eq!(normalize_changelog_lang("zh-CN"), "cn"); + assert_eq!(normalize_changelog_lang("zh-TW"), "cn"); + assert_eq!(normalize_changelog_lang("en"), "en"); + assert_eq!(normalize_changelog_lang("ja"), "en"); + } +} diff --git a/crates/dbx-core/src/lib.rs b/crates/dbx-core/src/lib.rs index 407143b89a..f153fe0575 100644 --- a/crates/dbx-core/src/lib.rs +++ b/crates/dbx-core/src/lib.rs @@ -11,6 +11,7 @@ pub mod agent_tools; pub mod ai; pub mod ai_cli_agent; pub mod ai_codex_cli; +pub mod changelog; pub mod cloud_sync; pub mod connection; pub mod connection_secrets; diff --git a/crates/dbx-web/src/main.rs b/crates/dbx-web/src/main.rs index 28a615d099..502f004951 100644 --- a/crates/dbx-web/src/main.rs +++ b/crates/dbx-web/src/main.rs @@ -556,6 +556,7 @@ async fn main() { // Update .route("/version", get(routes::update::get_version)) .route("/update/check", get(routes::update::check_for_updates)) + .route("/changelog", get(routes::update::fetch_changelog)) // Layout .route("/layout/sidebar", post(routes::layout::save_sidebar_layout).get(routes::layout::load_sidebar_layout)) // App settings diff --git a/crates/dbx-web/src/routes/update.rs b/crates/dbx-web/src/routes/update.rs index e3ca0bd4ee..1a3f261201 100644 --- a/crates/dbx-web/src/routes/update.rs +++ b/crates/dbx-web/src/routes/update.rs @@ -1,5 +1,5 @@ use axum::{extract::Query, Json}; -use dbx_core::update; +use dbx_core::{changelog, update}; use crate::error::AppError; @@ -19,3 +19,17 @@ pub async fn check_for_updates(Query(params): Query) -> Resul let info = update::build_update_info(release, env!("CARGO_PKG_VERSION")); Ok(Json(serde_json::to_value(info).map_err(|e| AppError(e.to_string()))?)) } + +#[derive(serde::Deserialize)] +pub struct ChangelogParams { + #[serde(default)] + pub lang: Option, +} + +pub async fn fetch_changelog( + Query(params): Query, +) -> Result, AppError> { + let lang = params.lang.unwrap_or_else(|| "en".to_string()); + let data = changelog::fetch_changelog(&lang).await.map_err(AppError)?; + Ok(Json(data)) +} diff --git a/src-tauri/src/commands/update.rs b/src-tauri/src/commands/update.rs index 222e8b444b..223394daa0 100644 --- a/src-tauri/src/commands/update.rs +++ b/src-tauri/src/commands/update.rs @@ -118,6 +118,12 @@ pub async fn check_for_updates(locale: Option) -> Result) -> Result { + let lang = lang.unwrap_or_else(|| "en".to_string()); + dbx_core::changelog::fetch_changelog(&lang).await +} + #[tauri::command] pub async fn get_system_proxy_url() -> Option { tauri::async_runtime::spawn_blocking(dbx_core::update::system_proxy_url).await.ok().flatten() diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index 6af956c05e..542af0a21d 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -1235,6 +1235,7 @@ pub fn run() { commands::mcp::check_mcp_server_status, commands::mcp::install_mcp_server, commands::update::check_for_updates, + commands::update::fetch_changelog, commands::update::get_system_proxy_url, commands::update::download_and_install_update, commands::transfer::start_transfer, From dd95aafc9e364b9a0c4674ac402f1a4d1da24afe Mon Sep 17 00:00:00 2001 From: Freedom <459102951@qq.com> Date: Thu, 16 Jul 2026 08:41:45 +0800 Subject: [PATCH 6/6] fix(mq): retry empty Kafka peek polls until caught up or deadline Empty first polls no longer abort early; stop only when partitions reach end offsets or the 5s deadline expires. Reject non-safe-integer peek partition/offset inputs instead of truncating decimals. Co-authored-by: Cursor --- .../java/com/dbx/agent/kafka/KafkaAgent.java | 88 ++++++++++++++++--- .../com/dbx/agent/kafka/KafkaAgentTest.java | 74 ++++++++++++++++ .../src/components/mq/SendMessagePanel.vue | 17 ++-- .../lib/__tests__/mq/mqPeekFilters.spec.ts | 26 ++++++ apps/desktop/src/lib/mq/mqPeekFilters.ts | 12 +++ 5 files changed, 196 insertions(+), 21 deletions(-) create mode 100644 apps/desktop/src/lib/__tests__/mq/mqPeekFilters.spec.ts create mode 100644 apps/desktop/src/lib/mq/mqPeekFilters.ts diff --git a/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java b/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java index 6e03d95020..2330d0ab2d 100644 --- a/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java +++ b/agents/drivers/kafka/src/main/java/com/dbx/agent/kafka/KafkaAgent.java @@ -851,20 +851,19 @@ private static Object peekMessages(JsonObject params) throws Exception { consumer.seek(entry.getKey(), entry.getValue()); } - List> messages = new ArrayList<>(); - long deadlineNs = System.nanoTime() + Duration.ofSeconds(5).toNanos(); - while (messages.size() < count && System.nanoTime() < deadlineNs) { - ConsumerRecords records = consumer.poll(Duration.ofMillis(500)); - if (records.isEmpty()) { - break; - } - for (ConsumerRecord record : records) { - messages.add(peekedMessageFromRecord(record)); - if (messages.size() >= count) { - break; + List> messages = collectPeekedMessages( + timeout -> consumer.poll(timeout), + () -> { + Map positions = new LinkedHashMap<>(); + for (TopicPartition tp : readablePartitions) { + positions.put(tp, consumer.position(tp)); } - } - } + return allPeekPartitionsCaughtUp(readablePartitions, positions, endOffsets); + }, + count, + System.nanoTime() + Duration.ofSeconds(5).toNanos(), + Duration.ofMillis(500) + ); sortPeekedMessages(messages); if (messages.size() > count) { messages = new ArrayList<>(messages.subList(0, count)); @@ -873,6 +872,69 @@ private static Object peekMessages(JsonObject params) throws Exception { } } + /** + * Poll until {@code count} messages are collected, every assigned partition has reached its + * end offset, or {@code deadlineNs} expires. Empty polls retry until caught-up or deadline — + * they must not abort early (broker / network / first-fetch latency can exceed one poll). + */ + static List> collectPeekedMessages( + PeekRecordPoller poller, + PeekCaughtUpChecker caughtUpChecker, + int count, + long deadlineNs, + Duration pollTimeout + ) { + List> messages = new ArrayList<>(); + while (messages.size() < count && System.nanoTime() < deadlineNs) { + long remainingNs = deadlineNs - System.nanoTime(); + if (remainingNs <= 0) { + break; + } + Duration timeout = pollTimeout.toNanos() > remainingNs + ? Duration.ofNanos(remainingNs) + : pollTimeout; + ConsumerRecords records = poller.poll(timeout); + if (records.isEmpty()) { + if (caughtUpChecker.allPartitionsCaughtUp()) { + break; + } + continue; + } + for (ConsumerRecord record : records) { + messages.add(peekedMessageFromRecord(record)); + if (messages.size() >= count) { + break; + } + } + } + return messages; + } + + static boolean allPeekPartitionsCaughtUp( + List partitions, + Map positions, + Map endOffsets + ) { + for (TopicPartition tp : partitions) { + long endOffset = endOffsets.getOrDefault(tp, 0L); + long position = positions.getOrDefault(tp, 0L); + if (position < endOffset) { + return false; + } + } + return true; + } + + @FunctionalInterface + interface PeekRecordPoller { + ConsumerRecords poll(Duration timeout); + } + + @FunctionalInterface + interface PeekCaughtUpChecker { + boolean allPartitionsCaughtUp(); + } + /** When partition is null, peek across every partition of the topic. */ static List resolvePeekPartitions( KafkaConsumer consumer, diff --git a/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java b/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java index 7003d1f574..6ebe51311a 100644 --- a/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java +++ b/agents/drivers/kafka/src/test/java/com/dbx/agent/kafka/KafkaAgentTest.java @@ -1,12 +1,21 @@ package com.dbx.agent.kafka; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; import com.google.gson.JsonParser; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Properties; +import java.util.concurrent.atomic.AtomicInteger; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.common.TopicPartition; import org.junit.jupiter.api.Test; class KafkaAgentTest { @@ -63,6 +72,71 @@ void sortPeekedMessagesOrdersByTimestampThenPartitionThenOffset() { assertEquals(20L, messages.get(3).get("timestamp")); } + @Test + void allPeekPartitionsCaughtUpRequiresEveryPartitionAtEndOffset() { + TopicPartition p0 = new TopicPartition("events", 0); + TopicPartition p1 = new TopicPartition("events", 1); + Map endOffsets = Map.of(p0, 10L, p1, 5L); + + assertFalse(KafkaAgent.allPeekPartitionsCaughtUp( + List.of(p0, p1), + Map.of(p0, 10L, p1, 4L), + endOffsets + )); + assertTrue(KafkaAgent.allPeekPartitionsCaughtUp( + List.of(p0, p1), + Map.of(p0, 10L, p1, 5L), + endOffsets + )); + } + + @Test + void collectPeekedMessagesRetriesAfterEmptyFirstPoll() { + TopicPartition tp = new TopicPartition("events", 0); + ConsumerRecord record = new ConsumerRecord<>( + "events", + 0, + 7L, + "k", + "hello".getBytes(StandardCharsets.UTF_8) + ); + Map>> batch = new HashMap<>(); + batch.put(tp, List.of(record)); + ConsumerRecords withData = new ConsumerRecords<>(batch); + + AtomicInteger polls = new AtomicInteger(); + List> messages = KafkaAgent.collectPeekedMessages( + timeout -> polls.getAndIncrement() == 0 ? ConsumerRecords.empty() : withData, + () -> false, + 1, + System.nanoTime() + Duration.ofSeconds(5).toNanos(), + Duration.ofMillis(1) + ); + + assertEquals(2, polls.get()); + assertEquals(1, messages.size()); + assertEquals(7L, messages.get(0).get("offset")); + assertEquals("hello", messages.get(0).get("payloadText")); + } + + @Test + void collectPeekedMessagesStopsOnEmptyPollWhenCaughtUp() { + AtomicInteger polls = new AtomicInteger(); + List> messages = KafkaAgent.collectPeekedMessages( + timeout -> { + polls.incrementAndGet(); + return ConsumerRecords.empty(); + }, + () -> true, + 10, + System.nanoTime() + Duration.ofSeconds(5).toNanos(), + Duration.ofMillis(1) + ); + + assertEquals(1, polls.get()); + assertTrue(messages.isEmpty()); + } + @Test void appliesKerberosKafkaProperties() { Properties props = new Properties(); diff --git a/apps/desktop/src/components/mq/SendMessagePanel.vue b/apps/desktop/src/components/mq/SendMessagePanel.vue index 0835ea9f59..d69c8d44d4 100644 --- a/apps/desktop/src/components/mq/SendMessagePanel.vue +++ b/apps/desktop/src/components/mq/SendMessagePanel.vue @@ -4,6 +4,7 @@ import { useI18n } from "vue-i18n"; import type { PeekedMessage, TopicInfo, TopicRef, SendMessageRequest, SendMessageResponse } from "@/types/mq"; import { mqSendMessage, mqListTopics, mqPeekMessages } from "@/lib/backend/api"; import { formatError } from "@/lib/backend/errorUtils"; +import { parseNonNegativeSafeInteger } from "@/lib/mq/mqPeekFilters"; interface Props { connectionId: string; @@ -176,20 +177,20 @@ async function loadMessages() { const partitionText = peekPartition.value.trim(); const offsetText = peekOffset.value.trim(); if (partitionText !== "") { - const partition = Number(partitionText); - if (!Number.isFinite(partition) || partition < 0) { + const partition = parseNonNegativeSafeInteger(partitionText); + if (partition == null) { throw new Error(t("mqMessages.partitionMustBeNonNegativeInt")); } - options.partition = Math.floor(partition); - peekPartition.value = String(options.partition); + options.partition = partition; + peekPartition.value = String(partition); } if (offsetText !== "") { - const offset = Number(offsetText); - if (!Number.isFinite(offset) || offset < 0) { + const offset = parseNonNegativeSafeInteger(offsetText); + if (offset == null) { throw new Error(t("mqMessages.offsetMustBeNonNegativeInt")); } - options.offset = Math.floor(offset); - peekOffset.value = String(options.offset); + options.offset = offset; + peekOffset.value = String(offset); } peekMessages.value = await mqPeekMessages(props.connectionId, topic, "__dbx_kafka_viewer__", count, options); } catch (e: unknown) { diff --git a/apps/desktop/src/lib/__tests__/mq/mqPeekFilters.spec.ts b/apps/desktop/src/lib/__tests__/mq/mqPeekFilters.spec.ts new file mode 100644 index 0000000000..e74b3ce77b --- /dev/null +++ b/apps/desktop/src/lib/__tests__/mq/mqPeekFilters.spec.ts @@ -0,0 +1,26 @@ +import { describe, expect, it } from "vitest"; +import { parseNonNegativeSafeInteger } from "@/lib/mq/mqPeekFilters"; + +describe("parseNonNegativeSafeInteger", () => { + it("accepts non-negative safe integers", () => { + expect(parseNonNegativeSafeInteger("0")).toBe(0); + expect(parseNonNegativeSafeInteger("20")).toBe(20); + expect(parseNonNegativeSafeInteger(String(Number.MAX_SAFE_INTEGER))).toBe(Number.MAX_SAFE_INTEGER); + }); + + it("rejects decimals and non-integers", () => { + expect(parseNonNegativeSafeInteger("1.9")).toBeNull(); + expect(parseNonNegativeSafeInteger("0.1")).toBeNull(); + expect(parseNonNegativeSafeInteger("abc")).toBeNull(); + }); + + it("rejects negatives and values outside the safe integer range", () => { + expect(parseNonNegativeSafeInteger("-1")).toBeNull(); + expect(parseNonNegativeSafeInteger(String(Number.MAX_SAFE_INTEGER + 1))).toBeNull(); + }); + + it("rejects empty input", () => { + expect(parseNonNegativeSafeInteger("")).toBeNull(); + expect(parseNonNegativeSafeInteger(" ")).toBeNull(); + }); +}); diff --git a/apps/desktop/src/lib/mq/mqPeekFilters.ts b/apps/desktop/src/lib/mq/mqPeekFilters.ts new file mode 100644 index 0000000000..dc130debda --- /dev/null +++ b/apps/desktop/src/lib/mq/mqPeekFilters.ts @@ -0,0 +1,12 @@ +/** Parse a non-negative safe integer from user input; rejects decimals and unsafe magnitudes. */ +export function parseNonNegativeSafeInteger(text: string): number | null { + const trimmed = text.trim(); + if (trimmed === "") { + return null; + } + const value = Number(trimmed); + if (!Number.isSafeInteger(value) || value < 0) { + return null; + } + return value; +}