Something went wrong. Try again.
Schedule posts and reposts to Bluesky (and other PDSes) with Cloudflare workers. skyscheduler.work
scheduling social-media cloudflare bsky service cloudflare-workers bsky-tool bluesky
Something went wrong. Try again.
24 kB · 591 lines
TypeScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592import { addHours, isAfter, isEqual } from "date-fns";import { and, asc, desc, eq, getTableColumns, gt, gte, ne, sql } from "drizzle-orm";import isEmpty from "just-is-empty";import { v4 as uuidv4 } from "uuid";import { APP_NAME } from "../appInfo";import { Post } from "../classes/post";import { RepostInfo } from "../classes/repost";import { mediaFiles, posts, repostCounts, reposts } from "../db/app.schema";import { accounts, users } from "../db/auth.schema";import { AccountStatus, PostLabel, RepostType, TimeShape } from "../enums";import { MAX_POSTS_PER_THREAD, MAX_REPOST_POSTS, MAX_REPOST_RULES_PER_POST } from "../limits";import type { AccountUpdatePayload, AllContext, BatchQuery, BatchQueryArray, CreateObjectResponse, CreatePostQueryResponse, DBProcessor, DeleteResponse, DeleteScheduleResponse, EditPostChanges, UserIdType} from "../types";import { PostSchema } from "../validation/postSchema";import { RepostSchema } from "../validation/repostSchema";import { getMentionsFromContent } from "./bsky/mentions";import { getChildPostsOfThread, getPostByCID, getPostThreadCount, getRepostCountQuery, updatePostForGivenUser} from "./db/data";import { getViolationsForUser, removeViolationsDB } from "./db/violations";import { has, isAltEditableType, isUUIDValid } from "./helpers";import { deleteEmbedsFromR2 } from "./r2Query";import { floorGivenTime } from "./time";
export const getPostsForUser = async (c: AllContext, attempts: number = 0): Promise<Post[]|null> => { if (attempts >= 3) { return null; } try { const userId: UserIdType = c.get("userId"); const db: DBProcessor = c.get("db"); if (userId && db) { const results = await db.select({ ...getTableColumns(posts), repostCount: repostCounts.count }) .from(posts).where(eq(posts.userId, userId)) .leftJoin(repostCounts, eq(posts.uuid, repostCounts.uuid)) .orderBy(desc(posts.scheduledDate), asc(posts.threadOrder), desc(posts.createdAt)).all();
if (isEmpty(results)) return null;
return results.map((itm) => new Post(itm)); } } catch (err: unknown) { console.error("Failed to get posts for user: " + String(err)); return await getPostsForUser(c, attempts + 1); } return null;};
export const updateUserData = async (c: AllContext, newData: AccountUpdatePayload): Promise<boolean> => { const userId: UserIdType = c.get("userId"); const db: DBProcessor = c.get("db"); try { if (!db) { return false; } if (userId) { const queriesToExecute: BatchQueryArray = [];
if (has(newData, "password")) { // cache out the new hash const newPassword = newData.password; // remove it from the original object delete newData.password;
// add the query to the db batch object queriesToExecute.push(db.update(accounts) .set({password: newPassword}) .where(eq(accounts.userId, userId))); }
// If we have new data about the username, pds, or password if (has(newData, "updatedSession") || has(newData, "username")) { await removeViolationsDB(db, userId, [AccountStatus.InvalidAccount, AccountStatus.Deactivated]); delete newData.updatedSession; }
if (!isEmpty(newData)) { queriesToExecute.push(db.update(users).set(newData) .where(eq(users.id, userId))); }
if (queriesToExecute.length > 0) await db.batch(queriesToExecute as BatchQuery); return true; } } catch(_err) { console.error(`Failed to update new user data for user ${userId}`); } return false;};
export const deletePost = async (c: AllContext, id: string): Promise<DeleteResponse> => { const userId: UserIdType = c.get("userId"); const returnObj: DeleteResponse = {success: false, isRepost: false}; if (!userId) { return returnObj; }
const db: DBProcessor = c.get("db"); if (!db) { console.error(`unable to delete post ${id}, db was null`); return returnObj; }
const postObj = await getPostById(c, id); if (postObj !== null) { const queriesToExecute: BatchQueryArray = []; // If the post has not been posted, that means we still have files for it, so // delete the files from R2 if (!postObj.posted) { await deleteEmbedsFromR2(c, postObj.embedContent).then(() => removeViolationsDB(db, userId, [AccountStatus.MediaTooBig])); } returnObj.isRepost = postObj.isRepost ?? false;
// If the parent post is not null, then attempt to find and update the post chain const parentPost = postObj.parentPost; if (parentPost !== undefined) { // set anyone who had this as their parent to this post chain queriesToExecute.push(db.update(posts).set({parentPost: parentPost, threadOrder: postObj.threadOrder}) .where(and(eq(posts.parentPost, postObj.uuid), eq(posts.rootPost, postObj.rootPost!))));
// Update the post order past here queriesToExecute.push(db.update(posts).set({threadOrder: sql`threadOrder - 1`}) .where( and(and(eq(posts.rootPost, postObj.rootPost!), ne(posts.threadOrder, -1)), gt(posts.threadOrder, postObj.threadOrder) ))); }
// We'll need to delete all of the child embeds then, a costly, annoying experience. if (postObj.isThreadRoot) { const childPosts = await getChildPostsOfThread(c, postObj.uuid); if (childPosts !== null) { for (const childPost of childPosts) { c.executionCtx.waitUntil(deleteEmbedsFromR2(c, childPost.embedContent)); queriesToExecute.push(db.delete(posts).where(eq(posts.uuid, childPost.uuid))); } } else { console.warn(`could not get child posts of thread ${postObj.uuid} during delete`); } } else if (postObj.isChildPost) { // this is not a thread root, so we should figure out how many children are left. const childPostCount = (await getPostThreadCount(db, postObj.userId, postObj.rootPost!)) - 1; if (childPostCount <= 0) { queriesToExecute.push(db.update(posts).set({threadOrder: -1}).where(eq(posts.uuid, postObj.rootPost!))); } }
// delete post queriesToExecute.push(db.delete(posts).where(eq(posts.uuid, id))); c.executionCtx.waitUntil(db.batch(queriesToExecute as BatchQuery)); returnObj.success = true; returnObj.wasThreadRoot = postObj.isThreadRoot; } return returnObj;};
export const createPost = async (c: AllContext, body: unknown): Promise<CreatePostQueryResponse> => { const db: DBProcessor = c.get("db"); const userId: UserIdType = c.get("userId"); if (!userId) return { ok: false, msg: "Your user session has expired, please login again"};
if (!db) { return { ok: false, msg: "An application error has occurred please refresh" }; }
const validation = PostSchema.safeParse(body); if (!validation.success) { return { ok: false, msg: validation.error.toString() }; }
const { content, scheduledDate, embeds, contentLabel, makePostNow, repostData, rootPost, parentPost } = validation.data; const scheduleDate = floorGivenTime((makePostNow) ? new Date() : new Date(scheduledDate), TimeShape.Post);
// Check if account is in violation const violationData = await getViolationsForUser(db, userId); if (violationData != null) { if (violationData.tosViolation) { return {ok: false, msg: `This account is unable to use ${APP_NAME} services at this time`}; } else if (violationData.userPassInvalid) { return {ok: false, msg: "The BSky account credentials is invalid, please update these in the settings"}; } }
// Check to see if this post already exists for thread let rootPostID:string|undefined = undefined; let parentPostID:string|undefined = undefined; let rootPostData: Post|null = null; let parentPostOrder: number = 0; if (isUUIDValid(rootPost)) { // returns null if the post doesn't appear on this account rootPostData = await getPostById(c, rootPost!); if (rootPostData !== null) { if (rootPostData.posted) { return { ok: false, msg: "You cannot make threads off already posted posts"}; } if (rootPostData.isChildPost) { return { ok: false, msg: "Subthreads of threads are not allowed." }; } if (rootPostData.isRepost) { return {ok: false, msg: "Threads cannot be made of repost actions"}; } rootPostID = rootPostData.rootPost ?? rootPostData.uuid; // If this isn't a direct reply, check directly underneath it if (rootPost !== parentPost) { if (isUUIDValid(parentPost)) { const parentPostData = await getPostById(c, parentPost!); if (parentPostData !== null) { parentPostID = parentPost!; parentPostOrder = parentPostData.threadOrder + 1; } else { return { ok: false, msg: "The given parent post cannot be found on your account"}; } } else { return { ok: false, msg: "The given parent post is invalid"}; } } else { parentPostID = rootPostData.uuid; parentPostOrder = 1; // Root will always be 0, so if this is root, go 1 up. } } else { return { ok: false, msg: "The given root post cannot be found on your account"}; } }
const isThreadedPost: boolean = (rootPostID !== undefined && parentPostID !== undefined); if (isThreadedPost) { const threadCount: number = await getPostThreadCount(db, userId, rootPostID!); if (threadCount >= MAX_POSTS_PER_THREAD) { return { ok: false, msg: `this thread has hit the limit of ${MAX_POSTS_PER_THREAD} posts per thread`}; } }
// Create repost metadata const scheduleGUID = (!isThreadedPost) ? uuidv4() : undefined; const repostInfo = (!isThreadedPost && repostData !== undefined) ? new RepostInfo(scheduleGUID!, scheduleDate, false, repostData) : undefined;
// Create the posts const postUUID = uuidv4(); const dbOperations: BatchQueryArray = [];
// if we're threaded, insert our post before the given parent if (isThreadedPost) { // Update the parent to our new post dbOperations.push(db.update(posts).set({parentPost: postUUID }) .where(and(eq(posts.parentPost, parentPostID!), eq(posts.rootPost, rootPostID!))));
// update all posts past this one to also update their order (we will take their id) dbOperations.push(db.update(posts).set({threadOrder: sql`threadOrder + 1`}) .where( and(and(eq(posts.rootPost, rootPostID!), ne(posts.threadOrder, -1)), gte(posts.threadOrder, parentPostOrder) )));
// Update the root post so that it has the correct flags set on it as well. if (!rootPostData!.isThreadRoot) { dbOperations.push(db.update(posts).set({threadOrder: 0, rootPost: rootPostData!.uuid}) .where(eq(posts.uuid, rootPostData!.uuid))); } } else { rootPostID = postUUID; }
// create the mentions cache for the initial post. const mentions = await getMentionsFromContent(content);
// Add the post to the DB dbOperations.push(db.insert(posts).values({ content, uuid: postUUID, postNow: makePostNow, scheduledDate: (!isThreadedPost) ? scheduleDate : new Date(rootPostData!.scheduledDate!), rootPost: rootPostID, parentPost: parentPostID, mentionsCache: mentions, repostInfo: (!isThreadedPost && repostInfo !== undefined) ? [repostInfo] : [], threadOrder: (!isThreadedPost) ? undefined : parentPostOrder, embedContent: embeds, contentLabel: contentLabel ?? PostLabel.None, userId: userId }));
if (!isEmpty(embeds)) { // Loop through all data within an embed blob so we can mark it as posted for (const embed of embeds!) { if (isAltEditableType(embed.type)) { dbOperations.push( db.update(mediaFiles).set({hasPost: true}).where(eq(mediaFiles.fileName, embed.content))); } } }
// Add repost data to the table if (repostData && !isThreadedPost) { for (let i = 1; i <= repostData.times; ++i) { dbOperations.push(db.insert(reposts).values({ uuid: postUUID, scheduleGuid: scheduleGUID, scheduledDate: addHours(scheduleDate, i*repostData.hours) })); } // Push the repost counts in dbOperations.push(db.insert(repostCounts) .values({uuid: postUUID, count: repostData.times})); }
// Batch the query const batchResponse: ProperD1Result[] = await db.batch(dbOperations as BatchQuery); const success = batchResponse.every((el) => el.success); return { ok: success, postNow: makePostNow, postId: postUUID, msg: success ? "success" : "fail" };};
export const createRepost = async (c: AllContext, body: unknown): Promise<CreateObjectResponse> => { const db: DBProcessor = c.get("db");
const userId: UserIdType = c.get("userId"); if (!userId) return { ok: false, msg: "Your user session has expired, please login again"};
if (!db) { return {ok: false, msg: "Invalid server operation occurred, please refresh"}; }
const validation = RepostSchema.safeParse(body); if (!validation.success) { return { ok: false, msg: validation.error.toString() }; } const { data, scheduledDate, repostData } = validation.data; const isScheduledPost = (data.type === RepostType.FuturePost); const scheduleDate = floorGivenTime(new Date(scheduledDate), TimeShape.Repost);
// Check if account is in violation const violationData = await getViolationsForUser(db, userId); if (violationData != null) { if (violationData.tosViolation) { return {ok: false, msg: `This account is unable to use ${APP_NAME} services at this time`}; } else if (violationData.userPassInvalid) { return {ok: false, msg: "The BSky account credentials is invalid, please update these in the settings"}; } } let postUUID; const dbOperations: BatchQueryArray = []; const scheduleGUID = uuidv4(); const repostInfo: RepostInfo = new RepostInfo(scheduleGUID, scheduleDate, true, repostData);
// Check to see if the post already exists // (check also against the userId here as well to avoid cross account data collisions) const existingPost = (isScheduledPost) ? await getPostById(c, data.id) : await getPostByCID(db, userId, data.cid);
// if we were expecting a future post, check to see if the post data is invalid, if so, then // give an error as a post that is not tied to the user's account attempted to be updated if (existingPost === null && isScheduledPost) { return { ok: false, msg: "Invalid post id"}; } if (existingPost !== null) { postUUID = existingPost.uuid; const existingPostDate = existingPost.scheduledDate!; // Ensure the date asked for is after what the post's schedule date is if (!isAfter(scheduleDate, existingPostDate) && !isEqual(scheduledDate, existingPostDate)) { return { ok: false, msg: "Scheduled date must be after the initial post's date" }; } // Make sure this isn't a thread post. // We could probably work around this but I don't think it's worth the effort. if (existingPost.isChildPost) { return {ok: false, msg: "Repost posts cannot be created from child thread posts"}; }
// Add repost info object to existing array const updatedRepostInfo: RepostInfo[] = isEmpty(existingPost.repostInfo) ? [] : existingPost.repostInfo!; if (updatedRepostInfo.length >= MAX_REPOST_RULES_PER_POST) { return {ok: false, msg: `Num of reposts rules for this post has exceeded the limit of ${MAX_REPOST_RULES_PER_POST} rules`}; }
// Cache a quick representation of the ISO string so we can check Date/str differences const repostInfoTimeStr = scheduleDate.toISOString(); // Check to see if we have an exact repost match. // If we do, do not update the repostInfo, as repost table will drop the duplicates for us anyways. const isNewInfoNotDuped = (el: RepostInfo) => { if (el.time === repostInfoTimeStr || el.time === repostInfo.time) { if (el.count == repostInfo.count) { return el.hours != repostInfo.hours; } } return true; }; // Check to see if existing data matches with our new data if (updatedRepostInfo.every(isNewInfoNotDuped)) { // it does not, so we can add it to the DB updatedRepostInfo.push(repostInfo);
const repostInfoUpdateQuery = db.update(posts).set({repostInfo: updatedRepostInfo}); // push record update to add to json array if (!isScheduledPost) { dbOperations.push(repostInfoUpdateQuery.where(and( eq(posts.userId, userId), eq(posts.cid, data.cid)))); } else { dbOperations.push(repostInfoUpdateQuery.where(and( eq(posts.userId, userId), eq(posts.uuid, data.id)))); } } } else { // Limit of post reposts on the user's account. const accountCurrentReposts = await db.$count(posts, and(eq(posts.userId, userId), eq(posts.isRepost, true))); if (MAX_REPOST_POSTS > 0 && accountCurrentReposts >= MAX_REPOST_POSTS) { return {ok: false, msg: `You've cannot create any more repost posts at this time. Using: (${accountCurrentReposts}/${MAX_REPOST_POSTS}) repost posts`}; }
if (isScheduledPost) { return {ok: false, msg: "Invalid data request"}; }
// Create the post base for this repost postUUID = uuidv4(); dbOperations.push(db.insert(posts).values({ content: !isEmpty(data.content) ? data.content! : `Repost of ${data.url}`, uuid: postUUID, cid: data.cid, uri: data.uri, posted: true, isRepost: true, repostInfo: [repostInfo], scheduledDate: scheduleDate, userId: userId })); }
// Push initial repost let totalRepostCount = 1; dbOperations.push(db.insert(reposts).values({ uuid: postUUID, scheduleGuid: scheduleGUID, scheduledDate: scheduleDate }).onConflictDoNothing());
// Push other repost times if we have them if (repostData) { for (let i = 1; i <= repostData.times; ++i) { dbOperations.push(db.insert(reposts).values({ uuid: postUUID, scheduleGuid: scheduleGUID, scheduledDate: addHours(scheduleDate, i*repostData.hours) }).onConflictDoNothing()); } totalRepostCount += repostData.times; } // Update repost counts if (existingPost !== null) { // update existing content posts (but only for reposts, no one else) if (existingPost.isRepost && !isScheduledPost && !isEmpty(data.content)) { dbOperations.push(db.update(posts).set({content: data.content!}).where(eq(posts.uuid, postUUID))); }
// Because there could be conflicts that drop, run a count on the entire list and use the value from that // we also don't know if the repost count table has repost values for this item, so we should // attempt to always insert and update if it already exists totalRepostCount = -1; }
// pushing any value under zero causes a full recount dbOperations.push(getRepostCountQuery(db, postUUID, totalRepostCount));
const batchResponse: ProperD1Result[] = await db.batch(dbOperations as BatchQuery); const success = batchResponse.every((el) => el.success); return { ok: success, msg: success ? "success" : "fail", postId: postUUID };};
export const updatePostForUser = async (c: AllContext, id: string, newData: EditPostChanges): Promise<boolean> => { const userId: UserIdType = c.get("userId"); if (!userId) return false; return updatePostForGivenUser(c, userId, id, newData);};
export const getPostById = async(c: AllContext|undefined, id: string): Promise<Post|null> => { if (c === undefined) return null;
const userId: UserIdType = c.get("userId"); if (!userId || !isUUIDValid(id)) return null;
const db: DBProcessor = c.get("db"); if (!db) { console.error(`unable to get post ${id}, db was null`); return null; }
const result = await db.select().from(posts) .where(and(eq(posts.uuid, id), eq(posts.userId, userId))) .limit(1).all();
if (!isEmpty(result)) return new Post(result[0]); return null;};
// used for post editing, acts very similar to getPostsForUserexport const getPostByIdWithReposts = async(c: AllContext, id: string): Promise<Post|null> => { const userId: UserIdType = c.get("userId"); if (!userId || !isUUIDValid(id)) return null;
const db: DBProcessor = c.get("db"); if (!db) { console.error(`unable to get post ${id} with reposts, db was null`); return null; }
const result = await db.select({ ...getTableColumns(posts), repostCount: repostCounts.count, }).from(posts) .where(and(eq(posts.uuid, id), eq(posts.userId, userId))) .leftJoin(repostCounts, eq(posts.uuid, repostCounts.uuid)) .limit(1).all();
if (!isEmpty(result)) return new Post(result[0]); return null;};
export const deleteRepostRule = async(c: AllContext, id: string, scheduleId: string): Promise<DeleteScheduleResponse> => { const db: DBProcessor = c.get("db"); if (!db) { console.error(`unable to delete schedule id ${scheduleId} from post ${id}, db was null`); return { success: false }; } if (!isUUIDValid(scheduleId)) { return { success: false }; }
// Get the post to make sure it's valid and update post json const currentPost = await getPostByIdWithReposts(c, id); if (currentPost?.repostInfo !== undefined) { const originalRuleLength: number = currentPost.repostInfo.length; // remove the schedule from the current json object set const newRepostInfo: RepostInfo[] = currentPost.repostInfo.filter((itm) => { return itm.guid !== scheduleId; });
// Was this schedule id in the repostInfo array originally? if (newRepostInfo.length == originalRuleLength) { // It was not, so don't do anything more. return {success: false }; }
const queriesToExecute: BatchQueryArray = []; // modify the current repost info queriesToExecute.push(db.update(posts).set({repostInfo: newRepostInfo}).where(and( eq(posts.userId, currentPost.userId), eq(posts.uuid, currentPost.uuid))));
// Delete batch schedule items // we don't bundle this one because we want to get a count to make the operation below it, better const deletedItems = await db.delete(reposts).where(eq(reposts.scheduleGuid, scheduleId)).returning({date: reposts.scheduledDate});
const newRepostCount: number = currentPost.repostCount! - deletedItems.length; // did we delete anything at all? if (deletedItems.length <= 0) { // Log this out, but allow for the bad data to be deleted anyways console.warn(`When trying to delete reposts for ${currentPost.uuid}, schedule id ${scheduleId} had empty items`); } else { // Force update the repost count :) queriesToExecute.push(getRepostCountQuery(db, id, newRepostCount)); }
// Batch push up everything const batchResponse: ProperD1Result[] = await db.batch(queriesToExecute as BatchQuery); if (batchResponse.every((el) => el.success)) { // This was successful, so we should write the new post data into our current object // and return it, so that it can be used for rendering currentPost.repostCount = newRepostCount; currentPost.repostInfo = newRepostInfo; return { success: true, postData: currentPost }; } } return { success: false };};