diff --git a/.gitignore b/.gitignore index b062044..aeea9f8 100644 --- a/.gitignore +++ b/.gitignore @@ -1,7 +1,7 @@ whatever.db system.db **/.DS_Store -docs/.vite/deps/_metadata.json -docs/.vite/deps/package.json -indexserver.ts +docs/.vite/ .gitignore +indexserver.ts +dbs/ \ No newline at end of file diff --git a/indexserver.ts b/indexserver.ts index 1362f98..4fd6e0b 100644 --- a/indexserver.ts +++ b/indexserver.ts @@ -1,6 +1,6 @@ import { indexHandlerContext } from "./index/types.ts"; -import { validateRecord } from "./utils/records.ts"; +import { assertRecord, validateRecord } from "./utils/records.ts"; import { searchParamsToJson, withCors } from "./utils/server.ts"; import * as IndexServerTypes from "./utils/indexservertypes.ts"; import { Database } from "jsr:@db/sqlite@0.11"; @@ -9,6 +9,7 @@ import { jetstreamurl } from "./main.ts"; import { JetstreamManager, SpacedustManager } from "./utils/sharders.ts"; import { handleSpacedust, SpacedustLinkMessage } from "./index/spacedust.ts"; import { handleJetstream } from "./index/jetstream.ts"; +import * as ATPAPI from "npm:@atproto/api"; import { AtUri } from "npm:@atproto/api"; export class IndexServerUserManager { @@ -44,24 +45,24 @@ export class IndexServerUserManager { getDbForDid(did: string): Database | null { if (!this.users.has(did)) { - return null + return null; } return this.users.get(did)?.db ?? null; } coldStart(db: Database) { - const rows = db.prepare("SELECT did FROM users").all(); - for (const row of rows) { - this.addUser(row.did); + const rows = db.prepare("SELECT did FROM users").all(); + for (const row of rows) { + this.addUser(row.did); + } } } -} class UserIndexServer { did: string; - db: Database;// | undefined; - jetstream: JetstreamManager;// | undefined; - spacedust: SpacedustManager;// | undefined; + db: Database; // | undefined; + jetstream: JetstreamManager; // | undefined; + spacedust: SpacedustManager; // | undefined; constructor(did: string) { this.did = did; @@ -80,6 +81,7 @@ class UserIndexServer { indexServerIndexer({ op, doer, + cid: msg.commit.cid, rev, aturi, value, @@ -90,7 +92,7 @@ class UserIndexServer { this.jetstream.start({ // for realsies pls get from db or something instead of this shit wantedDids: [ - this.did + this.did, // "did:plc:mn45tewwnse5btfftvd3powc", // "did:plc:yy6kbriyxtimkjqonqatv2rb", // "did:plc:zzhzjga3ab5fcs2vnsv2ist3", @@ -114,36 +116,46 @@ class UserIndexServer { //await connectToJetstream(this.did, this.db); this.spacedust = new SpacedustManager((msg: SpacedustLinkMessage) => { console.log("Received Spacedust message: ", msg); + const operation = msg.link.operation; + const sourceURI = new AtUri(msg.link.source_record); - const srcUri = msg.link.source_record - const srcDid = sourceURI.host - const srcField = msg.link.source - const srcCol = sourceURI.collection - const subjectURI = new AtUri(msg.link.subject) - const subUri = msg.link.subject - const subDid = subjectURI.host - const subCol = subjectURI.collection - - this.db.run( - `INSERT INTO backlink_skeleton ( - srcuri, - srcdid, - srcfield, - srccol, - suburi, - subdid, - subcol - ) VALUES (?, ?, ?, ?, ?, ?, ?)`, - [ - srcUri, // full AT URI of the source record - srcDid, // did: of the source - srcField, // e.g., "reply.parent.uri" or "facets.features.did" - srcCol, // e.g., "app.bsky.feed.post" - subUri, // full AT URI of the subject (linked record) - subDid, // did: of the subject - subCol, // subject collection (can be inferred or passed) - ] - ); + const srcUri = msg.link.source_record; + const srcDid = sourceURI.host; + const srcField = msg.link.source; + const srcCol = sourceURI.collection; + const subjectURI = new AtUri(msg.link.subject); + const subUri = msg.link.subject; + const subDid = subjectURI.host; + const subCol = subjectURI.collection; + + if (operation === "delete") { + this.db.run( + `DELETE FROM backlink_skeleton + WHERE srcuri = ? AND srcfield = ? AND suburi = ?`, + [srcUri, srcField, subUri] + ); + } else if (operation === "create") { + this.db.run( + `INSERT OR REPLACE INTO backlink_skeleton ( + srcuri, + srcdid, + srcfield, + srccol, + suburi, + subdid, + subcol + ) VALUES (?, ?, ?, ?, ?, ?, ?)`, + [ + srcUri, // full AT URI of the source record + srcDid, // did: of the source + srcField, // e.g., "reply.parent.uri" or "facets.features.did" + srcCol, // e.g., "app.bsky.feed.post" + subUri, // full AT URI of the subject (linked record) + subDid, // did: of the subject + subCol, // subject collection (can be inferred or passed) + ] + ); + } }); this.spacedust.start({ wantedSources: [ @@ -187,7 +199,7 @@ class UserIndexServer { } // initialize() { - + // } // async handleHttpRequest(route: string, req: Request): Promise { @@ -215,6 +227,7 @@ class UserIndexServer { } function openDbForDid(did: string): Database { + // TODO: we should disallow non users to open a db const path = `./dbs/${did}.sqlite`; const db = new Database(path); setupUserDb(db); @@ -477,22 +490,22 @@ const SQL = { links: ` SELECT srcuri, srcdid, srccol FROM backlink_skeleton - WHERE suburi = ? AND subcol = ? AND srcfield = ? + WHERE suburi = ? AND srccol = ? AND srcfield = ? `, distinctDids: ` SELECT DISTINCT srcdid FROM backlink_skeleton - WHERE suburi = ? AND subcol = ? AND srcfield = ? + WHERE suburi = ? AND srccol = ? AND srcfield = ? `, count: ` SELECT COUNT(*) as total FROM backlink_skeleton - WHERE suburi = ? AND subcol = ? AND srcfield = ? + WHERE suburi = ? AND srccol = ? AND srcfield = ? `, countDistinctDids: ` SELECT COUNT(DISTINCT srcdid) as total FROM backlink_skeleton - WHERE suburi = ? AND subcol = ? AND srcfield = ? + WHERE suburi = ? AND srccol = ? AND srcfield = ? `, all: ` SELECT suburi, srccol, COUNT(*) as records, COUNT(DISTINCT srcdid) as distinct_dids @@ -500,9 +513,13 @@ const SQL = { WHERE suburi = ? GROUP BY suburi, srccol `, +}; + +export function isDid(str: string): boolean { + return typeof str === "string" && str.startsWith("did:"); } -export async function constellationAPIHandler(req: Request, did: string): Promise { +export async function constellationAPIHandler(req: Request): Promise { const url = new URL(req.url); const pathname = url.pathname; // const bskyUrl = `https://api.bsky.app${pathname}${url.search}`; @@ -510,19 +527,29 @@ export async function constellationAPIHandler(req: Request, did: string): Promis // const constellationMethod = pathname.startsWith("/links") // ? pathname.slice("/links".length) // : null; - const searchParams = searchParamsToJson(url.searchParams); + const searchParams = searchParamsToJson(url.searchParams) as linksQuery; const jsonUntyped = searchParams; + const did = isDid(searchParams.target) + ? searchParams.target + : new AtUri(searchParams.target).host; const db = openDbForDid(did); switch (pathname) { case "/links": { const jsonTyped = jsonUntyped as linksQuery; // probably need to do pagination or something + console.log(JSON.stringify(jsonTyped, null, 2)); + const field = `${jsonTyped.collection}:${jsonTyped.path.replace( + /^\./, + "" + )}`; - const rows = db.prepare(SQL.links).all(jsonTyped.target, jsonTyped.collection, jsonTyped.path); + const rows = db + .prepare(SQL.links) + .all(jsonTyped.target, jsonTyped.collection, field); const linking_records: linksRecord[] = rows.map((row: any) => { - const rkey = row.srcuri.split('/').pop()!; + const rkey = row.srcuri.split("/").pop()!; return { did: row.srcdid, collection: row.srccol, @@ -590,12 +617,135 @@ export async function constellationAPIHandler(req: Request, did: string): Promis } } +function isImageEmbed(embed: unknown): embed is ATPAPI.AppBskyEmbedImages.Main { + return typeof embed === "object" && embed !== null && "$type" in embed && + (embed as any).$type === "app.bsky.embed.images"; +} + +function isVideoEmbed(embed: unknown): embed is ATPAPI.AppBskyEmbedVideo.Main { + return typeof embed === "object" && embed !== null && "$type" in embed && + (embed as any).$type === "app.bsky.embed.video"; +} + +function isRecordEmbed(embed: unknown): embed is ATPAPI.AppBskyEmbedRecord.Main { + return typeof embed === "object" && embed !== null && "$type" in embed && + (embed as any).$type === "app.bsky.embed.record"; +} + +function isRecordWithMediaEmbed(embed: unknown): embed is ATPAPI.AppBskyEmbedRecordWithMedia.Main { + return typeof embed === "object" && embed !== null && "$type" in embed && + (embed as any).$type === "app.bsky.embed.recordWithMedia"; +} + +function uncid(anything: any): string | null { + return (anything as Record)?.["$link"] as string | null || null; +} + +function extractImages(embed: unknown) { + if (isImageEmbed(embed)) return embed.images; + if (isRecordWithMediaEmbed(embed) && isImageEmbed(embed.media)) return embed.media.images; + return []; +} + +function extractVideo(embed: unknown) { + if (isVideoEmbed(embed)) return embed; + if (isRecordWithMediaEmbed(embed) && isVideoEmbed(embed.media)) return embed.media; + return null; +} + +function extractQuoteUri(embed: unknown): string | null { + if (isRecordEmbed(embed)) return embed.record.uri; + if (isRecordWithMediaEmbed(embed)) return embed.record.record.uri; + return null; +} + export function indexServerIndexer(ctx: indexHandlerContext) { - const record = validateRecord(ctx.value); + const record = assertRecord(ctx.value); + //const record = validateRecord(ctx.value); + const db = openDbForDid(ctx.doer); + console.log("indexering") switch (record?.$type) { case "app.bsky.feed.like": { return; } + case "app.bsky.feed.post": { + console.log("bsky post") + const stmt = db.prepare(` + INSERT OR IGNORE INTO app_bsky_feed_post ( + uri, did, cid, rev, createdat, indexedat, json, + text, replyroot, replyparent, quote, + imagecount, image1cid, image1mime, image1aspect, + image2cid, image2mime, image2aspect, + image3cid, image3mime, image3aspect, + image4cid, image4mime, image4aspect, + videocount, videocid, videomime, videoaspect + ) VALUES (?, ?, ?, ?, ?, ?, ?, + ?, ?, ?, ?, + ?, ?, ?, ?, + ?, ?, ?, + ?, ?, ?, + ?, ?, ?, + ?, ?, ?, ?) + `); + + const embed = record.embed; + + const images = extractImages(embed); + const video = extractVideo(embed); + const quoteUri = extractQuoteUri(embed); + try { + stmt.run( + ctx.aturi?? null, + ctx.doer?? null, + ctx.cid?? null, + ctx.rev?? null, + record.createdAt, + Date.now(), + JSON.stringify(record), + + record.text ?? null, + record.reply?.root?.uri ?? null, + record.reply?.parent?.uri ?? null, + + quoteUri, + + images.length, + uncid(images[0]?.image?.ref) ?? null, + images[0]?.image?.mimeType ?? null, + (images[0]?.aspectRatio && images[0].aspectRatio.width && images[0].aspectRatio.height) + ? `${images[0].aspectRatio.width}:${images[0].aspectRatio.height}` + : null, + + uncid(images[1]?.image?.ref) ?? null, + images[1]?.image?.mimeType ?? null, + (images[1]?.aspectRatio && images[1].aspectRatio.width && images[1].aspectRatio.height) + ? `${images[1].aspectRatio.width}:${images[1].aspectRatio.height}` + : null, + + uncid(images[2]?.image?.ref) ?? null, + images[2]?.image?.mimeType ?? null, + (images[2]?.aspectRatio && images[2].aspectRatio.width && images[2].aspectRatio.height) + ? `${images[2].aspectRatio.width}:${images[2].aspectRatio.height}` + : null, + + uncid(images[3]?.image?.ref) ?? null, + images[3]?.image?.mimeType ?? null, + (images[3]?.aspectRatio && images[3].aspectRatio.width && images[3].aspectRatio.height) + ? `${images[3].aspectRatio.width}:${images[3].aspectRatio.height}` + : null, + + uncid(video?.video) ? 1 : 0, + uncid(video?.video) ?? null, + uncid(video?.video) ? "video/mp4" : null, + video?.aspectRatio + ? `${video.aspectRatio.width}:${video.aspectRatio.height}` + : null + ); + } catch (err) { + console.error("stmt.run failed:", err); +} + return; + } default: { // what the hell return; diff --git a/main.ts b/main.ts index 54a3c0c..971a4dd 100644 --- a/main.ts +++ b/main.ts @@ -29,6 +29,25 @@ export const spacedusturl = Deno.env.get("SPACEDUST_URL"); export const systemDB = new Database("system.db"); setupSystemDb(systemDB); +// add me lol +systemDB.exec(` + INSERT OR IGNORE INTO users (did, role, registrationdate, onboardingstatus) + VALUES ( + 'did:plc:mn45tewwnse5btfftvd3powc', + 'admin', + datetime('now'), + 'ready' + ); + + INSERT OR IGNORE INTO users (did, role, registrationdate, onboardingstatus) + VALUES ( + 'did:web:did12.whey.party', + 'admin', + datetime('now'), + 'ready' + ); +`) + const userManager = new IndexServerUserManager(); userManager.coldStart(systemDB) @@ -127,11 +146,10 @@ Deno.serve( // const xrpcMethod = pathname.startsWith("/xrpc/") // ? pathname.slice("/xrpc/".length) // : null; - const constellationMethod = pathname.startsWith("/links") - ? pathname.slice("/links".length) - : null; + console.log(`request for "${pathname}"`) + const constellation = pathname.startsWith("/links") - if (constellationMethod) { + if (constellation) { return await constellationAPIHandler(req); } diff --git a/utils/dbsystem.ts b/utils/dbsystem.ts index 8a36d7d..3a1f0c8 100644 --- a/utils/dbsystem.ts +++ b/utils/dbsystem.ts @@ -33,21 +33,5 @@ export function setupSystemDb(db: Database) { handle TEXT ); ${createIndexINE} idx_did_handle ON did(handle); - - -- A global index for relationships between any two pieces of content - ${createTableINE} backlink_skeleton ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - srcuri TEXT, - srcdid TEXT, - srcfield TEXT, - srccol TEXT, - suburi TEXT, - subdid TEXT, - subcol TEXT - ); - ${createIndexINE} idx_backlink_subdid_mod ON backlink_skeleton(subdid, srcdid); - ${createIndexINE} idx_backlink_suburi_mod ON backlink_skeleton(suburi, srcdid); - ${createIndexINE} idx_backlink_subdid_filter_mod ON backlink_skeleton(subdid, srccol, srcdid); - ${createIndexINE} idx_backlink_suburi_filter_mod ON backlink_skeleton(suburi, srccol, srcdid); `); } \ No newline at end of file diff --git a/utils/dbuser.ts b/utils/dbuser.ts index e6e5183..8e4fc02 100644 --- a/utils/dbuser.ts +++ b/utils/dbuser.ts @@ -131,5 +131,21 @@ export function setupUserDb(db: Database) { -- User's notification settings declaration ${createTableINE} app_bsky_notification_declaration ( ${baseColumns}, allowSubscriptions TEXT ); ${createIndexINE} idx_notification_declaration_author ON app_bsky_notification_declaration(did); + + -- A global index for relationships between any two pieces of content + ${createTableINE} backlink_skeleton ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + srcuri TEXT, + srcdid TEXT, + srcfield TEXT, + srccol TEXT, + suburi TEXT, + subdid TEXT, + subcol TEXT + ); + ${createIndexINE} idx_backlink_subdid_mod ON backlink_skeleton(subdid, srcdid); + ${createIndexINE} idx_backlink_suburi_mod ON backlink_skeleton(suburi, srcdid); + ${createIndexINE} idx_backlink_subdid_filter_mod ON backlink_skeleton(subdid, srccol, srcdid); + ${createIndexINE} idx_backlink_suburi_filter_mod ON backlink_skeleton(suburi, srccol, srcdid); `); } \ No newline at end of file diff --git a/utils/records.ts b/utils/records.ts index 6cb2e2c..98f0fe3 100644 --- a/utils/records.ts +++ b/utils/records.ts @@ -75,7 +75,18 @@ export function validateRecord( if (result.success) return result.value; return undefined; } - +export function assertRecord( + record: unknown +): RecordTypeMap[T] | undefined { + if (typeof record !== 'object' || record === null) { + return undefined; + } + const type = (record as { $type?: string })?.$type; + if (typeof type !== 'string' || !(type in recordValidators)) { + return undefined; + } + return record as RecordTypeMap[T]; +} export async function resolveRecordFromURI({ did, uri,