Something went wrong. Try again.
Automations and webhooks for the AT Protocol airglow.run
automation webhook atproto atprotocol
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157// Executes one write-queue row against the owner's PDS, re-verifying intent// from the mirror at execution time (KTD-7). Returns what happened; the runner// persists the row state.
import { createArbitraryRecord, deleteArbitraryRecord, generateTid } from "../automations/pds.js";import { PdsHttpError } from "../automations/errors.js";import { profileStatus } from "../automations/profile-status.js";import { db } from "../db/index.js";import { syncQueue } from "../db/schema.js";import { eq } from "drizzle-orm";import { classifyPdsError, isRecordAlreadyExists, type SyncPdsFailure } from "./pds-errors.js";import { deleteMirrorRow, mirrorRowsForSubject, removeWaiting, takeMirrorByRkey, upsertMirror, upsertWaiting, type MirrorRow, type QueueRow, type SyncSet,} from "./store.js";import { POINTS, PROFILE_WATCH_APPS, followCollection } from "./targets.js";
const DAY_MS = 24 * 60 * 60 * 1000;/** Re-check spacing for a waiting subject. Apps whose profile creates are * watched live only need a slow safety net. */export function waitingRecheckMs(app: QueueRow["app"]): number { return PROFILE_WATCH_APPS.includes(app) ? 7 * DAY_MS : 2 * DAY_MS;}
export type ExecuteOutcome = | { kind: "written"; points: number } | { kind: "deleted"; points: number } | { kind: "skipped"; reason: "already_following" | "no_longer_followed" | "record_missing"; } | { kind: "waiting" } | { kind: "failed"; failure: SyncPdsFailure };
export async function executeRow(row: QueueRow, set: SyncSet, now: Date): Promise<ExecuteOutcome> { return row.kind === "follow" ? executeFollow(row, set, now) : executeUnfollow(row, now);}
/** Whether the follow is still wanted, read from the mirror. Returns the skip * outcome when it is not, else the row's own mirror row from an earlier * attempt (if any). */function followIntent( row: QueueRow, set: SyncSet,): { skip: ExecuteOutcome } | { own: MirrorRow | undefined } { const mirror = mirrorRowsForSubject(row.did, row.subject); // A mirror row carrying this row's own pre-generated rkey is an earlier // attempt of this same write, not an existing follow. const own = row.rkey ? mirror.find((m) => m.app === row.app && m.rkey === row.rkey) : undefined; if (mirror.some((m) => m.app === row.app && m.id !== own?.id)) { removeWaiting(row.did, row.app, row.subject); return { skip: { kind: "skipped", reason: "already_following" } }; } if (!mirror.some((m) => m.id !== own?.id && set.apps.includes(m.app))) { if (own) deleteMirrorRow(own.id); return { skip: { kind: "skipped", reason: "no_longer_followed" } }; } return { own };}
async function executeFollow(row: QueueRow, set: SyncSet, now: Date): Promise<ExecuteOutcome> { const before = followIntent(row, set); if ("skip" in before) return before.skip;
const status = await profileStatus(row.app, row.subject); if (status === "unknown") { return { kind: "failed", failure: { kind: "retryable", message: "profile lookup failed" } }; } // The lookup awaited the network: the owner may have unfollowed (or // followed by hand) meanwhile. Check again, with no await from here until // the mirror row lands, so a later unfollow always sees this write. const after = followIntent(row, set); if ("skip" in after) return after.skip; const { own } = after; if (status === "not-found") { if (own) deleteMirrorRow(own.id); upsertWaiting( row.did, row.app, row.subject, now, new Date(now.getTime() + waitingRecheckMs(row.app)), ); return { kind: "waiting" }; }
// Pre-generate the rkey and land the mirror row before the write, so the // Jetstream echo is always recognized as our own and a retry after a // timed-out success reuses the key instead of creating a second record. const rkey = row.rkey ?? generateTid(); if (!row.rkey) { db.update(syncQueue).set({ rkey }).where(eq(syncQueue.id, row.id)).run(); } try { upsertMirror({ did: row.did, app: row.app, rkey, subject: row.subject, createdBySync: true, seenAt: now, }); } catch { // The partial unique index says sync already wrote this follow under // another key: nothing to do. return { kind: "skipped", reason: "already_following" }; }
try { await createArbitraryRecord( row.did, followCollection(row.app), { subject: row.subject, createdAt: now.toISOString() }, rkey, ); } catch (err) { if (isRecordAlreadyExists(err)) { removeWaiting(row.did, row.app, row.subject); return { kind: "written", points: 0 }; } const failure = classifyPdsError(err); // Keep the mirror row on anything that may be retried (the write may // have landed); drop it only when the write definitely failed. if (failure.kind === "terminal") takeMirrorByRkey(row.did, row.app, rkey); return { kind: "failed", failure }; } removeWaiting(row.did, row.app, row.subject); return { kind: "written", points: POINTS.create };}
async function executeUnfollow(row: QueueRow, now: Date): Promise<ExecuteOutcome> { void now; if (!row.rkey) return { kind: "skipped", reason: "record_missing" }; try { await deleteArbitraryRecord(row.did, followCollection(row.app), row.rkey); } catch (err) { if ( err instanceof PdsHttpError && err.status === 400 && /RecordNotFound|not found/i.test(err.body) ) { takeMirrorByRkey(row.did, row.app, row.rkey); return { kind: "skipped", reason: "record_missing" }; } return { kind: "failed", failure: classifyPdsError(err) }; } takeMirrorByRkey(row.did, row.app, row.rkey); return { kind: "deleted", points: POINTS.delete };}