feat(redis): add stream consumer group monitoring

This commit is contained in:
onenewcode
2026-07-28 23:16:39 +08:00
committed by GitHub
parent 99cfa1b516
commit 31f0b7d6ef
23 changed files with 2092 additions and 44 deletions
@@ -1755,7 +1755,7 @@ defineExpose({ focusSearch, refreshData, refreshQueryEditorCompletionCache, hand
<!-- Redis mode: key browser -->
<template v-else-if="activeTab.mode === 'redis'">
<div class="flex-1 min-h-0">
<RedisKeyBrowser ref="redisKeyBrowserRef" :key="activeTab.id" :connection-id="activeTab.connectionId" :db="Number(activeTab.database)" :block-dangerous-redis-commands="props.blockDangerousRedisCommands" />
<RedisKeyBrowser ref="redisKeyBrowserRef" :key="`${activeTab.id}:${activeTab.connectionId}:${activeTab.database}`" :connection-id="activeTab.connectionId" :db="Number(activeTab.database)" :block-dangerous-redis-commands="props.blockDangerousRedisCommands" />
</div>
</template>
@@ -5,6 +5,7 @@ import { createApp, defineComponent, h, KeepAlive, nextTick, ref, type Component
import { createI18n } from "vue-i18n";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { calendarDateTimeToUnixSeconds } from "@/components/ui/date-time-picker/dateTimePicker";
import type { RedisKeyInfo } from "@/lib/backend/api";
const mocks = vi.hoisted(() => ({
redisScanKeysBatch: vi.fn(),
@@ -367,6 +368,32 @@ function mountBrowser() {
mountedApps.push({ unmount: () => app.unmount(), host });
}
function mountScopedBrowser() {
const connectionId = ref("connection");
const db = ref(0);
const host = document.createElement("div");
document.body.append(host);
const app = createApp(
defineComponent({
setup() {
return () => h(RedisKeyBrowser, { connectionId: connectionId.value, db: db.value, blockDangerousRedisCommands: false });
},
}),
);
app.use(createI18n({ legacy: false, locale: "en", messages: { en: {} }, missingWarn: false, fallbackWarn: false }));
app.mount(host);
mountedApps.push({ unmount: () => app.unmount(), host });
return {
host,
async setScope(nextConnectionId: string, nextDb: number) {
connectionId.value = nextConnectionId;
db.value = nextDb;
await settle();
},
};
}
function mountKeptAliveBrowser() {
const active = ref(true);
const host = document.createElement("div");
@@ -522,6 +549,34 @@ afterEach(() => {
resetLocalTimeZone();
});
describe("RedisKeyBrowser scope changes", () => {
it("reloads the new database and discards a late scan from the previous one", async () => {
const previousDatabase = deferred<{ cursor: number; keys: RedisKeyInfo[]; total_keys: number }>();
const currentKey = { key_display: "db1-key", key_raw: "ZGIxLWtleQ==", key_type: "string", ttl: -1 };
mocks.redisScanKeysBatch.mockImplementation((_connectionId: string, db: number) => {
if (db === 0) return previousDatabase.promise;
return Promise.resolve({ cursor: 0, keys: [currentKey], total_keys: 1 });
});
const browser = mountScopedBrowser();
await settle();
expect(mocks.redisScanKeysBatch).toHaveBeenCalledWith("connection", 0, 0, "*", 100, 8, false);
await browser.setScope("connection", 1);
expect(mocks.redisScanKeysBatch).toHaveBeenCalledWith("connection", 1, 0, "*", 100, 8, false);
previousDatabase.resolve({
cursor: 0,
keys: [{ key_display: "db0-key", key_raw: "ZGIwLWtleQ==", key_type: "string", ttl: -1 }],
total_keys: 1,
});
await settle();
expect(browser.host.textContent).toContain("db1-key");
expect(browser.host.textContent).not.toContain("db0-key");
});
});
describe("RedisKeyBrowser expiry creation", () => {
it.each(["string", "hash", "list", "set", "zset", "stream", "json"] as const)("writes %s before applying one relative TTL", async (type) => {
mountBrowser();
@@ -1432,9 +1432,19 @@ onDeactivated(pauseRedisBrowserBackgroundWork);
onUnmounted(pauseRedisBrowserBackgroundWork);
watch(
() => props.db,
(db) => {
() => [props.connectionId, props.db] as const,
async ([connectionId, db]) => {
// ContentArea remounts this browser for scope changes; keep embedded uses
// in sync as well so an old scan cannot populate the new scope.
commandDb.value = db;
resetLoadedKeys();
try {
await connectionStore.ensureConnected(connectionId);
} catch (error) {
console.warn("[DBX] ensureConnected failed for", connectionId, error);
}
if (connectionId !== props.connectionId || db !== props.db) return;
void loadKeys();
},
);
@@ -0,0 +1,426 @@
// @vitest-environment happy-dom
import { createApp, defineComponent, h, nextTick } from "vue";
import { createI18n } from "vue-i18n";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
const mocks = vi.hoisted(() => ({
redisGetValue: vi.fn(),
redisGetStreamGroups: vi.fn(),
redisGetStreamConsumers: vi.fn(),
redisGetStreamPending: vi.fn(),
redisSetTtl: vi.fn(),
redisSetExpireAt: vi.fn(),
toast: vi.fn(),
}));
vi.mock("@/lib/backend/api", () => ({
redisGetValue: mocks.redisGetValue,
redisGetStreamGroups: mocks.redisGetStreamGroups,
redisGetStreamConsumers: mocks.redisGetStreamConsumers,
redisGetStreamPending: mocks.redisGetStreamPending,
redisSetTtl: mocks.redisSetTtl,
redisSetExpireAt: mocks.redisSetExpireAt,
}));
vi.mock("@/composables/useEditorFontFamilyStyle", () => ({
useEditorFontFamilyStyle: () => ({}),
}));
vi.mock("@/composables/useToast", () => ({
useToast: () => ({ toast: mocks.toast }),
}));
vi.mock("@/lib/common/shikiJsonHighlighter", () => ({
createShikiJsonHighlighter: vi.fn().mockResolvedValue(() => ""),
}));
vi.mock("vue-virtual-scroller", async () => {
const { defineComponent, h } = await import("vue");
const DynamicScroller = defineComponent({
props: { items: { type: Array, default: () => [] } },
setup(props, { slots }) {
return () =>
h(
"div",
(props.items as unknown[]).map((item, index) => slots.default?.({ item, active: true, index })),
);
},
});
const DynamicScrollerItem = defineComponent({
setup(_, { slots }) {
return () => h("div", slots.default?.());
},
});
const RecycleScroller = defineComponent({
setup(_, { slots }) {
return () => h("div", slots.default?.());
},
});
return { DynamicScroller, DynamicScrollerItem, RecycleScroller };
});
import RedisValueViewer from "./RedisValueViewer.vue";
const mountedApps: Array<{ unmount: () => void; host: HTMLElement }> = [];
afterEach(() => {
for (const { unmount, host } of mountedApps.splice(0)) {
unmount();
host.remove();
}
});
beforeEach(() => {
vi.clearAllMocks();
});
async function settle() {
await nextTick();
await Promise.resolve();
await nextTick();
await Promise.resolve();
await nextTick();
}
function streamValue() {
return {
key_display: "orders",
key_raw: "b3JkZXJz",
ttl: -1,
redis_type: "stream",
data: { kind: "stream" as const, entries: [] },
};
}
function blob(raw_base64: string) {
return { raw_base64, encoding: "utf8" as const };
}
function deferred<T>() {
let resolve!: (value: T) => void;
const promise = new Promise<T>((next) => {
resolve = next;
});
return { promise, resolve };
}
function group(pending: number | string = 1) {
return {
name: blob("cGF5bWVudHM="),
consumers: 1,
pending,
last_delivered_id: "1714470000000-0",
entries_read: 12,
lag: 3,
};
}
function mountViewer() {
const host = document.createElement("div");
document.body.append(host);
const app = createApp(
defineComponent({
setup() {
return () =>
h(RedisValueViewer, {
connectionId: "redis-1",
db: 2,
keyDisplay: "orders",
keyRaw: "b3JkZXJz",
onDeleted: vi.fn(),
});
},
}),
);
app.use(
createI18n({
legacy: false,
locale: "en",
messages: {
en: {
common: { loading: "Loading", retry: "Retry" },
redis: { noConsumerGroups: "No consumer groups" },
},
},
missingWarn: false,
fallbackWarn: false,
}),
);
app.mount(host);
mountedApps.push({ unmount: () => app.unmount(), host });
return host;
}
function openGroups(host: HTMLElement) {
host.querySelector<HTMLButtonElement>("[data-redis-stream-groups-tab]")!.dispatchEvent(new MouseEvent("mousedown", { bubbles: true, button: 0 }));
}
describe("RedisValueViewer stream monitoring", () => {
it("renders large counter values transported as decimal strings", async () => {
const unsafeMetric = "9007199254740992";
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValue([group(unsafeMetric)]);
const host = mountViewer();
await settle();
openGroups(host);
await settle();
expect(host.querySelector<HTMLElement>("[data-redis-stream-group-row]")?.textContent).toContain(BigInt(unsafeMetric).toLocaleString());
});
it("loads groups lazily, then loads a selected consumer's pending entries", async () => {
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValue([group()]);
mocks.redisGetStreamConsumers.mockResolvedValue([{ name: blob("d29ya2VyLWE="), pending: 1, idle_ms: 1_200, inactive_ms: 800 }]);
mocks.redisGetStreamPending.mockResolvedValueOnce({
entries: [{ id: "1714470000000-0", consumer: blob("d29ya2VyLWE="), idle_ms: 2_400, deliveries: 2 }],
next_cursor: "1714470000000-0",
});
mocks.redisGetStreamPending.mockResolvedValueOnce({
entries: [{ id: "1714470000001-0", consumer: blob("d29ya2VyLWE="), idle_ms: 1_000, deliveries: 1 }],
});
const host = mountViewer();
await settle();
expect(mocks.redisGetValue).toHaveBeenCalledWith("redis-1", 2, "b3JkZXJz");
expect(mocks.redisGetStreamGroups).not.toHaveBeenCalled();
openGroups(host);
await settle();
expect(mocks.redisGetStreamGroups).toHaveBeenCalledWith("redis-1", 2, "b3JkZXJz");
const groupRow = host.querySelector<HTMLElement>("[data-redis-stream-group-row]")!;
expect(groupRow.textContent).toContain("payments");
groupRow.click();
await settle();
expect(mocks.redisGetStreamConsumers).toHaveBeenCalledWith("redis-1", 2, "b3JkZXJz", "cGF5bWVudHM=");
expect(host.querySelector("[data-redis-stream-group-detail]")?.textContent).toContain("worker-a");
expect(mocks.redisGetStreamPending).not.toHaveBeenCalled();
host.querySelector<HTMLButtonElement>("[data-redis-stream-consumer-row]")!.click();
await settle();
expect(mocks.redisGetStreamPending).toHaveBeenNthCalledWith(1, "redis-1", 2, "b3JkZXJz", "cGF5bWVudHM=", undefined, "d29ya2VyLWE=");
expect(host.querySelector("[data-redis-stream-consumer-crumb]")?.textContent).toContain("worker-a");
host.querySelector<HTMLButtonElement>("[data-redis-stream-pending-more]")!.click();
await settle();
expect(mocks.redisGetStreamPending).toHaveBeenNthCalledWith(2, "redis-1", 2, "b3JkZXJz", "cGF5bWVudHM=", "1714470000000-0", "d29ya2VyLWE=");
expect(host.querySelector("[data-redis-stream-group-detail]")?.textContent).toContain("1714470000001-0");
});
it("renders group metrics and consumer navigation without an aggregate pending list", async () => {
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValue([group()]);
mocks.redisGetStreamConsumers.mockResolvedValue([{ name: blob("d29ya2VyLWE="), pending: 1, idle_ms: 1_200, inactive_ms: 800 }]);
const host = mountViewer();
await settle();
openGroups(host);
await settle();
host.querySelector<HTMLElement>("[data-redis-stream-group-row]")!.click();
await settle();
const summary = host.querySelector<HTMLElement>("[data-redis-stream-group-summary]")!;
const consumerPanel = host.querySelector<HTMLElement>("[data-redis-stream-consumers]")!;
const consumerTable = host.querySelector<HTMLTableElement>("[data-redis-stream-consumers-table]")!;
const groupsTab = host.querySelector<HTMLElement>("[data-redis-stream-groups-tab]")!;
const groupCrumb = host.querySelector<HTMLElement>("[data-redis-stream-group-crumb]")!;
expect(summary.classList.contains("rounded-lg")).toBe(true);
expect(summary.classList.contains("border")).toBe(true);
expect(consumerPanel.classList.contains("rounded-lg")).toBe(true);
expect(consumerPanel.classList.contains("bg-card")).toBe(true);
expect(consumerTable.classList.contains("table-fixed")).toBe(true);
expect([...consumerTable.querySelectorAll("col")].map((column) => column.className)).toEqual(["w-1/4", "w-1/4", "w-1/4", "w-1/4"]);
const consumerHeaders = consumerTable.querySelectorAll<HTMLTableCellElement>("[data-redis-stream-consumer-header] th");
const consumerCells = consumerTable.querySelectorAll<HTMLTableCellElement>("tbody td");
expect(host.querySelector("[data-redis-stream-pending]")).toBeNull();
expect(groupsTab.className).toContain("group-data-[variant=line]/tabs-list:data-active:after:opacity-0");
expect(groupCrumb.classList.contains("border-foreground")).toBe(true);
expect(groupCrumb.classList.contains("border-transparent")).toBe(false);
expect(consumerHeaders[1].classList.contains("text-right")).toBe(true);
expect(consumerHeaders[2].classList.contains("text-right")).toBe(true);
expect(consumerHeaders[3].classList.contains("text-right")).toBe(true);
expect(consumerCells[0].classList.contains("text-left")).toBe(true);
expect(consumerCells[1].classList.contains("text-right")).toBe(true);
expect(consumerCells[2].classList.contains("text-right")).toBe(true);
expect(consumerCells[3].classList.contains("text-right")).toBe(true);
expect(consumerTable.querySelector("[data-redis-stream-consumer-row] svg")).toBeNull();
});
it("drills into a consumer and reloads its pending entries server-side", async () => {
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValue([group()]);
mocks.redisGetStreamConsumers.mockResolvedValue([
{ name: blob("d29ya2VyLWE="), pending: 1, idle_ms: 1_200, inactive_ms: 800 },
{ name: blob("d29ya2VyLWI="), pending: 2, idle_ms: 900, inactive_ms: 600 },
]);
mocks.redisGetStreamPending.mockResolvedValueOnce({ entries: [{ id: "1714470000001-0", consumer: blob("d29ya2VyLWI="), idle_ms: 1_000, deliveries: 1 }] });
const host = mountViewer();
await settle();
openGroups(host);
await settle();
host.querySelector<HTMLElement>("[data-redis-stream-group-row]")!.click();
await settle();
host.querySelectorAll<HTMLButtonElement>("[data-redis-stream-consumer-row]")[1].click();
await settle();
expect(mocks.redisGetStreamPending).toHaveBeenNthCalledWith(1, "redis-1", 2, "b3JkZXJz", "cGF5bWVudHM=", undefined, "d29ya2VyLWI=");
const groupCrumb = host.querySelector<HTMLElement>("[data-redis-stream-group-crumb]")!;
const consumerCrumb = host.querySelector<HTMLElement>("[data-redis-stream-consumer-crumb]")!;
expect(consumerCrumb.textContent).toContain("worker-b");
expect(groupCrumb.classList.contains("border-transparent")).toBe(true);
expect(consumerCrumb.classList.contains("border-foreground")).toBe(true);
expect(host.querySelector("[data-redis-stream-consumers]")).toBeNull();
expect(host.querySelectorAll("[data-redis-stream-pending-header] th")).toHaveLength(3);
expect([...host.querySelectorAll("[data-redis-stream-pending-table] col")].map((column) => column.className)).toEqual(["w-[44%]", "w-[36%]", "w-[20%]"]);
expect(host.querySelector("[data-redis-stream-pending]")?.textContent).toContain("1714470000001-0");
expect(host.querySelector<HTMLTableCellElement>("[data-redis-stream-pending-table] tbody td")?.textContent?.trim()).toBe("1714470000001-0");
host.querySelector<HTMLButtonElement>("[data-redis-stream-group-crumb]")!.click();
await settle();
expect(mocks.redisGetStreamPending).toHaveBeenCalledTimes(1);
expect(host.querySelector("[data-redis-stream-consumers]")).not.toBeNull();
expect(host.querySelector("[data-redis-stream-consumer-crumb]")).toBeNull();
});
it("shows an empty state and lets a failed group query be retried", async () => {
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockRejectedValueOnce(new Error("NOPERM XINFO is not allowed")).mockResolvedValueOnce([]);
const host = mountViewer();
await settle();
openGroups(host);
await settle();
expect(host.querySelector("[data-redis-stream-groups-retry]")).not.toBeNull();
expect(host.textContent).toContain("NOPERM XINFO is not allowed");
host.querySelector<HTMLButtonElement>("[data-redis-stream-groups-retry]")!.click();
await settle();
expect(mocks.redisGetStreamGroups).toHaveBeenCalledTimes(2);
expect(host.textContent).toContain("No consumer groups");
});
it("manually refreshes the active group and its dependent views", async () => {
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValue([group()]);
mocks.redisGetStreamConsumers.mockResolvedValue([{ name: blob("d29ya2VyLWE="), pending: 0, idle_ms: 0 }]);
mocks.redisGetStreamPending.mockResolvedValue({ entries: [] });
const host = mountViewer();
await settle();
openGroups(host);
await settle();
host.querySelector<HTMLElement>("[data-redis-stream-group-row]")!.click();
await settle();
host.querySelector<HTMLButtonElement>("[data-redis-stream-consumer-row]")!.click();
await settle();
host.querySelector<HTMLButtonElement>("[data-redis-value-refresh]")!.click();
await settle();
expect(mocks.redisGetValue).toHaveBeenCalledTimes(2);
expect(mocks.redisGetStreamGroups).toHaveBeenCalledTimes(2);
expect(mocks.redisGetStreamConsumers).toHaveBeenCalledTimes(2);
expect(mocks.redisGetStreamPending).toHaveBeenCalledTimes(2);
});
it("returns to the consumer list when refresh removes the selected consumer", async () => {
const refreshedPending = deferred<{ entries: Array<{ id: string; consumer: ReturnType<typeof blob>; idle_ms: number; deliveries: number }> }>();
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValue([group()]);
mocks.redisGetStreamConsumers.mockResolvedValueOnce([{ name: blob("d29ya2VyLWE="), pending: 1, idle_ms: 1_200, inactive_ms: 800 }]).mockResolvedValueOnce([]);
mocks.redisGetStreamPending.mockResolvedValueOnce({ entries: [{ id: "1714470000000-0", consumer: blob("d29ya2VyLWE="), idle_ms: 2_400, deliveries: 2 }] }).mockReturnValueOnce(refreshedPending.promise);
const host = mountViewer();
await settle();
openGroups(host);
await settle();
host.querySelector<HTMLElement>("[data-redis-stream-group-row]")!.click();
await settle();
host.querySelector<HTMLButtonElement>("[data-redis-stream-consumer-row]")!.click();
await settle();
host.querySelector<HTMLButtonElement>("[data-redis-value-refresh]")!.click();
await settle();
expect(mocks.redisGetStreamConsumers).toHaveBeenCalledTimes(2);
expect(host.querySelector("[data-redis-stream-consumer-crumb]")).toBeNull();
expect(host.querySelector("[data-redis-stream-consumers]")).not.toBeNull();
expect(host.querySelector("[data-redis-stream-pending]")).toBeNull();
refreshedPending.resolve({ entries: [{ id: "1714470000001-0", consumer: blob("d29ya2VyLWE="), idle_ms: 1_000, deliveries: 1 }] });
await settle();
expect(host.querySelector("[data-redis-stream-consumer-crumb]")).toBeNull();
expect(host.querySelector("[data-redis-stream-pending]")).toBeNull();
});
it("returns to the group list when refresh removes the selected group", async () => {
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValueOnce([group()]).mockResolvedValueOnce([]);
mocks.redisGetStreamConsumers.mockResolvedValue([{ name: blob("d29ya2VyLWE="), pending: 1, idle_ms: 1_200, inactive_ms: 800 }]);
const host = mountViewer();
await settle();
openGroups(host);
await settle();
host.querySelector<HTMLElement>("[data-redis-stream-group-row]")!.click();
await settle();
host.querySelector<HTMLButtonElement>("[data-redis-value-refresh]")!.click();
await settle();
expect(mocks.redisGetStreamGroups).toHaveBeenCalledTimes(2);
expect(host.querySelector("[data-redis-stream-group-detail]")).toBeNull();
expect(host.querySelector("[data-redis-stream-groups]")).toBeNull();
expect(host.textContent).toContain("No consumer groups");
});
it("retries failed consumer and pending reads without retaining stale data", async () => {
mocks.redisGetValue.mockResolvedValue(streamValue());
mocks.redisGetStreamGroups.mockResolvedValue([group()]);
mocks.redisGetStreamConsumers.mockRejectedValueOnce(new Error("NOPERM XINFO is not allowed")).mockResolvedValueOnce([{ name: blob("d29ya2VyLWE="), pending: 1, idle_ms: 1_200, inactive_ms: 800 }]);
mocks.redisGetStreamPending.mockRejectedValueOnce(new Error("NOPERM XPENDING is not allowed")).mockResolvedValueOnce({
entries: [{ id: "1714470000000-0", consumer: blob("d29ya2VyLWE="), idle_ms: 2_400, deliveries: 2 }],
});
const host = mountViewer();
await settle();
openGroups(host);
await settle();
host.querySelector<HTMLElement>("[data-redis-stream-group-row]")!.click();
await settle();
expect(host.querySelector("[data-redis-stream-consumers-retry]")).not.toBeNull();
expect(host.textContent).toContain("NOPERM XINFO is not allowed");
host.querySelector<HTMLButtonElement>("[data-redis-stream-consumers-retry]")!.click();
await settle();
expect(mocks.redisGetStreamConsumers).toHaveBeenCalledTimes(2);
host.querySelector<HTMLButtonElement>("[data-redis-stream-consumer-row]")!.click();
await settle();
expect(host.querySelector("[data-redis-stream-pending-retry]")).not.toBeNull();
expect(host.textContent).toContain("NOPERM XPENDING is not allowed");
host.querySelector<HTMLButtonElement>("[data-redis-stream-pending-retry]")!.click();
await settle();
expect(mocks.redisGetStreamPending).toHaveBeenCalledTimes(2);
expect(host.querySelector("[data-redis-stream-pending]")?.textContent).toContain("1714470000000-0");
});
});
@@ -9,6 +9,7 @@ import { Button } from "@/components/ui/button";
import { Input } from "@/components/ui/input";
import { Badge } from "@/components/ui/badge";
import { Switch } from "@/components/ui/switch";
import { Tabs, TabsContent, TabsList, TabsTrigger } from "@/components/ui/tabs";
import DateTimePicker from "@/components/ui/date-time-picker/DateTimePicker.vue";
import { Dialog, DialogContent, DialogFooter, DialogHeader, DialogTitle } from "@/components/ui/dialog";
import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue } from "@/components/ui/select";
@@ -16,7 +17,7 @@ import DangerConfirmDialog from "@/components/editor/DangerConfirmDialog.vue";
import JsonTree from "@/components/common/JsonTree.vue";
import RedisJsonEditor from "@/components/redis/RedisJsonEditor.vue";
import * as api from "@/lib/backend/api";
import type { RedisBlob, RedisHashItem, RedisKeyInfo, RedisListItem, RedisSetItem, RedisStreamEntry, RedisValue, RedisZsetItem } from "@/lib/backend/api";
import type { RedisBlob, RedisHashItem, RedisKeyInfo, RedisListItem, RedisSetItem, RedisStreamConsumer, RedisStreamEntry, RedisStreamGroup, RedisStreamPendingEntry, RedisValue, RedisZsetItem } from "@/lib/backend/api";
import { useToast } from "@/composables/useToast";
import { useTheme } from "@/composables/useTheme";
import { useEditorFontFamilyStyle } from "@/composables/useEditorFontFamilyStyle";
@@ -80,6 +81,37 @@ const data = ref<RedisValue | null>(null);
const loading = ref(false);
const loadingMore = ref(false);
let loadRequestId = 0;
const streamTab = ref<"entries" | "groups">("entries");
const streamGroups = ref<RedisStreamGroup[]>([]);
const streamGroupsLoaded = ref(false);
const streamGroupsLoading = ref(false);
const streamGroupsError = ref("");
const selectedStreamGroup = ref<RedisStreamGroup | null>(null);
const selectedStreamConsumer = ref<RedisStreamConsumer | null>(null);
const streamConsumers = ref<RedisStreamConsumer[]>([]);
const streamConsumersLoading = ref(false);
const streamConsumersError = ref("");
const streamPendingEntries = ref<RedisStreamPendingEntry[]>([]);
const streamPendingCursor = ref<string | undefined>();
const streamPendingLoading = ref(false);
const streamPendingLoadingMore = ref(false);
const streamPendingError = ref("");
const streamDateTimeFormatter = computed(
() =>
new Intl.DateTimeFormat(locale.value, {
year: "numeric",
month: "2-digit",
day: "2-digit",
hour: "2-digit",
minute: "2-digit",
second: "2-digit",
hour12: false,
}),
);
let streamGroupsRequestId = 0;
let streamGroupDetailRequestId = 0;
let streamConsumersRequestId = 0;
let streamPendingRequestId = 0;
const editValue = ref("");
const savingString = ref(false);
const savingJson = ref(false);
@@ -402,6 +434,256 @@ type RedisStreamRow = {
entry: RedisStreamEntry;
};
function isSelectedStreamGroup(group: RedisStreamGroup, requestId = streamGroupDetailRequestId): boolean {
return requestId === streamGroupDetailRequestId && selectedStreamGroup.value?.name.raw_base64 === group.name.raw_base64;
}
function resetStreamGroupDetail() {
streamGroupDetailRequestId++;
streamConsumersRequestId++;
streamPendingRequestId++;
selectedStreamGroup.value = null;
selectedStreamConsumer.value = null;
streamConsumers.value = [];
streamConsumersLoading.value = false;
streamConsumersError.value = "";
streamPendingEntries.value = [];
streamPendingCursor.value = undefined;
streamPendingLoading.value = false;
streamPendingLoadingMore.value = false;
streamPendingError.value = "";
}
function resetStreamMonitoring() {
streamGroupsRequestId++;
resetStreamGroupDetail();
streamTab.value = "entries";
streamGroups.value = [];
streamGroupsLoaded.value = false;
streamGroupsLoading.value = false;
streamGroupsError.value = "";
}
async function loadStreamGroups(force = false): Promise<boolean> {
if (redisKind.value !== "stream") return false;
if (!force && (streamGroupsLoaded.value || streamGroupsLoading.value)) return streamGroupsLoaded.value;
const requestId = ++streamGroupsRequestId;
streamGroupsLoading.value = true;
streamGroupsError.value = "";
try {
const groups = await api.redisGetStreamGroups(props.connectionId, props.db, props.keyRaw);
if (requestId !== streamGroupsRequestId || redisKind.value !== "stream") return false;
streamGroups.value = groups;
streamGroupsLoaded.value = true;
const selected = selectedStreamGroup.value;
if (selected && !groups.some((group) => group.name.raw_base64 === selected.name.raw_base64)) {
resetStreamGroupDetail();
}
return true;
} catch (error) {
if (requestId === streamGroupsRequestId) streamGroupsError.value = errorMessage(error);
return false;
} finally {
if (requestId === streamGroupsRequestId) streamGroupsLoading.value = false;
}
}
async function loadStreamConsumers(group: RedisStreamGroup, requestId: number) {
if (!isSelectedStreamGroup(group, requestId)) return;
const consumerRequestId = ++streamConsumersRequestId;
streamConsumersLoading.value = true;
streamConsumersError.value = "";
try {
const consumers = await api.redisGetStreamConsumers(props.connectionId, props.db, props.keyRaw, group.name.raw_base64);
if (!isSelectedStreamGroup(group, requestId) || consumerRequestId !== streamConsumersRequestId) return;
streamConsumers.value = consumers;
const selectedConsumerRaw = selectedStreamConsumer.value?.name.raw_base64;
if (selectedConsumerRaw) {
const selectedConsumer = consumers.find((consumer) => consumer.name.raw_base64 === selectedConsumerRaw);
if (selectedConsumer) {
selectedStreamConsumer.value = selectedConsumer;
} else {
// A refresh can race with XGROUP DELCONSUMER. Do not leave a stale
// consumer detail open or allow its pending request to update the view.
resetStreamConsumerDetail();
}
}
} catch (error) {
if (isSelectedStreamGroup(group, requestId) && consumerRequestId === streamConsumersRequestId) {
streamConsumersError.value = errorMessage(error);
}
} finally {
if (isSelectedStreamGroup(group, requestId) && consumerRequestId === streamConsumersRequestId) {
streamConsumersLoading.value = false;
}
}
}
async function loadStreamPendingPage(group: RedisStreamGroup, cursor?: string, append = false, requestId = streamGroupDetailRequestId) {
if (!isSelectedStreamGroup(group, requestId)) return;
const pendingRequestId = ++streamPendingRequestId;
const consumerRaw = selectedStreamConsumer.value?.name.raw_base64;
if (append) {
streamPendingLoadingMore.value = true;
} else {
streamPendingLoading.value = true;
streamPendingError.value = "";
streamPendingEntries.value = [];
streamPendingCursor.value = undefined;
}
try {
const page = consumerRaw ? await api.redisGetStreamPending(props.connectionId, props.db, props.keyRaw, group.name.raw_base64, cursor, consumerRaw) : await api.redisGetStreamPending(props.connectionId, props.db, props.keyRaw, group.name.raw_base64, cursor);
if (!isSelectedStreamGroup(group, requestId) || pendingRequestId !== streamPendingRequestId || selectedStreamConsumer.value?.name.raw_base64 !== consumerRaw) return;
streamPendingEntries.value = append ? [...streamPendingEntries.value, ...page.entries] : page.entries;
streamPendingCursor.value = page.next_cursor;
} catch (error) {
if (isSelectedStreamGroup(group, requestId) && pendingRequestId === streamPendingRequestId && selectedStreamConsumer.value?.name.raw_base64 === consumerRaw) {
streamPendingError.value = errorMessage(error);
}
} finally {
if (isSelectedStreamGroup(group, requestId) && pendingRequestId === streamPendingRequestId && selectedStreamConsumer.value?.name.raw_base64 === consumerRaw) {
if (append) streamPendingLoadingMore.value = false;
else streamPendingLoading.value = false;
}
}
}
function selectStreamGroup(group: RedisStreamGroup, reload = false) {
if (!reload && selectedStreamGroup.value?.name.raw_base64 === group.name.raw_base64) return;
const preservedConsumer = reload && selectedStreamGroup.value?.name.raw_base64 === group.name.raw_base64 ? selectedStreamConsumer.value : null;
streamGroupDetailRequestId++;
streamConsumersRequestId++;
streamPendingRequestId++;
const requestId = streamGroupDetailRequestId;
selectedStreamGroup.value = group;
selectedStreamConsumer.value = preservedConsumer;
streamConsumers.value = [];
streamConsumersLoading.value = false;
streamConsumersError.value = "";
streamPendingEntries.value = [];
streamPendingCursor.value = undefined;
streamPendingLoading.value = false;
streamPendingLoadingMore.value = false;
streamPendingError.value = "";
void loadStreamConsumers(group, requestId);
if (selectedStreamConsumer.value) void loadStreamPendingPage(group, undefined, false, requestId);
}
function selectStreamConsumer(consumer: RedisStreamConsumer) {
const group = selectedStreamGroup.value;
if (!group || selectedStreamConsumer.value?.name.raw_base64 === consumer.name.raw_base64) return;
streamPendingRequestId++;
selectedStreamConsumer.value = consumer;
streamPendingEntries.value = [];
streamPendingCursor.value = undefined;
streamPendingLoading.value = false;
streamPendingLoadingMore.value = false;
streamPendingError.value = "";
void loadStreamPendingPage(group, undefined, false, streamGroupDetailRequestId);
}
function resetStreamConsumerDetail() {
if (!selectedStreamConsumer.value) return;
streamPendingRequestId++;
selectedStreamConsumer.value = null;
streamPendingEntries.value = [];
streamPendingCursor.value = undefined;
streamPendingLoading.value = false;
streamPendingLoadingMore.value = false;
streamPendingError.value = "";
}
function retryStreamGroups() {
void loadStreamGroups(true);
}
function retryStreamConsumers() {
const group = selectedStreamGroup.value;
if (!group) return;
void loadStreamConsumers(group, streamGroupDetailRequestId);
}
function retryStreamPending() {
const group = selectedStreamGroup.value;
if (!group || !selectedStreamConsumer.value) return;
void loadStreamPendingPage(group, undefined, false, streamGroupDetailRequestId);
}
function loadMoreStreamPending() {
const group = selectedStreamGroup.value;
const cursor = streamPendingCursor.value;
if (!group || !selectedStreamConsumer.value || !cursor || streamPendingLoading.value || streamPendingLoadingMore.value) return;
void loadStreamPendingPage(group, cursor, true, streamGroupDetailRequestId);
}
async function refreshValueAndStreamGroups() {
const refreshGroups = redisKind.value === "stream" && streamTab.value === "groups";
const selectedGroupRaw = selectedStreamGroup.value?.name.raw_base64;
try {
const applied = await load();
if (!applied || !refreshGroups || redisKind.value !== "stream") return;
const loaded = await loadStreamGroups(true);
if (!loaded || !selectedGroupRaw) return;
const selected = streamGroups.value.find((group) => group.name.raw_base64 === selectedGroupRaw);
if (selected) selectStreamGroup(selected, true);
} catch (error) {
toast(errorMessage(error), 3000);
}
}
function streamMetricInteger(value: number | string | undefined): bigint | null {
if (typeof value === "number") return Number.isSafeInteger(value) && value >= 0 ? BigInt(value) : null;
if (typeof value !== "string" || !/^\d+$/.test(value)) return null;
try {
return BigInt(value);
} catch {
return null;
}
}
function formatStreamMetric(value: number | string | undefined): string {
const integer = streamMetricInteger(value);
return integer == null ? "-" : integer.toLocaleString();
}
function formatStreamDuration(value: number | string | undefined): string {
const integer = streamMetricInteger(value);
if (integer == null) return "-";
if (integer > BigInt(Number.MAX_SAFE_INTEGER)) return `${integer.toLocaleString()} ms`;
const milliseconds = Number(integer);
if (milliseconds < 1_000) return `${milliseconds.toLocaleString()} ms`;
const seconds = milliseconds / 1_000;
if (seconds < 60) return `${seconds.toFixed(seconds < 10 ? 1 : 0)} s`;
const minutes = Math.floor(seconds / 60);
if (minutes < 60) return `${minutes}m ${Math.floor(seconds % 60)}s`;
const hours = Math.floor(minutes / 60);
if (hours < 24) return `${hours}h ${minutes % 60}m`;
return `${Math.floor(hours / 24)}d ${hours % 24}h`;
}
function formatStreamDurationTitle(value: number | string | undefined): string {
const formatted = formatStreamMetric(value);
return formatted === "-" ? formatted : `${formatted} ms`;
}
function formatStreamDateTime(timestamp: number): string {
if (!Number.isFinite(timestamp) || timestamp <= 0) return "-";
return streamDateTimeFormatter.value.format(new Date(timestamp));
}
function formatStreamLastDelivery(value: number | string | undefined): string {
const idle = streamMetricInteger(value);
if (idle == null || idle > BigInt(Number.MAX_SAFE_INTEGER)) return "-";
return formatStreamDateTime(Date.now() - Number(idle));
}
function collectionCountLabel(kind: "items" | "fields" | "members", loaded: number, total?: number | null) {
if (total == null || total === loaded) return t(`redis.${kind}`, { count: loaded });
return t(`redis.loaded${kind[0].toUpperCase()}${kind.slice(1)}`, { loaded, total });
@@ -620,6 +902,7 @@ async function load(options: { selectDefaultMember?: boolean; preserveDraft?: bo
data.value = null;
collectionItems.value = [];
scanCursor.value = undefined;
resetStreamMonitoring();
stopAutoRefresh();
emit("deleted", props.keyRaw);
return true;
@@ -646,6 +929,7 @@ async function load(options: { selectDefaultMember?: boolean; preserveDraft?: bo
emit("loaded", loadedValue);
scanCursor.value = redisValueCollectionScanCursor(loadedValue);
collectionItems.value = redisValueCollectionItems(loadedValue);
if (loadedValue.data.kind !== "stream") resetStreamMonitoring();
// A foreground load replaces the current value, so it also starts a new
// draft lifecycle. Member saves opt out until selection is restored.
@@ -1484,13 +1768,18 @@ watch(contentSearchText, () => {
});
watch(
() => props.keyRaw,
() => [props.connectionId, props.db, props.keyRaw],
() => {
resetValueSearch();
valueViewerSearchActive.value = false;
resetStreamMonitoring();
},
);
watch(streamTab, (tab) => {
if (tab === "groups" && redisKind.value === "stream") void loadStreamGroups();
});
watch(stringValueView, () => {
if (!showMemberDetail.value) {
valueSearchMatchIndex.value = 0;
@@ -1562,7 +1851,7 @@ defineExpose({ focusSearch });
<div class="shrink-0 border-b bg-background">
<div class="flex h-9 items-center gap-2 px-4">
<span class="dbx-editor-font-family min-w-0 flex-1 truncate text-sm font-semibold">{{ formatValue(data.key_display) }}</span>
<Button variant="ghost" size="icon" class="h-7 w-7 shrink-0 animate-none" :disabled="hasUnsavedRedisDraft" @click="load"><RefreshCw class="h-3.5 w-3.5 animate-none" /></Button>
<Button data-redis-value-refresh variant="ghost" size="icon" class="h-7 w-7 shrink-0 animate-none" :disabled="hasUnsavedRedisDraft" @click="refreshValueAndStreamGroups"><RefreshCw class="h-3.5 w-3.5 animate-none" /></Button>
<Button variant="ghost" size="icon" class="h-7 w-7 shrink-0" :title="t('grid.copyValue')" :aria-label="t('grid.copyValue')" @click="copyValue"><Copy class="h-3.5 w-3.5" /></Button>
<Button variant="ghost" size="icon" class="h-7 w-7 shrink-0" :title="t('redis.copyInsertStatement')" :aria-label="t('redis.copyInsertStatement')" @click="copyInsertStatement"><ClipboardCopy class="h-3.5 w-3.5" /></Button>
<Button variant="ghost" size="icon" class="h-7 w-7 shrink-0 text-destructive" @click="requestDeleteKey"><Trash2 class="h-3.5 w-3.5" /></Button>
@@ -1904,36 +2193,272 @@ defineExpose({ focusSearch });
</RecycleScroller>
</div>
<!-- Stream (readonly) -->
<!-- Stream monitoring -->
<div v-else-if="redisKind === 'stream'" class="flex-1 flex flex-col overflow-hidden">
<div class="px-4 py-1 text-xs text-muted-foreground border-b shrink-0">
{{ t("redis.entries", { count: streamRows.length }) }}
</div>
<DynamicScroller class="flex-1 overflow-y-auto" :items="streamRows" :min-item-size="REDIS_STREAM_MIN_ROW_HEIGHT" :buffer="600" key-field="id">
<template #default="{ item: row, active }">
<DynamicScrollerItem :item="row" :active="active" :size-dependencies="[streamFieldCount(row)]" :data-index="row.index">
<div data-redis-stream-entry class="dbx-editor-font-family px-4 py-2 border-b text-sm hover:bg-accent/50">
<div class="mb-1 text-xs text-muted-foreground">{{ row.entry.id }}</div>
<div
v-for="(field, fieldIndex) in row.entry.fields"
:key="`${row.id}:${field.field}:${fieldIndex}`"
class="grid grid-cols-[minmax(6rem,0.35fr)_1fr_56px] gap-3 py-0.5 group cursor-pointer"
:class="{ 'bg-accent/60': isSelectedMember(field.field, field.value, streamFieldSelectionIdentity(row.entry.id, fieldIndex)) }"
@click="viewMember(field.field, field.value, { kind: 'stream', field: field.field, canEdit: false }, streamFieldSelectionIdentity(row.entry.id, fieldIndex))"
>
<span class="truncate text-blue-500">{{ field.field }}</span>
<span class="truncate text-muted-foreground">{{ field.value }}</span>
<span class="flex justify-end gap-1">
<Button variant="ghost" size="icon" class="h-5 w-5 opacity-0 group-hover:opacity-100" :title="t('redis.viewMember')" @click.stop="viewMember(field.field, field.value, { kind: 'stream', field: field.field, canEdit: false }, streamFieldSelectionIdentity(row.entry.id, fieldIndex))"
><Eye class="w-3 h-3"
/></Button>
<Button variant="ghost" size="icon" class="h-5 w-5 opacity-0 group-hover:opacity-100" :title="t('redis.copyMember')" @click.stop="copyMember(field.value)"><Copy class="w-3 h-3" /></Button>
</span>
<Tabs v-model="streamTab" :unmount-on-hide="false" class="h-full min-h-0 gap-0">
<div class="flex h-9 shrink-0 items-stretch overflow-x-auto border-b px-4">
<TabsList variant="line" class="h-full shrink-0 gap-0 p-0">
<TabsTrigger value="entries" data-redis-stream-entries-tab class="h-full flex-none rounded-none px-3 text-xs group-data-horizontal/tabs:after:bottom-0">{{ t("redis.streamData") }}</TabsTrigger>
<TabsTrigger
value="groups"
data-redis-stream-groups-tab
:class="['h-full flex-none rounded-none px-3 text-xs group-data-horizontal/tabs:after:bottom-0', selectedStreamGroup && 'data-active:text-foreground/60 group-data-[variant=line]/tabs-list:data-active:after:opacity-0']"
@click="resetStreamGroupDetail"
>
{{ t("redis.consumerGroups") }}
</TabsTrigger>
</TabsList>
<template v-if="streamTab === 'groups' && selectedStreamGroup">
<span class="mx-1 h-4 w-px shrink-0 self-center bg-border" aria-hidden="true" />
<button
data-redis-stream-group-crumb
type="button"
class="relative inline-flex h-full max-w-48 shrink-0 items-center border-b-2 px-3 text-left text-xs font-medium transition-colors focus-visible:outline-none focus-visible:ring-2 focus-visible:ring-ring focus-visible:ring-inset"
:class="selectedStreamConsumer ? 'border-transparent text-foreground/60 hover:text-foreground' : 'border-foreground text-foreground'"
:title="formatValue(selectedStreamGroup.name)"
@click="resetStreamConsumerDetail"
>
<span class="dbx-editor-font-family truncate">{{ formatValue(selectedStreamGroup.name) }}</span>
</button>
<span v-if="selectedStreamConsumer" data-redis-stream-consumer-crumb class="dbx-editor-font-family inline-flex h-full max-w-48 shrink-0 items-center truncate border-b-2 border-foreground px-3 text-xs font-medium text-foreground" :title="formatValue(selectedStreamConsumer.name)">
<span class="truncate">{{ formatValue(selectedStreamConsumer.name) }}</span>
</span>
</template>
</div>
<TabsContent value="entries" class="m-0 min-h-0 flex-1 flex flex-col">
<div class="px-4 py-1 text-xs text-muted-foreground border-b shrink-0">
{{ t("redis.entries", { count: streamRows.length }) }}
</div>
<DynamicScroller class="flex-1 overflow-y-auto" :items="streamRows" :min-item-size="REDIS_STREAM_MIN_ROW_HEIGHT" :buffer="600" key-field="id">
<template #default="{ item: row, active }">
<DynamicScrollerItem :item="row" :active="active" :size-dependencies="[streamFieldCount(row)]" :data-index="row.index">
<div data-redis-stream-entry class="dbx-editor-font-family px-4 py-2 border-b text-sm hover:bg-accent/50">
<div class="mb-1 text-xs text-muted-foreground">{{ row.entry.id }}</div>
<div
v-for="(field, fieldIndex) in row.entry.fields"
:key="`${row.id}:${field.field}:${fieldIndex}`"
class="grid grid-cols-[minmax(6rem,0.35fr)_1fr_56px] gap-3 py-0.5 group cursor-pointer"
:class="{ 'bg-accent/60': isSelectedMember(field.field, field.value, streamFieldSelectionIdentity(row.entry.id, fieldIndex)) }"
@click="viewMember(field.field, field.value, { kind: 'stream', field: field.field, canEdit: false }, streamFieldSelectionIdentity(row.entry.id, fieldIndex))"
>
<span class="truncate text-blue-500">{{ field.field }}</span>
<span class="truncate text-muted-foreground">{{ field.value }}</span>
<span class="flex justify-end gap-1">
<Button
variant="ghost"
size="icon"
class="h-5 w-5 opacity-0 group-hover:opacity-100"
:title="t('redis.viewMember')"
@click.stop="viewMember(field.field, field.value, { kind: 'stream', field: field.field, canEdit: false }, streamFieldSelectionIdentity(row.entry.id, fieldIndex))"
><Eye class="w-3 h-3"
/></Button>
<Button variant="ghost" size="icon" class="h-5 w-5 opacity-0 group-hover:opacity-100" :title="t('redis.copyMember')" @click.stop="copyMember(field.value)"><Copy class="w-3 h-3" /></Button>
</span>
</div>
</div>
</DynamicScrollerItem>
</template>
</DynamicScroller>
</TabsContent>
<TabsContent value="groups" class="m-0 min-h-0 flex-1 flex flex-col">
<template v-if="!selectedStreamGroup">
<div class="flex shrink-0 items-center gap-2 border-b px-4 py-1.5 text-xs text-muted-foreground">
<span>{{ t("redis.consumerGroups") }}</span>
<Loader2 v-if="streamGroupsLoading" class="h-3 w-3 animate-spin" />
</div>
<div v-if="streamGroupsLoading && !streamGroupsLoaded" class="flex flex-1 items-center justify-center text-xs text-muted-foreground">
<Loader2 class="mr-2 h-4 w-4 animate-spin" />
{{ t("common.loading") }}
</div>
<div v-else-if="streamGroupsError && streamGroups.length === 0" class="flex flex-1 flex-col items-center justify-center gap-3 p-4 text-center text-xs text-muted-foreground">
<span>{{ t("redis.streamGroupsLoadFailed") }}: {{ streamGroupsError }}</span>
<Button data-redis-stream-groups-retry variant="outline" size="sm" class="h-7 text-xs" @click="retryStreamGroups">{{ t("common.retry") }}</Button>
</div>
<div v-else class="min-h-0 flex-1 overflow-auto">
<div v-if="streamGroupsError" class="flex items-center gap-2 border-b bg-destructive/5 px-4 py-2 text-xs text-destructive">
<span class="min-w-0 flex-1 truncate">{{ t("redis.streamGroupsLoadFailed") }}: {{ streamGroupsError }}</span>
<Button data-redis-stream-groups-retry variant="ghost" size="sm" class="h-6 text-xs text-destructive" @click="retryStreamGroups">{{ t("common.retry") }}</Button>
</div>
<div v-if="streamGroups.length === 0" class="flex h-full min-h-40 items-center justify-center p-4 text-xs text-muted-foreground">
{{ t("redis.noConsumerGroups") }}
</div>
<table v-else data-redis-stream-groups class="w-full min-w-[780px] border-collapse text-left text-sm">
<thead class="sticky top-0 bg-muted/95 text-xs text-muted-foreground backdrop-blur">
<tr>
<th class="px-4 py-2 font-medium">{{ t("redis.group") }}</th>
<th class="px-3 py-2 text-right font-medium">{{ t("redis.consumers") }}</th>
<th class="px-3 py-2 text-right font-medium">{{ t("redis.pending") }}</th>
<th class="px-3 py-2 font-medium">{{ t("redis.lastDeliveredId") }}</th>
<th class="px-3 py-2 text-right font-medium">{{ t("redis.entriesRead") }}</th>
<th class="px-4 py-2 text-right font-medium">{{ t("redis.lag") }}</th>
</tr>
</thead>
<tbody>
<tr v-for="group in streamGroups" :key="group.name.raw_base64" data-redis-stream-group-row class="cursor-pointer border-t hover:bg-accent/50" @click="selectStreamGroup(group)">
<td class="dbx-editor-font-family max-w-72 truncate px-4 py-2" :title="formatValue(group.name)">{{ formatValue(group.name) }}</td>
<td class="px-3 py-2 text-right tabular-nums">{{ formatStreamMetric(group.consumers) }}</td>
<td class="px-3 py-2 text-right tabular-nums">{{ formatStreamMetric(group.pending) }}</td>
<td class="dbx-editor-font-family max-w-56 truncate px-3 py-2 text-muted-foreground" :title="group.last_delivered_id">{{ group.last_delivered_id }}</td>
<td class="px-3 py-2 text-right tabular-nums">{{ formatStreamMetric(group.entries_read) }}</td>
<td class="px-4 py-2 text-right tabular-nums">{{ formatStreamMetric(group.lag) }}</td>
</tr>
</tbody>
</table>
</div>
</template>
<template v-else>
<div data-redis-stream-group-detail class="min-h-0 flex-1 overflow-auto bg-muted/20 p-3">
<div class="mx-auto flex w-full max-w-[1120px] flex-col gap-3">
<div v-if="selectedStreamConsumer" data-redis-stream-consumer-summary class="grid grid-cols-3 gap-px overflow-hidden rounded-lg border bg-border">
<div class="bg-card px-3 py-2.5">
<div class="text-xs text-muted-foreground">{{ t("redis.pending") }}</div>
<div class="mt-1 text-lg font-semibold tabular-nums">{{ formatStreamMetric(selectedStreamConsumer.pending) }}</div>
</div>
<div class="bg-card px-3 py-2.5">
<div class="text-xs text-muted-foreground">{{ t("redis.idle") }}</div>
<div class="mt-1 text-lg font-semibold tabular-nums" :title="formatStreamDurationTitle(selectedStreamConsumer.idle_ms)">{{ formatStreamDuration(selectedStreamConsumer.idle_ms) }}</div>
</div>
<div class="bg-card px-3 py-2.5">
<div class="text-xs text-muted-foreground">{{ t("redis.inactive") }}</div>
<div class="mt-1 text-lg font-semibold tabular-nums" :title="formatStreamDurationTitle(selectedStreamConsumer.inactive_ms)">{{ formatStreamDuration(selectedStreamConsumer.inactive_ms) }}</div>
</div>
</div>
<div v-else data-redis-stream-group-summary class="grid grid-cols-3 gap-px overflow-hidden rounded-lg border bg-border">
<div class="bg-card px-3 py-2.5">
<div class="text-xs text-muted-foreground">{{ t("redis.consumers") }}</div>
<div class="mt-1 text-lg font-semibold tabular-nums">{{ formatStreamMetric(selectedStreamGroup.consumers) }}</div>
</div>
<div class="bg-card px-3 py-2.5">
<div class="text-xs text-muted-foreground">{{ t("redis.pending") }}</div>
<div class="mt-1 text-lg font-semibold tabular-nums">{{ formatStreamMetric(selectedStreamGroup.pending) }}</div>
</div>
<div class="bg-card px-3 py-2.5">
<div class="text-xs text-muted-foreground">{{ t("redis.lag") }}</div>
<div class="mt-1 text-lg font-semibold tabular-nums">{{ formatStreamMetric(selectedStreamGroup.lag) }}</div>
</div>
</div>
<section v-if="!selectedStreamConsumer" data-redis-stream-consumers class="overflow-hidden rounded-lg border bg-card text-card-foreground">
<div class="flex min-h-10 items-center gap-2 border-b px-3 py-2">
<span class="text-sm font-medium">{{ t("redis.streamConsumers") }}</span>
<Badge v-if="streamConsumers.length" variant="secondary" class="h-5 min-w-5 justify-center px-1.5 text-[10px] tabular-nums">{{ formatStreamMetric(streamConsumers.length) }}</Badge>
<Loader2 v-if="streamConsumersLoading" class="ml-auto h-3.5 w-3.5 animate-spin text-muted-foreground" />
</div>
<div v-if="streamConsumersLoading && streamConsumers.length === 0" class="flex items-center gap-2 px-3 py-6 text-xs text-muted-foreground">
<Loader2 class="h-3.5 w-3.5 animate-spin" />
{{ t("common.loading") }}
</div>
<div v-else-if="streamConsumersError && streamConsumers.length === 0" class="flex items-center gap-2 px-3 py-5 text-xs text-destructive">
<span class="min-w-0 flex-1 truncate">{{ t("redis.streamConsumersLoadFailed") }}: {{ streamConsumersError }}</span>
<Button data-redis-stream-consumers-retry variant="ghost" size="sm" class="h-7 text-xs text-destructive" @click="retryStreamConsumers">{{ t("common.retry") }}</Button>
</div>
<template v-else>
<div v-if="streamConsumersError" class="flex items-center gap-2 border-b bg-destructive/5 px-3 py-2 text-xs text-destructive">
<span class="min-w-0 flex-1 truncate">{{ t("redis.streamConsumersLoadFailed") }}: {{ streamConsumersError }}</span>
<Button data-redis-stream-consumers-retry variant="ghost" size="sm" class="h-6 text-xs text-destructive" @click="retryStreamConsumers">{{ t("common.retry") }}</Button>
</div>
<div v-if="streamConsumers.length === 0" class="px-3 py-6 text-center text-xs text-muted-foreground">{{ t("redis.noStreamConsumers") }}</div>
<table v-else data-redis-stream-consumers-table class="w-full table-fixed border-collapse text-sm">
<colgroup>
<col class="w-1/4" />
<col class="w-1/4" />
<col class="w-1/4" />
<col class="w-1/4" />
</colgroup>
<thead class="bg-muted/50 text-xs text-muted-foreground">
<tr data-redis-stream-consumer-header>
<th scope="col" class="px-3 py-1.5 text-left font-medium">{{ t("redis.consumer") }}</th>
<th scope="col" class="px-2 py-1.5 text-right font-medium">{{ t("redis.pending") }}</th>
<th scope="col" class="px-2 py-1.5 text-right font-medium">{{ t("redis.idle") }}</th>
<th scope="col" class="px-3 py-1.5 text-right font-medium">{{ t("redis.inactive") }}</th>
</tr>
</thead>
<tbody>
<tr v-for="consumer in streamConsumers" :key="consumer.name.raw_base64" class="border-t hover:bg-muted/20">
<td class="px-3 py-2 text-left">
<button
data-redis-stream-consumer-row
type="button"
class="block w-full min-w-0 rounded-sm text-left outline-none transition-colors hover:text-foreground focus-visible:ring-2 focus-visible:ring-ring focus-visible:ring-inset"
:title="t('redis.openConsumerDetails')"
@click="selectStreamConsumer(consumer)"
>
<span class="dbx-editor-font-family block truncate">{{ formatValue(consumer.name) }}</span>
</button>
</td>
<td class="px-2 py-2 text-right tabular-nums">{{ formatStreamMetric(consumer.pending) }}</td>
<td class="px-2 py-2 text-right tabular-nums" :title="formatStreamDurationTitle(consumer.idle_ms)">{{ formatStreamDuration(consumer.idle_ms) }}</td>
<td class="px-3 py-2 text-right tabular-nums" :title="formatStreamDurationTitle(consumer.inactive_ms)">{{ formatStreamDuration(consumer.inactive_ms) }}</td>
</tr>
</tbody>
</table>
</template>
</section>
<section v-if="selectedStreamConsumer" data-redis-stream-pending class="overflow-hidden rounded-lg border bg-card text-card-foreground">
<div class="flex min-h-10 items-center gap-2 border-b px-3 py-2">
<span class="text-sm font-medium">{{ t("redis.pendingEntries") }}</span>
<Badge v-if="selectedStreamConsumer" variant="outline" class="dbx-editor-font-family min-w-0 max-w-48 truncate px-1.5 text-[10px]" :title="formatValue(selectedStreamConsumer.name)">{{ formatValue(selectedStreamConsumer.name) }}</Badge>
<Badge v-if="streamPendingEntries.length" variant="secondary" class="h-5 min-w-5 justify-center px-1.5 text-[10px] tabular-nums">{{ formatStreamMetric(streamPendingEntries.length) }}</Badge>
<Loader2 v-if="streamPendingLoading || streamPendingLoadingMore" class="ml-auto h-3.5 w-3.5 animate-spin text-muted-foreground" />
</div>
<div v-if="streamPendingLoading && streamPendingEntries.length === 0" class="flex items-center gap-2 px-3 py-6 text-xs text-muted-foreground">
<Loader2 class="h-3.5 w-3.5 animate-spin" />
{{ t("common.loading") }}
</div>
<div v-else-if="streamPendingError && streamPendingEntries.length === 0" class="flex items-center gap-2 px-3 py-5 text-xs text-destructive">
<span class="min-w-0 flex-1 truncate">{{ t("redis.pendingEntriesLoadFailed") }}: {{ streamPendingError }}</span>
<Button data-redis-stream-pending-retry variant="ghost" size="sm" class="h-7 text-xs text-destructive" @click="retryStreamPending">{{ t("common.retry") }}</Button>
</div>
<template v-else>
<div v-if="streamPendingError" class="flex items-center gap-2 border-b bg-destructive/5 px-3 py-2 text-xs text-destructive">
<span class="min-w-0 flex-1 truncate">{{ t("redis.pendingEntriesLoadFailed") }}: {{ streamPendingError }}</span>
<Button data-redis-stream-pending-retry variant="ghost" size="sm" class="h-6 text-xs text-destructive" @click="retryStreamPending">{{ t("common.retry") }}</Button>
</div>
<div v-if="streamPendingEntries.length === 0" class="px-3 py-6 text-center text-xs text-muted-foreground">{{ t("redis.noPendingEntries") }}</div>
<table v-else data-redis-stream-pending-table class="w-full table-fixed border-collapse text-sm">
<colgroup>
<col class="w-[44%]" />
<col class="w-[36%]" />
<col class="w-[20%]" />
</colgroup>
<thead class="bg-muted/50 text-xs text-muted-foreground">
<tr data-redis-stream-pending-header>
<th scope="col" class="px-3 py-1.5 text-left font-medium">{{ t("redis.entryId") }}</th>
<th scope="col" class="px-2 py-1.5 text-left font-medium">{{ t("redis.lastDelivered") }}</th>
<th scope="col" class="px-3 py-1.5 text-right font-medium">{{ t("redis.deliveries") }}</th>
</tr>
</thead>
<tbody>
<tr v-for="entry in streamPendingEntries" :key="entry.id" class="border-t hover:bg-muted/20">
<td class="px-3 py-2 text-left" :title="entry.id">
<div class="dbx-editor-font-family truncate">{{ entry.id }}</div>
</td>
<td class="px-2 py-2 text-left tabular-nums" :title="formatStreamDurationTitle(entry.idle_ms)">{{ formatStreamLastDelivery(entry.idle_ms) }}</td>
<td class="px-3 py-2 text-right tabular-nums">{{ formatStreamMetric(entry.deliveries) }}</td>
</tr>
</tbody>
</table>
</template>
<div v-if="streamPendingCursor" class="border-t p-2">
<Button data-redis-stream-pending-more variant="outline" size="sm" class="h-7 w-full text-xs" :disabled="streamPendingLoadingMore" @click="loadMoreStreamPending">
<Loader2 v-if="streamPendingLoadingMore" class="mr-1.5 h-3 w-3 animate-spin" />
{{ t("redis.loadMorePending") }}
</Button>
</div>
</section>
</div>
</div>
</DynamicScrollerItem>
</template>
</DynamicScroller>
</template>
</TabsContent>
</Tabs>
</div>
<!-- Unknown -->
+25
View File
@@ -2940,6 +2940,31 @@ export default {
loadedFields: "{loaded} of {total} fields loaded",
loadedMembers: "{loaded} of {total} members loaded",
entries: "{count} entries",
streamData: "Stream Data",
consumerGroups: "Consumer Groups",
group: "Group",
consumers: "Consumers",
pending: "Pending",
lastDeliveredId: "Last Delivered ID",
lastDelivered: "Last Delivered",
entriesRead: "Entries Read",
lag: "Lag",
noConsumerGroups: "No consumer groups",
streamGroupsLoadFailed: "Failed to load consumer groups",
backToGroups: "Back to groups",
streamConsumers: "Consumers",
streamConsumersLoadFailed: "Failed to load consumers",
noStreamConsumers: "No consumers",
openConsumerDetails: "Open consumer details",
consumer: "Consumer",
idle: "Idle",
inactive: "Inactive",
pendingEntries: "Pending Entries",
pendingEntriesLoadFailed: "Failed to load pending entries",
noPendingEntries: "No pending entries",
entryId: "Entry ID",
deliveries: "Deliveries",
loadMorePending: "Load more pending entries",
noExpiry: "no expiry",
expiry: "Expiration",
expiryNone: "No expiry",
+25
View File
@@ -2801,6 +2801,31 @@ export default withEnglishFallback({
loadedFields: "{loaded} de {total} campos cargados",
loadedMembers: "{loaded} de {total} miembros cargados",
entries: "{count} entradas",
streamData: "Datos del stream",
consumerGroups: "Grupos de consumidores",
group: "Grupo",
consumers: "Consumidores",
pending: "Pendientes",
lastDeliveredId: "Último ID entregado",
lastDelivered: "Última entrega",
entriesRead: "Entradas leídas",
lag: "Retraso",
noConsumerGroups: "No hay grupos de consumidores",
streamGroupsLoadFailed: "No se pudieron cargar los grupos de consumidores",
backToGroups: "Volver a los grupos",
streamConsumers: "Consumidores",
streamConsumersLoadFailed: "No se pudieron cargar los consumidores",
noStreamConsumers: "No hay consumidores",
openConsumerDetails: "Abrir detalles del consumidor",
consumer: "Consumidor",
idle: "Inactivo",
inactive: "Inactivo",
pendingEntries: "Entradas pendientes",
pendingEntriesLoadFailed: "No se pudieron cargar las entradas pendientes",
noPendingEntries: "No hay entradas pendientes",
entryId: "ID de entrada",
deliveries: "Entregas",
loadMorePending: "Cargar más entradas pendientes",
noExpiry: "sin expiración",
expiry: "Expiración",
expiryNone: "Sin expiración",
+25
View File
@@ -2799,6 +2799,31 @@ export default withEnglishFallback({
loadedFields: "{loaded} di {total} campi caricati",
loadedMembers: "{loaded} di {total} membri caricati",
entries: "{count} voci",
streamData: "Dati stream",
consumerGroups: "Gruppi di consumer",
group: "Gruppo",
consumers: "Consumer",
pending: "In sospeso",
lastDeliveredId: "Ultimo ID consegnato",
lastDelivered: "Ultima consegna",
entriesRead: "Voci lette",
lag: "Ritardo",
noConsumerGroups: "Nessun gruppo di consumer",
streamGroupsLoadFailed: "Impossibile caricare i gruppi di consumer",
backToGroups: "Torna ai gruppi",
streamConsumers: "Consumer",
streamConsumersLoadFailed: "Impossibile caricare i consumer",
noStreamConsumers: "Nessun consumer",
openConsumerDetails: "Apri dettagli del consumer",
consumer: "Consumer",
idle: "Inattivo",
inactive: "Non attivo",
pendingEntries: "Voci in sospeso",
pendingEntriesLoadFailed: "Impossibile caricare le voci in sospeso",
noPendingEntries: "Nessuna voce in sospeso",
entryId: "ID voce",
deliveries: "Consegne",
loadMorePending: "Carica altre voci in sospeso",
noExpiry: "nessuna scadenza",
expiry: "Scadenza",
expiryNone: "Nessuna scadenza",
+25
View File
@@ -2800,6 +2800,31 @@ export default withEnglishFallback({
loadedFields: "{loaded}/{total}フィールド読み込み完了",
loadedMembers: "{loaded}/{total}メンバー読み込み完了",
entries: "{count}エントリ",
streamData: "ストリームデータ",
consumerGroups: "コンシューマーグループ",
group: "グループ",
consumers: "コンシューマー",
pending: "保留中",
lastDeliveredId: "最終配信 ID",
lastDelivered: "最終配信時刻",
entriesRead: "読み取り済みエントリ",
lag: "遅延",
noConsumerGroups: "コンシューマーグループはありません",
streamGroupsLoadFailed: "コンシューマーグループの読み込みに失敗しました",
backToGroups: "グループに戻る",
streamConsumers: "コンシューマー",
streamConsumersLoadFailed: "コンシューマーの読み込みに失敗しました",
noStreamConsumers: "コンシューマーはいません",
openConsumerDetails: "コンシューマーの詳細を開く",
consumer: "コンシューマー",
idle: "アイドル",
inactive: "非アクティブ",
pendingEntries: "保留中のエントリ",
pendingEntriesLoadFailed: "保留中のエントリの読み込みに失敗しました",
noPendingEntries: "保留中のエントリはありません",
entryId: "エントリ ID",
deliveries: "配信回数",
loadMorePending: "保留中のエントリをさらに読み込む",
noExpiry: "期限なし",
expiry: "有効期限",
expiryNone: "期限なし",
+25
View File
@@ -2801,6 +2801,31 @@ export default withEnglishFallback({
loadedFields: "{loaded} de {total} campos carregados",
loadedMembers: "{loaded} de {total} membros carregados",
entries: "{count} entradas",
streamData: "Dados do stream",
consumerGroups: "Grupos de consumidores",
group: "Grupo",
consumers: "Consumidores",
pending: "Pendentes",
lastDeliveredId: "Último ID entregue",
lastDelivered: "Última entrega",
entriesRead: "Entradas lidas",
lag: "Atraso",
noConsumerGroups: "Nenhum grupo de consumidores",
streamGroupsLoadFailed: "Não foi possível carregar os grupos de consumidores",
backToGroups: "Voltar aos grupos",
streamConsumers: "Consumidores",
streamConsumersLoadFailed: "Não foi possível carregar os consumidores",
noStreamConsumers: "Nenhum consumidor",
openConsumerDetails: "Abrir detalhes do consumidor",
consumer: "Consumidor",
idle: "Ocioso",
inactive: "Inativo",
pendingEntries: "Entradas pendentes",
pendingEntriesLoadFailed: "Não foi possível carregar as entradas pendentes",
noPendingEntries: "Nenhuma entrada pendente",
entryId: "ID da entrada",
deliveries: "Entregas",
loadMorePending: "Carregar mais entradas pendentes",
noExpiry: "sem expiração",
expiry: "Expiração",
expiryNone: "Sem expiração",
+25
View File
@@ -2940,6 +2940,31 @@ export default withEnglishFallback({
loadedFields: "已加载 {loaded} / 共 {total} 个字段",
loadedMembers: "已加载 {loaded} / 共 {total} 个成员",
entries: "{count} 条记录",
streamData: "Stream 数据",
consumerGroups: "消费组",
group: "组",
consumers: "消费者",
pending: "待处理",
lastDeliveredId: "最后投递 ID",
lastDelivered: "最后投递时间",
entriesRead: "已读取条目",
lag: "滞后",
noConsumerGroups: "暂无消费组",
streamGroupsLoadFailed: "加载消费组失败",
backToGroups: "返回消费组",
streamConsumers: "消费者",
streamConsumersLoadFailed: "加载消费者失败",
noStreamConsumers: "暂无消费者",
openConsumerDetails: "查看消费者详情",
consumer: "消费者",
idle: "空闲",
inactive: "非活跃",
pendingEntries: "待处理条目",
pendingEntriesLoadFailed: "加载待处理条目失败",
noPendingEntries: "暂无待处理条目",
entryId: "条目 ID",
deliveries: "投递次数",
loadMorePending: "加载更多待处理条目",
noExpiry: "永不过期",
expiry: "过期时间",
expiryNone: "永不过期",
+25
View File
@@ -2472,6 +2472,31 @@ export default withEnglishFallback({
loadedFields: "已載入 {loaded} / 共 {total} 個欄位",
loadedMembers: "已載入 {loaded} / 共 {total} 個成員",
entries: "{count} 筆項目",
streamData: "Stream 資料",
consumerGroups: "消費者群組",
group: "群組",
consumers: "消費者",
pending: "待處理",
lastDeliveredId: "最後投遞 ID",
lastDelivered: "最後投遞時間",
entriesRead: "已讀取項目",
lag: "延遲",
noConsumerGroups: "暫無消費者群組",
streamGroupsLoadFailed: "載入消費者群組失敗",
backToGroups: "返回消費者群組",
streamConsumers: "消費者",
streamConsumersLoadFailed: "載入消費者失敗",
noStreamConsumers: "暫無消費者",
openConsumerDetails: "檢視消費者詳細資料",
consumer: "消費者",
idle: "閒置",
inactive: "非活躍",
pendingEntries: "待處理項目",
pendingEntriesLoadFailed: "載入待處理項目失敗",
noPendingEntries: "暫無待處理項目",
entryId: "項目 ID",
deliveries: "投遞次數",
loadMorePending: "載入更多待處理項目",
noExpiry: "永不過期",
expiry: "到期時間",
expiryNone: "永不過期",
@@ -243,7 +243,7 @@ describe("native RedisJSON editor", () => {
const labels = templateElements(parsedViewer.descriptor.template!.ast as unknown as TemplateElement);
const stringTextarea = findTemplateElement((element) => element.tag === "textarea" && directiveExpression(element, "model") === "editValue");
const memberTextarea = findTemplateElement((element) => element.tag === "textarea" && directiveExpression(element, "model") === "memberEditValue");
const refreshButton = findTemplateElement((element) => element.tag === "Button" && directiveExpression(element, "on", "click") === "load");
const refreshButton = findTemplateElement((element) => element.tag === "Button" && directiveExpression(element, "on", "click") === "refreshValueAndStreamGroups");
expect(autoRefresh.getText()).toContain("if (hasUnsavedRedisDraft.value) return;");
expect(autoRefresh.getText()).toContain("load({ preserveDraft: true })");
@@ -0,0 +1,127 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
const mocks = vi.hoisted(() => ({
invoke: vi.fn(),
}));
vi.mock("@tauri-apps/api/core", () => ({
invoke: mocks.invoke,
}));
vi.mock("@tauri-apps/api/event", () => ({
listen: vi.fn(),
}));
describe("Redis stream monitoring Tauri API", () => {
beforeEach(() => {
vi.clearAllMocks();
vi.unstubAllGlobals();
});
it("invokes the group, consumer, and pending read commands with binary-safe raw names", async () => {
mocks.invoke.mockResolvedValue(undefined);
const { redisGetStreamConsumers, redisGetStreamGroups, redisGetStreamPending } = await import("@/lib/backend/tauri");
await redisGetStreamGroups("redis-1", 2, "b3JkZXJz");
await redisGetStreamConsumers("redis-1", 2, "b3JkZXJz", "Z3JvdXAALQ==");
await redisGetStreamPending("redis-1", 2, "b3JkZXJz", "Z3JvdXAALQ==", "1714470000000-17");
expect(mocks.invoke).toHaveBeenNthCalledWith(1, "redis_get_stream_groups", {
connectionId: "redis-1",
db: 2,
keyRaw: "b3JkZXJz",
});
expect(mocks.invoke).toHaveBeenNthCalledWith(2, "redis_get_stream_consumers", {
connectionId: "redis-1",
db: 2,
keyRaw: "b3JkZXJz",
groupRaw: "Z3JvdXAALQ==",
});
expect(mocks.invoke).toHaveBeenNthCalledWith(3, "redis_get_stream_pending", {
connectionId: "redis-1",
db: 2,
keyRaw: "b3JkZXJz",
groupRaw: "Z3JvdXAALQ==",
cursor: "1714470000000-17",
});
});
it("passes an optional binary-safe consumer name to the pending command", async () => {
mocks.invoke.mockResolvedValue(undefined);
const { redisGetStreamPending } = await import("@/lib/backend/tauri");
await redisGetStreamPending("redis-1", 2, "b3JkZXJz", "Z3JvdXAALQ==", "1714470000000-17", "d29ya2VyLWE=");
expect(mocks.invoke).toHaveBeenCalledWith("redis_get_stream_pending", {
connectionId: "redis-1",
db: 2,
keyRaw: "b3JkZXJz",
groupRaw: "Z3JvdXAALQ==",
cursor: "1714470000000-17",
consumerRaw: "d29ya2VyLWE=",
});
});
});
describe("Redis stream monitoring HTTP API", () => {
beforeEach(() => {
vi.clearAllMocks();
vi.unstubAllGlobals();
});
function stubFetch() {
const fetchMock = vi.fn().mockResolvedValue({
ok: true,
json: vi.fn().mockResolvedValue(undefined),
});
vi.stubGlobal("fetch", fetchMock);
return fetchMock;
}
function lastCall(fetchMock: ReturnType<typeof stubFetch>): { url: string; body: Record<string, unknown> } {
const [url, init] = fetchMock.mock.calls.at(-1) as [string, RequestInit];
return { url, body: JSON.parse(String(init.body)) as Record<string, unknown> };
}
it("posts group, consumer, and pending queries to their read-only endpoints", async () => {
const fetchMock = stubFetch();
const { redisGetStreamConsumers, redisGetStreamGroups, redisGetStreamPending } = await import("@/lib/backend/http");
await redisGetStreamGroups("redis-1", 2, "b3JkZXJz");
expect(lastCall(fetchMock)).toEqual({
url: "/api/redis/get-stream-groups",
body: { connectionId: "redis-1", db: 2, keyRaw: "b3JkZXJz" },
});
await redisGetStreamConsumers("redis-1", 2, "b3JkZXJz", "Z3JvdXAALQ==");
expect(lastCall(fetchMock)).toEqual({
url: "/api/redis/get-stream-consumers",
body: { connectionId: "redis-1", db: 2, keyRaw: "b3JkZXJz", groupRaw: "Z3JvdXAALQ==" },
});
await redisGetStreamPending("redis-1", 2, "b3JkZXJz", "Z3JvdXAALQ==", "1714470000000-17");
expect(lastCall(fetchMock)).toEqual({
url: "/api/redis/get-stream-pending",
body: { connectionId: "redis-1", db: 2, keyRaw: "b3JkZXJz", groupRaw: "Z3JvdXAALQ==", cursor: "1714470000000-17" },
});
});
it("posts a consumer-scoped pending query", async () => {
const fetchMock = stubFetch();
const { redisGetStreamPending } = await import("@/lib/backend/http");
await redisGetStreamPending("redis-1", 2, "b3JkZXJz", "Z3JvdXAALQ==", "1714470000000-17", "d29ya2VyLWE=");
expect(lastCall(fetchMock)).toEqual({
url: "/api/redis/get-stream-pending",
body: {
connectionId: "redis-1",
db: 2,
keyRaw: "b3JkZXJz",
groupRaw: "Z3JvdXAALQ==",
cursor: "1714470000000-17",
consumerRaw: "d29ya2VyLWE=",
},
});
});
});
+8
View File
@@ -387,6 +387,9 @@ export const redisScanKeys = forward("redisScanKeys");
export const redisScanKeysBatch = forward("redisScanKeysBatch");
export const redisScanValues = forward("redisScanValues");
export const redisGetValue = forward("redisGetValue");
export const redisGetStreamGroups = forward("redisGetStreamGroups");
export const redisGetStreamConsumers = forward("redisGetStreamConsumers");
export const redisGetStreamPending = forward("redisGetStreamPending");
export const redisSetString = forward("redisSetString");
export const redisDeleteKey = forward("redisDeleteKey");
export const redisHashSet = forward("redisHashSet");
@@ -627,8 +630,13 @@ export type {
RedisKeyInfo,
RedisListItem,
RedisSetItem,
RedisStreamConsumer,
RedisStreamEntry,
RedisStreamField,
RedisStreamGroup,
RedisStreamMetric,
RedisStreamPendingEntry,
RedisStreamPendingPage,
RedisValue,
RedisValueData,
RedisZsetItem,
+15
View File
@@ -68,6 +68,9 @@ import type {
UpdateDownloadSource,
RedisCollectionPage,
RedisDatabaseInfo,
RedisStreamConsumer,
RedisStreamGroup,
RedisStreamPendingPage,
RedisValue,
RedisScanResult,
RedisCommandResult,
@@ -1999,6 +2002,18 @@ export async function redisGetValue(connectionId: string, db: number, keyRaw: st
return post("/api/redis/get-value", { connectionId, db, keyRaw });
}
export async function redisGetStreamGroups(connectionId: string, db: number, keyRaw: string): Promise<RedisStreamGroup[]> {
return post("/api/redis/get-stream-groups", { connectionId, db, keyRaw });
}
export async function redisGetStreamConsumers(connectionId: string, db: number, keyRaw: string, groupRaw: string): Promise<RedisStreamConsumer[]> {
return post("/api/redis/get-stream-consumers", { connectionId, db, keyRaw, groupRaw });
}
export async function redisGetStreamPending(connectionId: string, db: number, keyRaw: string, groupRaw: string, cursor?: string, consumerRaw?: string): Promise<RedisStreamPendingPage> {
return post("/api/redis/get-stream-pending", { connectionId, db, keyRaw, groupRaw, cursor, ...(consumerRaw === undefined ? {} : { consumerRaw }) });
}
export async function redisSetString(connectionId: string, db: number, keyRaw: string, value: string, ttl?: number): Promise<void> {
return post("/api/redis/set-string", { connectionId, db, keyRaw, value, ttl });
}
+43
View File
@@ -1650,6 +1650,37 @@ export interface RedisStreamEntry {
fields: RedisStreamField[];
}
// Redis counters above Number.MAX_SAFE_INTEGER are transported as decimal strings.
export type RedisStreamMetric = number | string;
export interface RedisStreamGroup {
name: RedisBlob;
consumers: RedisStreamMetric;
pending: RedisStreamMetric;
last_delivered_id: string;
entries_read?: RedisStreamMetric;
lag?: RedisStreamMetric;
}
export interface RedisStreamConsumer {
name: RedisBlob;
pending: RedisStreamMetric;
idle_ms: RedisStreamMetric;
inactive_ms?: RedisStreamMetric;
}
export interface RedisStreamPendingEntry {
id: string;
consumer: RedisBlob;
idle_ms: RedisStreamMetric;
deliveries: RedisStreamMetric;
}
export interface RedisStreamPendingPage {
entries: RedisStreamPendingEntry[];
next_cursor?: string;
}
export type RedisValueData =
| { kind: "string"; content: RedisBlob }
| { kind: "json"; value: string }
@@ -1718,6 +1749,18 @@ export async function redisGetValue(connectionId: string, db: number, keyRaw: st
return invoke("redis_get_value", { connectionId, db, keyRaw });
}
export async function redisGetStreamGroups(connectionId: string, db: number, keyRaw: string): Promise<RedisStreamGroup[]> {
return invoke("redis_get_stream_groups", { connectionId, db, keyRaw });
}
export async function redisGetStreamConsumers(connectionId: string, db: number, keyRaw: string, groupRaw: string): Promise<RedisStreamConsumer[]> {
return invoke("redis_get_stream_consumers", { connectionId, db, keyRaw, groupRaw });
}
export async function redisGetStreamPending(connectionId: string, db: number, keyRaw: string, groupRaw: string, cursor?: string, consumerRaw?: string): Promise<RedisStreamPendingPage> {
return invoke("redis_get_stream_pending", { connectionId, db, keyRaw, groupRaw, cursor, ...(consumerRaw === undefined ? {} : { consumerRaw }) });
}
export async function redisSetString(connectionId: string, db: number, keyRaw: string, value: string, ttl?: number): Promise<void> {
return invoke("redis_set_string", { connectionId, db, keyRaw, value, ttl });
}
+463 -8
View File
@@ -15,9 +15,11 @@ use tokio::sync::{Mutex, MutexGuard};
use super::json_value_for_js;
const STREAM_ENTRY_LIMIT: usize = 100;
const STREAM_PENDING_PAGE_SIZE: usize = 100;
const COLLECTION_PAGE_SIZE: usize = 200;
const HASH_FILTER_SCAN_MAX_ITERATIONS: usize = 10;
const DEFAULT_REDIS_DATABASES: u32 = 16;
const JS_MAX_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
const CLUSTER_CURSOR_NODE_BITS: u64 = 16;
const CLUSTER_CURSOR_NODE_MASK: u64 = (1 << CLUSTER_CURSOR_NODE_BITS) - 1;
const CLUSTER_CURSOR_SCAN_MASK: u64 = (1 << (64 - CLUSTER_CURSOR_NODE_BITS)) - 1;
@@ -122,6 +124,70 @@ pub struct RedisStreamEntry {
pub fields: Vec<RedisStreamField>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RedisStreamGroup {
pub name: RedisBlob,
#[serde(serialize_with = "serialize_redis_u64_for_js")]
pub consumers: u64,
#[serde(serialize_with = "serialize_redis_u64_for_js")]
pub pending: u64,
pub last_delivered_id: String,
#[serde(default, skip_serializing_if = "Option::is_none", serialize_with = "serialize_optional_redis_u64_for_js")]
pub entries_read: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none", serialize_with = "serialize_optional_redis_u64_for_js")]
pub lag: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RedisStreamConsumer {
pub name: RedisBlob,
#[serde(serialize_with = "serialize_redis_u64_for_js")]
pub pending: u64,
#[serde(serialize_with = "serialize_redis_u64_for_js")]
pub idle_ms: u64,
#[serde(default, skip_serializing_if = "Option::is_none", serialize_with = "serialize_optional_redis_u64_for_js")]
pub inactive_ms: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RedisStreamPendingEntry {
pub id: String,
pub consumer: RedisBlob,
#[serde(serialize_with = "serialize_redis_u64_for_js")]
pub idle_ms: u64,
#[serde(serialize_with = "serialize_redis_u64_for_js")]
pub deliveries: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RedisStreamPendingPage {
pub entries: Vec<RedisStreamPendingEntry>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_cursor: Option<String>,
}
// JSON and Tauri IPC both deliver numbers to JavaScript, so preserve large Redis counters as text.
fn serialize_redis_u64_for_js<S>(value: &u64, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
if *value > JS_MAX_SAFE_INTEGER {
serializer.serialize_str(&value.to_string())
} else {
serializer.serialize_u64(*value)
}
}
fn serialize_optional_redis_u64_for_js<S>(value: &Option<u64>, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
match value {
Some(value) => serialize_redis_u64_for_js(value, serializer),
None => serializer.serialize_none(),
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum RedisValueData {
@@ -1888,8 +1954,9 @@ fn redis_raw_value_to_command_arg(v: &RedisRawValue) -> Option<String> {
/// Try to convert a RedisRawValue to a u64.
fn redis_value_to_u64(v: &RedisRawValue) -> Option<u64> {
match v {
RedisRawValue::Int(i) => Some(*i as u64),
RedisRawValue::Int(i) => u64::try_from(*i).ok(),
RedisRawValue::BulkString(bytes) => std::str::from_utf8(bytes).ok().and_then(|s| s.parse().ok()),
RedisRawValue::SimpleString(value) => value.parse().ok(),
_ => None,
}
}
@@ -2358,6 +2425,64 @@ where
Ok(parse_stream_entries(raw))
}
pub async fn get_stream_groups<C>(con: &mut C, key: &[u8]) -> Result<Vec<RedisStreamGroup>, String>
where
C: ConnectionLike + Send + Sync + Unpin,
{
let raw: RedisRawValue =
redis::cmd("XINFO").arg("GROUPS").arg(key).query_async(con).await.map_err(|e| e.to_string())?;
Ok(parse_stream_groups(raw))
}
pub async fn get_stream_consumers<C>(con: &mut C, key: &[u8], group: &[u8]) -> Result<Vec<RedisStreamConsumer>, String>
where
C: ConnectionLike + Send + Sync + Unpin,
{
let raw: RedisRawValue =
redis::cmd("XINFO").arg("CONSUMERS").arg(key).arg(group).query_async(con).await.map_err(|e| e.to_string())?;
Ok(parse_stream_consumers(raw))
}
pub async fn get_stream_pending_page<C>(
con: &mut C,
key: &[u8],
group: &[u8],
cursor: Option<&str>,
consumer: Option<&[u8]>,
) -> Result<RedisStreamPendingPage, String>
where
C: ConnectionLike + Send + Sync + Unpin,
{
let cursor = cursor.filter(|value| !value.is_empty());
let start = cursor.unwrap_or("-");
// Redis before 6.2 has no exclusive XPENDING ranges. A cursor page needs
// a duplicate candidate and a lookahead row to page without skipping IDs.
let requested_count = STREAM_PENDING_PAGE_SIZE + 1 + usize::from(cursor.is_some());
let mut command = redis::cmd("XPENDING");
command.arg(key).arg(group).arg(start).arg("+").arg(requested_count);
if let Some(consumer) = consumer {
command.arg(consumer);
}
let raw: RedisRawValue = command.query_async(con).await.map_err(|e| e.to_string())?;
let mut entries = parse_stream_pending_entries(raw);
if let Some(cursor) = cursor {
if entries.first().is_some_and(|entry| entry.id == cursor) {
entries.remove(0);
}
}
let next_cursor = if entries.len() > STREAM_PENDING_PAGE_SIZE {
entries.truncate(STREAM_PENDING_PAGE_SIZE);
entries.last().map(|entry| entry.id.clone())
} else {
None
};
Ok(RedisStreamPendingPage { entries, next_cursor })
}
fn parse_scan_keys(raw: RedisRawValue) -> Result<(u64, Vec<Vec<u8>>), String> {
let RedisRawValue::Array(parts) = raw else {
return Err("Invalid Redis SCAN response".to_string());
@@ -2417,6 +2542,81 @@ fn parse_stream_entry(entry: RedisRawValue) -> Option<RedisStreamEntry> {
Some(RedisStreamEntry { id, fields: parsed_fields })
}
fn parse_stream_groups(raw: RedisRawValue) -> Vec<RedisStreamGroup> {
match raw {
RedisRawValue::Array(groups) => groups.into_iter().filter_map(parse_stream_group).collect(),
_ => Vec::new(),
}
}
fn parse_stream_group(group: RedisRawValue) -> Option<RedisStreamGroup> {
let mut attributes = parse_stream_info_attributes(group)?;
Some(RedisStreamGroup {
name: redis_value_to_blob(attributes.remove("name")?)?,
consumers: redis_value_to_u64(&attributes.remove("consumers")?)?,
pending: redis_value_to_u64(&attributes.remove("pending")?)?,
last_delivered_id: redis_value_to_string(attributes.remove("last-delivered-id")?)?,
entries_read: attributes.remove("entries-read").and_then(|value| redis_value_to_u64(&value)),
lag: attributes.remove("lag").and_then(|value| redis_value_to_u64(&value)),
})
}
fn parse_stream_consumers(raw: RedisRawValue) -> Vec<RedisStreamConsumer> {
match raw {
RedisRawValue::Array(consumers) => consumers.into_iter().filter_map(parse_stream_consumer).collect(),
_ => Vec::new(),
}
}
fn parse_stream_consumer(consumer: RedisRawValue) -> Option<RedisStreamConsumer> {
let mut attributes = parse_stream_info_attributes(consumer)?;
Some(RedisStreamConsumer {
name: redis_value_to_blob(attributes.remove("name")?)?,
pending: redis_value_to_u64(&attributes.remove("pending")?)?,
idle_ms: redis_value_to_u64(&attributes.remove("idle")?)?,
inactive_ms: attributes.remove("inactive").and_then(|value| redis_value_to_u64(&value)),
})
}
fn parse_stream_pending_entries(raw: RedisRawValue) -> Vec<RedisStreamPendingEntry> {
match raw {
RedisRawValue::Array(entries) => entries.into_iter().filter_map(parse_stream_pending_entry).collect(),
_ => Vec::new(),
}
}
fn parse_stream_pending_entry(entry: RedisRawValue) -> Option<RedisStreamPendingEntry> {
let mut parts = match entry {
RedisRawValue::Array(parts) if parts.len() == 4 => parts.into_iter(),
_ => return None,
};
Some(RedisStreamPendingEntry {
id: redis_value_to_string(parts.next()?)?,
consumer: redis_value_to_blob(parts.next()?)?,
idle_ms: redis_value_to_u64(&parts.next()?)?,
deliveries: redis_value_to_u64(&parts.next()?)?,
})
}
fn parse_stream_info_attributes(raw: RedisRawValue) -> Option<HashMap<String, RedisRawValue>> {
let RedisRawValue::Array(parts) = raw else {
return None;
};
let mut attributes = HashMap::new();
let mut parts = parts.into_iter();
while let Some(name) = parts.next() {
let value = parts.next()?;
attributes.insert(redis_value_to_string(name)?, value);
}
Some(attributes)
}
fn redis_value_to_blob(value: RedisRawValue) -> Option<RedisBlob> {
redis_value_to_bytes(value).map(|bytes| redis_blob_from_bytes(&bytes))
}
fn redis_value_to_string(value: RedisRawValue) -> Option<String> {
match value {
RedisRawValue::BulkString(bytes) => Some(redis_bytes_to_display(&bytes)),
@@ -2975,13 +3175,14 @@ mod tests {
use super::{
classify_command, connection_info, decode_cluster_cursor, encode_cluster_cursor, is_redis_json_type,
parse_cluster_slots, parse_command_argv, parse_database_count, parse_redis_endpoint, parse_scan_keys,
parse_stream_entries, redis_auth_candidates, redis_blob_from_bytes, redis_cluster_slot,
redis_command_raw_to_json, redis_database_index, redis_key_bytes_to_display, redis_key_bytes_to_raw,
redis_key_matches_query, redis_key_raw_to_bytes, redis_key_value_preview, redis_sentinel_master_endpoint,
redis_value_matches_query, redis_value_to_bytes, standalone_connection_infos, RedisAuthCandidate, RedisBlob,
RedisBlobEncoding, RedisClusterSlotRange, RedisCollectionPage, RedisCommandSafety, RedisHashItem,
RedisNodeEndpoint, RedisNodeRoute, RedisRawValue, RedisSetItem, RedisStreamEntry, RedisStreamField, RedisValue,
RedisValueData,
parse_stream_consumers, parse_stream_entries, parse_stream_groups, parse_stream_pending_entries,
redis_auth_candidates, redis_blob_from_bytes, redis_cluster_slot, redis_command_raw_to_json,
redis_database_index, redis_key_bytes_to_display, redis_key_bytes_to_raw, redis_key_matches_query,
redis_key_raw_to_bytes, redis_key_value_preview, redis_sentinel_master_endpoint, redis_value_matches_query,
redis_value_to_bytes, standalone_connection_infos, RedisAuthCandidate, RedisBlob, RedisBlobEncoding,
RedisClusterSlotRange, RedisCollectionPage, RedisCommandSafety, RedisHashItem, RedisNodeEndpoint,
RedisNodeRoute, RedisRawValue, RedisSetItem, RedisStreamConsumer, RedisStreamEntry, RedisStreamField,
RedisStreamGroup, RedisStreamPendingEntry, RedisValue, RedisValueData,
};
use crate::models::connection::ConnectionConfig;
use redis::{aio::ConnectionLike, Cmd, ConnectionAddr, Pipeline, RedisFuture};
@@ -3167,6 +3368,260 @@ mod tests {
);
}
#[test]
fn parses_stream_groups_with_binary_names_and_compatible_optional_fields() {
let raw = RedisRawValue::Array(vec![
RedisRawValue::Array(vec![
bulk("name"),
RedisRawValue::BulkString(vec![0xFF, b'g']),
bulk("consumers"),
RedisRawValue::Int(2),
bulk("pending"),
bulk("3"),
bulk("last-delivered-id"),
bulk("1714470000000-0"),
bulk("entries-read"),
RedisRawValue::Int(11),
bulk("lag"),
RedisRawValue::Nil,
bulk("future-field"),
bulk("ignored"),
]),
RedisRawValue::Array(vec![
bulk("name"),
bulk("legacy-group"),
bulk("consumers"),
RedisRawValue::Int(0),
bulk("pending"),
RedisRawValue::Int(0),
bulk("last-delivered-id"),
bulk("0-0"),
]),
]);
let parsed = parse_stream_groups(raw);
assert_eq!(parsed.len(), 2);
assert_eq!(parsed[0].name, redis_blob_from_bytes(&[0xFF, b'g']));
assert_eq!(
parsed[0],
RedisStreamGroup {
name: redis_blob_from_bytes(&[0xFF, b'g']),
consumers: 2,
pending: 3,
last_delivered_id: "1714470000000-0".to_string(),
entries_read: Some(11),
lag: None,
}
);
assert_eq!(
parsed[1],
RedisStreamGroup {
name: text_blob("legacy-group"),
consumers: 0,
pending: 0,
last_delivered_id: "0-0".to_string(),
entries_read: None,
lag: None,
}
);
}
#[test]
fn parses_stream_consumers_and_pending_entries() {
let consumers = parse_stream_consumers(RedisRawValue::Array(vec![
RedisRawValue::Array(vec![
bulk("name"),
bulk("worker-a"),
bulk("pending"),
RedisRawValue::Int(4),
bulk("idle"),
bulk("1200"),
bulk("inactive"),
RedisRawValue::Int(800),
]),
RedisRawValue::Array(vec![
bulk("name"),
bulk("legacy-worker"),
bulk("pending"),
RedisRawValue::Int(0),
bulk("idle"),
RedisRawValue::Int(0),
]),
]));
let pending = parse_stream_pending_entries(RedisRawValue::Array(vec![
RedisRawValue::Array(vec![
bulk("1714470000000-0"),
RedisRawValue::BulkString(vec![0xFE, b'c']),
RedisRawValue::Int(2_400),
bulk("3"),
]),
RedisRawValue::Array(vec![bulk("malformed")]),
]));
assert_eq!(
consumers,
vec![
RedisStreamConsumer { name: text_blob("worker-a"), pending: 4, idle_ms: 1_200, inactive_ms: Some(800) },
RedisStreamConsumer { name: text_blob("legacy-worker"), pending: 0, idle_ms: 0, inactive_ms: None },
]
);
assert_eq!(
pending,
vec![RedisStreamPendingEntry {
id: "1714470000000-0".to_string(),
consumer: redis_blob_from_bytes(&[0xFE, b'c']),
idle_ms: 2_400,
deliveries: 3,
}]
);
}
#[test]
fn serializes_stream_metrics_without_losing_javascript_precision() {
let unsafe_value = super::JS_MAX_SAFE_INTEGER + 1;
let group = RedisStreamGroup {
name: text_blob("payments"),
consumers: unsafe_value,
pending: 2,
last_delivered_id: "1714470000000-0".to_string(),
entries_read: Some(unsafe_value),
lag: None,
};
let consumer = RedisStreamConsumer {
name: text_blob("worker-a"),
pending: 2,
idle_ms: unsafe_value,
inactive_ms: Some(unsafe_value),
};
let entry = RedisStreamPendingEntry {
id: "1714470000000-0".to_string(),
consumer: text_blob("worker-a"),
idle_ms: 2,
deliveries: unsafe_value,
};
let group_json = serde_json::to_value(group).unwrap();
let consumer_json = serde_json::to_value(consumer).unwrap();
let entry_json = serde_json::to_value(entry).unwrap();
assert_eq!(group_json["consumers"], unsafe_value.to_string());
assert_eq!(group_json["pending"].as_u64(), Some(2));
assert_eq!(group_json["entries_read"], unsafe_value.to_string());
assert!(group_json.get("lag").is_none());
assert_eq!(consumer_json["idle_ms"], unsafe_value.to_string());
assert_eq!(consumer_json["inactive_ms"], unsafe_value.to_string());
assert_eq!(entry_json["deliveries"], unsafe_value.to_string());
}
#[tokio::test]
async fn stream_monitoring_commands_are_read_only_and_pending_pagination_supports_redis_5() {
let groups = RedisRawValue::Array(vec![RedisRawValue::Array(vec![
bulk("name"),
bulk("payments"),
bulk("consumers"),
RedisRawValue::Int(1),
bulk("pending"),
RedisRawValue::Int(101),
bulk("last-delivered-id"),
bulk("1714470000000-0"),
])]);
let consumers = RedisRawValue::Array(vec![RedisRawValue::Array(vec![
bulk("name"),
bulk("worker-a"),
bulk("pending"),
RedisRawValue::Int(101),
bulk("idle"),
RedisRawValue::Int(10),
])]);
let pending = RedisRawValue::Array(
(17..=118)
.map(|index| {
RedisRawValue::Array(vec![
bulk(&format!("1714470000000-{index}")),
bulk("worker-a"),
RedisRawValue::Int(index),
RedisRawValue::Int(1),
])
})
.collect(),
);
let mut con = FakeRedisConnection::new(vec![groups, consumers, pending]);
let groups = super::get_stream_groups(&mut con, b"orders").await.unwrap();
let consumers = super::get_stream_consumers(&mut con, b"orders", b"payments").await.unwrap();
let page = super::get_stream_pending_page(&mut con, b"orders", b"payments", Some("1714470000000-17"), None)
.await
.unwrap();
assert_eq!(groups.len(), 1);
assert_eq!(consumers.len(), 1);
assert_eq!(page.entries.len(), 100);
assert_eq!(page.entries.first().map(|entry| entry.id.as_str()), Some("1714470000000-18"));
assert_eq!(page.entries.last().map(|entry| entry.id.as_str()), Some("1714470000000-117"));
assert_eq!(page.next_cursor.as_deref(), Some("1714470000000-117"));
assert_eq!(con.command_count("XINFO"), 2);
assert_eq!(con.command_count("XPENDING"), 1);
assert_eq!(con.command_count("XGROUP"), 0);
assert_eq!(con.command_count("XACK"), 0);
assert_eq!(con.command_count("XCLAIM"), 0);
assert!(con.commands[0].contains("\r\nGROUPS\r\n"));
assert!(con.commands[1].contains("\r\nCONSUMERS\r\n"));
assert!(con.commands[2].contains("\r\n1714470000000-17\r\n"));
assert!(!con.commands[2].contains("\r\n(1714470000000-17\r\n"));
assert!(con.commands[2].contains("\r\n102\r\n"));
}
#[tokio::test]
async fn stream_pending_pagination_keeps_the_first_entry_when_cursor_is_acknowledged() {
let pending = RedisRawValue::Array(
(18..=118)
.map(|index| {
RedisRawValue::Array(vec![
bulk(&format!("1714470000000-{index}")),
bulk("worker-a"),
RedisRawValue::Int(index),
RedisRawValue::Int(1),
])
})
.collect(),
);
let mut con = FakeRedisConnection::new(vec![pending]);
let page = super::get_stream_pending_page(&mut con, b"orders", b"payments", Some("1714470000000-17"), None)
.await
.unwrap();
assert_eq!(page.entries.len(), 100);
assert_eq!(page.entries.first().map(|entry| entry.id.as_str()), Some("1714470000000-18"));
assert_eq!(page.entries.last().map(|entry| entry.id.as_str()), Some("1714470000000-117"));
assert_eq!(page.next_cursor.as_deref(), Some("1714470000000-117"));
assert!(con.commands[0].contains("\r\n1714470000000-17\r\n"));
assert!(!con.commands[0].contains("\r\n(1714470000000-17\r\n"));
assert!(con.commands[0].contains("\r\n102\r\n"));
}
#[tokio::test]
async fn stream_pending_filters_by_consumer_server_side() {
let pending = RedisRawValue::Array(vec![RedisRawValue::Array(vec![
bulk("1714470000000-0"),
bulk("worker-a"),
RedisRawValue::Int(2_400),
RedisRawValue::Int(1),
])]);
let mut con = FakeRedisConnection::new(vec![pending]);
let page =
super::get_stream_pending_page(&mut con, b"orders", b"payments", None, Some(b"worker-a")).await.unwrap();
assert_eq!(page.entries.len(), 1);
assert_eq!(con.command_count("XPENDING"), 1);
assert!(con.commands[0].contains("\r\n-\r\n"));
assert!(con.commands[0].contains("\r\n+\r\n"));
assert!(con.commands[0].contains("\r\n101\r\n"));
assert!(con.commands[0].contains("\r\nworker-a\r\n"));
}
#[test]
fn parses_configured_database_count() {
let value = RedisRawValue::Array(vec![
+96 -1
View File
@@ -1,6 +1,7 @@
use crate::connection::{AppState, PoolKind};
use crate::db::redis_driver::{
self, RedisCollectionPage, RedisCommandResult, RedisConnection, RedisDatabaseInfo, RedisScanResult, RedisValue,
self, RedisCollectionPage, RedisCommandResult, RedisConnection, RedisDatabaseInfo, RedisScanResult,
RedisStreamConsumer, RedisStreamGroup, RedisStreamPendingPage, RedisValue,
};
async fn ensure_redis_pool(state: &AppState, connection_id: &str) -> Result<(), String> {
@@ -136,6 +137,100 @@ pub async fn redis_get_value_in_db_core(
}
}
pub async fn redis_stream_groups_in_db_core(
state: &AppState,
connection_id: &str,
db: u32,
key_raw: &str,
) -> Result<Vec<RedisStreamGroup>, String> {
ensure_redis_pool(state, connection_id).await?;
let connections = state.connections.read().await;
let pool = connections.get(connection_id).ok_or("Connection not found")?;
match pool {
PoolKind::Redis(redis) => {
let key = redis_driver::redis_key_raw_to_bytes(key_raw)?;
match redis {
RedisConnection::Direct(con) => {
let mut con = con.lock().await;
redis_driver::select_db(&mut *con, db).await?;
redis_driver::get_stream_groups(&mut *con, &key).await
}
RedisConnection::Cluster(cluster) => {
redis_driver::ensure_cluster_db(db)?;
let mut con = redis_driver::cluster_key_connection(cluster, &key).await?;
redis_driver::get_stream_groups(&mut con, &key).await
}
}
}
_ => Err("Not a Redis connection".to_string()),
}
}
pub async fn redis_stream_consumers_in_db_core(
state: &AppState,
connection_id: &str,
db: u32,
key_raw: &str,
group_raw: &str,
) -> Result<Vec<RedisStreamConsumer>, String> {
ensure_redis_pool(state, connection_id).await?;
let connections = state.connections.read().await;
let pool = connections.get(connection_id).ok_or("Connection not found")?;
match pool {
PoolKind::Redis(redis) => {
let key = redis_driver::redis_key_raw_to_bytes(key_raw)?;
let group = redis_driver::redis_key_raw_to_bytes(group_raw)?;
match redis {
RedisConnection::Direct(con) => {
let mut con = con.lock().await;
redis_driver::select_db(&mut *con, db).await?;
redis_driver::get_stream_consumers(&mut *con, &key, &group).await
}
RedisConnection::Cluster(cluster) => {
redis_driver::ensure_cluster_db(db)?;
let mut con = redis_driver::cluster_key_connection(cluster, &key).await?;
redis_driver::get_stream_consumers(&mut con, &key, &group).await
}
}
}
_ => Err("Not a Redis connection".to_string()),
}
}
pub async fn redis_stream_pending_in_db_core(
state: &AppState,
connection_id: &str,
db: u32,
key_raw: &str,
group_raw: &str,
cursor: Option<&str>,
consumer_raw: Option<&str>,
) -> Result<RedisStreamPendingPage, String> {
ensure_redis_pool(state, connection_id).await?;
let connections = state.connections.read().await;
let pool = connections.get(connection_id).ok_or("Connection not found")?;
match pool {
PoolKind::Redis(redis) => {
let key = redis_driver::redis_key_raw_to_bytes(key_raw)?;
let group = redis_driver::redis_key_raw_to_bytes(group_raw)?;
let consumer = consumer_raw.map(redis_driver::redis_key_raw_to_bytes).transpose()?;
match redis {
RedisConnection::Direct(con) => {
let mut con = con.lock().await;
redis_driver::select_db(&mut *con, db).await?;
redis_driver::get_stream_pending_page(&mut *con, &key, &group, cursor, consumer.as_deref()).await
}
RedisConnection::Cluster(cluster) => {
redis_driver::ensure_cluster_db(db)?;
let mut con = redis_driver::cluster_key_connection(cluster, &key).await?;
redis_driver::get_stream_pending_page(&mut con, &key, &group, cursor, consumer.as_deref()).await
}
}
}
_ => Err("Not a Redis connection".to_string()),
}
}
pub async fn redis_set_string_core(
state: &AppState,
connection_id: &str,
+3
View File
@@ -466,6 +466,9 @@ async fn main() {
.route("/redis/scan-keys-batch", post(routes::redis::scan_keys_batch))
.route("/redis/scan-values", post(routes::redis::scan_values))
.route("/redis/get-value", post(routes::redis::get_value))
.route("/redis/get-stream-groups", post(routes::redis::get_stream_groups))
.route("/redis/get-stream-consumers", post(routes::redis::get_stream_consumers))
.route("/redis/get-stream-pending", post(routes::redis::get_stream_pending))
.route("/redis/load-more", post(routes::redis::load_more))
.route("/redis/set-string", post(routes::redis::set_string))
.route("/redis/delete-key", post(routes::redis::delete_key))
+65
View File
@@ -76,6 +76,26 @@ pub struct RedisKeyRequest {
pub key_raw: String,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct RedisStreamGroupRequest {
pub connection_id: String,
pub db: u32,
pub key_raw: String,
pub group_raw: String,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct RedisStreamPendingRequest {
pub connection_id: String,
pub db: u32,
pub key_raw: String,
pub group_raw: String,
pub cursor: Option<String>,
pub consumer_raw: Option<String>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct RedisLoadMoreRequest {
@@ -302,6 +322,51 @@ pub async fn get_value(
Ok(Json(serde_json::to_value(result).map_err(|e| AppError::from(e.to_string()))?))
}
pub async fn get_stream_groups(
State(state): State<Arc<WebState>>,
Json(req): Json<RedisKeyRequest>,
) -> Result<Json<serde_json::Value>, AppError> {
let result =
dbx_core::redis_ops::redis_stream_groups_in_db_core(&state.app, &req.connection_id, req.db, &req.key_raw)
.await
.map_err(AppError::from)?;
Ok(Json(serde_json::to_value(result).map_err(|e| AppError::from(e.to_string()))?))
}
pub async fn get_stream_consumers(
State(state): State<Arc<WebState>>,
Json(req): Json<RedisStreamGroupRequest>,
) -> Result<Json<serde_json::Value>, AppError> {
let result = dbx_core::redis_ops::redis_stream_consumers_in_db_core(
&state.app,
&req.connection_id,
req.db,
&req.key_raw,
&req.group_raw,
)
.await
.map_err(AppError::from)?;
Ok(Json(serde_json::to_value(result).map_err(|e| AppError::from(e.to_string()))?))
}
pub async fn get_stream_pending(
State(state): State<Arc<WebState>>,
Json(req): Json<RedisStreamPendingRequest>,
) -> Result<Json<serde_json::Value>, AppError> {
let result = dbx_core::redis_ops::redis_stream_pending_in_db_core(
&state.app,
&req.connection_id,
req.db,
&req.key_raw,
&req.group_raw,
req.cursor.as_deref(),
req.consumer_raw.as_deref(),
)
.await
.map_err(AppError::from)?;
Ok(Json(serde_json::to_value(result).map_err(|e| AppError::from(e.to_string()))?))
}
pub async fn load_more(
State(state): State<Arc<WebState>>,
Json(req): Json<RedisLoadMoreRequest>,
+44 -1
View File
@@ -4,7 +4,7 @@ use tauri::State;
use crate::commands::connection::{ensure_connection_writable, AppState};
use dbx_core::db::redis_driver::{
classify_command, parse_command_argv, RedisCollectionPage, RedisCommandResult, RedisCommandSafety,
RedisDatabaseInfo, RedisScanResult, RedisValue,
RedisDatabaseInfo, RedisScanResult, RedisStreamConsumer, RedisStreamGroup, RedisStreamPendingPage, RedisValue,
};
#[tauri::command]
@@ -86,6 +86,49 @@ pub async fn redis_get_value(
dbx_core::redis_ops::redis_get_value_in_db_core(&state, &connection_id, db, &key_raw).await
}
#[tauri::command]
pub async fn redis_get_stream_groups(
state: State<'_, Arc<AppState>>,
connection_id: String,
db: u32,
key_raw: String,
) -> Result<Vec<RedisStreamGroup>, String> {
dbx_core::redis_ops::redis_stream_groups_in_db_core(&state, &connection_id, db, &key_raw).await
}
#[tauri::command]
pub async fn redis_get_stream_consumers(
state: State<'_, Arc<AppState>>,
connection_id: String,
db: u32,
key_raw: String,
group_raw: String,
) -> Result<Vec<RedisStreamConsumer>, String> {
dbx_core::redis_ops::redis_stream_consumers_in_db_core(&state, &connection_id, db, &key_raw, &group_raw).await
}
#[tauri::command]
pub async fn redis_get_stream_pending(
state: State<'_, Arc<AppState>>,
connection_id: String,
db: u32,
key_raw: String,
group_raw: String,
cursor: Option<String>,
consumer_raw: Option<String>,
) -> Result<RedisStreamPendingPage, String> {
dbx_core::redis_ops::redis_stream_pending_in_db_core(
&state,
&connection_id,
db,
&key_raw,
&group_raw,
cursor.as_deref(),
consumer_raw.as_deref(),
)
.await
}
#[tauri::command]
pub async fn redis_set_string(
state: State<'_, Arc<AppState>>,
+3
View File
@@ -1524,6 +1524,9 @@ pub fn run() {
commands::redis_cmd::redis_scan_keys_batch,
commands::redis_cmd::redis_scan_values,
commands::redis_cmd::redis_get_value,
commands::redis_cmd::redis_get_stream_groups,
commands::redis_cmd::redis_get_stream_consumers,
commands::redis_cmd::redis_get_stream_pending,
commands::redis_cmd::redis_set_string,
commands::redis_cmd::redis_delete_key,
commands::redis_cmd::redis_hash_set,