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

Commit fe66014

Browse files
Merge pull request #4 from thefringeninja/writing
Support Append and Delete
2 parents 5ba4d39 + a39c866 commit fe66014

13 files changed

Lines changed: 563 additions & 15 deletions
Lines changed: 198 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,198 @@
1+
namespace SqlStreamStore.HAL.Tests
2+
{
3+
using System;
4+
using System.Collections.Generic;
5+
using System.Net;
6+
using System.Net.Http;
7+
using System.Net.Http.Headers;
8+
using System.Threading.Tasks;
9+
using Newtonsoft.Json.Linq;
10+
using Shouldly;
11+
using SqlStreamStore.Streams;
12+
using Xunit;
13+
14+
public class StreamAppendTests : IDisposable
15+
{
16+
private const string StreamId = "a-stream";
17+
private readonly SqlStreamStoreHalMiddlewareFixture _fixture;
18+
19+
public StreamAppendTests()
20+
{
21+
_fixture = new SqlStreamStoreHalMiddlewareFixture();
22+
}
23+
24+
public static IEnumerable<object[]> AppendCases()
25+
{
26+
var messageId = Guid.NewGuid();
27+
28+
var jsonData = JObject.FromObject(new
29+
{
30+
property = "value"
31+
});
32+
33+
var jsonMetadata = JObject.FromObject(new
34+
{
35+
property = "metaValue"
36+
});
37+
38+
var bodies = new JToken[]
39+
{
40+
JObject.FromObject(new
41+
{
42+
messageId,
43+
type = "type",
44+
jsonData,
45+
jsonMetadata
46+
}),
47+
JArray.FromObject(new[]
48+
{
49+
new
50+
{
51+
messageId,
52+
type = "type",
53+
jsonData,
54+
jsonMetadata
55+
}
56+
})
57+
};
58+
59+
foreach(var body in bodies)
60+
{
61+
yield return new object[]
62+
{
63+
body,
64+
ExpectedVersion.Any,
65+
HttpStatusCode.OK,
66+
messageId,
67+
jsonData,
68+
jsonMetadata
69+
};
70+
yield return new object[]
71+
{
72+
body,
73+
default(int?),
74+
HttpStatusCode.OK,
75+
messageId,
76+
jsonData,
77+
jsonMetadata
78+
};
79+
yield return new object[]
80+
{
81+
body,
82+
ExpectedVersion.NoStream,
83+
HttpStatusCode.Created,
84+
messageId,
85+
jsonData,
86+
jsonMetadata
87+
};
88+
}
89+
}
90+
91+
[Theory, MemberData(nameof(AppendCases))]
92+
public async Task expected_version(
93+
JToken body,
94+
int? expectedVersion,
95+
HttpStatusCode statusCode,
96+
Guid messageId,
97+
JObject jsonData,
98+
JObject jsonMetadata)
99+
{
100+
var request = new HttpRequestMessage(HttpMethod.Post, $"/streams/{StreamId}")
101+
{
102+
Content = new StringContent(body.ToString())
103+
};
104+
105+
if(expectedVersion.HasValue)
106+
{
107+
request.Headers.Add(Constants.Headers.ExpectedVersion, $"{expectedVersion}");
108+
}
109+
110+
using(var response = await _fixture.HttpClient.SendAsync(request))
111+
{
112+
response.StatusCode.ShouldBe(statusCode);
113+
}
114+
115+
var page = await _fixture.StreamStore.ReadStreamForwards(StreamId, 0, int.MaxValue);
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+
[Theory]
126+
[InlineData(new[]{ExpectedVersion.NoStream, ExpectedVersion.NoStream})]
127+
[InlineData(new[]{ExpectedVersion.NoStream, 2})]
128+
public async Task wrong_expected_version(int[] expectedVersions)
129+
{
130+
var jsonData = JObject.FromObject(new
131+
{
132+
property = "value"
133+
});
134+
135+
var jsonMetadata = JObject.FromObject(new
136+
{
137+
property = "metaValue"
138+
});
139+
140+
for(var i = 0; i < expectedVersions.Length - 1; i++)
141+
{
142+
using(await _fixture.HttpClient.SendAsync(new HttpRequestMessage(HttpMethod.Post, $"/streams/{StreamId}")
143+
{
144+
Headers =
145+
{
146+
{Constants.Headers.ExpectedVersion, $"{expectedVersions[i]}"}
147+
},
148+
Content = new StringContent(JObject.FromObject(new
149+
{
150+
messageId = Guid.NewGuid(),
151+
type = "type",
152+
jsonData,
153+
jsonMetadata
154+
}).ToString())
155+
{
156+
Headers =
157+
{
158+
ContentType = new MediaTypeHeaderValue("application/json")
159+
}
160+
}
161+
}))
162+
{ }
163+
}
164+
165+
using(var response = await _fixture.HttpClient.SendAsync(new HttpRequestMessage(HttpMethod.Post, $"/streams/{StreamId}")
166+
{
167+
Headers =
168+
{
169+
{Constants.Headers.ExpectedVersion, $"{expectedVersions[expectedVersions.Length - 1]}"}
170+
},
171+
Content = new StringContent(JObject.FromObject(new
172+
{
173+
messageId = Guid.NewGuid(),
174+
type = "type",
175+
jsonData,
176+
jsonMetadata
177+
}).ToString())
178+
{
179+
Headers =
180+
{
181+
ContentType = new MediaTypeHeaderValue("application/json")
182+
}
183+
}
184+
}))
185+
{
186+
response.StatusCode.ShouldBe(HttpStatusCode.Conflict);
187+
response.Content.Headers.ContentType.ShouldBe(new MediaTypeHeaderValue(
188+
Constants.Headers.ContentTypes.ProblemDetails));
189+
}
190+
var page = await _fixture.StreamStore.ReadStreamForwards(StreamId, 0, int.MaxValue);
191+
192+
page.Status.ShouldBe(PageReadStatus.Success);
193+
page.Messages.Length.ShouldBe(expectedVersions.Length - 1);
194+
}
195+
196+
public void Dispose() => _fixture.Dispose();
197+
}
198+
}
Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,71 @@
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 Shouldly;
9+
using SqlStreamStore.Streams;
10+
using Xunit;
11+
12+
public class StreamDeleteTests : IDisposable
13+
{
14+
private const string StreamId = "a-stream";
15+
private readonly SqlStreamStoreHalMiddlewareFixture _fixture;
16+
17+
public StreamDeleteTests()
18+
{
19+
_fixture = new SqlStreamStoreHalMiddlewareFixture();
20+
}
21+
22+
[Theory, InlineData(ExpectedVersion.Any), InlineData(0), InlineData(null)]
23+
public async Task expected_version(int? expectedVersion)
24+
{
25+
await _fixture.WriteNMessages(StreamId, 1);
26+
27+
var request = new HttpRequestMessage(HttpMethod.Delete, $"/streams/{StreamId}");
28+
29+
if(expectedVersion.HasValue)
30+
{
31+
request.Headers.Add(Constants.Headers.ExpectedVersion, $"{expectedVersion}");
32+
}
33+
34+
using(var response = await _fixture.HttpClient.SendAsync(request))
35+
{
36+
response.StatusCode.ShouldBe(HttpStatusCode.OK);
37+
}
38+
39+
var page = await _fixture.StreamStore.ReadStreamForwards(StreamId, 0, 1);
40+
41+
page.Status.ShouldBe(PageReadStatus.StreamNotFound);
42+
}
43+
44+
[Theory, InlineData(ExpectedVersion.NoStream), InlineData(2)]
45+
public async Task wrong_expected_version(int expectedVersion)
46+
{
47+
await _fixture.WriteNMessages(StreamId, 1);
48+
var request = new HttpRequestMessage(HttpMethod.Delete, $"/streams/{StreamId}")
49+
{
50+
Headers =
51+
{
52+
{ Constants.Headers.ExpectedVersion, $"{expectedVersion}" }
53+
}
54+
};
55+
56+
using(var response = await _fixture.HttpClient.SendAsync(request))
57+
{
58+
response.StatusCode.ShouldBe(HttpStatusCode.Conflict);
59+
response.Content.Headers.ContentType.ShouldBe(new MediaTypeHeaderValue(
60+
Constants.Headers.ContentTypes.ProblemDetails));
61+
}
62+
63+
var page = await _fixture.StreamStore.ReadStreamForwards(StreamId, 0, 1);
64+
65+
page.Status.ShouldBe(PageReadStatus.Success);
66+
67+
}
68+
69+
public void Dispose() => _fixture.Dispose();
70+
}
71+
}
Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,57 @@
1+
namespace SqlStreamStore.HAL
2+
{
3+
using Microsoft.Owin;
4+
using Microsoft.Owin.Builder;
5+
using Owin;
6+
using SqlStreamStore.Streams;
7+
using MidFunc = System.Func<System.Func<System.Collections.Generic.IDictionary<string, object>,
8+
System.Threading.Tasks.Task
9+
>, System.Func<System.Collections.Generic.IDictionary<string, object>,
10+
System.Threading.Tasks.Task>
11+
>;
12+
13+
internal static class AppendStreamMiddleware
14+
{
15+
public static MidFunc UseStreamStore(IStreamStore streamStore)
16+
{
17+
var stream = new StreamResource(streamStore);
18+
19+
var builder = new AppBuilder()
20+
.MapWhen(IsStream, inner => inner.Use(AppendStream(stream)));
21+
22+
return next =>
23+
{
24+
builder.Run(ctx => next(ctx.Environment));
25+
26+
return builder.Build();
27+
};
28+
}
29+
30+
private static bool IsStream(IOwinContext context)
31+
=> context.IsPost() && context.Request.Path.Value?.Length > 1;
32+
33+
private static MidFunc AppendStream(StreamResource stream) => next => async env =>
34+
{
35+
var context = new OwinContext(env);
36+
37+
var options = await AppendStreamOptions.Create(context.Request, context.Request.CallCancelled);
38+
39+
try
40+
{
41+
var response = await stream.AppendMessages(options, context.Request.CallCancelled);
42+
43+
if(response.StatusCode == 201)
44+
{
45+
context.Response.ReasonPhrase = "Created";
46+
context.Response.Headers["Location"] = $"streams/{options.StreamId}";
47+
}
48+
49+
await context.WriteHalResponse(response);
50+
}
51+
catch(WrongExpectedVersionException ex)
52+
{
53+
await context.WriteProblemDetailsResponse(ex);
54+
}
55+
};
56+
}
57+
}

0 commit comments

Comments
 (0)