Skip to content
This repository was archived by the owner on Sep 3, 2024. It is now read-only.

Commit 1fd8d4e

Browse files
fixed inconsistent sorting
http client will have to re-order when reading backwards
1 parent 46d5720 commit 1fd8d4e

2 files changed

Lines changed: 38 additions & 30 deletions

File tree

src/SqlStreamStore.HAL/Resources/AllStreamResource.cs

Lines changed: 20 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,9 @@ public async Task<Response> GetPage(
3232
{
3333
var page = await operation.Invoke(_streamStore, cancellationToken);
3434

35-
var payloads = await Task.WhenAll(page.Messages
35+
var streamMessages = page.Messages.OrderByDescending(m => m.Position).ToArray();
36+
37+
var payloads = await Task.WhenAll(streamMessages
3638
.Select(message => operation.EmbedPayload
3739
? message.GetJsonData(cancellationToken)
3840
: Task.FromResult<string>(null))
@@ -50,29 +52,30 @@ public async Task<Response> GetPage(
5052
.AddLinks(Links.All.Feed(operation))
5153
.AddEmbeddedCollection(
5254
Constants.Relations.Message,
53-
page.Messages.Zip(payloads,
55+
streamMessages.Zip(
56+
payloads,
5457
(message, payload) => new HALResponse(new
55-
{
56-
message.MessageId,
57-
message.CreatedUtc,
58-
message.Position,
59-
message.StreamId,
60-
message.StreamVersion,
61-
message.Type,
62-
payload,
63-
metadata = message.JsonMetadata
64-
}).AddLinks(
65-
Links.All.Self(message)))));
58+
{
59+
message.MessageId,
60+
message.CreatedUtc,
61+
message.Position,
62+
message.StreamId,
63+
message.StreamVersion,
64+
message.Type,
65+
payload,
66+
metadata = message.JsonMetadata
67+
})
68+
.AddLinks(Links.All.Self(message)))));
6669

6770
if(operation.FromPositionInclusive == Position.End)
6871
{
69-
var headPosition = page.Messages.Length > 0
70-
? page.Messages[0].Position
72+
var headPosition = streamMessages.Length > 0
73+
? streamMessages[0].Position
7174
: Position.End;
72-
75+
7376
response.Headers[Constants.Headers.HeadPosition] = new[] { $"{headPosition}" };
7477
}
75-
78+
7679
return response;
7780
}
7881
}

src/SqlStreamStore.HAL/Resources/StreamResource.cs

Lines changed: 18 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -45,14 +45,17 @@ public async Task<Response> AppendMessages(
4545
{
4646
response.Headers[Constants.Headers.Location] = new[] { $"streams/{operation.StreamId}" };
4747
}
48+
4849
return response;
4950
}
5051

5152
public async Task<Response> GetPage(ReadStreamOperation operation, CancellationToken cancellationToken)
5253
{
5354
var page = await operation.Invoke(_streamStore, cancellationToken);
5455

55-
var payloads = await Task.WhenAll(page.Messages
56+
var streamMessages = page.Messages.OrderByDescending(m => m.Position).ToArray();
57+
58+
var payloads = await Task.WhenAll(streamMessages
5659
.Select(message => operation.EmbedPayload
5760
? message.GetJsonData(cancellationToken)
5861
: Task.FromResult<string>(null))
@@ -73,25 +76,27 @@ public async Task<Response> GetPage(ReadStreamOperation operation, CancellationT
7376
.AddLinks(Links.Stream.Metadata(operation))
7477
.AddEmbeddedCollection(
7578
Constants.Relations.Message,
76-
page.Messages.Zip(payloads,
79+
streamMessages.Zip(
80+
payloads,
7781
(message, payload) => new HALResponse(new
78-
{
79-
message.MessageId,
80-
message.CreatedUtc,
81-
message.Position,
82-
message.StreamId,
83-
message.StreamVersion,
84-
message.Type,
85-
payload,
86-
metadata = message.JsonMetadata
87-
}).AddLinks(Links.StreamMessage.Self(message)))),
82+
{
83+
message.MessageId,
84+
message.CreatedUtc,
85+
message.Position,
86+
message.StreamId,
87+
message.StreamVersion,
88+
message.Type,
89+
payload,
90+
metadata = message.JsonMetadata
91+
})
92+
.AddLinks(Links.StreamMessage.Self(message)))),
8893
page.Status == PageReadStatus.StreamNotFound ? 404 : 200);
8994
}
9095

9196
public async Task<Response> Delete(DeleteStreamOperation operation, CancellationToken cancellationToken)
9297
{
9398
await operation.Invoke(_streamStore, cancellationToken);
94-
99+
95100
return new Response(new HALResponse(new object()));
96101
}
97102
}

0 commit comments

Comments
 (0)