From bac51006bc1f566a367cc70ec4a6aad2d9f064b8 Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Wed, 19 Aug 2026 16:44:29 +0200 Subject: [PATCH] more fixes --- .gitattributes | 2 + .../@atmo-dev__contrail-appview@0.12.2.patch | 121 ++++++++++++++++++ pnpm-lock.yaml | 8 +- 3 files changed, 127 insertions(+), 4 deletions(-) create mode 100644 .gitattributes diff --git a/.gitattributes b/.gitattributes new file mode 100644 index 0000000..aa47290 --- /dev/null +++ b/.gitattributes @@ -0,0 +1,2 @@ +# Unified diffs require a leading space on blank context lines. +patches/*.patch -whitespace diff --git a/patches/@atmo-dev__contrail-appview@0.12.2.patch b/patches/@atmo-dev__contrail-appview@0.12.2.patch index 4cf2376..30d2786 100644 --- a/patches/@atmo-dev__contrail-appview@0.12.2.patch +++ b/patches/@atmo-dev__contrail-appview@0.12.2.patch @@ -328,3 +328,124 @@ index c97302002b9f1b701366b70ee8bb30881623f2df..eaf0a84a8a988422489087843f975d03 ftsTableName, getFeedPruneCursor, getLabelerState, +diff --git a/dist/index.d.ts b/dist/index.d.ts +index 2ba9f66..272643e 100644 +--- a/dist/index.d.ts ++++ b/dist/index.d.ts +@@ -62,6 +62,11 @@ declare function buildRecordSyncSchema(dialect: SqlDialect): string[]; + /** SchemaModule-shaped helper for `initSchema({ extraSchemas: [...] })`. */ + declare function applyRecordSyncSchema(db: Database): Promise; + ++/** A hot backlog can deliver tens of thousands of events during the drain ++ * window, leaving no time to apply them and persist the cursor before the ++ * caller's hard timeout. Bound each cycle so apply + save always make forward ++ * progress; subsequent cron ticks resume from the persisted safe cursor. */ ++declare const MAX_INGEST_EVENTS_PER_CYCLE = 250; + /** Distinct actors pruned per ingest tick by the rolling feed sweep. Each + * actor costs a handful of index-backed O(cap) deletes, so this bounds the + * prune's per-tick CPU regardless of how large feed_items grows. */ +@@ -796,4 +801,4 @@ declare function resetLabelerCursor(db: Database, did: string): Promise; + * per-labeler. */ + declare function buildLabelsSchema(dialect: SqlDialect): string[]; + +-export { type BackfillAllOptions, type BackfillOptions, type BackfillProgress, CONTRAIL_SCHEMA_VERSION, type CollectionStats, type CreateAppOptions, type ExistingRecordInfo, FEED_PRUNE_RECOVERY_INTERVAL_MS, FEED_PRUNE_SWEEP_ACTORS, type FeedSweepResult, type FormattedRecord, HostedAdapter, type HydrateResult, type HydratedLabel, type IncomingLabel, type IngestState, type InitSchemaOptions, type LabelerResolveOverrides, type LabelerState, type NotifyResult, OPTIMIZE_LAST_MS_KEY, type PersistentIngestOptions, type PersistentLabelsOptions, type ProfileEntry, type QueryOptions, type RealtimeRoutesOptions, type RecordHostSyncOptions, type RecordHostSyncSource, type ReferenceResult, type RefreshOptions, type RefreshProgress, type RefreshResult, type SchemaModule, type SelectedLabelers, type SortOption, type SpacesContext, type SpacesRoutesOptions, type TopicResolution, type TopicResolutionContext, type TopicResolutionError, addColumnIfNotExists, applyCountColumns, applyEvents, applyLabels, applyRecordSyncSchema, backfillFollowersFromConstellation, backfillPending, backfillUser, batchedInQuery, buildCollectionTables, buildCountColumns, buildDynamicIndexes, buildFtsContent, buildFtsTables, buildLabelsSchema, buildRecordSyncSchema, collectDids, createApp, createIngestState, discoverDIDs, fieldToParam, formatRecord, ftsRowTableName, ftsTableName, getFeedPruneCursor, getLabelerState, getLastCursor, getMeta, getMetaNumber, getSearchableFields, hydrateLabels, ingestEvents, initSchema, lookupExistingRecords, maybeOptimize, optimizeDatabase, parseAtUri, parseHydrateParams, parseIntParam, processNotifyUris, pruneActorFeed, pruneFeedItems, queryAcrossSources, queryRecords, refresh, registerAdminRoutes, registerCollectionRoutes, registerFeedRoutes, registerNotifyRoute, registerRealtimeRoutes, registerSpacesRoutes, resetLabelerCursor, resolveHydrates, resolveLabelerEndpoint, resolveProfiles, resolveReferences, resolveTopicForCaller, runFeedPruneSlice, runGatedFeedPrune, runIngestCycle, runLabelIngestCycle, runPersistent, runPersistentLabels, runPipeline, runRecordHostSync, saveCursor, saveFeedPruneCursor, saveLabelerCursor, selectAcceptedLabelers, setMeta, sweepFeedItems, validateEndpointUrl, wrapWithPublishing }; ++export { type BackfillAllOptions, type BackfillOptions, type BackfillProgress, CONTRAIL_SCHEMA_VERSION, type CollectionStats, type CreateAppOptions, type ExistingRecordInfo, FEED_PRUNE_RECOVERY_INTERVAL_MS, FEED_PRUNE_SWEEP_ACTORS, type FeedSweepResult, type FormattedRecord, HostedAdapter, type HydrateResult, type HydratedLabel, type IncomingLabel, type IngestState, type InitSchemaOptions, type LabelerResolveOverrides, type LabelerState, MAX_INGEST_EVENTS_PER_CYCLE, type NotifyResult, OPTIMIZE_LAST_MS_KEY, type PersistentIngestOptions, type PersistentLabelsOptions, type ProfileEntry, type QueryOptions, type RealtimeRoutesOptions, type RecordHostSyncOptions, type RecordHostSyncSource, type ReferenceResult, type RefreshOptions, type RefreshProgress, type RefreshResult, type SchemaModule, type SelectedLabelers, type SortOption, type SpacesContext, type SpacesRoutesOptions, type TopicResolution, type TopicResolutionContext, type TopicResolutionError, addColumnIfNotExists, applyCountColumns, applyEvents, applyLabels, applyRecordSyncSchema, backfillFollowersFromConstellation, backfillPending, backfillUser, batchedInQuery, buildCollectionTables, buildCountColumns, buildDynamicIndexes, buildFtsContent, buildFtsTables, buildLabelsSchema, buildRecordSyncSchema, collectDids, createApp, createIngestState, discoverDIDs, fieldToParam, formatRecord, ftsRowTableName, ftsTableName, getFeedPruneCursor, getLabelerState, getLastCursor, getMeta, getMetaNumber, getSearchableFields, hydrateLabels, ingestEvents, initSchema, lookupExistingRecords, maybeOptimize, optimizeDatabase, parseAtUri, parseHydrateParams, parseIntParam, processNotifyUris, pruneActorFeed, pruneFeedItems, queryAcrossSources, queryRecords, refresh, registerAdminRoutes, registerCollectionRoutes, registerFeedRoutes, registerNotifyRoute, registerRealtimeRoutes, registerSpacesRoutes, resetLabelerCursor, resolveHydrates, resolveLabelerEndpoint, resolveProfiles, resolveReferences, resolveTopicForCaller, runFeedPruneSlice, runGatedFeedPrune, runIngestCycle, runLabelIngestCycle, runPersistent, runPersistentLabels, runPipeline, runRecordHostSync, saveCursor, saveFeedPruneCursor, saveLabelerCursor, selectAcceptedLabelers, setMeta, sweepFeedItems, validateEndpointUrl, wrapWithPublishing }; +diff --git a/dist/index.js b/dist/index.js +index eaf0a84..4e42d8f 100644 +--- a/dist/index.js ++++ b/dist/index.js +@@ -28,6 +28,7 @@ __export(index_exports, { + FEED_PRUNE_SWEEP_ACTORS: () => FEED_PRUNE_SWEEP_ACTORS, + HostedAdapter: () => HostedAdapter, + InMemoryPubSub: () => in_memory_exports.InMemoryPubSub, ++ MAX_INGEST_EVENTS_PER_CYCLE: () => MAX_INGEST_EVENTS_PER_CYCLE, + OPTIMIZE_LAST_MS_KEY: () => OPTIMIZE_LAST_MS_KEY, + RealtimePubSubDO: () => durable_object_exports.RealtimePubSubDO, + TicketSigner: () => ticket_exports.TicketSigner, +@@ -1720,6 +1721,7 @@ async function backfillFollowersFromConstellation(db, config, subjectDid) { + + // src/core/jetstream.ts + var BATCH_SIZE = 50; ++var MAX_INGEST_EVENTS_PER_CYCLE = 250; + var FEED_PRUNE_SWEEP_ACTORS = 500; + var FEED_PRUNE_RECOVERY_INTERVAL_MS = 6 * 60 * 60 * 1e3; + var FEED_PRUNE_LAST_FULL_PASS_META = "feed_prune_last_full_pass_ms"; +@@ -1800,8 +1802,8 @@ async function ingestEvents(config, cursor, safetyTimeoutMs = 25e3, knownDids) { + let lastYieldedTimeUs = null; + let firstYieldedTimeUs = null; + let connectCount = 0; +- const seenUris = /* @__PURE__ */ new Map(); +- const duplicateUris = []; ++ const seenSourceEvents = /* @__PURE__ */ new Set(); ++ let duplicateEvents = 0; + const newlyKnownDids = /* @__PURE__ */ new Set(); + const identityUpdates = /* @__PURE__ */ new Map(); + const subscription = new JetstreamSubscription({ +@@ -1858,15 +1860,12 @@ async function ingestEvents(config, cursor, safetyTimeoutMs = 25e3, knownDids) { + } + if (!keep) return; + } +- const prev = seenUris.get(uri); +- if (prev !== void 0) { +- duplicateUris.push(uri); +- log.warn( +- `[ingest] DUPLICATE in cycle: ${uri} first time_us=${prev}, again=${event.time_us}, delta=${event.time_us - prev}us` +- ); +- } else { +- seenUris.set(uri, event.time_us); ++ const sourceEventKey = `${uri}\0${event.time_us}`; ++ if (seenSourceEvents.has(sourceEventKey)) { ++ duplicateEvents++; ++ return; + } ++ seenSourceEvents.add(sourceEventKey); + const now = Date.now(); + collected.push({ + uri, +@@ -1879,9 +1878,6 @@ async function ingestEvents(config, cursor, safetyTimeoutMs = 25e3, knownDids) { + record: commit.operation === "delete" ? null : JSON.stringify(commit.record), + indexed_at: now * 1e3 + }); +- log.log( +- `[ingest] keep: ${commit.operation} ${uri} time_us=${event.time_us}` +- ); + if (knownDids && !dependentCollections.has(commit.collection)) { + if (!knownDids.has(event.did)) { + knownDids.add(event.did); +@@ -1911,6 +1907,12 @@ async function ingestEvents(config, cursor, safetyTimeoutMs = 25e3, knownDids) { + if (firstYieldedTimeUs === null) firstYieldedTimeUs = event.time_us; + lastYieldedTimeUs = event.time_us; + handleEvent(event); ++ if (collected.length >= MAX_INGEST_EVENTS_PER_CYCLE) { ++ log.log( ++ `[ingest] event cap reached, stopping drain (cap=${MAX_INGEST_EVENTS_PER_CYCLE}, commits_seen=${totalCommits}, duplicates_dropped=${duplicateEvents})` ++ ); ++ break; ++ } + if (event.time_us >= startTimeUs) { + log.log( + `[ingest] caught up to present, stopping (last time_us=${event.time_us}, startTimeUs=${startTimeUs})` +@@ -1926,11 +1928,12 @@ async function ingestEvents(config, cursor, safetyTimeoutMs = 25e3, knownDids) { + `[ingest] ${filteredUnknownDid} events filtered (unknown did). sample dids: ${sample}` + ); + } +- const lastCursor = subscription.cursor || null; +- const cursorGap = lastCursor !== null && lastYieldedTimeUs !== null ? lastCursor - lastYieldedTimeUs : null; ++ const subscriptionCursor = subscription.cursor || null; ++ const lastCursor = lastYieldedTimeUs === null ? subscriptionCursor : subscriptionCursor === null ? lastYieldedTimeUs : Math.min(subscriptionCursor, lastYieldedTimeUs); ++ const cursorGap = subscriptionCursor !== null && lastYieldedTimeUs !== null ? subscriptionCursor - lastYieldedTimeUs : null; + const rolledBackUs = cursor !== null && firstYieldedTimeUs !== null && firstYieldedTimeUs < cursor ? cursor - firstYieldedTimeUs : 0; + log.log( +- `[ingest] jetstream loop done. commits_seen=${totalCommits}, filtered=${filteredUnknownDid}, kept=${collected.length}, dupes=${duplicateUris.length}, connects=${connectCount}, first_yielded=${firstYieldedTimeUs ?? "none"}, last_yielded=${lastYieldedTimeUs ?? "none"}, subscription_cursor=${lastCursor ?? "none"}, cursor_gap=${cursorGap ?? "n/a"}us, rolled_back=${rolledBackUs}us` ++ `[ingest] jetstream loop done. commits_seen=${totalCommits}, filtered=${filteredUnknownDid}, kept=${collected.length}, dupes_dropped=${duplicateEvents}, connects=${connectCount}, first_yielded=${firstYieldedTimeUs ?? "none"}, last_yielded=${lastYieldedTimeUs ?? "none"}, subscription_cursor=${subscriptionCursor ?? "none"}, safe_cursor=${lastCursor ?? "none"}, cursor_gap=${cursorGap ?? "n/a"}us, rolled_back=${rolledBackUs}us` + ); + if (cursorGap !== null && cursorGap > 1e3) { + log.warn( +@@ -5738,6 +5741,7 @@ export { + FEED_PRUNE_SWEEP_ACTORS, + HostedAdapter, + export_InMemoryPubSub as InMemoryPubSub, ++ MAX_INGEST_EVENTS_PER_CYCLE, + OPTIMIZE_LAST_MS_KEY, + export_RealtimePubSubDO as RealtimePubSubDO, + export_TicketSigner as TicketSigner, diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index cd912ab..da3acef 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -5,7 +5,7 @@ settings: excludeLinksFromLockfile: false patchedDependencies: - '@atmo-dev/contrail-appview@0.12.2': f4aa36f86a9f705620dcbd57e82d87abed38c36e5ff2c889e66ec9a0970f767f + '@atmo-dev/contrail-appview@0.12.2': aa343d8a2faef8d66a275aec3ddcb5d88bb4bb484e197032bfdf8604bacb96d8 '@atmo-dev/contrail-base@0.12.2': f13888cab253afcb9fe4fa92ace590d94e28ba71a58fea3919d6c2284daffa22 importers: @@ -141,7 +141,7 @@ importers: version: 0.1.12 '@atmo-dev/contrail-appview': specifier: 0.12.2 - version: 0.12.2(patch_hash=f4aa36f86a9f705620dcbd57e82d87abed38c36e5ff2c889e66ec9a0970f767f) + version: 0.12.2(patch_hash=aa343d8a2faef8d66a275aec3ddcb5d88bb4bb484e197032bfdf8604bacb96d8) '@atmo-dev/contrail-lexicons': specifier: ^0.4.4 version: 0.4.4(wrangler@4.77.0(@cloudflare/workers-types@4.20260317.1)) @@ -3856,7 +3856,7 @@ snapshots: '@badrap/valita': 0.4.6 nanoid: 5.1.7 - '@atmo-dev/contrail-appview@0.12.2(patch_hash=f4aa36f86a9f705620dcbd57e82d87abed38c36e5ff2c889e66ec9a0970f767f)': + '@atmo-dev/contrail-appview@0.12.2(patch_hash=aa343d8a2faef8d66a275aec3ddcb5d88bb4bb484e197032bfdf8604bacb96d8)': dependencies: '@atcute/atproto': 3.1.10 '@atcute/cbor': 2.3.2 @@ -3923,7 +3923,7 @@ snapshots: '@atcute/jetstream': 1.1.2 '@atcute/lexicons': 1.2.9 '@atcute/xrpc-server': 0.1.12 - '@atmo-dev/contrail-appview': 0.12.2(patch_hash=f4aa36f86a9f705620dcbd57e82d87abed38c36e5ff2c889e66ec9a0970f767f) + '@atmo-dev/contrail-appview': 0.12.2(patch_hash=aa343d8a2faef8d66a275aec3ddcb5d88bb4bb484e197032bfdf8604bacb96d8) '@atmo-dev/contrail-authority': 0.12.2 '@atmo-dev/contrail-base': 0.12.2(patch_hash=f13888cab253afcb9fe4fa92ace590d94e28ba71a58fea3919d6c2284daffa22) '@atmo-dev/contrail-record-host': 0.12.2 -- 2.51.2