|
| 1 | +# frozen_string_literal: true |
| 2 | + |
| 3 | +require 'digest' |
| 4 | + |
| 5 | +module UnitHub |
| 6 | + # Notifications about Unit Hub announcements and learning sessions. |
| 7 | + # |
| 8 | + # The models call the *_committed methods from after_commit, which decide |
| 9 | + # whether a save is worth telling anyone about and queue a job with only ids |
| 10 | + # and a version. The jobs call fan_out, which walks the unit's recipients in |
| 11 | + # batches and raises each notification through NotificationService, so the |
| 12 | + # category preference, email and push all behave like every other event. |
| 13 | + # |
| 14 | + # Every notification carries a dedupe key made of the event, the record and a |
| 15 | + # version of it, which the unique index scopes to the recipient. A retried or |
| 16 | + # duplicated job finds the row that is already there instead of sending again. |
| 17 | + module Notifications |
| 18 | + TYPE = 'unit_hub' |
| 19 | + |
| 20 | + ANNOUNCEMENT_PUBLISHED = 'unit_announcement_published' |
| 21 | + ANNOUNCEMENT_UPDATED = 'unit_announcement_updated' |
| 22 | + SESSION_CHANGED = 'unit_session_changed' |
| 23 | + SESSION_STARTING_SOON = 'unit_session_starting_soon' |
| 24 | + EVENTS = [ANNOUNCEMENT_PUBLISHED, ANNOUNCEMENT_UPDATED, SESSION_CHANGED, SESSION_STARTING_SOON].freeze |
| 25 | + |
| 26 | + BATCH_SIZE = 200 |
| 27 | + |
| 28 | + # At most one update notification per announcement per recipient in this |
| 29 | + # window. A publish inside the window counts, so fixing a line straight |
| 30 | + # after posting does not send a second notification. |
| 31 | + UPDATE_DEBOUNCE = 30.minutes |
| 32 | + |
| 33 | + # How long before a session its reminder goes out. |
| 34 | + REMINDER_LEAD = 30.minutes |
| 35 | + |
| 36 | + # An announcement created with a publication time older than this is an |
| 37 | + # import or a backfill, not news, so it does not raise a publish. |
| 38 | + BACKFILL_AGE = 1.day |
| 39 | + |
| 40 | + SESSION_SCHEDULE_FIELDS = %w[start_at end_at timezone recurrence recurrence_until location join_url].freeze |
| 41 | + |
| 42 | + # A one letter fix in a word of four or more letters is a typo. Anything |
| 43 | + # else that changes the words, or any change to a number, is not. |
| 44 | + TYPO_DISTANCE = 2 |
| 45 | + |
| 46 | + module_function |
| 47 | + |
| 48 | + # ---- deciding what a save means ------------------------------------------ |
| 49 | + |
| 50 | + def announcement_committed(record, created:) |
| 51 | + now = Time.current |
| 52 | + return unless announcement_source_allowed?(record) |
| 53 | + |
| 54 | + if created |
| 55 | + return if record.published_at.nil? |
| 56 | + return if record.published_at < record.created_at - BACKFILL_AGE |
| 57 | + |
| 58 | + return queue_publish(record, now) |
| 59 | + end |
| 60 | + |
| 61 | + changes = record.saved_changes |
| 62 | + was_visible = visible_with?( |
| 63 | + changes.key?('published_at') ? changes['published_at'].first : record.published_at, |
| 64 | + changes.key?('expires_at') ? changes['expires_at'].first : record.expires_at, |
| 65 | + now |
| 66 | + ) |
| 67 | + |
| 68 | + unless was_visible |
| 69 | + return unless changes.keys.intersect?(%w[published_at expires_at]) |
| 70 | + return if record.published_at.nil? || (record.expires_at && record.expires_at <= now) |
| 71 | + |
| 72 | + return queue_publish(record, now) |
| 73 | + end |
| 74 | + return unless visible_with?(record.published_at, record.expires_at, now) |
| 75 | + |
| 76 | + # Pinning a visible announcement puts it back in front of people, so it is |
| 77 | + # raised as a publish with its own version. Unpinning is not news. |
| 78 | + if changes.key?('pinned') && record.pinned |
| 79 | + UnitAnnouncementNotificationJob.perform_async(record.id, ANNOUNCEMENT_PUBLISHED, "pinned-#{record.updated_at.to_i}") |
| 80 | + return |
| 81 | + end |
| 82 | + |
| 83 | + old_title = changes.key?('title') ? changes['title'].first : record.title |
| 84 | + old_body = changes.key?('body') ? changes['body'].first : record.body |
| 85 | + return unless meaningful_edit?(old_title, record.title) || meaningful_edit?(old_body, record.body) |
| 86 | + |
| 87 | + UnitAnnouncementNotificationJob.perform_async(record.id, ANNOUNCEMENT_UPDATED, nil) |
| 88 | + end |
| 89 | + |
| 90 | + def session_committed(record) |
| 91 | + changes = record.saved_changes |
| 92 | + published_before = changes.key?('published') ? changes['published'].first : record.published |
| 93 | + return unless record.published && published_before |
| 94 | + |
| 95 | + cancelled_now = changes.key?('cancelled') && record.cancelled |
| 96 | + restored = changes.key?('cancelled') && !record.cancelled |
| 97 | + schedule_changed = changes.keys.intersect?(SESSION_SCHEDULE_FIELDS) |
| 98 | + return unless cancelled_now || (!record.cancelled && (schedule_changed || restored)) |
| 99 | + return if next_occurrence(record, Time.current).nil? |
| 100 | + |
| 101 | + UnitSessionNotificationJob.perform_async(record.id, record.lock_version, cancelled_now ? 'cancelled' : 'changed') |
| 102 | + end |
| 103 | + |
| 104 | + def queue_publish(record, now) |
| 105 | + version = "published-#{record.published_at.to_i}" |
| 106 | + if record.published_at > now |
| 107 | + UnitAnnouncementNotificationJob.perform_at(record.published_at, record.id, ANNOUNCEMENT_PUBLISHED, version) |
| 108 | + else |
| 109 | + UnitAnnouncementNotificationJob.perform_async(record.id, ANNOUNCEMENT_PUBLISHED, version) |
| 110 | + end |
| 111 | + end |
| 112 | + |
| 113 | + def visible_with?(published_at, expires_at, at) |
| 114 | + published_at.present? && published_at <= at && (expires_at.nil? || expires_at > at) |
| 115 | + end |
| 116 | + |
| 117 | + def announcement_source_allowed?(record) |
| 118 | + UnitAnnouncement.allowed_sources.exists?(id: record.id) |
| 119 | + end |
| 120 | + |
| 121 | + # Whether an edit changes what an announcement says, rather than fixing how |
| 122 | + # it is spelled. Case, spacing and punctuation never count. A number always |
| 123 | + # counts, because a changed date, time or room is the edit people need to |
| 124 | + # hear about. One word swapped for a near spelling, or a doubled word taken |
| 125 | + # out, is a typo. Any other change to the words counts. |
| 126 | + def meaningful_edit?(before, after) |
| 127 | + old_words = words(before) |
| 128 | + new_words = words(after) |
| 129 | + removed = multiset_difference(old_words, new_words) |
| 130 | + added = multiset_difference(new_words, old_words) |
| 131 | + changed = removed + added |
| 132 | + |
| 133 | + return false if changed.empty? |
| 134 | + return true if changed.any? { |word| word.match?(/\d/) } |
| 135 | + |
| 136 | + if removed.length == 1 && added.length == 1 |
| 137 | + return true if [removed.first.length, added.first.length].min < 4 |
| 138 | + |
| 139 | + return DidYouMean::Levenshtein.distance(removed.first, added.first) > TYPO_DISTANCE |
| 140 | + end |
| 141 | + |
| 142 | + return !(old_words.include?(changed.first) && new_words.include?(changed.first)) if changed.length == 1 |
| 143 | + |
| 144 | + true |
| 145 | + end |
| 146 | + |
| 147 | + def words(text) |
| 148 | + text.to_s.downcase.scan(/[[:alnum:]]+/) |
| 149 | + end |
| 150 | + |
| 151 | + def multiset_difference(left, right) |
| 152 | + remaining = right.tally |
| 153 | + left.each_with_object([]) do |word, result| |
| 154 | + if remaining[word].to_i.positive? |
| 155 | + remaining[word] -= 1 |
| 156 | + else |
| 157 | + result << word |
| 158 | + end |
| 159 | + end |
| 160 | + end |
| 161 | + |
| 162 | + # ---- who hears about it -------------------------------------------------- |
| 163 | + |
| 164 | + # Enrolled students and teaching staff of an active unit, with the Unit Hub |
| 165 | + # category on, never the person who wrote the thing. |
| 166 | + def recipients(unit, author_id: nil) |
| 167 | + return User.none unless unit&.active |
| 168 | + |
| 169 | + student_ids = Project.where(unit_id: unit.id, enrolled: true).select(:user_id) |
| 170 | + staff_ids = UnitRole.where(unit_id: unit.id, role_id: [Role.tutor.id, Role.convenor.id]).select(:user_id) |
| 171 | + scope = User.where(id: student_ids).or(User.where(id: staff_ids)).where(receive_unit_hub_notifications: true) |
| 172 | + author_id ? scope.where.not(id: author_id) : scope |
| 173 | + end |
| 174 | + |
| 175 | + # Raise one notification per recipient, a batch at a time. |
| 176 | + # |
| 177 | + # skip - called with a batch of user ids, returns the ids to leave out on |
| 178 | + # top of the ones that already hold this dedupe key. |
| 179 | + # Failures are collected and raised at the end so Sidekiq retries the whole |
| 180 | + # fan-out, which the dedupe key makes safe. |
| 181 | + def fan_out(scope:, event:, dedupe_key:, notifiable:, message:, link:, skip: nil) |
| 182 | + failed = [] |
| 183 | + scope.find_in_batches(batch_size: BATCH_SIZE) do |batch| |
| 184 | + ids = batch.map(&:id) |
| 185 | + skipped = Notification.where(user_id: ids, dedupe_key: dedupe_key).pluck(:user_id) |
| 186 | + skipped.concat(skip.call(ids)) if skip |
| 187 | + skipped = skipped.to_set |
| 188 | + |
| 189 | + batch.each do |user| |
| 190 | + next if skipped.include?(user.id) |
| 191 | + |
| 192 | + NotificationService.notify( |
| 193 | + user: user, type: TYPE, event: event, message: message, |
| 194 | + link: link, dedupe_key: dedupe_key, notifiable: notifiable |
| 195 | + ) |
| 196 | + rescue StandardError => e |
| 197 | + failed << user.id |
| 198 | + Rails.logger.error("Failed #{event} notification for User #{user.id}: #{e.class}") |
| 199 | + end |
| 200 | + end |
| 201 | + |
| 202 | + raise "#{event} notifications failed for users: #{failed.join(', ')}" if failed.any? |
| 203 | + end |
| 204 | + |
| 205 | + # ---- what it says -------------------------------------------------------- |
| 206 | + |
| 207 | + def announcement_link(record) |
| 208 | + "/unit-hub?unit=#{record.unit_id}&announcement=#{record.id}" |
| 209 | + end |
| 210 | + |
| 211 | + def session_link(record) |
| 212 | + "/unit-hub?unit=#{record.unit_id}&session=#{record.id}" |
| 213 | + end |
| 214 | + |
| 215 | + def announcement_message(record, event) |
| 216 | + prefix = |
| 217 | + if event == ANNOUNCEMENT_UPDATED |
| 218 | + 'Announcement updated' |
| 219 | + elsif record.pinned |
| 220 | + 'Pinned announcement' |
| 221 | + else |
| 222 | + 'New announcement' |
| 223 | + end |
| 224 | + "#{prefix} in #{record.unit.code}: #{record.title}".truncate(500) |
| 225 | + end |
| 226 | + |
| 227 | + def session_changed_message(record, change, at: Time.current) |
| 228 | + unit_code = record.unit.code |
| 229 | + if change == 'cancelled' |
| 230 | + return "The weekly #{record.title} sessions in #{unit_code} are cancelled.".truncate(500) if record.recurrence == 'weekly' |
| 231 | + |
| 232 | + occurrence = next_occurrence(record, at) |
| 233 | + when_text = occurrence ? " on #{format_time(occurrence[:start_at], record.timezone)}" : '' |
| 234 | + return "#{record.title} in #{unit_code}#{when_text} is cancelled.".truncate(500) |
| 235 | + end |
| 236 | + |
| 237 | + "#{record.title} in #{unit_code} has changed. #{where_and_when(record, at)}".truncate(500) |
| 238 | + end |
| 239 | + |
| 240 | + def session_starting_soon_message(record, start_at) |
| 241 | + place = session_place(record) |
| 242 | + "#{record.title} in #{record.unit.code} starts at #{format_clock(start_at, record.timezone)}#{" (#{place})" if place}.".truncate(500) |
| 243 | + end |
| 244 | + |
| 245 | + def where_and_when(record, at) |
| 246 | + occurrence = next_occurrence(record, at) |
| 247 | + return 'Open the Unit Hub for the details.' if occurrence.nil? |
| 248 | + |
| 249 | + cadence = record.recurrence == 'weekly' ? 'Next session' : 'Now' |
| 250 | + place = session_place(record) |
| 251 | + "#{cadence}: #{format_time(occurrence[:start_at], record.timezone)}#{", #{place}" if place}." |
| 252 | + end |
| 253 | + |
| 254 | + def session_place(record) |
| 255 | + return record.location if record.location.present? |
| 256 | + |
| 257 | + 'online' if record.join_url.present? |
| 258 | + end |
| 259 | + |
| 260 | + def next_occurrence(record, at) |
| 261 | + record.occurrences(from: at, to: at + 7.months).find { |occurrence| occurrence[:end_at] >= at } |
| 262 | + end |
| 263 | + |
| 264 | + def format_time(time, zone) |
| 265 | + local = time.in_time_zone(zone) |
| 266 | + "#{local.strftime('%a %-d %b')} at #{format_clock(local, zone)}" |
| 267 | + end |
| 268 | + |
| 269 | + def format_clock(time, zone) |
| 270 | + local = time.in_time_zone(zone) |
| 271 | + "#{local.strftime('%-l:%M%P')} #{local.zone}" |
| 272 | + end |
| 273 | + |
| 274 | + # The details an email shows for a session, in the session's own zone. |
| 275 | + def session_details(record, at:) |
| 276 | + occurrence = next_occurrence(record, at) |
| 277 | + details = { 'Unit' => "#{record.unit.code} #{record.unit.name}", 'Session' => record.title } |
| 278 | + if occurrence |
| 279 | + local_start = occurrence[:start_at].in_time_zone(record.timezone) |
| 280 | + local_end = occurrence[:end_at].in_time_zone(record.timezone) |
| 281 | + details['When'] = "#{format_time(local_start, record.timezone)} to #{local_end.strftime('%-l:%M%P')}" |
| 282 | + details['Repeats'] = "Weekly until #{record.recurrence_until.strftime('%-d %b %Y')}" if record.recurrence == 'weekly' && record.recurrence_until |
| 283 | + end |
| 284 | + details['Where'] = record.location if record.location.present? |
| 285 | + details['Online'] = 'A join link is on the Unit Hub' if record.join_url.present? && !record.cancelled |
| 286 | + details['Status'] = 'Cancelled' if record.cancelled |
| 287 | + details |
| 288 | + end |
| 289 | + |
| 290 | + def announcement_version(record) |
| 291 | + Digest::SHA256.hexdigest([record.title, record.body].join(" |