From faa63a68fa3c1c879919a8f058a00de7f98e04a7 Mon Sep 17 00:00:00 2001 From: George Madi Date: Mon, 20 Jul 2026 12:36:08 -0500 Subject: [PATCH] Add option to allow limiting key listings. Ref: RIAK-3044 --- .../riak/client/cli/CustomTLSConnection.java | 43 +++++++++++++------ .../riak/client/cli/InteractiveRiakCLI.java | 21 ++++++--- .../client/cli/RiakOperationsService.java | 13 ++++-- 3 files changed, 55 insertions(+), 22 deletions(-) diff --git a/src/main/java/org/openriak/riak/client/cli/CustomTLSConnection.java b/src/main/java/org/openriak/riak/client/cli/CustomTLSConnection.java index b76bc4b..0aeca2c 100644 --- a/src/main/java/org/openriak/riak/client/cli/CustomTLSConnection.java +++ b/src/main/java/org/openriak/riak/client/cli/CustomTLSConnection.java @@ -460,11 +460,11 @@ public void executeCustomListBuckets() throws Exception { } } - public void executeCustomListKeys(String bucket) throws Exception { + public void executeCustomListKeys(String bucket, int keysLimit) throws Exception { System.out.println("Listing keys in bucket '" + bucket + "'..."); // Create LIST_KEYS request message - byte[] listKeysMessage = createListKeysMessage(bucket); + byte[] listKeysMessage = createListKeysMessage(bucket, keysLimit); // Send message via custom TLS connection byte[] response = sendMessage(listKeysMessage); @@ -563,24 +563,29 @@ private List parseListBucketsResponse(byte[] response) { /** * Create a LIST_KEYS request message */ - private byte[] createListKeysMessage(String bucket) { + private byte[] createListKeysMessage(String bucket, int keysLimit) { // Riak protocol: 4-byte length + 1-byte message code + protobuf payload byte messageCode = (byte)17; // MSG_ListKeysReq = 17 try { - // Create protobuf message (simplified) byte[] bucketBytes = bucket.getBytes("UTF-8"); - // Simple protobuf-like encoding - byte[] payload = new byte[bucketBytes.length + 2]; - int pos = 0; + java.io.ByteArrayOutputStream payload = new java.io.ByteArrayOutputStream(); - // Write bucket (field 1, wire type 2) - payload[pos++] = (byte) (1 << 3 | 2); - payload[pos++] = (byte) bucketBytes.length; - System.arraycopy(bucketBytes, 0, payload, pos, bucketBytes.length); + // Write bucket (field 1, wire type 2 = length-delimited) + payload.write((1 << 3) | 2); + writeVarint(payload, bucketBytes.length); + payload.write(bucketBytes); - int totalLength = 1 + payload.length; // message code + payload + // Write keys_limit (field 4, wire type 0 = varint) when limited. + // 0 or negative means unlimited, so the field is omitted. + if (keysLimit > 0) { + payload.write((4 << 3) | 0); + writeVarint(payload, keysLimit); + } + + byte[] payloadBytes = payload.toByteArray(); + int totalLength = 1 + payloadBytes.length; // message code + payload byte[] message = new byte[4 + totalLength]; // Write length (big-endian) @@ -593,7 +598,7 @@ private byte[] createListKeysMessage(String bucket) { message[4] = messageCode; // Write payload - System.arraycopy(payload, 0, message, 5, payload.length); + System.arraycopy(payloadBytes, 0, message, 5, payloadBytes.length); return message; } catch (Exception e) { @@ -601,6 +606,18 @@ private byte[] createListKeysMessage(String bucket) { } } + /** + * Write an unsigned varint (protobuf base-128) to the given stream. + */ + private void writeVarint(java.io.ByteArrayOutputStream out, int value) { + int v = value; + while ((v & ~0x7F) != 0) { + out.write((v & 0x7F) | 0x80); + v >>>= 7; + } + out.write(v & 0x7F); + } + /** * Parse LIST_KEYS response to extract key names */ diff --git a/src/main/java/org/openriak/riak/client/cli/InteractiveRiakCLI.java b/src/main/java/org/openriak/riak/client/cli/InteractiveRiakCLI.java index d1a2c5e..cbca94f 100644 --- a/src/main/java/org/openriak/riak/client/cli/InteractiveRiakCLI.java +++ b/src/main/java/org/openriak/riak/client/cli/InteractiveRiakCLI.java @@ -696,9 +696,18 @@ private void runInteractive() { break; case "list-keys": if (parts.length >= 2) { - executeListKeys(parts[1]); + int keysLimit = 0; + if (parts.length >= 3) { + try { + keysLimit = Integer.parseInt(parts[2]); + } catch (NumberFormatException e) { + System.out.println("Usage: list-keys [limit]"); + break; + } + } + executeListKeys(parts[1], keysLimit); } else { - System.out.println("Usage: list-keys "); + System.out.println("Usage: list-keys [limit]"); } break; default: @@ -720,7 +729,7 @@ private void printInteractiveHelp() { System.out.println(" put - Store a value in Riak"); System.out.println(" delete - Delete a value from Riak"); System.out.println(" list-buckets - List all buckets"); - System.out.println(" list-keys - List all keys in a bucket"); + System.out.println(" list-keys [limit] - List keys in a bucket (optional max count)"); System.out.println(" help - Show this help message"); System.out.println(" quit - Exit the CLI"); } @@ -862,7 +871,7 @@ private void executeListBuckets() throws Exception { } } - private void executeListKeys(String bucket) throws Exception { + private void executeListKeys(String bucket, int keysLimit) throws Exception { if (!connected) { System.out.println("Not connected to Riak"); return; @@ -870,13 +879,13 @@ private void executeListKeys(String bucket) throws Exception { if (useTLS && USE_CUSTOM_TLS) { // Use custom TLS implementation - customTLS.executeCustomListKeys(bucket); + customTLS.executeCustomListKeys(bucket, keysLimit); } else { if (operationsService == null) { throw new IllegalStateException("Riak operations service is not initialized"); } System.out.println("Listing keys in bucket '" + bucket + "'..."); - RiakOperationResult result = operationsService.listKeys(bucket); + RiakOperationResult result = operationsService.listKeys(bucket, keysLimit); @SuppressWarnings("unchecked") List keys = (List) result.getData().getOrDefault("keys", List.of()); System.out.println("Keys:"); diff --git a/src/main/java/org/openriak/riak/client/cli/RiakOperationsService.java b/src/main/java/org/openriak/riak/client/cli/RiakOperationsService.java index e71410c..bdea998 100644 --- a/src/main/java/org/openriak/riak/client/cli/RiakOperationsService.java +++ b/src/main/java/org/openriak/riak/client/cli/RiakOperationsService.java @@ -209,10 +209,17 @@ public RiakOperationResult listBuckets() throws Exception { } public RiakOperationResult listKeys(String bucket) throws Exception { + return listKeys(bucket, 0); + } + + public RiakOperationResult listKeys(String bucket, int keysLimit) throws Exception { Namespace ns = new Namespace(defaultBucketType, bucket); - ListKeys.Builder listKeysBuilder = new ListKeys.Builder(ns).withTimeout(5000); - boolean allowListingConfigured = invokeNoArgIfPresent(listKeysBuilder, "withAllowListing"); - ListKeys lk = listKeysBuilder.build(); + ListKeys.Builder lkBuilder = new ListKeys.Builder(ns).withTimeout(5000); + boolean allowListingConfigured = invokeNoArgIfPresent(lkBuilder, "withAllowListing"); + if (keysLimit > 0) { + lkBuilder.withKeysLimit(keysLimit); + } + ListKeys lk = lkBuilder.build(); ListKeys.Response response; try { response = client.execute(lk);