#!/usr/bin/env ruby require 'bundler/setup' require 'fileutils' require 'json' require 'skyfall' require 'yaml' require_relative 'init' require_relative 'models' require_relative 'opts' ActiveRecord::Base.logger = nil def get_hosts(entries) if entries entries.map { |x| x.is_a?(Hash) ? x['host'] : x } else [] end end config = YAML.load(File.read(SOURCES)) options = parse_options(ARGV) if options[:relays] || options[:jetstreams] relays = options[:relays] || [] jetstreams = options[:jetstreams] || [] else relays = get_hosts(config['relays']) jetstreams = get_hosts(config['jetstreams']) end maxlen = (relays + jetstreams).map(&:length).max verbose = options[:verbose] duration = options[:duration] || config['duration']&.to_i || 60 * 15 test_start_time = Time.now log_dir = File.expand_path(File.join(__dir__, 'log', 'runs', test_start_time.getutc.iso8601.gsub(':', '-'))) FileUtils.mkdir_p(log_dir) Worker = Struct.new(:host, :type, :pid, :pipe) workers = [] sources = relays.map { |h| [:firehose, h] } + jetstreams.map { |h| [:jetstream, h] } sources.each do |type, host| input, output = IO.pipe pid = fork do input.close sky = (type == :firehose) ? Skyfall::Firehose.new(host) : Skyfall::Jetstream.new(host) events = 0 reconnects = 0 errors = 0 users = Set.new msg_times = [] post_times = [] minute = Time.now.to_i / verbose if verbose connected = false got_first_event = false sky.on_connecting do puts "[#{Time.now}] #{host}: Connecting..." end sky.on_connect do puts "[#{Time.now}] #{host}: Connected ✓" connected = true end sky.on_reconnect do puts "[#{Time.now}] #{host}: Connection lost, reconnecting..." reconnects += 1 end sky.on_error do |e| puts "[#{Time.now}] #{host}: ERROR: #{e.message}" errors += 1 end sky.on_message do |msg| if msg.operations.length > 0 || [:identity, :account].include?(msg.type) events += 1 users << msg.did end if type == :firehose msg_times << msg.time.to_f else msg_times << msg.time_us.to_f / 1_000_000 end msg.operations.each do |op| if op.type == :bsky_post && op.action == :create if timestamp = op.raw_record['createdAt'] begin post_times << Time.new(timestamp).to_f rescue StandardError # skip end end end end if verbose if !got_first_event puts "[#{Time.now}] #{host}: Starting at seq #{msg.seq}" got_first_event = true end now = Time.now.to_i / verbose if now > minute puts "[#{Time.now}] #{host.ljust(maxlen)} | events: #{events.to_s.ljust(8)} | users: #{users.size}" minute = now end end end trap('SIGINT') { sky.disconnect } sky.connect puts "[#{Time.now}] #{host}: Finished at seq #{sky.cursor}" output.puts(JSON.generate({ events: events, users: users.size, connected: connected, errors: errors, reconnects: reconnects })) log_file = File.join(log_dir, host + ".log") File.write(log_file, JSON.generate({ users: users.to_a, msg_times: msg_times, post_times: post_times })) end output.close workers << Worker.new(host, type, pid, input) end begin sleep(duration) Process.kill('SIGINT', *workers.map(&:pid)) unless options[:dont_save] test = TestRun.create!(start_time: test_start_time, duration: duration) end while !workers.empty? pid = Process.wait worker = workers.detect { |w| w.pid == pid } workers.delete(worker) line = worker.pipe.gets next if line.nil? result = JSON.parse(line) puts "#{worker.host}: #{result.inspect}" if verbose unless options[:dont_save] test.reports.create!( host: worker.host, source_type: worker.type, users: result['users'], events: result['events'], connected: result['connected'], error_count: result['errors'], reconnect_count: result['reconnects'] ) end end test&.update_max_users rescue Interrupt puts puts "Stopping..." Process.kill('SIGINT', *workers.map(&:pid)) while !workers.empty? pid = Process.wait worker = workers.detect { |w| w.pid == pid } workers.delete(worker) line = worker.pipe.gets end end