import { iterateAtpRepo } from "@atcute/car"; import * as TID from "@atcute/tid"; import type { SessionClient } from "../auth/client"; import { parseISO, isValid } from "date-fns"; import type { BookUtilContext } from "../context"; import { ids, Book as BookRecord, Buzz as BuzzRecord } from "../bsky/lexicon"; import type { HiveId, UserBook, UserBookRow } from "../types"; import { findBookIdentifiersByLookup } from "../bsky/bookLookup"; import { toBookIdentifiersOutput } from "./bookIdentifiers"; import { uploadImageBlob } from "./uploadImageBlob"; import { BOOK_STATUS } from "../constants"; import { hydrateUserBook, serializeUserBook } from "./bookProgress"; import { ensureBookCataloged } from "./ensureBookCataloged"; /** * Normalize a date string to ISO format at midnight UTC * Handles YYYY-MM-DD format and ISO strings, always returns UTC midnight */ function normalizeDate(dateString: string | undefined): string | undefined { if (!dateString || dateString === "") { return undefined; } try { if (/^\d{4}-\d{2}-\d{2}$/.test(dateString)) { return new Date(dateString + "T00:00:00.000Z").toISOString(); } const date = parseISO(dateString); if (!isValid(date)) { // If parseISO fails, try creating a new Date const fallbackDate = new Date(dateString); if (!isValid(fallbackDate)) { return undefined; } const year = fallbackDate.getUTCFullYear(); const month = fallbackDate.getUTCMonth(); const day = fallbackDate.getUTCDate(); return new Date(Date.UTC(year, month, day, 0, 0, 0, 0)).toISOString(); } const year = date.getUTCFullYear(); const month = date.getUTCMonth(); const day = date.getUTCDate(); return new Date(Date.UTC(year, month, day, 0, 0, 0, 0)).toISOString(); } catch { return undefined; } } /** * Infers book status and auto-sets dates based on user input */ function inferBookStatusAndDates(updates: { status?: string; startedAt?: string; finishedAt?: string; }): { status: string | undefined; startedAt: string | undefined; finishedAt: string | undefined; } { let autoStartedAt = updates.startedAt; let autoFinishedAt = updates.finishedAt; let autoStatus = updates.status; // If user sets startedAt and status is "want to read" or unset, infer they're reading if (updates.startedAt && (!updates.status || updates.status === BOOK_STATUS.WANTTOREAD)) { autoStatus = BOOK_STATUS.READING; } // If user sets finishedAt and status is "want to read", "reading", or unset, infer they're finished if ( updates.finishedAt && (!updates.status || updates.status === BOOK_STATUS.WANTTOREAD || updates.status === BOOK_STATUS.READING) ) { autoStatus = BOOK_STATUS.FINISHED; } // Auto-set dates based on status if not already provided. // Use UTC "today" so the date doesn't depend on server timezone (startOfDay uses server local TZ). if (autoStatus === BOOK_STATUS.READING && !updates.startedAt) { const now = new Date(); autoStartedAt = new Date( Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), now.getUTCDate(), 0, 0, 0, 0), ).toISOString(); } else if (autoStatus === BOOK_STATUS.FINISHED && !updates.finishedAt) { const now = new Date(); autoFinishedAt = new Date( Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), now.getUTCDate(), 0, 0, 0, 0), ).toISOString(); } // Normalize dates to ISO format at start of day autoStartedAt = normalizeDate(autoStartedAt); autoFinishedAt = normalizeDate(autoFinishedAt); return { status: autoStatus, startedAt: autoStartedAt, finishedAt: autoFinishedAt, }; } export async function getUserBook({ ctx, agent, hiveId, }: { ctx: Pick; agent: SessionClient; hiveId: HiveId; }): Promise { const rawUserBook = await ctx.db .selectFrom("user_book") .selectAll() .where("userDid", "=", agent.did) .where("hiveId", "=", hiveId) .executeTakeFirst(); if (!rawUserBook) { return null; } return hydrateUserBook(rawUserBook); } export async function updateUserBook({ ctx, userBook, }: { ctx: Pick; userBook: UserBook; }): Promise { const row: UserBookRow = serializeUserBook(userBook); await ctx.db .insertInto("user_book") .values(row) .onConflict((oc) => oc.column("uri").doUpdateSet((c) => ({ indexedAt: c.ref("excluded.indexedAt"), cid: c.ref("excluded.cid"), authors: c.ref("excluded.authors"), title: c.ref("excluded.title"), hiveId: c.ref("excluded.hiveId"), status: c.ref("excluded.status"), owned: c.ref("excluded.owned"), startedAt: c.ref("excluded.startedAt"), finishedAt: c.ref("excluded.finishedAt"), review: c.ref("excluded.review"), stars: c.ref("excluded.stars"), bookProgress: c.ref("excluded.bookProgress"), })), ) .execute(); } /** * Get a book from the user's PDS */ export async function getBookRecord({ agent, cid, uri, }: { agent: SessionClient; cid: string; uri: string; }): Promise { const res = await agent.get("com.atproto.repo.getRecord", { params: { repo: agent.did, collection: ids.BuzzBookhiveBook, rkey: uri.split("/").at(-1)!, cid, }, }); const payload = res.ok ? (res.data as { value?: unknown }) : null; const originalBook = payload?.value as BookRecord.Record | undefined; if (!originalBook) { return null; } return originalBook; } /** * Update a book in the user's PDS */ export async function updateBookRecord({ ctx, agent, hiveId, updates, }: { ctx: BookUtilContext; agent: SessionClient; hiveId: HiveId; updates: Partial & { coverImage?: string }; }): Promise<{ book: BookRecord.Record; userBook: UserBook }> { const userBook = await getUserBook({ ctx, agent, hiveId }); let originalBook: BookRecord.Record | null = null; if (userBook) { originalBook = await getBookRecord({ agent, cid: userBook.cid, uri: userBook.uri, }); } if (!originalBook && !userBook) { const hiveBook = await ctx.db .selectFrom("hive_book") .selectAll() .where("id", "=", hiveId) .executeTakeFirst(); if (hiveBook) { Object.assign(updates, { coverImage: (hiveBook.cover || hiveBook.thumbnail) as string, title: hiveBook.title, authors: hiveBook.authors, ...updates, }); } } // Infer status and auto-set dates based on user input const { status: autoStatus, startedAt: autoStartedAt, finishedAt: autoFinishedAt, } = inferBookStatusAndDates({ status: updates.status, startedAt: updates.startedAt, finishedAt: updates.finishedAt, }); if (autoStartedAt && autoFinishedAt) { const startedDateStr = autoStartedAt.split("T")[0]!; // Extract YYYY-MM-DD const finishedDateStr = autoFinishedAt.split("T")[0]!; // Extract YYYY-MM-DD if (finishedDateStr < startedDateStr) { throw new Error("Finished date must be on or after started date"); } } // Determine the final status (use auto-inferred status or original status) const finalStatus = autoStatus || originalBook?.status; const identifiersRow = await findBookIdentifiersByLookup({ ctx, hiveId }); const identifiers = toBookIdentifiersOutput(identifiersRow); // Ensure the book is cataloged before writing to the user's PDS so we can embed // hiveBookUri. Fast path (expected): already cataloged, one DB read, no network. // Slow path (last resort): enrich + catalog with timeouts if it slipped through. const hiveBookAtUri = await ensureBookCataloged(ctx, hiveId); const bookData = { $type: ids.BuzzBookhiveBook, // Always prefer original values title: originalBook?.title || updates.title, authors: originalBook?.authors || updates.authors, hiveId: originalBook?.hiveId || hiveId, createdAt: originalBook?.createdAt || new Date().toISOString(), cover: originalBook?.cover || (await uploadImageBlob(updates.coverImage, agent, 800)), // Always prefer new values (including auto-inferred status) status: finalStatus, startedAt: autoStartedAt !== undefined ? autoStartedAt === "" ? undefined : autoStartedAt : originalBook?.startedAt, finishedAt: autoFinishedAt !== undefined ? autoFinishedAt === "" ? undefined : autoFinishedAt : originalBook?.finishedAt, review: updates.review || originalBook?.review, stars: updates.stars || originalBook?.stars, // Clear bookProgress when marking as finished bookProgress: finalStatus === BOOK_STATUS.FINISHED ? undefined : updates.bookProgress !== undefined ? updates.bookProgress : originalBook?.bookProgress, // Default to owned when first adding a book to library owned: updates.owned ?? originalBook?.owned ?? true, identifiers: Object.keys(identifiers).length > 0 ? identifiers : undefined, hiveBookUri: hiveBookAtUri ?? originalBook?.hiveBookUri, }; const book = BookRecord.validateRecord(bookData); if (!book.success) { throw new Error("Book incomplete or invalid: " + book.error.message); } const record = book.value as BookRecord.Record; const response = await agent.post("com.atproto.repo.applyWrites", { input: { repo: agent.did, writes: [ { $type: originalBook ? "com.atproto.repo.applyWrites#update" : "com.atproto.repo.applyWrites#create", collection: ids.BuzzBookhiveBook, rkey: userBook ? userBook.uri.split("/").at(-1)! : TID.now(), value: record, }, ], }, }); type ApplyOut = { results?: Array<{ $type: string; uri?: string; cid?: string }>; }; const applyData = response.ok ? (response.data as ApplyOut) : null; const firstResult = applyData?.results?.[0]; if ( !response.ok || !applyData?.results || applyData.results.length === 0 || !firstResult || !( firstResult.$type === "com.atproto.repo.applyWrites#updateResult" || firstResult.$type === "com.atproto.repo.applyWrites#createResult" ) ) { throw new Error("Failed to record book"); } const nextUserBook = { uri: firstResult.uri!, cid: firstResult.cid!, userDid: agent.did, createdAt: record.createdAt, authors: record.authors, title: record.title, indexedAt: new Date().toISOString(), hiveId: record.hiveId as HiveId, status: record.status || null, owned: record.owned ? 1 : 0, startedAt: record.startedAt || null, finishedAt: record.finishedAt || null, review: record.review || null, stars: record.stars || null, bookProgress: record.bookProgress ?? null, }; await updateUserBook({ ctx, userBook: nextUserBook }); return { book: record, userBook: nextUserBook }; } /** * Update a book in the user's PDS */ export async function updateBookRecords({ ctx, agent, updates, bookRecords = getUserRepoRecords({ ctx, agent }), overwrite = false, }: { ctx: BookUtilContext; agent: SessionClient; updates: Map & { coverImage?: string }>; bookRecords?: Promise<{ books: Map; }>; overwrite?: boolean; }): Promise { const updatesToApply: Array<{ type: "create" | "update"; record: BookRecord.Record; rkey: string; userBook: Omit; originalUpdate: Partial & { coverImage?: string }; }> = []; const bookMap = (await bookRecords).books; const hiveIds = [...updates.keys()]; const idRows = await ctx.db .selectFrom("book_id_map") .where("hiveId", "in", hiveIds) .selectAll() .execute(); const identifiersByHiveId = new Map(idRows.map((r) => [r.hiveId, toBookIdentifiersOutput(r)])); // Ensure all books are cataloged before writing to the user's PDS. // Fast path (expected): already cataloged via backfill/searchBooks — just a DB read. // Runs in parallel across all books; failures are swallowed inside ensureBookCataloged. const hiveBookUriMap = new Map(); if (ctx.serviceAccountAgent) { await Promise.allSettled( hiveIds.map(async (hiveId) => { hiveBookUriMap.set(hiveId, await ensureBookCataloged(ctx, hiveId)); }), ); } for (const [hiveId, update] of updates.entries()) { const [rkey, originalBook] = bookMap.entries().find(([_rkey, book]) => book.hiveId === hiveId) ?? []; // TODO maybe overwrite can overwrite just those properties we allow if (!overwrite && originalBook) { // If we're not overwriting, and the book already exists, skip it continue; } // Infer status and auto-set dates based on user input const { status: autoStatus, startedAt: autoStartedAt, finishedAt: autoFinishedAt, } = inferBookStatusAndDates({ status: update.status, startedAt: update.startedAt, finishedAt: update.finishedAt, }); if (autoStartedAt && autoFinishedAt) { const startedDateStr = autoStartedAt.split("T")[0]!; // Extract YYYY-MM-DD const finishedDateStr = autoFinishedAt.split("T")[0]!; // Extract YYYY-MM-DD if (finishedDateStr < startedDateStr) { throw new Error("Finished date must be on or after started date"); } } // Determine the final status (use auto-inferred status or original status) const finalStatus = autoStatus || originalBook?.status; const idOutput = identifiersByHiveId.get(hiveId); const identifiers = idOutput && Object.keys(idOutput).length > 0 ? idOutput : undefined; const book = BookRecord.validateRecord({ $type: ids.BuzzBookhiveBook, // Always prefer original values title: originalBook?.title || update.title, authors: originalBook?.authors || update.authors, hiveId: originalBook?.hiveId || hiveId, createdAt: originalBook?.createdAt || new Date().toISOString(), cover: originalBook?.cover, // Always prefer new values (including auto-inferred status) status: finalStatus, startedAt: autoStartedAt !== undefined ? autoStartedAt === "" ? undefined : autoStartedAt : originalBook?.startedAt, finishedAt: autoFinishedAt !== undefined ? autoFinishedAt === "" ? undefined : autoFinishedAt : originalBook?.finishedAt, review: update.review || originalBook?.review, stars: update.stars || originalBook?.stars, // Clear bookProgress when marking as finished bookProgress: finalStatus === BOOK_STATUS.FINISHED ? undefined : update.bookProgress !== undefined ? update.bookProgress : originalBook?.bookProgress, // Default to owned when first adding a book to library owned: update.owned ?? originalBook?.owned ?? true, identifiers, hiveBookUri: hiveBookUriMap.get(hiveId) ?? originalBook?.hiveBookUri, }); if (!book.success) { throw new Error("Book incomplete or invalid: " + book.error.message); } const record = book.value as BookRecord.Record; updatesToApply.push({ type: originalBook ? "update" : "create", record: record, rkey: rkey ?? TID.now(), userBook: { userDid: agent.did, createdAt: record.createdAt, authors: record.authors, title: record.title, indexedAt: new Date().toISOString(), hiveId: record.hiveId as HiveId, status: record.status || null, owned: record.owned ? 1 : 0, startedAt: record.startedAt || null, finishedAt: record.finishedAt || null, review: record.review || null, stars: record.stars || null, bookProgress: record.bookProgress ?? null, }, originalUpdate: update, }); } if (updatesToApply.length === 0) { return; } // Upload the cover image in parallel if it is missing await Promise.all( updatesToApply.map(async (u) => { if (!u.record.cover) { u.record.cover = (await uploadImageBlob( u.originalUpdate.coverImage, agent, 800, )) as typeof u.record.cover; } return u; }), ); const response = await agent.post("com.atproto.repo.applyWrites", { input: { repo: agent.did, writes: updatesToApply.map(({ type, record, rkey }) => ({ $type: `com.atproto.repo.applyWrites#${type}`, collection: ids.BuzzBookhiveBook, rkey, value: record, })), }, }); type ApplyOutBulk = { results?: Array<{ $type: string; uri?: string; cid?: string }>; }; const applyData2 = response.ok ? (response.data as ApplyOutBulk) : null; if (!response.ok || !applyData2?.results || applyData2.results.length === 0) { throw new Error("Failed to record books"); } await applyData2.results.reduce( async ( acc: Promise, result: { $type: string; uri?: string; cid?: string }, index: number, ) => { await acc; const update = updatesToApply[index]!; if ( result.$type === "com.atproto.repo.applyWrites#updateResult" || result.$type === "com.atproto.repo.applyWrites#createResult" ) { await updateUserBook({ ctx, userBook: { ...update.userBook, uri: result.uri!, cid: result.cid! }, }); } }, Promise.resolve(), ); ctx.addWideEventContext({ event: "wrote_books", userDid: agent.did, book_count: updatesToApply.length, }); return; } /** * Get all the books and buzzes from a user's PDS repo */ export async function getUserRepoRecords({ ctx, agent, did = agent.did, }: { ctx: Pick; agent: SessionClient; did?: string; }): Promise<{ /** * key is the rkey of the book */ books: Map; /** * key is the rkey of the buzz */ buzzes: Map; }> { const res = await agent.get("com.atproto.sync.getRepo", { params: { did }, as: "bytes", }); const data: Uint8Array = res.ok ? (res.data as Uint8Array) : new Uint8Array(0); const books = new Map(); const buzzes = new Map(); for (const { collection, rkey: key, record: value } of iterateAtpRepo(data)) { switch (collection) { case ids.BuzzBookhiveBook: { // https://github.com/bluesky-social/atproto/issues/3866 to get the validation to pass // Need to parse the whole object into a JSON, then parse it back into a Lexicon object const book = BookRecord.validateRecord( JSON.parse(JSON.stringify(value)) as BookRecord.Record, ); if (book.success) { books.set(key, book.value); } break; } case ids.BuzzBookhiveBuzz: { // https://github.com/bluesky-social/atproto/issues/3866 to get the validation to pass // Need to parse the whole object into a JSON, then parse it back into a Lexicon object const buzz = BuzzRecord.validateRecord( JSON.parse(JSON.stringify(value)) as BuzzRecord.Record, ); if (buzz.success) { buzzes.set(key, buzz.value); } break; } } } ctx.addWideEventContext({ event: "fetched_repo", userDid: did, book_count: books.size, buzz_count: buzzes.size, }); return { books, buzzes }; }