diff --git a/app/lib/utils/records.ts b/app/lib/utils/records.ts new file mode 100644 index 0000000..4575d78 --- /dev/null +++ b/app/lib/utils/records.ts @@ -0,0 +1,120 @@ +import { callXrpc } from "$hatk/client"; +import { nextTid } from "./tid.ts"; + +/** One `dev.hatk.applyWrites#create` entry with a client-minted rkey. */ +export interface CreateWrite { + $type: "dev.hatk.applyWrites#create"; + collection: string; + rkey: string; + value: Record; +} + +/** One `dev.hatk.applyWrites#update` entry. The record must already exist. */ +export interface UpdateWrite { + $type: "dev.hatk.applyWrites#update"; + collection: string; + rkey: string; + value: Record; +} + +export interface DeleteWrite { + $type: "dev.hatk.applyWrites#delete"; + collection: string; + rkey: string; +} + +export type Write = CreateWrite | UpdateWrite | DeleteWrite; + +export function createWrite( + collection: string, + rkey: string, + value: Record, +): CreateWrite { + return { $type: "dev.hatk.applyWrites#create", collection, rkey, value }; +} + +export function updateWrite( + collection: string, + rkey: string, + value: Record, +): UpdateWrite { + return { $type: "dev.hatk.applyWrites#update", collection, rkey, value }; +} + +export function deleteWrite(collection: string, rkey: string): DeleteWrite { + return { $type: "dev.hatk.applyWrites#delete", collection, rkey }; +} + +export function atUri(did: string, collection: string, rkey: string): string { + return `at://${did}/${collection}/${rkey}`; +} + +export function rkeyOf(uri: string): string { + return uri.split("/").pop()!; +} + +/** + * Send every write in a single atomic PDS transaction. + * + * One request instead of one per record: a 10-photo gallery used to be ~31 + * sequential `createRecord` calls, each its own signed repo commit and firehose + * event, which is where the minute-plus publish time came from. It also means a + * failure can no longer leave a half-written gallery behind. + */ +export async function applyWrites(writes: Write[]): Promise { + if (writes.length === 0) return; + await callXrpc("dev.hatk.applyWrites", { writes } as never); +} + +/** True if the gallery is already in the index. Any failure reads as "no". */ +export async function galleryExists(galleryUri: string): Promise { + try { + const res = await callXrpc("social.grain.unspecced.getGallery", { gallery: galleryUri }); + return !!(res as { gallery?: unknown })?.gallery; + } catch { + return false; + } +} + +/** Run `fn` over `items`, at most `limit` at a time, preserving input order. */ +export async function mapWithConcurrency( + items: T[], + limit: number, + fn: (item: T, index: number) => Promise, +): Promise { + const results: R[] = []; + let next = 0; + const workers = Array.from({ length: Math.min(limit, items.length) }, async () => { + while (next < items.length) { + const index = next++; + results[index] = await fn(items[index], index); + } + }); + await Promise.all(workers); + return results; +} + +/** + * Blobs can't ride along in applyWrites — the PDS takes raw bytes — so they + * stay one request each. Overlapping a few keeps a full gallery from being a + * long chain of serial uploads, while staying gentle enough on a phone's + * connection and memory. + */ +const UPLOAD_CONCURRENCY = 3; + +export function dataUrlToBlob(dataUrl: string): Blob { + const base64 = dataUrl.split(",")[1]; + const binary = Uint8Array.from(atob(base64), (c) => c.charCodeAt(0)); + return new Blob([binary], { type: "image/jpeg" }); +} + +/** Upload each data URL as a blob and return the blob refs in input order. */ +export async function uploadPhotoBlobs(dataUrls: string[]): Promise { + return mapWithConcurrency(dataUrls, UPLOAD_CONCURRENCY, async (dataUrl) => { + const result = await callXrpc("dev.hatk.uploadBlob", dataUrlToBlob(dataUrl) as never); + return (result as { blob: unknown }).blob; + }); +} + +/** Mint a fresh record key. Re-exported so callers only import from one place. */ +export { nextTid }; diff --git a/app/lib/utils/tid.ts b/app/lib/utils/tid.ts new file mode 100644 index 0000000..a7c7a52 --- /dev/null +++ b/app/lib/utils/tid.ts @@ -0,0 +1,45 @@ +/** + * AT Protocol TID (timestamp identifier) generation. + * + * We mint record keys client-side rather than letting the PDS assign them so + * that writes are idempotent. A create with a server-assigned rkey turns every + * replayed request into a brand new record — that is how a single gallery ended + * up with the same photo at position 2 three times, ~20s apart, from one + * `await`. With the rkey fixed up front a replay collides with the record it + * already wrote and the PDS rejects it instead of duplicating it. + * + * A TID is 13 chars of base32-sortable encoding a 64-bit integer: one + * always-zero high bit, 53 bits of microseconds since the UNIX epoch, then a + * 10-bit clock id that keeps concurrent writers on different devices apart. + */ + +const S32 = "234567abcdefghijklmnopqrstuvwxyz"; + +/** Random per-session, per the TID spec's clock identifier. */ +const CLOCK_ID = BigInt(Math.floor(Math.random() * 1024)); + +let lastMicros = 0n; + +/** Encode microseconds + clock id as a 13-char base32-sortable TID. */ +export function encodeTid(micros: bigint, clockId: bigint): string { + let n = (micros << 10n) | (clockId & 1023n); + let out = ""; + for (let i = 0; i < 13; i++) { + out = S32[Number(n & 31n)] + out; + n >>= 5n; + } + return out; +} + +/** + * Mint the next TID. Strictly increasing, so a batch of writes minted in a + * loop sorts in call order. + */ +export function nextTid(): string { + // Date.now() is millisecond-granular and a whole gallery's worth of rkeys is + // minted inside one tick, so step forward by hand rather than handing back + // the same rkey 30 times. + const micros = BigInt(Date.now()) * 1000n; + lastMicros = micros > lastMicros ? micros : lastMicros + 1n; + return encodeTid(lastMicros, CLOCK_ID); +} diff --git a/app/routes/create/+page.svelte b/app/routes/create/+page.svelte index 44768a7..5472508 100644 --- a/app/routes/create/+page.svelte +++ b/app/routes/create/+page.svelte @@ -2,7 +2,14 @@ import { goto } from '$app/navigation' import { onMount } from 'svelte' import { useQueryClient } from '@tanstack/svelte-query' - import { callXrpc } from '$hatk/client' + import { + applyWrites, + atUri, + createWrite, + galleryExists, + nextTid, + uploadPhotoBlobs, + } from '$lib/utils/records' import { processPhotos, type ProcessedPhoto } from '$lib/utils/image-resize' import { reverseGeocode, formatLocationName, extractAddress } from '$lib/utils/nominatim' import { createBskyPost } from '$lib/utils/bsky-post' @@ -225,41 +232,14 @@ error = null try { - const now = new Date().toISOString() - const photoUris: string[] = [] - - // 1. Upload blobs + create photo records - for (const photo of photos) { - const base64 = photo.dataUrl.split(',')[1] - const binary = Uint8Array.from(atob(base64), (c) => c.charCodeAt(0)) - const blob = new Blob([binary], { type: 'image/jpeg' }) + const did = $viewer?.did + if (!did) throw new Error('You must be signed in to post a gallery.') - const uploadResult = await callXrpc('dev.hatk.uploadBlob', blob as any) + const now = new Date().toISOString() - const photoResult = await callXrpc('dev.hatk.createRecord', { - collection: 'social.grain.photo', - record: { - photo: (uploadResult as any).blob, - aspectRatio: { width: photo.width, height: photo.height }, - ...(photo.alt ? { alt: photo.alt } : {}), - createdAt: now, - }, - }) - const photoUri = (photoResult as any).uri as string - photoUris.push(photoUri) - - // Create EXIF record if we extracted metadata and user opted in - if (photo.exif && $includeExif) { - await callXrpc('dev.hatk.createRecord', { - collection: 'social.grain.photo.exif', - record: { - photo: photoUri, - ...photo.exif, - createdAt: now, - }, - }) - } - } + // 1. Upload blobs. These are the only per-photo requests left — the PDS + // takes raw bytes, so they can't be batched with the records. + const blobs = await uploadPhotoBlobs(photos.map((p) => p.dataUrl)) // 2. Parse facets from description let facets: any[] | undefined @@ -268,10 +248,36 @@ if (parsed.facets.length > 0) facets = parsed.facets } - // 3. Create gallery record - const galleryResult = await callXrpc('dev.hatk.createRecord', { - collection: 'social.grain.gallery', - record: { + // 3. Write every record in one atomic applyWrites. Minting the rkeys here + // means we know each URI before the call, so gallery items can point at + // a gallery created in the same transaction — and a replayed request + // collides with itself instead of duplicating photos. + const photoRkeys = photos.map(() => nextTid()) + const galleryRkey = nextTid() + const galleryUri = atUri(did, 'social.grain.gallery', galleryRkey) + const photoUris = photoRkeys.map((rkey) => atUri(did, 'social.grain.photo', rkey)) + + const writes = [ + ...photos.map((photo, i) => + createWrite('social.grain.photo', photoRkeys[i], { + photo: blobs[i], + aspectRatio: { width: photo.width, height: photo.height }, + ...(photo.alt ? { alt: photo.alt } : {}), + createdAt: now, + }), + ), + ...photos.flatMap((photo, i) => + photo.exif && $includeExif + ? [ + createWrite('social.grain.photo.exif', nextTid(), { + photo: photoUris[i], + ...photo.exif, + createdAt: now, + }), + ] + : [], + ), + createWrite('social.grain.gallery', galleryRkey, { title: title.trim(), ...(description.trim() ? { description: description.trim() } : {}), ...(facets ? { facets } : {}), @@ -293,27 +299,32 @@ } : {}), createdAt: now, - }, - }) - const galleryUri = (galleryResult as any).uri as string - - // 4. Create gallery items - for (let i = 0; i < photoUris.length; i++) { - await callXrpc('dev.hatk.createRecord', { - collection: 'social.grain.gallery.item', - record: { + }), + ...photoUris.map((photoUri, i) => + createWrite('social.grain.gallery.item', nextTid(), { gallery: galleryUri, - item: photoUris[i], + item: photoUri, position: i, createdAt: now, - }, - }) + }), + ), + ] + + try { + await applyWrites(writes) + } catch (err) { + // A replayed request re-sends writes the PDS already committed. Because + // the rkeys are fixed, the replay collides on the first duplicate and + // the whole transaction is rejected — nothing is written twice, but the + // post did go through. Surfacing the error here would send the user + // back to press Post again, which mints fresh rkeys and would leave two + // galleries behind. So: if the gallery is there, we're done. + if (!(await galleryExists(galleryUri))) throw err } - // 5. Create Bluesky post if opted in - if (postToBluesky && $viewer) { - const galleryRkey = galleryUri.split('/').pop() - const galleryUrl = `${window.location.origin}/profile/${$viewer.did}/gallery/${galleryRkey}` + // 4. Create Bluesky post if opted in + if (postToBluesky) { + const galleryUrl = `${window.location.origin}/profile/${did}/gallery/${galleryRkey}` await createBskyPost({ url: galleryUrl, title: title.trim() || undefined, diff --git a/app/routes/profile/[did]/gallery/[rkey]/edit/+page.svelte b/app/routes/profile/[did]/gallery/[rkey]/edit/+page.svelte index a0f172b..e40ef77 100644 --- a/app/routes/profile/[did]/gallery/[rkey]/edit/+page.svelte +++ b/app/routes/profile/[did]/gallery/[rkey]/edit/+page.svelte @@ -1,7 +1,17 @@