(eventRepository: IEventRepository)
| 89 | |
| 90 | export const createEventBatchPersister = |
| 91 | (eventRepository: IEventRepository) => |
| 92 | async (events: Event[]): Promise<number> => { |
| 93 | if (!events.length) { |
| 94 | return 0 |
| 95 | } |
| 96 | |
| 97 | let inserted = 0 |
| 98 | |
| 99 | const regularEvents: Event[] = [] |
| 100 | const replaceableEvents: Event[] = [] |
| 101 | |
| 102 | for (const event of events) { |
| 103 | if (isEphemeralEvent(event)) { |
| 104 | continue |
| 105 | } |
| 106 | |
| 107 | if (isDeleteEvent(event)) { |
| 108 | // flush pending batches before applying deletes |
| 109 | inserted += await eventRepository.createMany(regularEvents.splice(0)) |
| 110 | inserted += await eventRepository.upsertMany(replaceableEvents.splice(0)) |
| 111 | |
| 112 | const eventIdsToDelete = event.tags.reduce( |
| 113 | (ids, tag) => |
| 114 | tag.length >= 2 && tag[0] === EventTags.Event && /^[0-9a-f]{64}$/.test(tag[1]) ? [...ids, tag[1]] : ids, |
| 115 | [] as string[], |
| 116 | ) |
| 117 | |
| 118 | if (eventIdsToDelete.length) { |
| 119 | await eventRepository.deleteByPubkeyAndIds(event.pubkey, eventIdsToDelete) |
| 120 | } |
| 121 | |
| 122 | inserted += await eventRepository.create(enrichEventMetadata(event)) |
| 123 | continue |
| 124 | } |
| 125 | |
| 126 | const enrichedEvent = enrichEventMetadata(event) |
| 127 | |
| 128 | if (isReplaceableEvent(event) || isParameterizedReplaceableEvent(event)) { |
| 129 | replaceableEvents.push(enrichedEvent) |
| 130 | continue |
| 131 | } |
| 132 | |
| 133 | regularEvents.push(enrichedEvent) |
| 134 | } |
| 135 | |
| 136 | // flush remaining |
| 137 | inserted += await eventRepository.createMany(regularEvents) |
| 138 | inserted += await eventRepository.upsertMany(replaceableEvents) |
| 139 | |
| 140 | return inserted |
| 141 | } |
| 142 | |
| 143 | export class EventImportService { |
| 144 | public constructor( |
no test coverage detected