From 3d7d4068347e9ea14d34f9f8b4acdc35ae507c7a Mon Sep 17 00:00:00 2001 From: thepra Date: Wed, 7 Oct 2026 12:43:12 +0200 Subject: [PATCH] A persona's archive: exported in Mastodon's layout, imported without telling anyone A root exports one of its personas (a job) as Mastodon's account archive, so other servers' importers read it: the actor with its public key only, its own posts and boosts with their media, likes, bookmarks and Mastodon's CSV files, plus PrivaPub's filters, followed hashtags, notification policy, pins, scheduled and located posts; nothing of its root, its siblings, its keys or anyone's token. A ticket link downloads it, for a week. An archive (PrivaPub's or Mastodon's) uploaded in pieces is imported into a persona in a job, the parts the root picks, with progress and a stop. SafeArchive refuses links, escaping paths, duplicates, bombs and oversized items, and reads the outbox one item at a time. Imported posts are delivered to no one, put in no home and notify nobody, yet show on the profile, outbox, hashtags and search; back home a post keeps its id, from another actor it gets one of its date and ImportedFromURI, so importing twice changes nothing. Relationships go through the existing services; located, scheduled and likes only when asked; followers never. tools/pasture/scenarios/persona-archive.sh imports mastouser's real Mastodon archive (156 posts, 24 pictures) into a persona Mastodon follows: Mastodon receives none of it, and a second import changes nothing. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LsXgEaXee4GCU1hwYgPJXw --- CLAUDE.md | 22 + .../User/ViewPersonaArchive.cs | 48 + PrivaPub.Tests/Http/PersonaArchiveTests.cs | 342 +++++++ .../PersonaArchiveController.cs | 165 ++++ .../Domain/Portability/ExportArchiveJob.cs | 399 ++++++++ .../Domain/Portability/ImportArchiveJob.cs | 928 ++++++++++++++++++ .../Domain/Portability/PersonaArchives.cs | 79 ++ PrivaPub/Domain/Portability/SafeArchive.cs | 270 +++++ PrivaPub/Domain/Social/FollowService.cs | 4 + PrivaPub/Infrastructure/Backup/Backups.cs | 3 +- .../Infrastructure/Backup/ServerBackup.cs | 3 +- .../Infrastructure/Backup/TransferStore.cs | 28 + PrivaPub/Infrastructure/Data/Indexes.cs | 8 + .../Middleware/SocialPubConfigurations.cs | 2 + PrivaPub/Models/Jobs/Job.cs | 4 +- PrivaPub/Models/Post/Post.cs | 1 + deploy/nginx/privapub.thepra.dev.conf | 8 +- tools/pasture/scenarios/persona-archive.sh | 70 ++ 18 files changed, 2377 insertions(+), 7 deletions(-) create mode 100644 PrivaPub.ClientModels/User/ViewPersonaArchive.cs create mode 100644 PrivaPub.Tests/Http/PersonaArchiveTests.cs create mode 100644 PrivaPub/Controllers/ClientToServer/PersonaArchiveController.cs create mode 100644 PrivaPub/Domain/Portability/ExportArchiveJob.cs create mode 100644 PrivaPub/Domain/Portability/ImportArchiveJob.cs create mode 100644 PrivaPub/Domain/Portability/PersonaArchives.cs create mode 100644 PrivaPub/Domain/Portability/SafeArchive.cs create mode 100644 tools/pasture/scenarios/persona-archive.sh diff --git a/CLAUDE.md b/CLAUDE.md index fa08536..25597e0 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -541,6 +541,28 @@ group www-data and reaches the private mongod; `sudo -u www-data` works too. Every call checks the administrator against the database, not only the token. nginx streams downloads for an hour and takes upload pieces unbuffered (`deploy/nginx`, applied by `setup.sh`). +- **A persona's archive** (`Domain/Portability/`, `PersonaArchiveController` at `/clientapi/persona/{avatarId}/archive`; + owner decision 2026-10-07), for the persona's own root only, never a banned one. State per persona in `PersonaArchive` + (never in a backup), files in `/.personas//`, a week. + - **Export** (`ExportArchiveJob`): Mastodon's account-archive layout, so other servers' importers read it: `actor.json` + (its public key only), `outbox.json` (its own posts as Creates, addressed as delivered, but circles' and located + posts; its boosts as Announces), `media_attachments/files/` (what the attachments' urls name), avatar and header, + `likes.json`, `bookmarks.json` and Mastodon's CSV files; then `privapub/`: filters, followed hashtags, the + notification policy, pins, scheduled and located posts, `archive.json`. Nothing of its root, its siblings, its keys, + anyone's token. Downloaded through a ten-minute ticket in the link's path. + - **Import** (`ImportArchiveJob`, the parts the root picks, with progress, stoppable): uploaded in pieces like a backup, + read through `SafeArchive`, which refuses links, paths leaving it, names twice, too many or too large files, a file + compressed over 100:1, and streams the outbox an item at a time (none over 1 MB). Only the archive actor's own + Creates become posts; boosts, direct and circle posts are counted. **Imported posts are delivered to no one**, put in + no home and notify nobody (an exception next to located posts), yet show on the profile, the outbox, hashtags and + search. Back home (same actor) a post keeps its id and address, and one already here, even deleted, or tombstoned, + stays as it is; from another actor it gets an id of its date, its address here and `ImportedFromURI` (unique per + persona), so importing twice changes nothing. Replies and self-quotes inside the archive point at their copies; polls + come closed without votes; mentions stay links; HTML is sanitized; media go through `MediaService` and the quota. + The other parts go through the existing services: follows asked again (never on a blocked server), blocks, mutes, + blocked servers, lists (members only if followed), bookmarks (by signed fetch), filters, followed hashtags, the + notification policy, pins (no `Add` delivered), the profile (announced as Settings would). Located and scheduled + posts (PrivaPub's archives) and likes (each tells its author) only when asked. Followers are never imported. ## Code style diff --git a/PrivaPub.ClientModels/User/ViewPersonaArchive.cs b/PrivaPub.ClientModels/User/ViewPersonaArchive.cs new file mode 100644 index 0000000..a9bdd6a --- /dev/null +++ b/PrivaPub.ClientModels/User/ViewPersonaArchive.cs @@ -0,0 +1,48 @@ +namespace PrivaPub.ClientModels.User +{ + // A persona's archive (owner decision 2026-10-07): made for it to take away, in Mastodon's account-archive layout, and + // restored into it. States: none, queued, running, ready (export) or done (import), failed, stopped. + public class ViewPersonaArchive + { + public string ExportState { get; set; } + public DateTime? ExportAskedAt { get; set; } + public DateTime? ExportReadyAt { get; set; } + public DateTime? ExportExpiresAt { get; set; } + public long ExportBytes { get; set; } + public string ExportError { get; set; } + + public string ImportState { get; set; } + public DateTime? ImportAskedAt { get; set; } + public DateTime? ImportEndedAt { get; set; } + public List ImportParts { get; set; } = []; + public long ImportDone { get; set; } + public long ImportTotal { get; set; } + public Dictionary ImportCounts { get; set; } = [];//what was imported and what skipped, by why + public string ImportError { get; set; } + public string ImportFrom { get; set; }//the archive's actor + public bool ImportSameActor { get; set; } + } + + // the parts of an archive to import: posts (with their media), profile, follows, blocks, mutes, domainblocks, lists, + // bookmarks, filters, tags, notificationpolicy, pins; and only when asked, located, scheduled, likes (which tell their + // authors) + public class ImportPersonaArchiveForm + { + public string UploadId { get; set; } + public List Parts { get; set; } = []; + } + + public class ViewArchiveTicket + { + public string Url { get; set; } + public DateTime ExpiresAt { get; set; } + public long Bytes { get; set; } + } + + public class ViewArchiveUpload + { + public string Id { get; set; } + public long Received { get; set; } + public long ChunkBytes { get; set; } + } +} diff --git a/PrivaPub.Tests/Http/PersonaArchiveTests.cs b/PrivaPub.Tests/Http/PersonaArchiveTests.cs new file mode 100644 index 0000000..f5194ad --- /dev/null +++ b/PrivaPub.Tests/Http/PersonaArchiveTests.cs @@ -0,0 +1,342 @@ +using MongoDB.Entities; + +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Media; +using PrivaPub.Models.Social; +using PrivaPub.Models.User; +using PrivaPub.Tests.Support; +using PrivaPub.Tests.Support.Host; + +using System.IO.Compression; +using System.Net; +using System.Net.Http.Json; +using System.Text; +using System.Text.Json.Nodes; + +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Tests.Http +{ + // A persona's archive (owner decision 2026-10-07): exported in Mastodon's layout with nothing of its root, its siblings + // or its keys; imported back home without bringing anything back, or into another persona as posts nobody is told of, + // once; Mastodon's own archives read too; and hostile ones refused. + [Trait("Category", "Integration")] + public sealed class PersonaArchiveTests : 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; + + static string Base(Persona persona) => $"/clientapi/persona/{persona.Id}/archive"; + + async Task Status(Persona persona) + { + using var client = _host.As(persona.Root.Jwt); + return (await client.GetFromJsonAsync(Base(persona), Token))!; + } + + async Task Exported(Persona persona) + { + using var client = _host.As(persona.Root.Jwt); + Assert.Equal(HttpStatusCode.Accepted, (await client.PostAsync(Base(persona), null, Token)).StatusCode); + await _host.Run(j => j.Kind == JobKind.ExportArchive && j.Payload == persona.Id, Token); + Assert.Equal("ready", (await Status(persona))["exportState"]!.GetValue()); + var ticket = (await (await client.PostAsync(Base(persona) + "/ticket", null, Token)).Content.ReadFromJsonAsync(Token))!; + using var anonymous = _host.Client(); + var response = await anonymous.GetAsync(ticket["url"]!.GetValue(), Token); + Assert.Equal(HttpStatusCode.OK, response.StatusCode); + return await response.Content.ReadAsByteArrayAsync(Token); + } + + // uploaded in two pieces, then imported with the parts asked for, the job run: the import's state + async Task<(HttpStatusCode Status, JsonObject State)> Imported(Persona persona, byte[] archive, params string[] parts) + { + using var client = _host.As(persona.Root.Jwt); + var upload = (await (await client.PostAsync("/clientapi/persona/archive/uploads", null, Token)).Content.ReadFromJsonAsync(Token))!["id"]!.GetValue(); + var half = archive.Length / 2; + Assert.Equal(HttpStatusCode.OK, (await client.PutAsync($"/clientapi/persona/archive/uploads/{upload}?offset=0", new ByteArrayContent(archive[..half]), Token)).StatusCode); + Assert.Equal(HttpStatusCode.OK, (await client.PutAsync($"/clientapi/persona/archive/uploads/{upload}?offset={half}", new ByteArrayContent(archive[half..]), Token)).StatusCode); + var asked = await client.PostAsJsonAsync(Base(persona) + "/import", new { uploadId = upload, parts }, Token); + if (asked.StatusCode != HttpStatusCode.Accepted) + return (asked.StatusCode, default); + await _host.Run(j => j.Kind == JobKind.ImportArchive && j.Payload == persona.Id, Token); + return (asked.StatusCode, await Status(persona)); + } + + static int Count(JsonObject state, string what) => state["importCounts"]?[what]?.GetValue() ?? 0; + + static Dictionary Texts(byte[] zip) + { + using var archive = new ZipArchive(new MemoryStream(zip)); + return archive.Entries.Where(e => e.Name.EndsWith(".json") || e.Name.EndsWith(".csv")) + .ToDictionary(e => e.FullName, e => new StreamReader(e.Open()).ReadToEnd()); + } + + [Fact] + public async Task An_archive_holds_the_persona_and_nothing_of_its_root_its_siblings_or_its_keys() + { + var alice = await _host.Mastodon("alice"); + var sibling = await _host.Persona(alice.Persona.Root, "sibling"); + var bob = await _host.Mastodon("bob"); + var picture = (await alice.Client.Exchange(new HttpRequestMessage(HttpMethod.Post, "/api/v2/media") + { + Content = MastodonHelpers.Multipart(("file", MastodonHelpers.JpegWithMetadata(64, 48), "image/jpeg", "dot.jpg")) + })).Ok().Body.Text("id"); + var first = await alice.Status("a public post with a picture #archives", ("media_ids[]", picture)); + await alice.Status("for my followers only", ("visibility", "private")); + await alice.Status("a reply to myself", ("in_reply_to_id", first.Text("id"))); + await alice.Client.Post($"/api/v1/statuses/{first.Text("id")}/bookmark"); + await alice.Client.Post($"/api/v1/statuses/{first.Text("id")}/pin"); + await alice.Client.Post($"/api/v1/accounts/{bob.Persona.Id}/block"); + await alice.Client.Post("/api/v1/domain_blocks", ("domain", "spam.example")); + await alice.Client.Post("/api/v1/tags/archives/follow"); + + var zip = await Exported(alice.Persona); + var texts = Texts(zip); + + foreach (var file in new[] { "actor.json", "outbox.json", "likes.json", "bookmarks.json", "following_accounts.csv", "blocked_accounts.csv", + "muted_accounts.csv", "blocked_domains.csv", "lists.csv", "bookmarks.csv", "privapub/archive.json", "privapub/pins.json", + "privapub/followed_tags.json", "privapub/filters.json", "privapub/notification_policy.json" }) + Assert.True(texts.ContainsKey(file), $"{file} is missing"); + using (var archive = new ZipArchive(new MemoryStream(zip))) + Assert.Contains(archive.Entries, e => e.FullName.StartsWith("media_attachments/files/", StringComparison.Ordinal)); + var outbox = JsonNode.Parse(texts["outbox.json"])!["orderedItems"]!.AsArray(); + Assert.Equal(3, outbox.Count); + Assert.Contains(outbox, i => i!["object"]!["attachment"]?[0]?["url"]?.GetValue().StartsWith("/media_attachments/files/") == true); + Assert.Contains(bob.Persona.UserName, texts["blocked_accounts.csv"]); + Assert.Contains("spam.example", texts["blocked_domains.csv"]); + Assert.Contains(first.Text("uri"), texts["privapub/pins.json"]); + + var avatar = await DB.Default.Find().MatchID(alice.Persona.Id).ExecuteFirstAsync(Token); + var root = await DB.Default.Find().MatchID(alice.Persona.Root.Id).ExecuteFirstAsync(Token); + var all = string.Concat(texts.Values); + Assert.DoesNotContain(avatar.PrivateKey.Split('\n')[1], all); + Assert.DoesNotContain(avatar.SigningKey, all); + Assert.DoesNotContain(root.UserName, all); + Assert.DoesNotContain(root.HashedPassword, all); + Assert.DoesNotContain(sibling.UserName, all); + Assert.DoesNotContain(alice.Token, all); + } + + [Fact] + public async Task Back_home_an_archive_brings_nothing_back() + { + var alice = await _host.Mastodon("alice"); + await alice.Status("kept"); + var regretted = await alice.Status("regretted"); + var zip = await Exported(alice.Persona); + await alice.Client.Delete($"/api/v1/statuses/{regretted.Text("id")}"); + var before = await DB.Default.CountAsync(p => p.GroupUserId == alice.Persona.Id, Token); + + var (status, state) = await Imported(alice.Persona, zip, "posts", "pins"); + + Assert.Equal(HttpStatusCode.Accepted, status); + Assert.Equal("done", state["importState"]!.GetValue()); + Assert.True(state["importSameActor"]!.GetValue()); + Assert.Equal((0, 1, 1), (Count(state, "posts"), Count(state, "posts already here"), Count(state, "posts deleted here"))); + Assert.Equal(before, await DB.Default.CountAsync(p => p.GroupUserId == alice.Persona.Id, Token)); + Assert.NotNull((await DB.Default.Find().MatchID(regretted.Text("id")).ExecuteFirstAsync(Token)).DeletedAt); + } + + [Fact] + public async Task Into_another_persona_posts_arrive_once_and_nobody_is_told() + { + var alice = await _host.Mastodon("alice"); + var follower = await _host.Mastodon("follower"); + var picture = (await alice.Client.Exchange(new HttpRequestMessage(HttpMethod.Post, "/api/v2/media") + { + Content = MastodonHelpers.Multipart(("file", MastodonHelpers.JpegWithMetadata(64, 48), "image/jpeg", "dot.jpg")) + })).Ok().Body.Text("id"); + var first = await alice.Status("the first words #moving", ("media_ids[]", picture)); + await alice.Status("and a reply to them", ("in_reply_to_id", first.Text("id"))); + var zip = await Exported(alice.Persona); + + var carol = await _host.Mastodon("carol"); + await follower.Client.Post($"/api/v1/accounts/{carol.Persona.Id}/follow"); + var actor = $"/peasants/{carol.Persona.UserName}"; + var jobs = await DB.Default.CountAsync(j => j.Kind == JobKind.Deliver && j.Payload.Contains(actor), Token); + var (_, state) = await Imported(carol.Persona, zip, "posts", "pins", "tags"); + + Assert.Equal("done", state["importState"]!.GetValue()); + Assert.False(state["importSameActor"]!.GetValue()); + Assert.Equal(2, Count(state, "posts")); + Assert.Equal(1, Count(state, "media")); + var copies = await DB.Default.Find().Match(p => p.GroupUserId == carol.Persona.Id).ExecuteAsync(Token); + Assert.Equal(2, copies.Count); + var copy = copies.Single(p => p.ImportedFromURI == first.Text("uri")); + Assert.StartsWith($"{PrivaPubHost.Base}/peasants/{carol.Persona.UserName}/scribbles/", copy.ObjectURI); + Assert.Equal(["moving"], copy.Tags); + Assert.Single(copy.Media); + Assert.Equal(carol.Persona.Id, (await DB.Default.Find().MatchID(copy.Media[0].AttachmentId).ExecuteFirstAsync(Token)).OwnerAvatarId); + Assert.Equal(copy.ID, copies.Single(p => p.ID != copy.ID).AnsweringToPostId); + + // nobody told: no delivery, no home, no notification; on the profile all the same + Assert.Equal(jobs, await DB.Default.CountAsync(j => j.Kind == JobKind.Deliver && j.Payload.Contains(actor), Token)); + var ids = copies.Select(p => p.ID).ToList(); + Assert.False(await DB.Default.Find().Match(e => ids.Contains(e.PostId)).ExecuteAnyAsync(Token)); + Assert.Contains((await carol.Client.Get($"/api/v1/accounts/{carol.Persona.Id}/statuses")).Ok().Body!.AsArray(), s => s!.Text("id") == copy.ID); + + // twice changes nothing + var (_, again) = await Imported(carol.Persona, zip, "posts"); + Assert.Equal((0, 2), (Count(again, "posts"), Count(again, "posts already here"))); + Assert.Equal(2, await DB.Default.CountAsync(p => p.GroupUserId == carol.Persona.Id, Token)); + } + + // Mastodon's own archive: its posts, not its boosts nor its direct messages + [Fact] + public async Task A_mastodon_archive_is_read() + { + var dave = await _host.Mastodon("dave"); + const string actor = "https://mastodon.example/users/dave"; + const string Public = "https://www.w3.org/ns/activitystreams#Public"; + var outbox = new JsonObject + { + ["@context"] = "https://www.w3.org/ns/activitystreams", + ["id"] = "outbox.json", + ["type"] = "OrderedCollection", + ["orderedItems"] = new JsonArray( + Create(actor, "1", "

hello from Mastodon

", [Public], [actor + "/followers"], "2023-04-01T10:00:00Z"), + Create(actor, "2", "

only for one

", ["https://elsewhere.example/users/x"], [], "2023-04-02T10:00:00Z"), + new JsonObject { ["type"] = "Announce", ["actor"] = actor, ["object"] = "https://elsewhere.example/notes/9" }) + }; + var zip = Zip(("actor.json", new JsonObject { ["id"] = actor, ["type"] = "Person", ["followers"] = actor + "/followers" }.ToJsonString()), + ("outbox.json", outbox.ToJsonString())); + + var (_, state) = await Imported(dave.Persona, zip, "posts"); + + Assert.Equal("done", state["importState"]!.GetValue()); + Assert.Equal((1, 1, 1), (Count(state, "posts"), Count(state, "boosts skipped"), Count(state, "direct or circle posts skipped"))); + var post = await DB.Default.Find().Match(p => p.GroupUserId == dave.Persona.Id).ExecuteFirstAsync(Token); + Assert.Equal(actor + "/statuses/1", post.ImportedFromURI); + Assert.Equal(new DateTime(2023, 4, 1, 10, 0, 0, DateTimeKind.Utc), post.CreationDate); + Assert.DoesNotContain(" new() + { + ["type"] = "Create", + ["actor"] = actor, + ["object"] = new JsonObject + { + ["id"] = $"{actor}/statuses/{id}", + ["type"] = "Note", + ["attributedTo"] = actor, + ["content"] = content, + ["published"] = published, + ["to"] = new JsonArray(to.Select(t => (JsonNode)t).ToArray()), + ["cc"] = new JsonArray(cc.Select(c => (JsonNode)c).ToArray()) + } + }; + + static byte[] Zip(params (string Name, string Text)[] files) + { + var stream = new MemoryStream(); + using (var archive = new ZipArchive(stream, ZipArchiveMode.Create, leaveOpen: true)) + foreach (var (name, text) in files) + { + using var entry = archive.CreateEntry(name).Open(); + entry.Write(Encoding.UTF8.GetBytes(text)); + } + return stream.ToArray(); + } + + [Theory] + [InlineData("escape")] + [InlineData("twice")] + [InlineData("bomb")] + [InlineData("link")] + [InlineData("junk")] + public async Task A_hostile_archive_is_refused(string how) + { + var erin = await _host.Mastodon("erin"); + var actorJson = new JsonObject { ["id"] = "https://elsewhere.example/users/erin" }.ToJsonString(); + var zip = how switch + { + "escape" => Zip(("actor.json", actorJson), ("../escape.txt", "x")), + "twice" => Zip(("actor.json", actorJson), ("outbox.json", "{}"), ("outbox.json", "{}")), + "bomb" => Zip(("actor.json", actorJson), ("outbox.json", new string(' ', 30 * 1024 * 1024))), + "link" => Link(actorJson), + _ => "not a zip at all"u8.ToArray() + }; + + var (status, _) = await Imported(erin.Persona, zip, "posts"); + + Assert.Equal(HttpStatusCode.UnprocessableEntity, status); + Assert.False(await DB.Default.Find().Match(p => p.GroupUserId == erin.Persona.Id).ExecuteAnyAsync(Token)); + } + + // a symbolic link, as zip on Unix stores one + static byte[] Link(string actorJson) + { + var stream = new MemoryStream(); + using (var archive = new ZipArchive(stream, ZipArchiveMode.Create, leaveOpen: true)) + { + using (var actor = archive.CreateEntry("actor.json").Open()) + actor.Write(Encoding.UTF8.GetBytes(actorJson)); + var link = archive.CreateEntry("media_attachments/files/passwd"); + link.ExternalAttributes = unchecked((int)(0xA1FFu << 16)); + using var target = link.Open(); + target.Write("/etc/passwd"u8); + } + return stream.ToArray(); + } + + [Fact] + public async Task An_outbox_item_too_large_stops_the_import() + { + var frank = await _host.Mastodon("frank"); + const string actor = "https://elsewhere.example/users/frank"; + // (random, so it is no compression bomb: that is refused before the import starts) + var words = Convert.ToBase64String(System.Security.Cryptography.RandomNumberGenerator.GetBytes(1536 * 1024)); + var huge = Create(actor, "1", words, ["https://www.w3.org/ns/activitystreams#Public"], [], "2024-01-01T00:00:00Z"); + var zip = Zip(("actor.json", new JsonObject { ["id"] = actor }.ToJsonString()), + ("outbox.json", new JsonObject { ["orderedItems"] = new JsonArray(huge) }.ToJsonString())); + + var (_, state) = await Imported(frank.Persona, zip, "posts"); + + Assert.Equal("failed", state["importState"]!.GetValue()); + Assert.Contains("larger than", state["importError"]!.GetValue()); + } + + [Fact] + public async Task Only_the_personas_own_root_and_never_a_banned_one() + { + var grace = await _host.Mastodon("grace"); + var stranger = await _host.SignUp("stranger"); + using (var other = _host.As(stranger.Jwt)) + Assert.Equal(HttpStatusCode.NotFound, (await other.GetAsync(Base(grace.Persona), Token)).StatusCode); + Assert.Equal(HttpStatusCode.OK, (await _host.As(grace.Persona.Root.Jwt).GetAsync(Base(grace.Persona), Token)).StatusCode); + await DB.Default.Update().MatchID(grace.Persona.Root.Id).Modify(u => u.IsBanned, true).ExecuteAsync(Token); + Assert.NotEqual(HttpStatusCode.OK, (await _host.As(grace.Persona.Root.Jwt).GetAsync(Base(grace.Persona), Token)).StatusCode); + } + + [Fact] + public async Task An_import_asked_to_stop_stops() + { + var henry = await _host.Mastodon("henry"); + const string actor = "https://elsewhere.example/users/henry"; + var zip = Zip(("actor.json", new JsonObject { ["id"] = actor }.ToJsonString()), + ("outbox.json", new JsonObject { ["orderedItems"] = new JsonArray(Create(actor, "1", "one", ["https://www.w3.org/ns/activitystreams#Public"], [], "2024-01-01T00:00:00Z")) }.ToJsonString())); + using var client = _host.As(henry.Persona.Root.Jwt); + var upload = (await (await client.PostAsync("/clientapi/persona/archive/uploads", null, Token)).Content.ReadFromJsonAsync(Token))!["id"]!.GetValue(); + await client.PutAsync($"/clientapi/persona/archive/uploads/{upload}?offset=0", new ByteArrayContent(zip), Token); + Assert.Equal(HttpStatusCode.Accepted, (await client.PostAsJsonAsync(Base(henry.Persona) + "/import", new { uploadId = upload, parts = new[] { "posts" } }, Token)).StatusCode); + Assert.Equal(HttpStatusCode.Conflict, (await client.PostAsJsonAsync(Base(henry.Persona) + "/import", new { uploadId = upload, parts = new[] { "posts" } }, Token)).StatusCode); + Assert.Equal(HttpStatusCode.OK, (await client.PostAsync(Base(henry.Persona) + "/import/stop", null, Token)).StatusCode); + + await _host.Run(j => j.Kind == JobKind.ImportArchive && j.Payload == henry.Persona.Id, Token); + + Assert.Equal("stopped", (await Status(henry.Persona))["importState"]!.GetValue()); + Assert.False(await DB.Default.Find().Match(p => p.GroupUserId == henry.Persona.Id).ExecuteAnyAsync(Token)); + } + } +} diff --git a/PrivaPub/Controllers/ClientToServer/PersonaArchiveController.cs b/PrivaPub/Controllers/ClientToServer/PersonaArchiveController.cs new file mode 100644 index 0000000..141c517 --- /dev/null +++ b/PrivaPub/Controllers/ClientToServer/PersonaArchiveController.cs @@ -0,0 +1,165 @@ +using Microsoft.AspNetCore.Authorization; +using Microsoft.AspNetCore.Mvc; + +using MongoDB.Entities; + +using PrivaPub.ClientModels; +using PrivaPub.ClientModels.User; +using PrivaPub.Domain.Portability; +using PrivaPub.Extensions; +using PrivaPub.Infrastructure.Backup; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.User; + +namespace PrivaPub.Controllers.ClientToServer +{ + // A persona's archive for its root (owner decision 2026-10-07): made in a job and downloaded through a link good for + // ten minutes; another archive (its own, or Mastodon's) uploaded in pieces and imported into it in a job, the parts the + // root picks, with progress, until it asks to stop. Only the persona's own root, never a banned one. + [ApiController, + Route("clientapi/persona"), + Authorize(Policy = Policies.IsUser)] + public class PersonaArchiveController(Backups backups, TransferStore transfers, IJobQueue jobs, Federation.Actors.ILocalActorService localActors) + : ControllerBase + { + [HttpGet, Route("{avatarId}/archive")] + public async Task Status(string avatarId, CancellationToken token) + { + if (!await Owns(avatarId, token)) + return NotFound(); + await PersonaArchives.Forget(backups.Root, token); + return Ok(PersonaArchives.View(await PersonaArchives.Of(avatarId, token))); + } + + [HttpPost, Route("{avatarId}/archive")] + public async Task Export(string avatarId, CancellationToken token) + { + if (!await Owns(avatarId, token)) + return NotFound(); + var state = await PersonaArchives.Of(avatarId, token); + if (state.ExportState is "queued" or "running") + return Conflict(new WebResult().Invalidate("The archive is being made.")); + state.ExportState = "queued"; + state.ExportAskedAt = DateTime.UtcNow; + state.ExportError = default; + await DB.Default.SaveAsync(state, token); + await jobs.EnqueueMany([ExportArchiveJob.JobFor(avatarId, Host)], token); + return Accepted(PersonaArchives.View(state)); + } + + [HttpPost, Route("{avatarId}/archive/ticket")] + public async Task Ticket(string avatarId, CancellationToken token) + { + if (!await Owns(avatarId, token)) + return NotFound(); + var state = await PersonaArchives.Of(avatarId, token); + var file = state.ExportFile == default ? default : Path.Combine(PersonaArchives.Folder(backups.Root, avatarId), state.ExportFile); + if (state.ExportState != "ready" || file == default || !System.IO.File.Exists(file)) + return NotFound(); + var (ticket, expires) = transfers.Ticket($"persona/{avatarId}/{state.ExportFile}"); + return Ok(new ViewArchiveTicket { Url = $"/clientapi/persona/archive/download/{ticket}", ExpiresAt = expires, Bytes = new FileInfo(file).Length }); + } + + // the archive itself, for a ticket: a plain link, so the browser saves it to disk (and resumes it) + [HttpGet, Route("archive/download/{ticket}"), AllowAnonymous] + public IActionResult Download(string ticket) + { + if (transfers.Redeem(ticket) is not { } what || what.Split('/') is not ["persona", var avatarId, var name] || name.Contains("..")) + return NotFound(); + var file = Path.Combine(PersonaArchives.Folder(backups.Root, avatarId), name); + if (!System.IO.File.Exists(file)) + return NotFound(); + Response.Headers.CacheControl = "no-store"; + return PhysicalFile(file, "application/zip", name, enableRangeProcessing: true); + } + + [HttpPost, Route("archive/uploads")] + public IActionResult StartUpload() => + Ok(new ViewArchiveUpload { Id = transfers.Start(User.GetUserId()), ChunkBytes = TransferStore.ChunkBytes }); + + [HttpGet, Route("archive/uploads/{id}")] + public IActionResult Upload(string id) + { + if (transfers.Owner(id) != User.GetUserId()) + return NotFound(); + return Ok(new ViewArchiveUpload { Id = id, Received = transfers.Received(id), ChunkBytes = TransferStore.ChunkBytes }); + } + + [HttpPut, Route("archive/uploads/{id}"), RequestSizeLimit(TransferStore.ChunkBytes + 1024 * 1024)] + public async Task Append(string id, [FromQuery] long offset, CancellationToken token) + { + if (transfers.Owner(id) != User.GetUserId()) + return NotFound(); + var (received, taken) = await transfers.Append(id, offset, Request.Body, token); + var view = new ViewArchiveUpload { Id = id, Received = received, ChunkBytes = TransferStore.ChunkBytes }; + return taken ? Ok(view) : Conflict(view); + } + + [HttpDelete, Route("archive/uploads/{id}")] + public IActionResult CancelUpload(string id) + { + if (transfers.Owner(id) != User.GetUserId()) + return NotFound(); + transfers.Cancel(id); + return Ok(); + } + + // an uploaded archive imported into the persona: refused at once when it is not one, or holds what it must not + [HttpPost, Route("{avatarId}/archive/import")] + public async Task Import(string avatarId, ImportPersonaArchiveForm form, CancellationToken token) + { + if (!await Owns(avatarId, token)) + return NotFound(); + var state = await PersonaArchives.Of(avatarId, token); + if (state.ImportState is "queued" or "running") + return Conflict(new WebResult().Invalidate("An import is running.")); + if (form.Parts.Count == 0 || form.Parts.Any(p => !PersonaArchives.Parts.Contains(p))) + return UnprocessableEntity(new WebResult().Invalidate("Choose what to import.")); + if (transfers.Owner(form.UploadId) != User.GetUserId()) + return NotFound(); + var path = transfers.Take(form.UploadId, PersonaArchives.Folder(backups.Root, avatarId), $"import-{form.UploadId}.zip"); + if (path == default) + return NotFound(); + var (archive, refused) = SafeArchive.Open(path); + using (archive) + if (archive?.Has("actor.json") != true) + { + System.IO.File.Delete(path); + return UnprocessableEntity(new WebResult().Invalidate(refused ?? "the archive has no actor.json")); + } + state.ImportState = "queued"; + state.ImportAskedAt = DateTime.UtcNow; + state.ImportEndedAt = default; + state.ImportParts = [.. form.Parts.Distinct()]; + state.ImportFile = path; + state.ImportCounts = []; + (state.ImportDone, state.ImportTotal, state.ImportError, state.StopAsked) = (0, 0, default, false); + await DB.Default.SaveAsync(state, token); + await jobs.EnqueueMany([ImportArchiveJob.JobFor(avatarId, Host)], token); + return Accepted(PersonaArchives.View(state)); + } + + [HttpPost, Route("{avatarId}/archive/import/stop")] + public async Task Stop(string avatarId, CancellationToken token) + { + if (!await Owns(avatarId, token)) + return NotFound(); + await DB.Default.Update().MatchID(avatarId).Match(a => a.ImportState == "queued" || a.ImportState == "running") + .Modify(a => a.StopAsked, true).ExecuteAsync(token); + return Ok(); + } + + string Host => new Uri(localActors.BaseAddress).Host; + + // the root's own persona, and a root that may still act + async Task Owns(string avatarId, CancellationToken token) + { + var rootId = User.GetUserId(); + if (string.IsNullOrEmpty(avatarId) || !await DB.Default.Find().Match(r => r.RootId == rootId && r.AvatarId == avatarId).ExecuteAnyAsync(token)) + return false; + var root = await DB.Default.Find().MatchID(rootId).ExecuteFirstAsync(token); + return root is { IsBanned: false, DeletedAt: null } + && await DB.Default.Find().MatchID(avatarId).Match(a => !a.DeletionAt.HasValue).ExecuteAnyAsync(token); + } + } +} diff --git a/PrivaPub/Domain/Portability/ExportArchiveJob.cs b/PrivaPub/Domain/Portability/ExportArchiveJob.cs new file mode 100644 index 0000000..06fe603 --- /dev/null +++ b/PrivaPub/Domain/Portability/ExportArchiveJob.cs @@ -0,0 +1,399 @@ +using MongoDB.Entities; + +using PrivaPub.Domain.Media; +using PrivaPub.Domain.Social; +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Rendering; +using PrivaPub.Infrastructure.Backup; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Media; +using PrivaPub.Models.Post; +using PrivaPub.Models.Social; +using PrivaPub.Models.User; + +using System.IO.Compression; +using System.Text; +using System.Text.Encodings.Web; +using System.Text.Json; +using System.Text.Json.Nodes; + +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Domain.Portability +{ + // A persona's archive, in Mastodon's account-archive layout so other servers' importers read it: + // - actor.json (as other servers see the persona: its public key, never the private one), avatar and header; + // - outbox.json: every post of its own as a Create (any visibility but circles' and located posts, addressed as it was + // delivered) and its boosts as Announces; the media its posts hold under media_attachments/files/, which the + // attachments' urls name; + // - likes.json, bookmarks.json, and the CSV files Mastodon exports: follows, lists, blocks, mutes, blocked servers, + // bookmarks; + // - privapub/: what only PrivaPub has: filters, followed hashtags, the notification policy, pinned posts, scheduled and + // located posts, and archive.json saying whose it is. + // Nothing of its root, its siblings, its keys or anyone's tokens. Written whole to a .partial file, then renamed; kept a + // week (PersonaArchives.Forget). + public class ExportArchiveJob(ILocalActorService localActors, IMediaService media, Backups backups, ILogger logger) : IJobHandler + { + public const int Format = 1; + static readonly JsonSerializerOptions Options = new() { WriteIndented = true, Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping }; + static readonly JsonWriterOptions Writer = new() { Indented = true, Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping }; + const string MediaFolder = "media_attachments/files/"; + + public JobKind Kind => JobKind.ExportArchive; + public int Concurrency => 1; + public int MaxAttempts => 2; + public int PerHostLimit => 1; + + public static Job JobFor(string avatarId, string host) => new() + { + Kind = JobKind.ExportArchive, + Payload = avatarId, + Host = host, + DedupeKey = $"export|{avatarId}|{DateTime.UtcNow.Ticks}" + }; + + public async Task Handle(Job job, CancellationToken token) + { + var avatarId = job.Payload; + var actor = await localActors.FindById(LocalActorKind.Person, avatarId, token); + if (actor == default) + return JobOutcome.Done; + await DB.Default.Update().MatchID(avatarId).Modify(a => a.ExportState, "running").Modify(a => a.ExportError, null).ExecuteAsync(token); + var folder = PersonaArchives.Folder(backups.Root, avatarId); + ServerBackup.MakeDirectory(Path.GetDirectoryName(folder)!); + ServerBackup.MakeDirectory(folder); + var name = $"{actor.UserName}-{DateTime.UtcNow:yyyyMMdd-HHmmss}.zip"; + var partial = Path.Combine(folder, name + ServerBackup.PartialSuffix); + try + { + await using (var file = new FileStream(partial, FileMode.Create, FileAccess.Write)) + using (var zip = new ZipArchive(file, ZipArchiveMode.Create)) + await Write(zip, actor, token); + var final = Path.Combine(folder, name); + File.Move(partial, final); + foreach (var old in Directory.EnumerateFiles(folder, "*.zip").Where(f => f != final)) + File.Delete(old); + await DB.Default.Update().MatchID(avatarId) + .Modify(a => a.ExportState, "ready").Modify(a => a.ExportReadyAt, DateTime.UtcNow) + .Modify(a => a.ExportFile, name).Modify(a => a.ExportBytes, new FileInfo(final).Length).ExecuteAsync(token); + return JobOutcome.Done; + } + catch (Exception ex) when (ex is not OperationCanceledException) + { + logger.LogError(ex, "The archive of {Persona} could not be made", actor.Handle); + if (File.Exists(partial)) + File.Delete(partial); + await DB.Default.Update().MatchID(avatarId).Modify(a => a.ExportState, "failed") + .Modify(a => a.ExportError, ex.Message).ExecuteAsync(CancellationToken.None); + return JobOutcome.Done; + } + } + + async Task Write(ZipArchive zip, LocalActor actor, CancellationToken token) + { + var avatar = await DB.Default.Find().MatchID(actor.Id).ExecuteFirstAsync(token); + var files = new SortedSet(StringComparer.Ordinal); + var prefix = media.Url(string.Empty); + + var document = ActivityPubRenderer.Actor(actor); + foreach (var (field, kind) in new[] { ("icon", "avatar"), ("image", "header") }) + if (document[field] is JsonObject image && image["url"]?.GetValue() is { } url && url.StartsWith(prefix, StringComparison.Ordinal)) + { + var relative = url[prefix.Length..]; + var local = Path.Combine(media.Root, relative); + if (!File.Exists(local)) + continue; + var entry = kind + Path.GetExtension(relative); + await Copy(zip, local, entry, token); + image["url"] = entry; + } + await Text(zip, "actor.json", document.ToJsonString(Options), token); + + var counts = await Outbox(zip, actor, prefix, files, token); + foreach (var relative in files) + { + var local = Path.Combine(media.Root, relative); + if (File.Exists(local)) + await Copy(zip, local, MediaFolder + relative, token); + } + + var handles = new Handles(localActors); + await Collection(zip, "likes.json", await Uris(DB.Default.Find().Match(f => f.AccountId == actor.Id).Project(f => f.PostId), token), token); + var bookmarks = await Uris(DB.Default.Find().Match(b => b.AvatarId == actor.Id).Project(b => b.PostId), token); + await Collection(zip, "bookmarks.json", bookmarks, token); + await Text(zip, "bookmarks.csv", string.Concat(bookmarks.Select(b => b + "\n")), token); + + var follows = await DB.Default.Find().Match(f => f.AvatarId == actor.Id).ExecuteAsync(token); + var csv = new StringBuilder("Account address,Show boosts,Notify on new posts,Languages\n"); + foreach (var follow in follows) + csv.Append($"{Csv(await handles.Of(follow.TargetActorURI, token))},{(follow.ShowReblogs ? "true" : "false")},false,\n"); + await Text(zip, "following_accounts.csv", csv.ToString(), token); + + csv.Clear(); + var lists = await DB.Default.Find().Match(l => l.AvatarId == actor.Id).ExecuteAsync(token); + foreach (var list in lists) + foreach (var member in await DB.Default.Find().Match(m => m.ListId == list.ID).ExecuteAsync(token)) + if (follows.FirstOrDefault(f => f.TargetAccountId == member.AccountId) is { } followed) + csv.Append($"{Csv(list.Title)},{Csv(await handles.Of(followed.TargetActorURI, token))}\n"); + await Text(zip, "lists.csv", csv.ToString(), token); + + csv.Clear(); + foreach (var block in await DB.Default.Find().Match(b => b.AvatarId == actor.Id).ExecuteAsync(token)) + csv.Append(Csv(await handles.Of(block.TargetActorURI, token)) + "\n"); + await Text(zip, "blocked_accounts.csv", csv.ToString(), token); + + csv.Clear().Append("Account address,Hide notifications\n"); + foreach (var mute in await DB.Default.Find().Match(m => m.AvatarId == actor.Id).ExecuteAsync(token)) + csv.Append($"{Csv(await handles.Of(mute.TargetActorURI, token))},{(mute.HideNotifications ? "true" : "false")}\n"); + await Text(zip, "muted_accounts.csv", csv.ToString(), token); + + csv.Clear(); + foreach (var domain in await DB.Default.Find().Match(d => d.AvatarId == actor.Id).Project(d => d.Domain).ExecuteAsync(token)) + csv.Append(domain + "\n"); + await Text(zip, "blocked_domains.csv", csv.ToString(), token); + + await PrivaPubParts(zip, actor, avatar, prefix, counts, token); + } + + // the outbox, a post at a time, newest first as Mastodon writes it: what was published and boosted + async Task<(int Posts, int Boosts, int Located)> Outbox(ZipArchive zip, LocalActor actor, string prefix, ISet files, CancellationToken token) + { + var entry = zip.CreateEntry("outbox.json", CompressionLevel.Optimal); + await using var stream = entry.Open(); + await using var writer = new Utf8JsonWriter(stream, Writer); + writer.WriteStartObject(); + writer.WriteString("@context", "https://www.w3.org/ns/activitystreams"); + writer.WriteString("id", "outbox.json"); + writer.WriteString("type", "OrderedCollection"); + writer.WritePropertyName("orderedItems"); + writer.WriteStartArray(); + var (posts, boosts) = (0, 0); + var located = 0; + var before = default(string); + while (true) + { + var find = DB.Default.Find().Match(p => p.GroupUserId == actor.Id && !p.IsFederatedCopy && !p.DeletedAt.HasValue); + if (before != default) + find = find.Match(f => f.Lt(p => p.ID, before)); + var batch = await find.Sort(p => p.ID, Order.Descending).Limit(200).ExecuteAsync(token); + if (batch.Count == 0) + break; + before = batch[^1].ID; + var originals = (await DB.Default.Find().Match(p => batch.Select(b => b.ReblogOfPostId).Contains(p.ID)).ExecuteAsync(token)) + .ToDictionary(p => p.ID); + foreach (var post in batch) + { + if (post.Visibility == PostVisibility.LocalGeo) + { + located++; + continue; + } + if (post.Visibility == PostVisibility.Circle) + continue; + if (post.ReblogOfPostId != default) + { + if (!originals.TryGetValue(post.ReblogOfPostId, out var original) || string.IsNullOrEmpty(original.ObjectURI)) + continue; + new JsonObject + { + ["id"] = actor.ActivityUri($"announce-{post.ID}"), + ["type"] = "Announce", + ["actor"] = actor.Uri, + ["published"] = ActivityPubRenderer.Timestamp(post.CreationDate), + ["to"] = new JsonArray("https://www.w3.org/ns/activitystreams#Public"), + ["cc"] = new JsonArray(actor.Followers), + ["object"] = original.ObjectURI + }.WriteTo(writer); + boosts++; + continue; + } + var group = post.GroupId == default ? default : await localActors.FindById(LocalActorKind.Group, post.GroupId, token); + var note = ActivityPubRenderer.Note(post, actor, group, post.InReplyToURI); + if (post.Visibility is PostVisibility.Direct or PostVisibility.FollowersOnly && (post.To.Count > 0 || post.Cc.Count > 0)) + { + note["to"] = new JsonArray(post.To.Select(t => (JsonNode)t).ToArray()); + note["cc"] = new JsonArray(post.Cc.Select(c => (JsonNode)c).ToArray()); + } + if (note["attachment"] is JsonArray attachments) + foreach (var attachment in attachments.OfType()) + if (attachment["url"]?.GetValue() is { } url && url.StartsWith(prefix, StringComparison.Ordinal)) + { + var relative = url[prefix.Length..]; + files.Add(relative); + attachment["url"] = "/" + MediaFolder + relative; + } + var create = ActivityPubRenderer.Create(actor, note, $"create-{post.ID}"); + create["to"] = note["to"]?.DeepClone(); + create["cc"] = note["cc"]?.DeepClone(); + create.WriteTo(writer); + posts++; + } + await writer.FlushAsync(token); + } + writer.WriteEndArray(); + writer.WriteNumber("totalItems", posts + boosts); + writer.WriteEndObject(); + return (posts, boosts, located); + } + + async Task PrivaPubParts(ZipArchive zip, LocalActor actor, Avatar avatar, string prefix, (int Posts, int Boosts, int Located) counts, CancellationToken token) + { + var filters = await DB.Default.Find().Match(f => f.AvatarId == actor.Id).ExecuteAsync(token); + var filtered = await UriMap(filters.SelectMany(f => f.Statuses.Select(s => s.PostId)), token); + await Json(zip, "privapub/filters.json", filters.Select(f => new + { + title = f.Title, + context = f.Context, + expiresAt = f.ExpiresAt, + action = f.Action, + keywords = f.Keywords.Select(k => new { keyword = k.Keyword, wholeWord = k.WholeWord }), + statuses = f.Statuses.Select(s => filtered.GetValueOrDefault(s.PostId)).Where(u => u != default) + }), token); + await Json(zip, "privapub/followed_tags.json", + await DB.Default.Find().Match(t => t.AvatarId == actor.Id).Project(t => t.Name).ExecuteAsync(token), token); + var policy = await NotificationPolicies.Of(actor.Id, token); + await Json(zip, "privapub/notification_policy.json", new + { + forNotFollowing = policy.ForNotFollowing.ToString(), + forNotFollowers = policy.ForNotFollowers.ToString(), + forNewAccounts = policy.ForNewAccounts.ToString(), + forPrivateMentions = policy.ForPrivateMentions.ToString(), + forLimitedAccounts = policy.ForLimitedAccounts.ToString() + }, token); + await Json(zip, "privapub/pins.json", + await Uris(DB.Default.Find().Match(p => p.AvatarId == actor.Id).Project(p => p.PostId), token), token); + + var scheduled = await DB.Default.Find().Match(s => s.AvatarId == actor.Id).ExecuteAsync(token); + var scheduledMedia = await MediaFiles(scheduled.SelectMany(s => s.Params.MediaIds), token); + foreach (var file in scheduledMedia.Values) + await CopyMedia(zip, file.Path, token); + await Json(zip, "privapub/scheduled.json", scheduled.Select(s => new + { + scheduledAt = s.ScheduledAt, + text = s.Params.Text, + spoilerText = s.Params.SpoilerText, + sensitive = s.Params.Sensitive, + visibility = s.Params.Visibility, + language = s.Params.Language, + poll = s.Params.Poll, + media = s.Params.MediaIds.Where(scheduledMedia.ContainsKey).Select(id => new { url = MediaFolder + scheduledMedia[id].Path, description = scheduledMedia[id].Description }) + }), token); + + var located = await DB.Default.Find().Match(p => p.GroupUserId == actor.Id && !p.IsFederatedCopy && !p.DeletedAt.HasValue + && p.Visibility == PostVisibility.LocalGeo).ExecuteAsync(token); + foreach (var post in located) + foreach (var item in post.Media.Where(m => m.URL?.StartsWith(prefix, StringComparison.Ordinal) == true)) + await CopyMedia(zip, item.URL[prefix.Length..], token); + await Json(zip, "privapub/located.json", located.Select(p => new + { + id = p.ObjectURI, + published = p.CreationDate, + title = p.Title, + text = p.Text, + spoilerText = p.SpoilerText, + latitude = p.Geo?.Coordinates[1], + longitude = p.Geo?.Coordinates[0], + rangeKm = p.RangeKm, + media = p.Media.Where(m => m.URL?.StartsWith(prefix, StringComparison.Ordinal) == true) + .Select(m => new { url = MediaFolder + m.URL[prefix.Length..], description = m.Description }) + }), token); + + await Json(zip, "privapub/archive.json", new + { + format = Format, + software = "privapub", + actor = actor.Uri, + username = avatar?.UserName ?? actor.UserName, + exportedAt = DateTime.UtcNow, + posts = counts.Posts, + boosts = counts.Boosts, + located = counts.Located, + scheduled = scheduled.Count + }, token); + } + + async Task CopyMedia(ZipArchive zip, string relative, CancellationToken token) + { + var local = Path.Combine(media.Root, relative); + if (File.Exists(local) && zip.GetEntry(MediaFolder + relative) == default) + await Copy(zip, local, MediaFolder + relative, token); + } + + static async Task> MediaFiles(IEnumerable ids, CancellationToken token) + { + var wanted = ids.Distinct().ToList(); + if (wanted.Count == 0) + return []; + return (await DB.Default.Find().Match(m => wanted.Contains(m.ID) && m.TrashedAt == null).ExecuteAsync(token)) + .ToDictionary(m => m.ID, m => (m.FilePath, m.Description)); + } + + static async Task> UriMap(IEnumerable postIds, CancellationToken token) + { + var ids = postIds.Where(id => id != default).Distinct().ToList(); + if (ids.Count == 0) + return []; + return (await DB.Default.Find().Match(p => ids.Contains(p.ID)).Project(p => new PostEntity { ID = p.ID, ObjectURI = p.ObjectURI }).ExecuteAsync(token)) + .Where(p => !string.IsNullOrEmpty(p.ObjectURI)).ToDictionary(p => p.ID, p => p.ObjectURI); + } + + static async Task> Uris(Find postIds, CancellationToken token) where T : IEntity + { + var ids = await postIds.ExecuteAsync(token); + var map = await UriMap(ids, token); + return ids.Where(map.ContainsKey).Select(id => map[id]).ToList(); + } + + static async Task Collection(ZipArchive zip, string name, List items, CancellationToken token) => + await Text(zip, name, new JsonObject + { + ["@context"] = "https://www.w3.org/ns/activitystreams", + ["id"] = name, + ["type"] = "OrderedCollection", + ["totalItems"] = items.Count, + ["orderedItems"] = new JsonArray(items.Select(i => (JsonNode)i).ToArray()) + }.ToJsonString(Options), token); + + static Task Json(ZipArchive zip, string name, T value, CancellationToken token) => + Text(zip, name, JsonSerializer.Serialize(value, Options), token); + + static async Task Text(ZipArchive zip, string name, string text, CancellationToken token) + { + await using var stream = zip.CreateEntry(name, CompressionLevel.Optimal).Open(); + await stream.WriteAsync(Encoding.UTF8.GetBytes(text), token); + } + + static async Task Copy(ZipArchive zip, string path, string name, CancellationToken token) + { + await using var source = new FileStream(path, FileMode.Open, FileAccess.Read, FileShare.Read, 81920, FileOptions.SequentialScan); + await using var target = zip.CreateEntry(name, CompressionLevel.NoCompression).Open(); + await source.CopyToAsync(target, token); + } + + // a CSV cell, quoted when it must be + static string Csv(string value) => + value.IndexOfAny([',', '"', '\n', '\r']) < 0 ? value : "\"" + value.Replace("\"", "\"\"") + "\""; + + // an account as Mastodon's CSV files name it: user@host, or its address when the handle is not known + sealed class Handles(ILocalActorService localActors) + { + readonly Dictionary _known = new(StringComparer.Ordinal); + + public async Task Of(string actorUri, CancellationToken token) + { + if (_known.TryGetValue(actorUri, out var handle)) + return handle; + if (await localActors.FindByUri(actorUri, token) is { } local) + handle = local.Handle; + else if (await DB.Default.Find().Match(a => a.ActorURI == actorUri).ExecuteFirstAsync(token) is { UserName.Length: > 0 } remote + && Uri.TryCreate(actorUri, UriKind.Absolute, out var parsed)) + handle = $"{remote.UserName}@{(string.IsNullOrEmpty(remote.Domain) ? parsed.Authority : remote.Domain)}"; + else + handle = actorUri; + return _known[actorUri] = handle; + } + } + } +} diff --git a/PrivaPub/Domain/Portability/ImportArchiveJob.cs b/PrivaPub/Domain/Portability/ImportArchiveJob.cs new file mode 100644 index 0000000..73c7b33 --- /dev/null +++ b/PrivaPub/Domain/Portability/ImportArchiveJob.cs @@ -0,0 +1,928 @@ +using Microsoft.AspNetCore.Http; + +using MongoDB.Driver; +using MongoDB.Entities; + +using PrivaPub.Api.Mastodon.Controllers; +using PrivaPub.Domain.Media; +using PrivaPub.Domain.Relationships; +using PrivaPub.Domain.Social; +using PrivaPub.Domain.Statuses; +using PrivaPub.Federation.Actors; +using PrivaPub.Federation.Inbox; +using PrivaPub.Federation.Moderation; +using PrivaPub.Federation.Objects; +using PrivaPub.Federation.Outbox; +using PrivaPub.Federation.Rendering; +using PrivaPub.Infrastructure.Backup; +using PrivaPub.Infrastructure.Ids; +using PrivaPub.Infrastructure.Jobs; +using PrivaPub.Models.Federation; +using PrivaPub.Models.Jobs; +using PrivaPub.Models.Media; +using PrivaPub.Models.Post; +using PrivaPub.Models.Social; +using PrivaPub.Models.User; + +using System.Net; +using System.Text.Json.Nodes; +using System.Text.RegularExpressions; + +using PostEntity = PrivaPub.Models.Post.Post; + +namespace PrivaPub.Domain.Portability +{ + // An archive restored into one of a root's personas (owner decision 2026-10-07): PrivaPub's own, or Mastodon's account + // archive. Only the parts asked for, one at a time, with progress, until asked to stop. + // Posts: only the archive actor's own Creates (boosts, direct and circle posts are counted, not imported); delivered to + // no one, put in no home, notifying nobody, they show on the profile, the outbox, hashtags and search. From the same + // actor (a persona's own archive, back home) a post keeps its id and address, and one already here (even deleted) or + // whose address is a tombstone is left alone; from another actor it gets an id of its date and its address here, with + // the original in ImportedFromURI, so importing twice changes nothing. Replies and self-quotes inside the archive point + // at their copies; polls are closed, with no votes; mentions stay links; the HTML is sanitized as any remote post's; + // media go through MediaService, its checks and the root's quota. + // The rest: follows asked again (not on a blocked server), blocks, mutes, blocked servers, lists (with only accounts + // followed), bookmarks (read by signed fetch), filters, followed hashtags, the notification policy, pins (not + // announced), the profile; and only when asked, located and scheduled posts (PrivaPub's archives) and likes, which tell + // their authors. Followers are never imported: they choose. + public partial class ImportArchiveJob(IServiceScopeFactory scopes, ILocalActorService localActors, IMediaService media, IDomainBlocks domainBlocks, + IJobQueue jobs, ILogger logger) : IJobHandler + { + static readonly HashSet Public = ["https://www.w3.org/ns/activitystreams#Public", "as:Public", "Public"]; + static readonly HashSet Postable = ["Note", "Article", "Page", "Question"]; + const int MaxMedia = 4; + + public JobKind Kind => JobKind.ImportArchive; + public int Concurrency => 1; + public int MaxAttempts => 1; + public int PerHostLimit => 1; + + public static Job JobFor(string avatarId, string host) => new() + { + Kind = JobKind.ImportArchive, + Payload = avatarId, + Host = host, + DedupeKey = $"import|{avatarId}|{DateTime.UtcNow.Ticks}" + }; + + sealed class Run(PersonaArchive state, LocalActor me, SafeArchive archive, string archiveActor, bool sameActor, IServiceProvider services) + { + public PersonaArchive State { get; } = state; + public LocalActor Me { get; } = me; + public SafeArchive Archive { get; } = archive; + public string ArchiveActor { get; } = archiveActor; + public bool SameActor { get; } = sameActor; + public IServiceProvider Services { get; } = services; + public Dictionary Posts { get; } = new(StringComparer.Ordinal);//archive object id → here + public DateTime LastSaved { get; set; } + + public void Count(string what, int by = 1) => State.ImportCounts[what] = State.ImportCounts.GetValueOrDefault(what) + by; + } + + public async Task Handle(Job job, CancellationToken token) + { + var state = await DB.Default.Find().MatchID(job.Payload).ExecuteFirstAsync(token); + if (state is not { ImportState: "queued" or "running" }) + return JobOutcome.Done; + var me = await localActors.FindById(LocalActorKind.Person, job.Payload, token); + var path = state.ImportFile; + try + { + if (me == default || !await RootMayImport(me.Id, token)) + { + await End(state, "failed", "the persona or its login can no longer import", token); + return JobOutcome.Done; + } + var (archive, refused) = path == default ? (default, "no archive") : SafeArchive.Open(path); + if (archive == default) + { + await End(state, "failed", refused, token); + return JobOutcome.Done; + } + using (archive) + { + var actor = archive.Json("actor.json") as JsonObject; + var actorId = Text(actor?["id"]); + if (actorId == default) + { + await End(state, "failed", "the archive has no actor.json", token); + return JobOutcome.Done; + } + using var scope = scopes.CreateScope(); + var run = new Run(state, me, archive, actorId, actorId == me.Uri, scope.ServiceProvider); + state.ImportState = "running"; + state.ImportFrom = actorId; + state.ImportSameActor = run.SameActor; + state.ImportCounts = []; + await Save(run, token, force: true); + + var parts = state.ImportParts.ToHashSet(StringComparer.Ordinal); + var privapub = Text((archive.Json("privapub/archive.json") as JsonObject)?["software"]) == "privapub"; + if (parts.Contains("profile") && !await Stopped(run, token)) + await Profile(run, actor, token); + if (parts.Contains("posts") && !await Stopped(run, token)) + await Posts(run, token); + else + await KnownPosts(run, token); + if (parts.Contains("located") && privapub && !await Stopped(run, token)) + await Located(run, token); + if (parts.Contains("scheduled") && privapub && !await Stopped(run, token)) + await Scheduled(run, token); + if (parts.Contains("tags") && !await Stopped(run, token)) + await Tags(run, token); + if (parts.Contains("notificationpolicy") && !await Stopped(run, token)) + await Policy(run, token); + if (parts.Contains("filters") && !await Stopped(run, token)) + await Filters(run, token); + if (parts.Contains("pins") && !await Stopped(run, token)) + await Pins(run, token); + if (parts.Contains("domainblocks") && !await Stopped(run, token)) + await DomainBlocks(run, token); + if (parts.Contains("blocks") && !await Stopped(run, token)) + await Accounts(run, "blocked_accounts.csv", "blocks", Block, token); + if (parts.Contains("mutes") && !await Stopped(run, token)) + await Accounts(run, "muted_accounts.csv", "mutes", Mute, token); + if (parts.Contains("follows") && !await Stopped(run, token)) + await Accounts(run, "following_accounts.csv", "follows", Follow, token); + if (parts.Contains("lists") && !await Stopped(run, token)) + await Lists(run, token); + if (parts.Contains("bookmarks") && !await Stopped(run, token)) + await Bookmarks(run, token); + if (parts.Contains("likes") && !await Stopped(run, token)) + await Likes(run, token); + await End(state, await Stopped(run, token) ? "stopped" : "done", default, token); + } + } + catch (Exception ex) when (ex is not OperationCanceledException) + { + logger.LogError(ex, "The import into {Persona} failed", me?.Handle ?? job.Payload); + await End(state, "failed", ex.Message, CancellationToken.None); + } + finally + { + if (path != default && File.Exists(path)) + File.Delete(path); + } + return JobOutcome.Done; + } + + static async Task RootMayImport(string avatarId, CancellationToken token) + { + var link = await DB.Default.Find().Match(r => r.AvatarId == avatarId).ExecuteFirstAsync(token); + var root = link == default ? default : await DB.Default.Find().MatchID(link.RootId).ExecuteFirstAsync(token); + return root is { IsBanned: false, DeletedAt: null }; + } + + static async Task End(PersonaArchive state, string how, string error, CancellationToken token) => + await DB.Default.Update().MatchID(state.ID) + .Modify(a => a.ImportState, how).Modify(a => a.ImportError, error).Modify(a => a.ImportEndedAt, DateTime.UtcNow) + .Modify(a => a.ImportCounts, state.ImportCounts).Modify(a => a.ImportDone, state.ImportDone).Modify(a => a.ImportTotal, state.ImportTotal) + .Modify(a => a.ImportFile, null).Modify(a => a.StopAsked, false) + .ExecuteAsync(token); + + // progress written at most every two seconds; the answer to whether the root asked to stop + static async Task Save(Run run, CancellationToken token, bool force = false) + { + if (!force && DateTime.UtcNow - run.LastSaved < TimeSpan.FromSeconds(2)) + return; + run.LastSaved = DateTime.UtcNow; + await DB.Default.Update().MatchID(run.State.ID) + .Modify(a => a.ImportState, run.State.ImportState).Modify(a => a.ImportFrom, run.State.ImportFrom) + .Modify(a => a.ImportSameActor, run.State.ImportSameActor).Modify(a => a.ImportCounts, run.State.ImportCounts) + .Modify(a => a.ImportDone, run.State.ImportDone).Modify(a => a.ImportTotal, run.State.ImportTotal) + .ExecuteAsync(token); + } + + static async Task Stopped(Run run, CancellationToken token) + { + await Save(run, token); + return await DB.Default.Find().MatchID(run.State.ID).Match(a => a.StopAsked).ExecuteAnyAsync(token); + } + + // ---- posts + + // what an outbox item is to this import: the object to import, or why not + (JsonObject Note, string Skipped) Importable(Run run, JsonObject item) + { + var type = Text(item["type"]); + if (type == "Announce") + return (default, "boosts skipped"); + if (type != "Create" || item["object"] is not JsonObject note || !Postable.Contains(Text(note["type"]) ?? string.Empty)) + return (default, "not a post"); + if (Text(item["actor"]) != run.ArchiveActor || Id(note["attributedTo"]) is { } author && author != run.ArchiveActor) + return (default, "not the archive's own"); + if (Text(note["id"]) == default) + return (default, "not a post"); + return Visibility(run, note) == default ? (default, "direct or circle posts skipped") : (note, default); + } + + PostVisibility? Visibility(Run run, JsonObject note) + { + var to = Strings(note["to"]); + var cc = Strings(note["cc"]); + if (to.Any(Public.Contains)) + return PostVisibility.Public; + if (cc.Any(Public.Contains)) + return PostVisibility.Unlisted; + var followers = run.ArchiveActor.TrimEnd('/') is var actor && run.Archive.Json("actor.json") is JsonObject document + ? Text(document["followers"]) ?? actor + "/followers" + : default; + if (followers != default && to.Concat(cc).Contains(followers) && Text(note["audience"]) == default) + return PostVisibility.FollowersOnly; + return default; + } + + // what is already here of the archive's posts, for pins, filters and bookmarks when posts are not imported again + async Task KnownPosts(Run run, CancellationToken token) + { + await foreach (var item in run.Archive.Outbox(token)) + { + if (Importable(run, item).Note is not { } note) + continue; + var original = Text(note["id"]); + var here = run.SameActor + ? await DB.Default.Find().Match(p => p.ObjectURI == original && p.GroupUserId == run.Me.Id).ExecuteFirstAsync(token) + : await DB.Default.Find().Match(p => p.GroupUserId == run.Me.Id && p.ImportedFromURI == original).ExecuteFirstAsync(token); + if (here is { DeletedAt: null }) + run.Posts[original] = (here.ID, here.ObjectURI); + } + } + + async Task Posts(Run run, CancellationToken token) + { + // first the ids each post will have here, so a reply or a quote anywhere in the archive finds its copy + var planned = new Dictionary(StringComparer.Ordinal); + var total = 0L; + await foreach (var item in run.Archive.Outbox(token)) + { + total++; + if (Importable(run, item).Note is not { } note) + continue; + var original = Text(note["id"]); + var published = Date(note["published"]) ?? DateTime.UtcNow; + if (run.SameActor) + { + var id = original.TrimEnd('/').Split('/')[^1]; + if (!IsObjectId(id) || original != run.Me.PostUri(id)) + id = PrivacyIds.At(Clamp(published)); + var here = await DB.Default.Find().Match(p => p.ID == id || p.ObjectURI == original).ExecuteFirstAsync(token); + var tombstone = await DB.Default.Find().Match(d => d.ObjectURI == original).ExecuteAnyAsync(token); + if (here != default || tombstone) + { + if (here is { DeletedAt: null, GroupUserId: var owner } && owner == run.Me.Id) + run.Posts[original] = (here.ID, here.ObjectURI); + run.Count(here is { DeletedAt: not null } || tombstone ? "posts deleted here" : "posts already here"); + continue; + } + planned[original] = (id, original, true); + } + else + { + var here = await DB.Default.Find().Match(p => p.GroupUserId == run.Me.Id && p.ImportedFromURI == original).ExecuteFirstAsync(token); + if (here != default) + { + if (here.DeletedAt == default) + run.Posts[original] = (here.ID, here.ObjectURI); + run.Count("posts already here"); + continue; + } + var id = PrivacyIds.At(Clamp(published)); + planned[original] = (id, run.Me.PostUri(id), true); + } + } + foreach (var (original, (id, uri, _)) in planned) + run.Posts[original] = (id, uri); + run.State.ImportTotal = total; + run.State.ImportDone = 0; + + await foreach (var item in run.Archive.Outbox(token)) + { + run.State.ImportDone++; + if (run.State.ImportDone % 50 == 0 && await Stopped(run, token)) + return; + var (note, skipped) = Importable(run, item); + if (note == default) + { + run.Count(skipped); + continue; + } + var original = Text(note["id"]); + if (!planned.TryGetValue(original, out var target)) + continue; + await Post(run, note, original, target.Id, target.Uri, Visibility(run, note).Value, token); + run.Count("posts"); + } + } + + async Task Post(Run run, JsonObject note, string original, string id, string uri, PostVisibility visibility, CancellationToken token) + { + var me = run.Me; + var summary = Clean(Text(note["summary"])); + var content = Text(note["content"]) ?? (note["contentMap"] as JsonObject)?.Select(c => Text(c.Value)).FirstOrDefault(); + var source = note["source"] as JsonObject; + var post = new PostEntity + { + ID = id, + GroupUserId = me.Id, + AuthorAccountId = me.Id, + ActorURI = me.Uri, + Visibility = visibility, + Title = Clean(Text(note["name"])), + SpoilerText = summary, + HasContentWarning = summary != default || note["sensitive"]?.GetValueKind() == System.Text.Json.JsonValueKind.True, + ContentHtml = content == default ? default : ContentSanitizer.Html(content), + Text = Text(source?["mediaType"]) is "text/plain" or "text/markdown" ? Text(source["content"]) : default, + ContentFormat = ContentFormat.Html, + Language = (note["contentMap"] as JsonObject)?.Select(c => c.Key).FirstOrDefault(), + Tags = (note["tag"] as JsonArray ?? []) + .OfType().Where(t => Text(t["type"]) == "Hashtag").Select(t => TagsController.Normalise(Text(t["name"]))) + .Where(t => t.Length > 0).Distinct().ToList(), + ObjectURI = uri, + Url = me.PostHtmlUrl(id), + ActivityURI = me.ActivityUri("create-" + id), + ImportedFromURI = run.SameActor ? default : original, + CreationDate = Clamp(Date(note["published"]) ?? DateTime.UtcNow), + LocalQuotePolicy = me.Settings.QuotePolicy ?? QuotePolicies.Public + }; + if (Id(note["inReplyTo"]) is { } inReplyTo) + { + if (run.Posts.TryGetValue(inReplyTo, out var parent)) + (post.AnsweringToPostId, post.InReplyToURI, post.InReplyToAccountId) = (parent.Id, parent.Uri, me.Id); + else + { + post.InReplyToURI = inReplyTo; + if (await DB.Default.Find().Match(p => p.ObjectURI == inReplyTo).ExecuteFirstAsync(token) is { } held) + (post.AnsweringToPostId, post.InReplyToAccountId) = (held.ID, held.AuthorAccountId ?? held.GroupUserId); + } + } + foreach (var key in new[] { "quote", "quoteUrl", "quoteUri", "_misskey_quote" }) + if (Text(note[key]) is { } quoted && run.Posts.TryGetValue(quoted, out var copy)) + { + (post.QuotedPostId, post.QuoteURI, post.QuoteState, post.QuoteByConsent) = (copy.Id, copy.Uri, QuoteState.Accepted, true); + break; + } + if (Text(note["type"]) == "Question") + post.Poll = Poll(note); + if (visibility is PostVisibility.Public or PostVisibility.Unlisted) + post.ContextURI = post.AnsweringToPostId != default && run.Posts.ContainsValue((post.AnsweringToPostId, post.InReplyToURI)) + ? post.InReplyToURI + "/context" + : post.ObjectURI + "/context"; + post.Media = await Media(run, note["attachment"] as JsonArray, id, token); + var rendered = ActivityPubRenderer.Note(post, me, default, post.InReplyToURI); + post.To = Strings(rendered["to"]); + post.Cc = Strings(rendered["cc"]); + await DB.Default.SaveAsync(post, token); + if (post.AnsweringToPostId != default) + await DB.Default.Update().MatchID(post.AnsweringToPostId).Modify(b => b.Inc(p => p.RepliesCount, 1)).ExecuteAsync(token); + } + + static PostPoll Poll(JsonObject note) + { + var multiple = note["anyOf"] is JsonArray; + var options = (note["anyOf"] ?? note["oneOf"]) as JsonArray ?? []; + var ends = Date(note["endTime"]) ?? Date(note["closed"]) ?? DateTime.UtcNow; + var now = DateTime.UtcNow; + return new PostPoll + { + Options = options.OfType().Select(o => new PollOption { Title = Text(o["name"]) }).Where(o => o.Title != default).ToList(), + Multiple = multiple, + ExpiresAt = ends < now ? ends : now, + ClosedAt = ends < now ? ends : now, + VotersCount = 0 + }; + } + + // the attachments the archive holds, uploaded as the persona's own (checked, processed, counted against its quota) + async Task> Media(Run run, JsonArray attachments, string postId, CancellationToken token) + { + var made = new List(); + foreach (var attachment in (attachments ?? []).OfType().Take(MaxMedia)) + { + var url = Text(attachment["url"]) ?? Text((attachment["url"] as JsonArray)?.OfType().FirstOrDefault()?["href"]); + var file = url?.TrimStart('/'); + if (file == default || file.Contains("://") || !run.Archive.Has(file)) + { + run.Count("media not in the archive"); + continue; + } + var focus = attachment["focalPoint"] is JsonArray { Count: 2 } point ? $"{point[0]},{point[1]}" : default; + var outcome = await Upload(run, file, Text(attachment["mediaType"]), Clean(Text(attachment["name"])), focus, token); + if (outcome?.Ok != true) + { + run.Count("media refused"); + continue; + } + var row = outcome.Attachment; + await DB.Default.Update().MatchID(row.ID).Modify(m => m.PostId, postId).Modify(m => m.AttachedAt, DateTime.UtcNow).ExecuteAsync(token); + made.Add(new PostMedia + { + AttachmentId = row.ID, + ContentType = row.ContentType, + Kind = row.Kind, + DurationSeconds = row.DurationSeconds, + URL = media.Url(row.FilePath), + PreviewURL = media.Url(row.PreviewPath ?? row.FilePath), + Description = row.Description, + Blurhash = row.Blurhash, + Width = row.Width, + Height = row.Height, + Focus = row.Focus + }); + run.Count("media"); + } + return made; + } + + // a file of the archive given to MediaService as an upload would be + async Task Upload(Run run, string file, string contentType, string description, string focus, CancellationToken token) + { + var temporary = Path.Combine(Path.GetTempPath(), $"privapub-import-{Guid.NewGuid():N}{Path.GetExtension(file)}"); + try + { + await using (var source = run.Archive.Read(file)) + await using (var target = new FileStream(temporary, FileMode.CreateNew, FileAccess.Write)) + await source.CopyToAsync(target, token); + await using var stream = new FileStream(temporary, FileMode.Open, FileAccess.Read); + var form = new FormFile(stream, 0, stream.Length, "file", Path.GetFileName(file)) + { + Headers = new HeaderDictionary(), + ContentType = contentType ?? (MediaService.ContentTypes.TryGetContentType(file, out var guessed) ? guessed : "application/octet-stream") + }; + return await media.Upload(run.Me, form, description, focus, later: false, token); + } + finally + { + File.Delete(temporary); + } + } + + // ---- PrivaPub's own parts + + async Task Located(Run run, CancellationToken token) + { + foreach (var located in (run.Archive.Json("privapub/located.json") as JsonArray ?? []).OfType()) + { + var original = Text(located["id"]); + if (original != default && await DB.Default.Find().Match(p => p.GroupUserId == run.Me.Id + && (p.ObjectURI == original || p.ImportedFromURI == original)).ExecuteAnyAsync(token)) + { + run.Count("located posts already here"); + continue; + } + if (located["latitude"]?.GetValue() is not { } latitude || located["longitude"]?.GetValue() is not { } longitude) + continue; + var id = PrivacyIds.At(Clamp(Date(located["published"]) ?? DateTime.UtcNow)); + var post = new PostEntity + { + ID = id, + GroupUserId = run.Me.Id, + AuthorAccountId = run.Me.Id, + ActorURI = run.Me.Uri, + Visibility = PostVisibility.LocalGeo, + IsLocalOnly = true, + Title = Clean(Text(located["title"])), + SpoilerText = Clean(Text(located["spoilerText"])), + HasContentWarning = Clean(Text(located["spoilerText"])) != default, + Text = Text(located["text"]), + ContentHtml = ActivityPubRenderer.Html(Text(located["text"]) ?? string.Empty), + ContentFormat = ContentFormat.Markdown, + Geo = new MongoDB.Entities.Coordinates2D(Math.Round(longitude, 2), Math.Round(latitude, 2)), + RangeKm = (float)Math.Clamp(located["rangeKm"]?.GetValue() ?? 5, 1, 50), + ObjectURI = run.Me.PostUri(id), + Url = run.Me.PostHtmlUrl(id), + ImportedFromURI = original == run.Me.PostUri(original?.Split('/')[^1] ?? string.Empty) ? default : original, + CreationDate = Clamp(Date(located["published"]) ?? DateTime.UtcNow) + }; + var files = (located["media"] as JsonArray ?? []).OfType() + .Select(m => new JsonObject { ["url"] = Text(m["url"]), ["name"] = Text(m["description"]) }).ToArray(); + post.Media = await Media(run, new JsonArray(files), id, token); + await DB.Default.SaveAsync(post, token); + run.Count("located posts"); + } + } + + async Task Scheduled(Run run, CancellationToken token) + { + var made = new List(); + foreach (var scheduled in (run.Archive.Json("privapub/scheduled.json") as JsonArray ?? []).OfType()) + { + var at = Date(scheduled["scheduledAt"]); + if (at is not { } when || when < DateTime.UtcNow.AddMinutes(5) || await ScheduledStatuses.Refusal(run.Me.Id, when, default, token) != default) + { + run.Count("scheduled posts past or refused"); + continue; + } + var status = new ScheduledStatus + { + AvatarId = run.Me.Id, + ScheduledAt = when, + Params = new ScheduledParams + { + Text = Text(scheduled["text"]), + SpoilerText = Text(scheduled["spoilerText"]), + Sensitive = scheduled["sensitive"]?.GetValueKind() == System.Text.Json.JsonValueKind.True, + Visibility = Text(scheduled["visibility"]) is { } visibility && visibility is "public" or "unlisted" or "private" or "direct" ? visibility : "public", + Language = Text(scheduled["language"]) + } + }; + status.ID = (string)status.GenerateNewID(); + foreach (var file in (scheduled["media"] as JsonArray ?? []).OfType().Take(MaxMedia)) + { + var path = Text(file["url"])?.TrimStart('/'); + if (path == default || !run.Archive.Has(path)) + continue; + var outcome = await Upload(run, path, default, Clean(Text(file["description"])), default, token); + if (outcome?.Ok != true) + continue; + await DB.Default.Update().MatchID(outcome.Attachment.ID).Modify(m => m.ScheduledStatusId, status.ID).ExecuteAsync(token); + status.Params.MediaIds.Add(outcome.Attachment.ID); + } + await DB.Default.SaveAsync(status, token); + made.Add(status); + run.Count("scheduled posts"); + } + if (made.Count > 0) + await jobs.EnqueueMany(made.Select(s => ScheduledStatuses.JobFor(s, new Uri(localActors.BaseAddress).Host)), token); + } + + async Task Tags(Run run, CancellationToken token) + { + var mine = (await DB.Default.Find().Match(t => t.AvatarId == run.Me.Id).Project(t => t.Name).ExecuteAsync(token)).ToHashSet(); + foreach (var name in (run.Archive.Json("privapub/followed_tags.json") as JsonArray ?? []).Select(t => TagsController.Normalise(Text(t)))) + { + if (name.Length == 0 || mine.Contains(name) || mine.Count >= 500) + continue; + await DB.Default.SaveAsync(new FollowedTag { AvatarId = run.Me.Id, Name = name }, token); + mine.Add(name); + run.Count("followed hashtags"); + } + } + + async Task Policy(Run run, CancellationToken token) + { + if (run.Archive.Json("privapub/notification_policy.json") is not JsonObject json) + return; + var policy = await NotificationPolicies.Of(run.Me.Id, token); + NotificationAction Action(string field, NotificationAction kept) => Enum.TryParse(Text(json[field]), out var action) ? action : kept; + policy.ForNotFollowing = Action("forNotFollowing", policy.ForNotFollowing); + policy.ForNotFollowers = Action("forNotFollowers", policy.ForNotFollowers); + policy.ForNewAccounts = Action("forNewAccounts", policy.ForNewAccounts); + policy.ForPrivateMentions = Action("forPrivateMentions", policy.ForPrivateMentions); + policy.ForLimitedAccounts = Action("forLimitedAccounts", policy.ForLimitedAccounts); + await NotificationPolicies.Save(policy, token); + run.Count("notification policy"); + } + + async Task Filters(Run run, CancellationToken token) + { + var titles = (await DB.Default.Find().Match(f => f.AvatarId == run.Me.Id).Project(f => f.Title).ExecuteAsync(token)).ToHashSet(); + foreach (var filter in (run.Archive.Json("privapub/filters.json") as JsonArray ?? []).OfType()) + { + var title = Clean(Text(filter["title"])); + if (title == default || titles.Contains(title)) + continue; + var statuses = new List(); + foreach (var uri in Strings(filter["statuses"])) + if (await Held(run, uri, token) is { } post) + statuses.Add(new FilterStatus { Id = MongoDB.Bson.ObjectId.GenerateNewId().ToString(), PostId = post.ID }); + await DB.Default.SaveAsync(new PersonaFilter + { + AvatarId = run.Me.Id, + Title = title, + Context = Strings(filter["context"]).Where(c => c is "home" or "notifications" or "public" or "thread" or "account").ToList(), + ExpiresAt = Date(filter["expiresAt"]), + Action = Text(filter["action"]) is { } action && action is "warn" or "hide" or "blur" ? action : "warn", + Keywords = (filter["keywords"] as JsonArray ?? []).OfType().Where(k => Clean(Text(k["keyword"])) != default) + .Select(k => new FilterKeyword { Id = MongoDB.Bson.ObjectId.GenerateNewId().ToString(), Keyword = Text(k["keyword"]).Trim(), + WholeWord = k["wholeWord"]?.GetValueKind() == System.Text.Json.JsonValueKind.True }).ToList(), + Statuses = statuses + }, token); + titles.Add(title); + run.Count("filters"); + } + } + + // pinned again without announcing it (an Add would tell every follower's server) + async Task Pins(Run run, CancellationToken token) + { + var pinned = (await DB.Default.Find().Match(p => p.AvatarId == run.Me.Id).Project(p => p.PostId).ExecuteAsync(token)).ToHashSet(); + foreach (var uri in Strings(run.Archive.Json("privapub/pins.json"))) + { + if (pinned.Count >= 5) + break; + if (!run.Posts.TryGetValue(uri, out var copy) || pinned.Contains(copy.Id)) + continue; + await DB.Default.SaveAsync(new Pin { AvatarId = run.Me.Id, PostId = copy.Id }, token); + pinned.Add(copy.Id); + run.Count("pins"); + } + } + + // ---- relationships + + async Task DomainBlocks(Run run, CancellationToken token) + { + var relationships = run.Services.GetRequiredService(); + var mine = (await DB.Default.Find().Match(d => d.AvatarId == run.Me.Id).Project(d => d.Domain).ExecuteAsync(token)).ToHashSet(); + foreach (var line in run.Archive.Lines("blocked_domains.csv")) + { + var domain = Cells(line).FirstOrDefault()?.Trim().ToLowerInvariant(); + if (string.IsNullOrEmpty(domain) || domain == "#domain" || Uri.CheckHostName(domain) != UriHostNameType.Dns || mine.Contains(domain)) + continue; + await relationships.BlockDomain(run.Me, domain, token); + mine.Add(domain); + run.Count("blocked servers"); + } + } + + // one CSV of accounts, each resolved (WebFinger or its address) and acted on; a header line is skipped + async Task Accounts(Run run, string file, string what, Func> act, + CancellationToken token) + { + var follows = run.Services.GetRequiredService(); + var lines = run.Archive.Lines(file); + foreach (var line in lines) + { + if (await Stopped(run, token)) + return; + var cells = Cells(line); + var address = cells.FirstOrDefault()?.Trim(); + if (string.IsNullOrEmpty(address) || address.Equals("Account address", StringComparison.OrdinalIgnoreCase)) + continue; + var target = await follows.Resolve(address, token); + if (target.Local == default && target.Remote == default) + { + run.Count($"{what} not found"); + continue; + } + if (target.Local?.Id == run.Me.Id) + continue; + run.Count(await act(run, target, cells, token) ? what : $"{what} skipped"); + } + } + + static async Task Block(Run run, (LocalActor Local, ForeignAvatar Remote) target, string[] cells, CancellationToken token) + { + var (uri, id) = target.Local != default ? (target.Local.Uri, target.Local.Id) : (target.Remote.ActorURI, target.Remote.ID); + if (await DB.Default.Find().Match(b => b.AvatarId == run.Me.Id && b.TargetActorURI == uri).ExecuteAnyAsync(token)) + return false; + await run.Services.GetRequiredService().Block(run.Me, uri, id, token); + return true; + } + + static async Task Mute(Run run, (LocalActor Local, ForeignAvatar Remote) target, string[] cells, CancellationToken token) + { + var (uri, id) = target.Local != default ? (target.Local.Uri, target.Local.Id) : (target.Remote.ActorURI, target.Remote.ID); + var hide = cells.Length < 2 || !cells[1].Trim().Equals("false", StringComparison.OrdinalIgnoreCase); + await run.Services.GetRequiredService().Mute(run.Me, uri, id, hide, default, token); + return true; + } + + async Task Follow(Run run, (LocalActor Local, ForeignAvatar Remote) target, string[] cells, CancellationToken token) + { + var uri = target.Local?.Uri ?? target.Remote.ActorURI; + var host = new Uri(uri).Host; + if (domainBlocks.IsSuspended(host) + || await DB.Default.Find().Match(d => d.AvatarId == run.Me.Id && d.Domain == host).ExecuteAnyAsync(token) + || await DB.Default.Find().Match(f => f.AvatarId == run.Me.Id && f.TargetActorURI == uri).ExecuteAnyAsync(token)) + return false; + var boosts = cells.Length < 2 || !cells[1].Trim().Equals("false", StringComparison.OrdinalIgnoreCase); + return await run.Services.GetRequiredService().FollowAs(run.Me, uri, boosts, token) != default; + } + + // lists by title, with only the accounts the persona follows + async Task Lists(Run run, CancellationToken token) + { + var follows = run.Services.GetRequiredService(); + var lists = (await DB.Default.Find().Match(l => l.AvatarId == run.Me.Id).ExecuteAsync(token)).ToDictionary(l => l.Title, StringComparer.Ordinal); + foreach (var line in run.Archive.Lines("lists.csv")) + { + var cells = Cells(line); + if (cells.Length < 2 || Clean(cells[0]) is not { } title) + continue; + if (!lists.TryGetValue(title, out var list)) + { + if (lists.Count >= 50) + continue; + list = new PersonaList { AvatarId = run.Me.Id, Title = title, RepliesPolicy = "list" }; + await DB.Default.SaveAsync(list, token); + lists[title] = list; + run.Count("lists"); + } + var target = await follows.Resolve(cells[1].Trim(), token); + var uri = target.Local?.Uri ?? target.Remote?.ActorURI; + var following = uri == default ? default + : await DB.Default.Find().Match(f => f.AvatarId == run.Me.Id && f.TargetActorURI == uri && f.State == FollowState.Accepted).ExecuteFirstAsync(token); + if (following == default) + { + run.Count("list members not followed"); + continue; + } + if (await DB.Default.Find().Match(m => m.ListId == list.ID && m.AccountId == following.TargetAccountId).ExecuteAnyAsync(token)) + continue; + await DB.Default.SaveAsync(new PersonaListMember { ListId = list.ID, AvatarId = run.Me.Id, AccountId = following.TargetAccountId }, token); + run.Count("list members"); + } + } + + async Task Bookmarks(Run run, CancellationToken token) + { + var uris = Strings(run.Archive.Json("bookmarks.json")?["orderedItems"]); + if (uris.Count == 0) + uris = run.Archive.Lines("bookmarks.csv").Select(l => Cells(l).FirstOrDefault()?.Trim()).Where(u => u?.StartsWith("http") == true).ToList(); + foreach (var uri in uris.Distinct()) + { + if (await Stopped(run, token)) + return; + var post = await Fetched(run, uri, token); + if (post == default) + { + run.Count("bookmarks not found"); + continue; + } + if (await DB.Default.Find().Match(b => b.AvatarId == run.Me.Id && b.PostId == post.ID).ExecuteAnyAsync(token)) + continue; + await DB.Default.SaveAsync(new Bookmark { AvatarId = run.Me.Id, PostId = post.ID }, token); + run.Count("bookmarks"); + } + } + + // asked for on purpose: each like tells the post's author + async Task Likes(Run run, CancellationToken token) + { + var statuses = run.Services.GetRequiredService(); + foreach (var uri in Strings(run.Archive.Json("likes.json")?["orderedItems"]).Distinct()) + { + if (await Stopped(run, token)) + return; + var post = await Fetched(run, uri, token); + if (post == default) + { + run.Count("likes not found"); + continue; + } + run.Count((await statuses.Favourite(run.Me, post.ID, true, token)).Ok ? "likes" : "likes not found"); + } + } + + // a post here by its address, or read from its server (signed) and kept + static async Task Fetched(Run run, string uri, CancellationToken token) + { + if (await Held(run, uri, token) is { } held) + return held; + return await run.Services.GetRequiredService().StoreContext(uri, 0, token); + } + + static async Task Held(Run run, string uri, CancellationToken token) + { + if (uri == default) + return default; + if (run.Posts.TryGetValue(uri, out var copy)) + return await DB.Default.Find().MatchID(copy.Id).ExecuteFirstAsync(token); + return await DB.Default.Find().Match(p => (p.ObjectURI == uri || p.Url == uri) && !p.DeletedAt.HasValue).ExecuteFirstAsync(token); + } + + // ---- the profile: name, words, fields, flags and pictures, saved as Settings saves them (and announced the same way) + + async Task Profile(Run run, JsonObject actor, CancellationToken token) + { + var avatar = await DB.Default.Find().MatchID(run.Me.Id).ExecuteFirstAsync(token); + if (Clean(Text(actor["name"])) is { } name) + avatar.Name = name; + if (Text(actor["summary"]) is { } summary) + avatar.Biography = Plain(summary); + var fields = (actor["attachment"] as JsonArray ?? []).OfType().Where(f => Text(f["type"]) == "PropertyValue" && Clean(Text(f["name"])) != default) + .Take(4).ToDictionary(f => Text(f["name"]).Trim(), f => Plain(Text(f["value"]) ?? string.Empty)); + if (fields.Count > 0) + avatar.Fields = fields; + if (actor["manuallyApprovesFollowers"]?.GetValueKind() is { } locked and (System.Text.Json.JsonValueKind.True or System.Text.Json.JsonValueKind.False)) + avatar.Settings.IsLocked = locked == System.Text.Json.JsonValueKind.True; + if (actor["discoverable"]?.GetValueKind() is { } discoverable and (System.Text.Json.JsonValueKind.True or System.Text.Json.JsonValueKind.False)) + avatar.Settings.IsDiscoverable = discoverable == System.Text.Json.JsonValueKind.True; + if (actor["indexable"]?.GetValueKind() is { } indexable and (System.Text.Json.JsonValueKind.True or System.Text.Json.JsonValueKind.False)) + avatar.Settings.IsIndexable = indexable == System.Text.Json.JsonValueKind.True; + foreach (var (field, kind, width, height) in new[] { ("icon", "avatar", 400, 400), ("image", "header", 1500, 500) }) + { + var file = (Text((actor[field] as JsonObject)?["url"]) ?? Text(actor[field]))?.TrimStart('/'); + if (file == default || file.Contains("://") || !run.Archive.Has(file)) + continue; + var outcome = await Picture(run, file, kind, width, height, token); + if (outcome?.Ok != true) + { + run.Count("pictures refused"); + continue; + } + if (kind == "avatar") + avatar.PictureURL = media.Url(outcome.Attachment.FilePath); + else + avatar.ThumbnailURL = media.Url(outcome.Attachment.FilePath); + await media.Trash(m => m.ProfileOfAvatarId == avatar.ID && m.Kind == kind && m.ID != outcome.Attachment.ID, "replaced", token); + } + avatar.UpdatedAt = DateTime.UtcNow; + await DB.Default.SaveAsync(avatar, token); + await run.Services.GetRequiredService().PublishProfile(localActors.FromAvatar(avatar), token); + run.Count("profile"); + } + + async Task Picture(Run run, string file, string kind, int width, int height, CancellationToken token) + { + var temporary = Path.Combine(Path.GetTempPath(), $"privapub-import-{Guid.NewGuid():N}{Path.GetExtension(file)}"); + try + { + await using (var source = run.Archive.Read(file)) + await using (var target = new FileStream(temporary, FileMode.CreateNew, FileAccess.Write)) + await source.CopyToAsync(target, token); + await using var stream = new FileStream(temporary, FileMode.Open, FileAccess.Read); + var form = new FormFile(stream, 0, stream.Length, kind, Path.GetFileName(file)) + { + Headers = new HeaderDictionary(), + ContentType = MediaService.ContentTypes.TryGetContentType(file, out var guessed) ? guessed : "application/octet-stream" + }; + return await media.ProfileImage(run.Me.Id, kind, form, width, height, token); + } + finally + { + File.Delete(temporary); + } + } + + // ---- helpers + + static string Text(JsonNode node) => node is JsonValue value && value.TryGetValue(out var text) ? text : default; + + static string Id(JsonNode node) => node switch + { + JsonValue => Text(node), + JsonObject json => Text(json["id"]), + JsonArray list => list.Select(Id).FirstOrDefault(i => i != default), + _ => default + }; + + static List Strings(JsonNode node) => node switch + { + JsonArray list => list.Select(Id).Where(i => i != default).ToList(), + JsonValue => Text(node) is { } one ? [one] : [], + JsonObject json when Text(json["type"]) is "OrderedCollection" or "Collection" => Strings(json["orderedItems"] ?? json["items"]), + _ => [] + }; + + static DateTime? Date(JsonNode node) => + DateTime.TryParse(Text(node), System.Globalization.CultureInfo.InvariantCulture, + System.Globalization.DateTimeStyles.AdjustToUniversal | System.Globalization.DateTimeStyles.AssumeUniversal, out var date) ? date : default; + + static DateTime Clamp(DateTime when) => when > DateTime.UtcNow ? DateTime.UtcNow : when; + + static string Clean(string text) => string.IsNullOrWhiteSpace(text) ? default : text.Trim(); + + static bool IsObjectId(string id) => id is { Length: 24 } && id.All(char.IsAsciiHexDigitLower); + + // a CSV line's cells (quotes as Mastodon writes them) + static string[] Cells(string line) + { + var cells = new List(); + var cell = new System.Text.StringBuilder(); + var quoted = false; + for (var i = 0; i < line.Length; i++) + { + var c = line[i]; + if (quoted) + { + if (c == '"' && i + 1 < line.Length && line[i + 1] == '"') + { + cell.Append('"'); + i++; + } + else if (c == '"') + quoted = false; + else + cell.Append(c); + } + else if (c == '"') + quoted = true; + else if (c == ',') + { + cells.Add(cell.ToString()); + cell.Clear(); + } + else + cell.Append(c); + } + cells.Add(cell.ToString()); + return [.. cells]; + } + + // HTML to the plain words a persona writes in its profile + static string Plain(string html) + { + var text = Breaks().Replace(html, "\n"); + text = Tags().Replace(text, string.Empty); + return WebUtility.HtmlDecode(text).Trim(); + } + + [GeneratedRegex(@"|

\s*]*>", RegexOptions.IgnoreCase)] + private static partial Regex Breaks(); + + [GeneratedRegex(@"<[^>]+>")] + private static partial Regex Tags(); + } +} diff --git a/PrivaPub/Domain/Portability/PersonaArchives.cs b/PrivaPub/Domain/Portability/PersonaArchives.cs new file mode 100644 index 0000000..4845994 --- /dev/null +++ b/PrivaPub/Domain/Portability/PersonaArchives.cs @@ -0,0 +1,79 @@ +using MongoDB.Entities; + +using PrivaPub.ClientModels.User; + +namespace PrivaPub.Domain.Portability +{ + // A persona's archive, made and restored in jobs (owner decision 2026-10-07): where each stands, one row per persona + // (its ID is the persona's). The files live in /.personas//: the export, kept a week, and the upload + // being imported. + public class PersonaArchive : Entity + { + public string ExportState { get; set; } = "none"; + public DateTime? ExportAskedAt { get; set; } + public DateTime? ExportReadyAt { get; set; } + public string ExportFile { get; set; } + public long ExportBytes { get; set; } + public string ExportError { get; set; } + + public string ImportState { get; set; } = "none"; + public DateTime? ImportAskedAt { get; set; } + public DateTime? ImportEndedAt { get; set; } + public string ImportFile { get; set; } + public List ImportParts { get; set; } = []; + public long ImportDone { get; set; } + public long ImportTotal { get; set; } + public Dictionary ImportCounts { get; set; } = []; + public string ImportError { get; set; } + public string ImportFrom { get; set; } + public bool ImportSameActor { get; set; } + public bool StopAsked { get; set; } + } + + public static class PersonaArchives + { + public static readonly TimeSpan Kept = TimeSpan.FromDays(7); + + /// Every part an import can take; the last three only when asked (they act, or tell others). + public static readonly string[] Parts = ["posts", "profile", "follows", "blocks", "mutes", "domainblocks", "lists", "bookmarks", "filters", "tags", + "notificationpolicy", "pins", "located", "scheduled", "likes"]; + + public static string Folder(string backupsRoot, string avatarId) => Path.Combine(backupsRoot, ".personas", avatarId); + + public static async Task Of(string avatarId, CancellationToken token) => + await DB.Default.Find().MatchID(avatarId).ExecuteFirstAsync(token) ?? new PersonaArchive { ID = avatarId }; + + public static ViewPersonaArchive View(PersonaArchive archive) => new() + { + ExportState = archive.ExportState, + ExportAskedAt = archive.ExportAskedAt, + ExportReadyAt = archive.ExportReadyAt, + ExportExpiresAt = archive.ExportReadyAt + Kept, + ExportBytes = archive.ExportBytes, + ExportError = archive.ExportError, + ImportState = archive.ImportState, + ImportAskedAt = archive.ImportAskedAt, + ImportEndedAt = archive.ImportEndedAt, + ImportParts = archive.ImportParts, + ImportDone = archive.ImportDone, + ImportTotal = archive.ImportTotal, + ImportCounts = archive.ImportCounts, + ImportError = archive.ImportError, + ImportFrom = archive.ImportFrom, + ImportSameActor = archive.ImportSameActor + }; + + /// Exports older than a week go, as do imports' files left behind. + public static async Task Forget(string backupsRoot, CancellationToken token) + { + var folder = Path.Combine(backupsRoot, ".personas"); + if (!Directory.Exists(folder)) + return; + var before = DateTime.UtcNow - Kept; + foreach (var file in Directory.EnumerateFiles(folder, "*", SearchOption.AllDirectories).Where(f => File.GetLastWriteTimeUtc(f) < before)) + File.Delete(file); + await DB.Default.Update().Match(a => a.ExportState == "ready" && a.ExportReadyAt < before) + .Modify(a => a.ExportState, "none").Modify(a => a.ExportFile, null).ExecuteAsync(token); + } + } +} diff --git a/PrivaPub/Domain/Portability/SafeArchive.cs b/PrivaPub/Domain/Portability/SafeArchive.cs new file mode 100644 index 0000000..0da4b14 --- /dev/null +++ b/PrivaPub/Domain/Portability/SafeArchive.cs @@ -0,0 +1,270 @@ +using System.IO.Compression; +using System.Text.Json; +using System.Text.Json.Nodes; + +namespace PrivaPub.Domain.Portability +{ + // An archive someone uploads, opened only if nothing in it can harm the server: no link, no path leaving it (.. or + // absolute), no name twice, at most MaxEntries files and MaxBytes in all (and what the disk has), and no file + // compressed more than MaxRatio to one. Every file is read through a stream that stops at its stated size, and the + // outbox one item at a time, none larger than MaxItemBytes, so a large archive never sits in memory. + public sealed class SafeArchive : IDisposable + { + public const int MaxEntries = 200_000; + public const long MaxBytes = 64L * 1024 * 1024 * 1024; + public const int MaxRatio = 100; + public const int MaxItemBytes = 1024 * 1024; + public const int MaxJsonBytes = 64 * 1024 * 1024; + + readonly ZipArchive _zip; + readonly Dictionary _entries; + readonly string _root;//the folder the archive's files sit in, when it wraps them in one + + SafeArchive(ZipArchive zip, Dictionary entries, string root) + { + _zip = zip; + _entries = entries; + _root = root; + } + + public static (SafeArchive Archive, string Error) Open(string path) + { + ZipArchive zip; + try + { + zip = ZipFile.OpenRead(path); + } + catch (InvalidDataException) + { + return (default, "not a zip archive"); + } + var entries = new Dictionary(StringComparer.Ordinal); + var total = 0L; + var free = new DriveInfo(Path.GetFullPath(path)).AvailableFreeSpace; + string Refuse(string why) + { + zip.Dispose(); + return why; + } + if (zip.Entries.Count > MaxEntries) + return (default, Refuse("too many files")); + foreach (var entry in zip.Entries) + { + var name = entry.FullName; + if (name.Contains('\\') || name.StartsWith('/') || name.Contains(':') || name.Split('/').Any(part => part is ".." or ".")) + return (default, Refuse($"{name}: a path that leaves the archive")); + if (((entry.ExternalAttributes >> 16) & 0xF000) == 0xA000) + return (default, Refuse($"{name}: a link")); + if (name.EndsWith('/')) + continue; + if (!entries.TryAdd(name, entry)) + return (default, Refuse($"{name}: twice")); + total += entry.Length; + if (total > MaxBytes || total > free) + return (default, Refuse("larger than the server can take")); + if (entry.Length > 1024 * 1024 && entry.Length > (long)MaxRatio * Math.Max(1, entry.CompressedLength)) + return (default, Refuse($"{name}: compressed more than {MaxRatio} to one")); + } + // some archivers wrap everything in one folder + var root = string.Empty; + if (!entries.ContainsKey("actor.json") && entries.Keys.FirstOrDefault(k => k.EndsWith("/actor.json", StringComparison.Ordinal)) is { } nested + && nested.Count(c => c == '/') == 1) + root = nested[..(nested.IndexOf('/') + 1)]; + return (new SafeArchive(zip, entries, root), default); + } + + public bool Has(string name) => _entries.ContainsKey(_root + name); + + public long Length(string name) => _entries.TryGetValue(_root + name, out var entry) ? entry.Length : -1; + + /// A file's contents, never more than its stated size. + public Stream Read(string name) => + _entries.TryGetValue(_root + name, out var entry) ? new Bounded(entry.Open(), entry.Length) : default; + + /// A JSON file, or null when it is missing, too large or not JSON. + public JsonNode Json(string name) + { + if (!_entries.TryGetValue(_root + name, out var entry) || entry.Length > MaxJsonBytes) + return default; + try + { + using var stream = Read(name); + return JsonNode.Parse(stream); + } + catch (JsonException) + { + return default; + } + } + + /// The lines of a text file (a CSV), none when it is missing or too large. + public List Lines(string name) + { + if (!_entries.TryGetValue(_root + name, out var entry) || entry.Length > MaxJsonBytes) + return []; + using var reader = new StreamReader(Read(name)); + var lines = new List(); + while (reader.ReadLine() is { } line) + if (line.Length > 0) + lines.Add(line); + return lines; + } + + /// outbox.json's orderedItems, one object at a time; an item larger than MaxItemBytes stops it. + public async IAsyncEnumerable Outbox([System.Runtime.CompilerServices.EnumeratorCancellation] CancellationToken token) + { + await using var stream = Read("outbox.json"); + if (stream == default) + yield break; + var buffer = new byte[64 * 1024]; + var filled = 0; + var final = false; + var state = new JsonReaderState(); + var phase = Phase.Seeking; + while (phase != Phase.Done) + { + if (!final) + { + if (filled == buffer.Length) + { + if (buffer.Length > MaxItemBytes + 64 * 1024) + throw new InvalidDataException($"an outbox item larger than {MaxItemBytes / 1024} KiB"); + Array.Resize(ref buffer, buffer.Length * 2); + } + var read = await stream.ReadAsync(buffer.AsMemory(filled), token); + filled += read; + final = read == 0; + } + var (items, consumed, next, after) = Scan(buffer.AsSpan(0, filled), final, state, phase); + foreach (var item in items) + yield return item; + (state, phase) = (next, after); + Buffer.BlockCopy(buffer, consumed, buffer, 0, filled - consumed); + filled -= consumed; + if (final && consumed == 0 && items.Count == 0) + break; + } + } + + enum Phase { Seeking, Array, Done } + + static (List Items, int Consumed, JsonReaderState State, Phase Phase) Scan(ReadOnlySpan data, bool final, JsonReaderState state, Phase phase) + { + var items = new List(); + var reader = new Utf8JsonReader(data, final, state); + while (phase != Phase.Done) + { + var mark = reader; + if (!reader.Read()) + { + reader = mark; + break; + } + if (phase == Phase.Seeking) + { + if (reader.TokenType == JsonTokenType.PropertyName && reader.CurrentDepth == 1 && reader.ValueTextEquals("orderedItems")) + { + var before = reader; + if (!reader.Read()) + { + reader = mark; + break; + } + if (reader.TokenType == JsonTokenType.StartArray) + phase = Phase.Array; + else + { + reader = before; + if (!reader.TrySkip()) + { + reader = mark; + break; + } + } + continue; + } + if (reader.TokenType is JsonTokenType.StartObject or JsonTokenType.StartArray && reader.CurrentDepth >= 1 && !reader.TrySkip()) + { + reader = mark; + break; + } + continue; + } + if (reader.TokenType == JsonTokenType.EndArray) + { + phase = Phase.Done; + break; + } + if (reader.TokenType != JsonTokenType.StartObject) + { + if (reader.TokenType == JsonTokenType.StartArray && !reader.TrySkip()) + { + reader = mark; + break; + } + continue; + } + var start = (int)reader.TokenStartIndex; + if (!reader.TrySkip()) + { + reader = mark; + break; + } + var length = (int)reader.BytesConsumed - start; + if (length > MaxItemBytes) + throw new InvalidDataException($"an outbox item larger than {MaxItemBytes / 1024} KiB"); + if (JsonNode.Parse(data.Slice(start, length)) is JsonObject item) + items.Add(item); + } + return (items, (int)reader.BytesConsumed, reader.CurrentState, phase); + } + + public void Dispose() => _zip.Dispose(); + + // a stream that ends at the size the archive states, whatever the compressed data says + sealed class Bounded(Stream inner, long length) : Stream + { + long _left = length; + + public override int Read(byte[] buffer, int offset, int count) => Read(buffer.AsSpan(offset, count)); + + public override int Read(Span buffer) + { + if (_left <= 0) + return 0; + var read = inner.Read(buffer[..(int)Math.Min(buffer.Length, _left)]); + _left -= read; + return read; + } + + public override async ValueTask ReadAsync(Memory buffer, CancellationToken token = default) + { + if (_left <= 0) + return 0; + var read = await inner.ReadAsync(buffer[..(int)Math.Min(buffer.Length, _left)], token); + _left -= read; + return read; + } + + public override Task ReadAsync(byte[] buffer, int offset, int count, CancellationToken token) => + ReadAsync(buffer.AsMemory(offset, count), token).AsTask(); + + public override bool CanRead => true; + public override bool CanSeek => false; + public override bool CanWrite => false; + public override long Length => length; + public override long Position { get => length - _left; set => throw new NotSupportedException(); } + public override void Flush() { } + public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException(); + public override void SetLength(long value) => throw new NotSupportedException(); + public override void Write(byte[] buffer, int offset, int count) => throw new NotSupportedException(); + + protected override void Dispose(bool disposing) + { + if (disposing) + inner.Dispose(); + base.Dispose(disposing); + } + } + } +} diff --git a/PrivaPub/Domain/Social/FollowService.cs b/PrivaPub/Domain/Social/FollowService.cs index edf308e..4820f06 100644 --- a/PrivaPub/Domain/Social/FollowService.cs +++ b/PrivaPub/Domain/Social/FollowService.cs @@ -30,6 +30,8 @@ namespace PrivaPub.Domain.Social Task UnfollowAs(LocalActor follower, string target, CancellationToken token); Task Decide(LocalActor me, string followerAccountId, bool accept, CancellationToken token); Task ResendPending(DateTime now, CancellationToken token); + /// An account named by its address or its handle (user@host): a persona here, or one elsewhere. + Task<(LocalActor Local, ForeignAvatar Remote)> Resolve(string target, CancellationToken token); } public class FollowService : IFollowService @@ -257,6 +259,8 @@ namespace PrivaPub.Domain.Social follower.Id, follower.Uri, default, token); } + public Task<(LocalActor Local, ForeignAvatar Remote)> Resolve(string target, CancellationToken token) => ResolveTarget(target, token); + async Task<(LocalActor Local, ForeignAvatar Remote)> ResolveTarget(string target, CancellationToken token) { target = target?.Trim(); diff --git a/PrivaPub/Infrastructure/Backup/Backups.cs b/PrivaPub/Infrastructure/Backup/Backups.cs index 97d0c9b..44abc6f 100644 --- a/PrivaPub/Infrastructure/Backup/Backups.cs +++ b/PrivaPub/Infrastructure/Backup/Backups.cs @@ -32,7 +32,7 @@ namespace PrivaPub.Infrastructure.Backup } // A nightly backup at Backups:NightlyAt (UTC), or as soon as the service is up after missing it; the old ones rotated - // out as each is made. + // out as each is made. Personas' archives a week old go too. public class BackupScheduler(Backups backups, ILogger logger) : BackgroundService { static readonly TimeSpan Interval = TimeSpan.FromMinutes(10); @@ -52,6 +52,7 @@ namespace PrivaPub.Infrastructure.Backup try { await RunIfDue(DateTime.UtcNow, stoppingToken); + await Domain.Portability.PersonaArchives.Forget(backups.Root, stoppingToken); } catch (Exception ex) when (ex is not OperationCanceledException) { diff --git a/PrivaPub/Infrastructure/Backup/ServerBackup.cs b/PrivaPub/Infrastructure/Backup/ServerBackup.cs index aa42087..4f73286 100644 --- a/PrivaPub/Infrastructure/Backup/ServerBackup.cs +++ b/PrivaPub/Infrastructure/Backup/ServerBackup.cs @@ -39,7 +39,8 @@ namespace PrivaPub.Infrastructure.Backup ["openiddict.authorizations"] = "sessions: a restore ends every one", [nameof(MaintenanceLock)] = "who is backing up or restoring now", ["AppConfiguration"] = "the configuration's copy holds the SMTP password, and is made again from appsettings at boot", - [nameof(RestoreRecord)] = "what restores did outlives what they restore" + [nameof(RestoreRecord)] = "what restores did outlives what they restore", + ["PersonaArchive"] = "personas' archives being made or imported belong to the moment" }; static readonly JsonWriterSettings Json = new() { OutputMode = JsonOutputMode.CanonicalExtendedJson, Indent = false }; diff --git a/PrivaPub/Infrastructure/Backup/TransferStore.cs b/PrivaPub/Infrastructure/Backup/TransferStore.cs index d4e262d..6dd3eb6 100644 --- a/PrivaPub/Infrastructure/Backup/TransferStore.cs +++ b/PrivaPub/Infrastructure/Backup/TransferStore.cs @@ -55,6 +55,34 @@ namespace PrivaPub.Infrastructure.Backup return id; } + /// Who started an upload (an administrator's name, or the root a persona's archive belongs to). + public string Owner(string id) + { + if (!IsId(id) || !System.IO.File.Exists(Path.Combine(Folder, id + ".json"))) + return default; + try + { + return JsonDocument.Parse(System.IO.File.ReadAllText(Path.Combine(Folder, id + ".json"))).RootElement.GetProperty("by").GetString(); + } + catch (JsonException) + { + return default; + } + } + + /// A whole upload moved where it is read from: its new path, or null when there is no such upload. + public string Take(string id, string folder, string name) + { + if (!IsId(id) || !System.IO.File.Exists(File(id))) + return default; + ServerBackup.MakeDirectory(folder); + var path = Path.Combine(folder, name); + System.IO.File.Move(File(id), path, overwrite: true); + System.IO.File.Delete(Path.Combine(Folder, id + ".json")); + _appending.TryRemove(id, out _); + return path; + } + /// How much of an upload has arrived; -1 when there is no such upload. public long Received(string id) => IsId(id) && System.IO.File.Exists(File(id)) ? new FileInfo(File(id)).Length : -1; diff --git a/PrivaPub/Infrastructure/Data/Indexes.cs b/PrivaPub/Infrastructure/Data/Indexes.cs index 6840e5c..4eda5cc 100644 --- a/PrivaPub/Infrastructure/Data/Indexes.cs +++ b/PrivaPub/Infrastructure/Data/Indexes.cs @@ -31,6 +31,14 @@ namespace PrivaPub.Infrastructure.Data await Plain(token, p => p.ReblogOfPostId); await Plain(token, p => p.AnsweringToPostId); await DB.Default.Index().Key(p => p.Geo, KeyType.Geo2DSphere).CreateAsync(token); + // a post imported from another actor's archive, once per persona: importing it twice changes nothing + await DB.Default.Index().Key(p => p.GroupUserId, KeyType.Ascending).Key(p => p.ImportedFromURI, KeyType.Ascending) + .Option(o => + { + o.Unique = true; + o.PartialFilterExpression = Builders.Filter.Gt(p => p.ImportedFromURI, ""); + }) + .CreateAsync(token); // media: a post's, a persona's, what a scheduled post holds, the janitor's never-posted uploads and its trash await Plain(token, m => m.PostId, m => m.ScheduledStatusId, m => m.CreatedAt); diff --git a/PrivaPub/Middleware/SocialPubConfigurations.cs b/PrivaPub/Middleware/SocialPubConfigurations.cs index 38fb3c3..958f245 100644 --- a/PrivaPub/Middleware/SocialPubConfigurations.cs +++ b/PrivaPub/Middleware/SocialPubConfigurations.cs @@ -122,6 +122,8 @@ namespace PrivaPub.Middleware .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() + .AddSingleton() .AddSingleton() .AddSingleton() .AddSingleton(services => services.GetRequiredService()) diff --git a/PrivaPub/Models/Jobs/Job.cs b/PrivaPub/Models/Jobs/Job.cs index dd38d69..c300619 100644 --- a/PrivaPub/Models/Jobs/Job.cs +++ b/PrivaPub/Models/Jobs/Job.cs @@ -36,7 +36,9 @@ namespace PrivaPub.Models.Jobs FetchReplies, BackfillOutbox, SynchronizeFollowing, - ProcessMedia + ProcessMedia, + ExportArchive, + ImportArchive } public enum JobState diff --git a/PrivaPub/Models/Post/Post.cs b/PrivaPub/Models/Post/Post.cs index 714864f..9c7e04c 100644 --- a/PrivaPub/Models/Post/Post.cs +++ b/PrivaPub/Models/Post/Post.cs @@ -58,6 +58,7 @@ namespace PrivaPub.Models.Post public bool IsFederatedCopy { get; set; } public string ObjectURI { get; set; }//the Note's id + public string ImportedFromURI { get; set; }//a post imported from another actor's archive: the original's id (PersonaImport) public string ObjectType { get; set; }//the remote object's type: Note, Article, Page, Video, Event…; null for local posts public bool AsChatMessage { get; set; }//a local direct message sent as a ChatMessage: its one recipient writes to us that way public string ActivityURI { get; set; }//the Create's id diff --git a/deploy/nginx/privapub.thepra.dev.conf b/deploy/nginx/privapub.thepra.dev.conf index f885bfc..80d829b 100644 --- a/deploy/nginx/privapub.thepra.dev.conf +++ b/deploy/nginx/privapub.thepra.dev.conf @@ -51,9 +51,9 @@ server { proxy_read_timeout 300s; } - # the administrator's backups (owner decision 2026-10-07): a download streams a whole backup for as long as it takes; - # an upload arrives in pieces of at most 32 MB, unbuffered - location ^~ /clientapi/admin/backups/download/ { + # the administrator's backups and the personas' archives (owner decision 2026-10-07): a download streams a whole one for + # as long as it takes; an upload arrives in pieces of at most 32 MB, unbuffered + location ~ ^/clientapi/(admin/backups|persona/archive)/download/ { proxy_buffering off; proxy_pass http://127.0.0.1:6970; proxy_http_version 1.1; @@ -65,7 +65,7 @@ server { proxy_read_timeout 3600s; proxy_send_timeout 3600s; } - location ^~ /clientapi/admin/backups/uploads { + location ~ ^/clientapi/(admin/backups|persona/archive)/uploads { client_max_body_size 40m; proxy_request_buffering off; limit_req zone=privapub_backup_uploads burst=30 nodelay; diff --git a/tools/pasture/scenarios/persona-archive.sh b/tools/pasture/scenarios/persona-archive.sh new file mode 100644 index 0000000..0dcabb4 --- /dev/null +++ b/tools/pasture/scenarios/persona-archive.sh @@ -0,0 +1,70 @@ +# A persona's archive (owner decision 2026-10-07), with a real one: mastouser's own account archive, made by Mastodon's +# backup service, is imported into mover_archive, a PrivaPub persona mastouser follows. Its posts arrive, with their +# media, on the persona's profile, and nobody is told: Mastodon, which would get anything delivered to the persona's +# followers, holds none of the copies. Importing it again changes nothing. Then the persona's own archive is exported and +# downloaded through its link, in Mastodon's layout. Needs the mastodon peer. +M=https://mastodon.test:6443 +mcurl() { curl -sk --resolve mastodon.test:6443:127.0.0.1 "$@"; } +. "$here/peers/mastodon.sh" +m_rails() { podman exec pasture-mastodon bin/rails runner "$1" 2>/dev/null | tail -1; } + +echo "persona-archive" +JWT=$(privapub_root) +RH="Authorization: Bearer $JWT" +AT=$(privapub_token mover_archive) +AH="Authorization: Bearer $AT" +[ -n "$AT" ] && ok "PrivaPub token for mover_archive" || { ko "PrivaPub token for mover_archive"; return 1; } +me=$(curl -s -H "$AH" "$P/api/v1/accounts/verify_credentials") +me_id=$(echo "$me" | j "print(d['id'])") +me_uri=$(echo "$me" | j "print(d['url'])" | sed 's|/@|/peasants/|') +A="$P/clientapi/persona/$me_id/archive" +# each run starts from a persona that has imported nothing +podman exec pasture-mongo mongosh --quiet PrivaPub --eval "db.Post.deleteMany({GroupUserId:'$me_id', ImportedFromURI:{\$exists:true, \$ne:null}}); db.MediaAttachment.deleteMany({OwnerAvatarId:'$me_id'})" >/dev/null + +m_follows() { m_rails "puts Follow.exists?(account: Account.find_local(\"mastouser\"), target_account: Account.find_by(uri: \"$me_uri\")) ? \"True\" : \"False\""; } +if [ "$(m_follows)" != "True" ]; then + on_m=$(mcurl -H "Authorization: Bearer $(mastodon_token)" "$M/api/v2/search?q=@mover_archive@privapub.test&resolve=true&type=accounts" | j "print(d['accounts'][0]['id'])") + mcurl -o /dev/null -X POST -H "Authorization: Bearer $(mastodon_token)" "$M/api/v1/accounts/$on_m/follow" +fi +until_true 45 '[ "$(m_follows)" = "True" ]' && ok "mastouser follows mover_archive" || ko "mastouser never followed mover_archive" + +dump=$(m_rails 'u = Account.find_local("mastouser").user; b = u.backups.create!; BackupService.new.call(b); puts b.reload.dump.path') +podman cp "pasture-mastodon:$dump" "$work/mastodon.zip" 2>/dev/null +[ -s "$work/mastodon.zip" ] && ok "Mastodon made mastouser's archive" || { ko "Mastodon made no archive"; return 1; } + +import() { + local upload + upload=$(curl -s -X POST -H "$RH" "$P/clientapi/persona/archive/uploads" | j "print(d['id'])") + curl -s -o /dev/null -X PUT -H "$RH" -H "Content-Type: application/octet-stream" --data-binary "@$work/mastodon.zip" "$P/clientapi/persona/archive/uploads/$upload?offset=0" + curl -s -o /dev/null -w '%{http_code}' -X POST -H "$RH" -H 'Content-Type: application/json' "$A/import" -d "{\"uploadId\":\"$upload\",\"parts\":[\"posts\"]}" +} +state() { curl -s -H "$RH" "$A" | j "print(d['$1'])"; } +count() { curl -s -H "$RH" "$A" | j "print(d['importCounts'].get('$1', 0))"; } + +[ "$(import)" = "202" ] && ok "the import is asked for" || { ko "the import was refused"; return 1; } +until_true 150 '[ "$(state importState)" = "done" ]' && ok "the import is done" || { ko "the import is $(state importState): $(state importError)"; return 1; } +posts=$(count posts) +[ "$posts" -gt 0 ] && ok "$posts of mastouser's posts imported ($(count 'boosts skipped') boosts and $(count 'direct or circle posts skipped') private posts left out)" \ + || ko "no post was imported" +[ "$(count media)" -gt 0 ] && ok "$(count media) pictures came with them" || ko "no picture came with them" +copies=$(curl -s -H "$AH" "$P/api/v1/accounts/$me_id/statuses?limit=40") +[ "$(echo "$copies" | j "print(len(d))")" -gt 0 ] && ok "they are on mover_archive's profile" || ko "the profile shows none of them" + +# nobody told: what PrivaPub would have delivered would be at Mastodon by now +sleep 20 +uris=$(echo "$copies" | j "print(' '.join(s['uri'] for s in d[:5]))") +held=$("$here/town.sh" stored mastodon $uris | j "print(sum(1 for v in d.values() if v['exists']))") +[ "$held" = "0" ] && ok "Mastodon holds none of the copies" || ko "Mastodon holds $held of the copies" + +[ "$(import)" = "202" ] && until_true 150 '[ "$(state importState)" = "done" ]' \ + && [ "$(count posts)" = "0" ] && [ "$(count 'posts already here')" = "$posts" ] \ + && ok "importing it again changes nothing" || ko "the second import: $(count posts) posts, $(count 'posts already here') already here" + +echo " exported" +[ "$(curl -s -o /dev/null -w '%{http_code}' -X POST -H "$RH" "$A")" = "202" ] && ok "the persona's archive is asked for" || ko "the archive was refused" +until_true 90 '[ "$(state exportState)" = "ready" ]' && ok "the archive is ready" || ko "the archive is $(state exportState): $(state exportError)" +link=$(curl -s -X POST -H "$RH" "$A/ticket" | j "print(d['url'])") +curl -s -o "$work/mover.zip" "$P$link" +names=$(unzip -Z1 "$work/mover.zip" 2>/dev/null) +grep -qx outbox.json <<< "$names" && grep -qx actor.json <<< "$names" \ + && ok "it downloads, in Mastodon's layout" || ko "the download is not an archive"