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"