diff --git a/app/at_uri.rb b/app/at_uri.rb index 1613476..531b2bc 100644 --- a/app/at_uri.rb +++ b/app/at_uri.rb @@ -7,6 +7,7 @@ class AT_URI attr_reader :repo, :collection, :rkey def initialize(uri) + uri = uri.to_s raise InvalidURIError, "Invalid AT URI: #{uri}" if uri.include?(' ') || !uri.start_with?('at://') parts = uri.split('/') @@ -38,5 +39,6 @@ class AT_URI end def AT_URI(uri) + return uri if uri.is_a?(AT_URI) AT_URI.new(uri) end diff --git a/app/importers/base_importer.rb b/app/importers/base_importer.rb index 293d1c6..126b780 100644 --- a/app/importers/base_importer.rb +++ b/app/importers/base_importer.rb @@ -43,22 +43,4 @@ class BaseImporter def import_items raise NotImplementedError end - - def create_item_for_post(uri) - post_uri = AT_URI(uri) - return unless post_uri.is_post? - - post = Post.find_by_at_uri(post_uri) - - if post - yield({ post: post }) - else - item_stub = yield({ post_uri: post_uri }) - - if @item_queue - @item_queue.push(item_stub) - @report&.update(queue: { length: @item_queue.length }) - end - end - end end diff --git a/app/importers/likes_importer.rb b/app/importers/likes_importer.rb index a1e7955..a90acee 100644 --- a/app/importers/likes_importer.rb +++ b/app/importers/likes_importer.rb @@ -1,6 +1,4 @@ require 'time' - -require_relative '../at_uri' require_relative 'base_importer' class LikesImporter < BaseImporter @@ -20,14 +18,11 @@ class LikesImporter < BaseImporter records.each do |record| begin - like_rkey = AT_URI(record['uri']).rkey - next if @user.likes.where(rkey: like_rkey).exists? - - like_time = Time.parse(record['value']['createdAt']) - post_uri = record['value']['subject']['uri'] + like = @user.likes.import_from_record(record['uri'], record['value']) - create_item_for_post(post_uri) do |args| - @user.likes.create!(args.merge(rkey: like_rkey, time: like_time)) + if like && like.pending? && @item_queue + @item_queue.push(like) + @report&.update(queue: { length: @item_queue.length }) end rescue StandardError => e puts "Error in LikesImporter: #{record['uri']}: #{e}" diff --git a/app/importers/posts_importer.rb b/app/importers/posts_importer.rb index 2915b4b..38ec698 100644 --- a/app/importers/posts_importer.rb +++ b/app/importers/posts_importer.rb @@ -1,6 +1,4 @@ require 'time' - -require_relative '../at_uri' require_relative 'base_importer' class PostsImporter < BaseImporter @@ -20,12 +18,19 @@ class PostsImporter < BaseImporter records.each do |record| begin - if record['value']['embed'] && record['value']['embed']['record'] - save_quote(record) - end + quote = @user.quotes.import_from_record(record['uri'], record['value']) + pin = @user.pins.import_from_record(record['uri'], record['value']) + + if @item_queue + if quote && quote.pending? + @item_queue.push(quote) + end + + if pin && pin.pending? + @item_queue.push(pin) + end - if record['value']['reply'] && record['value']['text'].include?('📌') - save_pin(record) + @report&.update(queue: { length: @item_queue.length }) end rescue StandardError => e puts "Error in LikesImporter: #{record['uri']}: #{e}" @@ -39,36 +44,4 @@ class PostsImporter < BaseImporter break if @time_limit && records.any? { |x| Time.parse(x['value']['createdAt']) < @time_limit } end end - - def save_quote(record) - post_rkey = AT_URI(record['uri']).rkey - return if @user.quotes.where(rkey: post_rkey).exists? - - post_time = Time.parse(record['value']['createdAt']) - - quoted_post_uri = case record['value']['embed']['$type'] - when 'app.bsky.embed.record' - record['value']['embed']['record']['uri'] - when 'app.bsky.embed.recordWithMedia' - record['value']['embed']['record']['record']['uri'] - else - return - end - - create_item_for_post(quoted_post_uri) do |args| - @user.quotes.create!(args.merge(rkey: post_rkey, time: post_time, quote_text: record['value']['text'])) - end - end - - def save_pin(record) - post_rkey = AT_URI(record['uri']).rkey - return if @user.pins.where(rkey: post_rkey).exists? - - post_time = Time.parse(record['value']['createdAt']) - parent_post_uri = record['value']['reply']['parent']['uri'] - - create_item_for_post(parent_post_uri) do |args| - @user.pins.create!(args.merge(rkey: post_rkey, time: post_time, pin_text: record['value']['text'])) - end - end end diff --git a/app/importers/reposts_importer.rb b/app/importers/reposts_importer.rb index bfd3e03..e173205 100644 --- a/app/importers/reposts_importer.rb +++ b/app/importers/reposts_importer.rb @@ -1,6 +1,4 @@ require 'time' - -require_relative '../at_uri' require_relative 'base_importer' class RepostsImporter < BaseImporter @@ -20,14 +18,11 @@ class RepostsImporter < BaseImporter records.each do |record| begin - repost_rkey = AT_URI(record['uri']).rkey - next if @user.reposts.where(rkey: repost_rkey).exists? - - repost_time = Time.parse(record['value']['createdAt']) - post_uri = record['value']['subject']['uri'] + repost = @user.reposts.import_from_record(record['uri'], record['value']) - create_item_for_post(post_uri) do |args| - @user.reposts.create!(args.merge(rkey: repost_rkey, time: repost_time)) + if repost && repost.pending? && @item_queue + @item_queue.push(repost) + @report&.update(queue: { length: @item_queue.length }) end rescue StandardError => e puts "Error in RepostsImporter: #{record['uri']}: #{e}" diff --git a/app/models/importable.rb b/app/models/importable.rb new file mode 100644 index 0000000..193c97e --- /dev/null +++ b/app/models/importable.rb @@ -0,0 +1,29 @@ +require 'active_support/concern' + +require_relative '../at_uri' +require_relative 'post' + +module Importable + extend ActiveSupport::Concern + + included do + scope :pending, -> { where(post: nil) } + + def pending? + post_uri != nil + end + + def import_item! + post_uri = AT_URI(self.post_uri) + return nil if !post_uri.is_post? + + if post = Post.find_by_at_uri(post_uri) + self.post = post + self.post_uri = nil + end + + self.save! + self + end + end +end diff --git a/app/models/like.rb b/app/models/like.rb index 492f3de..bd4fb85 100644 --- a/app/models/like.rb +++ b/app/models/like.rb @@ -1,11 +1,15 @@ require 'active_record' +require 'time' +require_relative '../at_uri' require_relative 'post' +require_relative 'importable' require_relative 'searchable' require_relative 'user' class Like < ActiveRecord::Base include Searchable + include Importable validates_presence_of :time, :rkey validates_length_of :rkey, is: 13 @@ -14,4 +18,12 @@ class Like < ActiveRecord::Base belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true + + def self.new_from_record(uri, record) + self.new( + rkey: AT_URI(uri).rkey, + time: Time.parse(record['createdAt']), + post_uri: record['subject']['uri'] + ) + end end diff --git a/app/models/pin.rb b/app/models/pin.rb index bb0d05a..513191d 100644 --- a/app/models/pin.rb +++ b/app/models/pin.rb @@ -1,11 +1,16 @@ require 'active_record' +require 'time' +require_relative '../at_uri' require_relative 'post' require_relative 'searchable' require_relative 'user' class Pin < ActiveRecord::Base include Searchable + include Importable + + PIN_SIGN = '📌' validates_presence_of :time, :rkey validates_length_of :rkey, is: 13 @@ -15,4 +20,15 @@ class Pin < ActiveRecord::Base belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true + + def self.new_from_record(uri, record) + return nil unless record['reply'] && record['text'].include?(PIN_SIGN) + + self.new( + rkey: AT_URI(uri).rkey, + time: Time.parse(record['createdAt']), + post_uri: record['reply']['parent']['uri'], + pin_text: record['text'] + ) + end end diff --git a/app/models/quote.rb b/app/models/quote.rb index ea8fb3b..adb75a5 100644 --- a/app/models/quote.rb +++ b/app/models/quote.rb @@ -1,11 +1,14 @@ require 'active_record' +require 'time' +require_relative '../at_uri' require_relative 'post' require_relative 'searchable' require_relative 'user' class Quote < ActiveRecord::Base include Searchable + include Importable validates_presence_of :time, :rkey validates_length_of :rkey, is: 13 @@ -15,4 +18,24 @@ class Quote < ActiveRecord::Base belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true + + def self.new_from_record(uri, record) + return nil unless record['embed'] + + quoted_post_uri = case record['embed']['$type'] + when 'app.bsky.embed.record' + record['embed']['record']['uri'] + when 'app.bsky.embed.recordWithMedia' + record['embed']['record']['record']['uri'] + else + return nil + end + + self.new( + rkey: AT_URI(uri).rkey, + time: Time.parse(record['createdAt']), + post_uri: quoted_post_uri, + quote_text: record['text'] + ) + end end diff --git a/app/models/repost.rb b/app/models/repost.rb index 189eed4..f2165c9 100644 --- a/app/models/repost.rb +++ b/app/models/repost.rb @@ -1,11 +1,14 @@ require 'active_record' +require 'time' +require_relative '../at_uri' require_relative 'post' require_relative 'searchable' require_relative 'user' class Repost < ActiveRecord::Base include Searchable + include Importable validates_presence_of :time, :rkey validates_length_of :rkey, is: 13 @@ -14,4 +17,12 @@ class Repost < ActiveRecord::Base belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true + + def self.new_from_record(uri, record) + self.new( + rkey: AT_URI(uri).rkey, + time: Time.parse(record['createdAt']), + post_uri: record['subject']['uri'] + ) + end end diff --git a/app/models/searchable.rb b/app/models/searchable.rb index e0c67ea..ec1a802 100644 --- a/app/models/searchable.rb +++ b/app/models/searchable.rb @@ -41,8 +41,6 @@ module Searchable ) } - scope :pending, -> { where(post: nil) } - def cursor "#{self.time.to_f}:#{self.id}" end diff --git a/app/models/user.rb b/app/models/user.rb index 0ffac51..eade84e 100644 --- a/app/models/user.rb +++ b/app/models/user.rb @@ -12,12 +12,44 @@ class User < ActiveRecord::Base validates_length_of :did, maximum: 260 has_many :posts - has_many :likes, foreign_key: 'actor_id' - has_many :reposts, foreign_key: 'actor_id' - has_many :quotes, foreign_key: 'actor_id' - has_many :pins, foreign_key: 'actor_id' has_many :imports + has_many :likes, foreign_key: 'actor_id' do + def import_from_record(like_uri, record) + like = self.new_from_record(like_uri, record) + return nil if like.nil? || self.where(rkey: like.rkey).exists? + + like.import_item! + end + end + + has_many :reposts, foreign_key: 'actor_id' do + def import_from_record(repost_uri, record) + repost = self.new_from_record(repost_uri, record) + return nil if repost.nil? || self.where(rkey: repost.rkey).exists? + + repost.import_item! + end + end + + has_many :quotes, foreign_key: 'actor_id' do + def import_from_record(post_uri, record) + quote = self.new_from_record(post_uri, record) + return nil if quote.nil? || self.where(rkey: quote.rkey).exists? + + quote.import_item! + end + end + + has_many :pins, foreign_key: 'actor_id' do + def import_from_record(post_uri, record) + pin = self.new_from_record(post_uri, record) + return nil if pin.nil? || self.where(rkey: pin.rkey).exists? + + pin.import_item! + end + end + def all_pending_items [:likes, :reposts, :quotes, :pins].map { |x| self.send(x).pending.to_a }.reduce(&:+) end