diff --git a/app/api/inngest/client.ts b/app/api/inngest/client.ts index e9bd13fa..ffce505c 100644 --- a/app/api/inngest/client.ts +++ b/app/api/inngest/client.ts @@ -1,6 +1,6 @@ import { Inngest } from "inngest"; + import { EventSchemas } from "inngest"; -import { Json } from "supabase/database.types"; export type Events = { "feeds/index-follows": { @@ -51,13 +51,10 @@ export type Events = { documentUris?: string[]; }; }; - "appview/index-document": { + "appview/sync-document-metadata": { data: { document_uri: string; - document_data: Json; bsky_post_uri?: string; - publication: string | null; - did: string; }; }; "user/write-records-to-pds": { diff --git a/app/api/inngest/functions/index_document.ts b/app/api/inngest/functions/sync_document_metadata.ts similarity index 60% rename from app/api/inngest/functions/index_document.ts rename to app/api/inngest/functions/sync_document_metadata.ts index 295f97f5..af2c3d4c 100644 --- a/app/api/inngest/functions/index_document.ts +++ b/app/api/inngest/functions/sync_document_metadata.ts @@ -1,27 +1,16 @@ import { inngest } from "../client"; import { supabaseServerClient } from "supabase/serverClient"; -import { AtpAgent } from "@atproto/api"; +import { AtpAgent, AtUri } from "@atproto/api"; import { idResolver } from "app/(home-pages)/reader/idResolver"; // 1m, 2m, 4m, 8m, 16m, 32m, 1h, 2h, 4h, 8h, 8h, 8h (~37h total) const SLEEP_INTERVALS = [ - "1m", - "2m", - "4m", - "8m", - "16m", - "32m", - "1h", - "2h", - "4h", - "8h", - "8h", - "8h", + "1m", "2m", "4m", "8m", "16m", "32m", "1h", "2h", "4h", "8h", "8h", "8h", ]; -export const index_document = inngest.createFunction( +export const sync_document_metadata = inngest.createFunction( { - id: "index_document_v2", + id: "sync_document_metadata_v2", debounce: { key: "event.data.document_uri", period: "60s", @@ -29,10 +18,11 @@ export const index_document = inngest.createFunction( }, concurrency: [{ key: "event.data.document_uri", limit: 1 }], }, - { event: "appview/index-document" }, + { event: "appview/sync-document-metadata" }, async ({ event, step }) => { - const { document_uri, document_data, bsky_post_uri, publication, did } = - event.data; + const { document_uri, bsky_post_uri } = event.data; + + const did = new AtUri(document_uri).host; const handleResult = await step.run("resolve-handle", async () => { const doc = await idResolver.did.resolve(did); @@ -49,37 +39,15 @@ export const index_document = inngest.createFunction( }); if (!handleResult) return { error: "No Handle" }; - if (handleResult.isBridgy) { - return { handle: handleResult.handle, skipped: true }; - } - - await step.run("write-document", async () => { - const docResult = await supabaseServerClient + await step.run("set-indexed", async () => { + return await supabaseServerClient .from("documents") - .upsert({ - uri: document_uri, - data: document_data, - indexed: true, - }); - if (docResult.error) console.log(docResult.error); - - if (publication) { - const docInPubResult = await supabaseServerClient - .from("documents_in_publications") - .upsert({ - publication, - document: document_uri, - }); - await supabaseServerClient - .from("documents_in_publications") - .delete() - .neq("publication", publication) - .eq("document", document_uri); - if (docInPubResult.error) console.log(docInPubResult.error); - } + .update({ indexed: !handleResult.isBridgy }) + .eq("uri", document_uri) + .select(); }); - if (!bsky_post_uri) { + if (!bsky_post_uri || handleResult.isBridgy) { return { handle: handleResult.handle }; } diff --git a/app/api/inngest/route.tsx b/app/api/inngest/route.tsx index f3e7ac1f..a74f1d98 100644 --- a/app/api/inngest/route.tsx +++ b/app/api/inngest/route.tsx @@ -13,7 +13,7 @@ import { check_oauth_session, } from "./functions/cleanup_expired_oauth_sessions"; import { write_records_to_pds } from "./functions/write_records_to_pds"; -import { index_document } from "./functions/index_document"; +import { sync_document_metadata } from "./functions/sync_document_metadata"; export const { GET, POST, PUT } = serve({ client: inngest, @@ -29,6 +29,6 @@ export const { GET, POST, PUT } = serve({ cleanup_expired_oauth_sessions, check_oauth_session, write_records_to_pds, - index_document, + sync_document_metadata, ], }); diff --git a/appview/index.ts b/appview/index.ts index dcdf00ff..31ead5fb 100644 --- a/appview/index.ts +++ b/appview/index.ts @@ -104,25 +104,40 @@ async function handleEvent(evt: Event) { console.log(record.error); return; } - let publication: string | null = null; + let docResult = await supabase.from("documents").upsert({ + uri: evt.uri.toString(), + data: record.value as Json, + }); + if (docResult.error) console.log(docResult.error); + await inngest.send({ + name: "appview/sync-document-metadata", + data: { + document_uri: evt.uri.toString(), + bsky_post_uri: record.value.postRef?.uri, + }, + }); if (record.value.publication) { let publicationURI = new AtUri(record.value.publication); + if (publicationURI.host !== evt.uri.host) { console.log("Unauthorized to create post!"); return; } - publication = record.value.publication; + let docInPublicationResult = await supabase + .from("documents_in_publications") + .upsert({ + publication: record.value.publication, + document: evt.uri.toString(), + }); + await supabase + .from("documents_in_publications") + .delete() + .neq("publication", record.value.publication) + .eq("document", evt.uri.toString()); + + if (docInPublicationResult.error) + console.log(docInPublicationResult.error); } - await inngest.send({ - name: "appview/index-document", - data: { - document_uri: evt.uri.toString(), - document_data: record.value as Json, - bsky_post_uri: record.value.postRef?.uri, - publication, - did: evt.did, - }, - }); } if (evt.event === "delete") { await supabase.from("documents").delete().eq("uri", evt.uri.toString()); @@ -256,29 +271,45 @@ async function handleEvent(evt: Event) { console.log(record.error); return; } + let docResult = await supabase.from("documents").upsert({ + uri: evt.uri.toString(), + data: record.value as Json, + }); + if (docResult.error) console.log(docResult.error); + await inngest.send({ + name: "appview/sync-document-metadata", + data: { + document_uri: evt.uri.toString(), + bsky_post_uri: record.value.bskyPostRef?.uri, + }, + }); + // site.standard.document uses "site" field to reference the publication // For documents in publications, site is an AT-URI (at://did:plc:xxx/site.standard.publication/rkey) // For standalone documents, site is an HTTPS URL (https://leaflet.pub/p/did:plc:xxx) // Only link to publications table for AT-URI sites - let publication: string | null = null; if (record.value.site && record.value.site.startsWith("at://")) { let siteURI = new AtUri(record.value.site); + if (siteURI.host !== evt.uri.host) { console.log("Unauthorized to create document in site!"); return; } - publication = record.value.site; + let docInPublicationResult = await supabase + .from("documents_in_publications") + .upsert({ + publication: record.value.site, + document: evt.uri.toString(), + }); + await supabase + .from("documents_in_publications") + .delete() + .neq("publication", record.value.site) + .eq("document", evt.uri.toString()); + + if (docInPublicationResult.error) + console.log(docInPublicationResult.error); } - await inngest.send({ - name: "appview/index-document", - data: { - document_uri: evt.uri.toString(), - document_data: record.value as Json, - bsky_post_uri: record.value.bskyPostRef?.uri, - publication, - did: evt.did, - }, - }); } if (evt.event === "delete") { await supabase.from("documents").delete().eq("uri", evt.uri.toString());