diff --git a/.gitignore b/.gitignore index c4a27a5..b038b6f 100644 --- a/.gitignore +++ b/.gitignore @@ -402,6 +402,7 @@ FodyWeavers.xsd PrivaPub/media-store/ PrivaPub/media-store-proxy/ tools/pasture/.publish/ +tools/pasture/.flood/ .claude/worktrees/ tools/pasture/.ca/ tools/pasture/.state/ diff --git a/CLAUDE.md b/CLAUDE.md index 941e333..e0ce4ff 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -553,6 +553,10 @@ tools/pasture/run.sh down # removes e `connection`) and its API is v3 (`/api/v3`, `sort=New`, `resolve_object` answering views). It takes private messages only as `ChatMessage`: a reply goes out as one, a first message cannot (G-0008). `scenarios/lemmy19.sh`, 29 checks and that gap. +- **Load (`load.sh`, needs the `flood` peer):** `flood` (`flood/flood.cs`, published once into `.flood`) answers as + twenty fake servers and sends signed activities at a set rate; `load.sh --rate=N --seconds=N` measures the answers, + the queue's wait and processing times, its drain, and a persona's home timeline meanwhile, and keeps each run in + `out/load/`. `docs/LOAD.md` has the method and the runs. - **Account moves (`scenarios/moves.sh`, needs gts):** two fresh GoToSocial accounts made by its admin CLI; alice follows the old one, the new one names it as an alias (`/api/v1/accounts/alias`), the old one moves (`/api/v1/accounts/move`), and PrivaPub shows it `moved` while alice's follow stays. 6 checks. diff --git a/PrivaPub/Federation/Inbox/InboxProcessor.cs b/PrivaPub/Federation/Inbox/InboxProcessor.cs index e1b97cf..0d1789a 100644 --- a/PrivaPub/Federation/Inbox/InboxProcessor.cs +++ b/PrivaPub/Federation/Inbox/InboxProcessor.cs @@ -28,9 +28,13 @@ namespace PrivaPub.Federation.Inbox readonly IInteractionLedger _ledger; readonly ILocalActorService _localActors; + readonly int _concurrency; + public InboxProcessor(IRemoteActorService remoteActors, IEnumerable handlers, ILogger logger, - IInteractionLedger ledger = default, ILocalActorService localActors = default) + IInteractionLedger ledger = default, ILocalActorService localActors = default, + Microsoft.Extensions.Options.IOptions federation = default) { + _concurrency = Math.Max(1, federation?.Value.InboxConcurrency ?? 8); _localActors = localActors; _remoteActors = remoteActors; _handlers = handlers.ToDictionary(h => h.Type, StringComparer.Ordinal); @@ -39,7 +43,7 @@ namespace PrivaPub.Federation.Inbox } public JobKind Kind => JobKind.ProcessInbox; - public int Concurrency => 2; + public int Concurrency => _concurrency; public int MaxAttempts => 8; public int PerHostLimit => 2; diff --git a/PrivaPub/Federation/Outbox/DeliveryService.cs b/PrivaPub/Federation/Outbox/DeliveryService.cs index 7040608..8fad02f 100644 --- a/PrivaPub/Federation/Outbox/DeliveryService.cs +++ b/PrivaPub/Federation/Outbox/DeliveryService.cs @@ -89,9 +89,12 @@ namespace PrivaPub.Federation.Outbox readonly ILogger _logger; readonly IInteractionLedger _ledger; + readonly int _concurrency; + public DeliveryJobHandler(ILocalActorService actors, IFederationHttp http, IHostCircuitBreaker breaker, ILogger logger, - IInteractionLedger ledger = default) + IInteractionLedger ledger = default, Microsoft.Extensions.Options.IOptions federation = default) { + _concurrency = Math.Max(1, federation?.Value.DeliveryConcurrency ?? 8); _actors = actors; _http = http; _breaker = breaker; @@ -109,7 +112,7 @@ namespace PrivaPub.Federation.Outbox } public JobKind Kind => JobKind.Deliver; - public int Concurrency => 8; + public int Concurrency => _concurrency; public int MaxAttempts => 16; public int PerHostLimit => 2; diff --git a/PrivaPub/Infrastructure/Data/Indexes.cs b/PrivaPub/Infrastructure/Data/Indexes.cs index dbcd1ce..a0d8179 100644 --- a/PrivaPub/Infrastructure/Data/Indexes.cs +++ b/PrivaPub/Infrastructure/Data/Indexes.cs @@ -18,18 +18,23 @@ namespace PrivaPub.Infrastructure.Data public static async Task Create(CancellationToken token = default) { - await Unique(p => p.ObjectURI, Builders.Filter.Type(p => p.ObjectURI, BsonType.String), token); + await Unique(p => p.ObjectURI, Builders.Filter.Gt(p => p.ObjectURI, ""), token); await Plain(token, p => p.GroupUserId, p => p.ID); await Plain(token, p => p.GroupId, p => p.ID); await Plain(token, p => p.ActorURI); await Plain(token, p => p.ConversationId, p => p.ID); await Plain(token, p => p.InReplyToURI); + // who boosted what: a persona's boosts of the posts a page shows were a scan of every post under load + // (tools/pasture/load.sh), and a post's boosts and replies are read for its counts and context + await Plain(token, p => p.AuthorAccountId, p => p.ReblogOfPostId); + await Plain(token, p => p.ReblogOfPostId); + await Plain(token, p => p.AnsweringToPostId); await DB.Default.Index().Key(p => p.Geo, KeyType.Geo2DSphere).CreateAsync(token); - await Unique(p => p.ObjectURI, Builders.Filter.Type(p => p.ObjectURI, BsonType.String), token); + await Unique(p => p.ObjectURI, Builders.Filter.Gt(p => p.ObjectURI, ""), token); await Plain(token, p => p.GroupId, p => p.ID); - await Unique(a => a.ActorURI, Builders.Filter.Type(a => a.ActorURI, BsonType.String), token); + await Unique(a => a.ActorURI, Builders.Filter.Gt(a => a.ActorURI, ""), token); await Plain(token, a => a.PublicKeyId); await DB.Default.Index() @@ -47,8 +52,8 @@ namespace PrivaPub.Infrastructure.Data .CreateAsync(token); await Plain(token, r => r.AvatarId); - await Unique(r => r.Name, Builders.Filter.Type(r => r.Name, BsonType.String), token); - await Unique(u => u.UserName, Builders.Filter.Type(u => u.UserName, BsonType.String), token); + await Unique(r => r.Name, Builders.Filter.Gt(r => r.Name, ""), token); + await Unique(u => u.UserName, Builders.Filter.Gt(u => u.UserName, ""), token); await Plain(token, a => a.UserName); await Plain(token, g => g.UserName); await Plain(token, g => g.InvitationCode); @@ -60,13 +65,13 @@ namespace PrivaPub.Infrastructure.Data .Key(s => s.ConversationId, KeyType.Ascending) .Option(o => o.Unique = true) .CreateAsync(token); - await Unique(a => a.ActivityURI, Builders.Filter.Type(a => a.ActivityURI, BsonType.String), token); + await Unique(a => a.ActivityURI, Builders.Filter.Gt(a => a.ActivityURI, ""), token); await Plain(token, d => d.DeliveredAt, d => d.AbandonedAt, d => d.NextAttemptAt); await Plain(token, j => j.Kind, j => j.State, j => j.RunAt); await Plain(token, j => j.State, j => j.LeasedUntil); - await Unique(j => j.DedupeKey, Builders.Filter.Type(j => j.DedupeKey, BsonType.String), token); + await Unique(j => j.DedupeKey, Builders.Filter.Gt(j => j.DedupeKey, ""), token); await DB.Default.Index() .Key(j => j.FinishedAt, KeyType.Ascending) .Option(o => o.ExpireAfter = TimeSpan.FromDays(7)) @@ -78,6 +83,7 @@ namespace PrivaPub.Infrastructure.Data .CreateAsync(token); await Plain(token, f => f.TargetActorURI, f => f.State); await Plain(token, f => f.TargetAccountId, f => f.State); + await Plain(token, f => f.TargetActorURI, f => f.State);//who here follows the author of what arrives await Plain(token, f => f.FollowActivityURI); await DB.Default.Index() .Key(e => e.AvatarId, KeyType.Ascending) @@ -85,7 +91,8 @@ namespace PrivaPub.Infrastructure.Data .Option(o => o.Unique = true) .CreateAsync(token); await Plain(token, e => e.PostId); - await Unique(n => n.DedupeKey, Builders.Filter.Type(n => n.DedupeKey, BsonType.String), token); + await Plain(token, e => e.AuthorAccountId); + await Unique(n => n.DedupeKey, Builders.Filter.Gt(n => n.DedupeKey, ""), token); await Plain(token, n => n.AvatarId, n => n.ID); await DB.Default.Index() .Key(f => f.AccountId, KeyType.Ascending) @@ -93,6 +100,7 @@ namespace PrivaPub.Infrastructure.Data .Option(o => o.Unique = true) .CreateAsync(token); await Plain(token, f => f.ActivityURI); + await Plain(token, f => f.PostId); await DB.Default.Index() .Key(p => p.AvatarId, KeyType.Ascending) .Key(p => p.PostId, KeyType.Ascending) @@ -117,6 +125,7 @@ namespace PrivaPub.Infrastructure.Data }) await pair.Item1(); await Plain(token, b => b.TargetActorURI); + await Plain(token, p => p.PostId); await Plain(token, b => b.ActorURI); await Plain(token, l => l.AvatarId); await Plain(token, f => f.AvatarId); @@ -130,11 +139,11 @@ namespace PrivaPub.Infrastructure.Data await Plain(token, m => m.AvatarId); await Plain(token, m => m.TargetActorURI); await Plain(token, r => r.IsResolved, r => r.ID); - await Unique(b => b.Domain, Builders.Filter.Type(b => b.Domain, BsonType.String), token); - await Unique(i => i.Host, Builders.Filter.Type(i => i.Host, BsonType.String), token); - await Unique(r => r.ObjectURI, Builders.Filter.Type(r => r.ObjectURI, BsonType.String), token); + await Unique(b => b.Domain, Builders.Filter.Gt(b => b.Domain, ""), token); + await Unique(i => i.Host, Builders.Filter.Gt(i => i.Host, ""), token); + await Unique(r => r.ObjectURI, Builders.Filter.Gt(r => r.ObjectURI, ""), token); await Plain(token, r => r.PostId); - await Unique(d => d.ObjectURI, Builders.Filter.Type(d => d.ObjectURI, BsonType.String), token); + await Unique(d => d.ObjectURI, Builders.Filter.Gt(d => d.ObjectURI, ""), token); await DB.Default.Index().Key(d => d.DeletedAt, KeyType.Ascending).Option(o => o.ExpireAfter = TombstoneLifetime).CreateAsync(token); await DB.Default.Index().Key(d => d.PostId, KeyType.Ascending).Key(d => d.ActorURI, KeyType.Ascending).Option(o => o.Unique = true).CreateAsync(token); await Plain(token, d => d.ActivityURI); @@ -145,7 +154,7 @@ namespace PrivaPub.Infrastructure.Data .Option(o => o.Unique = true).CreateAsync(token); await Plain(token, r => r.ActivityURI); await DB.Default.Index().Key(l => l.PostId, KeyType.Ascending).Key(l => l.QuotingObjectURI, KeyType.Ascending).CreateAsync(token); - await Unique(p => p.Url, Builders.Filter.Type(p => p.Url, BsonType.String), token); + await Unique(p => p.Url, Builders.Filter.Gt(p => p.Url, ""), token); await DB.Default.Index().Key(e => e.At, KeyType.Ascending).Option(o => o.ExpireAfter = InteractionEvent.Retention).CreateAsync(token); await Plain(token, e => e.Host, e => e.At); @@ -164,6 +173,10 @@ namespace PrivaPub.Infrastructure.Data await Plain(token, i => i.LastCrawledAt); } + // unique among the documents that have the key: the partial filter is `$gt: ""` (any non-empty string) rather than + // `$type: "string"`, because MongoDB uses a partial index for a lookup only when the lookup implies its filter, and an + // equality on a string implies the first, never the second: every lookup by these keys scanned the collection + // (tools/pasture/load.sh, 280 ms a lookup at 30 000 posts) static async Task Unique(System.Linq.Expressions.Expression> key, FilterDefinition partial, CancellationToken token) where T : IEntity => await DB.Default.Index() diff --git a/PrivaPub/Infrastructure/Http/FederationOptions.cs b/PrivaPub/Infrastructure/Http/FederationOptions.cs index 59ab336..7ed10ec 100644 --- a/PrivaPub/Infrastructure/Http/FederationOptions.cs +++ b/PrivaPub/Infrastructure/Http/FederationOptions.cs @@ -7,5 +7,9 @@ namespace PrivaPub.Infrastructure.Http public bool SecureMode { get; set; } public bool AcceptAnyCertificate { get; set; } public bool FetchLinkPreviews { get; set; } = true;//owner decision 1: the server reads linked pages of public posts + // how many received activities are processed at once, and deliveries made at once (each still at most two per host). + // At two the inbox kept up with about 110 activities a second (tools/pasture/load.sh) + public int InboxConcurrency { get; set; } = 8; + public int DeliveryConcurrency { get; set; } = 8; } } diff --git a/docs/LOAD.md b/docs/LOAD.md new file mode 100644 index 0000000..7c4406d --- /dev/null +++ b/docs/LOAD.md @@ -0,0 +1,49 @@ +# PrivaPub under load + +How PrivaPub holds up when the fediverse sends it a lot at once, measured in the pasture with `tools/pasture/load.sh` +against the `flood` peer, and what each run changed. Numbers are from the workstation (32 cores, PrivaPub and Mongo +8 in podman), so they compare runs with each other; they do not promise a production figure. + +## The method + +- **The crowd:** `flood` (`tools/pasture/flood/flood.cs`, a .NET file-based app) answers as twenty servers, + `flood1..20.test`, with ten actors each: their actor documents and keys, WebFinger, NodeInfo, and inboxes that accept + any Follow with a signed Accept. `flood run` sends signed activities from random actors to PrivaPub's shared inbox at + a set rate: `Create`s of public notes that mention a persona, `Like`s of the personas' posts, `Follow`s of them. +- **The audience:** `load.sh` makes the personas `load0..N` under the pasture root; each follows its own share of the + crowd (so a note reaches some homes and not others) and posts twice for the crowd to like. +- **What is measured:** + - the crowd's view: what PrivaPub answered (202, or 429 when an origin is over its limit) and how fast; + - for every activity PrivaPub took, how long it waited in the job queue and how long processing took (its + `InteractionEvent`s); + - the queue's peak of due jobs, and how long it took to drain once the crowd stopped; + - a persona's home timeline (`/api/v1/timelines/home?limit=40`), read four times a second throughout. +- Each run is printed and kept in `tools/pasture/out/load/`. + +`tools/pasture/load.sh --rate=300 --seconds=60 --personas=20 --follows=50` is the reference run below. PrivaPub lets +each sending origin burst 300 deliveries and earn back 5 a second (`RateLimits:Inbox*`), so twenty origins at 300 a +second leave about 200 a second through and are answered 429 for the rest, by design. + +## Runs (2026-10-05) + +| Run | Posts held | Queue wait p50 / p95 | Processing p50 / p95 | Queue peak, drain | Home p95 | +|---|---|---|---|---|---| +| 100/s, 2 workers | 17 000 | 12 / 17 ms | 11 / 15 ms | 4, 0 s | 18 ms | +| 300/s, 2 workers | 23 000 | 29 / 47 s | 17 / 23 ms | 5 104, 49 s | 27 ms | +| 300/s, 8 workers | 29 000 | 38 / 60 s | 81 / 123 ms | 5 786, 63 s | 119 ms | +| 300/s, 8 workers, unique indexes fixed | 45 000 | 6 / 9 ms | 6 / 8 ms | 0, 0 s | 14 ms | +| 600/s, origin limits lifted | 70 000 | 7 / 80 ms | 6 / 9 ms | 0, 0 s | 17 ms | +| 1000/s, origin limits lifted | 110 000 | 3.1 / 4.1 s | 7 / 11 ms | 3 933, 6 s | 18 ms | + +What the runs found: +- **Two inbox workers** kept up with about 110 activities a second. `Federation:InboxConcurrency` and + `Federation:DeliveryConcurrency` (default 8) now set them. +- **More workers made it worse**, because every lookup by a unique key scanned its whole collection: those indexes + were partial on `$type: "string"`, which MongoDB never uses for an equality lookup. Every post by `ObjectURI` (each + arriving `Create` checks for a duplicate), every actor by `ActorURI`, deleted objects, domain blocks, remote + instances and the rest. They are partial on `$gt: ""` now (`Indexes.Unique`), which an equality on a string implies, + and MongoDB.Entities rebuilt them in place at the next start. +- The indexes the plan listed as missing were added on the way (a post's boosts and replies, a persona's boosts, who + follows an actor, a timeline's rows by author, a post's likes and pins). +- With both fixed, eight workers process about 900 activities a second here, each in under 10 ms, and the home + timeline does not notice. Above that the queue grows and drains as soon as the burst ends. diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index bbe8315..efe3ff0 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -72,6 +72,10 @@ Written 2026-10-01 from the original 2023 code, the decePubClient UI, a federati - groups run from the client (2026-10-05): `/clientapi/group/members` shows a group's members and requests to its owner and moderators only, `reject` declines a request and `remove` takes a member out (a `Reject{Follow}` to a member elsewhere); decePubClient's Groups page makes, joins, edits and runs them. + - load (S10, 2026-10-05): the `flood` peer and `load.sh` (`docs/LOAD.md`). The first runs found every lookup by a + unique key scanning its collection (the partial indexes were `$type: "string"`, which MongoDB never uses for an + equality) and two inbox workers capping intake at about 110 activities a second; fixed, with the missing indexes, + PrivaPub processes about 900 a second here, each in under 10 ms. - wave 2, under way: PieFed, Mbin, NodeBB and Lemmy 0.19 in the pasture with scenarios (2026-10-05). What they showed and was fixed: the instance actor answers at the server's root, where PieFed looks for the inbox it announces to; a community's removal of a post on its own server is believed at once; a followed group's post its own server sends without announcing diff --git a/tools/pasture/Caddyfile b/tools/pasture/Caddyfile index bd47373..546e5ee 100644 --- a/tools/pasture/Caddyfile +++ b/tools/pasture/Caddyfile @@ -112,3 +112,8 @@ lemmy19.test { tls internal reverse_proxy pasture-lemmy19:8536 } + +flood1.test, flood2.test, flood3.test, flood4.test, flood5.test, flood6.test, flood7.test, flood8.test, flood9.test, flood10.test, flood11.test, flood12.test, flood13.test, flood14.test, flood15.test, flood16.test, flood17.test, flood18.test, flood19.test, flood20.test { + tls internal + reverse_proxy pasture-flood:8080 +} diff --git a/tools/pasture/flood/flood.cs b/tools/pasture/flood/flood.cs new file mode 100644 index 0000000..3e753a4 --- /dev/null +++ b/tools/pasture/flood/flood.cs @@ -0,0 +1,295 @@ +#:sdk Microsoft.NET.Sdk.Web +#:property PublishAot=false +#:property InvariantGlobalization=true +#:property Nullable=disable + +// flood: a crowd of fake servers for the pasture's load runs (tools/pasture/load.sh). +// flood serve answers as flood1..N.test: each host's actors (/users/u, with their keys), WebFinger, +// NodeInfo, empty collections, and inboxes that answer any Follow with a signed Accept +// flood run [options] sends signed activities from those actors to PrivaPub's shared inbox at a set rate and +// prints, as JSON, what it sent and how PrivaPub answered +// Both read the actors' keys from the state directory (made by the first `serve`), so they speak as the same actors. +using System.Collections.Concurrent; +using System.Diagnostics; +using System.Globalization; +using System.Net; +using System.Security.Cryptography; +using System.Text; +using System.Text.Json; +using System.Text.Json.Nodes; + +var state = Environment.GetEnvironmentVariable("FLOOD_STATE") ?? "/state"; +var hosts = int.Parse(Environment.GetEnvironmentVariable("FLOOD_HOSTS") ?? "20", CultureInfo.InvariantCulture); +var perHost = int.Parse(Environment.GetEnvironmentVariable("FLOOD_ACTORS") ?? "10", CultureInfo.InvariantCulture); +var keys = Keys.Load(state, hosts, perHost); + +if (args.Length > 0 && args[0] == "run") + return await Run.Start(args[1..], keys); + +var builder = WebApplication.CreateBuilder(); +builder.Logging.SetMinimumLevel(LogLevel.Warning); +builder.WebHost.UseUrls("http://0.0.0.0:8080"); +var app = builder.Build(); +var http = Signer.Client(); + +app.MapGet("/users/{name}", (HttpContext context, string name) => + keys.Find(context.Request.Host.Host, name) is { } actor ? Json(Actor(actor)) : Results.NotFound()); +app.MapGet("/users/{name}/{collection}", (HttpContext context, string name, string collection) => + keys.Find(context.Request.Host.Host, name) is { } actor && collection is "followers" or "following" or "outbox" + ? Json(new JsonObject { ["@context"] = "https://www.w3.org/ns/activitystreams", ["id"] = $"{actor.Id}/{collection}", ["type"] = "OrderedCollection", ["totalItems"] = 0 }) + : Results.NotFound()); +app.MapGet("/.well-known/webfinger", (HttpContext context, string resource) => +{ + var handle = resource?.StartsWith("acct:", StringComparison.Ordinal) == true ? resource[5..] : resource; + var parts = handle?.Split('@') ?? []; + return parts.Length == 2 && keys.Find(parts[1], parts[0]) is { } actor + ? Results.Json(new { subject = "acct:" + handle, links = new[] { new { rel = "self", type = "application/activity+json", href = actor.Id } } }, + contentType: "application/jrd+json") + : Results.NotFound(); +}); +app.MapGet("/.well-known/nodeinfo", (HttpContext context) => Results.Json(new +{ + links = new[] { new { rel = "http://nodeinfo.diaspora.software/ns/schema/2.0", href = $"https://{context.Request.Host.Host}/nodeinfo/2.0" } } +})); +app.MapGet("/nodeinfo/2.0", () => Results.Json(new +{ + version = "2.0", software = new { name = "flood", version = "1" }, protocols = new[] { "activitypub" }, + usage = new { users = new { total = perHost } }, openRegistrations = false +})); +app.MapPost("/inbox", (Delegate)Inbox); +app.MapPost("/users/{name}/inbox", (Delegate)Inbox); +app.Run(); +return 0; + +// a Follow of one of ours is accepted, signed by the followed actor, at the follower's inbox; anything else is taken +async Task Inbox(HttpContext context) +{ + var activity = await JsonNode.ParseAsync(context.Request.Body); + if (activity?["type"]?.GetValue() == "Follow" && activity["object"]?.GetValue() is { } followed + && keys.ById(followed) is { } actor && activity["actor"]?.GetValue() is { } follower) + _ = Task.Run(async () => + { + var inbox = (await Signer.Get(http, follower, actor))?["inbox"]?.GetValue(); + if (inbox == null) + return; + var accept = new JsonObject + { + ["@context"] = "https://www.w3.org/ns/activitystreams", ["id"] = $"{actor.Id}/accepts/{Guid.NewGuid():N}", + ["type"] = "Accept", ["actor"] = actor.Id, ["object"] = activity.DeepClone() + }; + await Signer.Post(http, inbox, accept, actor); + }); + return Results.Accepted(); +} + +static IResult Json(JsonNode document) => Results.Text(document.ToJsonString(), "application/activity+json"); + +static JsonObject Actor(FloodActor actor) => new() +{ + ["@context"] = new JsonArray("https://www.w3.org/ns/activitystreams", "https://w3id.org/security/v1"), + ["id"] = actor.Id, + ["type"] = "Person", + ["preferredUsername"] = actor.Name, + ["name"] = $"Flood {actor.Name} of {actor.Host}", + ["inbox"] = actor.Id + "/inbox", + ["outbox"] = actor.Id + "/outbox", + ["followers"] = actor.Id + "/followers", + ["following"] = actor.Id + "/following", + ["endpoints"] = new JsonObject { ["sharedInbox"] = $"https://{actor.Host}/inbox" }, + ["publicKey"] = new JsonObject { ["id"] = actor.Id + "#main-key", ["owner"] = actor.Id, ["publicKeyPem"] = actor.PublicPem } +}; + +sealed record FloodActor(string Host, string Name, string PrivatePem, string PublicPem) +{ + public string Id => $"https://{Host}/users/{Name}"; +} + +sealed class Keys(List actors) +{ + public IReadOnlyList All => actors; + + public FloodActor Find(string host, string name) => actors.FirstOrDefault(a => a.Host == host && a.Name == name); + + public FloodActor ById(string id) => actors.FirstOrDefault(a => a.Id == id); + + public static Keys Load(string state, int hosts, int perHost) + { + var path = Path.Combine(state, "keys.json"); + if (File.Exists(path)) + return new Keys(JsonSerializer.Deserialize>(File.ReadAllText(path))); + var actors = new List(); + for (var h = 1; h <= hosts; h++) + for (var i = 0; i < perHost; i++) + { + using var rsa = RSA.Create(2048); + actors.Add(new FloodActor($"flood{h}.test", $"u{i}", rsa.ExportPkcs8PrivateKeyPem(), rsa.ExportSubjectPublicKeyInfoPem())); + } + Directory.CreateDirectory(state); + File.WriteAllText(path, JsonSerializer.Serialize(actors)); + return new Keys(actors); + } +} + +// draft-cavage rsa-sha256, as PrivaPub signs and checks it: (request-target) host date [digest] +static class Signer +{ + public static HttpClient Client() => new(new SocketsHttpHandler { MaxConnectionsPerServer = 128, PooledConnectionLifetime = TimeSpan.FromMinutes(5) }) + { + Timeout = TimeSpan.FromSeconds(60) + }; + + static string Signature(FloodActor actor, string signingString, string headers) + { + using var rsa = RSA.Create(); + rsa.ImportFromPem(actor.PrivatePem); + var signature = Convert.ToBase64String(rsa.SignData(Encoding.UTF8.GetBytes(signingString), HashAlgorithmName.SHA256, RSASignaturePadding.Pkcs1)); + return $"keyId=\"{actor.Id}#main-key\",algorithm=\"rsa-sha256\",headers=\"{headers}\",signature=\"{signature}\""; + } + + public static async Task Post(HttpClient http, string inbox, JsonNode activity, FloodActor actor) + { + var uri = new Uri(inbox); + var body = Encoding.UTF8.GetBytes(activity.ToJsonString()); + var date = DateTime.UtcNow.ToString("r", CultureInfo.InvariantCulture); + var digest = "SHA-256=" + Convert.ToBase64String(SHA256.HashData(body)); + using var request = new HttpRequestMessage(HttpMethod.Post, uri) { Content = new ByteArrayContent(body) }; + request.Content.Headers.TryAddWithoutValidation("Content-Type", "application/activity+json"); + request.Headers.TryAddWithoutValidation("Date", date); + request.Headers.TryAddWithoutValidation("Digest", digest); + request.Headers.TryAddWithoutValidation("Signature", + Signature(actor, $"(request-target): post {uri.PathAndQuery}\nhost: {uri.Host}\ndate: {date}\ndigest: {digest}", "(request-target) host date digest")); + try + { + using var response = await http.SendAsync(request); + return response.StatusCode; + } + catch (Exception) + { + return 0; + } + } + + public static async Task Get(HttpClient http, string url, FloodActor actor) + { + var uri = new Uri(url); + var date = DateTime.UtcNow.ToString("r", CultureInfo.InvariantCulture); + using var request = new HttpRequestMessage(HttpMethod.Get, uri); + request.Headers.TryAddWithoutValidation("Accept", "application/activity+json"); + request.Headers.TryAddWithoutValidation("Date", date); + request.Headers.TryAddWithoutValidation("Signature", Signature(actor, $"(request-target): get {uri.PathAndQuery}\nhost: {uri.Host}\ndate: {date}", "(request-target) host date")); + try + { + using var response = await http.SendAsync(request); + return response.IsSuccessStatusCode ? JsonNode.Parse(await response.Content.ReadAsStringAsync()) : null; + } + catch (Exception) + { + return null; + } + } +} + +// the load: activities at a rate, each from a random actor, to PrivaPub's shared inbox +static class Run +{ + public static async Task Start(string[] args, Keys keys) + { + var options = args.Select(a => a.Split('=', 2)).Where(p => p.Length == 2).ToDictionary(p => p[0].TrimStart('-'), p => p[1]); + var target = options.GetValueOrDefault("target", "https://privapub.test"); + var rate = double.Parse(options.GetValueOrDefault("rate", "20"), CultureInfo.InvariantCulture); + var seconds = double.Parse(options.GetValueOrDefault("seconds", "60"), CultureInfo.InvariantCulture); + var personas = options.GetValueOrDefault("personas", "").Split(',', StringSplitOptions.RemoveEmptyEntries); + var posts = options.TryGetValue("posts", out var postsFile) && File.Exists(postsFile) + ? File.ReadAllLines(postsFile).Where(l => l.Length > 0).ToArray() : []; + var mix = options.GetValueOrDefault("mix", "create=80,like=20").Split(',') + .Select(p => p.Split('=')).ToDictionary(p => p[0], p => double.Parse(p[1], CultureInfo.InvariantCulture)); + var inbox = target.TrimEnd('/') + "/human-centipede"; + var http = Signer.Client(); + var random = new Random(4127); + var gate = new SemaphoreSlim(int.Parse(options.GetValueOrDefault("concurrency", "64"), CultureInfo.InvariantCulture)); + var results = new ConcurrentBag<(string Kind, int Status, double Ms)>(); + var pending = new List(); + var clock = Stopwatch.StartNew(); + var total = (int)(rate * seconds); + for (var n = 0; n < total; n++) + { + var due = TimeSpan.FromSeconds(n / rate); + if (due > clock.Elapsed) + await Task.Delay(due - clock.Elapsed); + var actor = keys.All[random.Next(keys.All.Count)]; + var kind = Pick(mix, random); + var activity = kind switch + { + "like" when posts.Length > 0 => Like(actor, posts[random.Next(posts.Length)]), + "follow" when personas.Length > 0 => Follow(actor, personas[random.Next(personas.Length)]), + _ => Create(actor, personas, n) + }; + await gate.WaitAsync(); + pending.Add(Task.Run(async () => + { + var sent = Stopwatch.StartNew(); + var status = await Signer.Post(http, inbox, activity, actor); + results.Add((activity["type"]!.GetValue(), (int)status, sent.Elapsed.TotalMilliseconds)); + gate.Release(); + })); + } + await Task.WhenAll(pending); + var elapsed = clock.Elapsed.TotalSeconds; + var latencies = results.Select(r => r.Ms).Order().ToArray(); + double Percentile(double p) => latencies.Length == 0 ? 0 : latencies[Math.Min(latencies.Length - 1, (int)Math.Ceiling(p * latencies.Length) - 1)]; + Console.WriteLine(new JsonObject + { + ["sent"] = results.Count, + ["seconds"] = Math.Round(elapsed, 1), + ["rate"] = Math.Round(results.Count / elapsed, 1), + ["byType"] = new JsonObject(results.GroupBy(r => r.Kind).Select(g => KeyValuePair.Create(g.Key, (JsonNode)g.Count()))), + ["byStatus"] = new JsonObject(results.GroupBy(r => r.Status.ToString(CultureInfo.InvariantCulture)).Select(g => KeyValuePair.Create(g.Key, (JsonNode)g.Count()))), + ["p50Ms"] = Math.Round(Percentile(0.5), 1), + ["p95Ms"] = Math.Round(Percentile(0.95), 1), + ["p99Ms"] = Math.Round(Percentile(0.99), 1), + ["maxMs"] = Math.Round(latencies.LastOrDefault(), 1) + }.ToJsonString()); + return 0; + } + + static string Pick(Dictionary mix, Random random) + { + var roll = random.NextDouble() * mix.Values.Sum(); + foreach (var (kind, weight) in mix) + if ((roll -= weight) <= 0) + return kind; + return mix.Keys.First(); + } + + static JsonObject Create(FloodActor actor, string[] personas, int n) + { + var id = $"{actor.Id}/notes/{Guid.NewGuid():N}"; + var mentioned = personas.Length == 0 ? null : personas[n % personas.Length]; + var note = new JsonObject + { + ["id"] = id, ["type"] = "Note", ["attributedTo"] = actor.Id, + ["content"] = $"

flood note {n}{(mentioned == null ? "" : $" for @someone")}

", + ["published"] = DateTime.UtcNow.ToString("o", CultureInfo.InvariantCulture), + ["to"] = new JsonArray("https://www.w3.org/ns/activitystreams#Public"), + ["cc"] = mentioned == null ? new JsonArray(actor.Id + "/followers") : new JsonArray(actor.Id + "/followers", mentioned), + ["tag"] = mentioned == null ? new JsonArray() : new JsonArray(new JsonObject { ["type"] = "Mention", ["href"] = mentioned }) + }; + return new JsonObject + { + ["@context"] = "https://www.w3.org/ns/activitystreams", ["id"] = id + "/activity", ["type"] = "Create", + ["actor"] = actor.Id, ["to"] = note["to"]!.DeepClone(), ["cc"] = note["cc"]!.DeepClone(), ["object"] = note + }; + } + + static JsonObject Like(FloodActor actor, string post) => new() + { + ["@context"] = "https://www.w3.org/ns/activitystreams", ["id"] = $"{actor.Id}/likes/{Guid.NewGuid():N}", + ["type"] = "Like", ["actor"] = actor.Id, ["object"] = post + }; + + static JsonObject Follow(FloodActor actor, string persona) => new() + { + ["@context"] = "https://www.w3.org/ns/activitystreams", ["id"] = $"{actor.Id}/follows/{Guid.NewGuid():N}", + ["type"] = "Follow", ["actor"] = actor.Id, ["object"] = persona + }; +} diff --git a/tools/pasture/load.sh b/tools/pasture/load.sh new file mode 100755 index 0000000..2cca72e --- /dev/null +++ b/tools/pasture/load.sh @@ -0,0 +1,82 @@ +#!/usr/bin/env bash +# PrivaPub under a crowd: the flood peer's twenty servers send signed activities at a set rate while a persona reads its +# home timeline, then the queue drains. Prints, and keeps in out/load/, what the crowd saw (PrivaPub's answers and how +# long they took), how long the activities waited in PrivaPub's queue and took to process, how long the queue took to +# drain, and how the home timeline answered meanwhile. Needs the pasture with the flood peer (run.sh add flood). +# usage: tools/pasture/load.sh [--rate=40] [--seconds=60] [--personas=5] [--follows=20] [--mix=create=70,like=20,follow=10] +set -uo pipefail +here="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +. "$here/lib/interop.sh" +rate=40; seconds=60; personas=5; follows=20; mix="create=70,like=20,follow=10" +for arg in "$@"; do + case "$arg" in + --rate=*) rate=${arg#*=} ;; --seconds=*) seconds=${arg#*=} ;; --personas=*) personas=${arg#*=} ;; + --follows=*) follows=${arg#*=} ;; --mix=*) mix=${arg#*=} ;; + *) echo "usage: $0 [--rate=N] [--seconds=N] [--personas=N] [--follows=N] [--mix=...]" >&2; exit 2 ;; + esac +done +podman container exists pasture-flood || { echo "the flood peer is not up: run.sh add flood" >&2; exit 1; } +mongo() { podman exec pasture-mongo mongosh --quiet PrivaPub --eval "$1"; } + +echo "load: $personas personas following $follows flood actors each; $rate activities a second for ${seconds}s ($mix)" +uris=(); tokens=(); : > "$here/.state/load-posts.txt" +for i in $(seq 0 $((personas - 1))); do + token=$(privapub_token "load$i") + [ -n "$token" ] || { echo "no token for load$i" >&2; exit 1; } + tokens+=("$token") + uris+=("$(curl -s -H "Authorization: Bearer $token" "$P/api/v1/accounts/verify_credentials" | j "print(d['url'].replace('/@', '/peasants/'))")") + # each follows its own share of the crowd, so a note reaches some homes and not others + for f in $(seq 0 $((follows - 1))); do + n=$(( (i * follows + f) % 200 )) + id=$(curl -s -H "Authorization: Bearer $token" "$P/api/v2/search?q=u$((n % 10))@flood$((n / 10 + 1)).test&resolve=true&type=accounts" | j "print(d['accounts'][0]['id'])") + [ -n "$id" ] && curl -s -o /dev/null -X POST -H "Authorization: Bearer $token" "$P/api/v1/accounts/$id/follow" + done + for p in 1 2; do + curl -s -X POST -H "Authorization: Bearer $token" "$P/api/v1/statuses" -d "status=a post for the crowd to like ($p)&visibility=public" | j "print(d['uri'])" >> "$here/.state/load-posts.txt" + done +done +podman cp "$here/.state/load-posts.txt" pasture-flood:/state/posts.txt +# the crowd's Accepts arrive before the run starts +until_true 60 '[ "$(mongo "print(db.Following.countDocuments({TargetActorURI: /flood/, State: 1}))")" -ge $((personas * follows * 9 / 10)) ]' || true +echo " follows accepted: $(mongo "print(db.Following.countDocuments({TargetActorURI: /flood/, State: 1}))")" + +started=$(date -u +%Y-%m-%dT%H:%M:%S.000Z) +reader="${tokens[0]}" +sampler="$here/.state/load-api.txt"; : > "$sampler" +rm -f "$here/.state/load-stop"; echo 0 > "$here/.state/load-peak" +( while [ ! -e "$here/.state/load-stop" ]; do + curl -s -o /dev/null -w '%{http_code} %{time_total}\n' -H "Authorization: Bearer $reader" "$P/api/v1/timelines/home?limit=40" >> "$sampler" + sleep 0.25 +done ) & +( peak=0; while [ ! -e "$here/.state/load-stop" ]; do + q=$(mongo 'print(db.Job.countDocuments({State: {$in: [0, 1]}, RunAt: {$lte: new Date()}}))'); [ "$q" -gt "$peak" ] && peak=$q && echo "$peak" > "$here/.state/load-peak"; sleep 2 +done ) & +crowd=$(podman exec pasture-flood /app/flood run --target=https://privapub.test --rate="$rate" --seconds="$seconds" \ + --personas="$(IFS=,; echo "${uris[*]}")" --posts=/state/posts.txt --mix="$mix") +sent_at=$(date +%s) +# (jobs that are due: retries waiting for a later attempt, a dead host's, are not the crowd's) +until_true 300 '[ "$(mongo "print(db.Job.countDocuments({State: {\$in: [0, 1]}, RunAt: {\$lte: new Date()}}))")" = "0" ]' +drain=$(( $(date +%s) - sent_at )) +touch "$here/.state/load-stop"; wait 2>/dev/null; rm -f "$here/.state/load-stop" + +inbox=$(mongo " + const events = db.InteractionEvent.find({Channel: 'in', Host: /^flood/, At: {\$gte: ISODate('$started')}}, {Outcome: 1, Reason: 1, WaitMs: 1, LatencyMs: 1}).toArray(); + const pct = (xs, p) => { const s = xs.filter(x => x != null).sort((a, b) => a - b); return s.length ? s[Math.min(s.length - 1, Math.ceil(p * s.length) - 1)] : null; }; + const processed = events.filter(e => e.Outcome != 'queued' && e.Outcome != 'refused'); + const outcomes = {}; events.forEach(e => outcomes[e.Outcome] = (outcomes[e.Outcome] || 0) + 1); + print(JSON.stringify({ events: events.length, outcomes, waitP50Ms: pct(processed.map(e => e.WaitMs), 0.5), waitP95Ms: pct(processed.map(e => e.WaitMs), 0.95), + processP50Ms: pct(processed.map(e => e.LatencyMs), 0.5), processP95Ms: pct(processed.map(e => e.LatencyMs), 0.95) }));") +api=$(python3 - "$sampler" <<'PY' +import json, sys +rows = [l.split() for l in open(sys.argv[1]) if l.strip()] +times = sorted(float(t) * 1000 for s, t in rows if s == "200") +pct = lambda p: round(times[min(len(times) - 1, max(0, int(-(-p * len(times) // 1)) - 1))], 1) if times else None +print(json.dumps({"requests": len(rows), "ok": len(times), "p50Ms": pct(0.5), "p95Ms": pct(0.95), "maxMs": round(times[-1], 1) if times else None})) +PY +) +mkdir -p "$here/out/load" +result=$(python3 -c 'import json, sys; print(json.dumps({"started": sys.argv[1], "rate": float(sys.argv[2]), "seconds": float(sys.argv[3]), "personas": int(sys.argv[4]), + "follows": int(sys.argv[5]), "mix": sys.argv[6], "crowd": json.loads(sys.argv[7]), "inbox": json.loads(sys.argv[8]), "queuePeak": int(sys.argv[9]), + "drainSeconds": int(sys.argv[10]), "homeTimeline": json.loads(sys.argv[11])}, indent=1))' \ + "$started" "$rate" "$seconds" "$personas" "$follows" "$mix" "$crowd" "$inbox" "$(cat "$here/.state/load-peak" 2>/dev/null || echo 0)" "$drain" "$api") +echo "$result" | tee "$here/out/load/$(date +%Y%m%d-%H%M%S).json" diff --git a/tools/pasture/peers/flood.sh b/tools/pasture/peers/flood.sh new file mode 100644 index 0000000..8531727 --- /dev/null +++ b/tools/pasture/peers/flood.sh @@ -0,0 +1,16 @@ +# flood: twenty fake servers (flood1..20.test, ten actors each) for load runs (load.sh), from tools/pasture/flood/flood.cs, +# a .NET file-based app published once into .flood and run on the runtime image PrivaPub uses. Its keys live in the +# pasture-flood volume, so `flood run` (podman exec) speaks as the actors `flood serve` answers for. +flood_up() { + local dotnet=${DOTNET:-$(command -v dotnet || echo ~/.dotnet/dotnet)} + [ -x "$here/.flood/flood" ] || "$dotnet" publish "$here/flood/flood.cs" -r linux-x64 --self-contained true -o "$here/.flood" -v quiet + podman volume exists pasture-flood || podman volume create --label pasture=1 pasture-flood >/dev/null + podman run -d --replace --name pasture-flood --network $net --label pasture=1 -w /app -v "$here/.flood:/app:Z,ro" \ + -v pasture-flood:/state -v "$ca/bundle.pem:/etc/ssl/certs/ca-certificates.crt:z,ro" -e SSL_CERT_FILE=/etc/ssl/certs/ca-certificates.crt \ + mcr.microsoft.com/dotnet/runtime-deps:10.0 /app/flood serve >/dev/null + for _ in $(seq 1 60); do + site flood1.test -s -o /dev/null -w '%{http_code}' https://flood1.test:6443/users/u0 2>/dev/null | grep -q 200 && break + sleep 2 + done + echo "flood: https://flood1.test:6443 .. flood20.test" +}