diff --git a/PrivaPub.Tests/Http/MastodonInstanceTests.cs b/PrivaPub.Tests/Http/MastodonInstanceTests.cs index 4f16299..3231071 100644 --- a/PrivaPub.Tests/Http/MastodonInstanceTests.cs +++ b/PrivaPub.Tests/Http/MastodonInstanceTests.cs @@ -41,7 +41,8 @@ namespace PrivaPub.Tests.Http var v2 = (await anonymous.Get("/api/v2/instance")).Ok(); Assert.Equal(InstanceController.Version, v2.Body.Text("version")); Assert.Equal(PrivaPubHost.Host, v2.Body.Text("domain")); - Assert.Null(v2.Body["configuration"]!["urls"]!["streaming"]);//nothing streams yet, so nothing is advertised + Assert.Equal($"wss://{PrivaPubHost.Host}", v2.Body["configuration"]!["urls"]!.Text("streaming")); + Assert.Equal($"wss://{PrivaPubHost.Host}", v1.Body["urls"]!.Text("streaming_api")); Assert.Equal(PrivaPub.Api.Mastodon.Controllers.StatusesController.MaxPins, v2.Body["configuration"]!["accounts"].Number("max_pinned_statuses")); Assert.Contains("image/avif", v2.Body["configuration"]!["media_attachments"]!["supported_mime_types"]!.AsArray().Select(t => t!.GetValue())); Assert.True(v2.Body["registrations"].Flag("enabled")); diff --git a/PrivaPub.Tests/Http/MastodonStreamingTests.cs b/PrivaPub.Tests/Http/MastodonStreamingTests.cs new file mode 100644 index 0000000..2dd0011 --- /dev/null +++ b/PrivaPub.Tests/Http/MastodonStreamingTests.cs @@ -0,0 +1,165 @@ +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")); + } + } +} diff --git a/PrivaPub/Api/Mastodon/Controllers/InstanceController.cs b/PrivaPub/Api/Mastodon/Controllers/InstanceController.cs index eb08e5f..3ecb7df 100644 --- a/PrivaPub/Api/Mastodon/Controllers/InstanceController.cs +++ b/PrivaPub/Api/Mastodon/Controllers/InstanceController.cs @@ -40,6 +40,10 @@ namespace PrivaPub.Api.Mastodon.Controllers string Domain => new Uri(_localActors.BaseAddress).Authority; + // where clients open their streams (StreamingController): this server, as ws or wss + string Streaming => new UriBuilder(_localActors.BaseAddress) { Scheme = _localActors.BaseAddress.StartsWith("https", StringComparison.OrdinalIgnoreCase) ? "wss" : "ws" } + .Uri.GetLeftPart(UriPartial.Authority); + // every limit is what the server enforces, read from where it is enforced static object Statuses => new { max_characters = PrivaPub.Domain.Statuses.StatusService.MaxCharacters, max_media_attachments = 4, characters_reserved_per_url = 23 }; @@ -82,7 +86,7 @@ namespace PrivaPub.Api.Mastodon.Controllers description = Description, email = string.Empty, version = Version, - urls = new { },//no streaming API yet, so none is advertised + urls = new { streaming_api = Streaming }, stats = new { user_count = users, status_count = statuses, domain_count = domains }, thumbnail = $"{_localActors.BaseAddress}/media/missing-header.png", languages = LanguageCodes(HttpContext.RequestServices), @@ -115,7 +119,7 @@ namespace PrivaPub.Api.Mastodon.Controllers languages = LanguageCodes(HttpContext.RequestServices), configuration = new { - urls = new { },//no streaming API yet, so none is advertised + urls = new { streaming = Streaming }, accounts = new { max_featured_tags = 0, max_pinned_statuses = StatusesController.MaxPins }, statuses = Statuses, media_attachments = MediaAttachments, diff --git a/PrivaPub/Api/Mastodon/Controllers/StreamingController.cs b/PrivaPub/Api/Mastodon/Controllers/StreamingController.cs new file mode 100644 index 0000000..ded56d5 --- /dev/null +++ b/PrivaPub/Api/Mastodon/Controllers/StreamingController.cs @@ -0,0 +1,223 @@ +using System.Net.WebSockets; +using System.Text; +using System.Text.Json; +using System.Text.Json.Nodes; + +using Microsoft.AspNetCore.Mvc; + +using MongoDB.Entities; + +using PrivaPub.Api.Mastodon.Infrastructure; +using PrivaPub.Api.Mastodon.Mappers; +using PrivaPub.Domain.Privacy; +using PrivaPub.Domain.Timelines; +using PrivaPub.Models.Social; +using PrivaPub.StaticServices; + +namespace PrivaPub.Api.Mastodon.Controllers +{ + // Mastodon's streaming API: one WebSocket on /api/v1/streaming that subscribes to streams by message, or one stream + // per request as server-sent events on /api/v1/streaming/{stream}. Every stream needs a signed-in persona (the token + // comes as access_token, in the Authorization header, or as the WebSocket's protocol). What a stream tells is mapped + // for that persona as it is sent: a post it may not see, or whose author it blocked or muted, is never sent. + public class StreamingController : MastodonController + { + static readonly TimeSpan KeepAlive = TimeSpan.FromSeconds(15); + + readonly MastodonMapper _mapper; + readonly DbEntities _dbEntities; + + public StreamingController(MastodonMapper mapper, DbEntities dbEntities) + { + _mapper = mapper; + _dbEntities = dbEntities; + } + + [HttpGet("/api/v1/streaming/health"), Microsoft.AspNetCore.Authorization.AllowAnonymous] + public IActionResult Health() => Content("OK", "text/plain"); + + // the stream a name and its parameter make, or why not: a hashtag needs its tag, a list must be the persona's own + async Task<(Streams.Key Key, string Error)> Resolve(string name, string tag, string list, CancellationToken token) + { + switch (name) + { + case Streams.User or Streams.UserNotification or Streams.Public or Streams.PublicLocal or Streams.PublicRemote: + return (new Streams.Key(name), default); + case Streams.Hashtag or Streams.HashtagLocal: + var normal = TagsController.Normalise(tag); + return normal.Length == 0 ? (default, "Missing tag name parameter") : (new Streams.Key(name, normal), default); + case Streams.List: + return !string.IsNullOrEmpty(list) && await DB.Default.Find().Match(l => l.ID == list && l.AvatarId == MyId).ExecuteAnyAsync(token) + ? (new Streams.Key(name, list), default) + : (default, "Unknown list"); + default: + return (default, "Unknown stream type"); + } + } + + // what the event says, for this persona: a status or a notification as JSON, a deleted id as itself; null when it + // is not the persona's to see + async Task Payload(Streams.Event e, CancellationToken token) + { + if (e.Name == "delete") + return e.Id; + if (e.Name == "notification") + { + var notification = await _dbEntities.Notifications.Match(n => n.ID == e.Id && n.AvatarId == MyId).ExecuteFirstAsync(token); + var mapped = notification == default ? default : (await _mapper.Notifications(new List { notification }, MyId, token)).FirstOrDefault(); + return mapped == default ? default : JsonSerializer.Serialize(mapped, MastodonJson.Options); + } + var post = await _dbEntities.Posts.MatchID(e.Id).ExecuteFirstAsync(token); + if (post == default || !await VisibilityPolicy.CanSee(post, MyId, token)) + return default; + var status = (await _mapper.Statuses(new[] { post }, MyId, token)).FirstOrDefault(); + return status == default ? default : JsonSerializer.Serialize(status, MastodonJson.Options); + } + + [HttpGet("/api/v1/streaming"), Scope("read")] + public async Task Socket(CancellationToken token) + { + if (!HttpContext.WebSockets.IsWebSocketRequest) + return Error(StatusCodes.Status400BadRequest, "This is a WebSocket endpoint; server-sent events are at /api/v1/streaming/{stream}"); + var protocol = HttpContext.WebSockets.WebSocketRequestedProtocols.FirstOrDefault(); + using var socket = await HttpContext.WebSockets.AcceptWebSocketAsync(protocol); + var listener = Streams.Open(MyId); + var sending = new SemaphoreSlim(1, 1); + async Task Send(JsonObject message) + { + var bytes = Encoding.UTF8.GetBytes(message.ToJsonString()); + await sending.WaitAsync(token); + try + { + await socket.SendAsync(bytes, WebSocketMessageType.Text, true, token); + } + finally + { + sending.Release(); + } + } + async Task Subscribe(string type, string name, string tag, string list) + { + var (key, error) = await Resolve(name, tag, list, token); + if (error != default) + { + await Send(new JsonObject { ["error"] = error, ["status"] = 400 }); + return; + } + if (type == "unsubscribe") + listener.Stop(key); + else + listener.Listen(key); + } + try + { + if (Params.Get("stream") is { Length: > 0 } first) + await Subscribe("subscribe", first, Params.Get("tag"), Params.Get("list")); + using var done = CancellationTokenSource.CreateLinkedTokenSource(token); + var reading = Task.Run(async () => + { + var buffer = new byte[4096]; + while (socket.State == WebSocketState.Open && !done.IsCancellationRequested) + { + var text = new StringBuilder(); + WebSocketReceiveResult received; + do + { + received = await socket.ReceiveAsync(buffer, done.Token); + if (received.MessageType == WebSocketMessageType.Close) + { + done.Cancel(); + return; + } + text.Append(Encoding.UTF8.GetString(buffer, 0, received.Count)); + } + while (!received.EndOfMessage && text.Length < 16384); + JsonNode message; + try + { + message = JsonNode.Parse(text.ToString()); + } + catch (JsonException) + { + continue; + } + if (message?["type"]?.GetValue() is "subscribe" or "unsubscribe" && message["stream"]?.GetValue() is { } stream) + await Subscribe(message["type"]!.GetValue(), stream, message["tag"]?.GetValue(), message["list"]?.GetValue()); + } + }, done.Token); + await foreach (var (key, e) in listener.Queue.Reader.ReadAllAsync(done.Token)) + { + if (!listener.Hears(key) || await Payload(e, done.Token) is not { } payload) + continue; + await Send(new JsonObject { ["stream"] = new JsonArray(key.Wire.Select(w => (JsonNode)w).ToArray()), ["event"] = e.Name, ["payload"] = payload }); + } + await reading; + } + catch (Exception ex) when (ex is OperationCanceledException or WebSocketException) + { + } + finally + { + Streams.Close(listener); + if (socket.State == WebSocketState.Open) + await socket.CloseAsync(WebSocketCloseStatus.NormalClosure, default, CancellationToken.None); + } + return new EmptyResult(); + } + + [HttpGet("/api/v1/streaming/{*stream}"), Scope("read")] + public async Task Events(string stream, CancellationToken token) + { + var name = stream.Replace('/', ':'); + var (key, error) = await Resolve(name, Params.Get("tag"), Params.Get("list"), token); + if (error != default) + return Error(StatusCodes.Status400BadRequest, error); + Response.ContentType = "text/event-stream"; + Response.Headers.CacheControl = "no-cache"; + Response.Headers["X-Accel-Buffering"] = "no"; + var listener = Streams.Open(MyId); + listener.Listen(key); + var writing = new SemaphoreSlim(1, 1); + async Task Write(string text) + { + await writing.WaitAsync(token); + try + { + await Response.WriteAsync(text, token); + await Response.Body.FlushAsync(token); + } + finally + { + writing.Release(); + } + } + try + { + await Write(":)\n\n"); + var beating = Task.Run(async () => + { + while (!token.IsCancellationRequested) + { + await Task.Delay(KeepAlive, token); + await Write(":thump\n\n"); + } + }, token); + await foreach (var (_, e) in listener.Queue.Reader.ReadAllAsync(token)) + { + if (await Payload(e, token) is not { } payload) + continue; + await Write($"event: {e.Name}\ndata: {payload}\n\n"); + } + await beating; + } + catch (OperationCanceledException) + { + } + finally + { + Streams.Close(listener); + } + return new EmptyResult(); + } + } +} diff --git a/PrivaPub/Api/Mastodon/Controllers/TimelinesController.cs b/PrivaPub/Api/Mastodon/Controllers/TimelinesController.cs index 5c237b1..d629e73 100644 --- a/PrivaPub/Api/Mastodon/Controllers/TimelinesController.cs +++ b/PrivaPub/Api/Mastodon/Controllers/TimelinesController.cs @@ -160,7 +160,7 @@ namespace PrivaPub.Api.Mastodon.Controllers public class NotificationsController : MastodonController { - static readonly Dictionary Names = new() + internal static readonly Dictionary Names = new() { [NotificationType.Mention] = "mention", [NotificationType.Follow] = "follow", @@ -351,29 +351,8 @@ namespace PrivaPub.Api.Mastodon.Controllers return Json(new { count = (await Map(unread, token)).Count }); } - async Task> Map(List notifications, CancellationToken token) - { - var accounts = await _mapper.Accounts(notifications.Select(n => n.FromAccountId), token); - var postIds = notifications.Where(n => n.PostId != default).Select(n => n.PostId).Distinct().ToList(); - var posts = postIds.Count == 0 - ? new List() - : await _dbEntities.Posts.Match(p => postIds.Contains(p.ID)).Match(VisibilityPolicy.IsShown).ExecuteAsync(token); - var statuses = (await _mapper.Statuses(posts, MyId, token)).ToDictionary(s => s.Id); - return notifications - .Where(n => accounts.ContainsKey(n.FromAccountId ?? string.Empty) && (n.PostId == default || statuses.ContainsKey(n.PostId))) - .Select(n => new Entities.Notification - { - Id = n.ID, - Type = Names[n.Type], - CreatedAt = MastodonJson.Time(n.CreatedAt), - GroupKey = $"ungrouped-{n.ID}", - Account = accounts[n.FromAccountId], - Status = n.PostId == default ? default : statuses[n.PostId], - Emoji = n.Emoji == default ? default : n.EmojiURL == default ? n.Emoji : $":{n.Emoji}:", - EmojiUrl = _mapper.ProxiedUrl(n.EmojiURL) - }) - .ToList(); - } + Task> Map(List notifications, CancellationToken token) => + _mapper.Notifications(notifications, MyId, token); static NotificationType? Parse(string name) => Names.FirstOrDefault(n => n.Value == name) is { Value: not null } pair ? pair.Key : (NotificationType?)null; } diff --git a/PrivaPub/Api/Mastodon/Mappers/MastodonMapper.cs b/PrivaPub/Api/Mastodon/Mappers/MastodonMapper.cs index ccffd91..be5ebfc 100644 --- a/PrivaPub/Api/Mastodon/Mappers/MastodonMapper.cs +++ b/PrivaPub/Api/Mastodon/Mappers/MastodonMapper.cs @@ -40,6 +40,31 @@ namespace PrivaPub.Api.Mastodon.Mappers public string ProxiedUrl(string url) => Proxied(url); + // notifications as the persona sees them: one whose account or post it can no longer see is left out + public async Task> Notifications(List notifications, string viewerId, CancellationToken token) + { + var accounts = await Accounts(notifications.Select(n => n.FromAccountId), token); + var postIds = notifications.Where(n => n.PostId != default).Select(n => n.PostId).Distinct().ToList(); + var posts = postIds.Count == 0 + ? new List() + : await _dbEntities.Posts.Match(p => postIds.Contains(p.ID)).Match(VisibilityPolicy.IsShown).ExecuteAsync(token); + var statuses = (await Statuses(posts, viewerId, token)).ToDictionary(s => s.Id); + return notifications + .Where(n => accounts.ContainsKey(n.FromAccountId ?? string.Empty) && (n.PostId == default || statuses.ContainsKey(n.PostId))) + .Select(n => new Entities.Notification + { + Id = n.ID, + Type = Controllers.NotificationsController.Names[n.Type], + CreatedAt = MastodonJson.Time(n.CreatedAt), + GroupKey = $"ungrouped-{n.ID}", + Account = accounts[n.FromAccountId], + Status = n.PostId == default ? default : statuses[n.PostId], + Emoji = n.Emoji == default ? default : n.EmojiURL == default ? n.Emoji : $":{n.Emoji}:", + EmojiUrl = Proxied(n.EmojiURL) + }) + .ToList(); + } + string MissingAvatar => $"{_localActors.BaseAddress}/media/missing-avatar.png"; string MissingHeader => $"{_localActors.BaseAddress}/media/missing-header.png"; diff --git a/PrivaPub/Domain/Social/Notifications.cs b/PrivaPub/Domain/Social/Notifications.cs index 68afeea..12ad1a8 100644 --- a/PrivaPub/Domain/Social/Notifications.cs +++ b/PrivaPub/Domain/Social/Notifications.cs @@ -16,7 +16,7 @@ namespace PrivaPub.Domain.Social return; try { - await DB.Default.SaveAsync(new Notification + var notification = new Notification { AvatarId = avatarId, Type = type, @@ -26,7 +26,11 @@ namespace PrivaPub.Domain.Social Emoji = emoji, EmojiURL = emojiUrl, DedupeKey = emoji == default ? $"{type}|{avatarId}|{fromActorUri}|{postId}" : $"{type}|{avatarId}|{fromActorUri}|{postId}|{emoji}" - }, token); + }; + await DB.Default.SaveAsync(notification, token); + var told = new Timelines.Streams.Event("notification", notification.ID); + Timelines.Streams.ToPersona(avatarId, new Timelines.Streams.Key(Timelines.Streams.User), told); + Timelines.Streams.ToPersona(avatarId, new Timelines.Streams.Key(Timelines.Streams.UserNotification), told); } catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey) { diff --git a/PrivaPub/Domain/Statuses/StatusService.cs b/PrivaPub/Domain/Statuses/StatusService.cs index 47c8c61..433b5df 100644 --- a/PrivaPub/Domain/Statuses/StatusService.cs +++ b/PrivaPub/Domain/Statuses/StatusService.cs @@ -310,6 +310,7 @@ namespace PrivaPub.Domain.Statuses post.EditedAt = DateTime.UtcNow; post.UpdateDate = post.EditedAt; await DB.Default.SaveAsync(post, token); + await Fanout.Edited(post, token); if (!post.IsLocalOnly) { @@ -354,6 +355,7 @@ namespace PrivaPub.Domain.Statuses .Modify(p => p.Media, new List()) .Modify(p => p.Revisions, new List()) .ExecuteAsync(token); + await Fanout.Deleting(post, token); await DB.Default.DeleteAsync(e => e.PostId == post.ID || e.ReblogOfPostId == post.ID); // boosts of it end with it, as on Mastodon, so no count keeps them await DB.Default.Update().Match(p => p.ReblogOfPostId == post.ID && !p.DeletedAt.HasValue) @@ -493,6 +495,7 @@ namespace PrivaPub.Domain.Statuses if (existing == default) return new StatusOutcome(original); + await Fanout.Deleting(existing, token); await DB.Default.DeleteAsync(existing.ID); await DB.Default.DeleteAsync(e => e.PostId == existing.ID); await DB.Default.Update().MatchID(original.ID).Modify(b => b.Inc(p => p.ReblogsCount, -1)).ExecuteAsync(token); diff --git a/PrivaPub/Domain/Timelines/Fanout.cs b/PrivaPub/Domain/Timelines/Fanout.cs index 5373dc6..fa0682d 100644 --- a/PrivaPub/Domain/Timelines/Fanout.cs +++ b/PrivaPub/Domain/Timelines/Fanout.cs @@ -85,6 +85,90 @@ namespace PrivaPub.Domain.Timelines if (string.IsNullOrEmpty(post.ReblogOfPostId)) foreach (var mention in post.Mentions.Where(m => m.IsLocal && !string.IsNullOrEmpty(m.AccountId))) await Notifications.Add(mention.AccountId, NotificationType.Mention, post.AuthorAccountId ?? post.GroupUserId, post.ActorURI, post.ID, token); + if (Streams.Anyone) + await Stream(post, recipients, token); + } + + // the post to the streams it belongs in: the homes it reached (but not where an exclusive list keeps its author + // apart), the lists of those homes that hold its author, and the public feeds and hashtags when it is public + static async Task Stream(Post post, IReadOnlyCollection recipients, CancellationToken token) + { + var update = new Streams.Event("update", post.ID); + var author = post.AuthorAccountId ?? post.GroupUserId; + var listening = recipients.Where(Streams.Listening).ToList(); + if (listening.Count > 0) + { + var memberships = await DB.Default.Find().Match(m => listening.Contains(m.AvatarId) && m.AccountId == author).ExecuteAsync(token); + var listIds = memberships.Select(m => m.ListId).Distinct().ToList(); + var exclusive = listIds.Count == 0 + ? new HashSet() + : (await DB.Default.Find().Match(l => listIds.Contains(l.ID) && l.Exclusive).ExecuteAsync(token)).Select(l => l.ID).ToHashSet(); + foreach (var avatarId in listening) + { + var lists = memberships.Where(m => m.AvatarId == avatarId).Select(m => m.ListId).ToList(); + if (!lists.Any(exclusive.Contains)) + Streams.ToPersona(avatarId, new Streams.Key(Streams.User), update); + foreach (var listId in lists) + Streams.ToPersona(avatarId, new Streams.Key(Streams.List, listId), update); + } + } + if (post.Visibility != PostVisibility.Public || !string.IsNullOrEmpty(post.ReblogOfPostId) || !Privacy.VisibilityPolicy.Shown(post)) + return; + var local = !post.IsFederatedCopy; + if (!post.IsLocalOnly) + Streams.ToAll(new Streams.Key(Streams.Public), update); + Streams.ToAll(new Streams.Key(local ? Streams.PublicLocal : Streams.PublicRemote), update); + foreach (var tag in post.Tags.Distinct()) + { + if (!post.IsLocalOnly) + Streams.ToAll(new Streams.Key(Streams.Hashtag, tag), update); + if (local) + Streams.ToAll(new Streams.Key(Streams.HashtagLocal, tag), update); + } + } + + // a post's deletion, told before its home entries go: to the homes holding it or a boost of it (the boost's id + // there), to its author's own connections, and to the public feeds and hashtags when it was public. Nobody else + // hears of it, not even its id. + public static async Task Deleting(Post post, CancellationToken token) + { + if (!Streams.Anyone || post == default) + return; + var entries = await DB.Default.Find().Match(e => e.PostId == post.ID || e.ReblogOfPostId == post.ID).ExecuteAsync(token); + foreach (var entry in entries.Where(e => Streams.Listening(e.AvatarId))) + Streams.ToPersonaEverywhere(entry.AvatarId, new Streams.Event("delete", entry.PostId)); + var author = post.AuthorAccountId ?? post.GroupUserId; + if (!post.IsFederatedCopy && !string.IsNullOrEmpty(author)) + Streams.ToPersonaEverywhere(author, new Streams.Event("delete", post.ID)); + if (post.Visibility != PostVisibility.Public || !string.IsNullOrEmpty(post.ReblogOfPostId)) + return; + var gone = new Streams.Event("delete", post.ID); + var local = !post.IsFederatedCopy; + Streams.ToAll(new Streams.Key(Streams.Public), gone); + Streams.ToAll(new Streams.Key(local ? Streams.PublicLocal : Streams.PublicRemote), gone); + foreach (var tag in post.Tags.Distinct()) + { + Streams.ToAll(new Streams.Key(Streams.Hashtag, tag), gone); + if (local) + Streams.ToAll(new Streams.Key(Streams.HashtagLocal, tag), gone); + } + } + + // an edit to the homes holding the post and, when it is public, to the public feeds (status.update) + public static async Task Edited(Post post, CancellationToken token) + { + if (!Streams.Anyone) + return; + var update = new Streams.Event("status.update", post.ID); + var homes = (await DB.Default.Find().Match(e => e.PostId == post.ID).ExecuteAsync(token)).Select(e => e.AvatarId).Distinct(); + foreach (var avatarId in homes.Where(Streams.Listening)) + Streams.ToPersona(avatarId, new Streams.Key(Streams.User), update); + if (post.Visibility == PostVisibility.Public && Privacy.VisibilityPolicy.Shown(post)) + { + if (!post.IsLocalOnly) + Streams.ToAll(new Streams.Key(Streams.Public), update); + Streams.ToAll(new Streams.Key(post.IsFederatedCopy ? Streams.PublicRemote : Streams.PublicLocal), update); + } } async Task Shows(Following following, Post post, CancellationToken token) diff --git a/PrivaPub/Domain/Timelines/Streams.cs b/PrivaPub/Domain/Timelines/Streams.cs new file mode 100644 index 0000000..bdef98d --- /dev/null +++ b/PrivaPub/Domain/Timelines/Streams.cs @@ -0,0 +1,107 @@ +using System.Collections.Concurrent; +using System.Threading.Channels; + +namespace PrivaPub.Domain.Timelines +{ + // Mastodon's streaming, inside this one process: what happens (a post arriving in someone's home or a public feed, a + // notification, an edit, a deletion) is told to the connections listening for it, by id only. Each connection maps a + // post or a notification for its own persona when it sends it (StreamingController), so blocks, mutes and visibility + // are decided there, never here. Nothing is kept: a connection that is gone misses what happened meanwhile, and a + // client reads the timeline again when it comes back. + public static class Streams + { + public const string User = "user"; + public const string UserNotification = "user:notification"; + public const string Public = "public"; + public const string PublicLocal = "public:local"; + public const string PublicRemote = "public:remote"; + public const string Hashtag = "hashtag"; + public const string HashtagLocal = "hashtag:local"; + public const string List = "list"; + public const string Direct = "direct"; + + public static readonly string[] Names = { User, UserNotification, Public, PublicLocal, PublicRemote, Hashtag, HashtagLocal, List, Direct }; + + // what happened: Mastodon's event name, and the id of what it is about + public sealed record Event(string Name, string Id); + + // one stream a connection listens to: its name and, for a hashtag or a list, which one + public sealed record Key(string Name, string Parameter = default) + { + public string[] Wire => Parameter == default ? new[] { Name } : new[] { Name, Parameter }; + } + + public sealed class Listener + { + readonly ConcurrentDictionary _keys = new(); + + public Listener(string avatarId) + { + AvatarId = avatarId; + } + + public string AvatarId { get; } + + // what reached it and was not sent yet; the oldest goes first when a slow connection falls behind + public Channel<(Key Key, Event Event)> Queue { get; } = + Channel.CreateBounded<(Key, Event)>(new BoundedChannelOptions(1000) { FullMode = BoundedChannelFullMode.DropOldest, SingleReader = true }); + + public void Listen(Key key) => _keys[key] = 0; + + public void Stop(Key key) => _keys.TryRemove(key, out _); + + public bool Hears(Key key) => _keys.ContainsKey(key); + + public IEnumerable Keys => _keys.Keys; + + public void Tell(Key key, Event e) => Queue.Writer.TryWrite((key, e)); + } + + static readonly ConcurrentDictionary Listeners = new(); + + public static Listener Open(string avatarId) + { + var listener = new Listener(avatarId); + Listeners[listener] = 0; + return listener; + } + + public static void Close(Listener listener) + { + Listeners.TryRemove(listener, out _); + listener.Queue.Writer.TryComplete(); + } + + public static int Count => Listeners.Count; + + // whether the persona has a connection open, so the work of finding its lists can be skipped when it has none + public static bool Listening(string avatarId) => Listeners.Keys.Any(l => l.AvatarId == avatarId); + + // to the persona's own streams: its home (user) and, for a notification, user:notification too + public static void ToPersona(string avatarId, Key key, Event e) + { + foreach (var listener in Listeners.Keys) + if (listener.AvatarId == avatarId && listener.Hears(key)) + listener.Tell(key, e); + } + + // to everyone listening to a shared stream (public, a hashtag) + public static void ToAll(Key key, Event e) + { + foreach (var listener in Listeners.Keys) + if (listener.Hears(key)) + listener.Tell(key, e); + } + + // to every stream the persona listens to (a deletion: wherever its client may show the post) + public static void ToPersonaEverywhere(string avatarId, Event e) + { + foreach (var listener in Listeners.Keys) + if (listener.AvatarId == avatarId) + foreach (var key in listener.Keys) + listener.Tell(key, e); + } + + public static bool Anyone => !Listeners.IsEmpty; + } +} diff --git a/PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs b/PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs index c5545f8..01e3b83 100644 --- a/PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs +++ b/PrivaPub/Federation/Inbox/Handlers/UndoHandler.cs @@ -126,6 +126,7 @@ namespace PrivaPub.Federation.Inbox.Handlers .ExecuteFirstAsync(token); if (reblog == default) return false; + await Domain.Timelines.Fanout.Deleting(reblog, token); await DB.Default.DeleteAsync(reblog.ID); await DB.Default.DeleteAsync(e => e.PostId == reblog.ID); await DB.Default.Update().MatchID(reblog.ReblogOfPostId).Modify(b => b.Inc(p => p.ReblogsCount, -1)).ExecuteAsync(token); diff --git a/PrivaPub/Federation/Inbox/RemoteEdits.cs b/PrivaPub/Federation/Inbox/RemoteEdits.cs index 8fba836..afb9327 100644 --- a/PrivaPub/Federation/Inbox/RemoteEdits.cs +++ b/PrivaPub/Federation/Inbox/RemoteEdits.cs @@ -73,6 +73,7 @@ namespace PrivaPub.Federation.Inbox post.EditedAt = note.Updated ?? DateTime.UtcNow; post.UpdateDate = DateTime.UtcNow; await DB.Default.SaveAsync(post, token); + await Domain.Timelines.Fanout.Edited(post, token); await records.Revise(note, activityId, token); await Requote(post, note, quotes, token); return true; @@ -97,6 +98,7 @@ namespace PrivaPub.Federation.Inbox public static async Task Remove(PostEntity post, string objectUri, CancellationToken token) { await Tombstone(objectUri, token); + await Domain.Timelines.Fanout.Deleting(post, token); await DB.Default.DeleteAsync(post.ID); await DB.Default.DeleteAsync(r => r.PostId == post.ID || r.ObjectURI == objectUri); await DB.Default.DeleteAsync(e => e.PostId == post.ID || e.ReblogOfPostId == post.ID); diff --git a/PrivaPub/Program.cs b/PrivaPub/Program.cs index 9382177..e63f31e 100644 --- a/PrivaPub/Program.cs +++ b/PrivaPub/Program.cs @@ -154,10 +154,25 @@ try app.UseRequestLocalization(await localizationService.Get()); + app.UseWebSockets(new WebSocketOptions { KeepAliveInterval = TimeSpan.FromSeconds(30) }); app.UseRouting(); app.UseTrafficMeter(); app.UseRateLimiter(); + // a stream's token may come as access_token or as the WebSocket's protocol, as Mastodon's clients send it + app.Use(async (context, next) => + { + if (context.Request.Path.StartsWithSegments("/api/v1/streaming") && !context.Request.Headers.ContainsKey("Authorization")) + { + var token = context.Request.Query["access_token"].ToString(); + if (string.IsNullOrEmpty(token)) + token = context.Request.Headers["Sec-WebSocket-Protocol"].ToString(); + if (!string.IsNullOrEmpty(token)) + context.Request.Headers.Authorization = "Bearer " + token; + } + await next(); + }); + app.UseAuthentication(); app.UseAuthorization(); //app.UseWhen(context => context.Request.Path.StartsWithSegments("/peasants") || diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index 26a48fd..ac1d9da 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -470,6 +470,10 @@ The first refactor commit is a pure move with namespaces only. Logic changes fol - **Scheduled statuses:** `scheduled_at` on posting, `scheduled_statuses` (P9). - **Tags:** `tags/:name` with history, follow and unfollow, `followed_tags` (P9). - **Trends and directory:** `trends/tags`, `trends/statuses`, `trends/links`, `directory` (P9). + - **Streaming:** `/api/v1/streaming` (WebSocket, subscribe and unsubscribe messages) and `/api/v1/streaming/{stream}` + (server-sent events): `user`, `user:notification`, `public`, `public:local`, `public:remote`, `hashtag`, + `hashtag:local`, `list`; events `update`, `notification`, `status.update`, `delete`; `health`. Advertised as + `urls.streaming`. - **Stubs:** `custom_emojis`, announcements, suggestions, preferences. - **Mapping:** - Avatar, ForeignAvatar and Group all map to Account (`group: true` for groups). `acct` uses a WebFinger-verified @@ -634,7 +638,10 @@ it, raw where it doesn't. - a relay client for both relay styles; - instance actor discovery (FEP-d556, FEP-2677); - `implements` (FEP-844e). -- **Mastodon API:** streaming WebSocket, Web Push, grouped notifications. Then advertise an honest version. +- **Mastodon API:** ~~streaming WebSocket~~ (done 2026-10-05: `/api/v1/streaming` as a WebSocket and as server-sent + events, user, notification, public, hashtag and list streams, each event mapped for its persona; a deletion reaches + only the streams that showed the post), Web Push (gated: outbound traffic to push services), grouped notifications. + Then advertise an honest version. - **Backfill:** an author's outbox after following them. - **Long tail:** - MFM rendering data;