tools/pasture/flood/flood.cs answers as twenty fake servers (flood1..20.test) and sends signed Creates, Likes and Follows at a set rate; load.sh measures the answers, the queue's wait and processing times, its drain and a persona's home timeline meanwhile (docs/LOAD.md has the method and the runs). What the runs found: - Every unique index was partial on $type: "string", which MongoDB never uses for an equality lookup, so every post by ObjectURI, actor by ActorURI, deleted object, domain block, remote instance and the rest was a collection scan (280 ms a post lookup at 30 000 posts). They are partial on $gt: "" now, which an equality on a string implies; MongoDB.Entities rebuilds them in place at the next start. - Two inbox workers capped intake near 110 activities a second: Federation:InboxConcurrency and DeliveryConcurrency (default 8) set them. - The indexes the plan listed as missing: a post's boosts and replies, a persona's boosts, who follows an actor, timeline rows by author, a post's likes and pins. At 300 activities a second (200 let through, the rest 429 by the per-origin limit) the queue wait went from 29 s to 6 ms at p50; with the limits lifted PrivaPub processes about 900 a second, each in under 10 ms, and the home timeline stays under 20 ms. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw
296 lines
13 KiB
C#
296 lines
13 KiB
C#
#: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<i>, 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<IResult> Inbox(HttpContext context)
|
|
{
|
|
var activity = await JsonNode.ParseAsync(context.Request.Body);
|
|
if (activity?["type"]?.GetValue<string>() == "Follow" && activity["object"]?.GetValue<string>() is { } followed
|
|
&& keys.ById(followed) is { } actor && activity["actor"]?.GetValue<string>() is { } follower)
|
|
_ = Task.Run(async () =>
|
|
{
|
|
var inbox = (await Signer.Get(http, follower, actor))?["inbox"]?.GetValue<string>();
|
|
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<FloodActor> actors)
|
|
{
|
|
public IReadOnlyList<FloodActor> 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<List<FloodActor>>(File.ReadAllText(path)));
|
|
var actors = new List<FloodActor>();
|
|
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<HttpStatusCode> 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<JsonNode> 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<int> 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<Task>();
|
|
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<string>(), (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<string, double> 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"] = $"<p>flood note {n}{(mentioned == null ? "" : $" for <a href=\"{mentioned}\" class=\"mention\">@someone</a>")}</p>",
|
|
["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
|
|
};
|
|
}
|