Skip to content

Commit c0a07db

Browse files
Fix large string handling in CLUSTER|NODES and CLUSTER|SHARDS (#1858)
* fix CLUSTER|NODES and CLUSTER|SHARDS to handle larger string responses; another example of the issue we should cleanup up re. #1781 * fixes * more fixes * even more fixes; tricky to actually test this
1 parent c5ec666 commit c0a07db

2 files changed

Lines changed: 89 additions & 6 deletions

File tree

libs/cluster/Session/RespClusterBasicCommands.cs

Lines changed: 81 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
// Licensed under the MIT license.
33

44
using System;
5+
using System.Buffers;
56
using System.Diagnostics;
67
using System.Text;
78
using Garnet.common;
@@ -274,8 +275,8 @@ private bool NetworkClusterNodes(out bool invalidParameters)
274275
}
275276

276277
var nodes = clusterProvider.clusterManager.CurrentConfig.GetClusterInfo(clusterProvider);
277-
while (!RespWriteUtils.TryWriteAsciiBulkString(nodes, ref dcurr, dend))
278-
SendAndReset();
278+
279+
WriteAsciiLargeRespString(nodes);
279280

280281
return true;
281282
}
@@ -343,8 +344,8 @@ private bool NetworkClusterShards(out bool invalidParameters)
343344

344345
var preferredType = clusterProvider.serverOptions.ClusterPreferredEndpointType;
345346
var shardsInfo = clusterProvider.clusterManager.CurrentConfig.GetShardsInfo(clusterProvider.clusterManager.clusterConnectionStore, preferredType);
346-
while (!RespWriteUtils.TryWriteAsciiDirect(shardsInfo, ref dcurr, dend))
347-
SendAndReset();
347+
348+
WriteLargeAsciiDirectString(shardsInfo);
348349

349350
return true;
350351
}
@@ -523,5 +524,81 @@ private bool NetworkClusterPublish(out bool invalidParameters)
523524
clusterProvider.storeWrapper.subscribeBroker.Publish(parseState.GetArgSliceByRef(0), parseState.GetArgSliceByRef(1));
524525
return true;
525526
}
527+
528+
529+
/// <summary>
530+
/// Handle a potentially large already encoded RESP response, breaking it into pieces if necessary.
531+
/// </summary>
532+
private void WriteLargeAsciiDirectString(ReadOnlySpan<char> message)
533+
{
534+
// Attempt to write w/o any buffering
535+
if (RespWriteUtils.TryWriteAsciiDirect(message, ref dcurr, dend))
536+
{
537+
return;
538+
}
539+
540+
var buffer = ArrayPool<byte>.Shared.Rent(message.Length);
541+
try
542+
{
543+
var len = Encoding.ASCII.GetBytes(message, buffer);
544+
var remaining = buffer.AsSpan()[..len];
545+
546+
while (!remaining.IsEmpty)
547+
{
548+
var space = Math.Min((int)(dend - dcurr), remaining.Length);
549+
_ = RespWriteUtils.TryWriteDirect(remaining[..space], ref dcurr, dend);
550+
551+
SendAndReset();
552+
553+
remaining = remaining[space..];
554+
}
555+
}
556+
finally
557+
{
558+
ArrayPool<byte>.Shared.Return(buffer);
559+
}
560+
}
561+
562+
/// <summary>
563+
/// Handle a potentially large bulk string by breaking it into pieces if necessary.
564+
/// </summary>
565+
private void WriteAsciiLargeRespString(ReadOnlySpan<char> message)
566+
{
567+
while (!RespWriteUtils.TryWriteBulkStringLength(message.Length, ref dcurr, dend))
568+
SendAndReset();
569+
570+
// Attempt to write w/o any buffering
571+
if (RespWriteUtils.TryWriteAsciiDirect(message, ref dcurr, dend))
572+
{
573+
while (!RespWriteUtils.TryWriteNewLine(ref dcurr, dend))
574+
SendAndReset();
575+
576+
return;
577+
}
578+
579+
var buffer = ArrayPool<byte>.Shared.Rent(message.Length);
580+
try
581+
{
582+
var len = Encoding.ASCII.GetBytes(message, buffer);
583+
var remaining = buffer.AsSpan()[..len];
584+
585+
while (!remaining.IsEmpty)
586+
{
587+
var space = Math.Min((int)(dend - dcurr), remaining.Length);
588+
_ = RespWriteUtils.TryWriteDirect(remaining[..space], ref dcurr, dend);
589+
590+
SendAndReset();
591+
592+
remaining = remaining[space..];
593+
}
594+
595+
while (!RespWriteUtils.TryWriteNewLine(ref dcurr, dend))
596+
SendAndReset();
597+
}
598+
finally
599+
{
600+
ArrayPool<byte>.Shared.Return(buffer);
601+
}
602+
}
526603
}
527604
}

libs/common/RespWriteUtils.cs

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -321,14 +321,20 @@ public static bool TryWriteDirect<T>(ref T item, ref byte* curr, byte* end) wher
321321
/// Write length header of bulk string
322322
/// </summary>
323323
public static bool TryWriteBulkStringLength(ReadOnlySpan<byte> item, ref byte* curr, byte* end)
324+
=> TryWriteBulkStringLength(item.Length, ref curr, end);
325+
326+
/// <summary>
327+
/// Write length header of bulk string
328+
/// </summary>
329+
public static bool TryWriteBulkStringLength(int len, ref byte* curr, byte* end)
324330
{
325-
var itemDigits = NumUtils.CountDigits(item.Length);
331+
var itemDigits = NumUtils.CountDigits(len);
326332
var totalLen = 1 + itemDigits + 2;
327333
if (totalLen > (int)(end - curr))
328334
return false;
329335

330336
*curr++ = (byte)'$';
331-
NumUtils.WriteInt32(item.Length, itemDigits, ref curr);
337+
NumUtils.WriteInt32(len, itemDigits, ref curr);
332338
WriteNewline(ref curr);
333339
return true;
334340
}

0 commit comments

Comments
 (0)