diff --git a/app/firehose_client.rb b/app/firehose_client.rb index 391901c..a05f700 100644 --- a/app/firehose_client.rb +++ b/app/firehose_client.rb @@ -153,7 +153,7 @@ class FirehoseClient return unless @current_user if op.action == :create - @current_user.likes.import_from_record(op.uri, op.raw_record) + @current_user.likes.import_from_record(op.uri, op.raw_record, queue: :firehose) elsif op.action == :delete @current_user.likes.where(rkey: op.rkey).delete_all end @@ -163,7 +163,7 @@ class FirehoseClient return unless @current_user if op.action == :create - @current_user.reposts.import_from_record(op.uri, op.raw_record) + @current_user.reposts.import_from_record(op.uri, op.raw_record, queue: :firehose) elsif op.action == :delete @current_user.reposts.where(rkey: op.rkey).delete_all end @@ -172,8 +172,8 @@ class FirehoseClient def process_post(msg, op) if op.action == :create if @current_user - @current_user.quotes.import_from_record(op.uri, op.raw_record) - @current_user.pins.import_from_record(op.uri, op.raw_record) + @current_user.quotes.import_from_record(op.uri, op.raw_record, queue: :firehose) + @current_user.pins.import_from_record(op.uri, op.raw_record, queue: :firehose) end elsif op.action == :delete if @current_user diff --git a/app/import_manager.rb b/app/import_manager.rb index 7163959..cbc8ba6 100644 --- a/app/import_manager.rb +++ b/app/import_manager.rb @@ -11,7 +11,7 @@ class ImportManager @user = user end - def start(sets, include_pending) + def start(sets) queued_items = [] importers = [] sets = [sets] unless sets.is_a?(Array) @@ -19,16 +19,16 @@ class ImportManager sets.each do |set| case set when 'likes' - queued_items += @user.likes.pending.to_a if include_pending + queued_items += @user.likes.in_queue(:import).to_a importers << LikesImporter.new(@user) when 'reposts' - queued_items += @user.reposts.pending.to_a if include_pending + queued_items += @user.reposts.in_queue(:import).to_a importers << RepostsImporter.new(@user) when 'posts' - queued_items += @user.quotes.pending.to_a + @user.pins.pending.to_a if include_pending + queued_items += @user.quotes.in_queue(:import).to_a + @user.pins.in_queue(:import).to_a importers << PostsImporter.new(@user) when 'all' - queued_items += @user.all_pending_items if include_pending + queued_items += @user.all_items_in_queue(:import) importers += [ LikesImporter.new(@user), RepostsImporter.new(@user), diff --git a/app/import_worker.rb b/app/import_worker.rb index 7974216..a678b84 100644 --- a/app/import_worker.rb +++ b/app/import_worker.rb @@ -23,7 +23,7 @@ class ImportWorker def run(collections) import = ImportManager.new(@user) import.report = BasicReport.new if @verbose - import.start(collections, false) + import.start(collections) end end diff --git a/app/importers/likes_importer.rb b/app/importers/likes_importer.rb index a90acee..69dc7bc 100644 --- a/app/importers/likes_importer.rb +++ b/app/importers/likes_importer.rb @@ -18,7 +18,7 @@ class LikesImporter < BaseImporter records.each do |record| begin - like = @user.likes.import_from_record(record['uri'], record['value']) + like = @user.likes.import_from_record(record['uri'], record['value'], queue: :import) if like && like.pending? && @item_queue @item_queue.push(like) diff --git a/app/importers/posts_importer.rb b/app/importers/posts_importer.rb index 38ec698..c47ae14 100644 --- a/app/importers/posts_importer.rb +++ b/app/importers/posts_importer.rb @@ -18,8 +18,8 @@ class PostsImporter < BaseImporter records.each do |record| begin - quote = @user.quotes.import_from_record(record['uri'], record['value']) - pin = @user.pins.import_from_record(record['uri'], record['value']) + quote = @user.quotes.import_from_record(record['uri'], record['value'], queue: :import) + pin = @user.pins.import_from_record(record['uri'], record['value'], queue: :import) if @item_queue if quote && quote.pending? diff --git a/app/importers/reposts_importer.rb b/app/importers/reposts_importer.rb index e173205..f845b04 100644 --- a/app/importers/reposts_importer.rb +++ b/app/importers/reposts_importer.rb @@ -18,7 +18,7 @@ class RepostsImporter < BaseImporter records.each do |record| begin - repost = @user.reposts.import_from_record(record['uri'], record['value']) + repost = @user.reposts.import_from_record(record['uri'], record['value'], queue: :import) if repost && repost.pending? && @item_queue @item_queue.push(repost) diff --git a/app/models/importable.rb b/app/models/importable.rb index 193c97e..1e192fb 100644 --- a/app/models/importable.rb +++ b/app/models/importable.rb @@ -8,11 +8,21 @@ module Importable included do scope :pending, -> { where(post: nil) } + scope :in_queue, ->(q) { where(queue: q) } + + enum :queue, { firehose: 0, import: 1 } + + validates_presence_of :post_uri, if: -> { post_id.nil? } + validate :check_queue def pending? post_uri != nil end + def check_queue + errors.add(:queue, 'must be nil if already processed') if queue && post + end + def import_item! post_uri = AT_URI(self.post_uri) return nil if !post_uri.is_post? @@ -20,6 +30,7 @@ module Importable if post = Post.find_by_at_uri(post_uri) self.post = post self.post_uri = nil + self.queue = nil end self.save! diff --git a/app/models/like.rb b/app/models/like.rb index bd4fb85..77435f0 100644 --- a/app/models/like.rb +++ b/app/models/like.rb @@ -14,8 +14,6 @@ class Like < ActiveRecord::Base validates_presence_of :time, :rkey validates_length_of :rkey, is: 13 - validates_presence_of :post_uri, if: -> { post_id.nil? } - belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true diff --git a/app/models/pin.rb b/app/models/pin.rb index 513191d..c98c00e 100644 --- a/app/models/pin.rb +++ b/app/models/pin.rb @@ -16,8 +16,6 @@ class Pin < ActiveRecord::Base validates_length_of :rkey, is: 13 validates :pin_text, length: { minimum: 0, maximum: 1000, allow_nil: false } - validates_presence_of :post_uri, if: -> { post_id.nil? } - belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true diff --git a/app/models/quote.rb b/app/models/quote.rb index adb75a5..0df21b0 100644 --- a/app/models/quote.rb +++ b/app/models/quote.rb @@ -14,8 +14,6 @@ class Quote < ActiveRecord::Base validates_length_of :rkey, is: 13 validates :quote_text, length: { minimum: 0, maximum: 1000, allow_nil: false } - validates_presence_of :post_uri, if: -> { post_id.nil? } - belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true diff --git a/app/models/repost.rb b/app/models/repost.rb index f2165c9..937ad62 100644 --- a/app/models/repost.rb +++ b/app/models/repost.rb @@ -13,8 +13,6 @@ class Repost < ActiveRecord::Base validates_presence_of :time, :rkey validates_length_of :rkey, is: 13 - validates_presence_of :post_uri, if: -> { post_id.nil? } - belongs_to :user, foreign_key: 'actor_id' belongs_to :post, optional: true diff --git a/app/models/user.rb b/app/models/user.rb index ba50407..d6d0718 100644 --- a/app/models/user.rb +++ b/app/models/user.rb @@ -19,37 +19,41 @@ class User < ActiveRecord::Base before_destroy :delete_posts_cascading has_many :likes, foreign_key: 'actor_id', dependent: :delete_all do - def import_from_record(like_uri, record) + def import_from_record(like_uri, record, **args) like = self.new_from_record(like_uri, record) return nil if like.nil? || self.where(rkey: like.rkey).exists? + like.assign_attributes(args) like.import_item! end end has_many :reposts, foreign_key: 'actor_id', dependent: :delete_all do - def import_from_record(repost_uri, record) + def import_from_record(repost_uri, record, **args) repost = self.new_from_record(repost_uri, record) return nil if repost.nil? || self.where(rkey: repost.rkey).exists? + repost.assign_attributes(args) repost.import_item! end end has_many :quotes, foreign_key: 'actor_id', dependent: :delete_all do - def import_from_record(post_uri, record) + def import_from_record(post_uri, record, **args) quote = self.new_from_record(post_uri, record) return nil if quote.nil? || self.where(rkey: quote.rkey).exists? + quote.assign_attributes(args) quote.import_item! end end has_many :pins, foreign_key: 'actor_id', dependent: :delete_all do - def import_from_record(post_uri, record) + def import_from_record(post_uri, record, **args) pin = self.new_from_record(post_uri, record) return nil if pin.nil? || self.where(rkey: pin.rkey).exists? + pin.assign_attributes(args) pin.import_item! end end @@ -70,6 +74,10 @@ class User < ActiveRecord::Base [:likes, :reposts, :quotes, :pins].map { |x| self.send(x).pending.to_a }.reduce(&:+) end + def all_items_in_queue(queue) + [:likes, :reposts, :quotes, :pins].map { |x| self.send(x).in_queue(queue).to_a }.reduce(&:+) + end + def delete_posts_cascading posts_subquery = self.posts.select(:id) diff --git a/app/post_downloader.rb b/app/post_downloader.rb index fd329e9..31be37d 100644 --- a/app/post_downloader.rb +++ b/app/post_downloader.rb @@ -88,7 +88,7 @@ class PostDownloader end def update_item(item, post) - item.update!(post: post, post_uri: nil) + item.update!(post: post, post_uri: nil, queue: nil) @total_count += 1 @oldest_imported = [@oldest_imported, item.time].min @@ -148,6 +148,10 @@ class PostDownloader item.destroy if hostname == 'bsky.social' end end + + if !item.destroyed? + item.update!(queue: nil) + end end end end diff --git a/db/migrate/20250923014702_add_queued_field.rb b/db/migrate/20250923014702_add_queued_field.rb new file mode 100644 index 0000000..53dd936 --- /dev/null +++ b/db/migrate/20250923014702_add_queued_field.rb @@ -0,0 +1,7 @@ +class AddQueuedField < ActiveRecord::Migration[7.2] + def change + [:likes, :reposts, :quotes, :pins].each do |table| + add_column(table, :queue, :smallint, null: true) + end + end +end diff --git a/db/schema.rb b/db/schema.rb index 93f7361..a1f29ad 100644 --- a/db/schema.rb +++ b/db/schema.rb @@ -10,7 +10,7 @@ # # It's strongly recommended that you check this file into your version control system. -ActiveRecord::Schema[7.2].define(version: 2025_09_20_182018) do +ActiveRecord::Schema[7.2].define(version: 2025_09_23_014702) do # These are extensions that must be enabled in order to support this database enable_extension "plpgsql" @@ -34,6 +34,7 @@ ActiveRecord::Schema[7.2].define(version: 2025_09_20_182018) do t.datetime "time", null: false t.bigint "post_id" t.string "post_uri" + t.integer "queue", limit: 2 t.index ["actor_id", "rkey"], name: "index_likes_on_actor_id_and_rkey", unique: true t.index ["actor_id", "time", "id"], name: "index_likes_on_actor_id_and_time_and_id", order: { time: :desc, id: :desc } end @@ -45,6 +46,7 @@ ActiveRecord::Schema[7.2].define(version: 2025_09_20_182018) do t.text "pin_text", null: false t.bigint "post_id" t.string "post_uri" + t.integer "queue", limit: 2 t.index ["actor_id", "rkey"], name: "index_pins_on_actor_id_and_rkey", unique: true t.index ["actor_id", "time", "id"], name: "index_pins_on_actor_id_and_time_and_id", order: { time: :desc, id: :desc } end @@ -66,6 +68,7 @@ ActiveRecord::Schema[7.2].define(version: 2025_09_20_182018) do t.text "quote_text", null: false t.bigint "post_id" t.string "post_uri" + t.integer "queue", limit: 2 t.index ["actor_id", "rkey"], name: "index_quotes_on_actor_id_and_rkey", unique: true t.index ["actor_id", "time", "id"], name: "index_quotes_on_actor_id_and_time_and_id", order: { time: :desc, id: :desc } end @@ -76,6 +79,7 @@ ActiveRecord::Schema[7.2].define(version: 2025_09_20_182018) do t.datetime "time", null: false t.bigint "post_id" t.string "post_uri" + t.integer "queue", limit: 2 t.index ["actor_id", "rkey"], name: "index_reposts_on_actor_id_and_rkey", unique: true t.index ["actor_id", "time", "id"], name: "index_reposts_on_actor_id_and_time_and_id", order: { time: :desc, id: :desc } end diff --git a/lib/tasks/import.rake b/lib/tasks/import.rake index a077eb6..af0556b 100644 --- a/lib/tasks/import.rake +++ b/lib/tasks/import.rake @@ -27,7 +27,6 @@ task :import_user do end user = User.find_or_create_by!(did: ENV['DID']) - pending = !ENV['SKIP_PENDING'] unless ENV['COLLECTION'] raise "Required COLLECTION parameter missing" @@ -42,7 +41,7 @@ task :import_user do exit } - import.start(ENV['COLLECTION'], pending) + import.start(ENV['COLLECTION']) puts "\n\n\n\n\n" end