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; } }