Skip to content

Commit a61780f

Browse files
committed
disable InProcPubSubTests for now, server is brittle
1 parent 04031cf commit a61780f

3 files changed

Lines changed: 120 additions & 44 deletions

File tree

tests/StackExchange.Redis.Tests/PubSubTests.cs

Lines changed: 31 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -16,11 +16,14 @@ public class PubSubTests(ITestOutputHelper output, SharedConnectionFixture fixtu
1616
{
1717
}
1818

19+
/*
1920
[RunPerProtocol]
2021
public class InProcPubSubTests(ITestOutputHelper output, InProcServerFixture fixture)
2122
: PubSubTestBase(output, null, fixture)
2223
{
24+
protected override bool UseDedicatedInProcessServer => false;
2325
}
26+
*/
2427

2528
[RunPerProtocol]
2629
public abstract class PubSubTestBase(
@@ -32,7 +35,7 @@ public abstract class PubSubTestBase(
3235
[Fact]
3336
public async Task ExplicitPublishMode()
3437
{
35-
await using var conn = Create(channelPrefix: "foo:", log: Writer);
38+
await using var conn = ConnectFactory(channelPrefix: "foo:");
3639

3740
var pub = conn.GetSubscriber();
3841
int a = 0, b = 0, c = 0, d = 0;
@@ -70,13 +73,9 @@ await UntilConditionAsync(
7073
[InlineData("Foo:", true, "f")]
7174
public async Task TestBasicPubSub(string? channelPrefix, bool wildCard, string breaker)
7275
{
73-
// await using var conn = Create(channelPrefix: channelPrefix, shared: false, log: Writer);
74-
using var server = new InProcessTestServer(Output);
75-
var options = server.GetClientConfig();
76-
if (channelPrefix is not null) options.ChannelPrefix = RedisChannel.Literal(channelPrefix);
77-
await using var conn = await ConnectionMultiplexer.ConnectAsync(options);
76+
await using var conn = ConnectFactory(channelPrefix: channelPrefix, shared: false);
7877

79-
var pub = GetAnyPrimary(conn);
78+
var pub = GetAnyPrimary(conn.DefaultClient);
8079
var sub = conn.GetSubscriber();
8180
await PingAsync(pub, sub).ForAwait();
8281
HashSet<string?> received = [];
@@ -159,10 +158,10 @@ public async Task TestBasicPubSub(string? channelPrefix, bool wildCard, string b
159158
[Fact]
160159
public async Task TestBasicPubSubFireAndForget()
161160
{
162-
await using var conn = Create(shared: false, log: Writer);
161+
await using var conn = ConnectFactory(shared: false);
163162

164-
var profiler = conn.AddProfiler();
165-
var pub = GetAnyPrimary(conn);
163+
var profiler = conn.DefaultClient.AddProfiler();
164+
var pub = GetAnyPrimary(conn.DefaultClient);
166165
var sub = conn.GetSubscriber();
167166

168167
RedisChannel key = RedisChannel.Literal(Me() + Guid.NewGuid());
@@ -234,9 +233,9 @@ private async Task PingAsync(IServer pub, ISubscriber sub, int times = 1)
234233
[Fact]
235234
public async Task TestPatternPubSub()
236235
{
237-
await using var conn = Create(shared: false, log: Writer);
236+
await using var conn = ConnectFactory(shared: false);
238237

239-
var pub = GetAnyPrimary(conn);
238+
var pub = GetAnyPrimary(conn.DefaultClient);
240239
var sub = conn.GetSubscriber();
241240

242241
HashSet<string?> received = [];
@@ -293,7 +292,7 @@ public async Task TestPatternPubSub()
293292
[Fact]
294293
public async Task TestPublishWithNoSubscribers()
295294
{
296-
await using var conn = Create();
295+
await using var conn = ConnectFactory();
297296

298297
var sub = conn.GetSubscriber();
299298
#pragma warning disable CS0618
@@ -305,7 +304,7 @@ public async Task TestPublishWithNoSubscribers()
305304
public async Task TestMassivePublishWithWithoutFlush_Local()
306305
{
307306
Skip.UnlessLongRunning();
308-
await using var conn = Create();
307+
await using var conn = ConnectFactory();
309308

310309
var sub = conn.GetSubscriber();
311310
TestMassivePublish(sub, Me(), "local");
@@ -355,7 +354,7 @@ private void TestMassivePublish(ISubscriber sub, string channel, string caption)
355354
[Fact]
356355
public async Task SubscribeAsyncEnumerable()
357356
{
358-
await using var conn = Create(syncTimeout: 20000, shared: false, log: Writer);
357+
await using var conn = ConnectFactory(shared: false);
359358

360359
var sub = conn.GetSubscriber();
361360
RedisChannel channel = RedisChannel.Literal(Me());
@@ -390,7 +389,7 @@ public async Task SubscribeAsyncEnumerable()
390389
[Fact]
391390
public async Task PubSubGetAllAnyOrder()
392391
{
393-
await using var conn = Create(syncTimeout: 20000, shared: false, log: Writer);
392+
await using var conn = ConnectFactory(shared: false);
394393

395394
var sub = conn.GetSubscriber();
396395
RedisChannel channel = RedisChannel.Literal(Me());
@@ -645,9 +644,10 @@ public async Task PubSubGetAllCorrectOrder_OnMessage_Async()
645644
[Fact]
646645
public async Task TestPublishWithSubscribers()
647646
{
648-
await using var connA = Create(shared: false, log: Writer);
649-
await using var connB = Create(shared: false, log: Writer);
650-
await using var connPub = Create();
647+
await using var pair = ConnectFactory(shared: false);
648+
await using var connA = pair.DefaultClient;
649+
await using var connB = pair.CreateClient();
650+
await using var connPub = pair.CreateClient();
651651

652652
var channel = Me();
653653
var listenA = connA.GetSubscriber();
@@ -672,9 +672,10 @@ public async Task TestPublishWithSubscribers()
672672
[Fact]
673673
public async Task TestMultipleSubscribersGetMessage()
674674
{
675-
await using var connA = Create(shared: false, log: Writer);
676-
await using var connB = Create(shared: false, log: Writer);
677-
await using var connPub = Create();
675+
await using var pair = ConnectFactory(shared: false);
676+
await using var connA = pair.DefaultClient;
677+
await using var connB = pair.CreateClient();
678+
await using var connPub = pair.CreateClient();
678679

679680
var channel = RedisChannel.Literal(Me());
680681
var listenA = connA.GetSubscriber();
@@ -702,7 +703,7 @@ public async Task TestMultipleSubscribersGetMessage()
702703
[Fact]
703704
public async Task Issue38()
704705
{
705-
await using var conn = Create(log: Writer);
706+
await using var conn = ConnectFactory();
706707

707708
var sub = conn.GetSubscriber();
708709
int count = 0;
@@ -737,9 +738,10 @@ public async Task Issue38()
737738
[Fact]
738739
public async Task TestPartialSubscriberGetMessage()
739740
{
740-
await using var connA = Create();
741-
await using var connB = Create();
742-
await using var connPub = Create();
741+
await using var pair = ConnectFactory();
742+
await using var connA = pair.DefaultClient;
743+
await using var connB = pair.CreateClient();
744+
await using var connPub = pair.CreateClient();
743745

744746
int gotA = 0, gotB = 0;
745747
var listenA = connA.GetSubscriber();
@@ -770,8 +772,9 @@ public async Task TestPartialSubscriberGetMessage()
770772
[Fact]
771773
public async Task TestSubscribeUnsubscribeAndSubscribeAgain()
772774
{
773-
await using var connPub = Create();
774-
await using var connSub = Create();
775+
await using var pair = ConnectFactory();
776+
await using var connPub = pair.DefaultClient;
777+
await using var connSub = pair.CreateClient();
775778

776779
var prefix = Me();
777780
var pub = connPub.GetSubscriber();

tests/StackExchange.Redis.Tests/TestBase.cs

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -585,4 +585,74 @@ protected static async Task UntilConditionAsync(TimeSpan maxWaitTime, Func<bool>
585585
spent += wait;
586586
}
587587
}
588+
589+
// simplified usage to get an interchangeable dedicated vs shared in-process server, useful for debugging
590+
protected virtual bool UseDedicatedInProcessServer => false; // use the shared server by default
591+
internal ClientFactory ConnectFactory(bool allowAdmin = false, string? channelPrefix = null, bool shared = true)
592+
{
593+
if (UseDedicatedInProcessServer)
594+
{
595+
var server = new InProcessTestServer(Output);
596+
return new ClientFactory(this, allowAdmin, channelPrefix, shared, server);
597+
}
598+
return new ClientFactory(this, allowAdmin, channelPrefix, shared, null);
599+
}
600+
601+
internal sealed class ClientFactory : IDisposable, IAsyncDisposable
602+
{
603+
private readonly TestBase _testBase;
604+
private readonly bool _allowAdmin;
605+
private readonly string? _channelPrefix;
606+
private readonly bool _shared;
607+
private readonly InProcessTestServer? _server;
608+
private IInternalConnectionMultiplexer? _defaultClient;
609+
610+
internal ClientFactory(TestBase testBase, bool allowAdmin, string? channelPrefix, bool shared, InProcessTestServer? server)
611+
{
612+
_testBase = testBase;
613+
_allowAdmin = allowAdmin;
614+
_channelPrefix = channelPrefix;
615+
_shared = shared;
616+
_server = server;
617+
}
618+
619+
public IInternalConnectionMultiplexer DefaultClient => _defaultClient ??= CreateClient();
620+
621+
public InProcessTestServer? Server => _server;
622+
623+
public IInternalConnectionMultiplexer CreateClient()
624+
{
625+
if (_server is not null)
626+
{
627+
var config = _server.GetClientConfig();
628+
config.AllowAdmin = _allowAdmin;
629+
if (_channelPrefix is not null)
630+
{
631+
config.ChannelPrefix = RedisChannel.Literal(_channelPrefix);
632+
}
633+
return ConnectionMultiplexer.ConnectAsync(config).Result;
634+
}
635+
return _testBase.Create(allowAdmin: _allowAdmin, channelPrefix: _channelPrefix, shared: _shared);
636+
}
637+
638+
public IDatabase GetDatabase(int db = -1) => DefaultClient.GetDatabase(db);
639+
640+
public ISubscriber GetSubscriber() => DefaultClient.GetSubscriber();
641+
642+
public void Dispose()
643+
{
644+
_server?.Dispose();
645+
_defaultClient?.Dispose();
646+
}
647+
648+
public ValueTask DisposeAsync()
649+
{
650+
_server?.Dispose();
651+
if (_defaultClient is not null)
652+
{
653+
return _defaultClient.DisposeAsync();
654+
}
655+
return default;
656+
}
657+
}
588658
}

toys/StackExchange.Redis.Server/RedisServer.cs

Lines changed: 19 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -155,19 +155,21 @@ protected override void AppendStats(StringBuilder sb)
155155

156156
public override TypedRedisValue Execute(RedisClient client, in RedisRequest request)
157157
{
158-
var pw = Password;
159-
if (pw.Length != 0 & !client.IsAuthenticated)
158+
if (request.Count != 0)
160159
{
161-
if (!Literals.IsAuthCommand(in request))
162-
return TypedRedisValue.Error("NOAUTH Authentication required.");
163-
}
164-
else if (client.Protocol is RedisProtocol.Resp2 && client.IsSubscriber &&
165-
!Literals.IsPubSubCommand(in request, out var cmd))
166-
{
167-
return TypedRedisValue.Error(
168-
$"ERR only (P|S)SUBSCRIBE / (P|S)UNSUBSCRIBE / PING / QUIT allowed in this context (got: '{cmd}')");
160+
var pw = Password;
161+
if (pw.Length != 0 & !client.IsAuthenticated)
162+
{
163+
if (!Literals.IsAuthCommand(in request))
164+
return TypedRedisValue.Error("NOAUTH Authentication required.");
165+
}
166+
else if (client.Protocol is RedisProtocol.Resp2 && client.IsSubscriber &&
167+
!Literals.IsPubSubCommand(in request, out var cmd))
168+
{
169+
return TypedRedisValue.Error(
170+
$"ERR only (P|S)SUBSCRIBE / (P|S)UNSUBSCRIBE / PING / QUIT allowed in this context (got: '{cmd}')");
171+
}
169172
}
170-
171173
return base.Execute(client, request);
172174
}
173175

@@ -183,24 +185,25 @@ public static readonly CommandBytes
183185
PSUBSCRIBE = new("PSUBSCRIBE"u8),
184186
SSUBSCRIBE = new("SSUBSCRIBE"u8),
185187
UNSUBSCRIBE = new("UNSUBSCRIBE"u8),
186-
PUNUBSCRIBE = new("PUNUBSCRIBE"u8),
188+
PUNSUBSCRIBE = new("PUNSUBSCRIBE"u8),
187189
SUNSUBSCRIBE = new("SUNSUBSCRIBE"u8);
188190

189191
public static bool IsAuthCommand(in RedisRequest request) =>
190192
request.Count != 0 && request.TryGetCommandBytes(0, out var command)
191193
&& (command.Equals(AUTH) || command.Equals(HELLO));
192194

193-
public static bool IsPubSubCommand(in RedisRequest request, out CommandBytes command)
195+
public static bool IsPubSubCommand(in RedisRequest request, out string badCommand)
194196
{
195-
if (request.Count == 0 || !request.TryGetCommandBytes(0, out command))
197+
badCommand = "";
198+
if (request.Count == 0 || !request.TryGetCommandBytes(0, out var command))
196199
{
197-
command = default;
200+
if (request.Count != 0) badCommand = request.GetString(0);
198201
return false;
199202
}
200203

201204
return command.Equals(SUBSCRIBE) || command.Equals(UNSUBSCRIBE)
202205
|| command.Equals(SSUBSCRIBE) || command.Equals(SUNSUBSCRIBE)
203-
|| command.Equals(PSUBSCRIBE) || command.Equals(PUNUBSCRIBE)
206+
|| command.Equals(PSUBSCRIBE) || command.Equals(PUNSUBSCRIBE)
204207
|| command.Equals(PING) || command.Equals(QUIT);
205208
}
206209
}

0 commit comments

Comments
 (0)