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.
4.2 kB · 102 lines
TypeScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103import unique from "just-unique";import { AgentMap } from "../../classes/bskyAgents";import { ScheduledContext } from "../../classes/context";import { Post } from "../../classes/post";import { Repost } from "../../classes/repost";import { TaskType } from "../../enums";import type { QueueTaskData } from "../../types";import { userHasViolations } from "../db/violations";import { getEnumKeyByValue, has } from "../helpers";import { handlePostTask, handleRepostTask } from "../scheduler";import { enqueueEmptyWork } from "./queuePublisher";
interface BufferBlast { type: TaskType; time: number;}
export async function processQueue(batch: MessageBatch<QueueTaskData>, env: Env, ctx: ExecutionContext) { // runtime overhead const runtimeWrapper = new ScheduledContext(env, ctx); const agency = new AgentMap(env.TASK_SETTINGS);
// Retry settings const delay: number = env.QUEUE_SETTINGS.delay_val; const maxRetries: number = env.QUEUE_SETTINGS.max_retries; let bufferBlasts: BufferBlast[] = [];
for (const message of batch.messages) { let wasSuccess: boolean = false; const taskType: TaskType = message.body.type; if (taskType == TaskType.Post || taskType == TaskType.Repost) { if (message.body.data == null) { console.error(`task type ${getEnumKeyByValue(TaskType, taskType)} with no message body. cannot be processed!`); // maybe this was a bad send, so try it again later. Do not backblast as it was not an upstream failure. message.retry(); continue; }
// how we currently tell the difference is by searching for content label // because this object might be a copy of the class, or it could just be some arbitrary json because // the class didn't copy properly. const postDataObj: Post | Repost = has(message.body.data, "contentLabel") ? new Post(message.body.data as Post) : new Repost(message.body.data);
const userId = postDataObj.getUser(); const agent = await agency.getOrAddAgent(runtimeWrapper, userId, taskType); if (agent == null) { // if we could not get an agent for you, we should check to see if you have violations // if you do, we stop processing you. if (await userHasViolations(runtimeWrapper, userId)) { console.log(`User ${userId} has violations, dropping them from the queue`); message.ack(); continue; } else { console.warn(`Could not make an agent for ${userId}, got null.`); } } else { switch (taskType) { case TaskType.Post: wasSuccess = await handlePostTask(runtimeWrapper, postDataObj as Post, agent); break; case TaskType.Repost: wasSuccess = await handleRepostTask(runtimeWrapper, postDataObj, agent); break; } } } else if (taskType == TaskType.Blast) { console.log(`Got a blast message with ${batch.messages.length} messages in batch`); wasSuccess = true; } else { message.ack(); return; } // Handle queue acknowledgement on success/failure if (!wasSuccess) { const currentAttempts: number = message.attempts; const delaySeconds = delay * (currentAttempts + 1); console.log(`attempting to retry message ${taskType} in ${delaySeconds}`); message.retry({ delaySeconds: delaySeconds });
// if the attempts are over the maximum amount of retries then do not backblast if (currentAttempts > maxRetries) continue;
// push a backblast so that this item will retry in the future. // it basically just writes null in the buffer, which is silly but w/e bufferBlasts.push({ type: taskType, time: delaySeconds }); } else { message.ack(); } } // If we have any retries, they'll only get delivered on next batch // so we're going to back blast the buffer queue so that we can make sure the retries go. if (bufferBlasts.length > 0) { bufferBlasts = unique(bufferBlasts); console.log(`Attempting to backblast ${bufferBlasts.length} items`); for (const blast of bufferBlasts) { await enqueueEmptyWork(runtimeWrapper, blast.type, blast.time + 10); } }}