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.
12 kB · 320 lines
TypeScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320import { and, asc, desc, eq, inArray, isNotNull, lte, ne, notInArray, or, sql} from "drizzle-orm";import type { DrizzleD1Database } from "drizzle-orm/d1";import isEmpty from "just-is-empty";import { Post } from "../../classes/post";import { Repost } from "../../classes/repost";import { TRUNCATE_POSTED_CONTENT } from "../../config";import { posts, repostCounts, reposts } from "../../db/app.schema";import { violations } from "../../db/enforcement.schema";import { TimeShape } from "../../enums";import { MAX_HOLD_DAYS_BEFORE_PURGE, MAX_POSTED_LENGTH } from "../../limits";import type { AllContext, BatchQuery, BatchQueryArray, DBProcessor, EditPostChanges, GetAllPostedBatch, PostRecordResponse} from "../../types";import { isUUIDValid } from "../helpers";import { floorCurrentTime } from "../time";
export const getAllPostsForCurrentTime = async (c: AllContext, removeThreads: boolean = false): Promise<Post[]> => { // Get all scheduled posts for current time const db: DBProcessor = c.get("db"); if (!db) { return []; } const currentTime: Date = floorCurrentTime(TimeShape.Post);
const violationUsers = db.select({violators: violations.userId}).from(violations); const postsToMake = db.$with('scheduledPosts').as(db.select().from(posts) .where( and( and( and( eq(posts.posted, false), ne(posts.postNow, true) // Ignore any posts that are marked for post now ), lte(posts.scheduledDate, currentTime) ), // ignore threads, we'll create this one later. removeThreads ? eq(posts.threadOrder, -1) : lte(posts.threadOrder, 0) ) )); const results = await db.with(postsToMake).select().from(postsToMake) .where(notInArray(postsToMake.userId, violationUsers)).orderBy(asc(postsToMake.createdAt)).all(); return results.map((item) => new Post(item));};
export const getAllRepostsForGivenTime = async (c: AllContext, givenDate: Date): Promise<Repost[]> => { // Get all scheduled posts for the given time const db: DBProcessor = c.get("db"); if (!db) { return []; } const query = db.select({uuid: reposts.uuid}).from(reposts) .where(lte(reposts.scheduledDate, givenDate)); const violationsQuery = db.select({data: violations.userId}).from(violations); const results = await db.select({uuid: posts.uuid, uri: posts.uri, cid: posts.cid, userId: posts.userId }) .from(posts) .where(and(inArray(posts.uuid, query), notInArray(posts.userId, violationsQuery))) .all();
return results.map((item) => new Repost(item));};
export const getAllRepostsForCurrentTime = async (c: AllContext): Promise<Repost[]> => { return getAllRepostsForGivenTime(c, floorCurrentTime(TimeShape.Repost));};
export const deleteAllRepostsBeforeCurrentTime = async (c: AllContext) => { const db: DBProcessor = c.get("db"); if (!db) { return; } const currentTime = floorCurrentTime(TimeShape.Repost); const deletedPosts = await db.delete(reposts).where(lte(reposts.scheduledDate, currentTime)) .returning({id: reposts.uuid, scheduleGuid: reposts.scheduleGuid});
// This is really stupid and I hate it, but someone has to update repost counts once posted if (deletedPosts.length > 0) { const batchedQueries: BatchQueryArray = []; for (const deleted of deletedPosts) { // Update counts const newCount = db.$count(reposts, eq(reposts.uuid, deleted.id)); batchedQueries.push(db.update(repostCounts) .set({count: newCount}) .where(eq(repostCounts.uuid, deleted.id)));
// check if the repost data needs to be killed if (!isEmpty(deleted.scheduleGuid)) { // do a search to find if there are any reposts with the same scheduleguid. // if there are none, this schedule should get removed from the repostInfo array const stillHasSchedule = await db.select().from(reposts) .where(and( eq(reposts.scheduleGuid, deleted.scheduleGuid!), eq(reposts.uuid, deleted.id))) .limit(1).all();
// if this is empty, then we need to update the repost info. if (isEmpty(stillHasSchedule)) { // get the existing repost info to filter out this old data const existingRepostInfoArr = (await db.select({repostInfo: posts.repostInfo}).from(posts) .where(eq(posts.uuid, deleted.id)).limit(1).all())[0]; // check to see if there is anything in the repostInfo array if (!isEmpty(existingRepostInfoArr)) { // create a new array with the deleted out object const newRepostInfoArr = existingRepostInfoArr.repostInfo!.filter((obj) => { return obj.guid !== deleted.scheduleGuid!; }); // push the new repost info array batchedQueries.push(db.update(posts).set({repostInfo: newRepostInfoArr}).where(eq(posts.uuid, deleted.id))); } } } } if (batchedQueries.length > 0) await db.batch(batchedQueries as BatchQuery); }};
export const bulkUpdatePostedData = async (c: AllContext, records: PostRecordResponse[], allPosted: boolean) => { const db: DBProcessor = c.get("db"); if (!db) { return; } const dbOperations: BatchQueryArray = [];
for (let i = 0; i < records.length; ++i) { const record = records[i]; // skip over invalid records if (record.postID === null) continue;
const wasPosted = (i == 0 && !allPosted) ? false : true; dbOperations.push(db.update(posts).set( {content: TRUNCATE_POSTED_CONTENT ? sql`substr(posts.content, 0, ${MAX_POSTED_LENGTH+1})` : posts.content, posted: wasPosted, uri: record.uri, cid: record.cid, embedContent: [], mentionsCache: []}) .where(eq(posts.uuid, record.postID))); }
if (dbOperations.length > 0) await db.batch(dbOperations as BatchQuery);};
export const setPostNowOffForPost = async (c: AllContext, id: string) => { const db: DBProcessor = c.get("db"); if (!isUUIDValid(id)) return false;
if (!db) { console.warn(`cannot set off post now for post ${id}`); return false; }
const result = await db.update(posts).set({postNow: false}).where(eq(posts.uuid, id)).limit(1).returning({updated_id: posts.uuid}); if (!isUUIDValid(result[0].updated_id)) console.error(`Unable to set PostNow to off for post ${id}`);};
export const updatePostForGivenUser = async (c: AllContext, userId: string, id: string, newData: EditPostChanges) => { const db: DBProcessor = c.get("db"); if (isEmpty(userId) || !isUUIDValid(id)) return false;
if (!db) { console.error(`unable to update post ${id} for user ${userId}, db was null`); return false; }
const {success} = await db.update(posts).set(newData).where( and(eq(posts.uuid, id), eq(posts.userId, userId))); return success;};
export const getAllPostedPostsOfUser = async (c: AllContext, userId: string): Promise<GetAllPostedBatch[]> => { const db: DBProcessor = c.get("db"); if (isEmpty(userId)) return [];
if (!db) { console.error(`unable to get all posted posts of user ${userId}, db was null`); return []; }
return db.select({id: posts.uuid, uri: posts.uri}) .from(posts) .where(and(eq(posts.userId, userId), eq(posts.posted, true))) .all();};
export const getAllPostedPosts = async (c: AllContext): Promise<GetAllPostedBatch[]> => { const db: DBProcessor = c.get("db"); if (!db) { return []; } return db.select({id: posts.uuid, uri: posts.uri}) .from(posts) .where(eq(posts.posted, true)) .all();};
export const isPostAlreadyPosted = async (c: AllContext, postId: string): Promise<boolean> => { const db: DBProcessor = c.get("db"); if (!isUUIDValid(postId)) return true;
if (!db) { console.error(`unable to get database to tell if ${postId} has been posted`); return true; }
const query = await db.select({posted: posts.posted}).from(posts).where(eq(posts.uuid, postId)).all(); if (isEmpty(query) || query[0].posted === null) { // if the post does not exist, return true anyways return true; } return query[0].posted;};
export const getChildPostsOfThread = async (c: AllContext, rootId: string): Promise<Post[]|null> => { const db: DBProcessor = c.get("db"); if (!isUUIDValid(rootId)) return null;
if (!db) { console.error(`unable to get child posts of root ${rootId}, db was null`); return null; }
const query = await db.select().from(posts) .where(and(isNotNull(posts.parentPost), eq(posts.rootPost, rootId))) .orderBy(asc(posts.threadOrder), desc(posts.createdAt)).all(); if (query.length > 0) { return query.map((child) => new Post(child)); } return null;};
export const getPostThreadCount = async (db: DrizzleD1Database, userId: string, rootId: string): Promise<number> => { if (!isUUIDValid(rootId)) return 0;
return db.$count(posts, and( eq(posts.rootPost, rootId), eq(posts.userId, userId)));};
// deletes multiple posted posts from a database. Posts must be already posted as this does// no R2 db queries to cleanexport const deletePosts = async (c: AllContext, postsToDelete: string[]): Promise<number> => { // Don't do anything on empty arrays. if (isEmpty(postsToDelete)) return 0;
const db: DBProcessor = c.get("db"); if (!db) { console.error(`could not delete posts ${postsToDelete.toString()}, db was null`); return 0; } const deleteQueries: BatchQueryArray = []; postsToDelete.forEach((itm) => { // this will wipe out any posts and their children if they are marked for delete deleteQueries.push(db.delete(posts).where( and( or(eq(posts.uuid, itm), eq(posts.rootPost, itm)), eq(posts.posted, true)))); });
// Batching this should improve db times if (deleteQueries.length > 0) { const batchResponse: ProperD1Result[] = await db.batch(deleteQueries as BatchQuery); // Return the number of items that have been deleted return batchResponse.reduce((val: number, item: ProperD1Result) => val + (item.success ? 1 : 0), 0); } return 0;};
export const purgePostedPosts = async (c: AllContext): Promise<number> => { const db: DBProcessor = c.get("db"); if (!db) { return 0; } const positiveDays: number = Math.abs(MAX_HOLD_DAYS_BEFORE_PURGE); const dbQuery = await db.select({ data: posts.uuid }).from(posts) .leftJoin(repostCounts, eq(posts.uuid, repostCounts.uuid)) .where( and( and( eq(posts.posted, true), lte(posts.updatedAt, sql.raw(`datetime('now', '-${positiveDays} days')`)) //eq(posts.posted, true), lte(posts.updatedAt, sql`datetime('now', '-7 days')`) ), // skip child posts objects, only get us root posts and non-threads and(lte(posts.threadOrder, 0), lte(repostCounts.count, 0)) ) ).all(); const postsToDelete = dbQuery.map((item) => { return item.data }); if (isEmpty(postsToDelete)) return 0;
return deletePosts(c, postsToDelete);};
export const getPostByCID = async(db: DrizzleD1Database, userId: string, cid: string): Promise<Post|null> => { const result = await db.select().from(posts) .where(and(eq(posts.userId, userId), eq(posts.cid, cid))) .limit(1).all();
if (!isEmpty(result)) return new Post(result[0]); return null;};
export const getRepostCountQuery = (db: DrizzleD1Database, postUUID: string, newValue: number = -1) => { // if we're given any value underneath 0, we need to recount for the entire post const newCount = (newValue < 0) ? db.$count(reposts, eq(reposts.uuid, postUUID)) : newValue; return db.insert(repostCounts) .values({uuid: postUUID, count: newCount}) .onConflictDoUpdate({target: repostCounts.uuid, set: {count: newCount}});};