Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 30 additions & 13 deletions src/main/java/org/openriak/riak/client/cli/CustomTLSConnection.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down Expand Up @@ -563,24 +563,29 @@ private List<String> 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)
Expand All @@ -593,14 +598,26 @@ 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) {
throw new RuntimeException("Failed to create list keys message: " + e.getMessage());
}
}

/**
* 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
*/
Expand Down
21 changes: 15 additions & 6 deletions src/main/java/org/openriak/riak/client/cli/InteractiveRiakCLI.java
Original file line number Diff line number Diff line change
Expand Up @@ -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 <bucket> [limit]");
break;
}
}
executeListKeys(parts[1], keysLimit);
} else {
System.out.println("Usage: list-keys <bucket>");
System.out.println("Usage: list-keys <bucket> [limit]");
}
break;
default:
Expand All @@ -720,7 +729,7 @@ private void printInteractiveHelp() {
System.out.println(" put <bucket> <key> <value> - Store a value in Riak");
System.out.println(" delete <bucket> <key> - Delete a value from Riak");
System.out.println(" list-buckets - List all buckets");
System.out.println(" list-keys <bucket> - List all keys in a bucket");
System.out.println(" list-keys <bucket> [limit] - List keys in a bucket (optional max count)");
System.out.println(" help - Show this help message");
System.out.println(" quit - Exit the CLI");
}
Expand Down Expand Up @@ -862,21 +871,21 @@ 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;
}

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<String> keys = (List<String>) result.getData().getOrDefault("keys", List.of());
System.out.println("Keys:");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down