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

Commit b4cadc0

Browse files
basic support for writing
1 parent 75964ab commit b4cadc0

7 files changed

Lines changed: 251 additions & 5 deletions

File tree

Lines changed: 127 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,127 @@
1+
namespace SqlStreamStore.HAL.Tests
2+
{
3+
using System;
4+
using System.Net;
5+
using System.Net.Http;
6+
using System.Net.Http.Headers;
7+
using System.Threading.Tasks;
8+
using Newtonsoft.Json.Linq;
9+
using Shouldly;
10+
using SqlStreamStore.Streams;
11+
using Xunit;
12+
13+
public class StreamAppendTests : IDisposable
14+
{
15+
private readonly SqlStreamStoreHalMiddlewareFixture _fixture;
16+
17+
public StreamAppendTests()
18+
{
19+
_fixture = new SqlStreamStoreHalMiddlewareFixture();
20+
}
21+
22+
[Fact]
23+
public async Task append_expected_version_any()
24+
{
25+
var messageId = Guid.NewGuid();
26+
27+
var jsonData = JObject.FromObject(new
28+
{
29+
property = "value"
30+
});
31+
32+
var jsonMetadata = JObject.FromObject(new
33+
{
34+
property = "metaValue"
35+
});
36+
37+
using(var response = await _fixture.HttpClient.PostAsync(
38+
"/streams/a-stream",
39+
new StringContent(JObject.FromObject(new
40+
{
41+
expectedVersion = ExpectedVersion.Any,
42+
messages = new[]
43+
{
44+
new
45+
{
46+
messageId,
47+
type = "type",
48+
jsonData,
49+
jsonMetadata
50+
}
51+
}
52+
}).ToString())
53+
{
54+
Headers =
55+
{
56+
ContentType = new MediaTypeHeaderValue("application/json")
57+
}
58+
}))
59+
{
60+
response.StatusCode.ShouldBe(HttpStatusCode.OK);
61+
}
62+
63+
var page = await _fixture.StreamStore.ReadStreamForwards("a-stream", 0, 1);
64+
65+
page.Status.ShouldBe(PageReadStatus.Success);
66+
page.Messages.Length.ShouldBe(1);
67+
page.Messages[0].MessageId.ShouldBe(messageId);
68+
page.Messages[0].Type.ShouldBe("type");
69+
JToken.DeepEquals(JObject.Parse(await page.Messages[0].GetJsonData()), jsonData).ShouldBeTrue();
70+
JToken.DeepEquals(JObject.Parse(page.Messages[0].JsonMetadata), jsonMetadata).ShouldBeTrue();
71+
}
72+
73+
[Fact]
74+
public async Task append_expected_version_no_stream()
75+
{
76+
var messageId = Guid.NewGuid();
77+
78+
var jsonData = JObject.FromObject(new
79+
{
80+
property = "value"
81+
});
82+
83+
var jsonMetadata = JObject.FromObject(new
84+
{
85+
property = "metaValue"
86+
});
87+
88+
using(var response = await _fixture.HttpClient.PostAsync(
89+
"/streams/a-stream",
90+
new StringContent(JObject.FromObject(new
91+
{
92+
expectedVersion = ExpectedVersion.NoStream,
93+
messages = new[]
94+
{
95+
new
96+
{
97+
messageId,
98+
type = "type",
99+
jsonData,
100+
jsonMetadata
101+
}
102+
}
103+
}).ToString())
104+
{
105+
Headers =
106+
{
107+
ContentType = new MediaTypeHeaderValue("application/json")
108+
}
109+
}))
110+
{
111+
response.StatusCode.ShouldBe(HttpStatusCode.Created);
112+
response.Headers.Location.ToString().ShouldBe("streams/a-stream");
113+
}
114+
115+
var page = await _fixture.StreamStore.ReadStreamForwards("a-stream", 0, 1);
116+
117+
page.Status.ShouldBe(PageReadStatus.Success);
118+
page.Messages.Length.ShouldBe(1);
119+
page.Messages[0].MessageId.ShouldBe(messageId);
120+
page.Messages[0].Type.ShouldBe("type");
121+
JToken.DeepEquals(JObject.Parse(await page.Messages[0].GetJsonData()), jsonData).ShouldBeTrue();
122+
JToken.DeepEquals(JObject.Parse(page.Messages[0].JsonMetadata), jsonMetadata).ShouldBeTrue();
123+
}
124+
125+
public void Dispose() => _fixture.Dispose();
126+
}
127+
}
Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
1+
namespace SqlStreamStore.HAL
2+
{
3+
using Microsoft.Owin;
4+
using Microsoft.Owin.Builder;
5+
using Owin;
6+
using MidFunc = System.Func<System.Func<System.Collections.Generic.IDictionary<string, object>,
7+
System.Threading.Tasks.Task
8+
>, System.Func<System.Collections.Generic.IDictionary<string, object>,
9+
System.Threading.Tasks.Task>
10+
>;
11+
12+
internal static class AppendStreamMiddleware
13+
{
14+
public static MidFunc UseStreamStore(IStreamStore streamStore)
15+
{
16+
var stream = new StreamResource(streamStore);
17+
18+
var builder = new AppBuilder()
19+
.MapWhen(IsStream, inner => inner.Use(AppendStream(stream)));
20+
21+
return next =>
22+
{
23+
builder.Run(ctx => next(ctx.Environment));
24+
25+
return builder.Build();
26+
};
27+
}
28+
29+
private static bool IsStream(IOwinContext context)
30+
=> context.IsPost() && context.Request.Path.Value?.Length > 1;
31+
32+
private static MidFunc AppendStream(StreamResource stream) => next => async env =>
33+
{
34+
var context = new OwinContext(env);
35+
36+
var options = await AppendStreamOptions.Create(context.Request, context.Request.CallCancelled);
37+
38+
var response = await stream.AppendMessages(options, context.Request.CallCancelled);
39+
40+
if(response.StatusCode == 201)
41+
{
42+
context.Response.ReasonPhrase = "Created";
43+
context.Response.Headers["Location"] = $"streams/{options.StreamId}";
44+
}
45+
46+
await context.WriteHalResponse(response);
47+
};
48+
49+
}
50+
}
Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
1+
namespace SqlStreamStore.HAL
2+
{
3+
using System;
4+
using System.IO;
5+
using System.Linq;
6+
using System.Threading;
7+
using System.Threading.Tasks;
8+
using Microsoft.Owin;
9+
using Newtonsoft.Json;
10+
using Newtonsoft.Json.Linq;
11+
using SqlStreamStore.Streams;
12+
13+
internal class AppendStreamOptions
14+
{
15+
public static async Task<AppendStreamOptions> Create(IOwinRequest request, CancellationToken ct)
16+
{
17+
using(var reader = new JsonTextReader(new StreamReader(request.Body))
18+
{
19+
CloseInput = false
20+
})
21+
{
22+
var body = await JObject.LoadAsync(reader, ct);
23+
24+
return new AppendStreamOptions(request, body);
25+
}
26+
}
27+
28+
private AppendStreamOptions(IOwinRequest request, JObject body)
29+
{
30+
StreamId = request.Path.Value.Remove(0, 1);
31+
32+
ExpectedVersion = body.Value<int>("expectedVersion");
33+
34+
NewStreamMessages = body.Value<JArray>("messages")
35+
.Select(newStreamMessage => new NewStreamMessage(
36+
Guid.Parse(newStreamMessage.Value<string>("messageId")),
37+
newStreamMessage.Value<string>("type"),
38+
newStreamMessage.Value<JObject>("jsonData").ToString(),
39+
newStreamMessage.Value<JObject>("jsonMetadata")?.ToString()))
40+
.ToArray();
41+
}
42+
43+
public string StreamId { get; }
44+
public int ExpectedVersion { get; }
45+
public NewStreamMessage[] NewStreamMessages { get; }
46+
47+
public Func<IStreamStore, CancellationToken, Task<AppendResult>> GetAppendOperation()
48+
=> (streamStore, ct) => streamStore.AppendToStream(StreamId, ExpectedVersion, NewStreamMessages, ct);
49+
}
50+
}

src/SqlStreamStore.HAL/OwinContextExtensions.cs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,5 +41,8 @@ public static async Task WriteHalResponse(this IOwinContext context, Response re
4141

4242
public static bool IsGetOrHead(this IOwinContext context)
4343
=> context.Request.Method == "GET" || context.Request.Method == "HEAD";
44+
45+
public static bool IsPost(this IOwinContext context)
46+
=> context.Request.Method == "POST";
4447
}
4548
}

src/SqlStreamStore.HAL/ReadStreamMiddleware.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ namespace SqlStreamStore.HAL
1212

1313
internal static class ReadStreamMiddleware
1414
{
15-
public static MidFunc UseStreamStore(IReadonlyStreamStore streamStore)
15+
public static MidFunc UseStreamStore(IStreamStore streamStore)
1616
{
1717
var streams = new StreamResource(streamStore);
1818

src/SqlStreamStore.HAL/SqlStreamStoreHalMiddleware.cs

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ public static class SqlStreamStoreHalMiddleware
4545
return context.WriteHalResponse(response);
4646
};
4747

48-
public static MidFunc UseSqlStreamStoreHal(IReadonlyStreamStore streamStore)
48+
public static MidFunc UseSqlStreamStoreHal(IStreamStore streamStore)
4949
{
5050
if(streamStore == null)
5151
throw new ArgumentNullException(nameof(streamStore));
@@ -58,7 +58,8 @@ public static MidFunc UseSqlStreamStoreHal(IReadonlyStreamStore streamStore)
5858
.Use(MethodsNotAllowed("POST", "PUT", "DELETE", "TRACE", "PATCH", "OPTIONS")))
5959
.Map("/streams", inner => inner
6060
.Use(ReadStreamMiddleware.UseStreamStore(streamStore))
61-
.Use(MethodsNotAllowed("POST", "PUT", "DELETE", "TRACE", "PATCH", "OPTIONS")));
61+
.Use(AppendStreamMiddleware.UseStreamStore(streamStore))
62+
.Use(MethodsNotAllowed("PUT", "DELETE", "TRACE", "PATCH", "OPTIONS")));
6263

6364
return next =>
6465
{

src/SqlStreamStore.HAL/StreamResource.cs

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,15 +10,30 @@ namespace SqlStreamStore.HAL
1010

1111
internal class StreamResource
1212
{
13-
private readonly IReadonlyStreamStore _streamStore;
13+
private readonly IStreamStore _streamStore;
1414

15-
public StreamResource(IReadonlyStreamStore streamStore)
15+
public StreamResource(IStreamStore streamStore)
1616
{
1717
if(streamStore == null)
1818
throw new ArgumentNullException(nameof(streamStore));
1919
_streamStore = streamStore;
2020
}
2121

22+
public async Task<Response> AppendMessages(
23+
AppendStreamOptions options,
24+
CancellationToken cancellationToken)
25+
{
26+
var operation = options.GetAppendOperation();
27+
28+
var result = await operation.Invoke(_streamStore, cancellationToken);
29+
30+
return new Response(
31+
new HALResponse(new object()),
32+
options.ExpectedVersion == ExpectedVersion.NoStream
33+
? 201
34+
: 200);
35+
}
36+
2237
public async Task<Response> GetMessage(
2338
ReadStreamMessageOptions options,
2439
CancellationToken cancellationToken)

0 commit comments

Comments
 (0)