using PrivaPub.Tests.Support; using PrivaPub.Tests.Support.Host; using System.Net; using System.Net.WebSockets; using System.Text; using System.Text.Json.Nodes; namespace PrivaPub.Tests.Http { // Mastodon's streaming API: a WebSocket whose streams are subscribed by message or in the URL, and server-sent events // per stream; each event mapped for its persona, never what it may not see [Trait("Category", "Integration")] public sealed class MastodonStreamingTests : IAsyncLifetime { PrivaPubHost _host; static CancellationToken Token => TestContext.Current.CancellationToken; public async ValueTask InitializeAsync() { Assert.SkipUnless(MongoFixture.Enabled, MongoFixture.Skip); _host = await PrivaPubHost.Shared(); } public ValueTask DisposeAsync() => ValueTask.CompletedTask; async Task Connect(Mastodon account, string query = "") { var client = _host.Server.CreateWebSocketClient(); var uri = new Uri($"ws://{PrivaPubHost.Host}/api/v1/streaming?access_token={account.Token}{query}"); var socket = await client.ConnectAsync(uri, Token); // the subscription is made as the socket opens: give the server a moment to hear it await Task.Delay(200, Token); return socket; } static async Task Say(WebSocket socket, JsonObject message) => await socket.SendAsync(Encoding.UTF8.GetBytes(message.ToJsonString()), WebSocketMessageType.Text, true, Token); // the next message for which `wanted` holds, within a few seconds, or null static async Task Next(WebSocket socket, Func wanted, int seconds = 5) { using var timeout = CancellationTokenSource.CreateLinkedTokenSource(Token); timeout.CancelAfter(TimeSpan.FromSeconds(seconds)); var buffer = new byte[65536]; try { while (true) { var text = new StringBuilder(); WebSocketReceiveResult received; do { received = await socket.ReceiveAsync(buffer, timeout.Token); text.Append(Encoding.UTF8.GetString(buffer, 0, received.Count)); } while (!received.EndOfMessage); if (JsonNode.Parse(text.ToString()) is JsonObject message && wanted(message)) return message; } } catch (OperationCanceledException) { return null; } } static string Id(JsonObject message) => JsonNode.Parse(message["payload"]!.GetValue())!["id"]!.GetValue(); [Fact] public async Task Health_is_answered_and_a_stream_needs_a_persona() { var health = await _host.Client().GetAsync("/api/v1/streaming/health", Token); Assert.Equal("OK", await health.Content.ReadAsStringAsync(Token)); Assert.Equal(HttpStatusCode.Unauthorized, (await _host.Client().GetAsync("/api/v1/streaming/public", Token)).StatusCode); var instance = (await _host.Client().Get("/api/v2/instance")).Ok().Body; Assert.StartsWith("ws", instance["configuration"]!["urls"]!["streaming"]!.GetValue()); } [Fact] public async Task The_user_stream_tells_home_posts_notifications_edits_and_deletions() { var alice = await _host.Mastodon("alice"); var bob = await _host.Mastodon("bob"); var carol = await _host.Mastodon("carol"); var stranger = await _host.Mastodon("stranger"); (await alice.Client.Post($"/api/v1/accounts/{bob.Id}/follow")).Ok(); using var socket = await Connect(alice, "&stream=user"); using var elsewhere = await Connect(stranger, "&stream=user"); var post = (await bob.Status("streamed to alice")).Text("id"); var update = await Next(socket, m => m.Text("event") == "update"); Assert.NotNull(update); Assert.Equal(new[] { "user" }, update!["stream"]!.AsArray().Select(s => s!.GetValue())); Assert.Equal(post, Id(update)); await carol.Status($"@{alice.UserName} hello there", ("visibility", "direct")); var notification = await Next(socket, m => m.Text("event") == "notification"); Assert.Equal("mention", JsonNode.Parse(notification!["payload"]!.GetValue())!.Text("type")); (await bob.Client.Put($"/api/v1/statuses/{post}", ("status", "streamed, then edited"))).Ok(); var edited = await Next(socket, m => m.Text("event") == "status.update"); Assert.Contains("then edited", JsonNode.Parse(edited!["payload"]!.GetValue())!.Text("content")); (await bob.Client.Delete($"/api/v1/statuses/{post}")).Ok(); Assert.NotNull(await Next(socket, m => m.Text("event") == "delete" && m["payload"]!.GetValue() == post)); // a persona whose home never held it hears nothing of it, not even its id Assert.Null(await Next(elsewhere, m => m.Text("event") is "update" or "delete" && m["payload"]!.GetValue().Contains(post), seconds: 2)); } [Fact] public async Task A_hashtag_stream_subscribed_by_message_tells_public_posts_only_and_never_a_blocked_authors() { var alice = await _host.Mastodon("alice"); var bob = await _host.Mastodon("bob"); var dave = await _host.Mastodon("dave"); (await alice.Client.Post($"/api/v1/accounts/{dave.Id}/block")).Ok(); var tag = $"st{Guid.NewGuid():N}"[..12]; using var socket = await Connect(alice); await Say(socket, new JsonObject { ["type"] = "subscribe", ["stream"] = "hashtag", ["tag"] = tag }); await Task.Delay(200, Token); await dave.Status($"blocked #{tag}"); await bob.Status($"quiet #{tag}", ("visibility", "unlisted")); var open = (await bob.Status($"open #{tag}")).Text("id"); var told = await Next(socket, m => m.Text("event") == "update"); Assert.Equal(open, Id(told!)); Assert.Equal(new[] { "hashtag", tag }, told["stream"]!.AsArray().Select(s => s!.GetValue())); Assert.Null(await Next(socket, m => m.Text("event") == "update", seconds: 2)); await Say(socket, new JsonObject { ["type"] = "unsubscribe", ["stream"] = "hashtag", ["tag"] = tag }); await Task.Delay(200, Token); await bob.Status($"again #{tag}"); Assert.Null(await Next(socket, m => m.Text("event") == "update", seconds: 2)); await Say(socket, new JsonObject { ["type"] = "subscribe", ["stream"] = "list", ["list"] = "000000000000000000000000" }); Assert.Equal("Unknown list", (await Next(socket, m => m["error"] != null))!.Text("error")); } [Fact] public async Task Server_sent_events_tell_a_personas_notifications() { var alice = await _host.Mastodon("alice"); var bob = await _host.Mastodon("bob"); using var request = new HttpRequestMessage(HttpMethod.Get, "/api/v1/streaming/user/notification"); request.Headers.Authorization = new System.Net.Http.Headers.AuthenticationHeaderValue("Bearer", alice.Token); using var response = await _host.Client().SendAsync(request, HttpCompletionOption.ResponseHeadersRead, Token); Assert.Equal("text/event-stream", response.Content.Headers.ContentType!.MediaType); using var reader = new StreamReader(await response.Content.ReadAsStreamAsync(Token)); Assert.Equal(":)", await reader.ReadLineAsync(Token)); (await bob.Client.Post($"/api/v1/accounts/{alice.Id}/follow")).Ok(); using var timeout = CancellationTokenSource.CreateLinkedTokenSource(Token); timeout.CancelAfter(TimeSpan.FromSeconds(5)); string line; while ((line = await reader.ReadLineAsync(timeout.Token)) != null && line != "event: notification") { } Assert.Equal("event: notification", line); var data = await reader.ReadLineAsync(timeout.Token); Assert.Equal("follow", JsonNode.Parse(data!["data: ".Length..])!.Text("type")); } } }