diff --git a/.changeset/persist-cursor-before-identity-refresh.md b/.changeset/persist-cursor-before-identity-refresh.md new file mode 100644 index 0000000..6348578 --- /dev/null +++ b/.changeset/persist-cursor-before-identity-refresh.md @@ -0,0 +1,7 @@ +--- +"@atmo-dev/contrail-appview": patch +--- + +Persist the jetstream ingest cursor before the identity-refresh tail in `runIngestCycle`. + +`saveCursor` previously ran after `refreshStaleIdentities`, whose per-DID network calls can run long. If the ingest isolate was aborted (e.g. a scheduled-invocation deadline) before the save, the cursor never advanced and the next cycle re-drained the same jetstream window indefinitely. Records are durably applied before this point, so the cursor is now saved first; identity refresh is idempotent and staleness-driven, so deferring it past the save is safe. diff --git a/packages/contrail-appview/src/core/jetstream.ts b/packages/contrail-appview/src/core/jetstream.ts index f96c78d..ea0b479 100644 --- a/packages/contrail-appview/src/core/jetstream.ts +++ b/packages/contrail-appview/src/core/jetstream.ts @@ -385,16 +385,14 @@ export async function runIngestCycle( log.log(`[ingest] applied ${identityUpdates.size} identity event(s)`); } - // Refresh stale/missing identities for DIDs in this batch - const uniqueDids = [...new Set(events.map((e) => e.did))]; - if (uniqueDids.length > 0) { - try { - await refreshStaleIdentities(db, uniqueDids, config); - } catch (err) { - log.warn(`Identity refresh failed: ${err}`); - } - } - + // Persist the cursor BEFORE the best-effort enrichment tail below. Records are + // already durably applied (applyEvents) and handle changes recorded, so the + // cursor's forward progress is real and must be committed now. The steps that + // follow — refreshStaleIdentities especially — make per-DID network calls and + // can run long; if the cron isolate is aborted (e.g. a scheduled-invocation + // deadline) while they run, an un-saved cursor makes the next cycle re-drain the + // identical window forever. Identity refresh is idempotent and staleness-driven, + // so deferring it past the save costs nothing. if (lastCursor !== null) { await saveCursor(db, lastCursor); log.log( @@ -406,6 +404,17 @@ export async function runIngestCycle( log.log(`[ingest] no cursor returned from subscription; not saving`); } + // Refresh stale/missing identities for DIDs in this batch (best-effort; runs + // after the cursor save so its network latency can't strand forward progress). + const uniqueDids = [...new Set(events.map((e) => e.did))]; + if (uniqueDids.length > 0) { + try { + await refreshStaleIdentities(db, uniqueDids, config); + } catch (err) { + log.warn(`Identity refresh failed: ${err}`); + } + } + // Newly-discovered DIDs: ask Constellation for back-edges so they // immediately appear in existing followers' feeds (best-effort, opt-out). if (config.feeds && newlyKnownDids.length > 0) {