Mastodon's streaming API: /api/v1/streaming as a WebSocket (streams
subscribed in the URL or by message) and /api/v1/streaming/{stream} as
server-sent events, with health and the URL advertised. The user stream
tells posts reaching the persona's home (not those an exclusive list keeps
apart, which its list stream tells), notifications, edits and deletions;
public, hashtag and list streams tell what belongs in them. An in-process
hub carries ids only; each connection maps a post or a notification for its
own persona as it sends it, so nothing it may not see, or whose author it
blocked or muted, goes out. A deletion reaches only the streams that showed
the post. The token comes as access_token, header or WebSocket protocol.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
108 lines
3.9 KiB
C#
108 lines
3.9 KiB
C#
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<Key, byte> _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<Key> Keys => _keys.Keys;
|
|
|
|
public void Tell(Key key, Event e) => Queue.Writer.TryWrite((key, e));
|
|
}
|
|
|
|
static readonly ConcurrentDictionary<Listener, byte> 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;
|
|
}
|
|
}
|