diff --git a/bin/sim_job b/bin/sim_job index 7a4ef2486ddfe2f3ecf0e49e1297f2b8f566bd67..4ce1777d35ce286b8d649d05af1b43bf33c13061 100755 --- a/bin/sim_job +++ b/bin/sim_job @@ -2,30 +2,9 @@ # frozen_string_literal: true require "bigdecimal" -require "pg/em/connection_pool" -require "eventmachine" +require "delegate" require "em_promise" -require "em-hiredis" -require "dhall" -require "ougai" -require "sentry-ruby" - -$stdout.sync = true -LOG = Ougai::Logger.new($stdout) -LOG.level = ENV.fetch("LOG_LEVEL", "info") -LOG.formatter = Ougai::Formatters::Readable.new( - nil, - nil, - plain: !$stdout.isatty -) - -def log - LOG -end - -Sentry.init do |config| - config.background_worker_threads = 0 -end +require "set" SCHEMA = "{ xmpp: { jid: Text, password: Text }, @@ -53,34 +32,31 @@ SCHEMA = "{ } }" -raise "Need a Dhall config" unless ARGV[0] - -CONFIG = Dhall::Coder - .new(safe: Dhall::Coder::JSON_LIKE + [Symbol, Proc]) - .load("#{ARGV.first} : #{SCHEMA}", transform_keys: :to_sym) - -CONFIG[:keep_area_codes_in] = {} -CONFIG[:creds] = {} +class JobCustomer < SimpleDelegator + def initialize(customer, billing_customer: nil) + @billing_customer = billing_customer + super(customer) + end -require_relative "../lib/async_braintree" -require_relative "../lib/blather_notify" -require_relative "../lib/customer_repo" -require_relative "../lib/expiring_lock" -require_relative "../lib/low_balance" -require_relative "../lib/postgres" -require_relative "../lib/sim_repo" -require_relative "../lib/sim_pricing" -require_relative "../lib/transaction" + def auto_top_up_amount + @billing_customer&.auto_top_up_amount || super + end -BRAINTREE = AsyncBraintree.new(**CONFIG[:braintree]) -CUSTOMER_REPO = CustomerRepo.new -SIM_REPO = SIMRepo.new + def billing_customer(*) + if @billing_customer + return EMPromise.resolve(self.class.new(@billing_customer)) + end -class JobCustomer < SimpleDelegator - def billing_customer super.then(&self.class.method(:new)) end + def with(balance:) + self.class.new( + __getobj__.with(balance: balance), + billing_customer: @billing_customer&.with(balance: balance) + ) + end + def stanza_to(stanza) stanza = stanza.dup stanza.from = nil # It's a client connection, use default @@ -157,15 +133,23 @@ class SimTopUp def call EMPromise.all([amount_spent, monthly_limit]).then do |(spent, limit)| - if spent < limit + if monthly_limit_allows?(spent, limit) next low_balance if customer.balance < refill_price refill_and_bill(5120, refill_price, "5GB Data Topup for #{iccid}") else - SimWarn.new(@sim, customer: customer).call + warn end end end + + def warn + SimWarn.new(@sim, customer: customer).notify + end + + def monthly_limit_allows?(spent, limit) + spent + refill_price <= limit + end end class SimWarn @@ -177,18 +161,22 @@ class SimWarn m.body = "Your SIM #{iccid} only has " \ "#{(@sim.remaining_usage_kb / 1024).to_i} MB left" customer.stanza_to(m) + LOG.info "Data warning #{customer.customer_id} #{@sim.iccid}" end end def call EMPromise.all([amount_spent, monthly_limit]).then do |(spent, limit)| - next unless spent >= limit || low_balance_and_not_auto_top_up + next unless notify?(spent, limit) notify - LOG.info "Data warning #{customer.customer_id} #{@sim.iccid}" end end + def notify?(spent, limit) + spent + refill_price > limit || low_balance_and_not_auto_top_up + end + def low_balance_and_not_auto_top_up customer.balance < refill_price && !customer.auto_top_up_amount&.positive? end @@ -202,8 +190,8 @@ class SimAnnual m = Blather::Stanza::Message.new m.body = "Your SIM #{iccid} only has #{@sim.remaining_days} days left" customer.stanza_to(m) + LOG.info "Annual warning #{customer.customer_id} #{@sim.iccid}" end - LOG.info "Annual warning #{customer.customer_id} #{@sim.iccid}" end def annual_price @@ -214,22 +202,27 @@ class SimAnnual if customer.balance >= annual_price refill_and_bill(1024, annual_price, "Annual fee for #{iccid}") else - LowBalance.for(customer, annual_price).then(&:notify!).then do |result| - next call if result.positive? + low_balance + end + end - notify - end + def low_balance + LowBalance.for(customer, annual_price).then(&:notify!).then do |result| + @customer = customer.with(balance: customer.balance + result) + next call if result.positive? + + notify end end end -def fetch_customers(cids) +def fetch_customers(cids, customer_repo: CUSTOMER_REPO) # This is gross N+1 against the DB, but also does a buch of Redis work # We expect the set to be very small for the forseeable future, # hundreds at most EMPromise.all( Set.new(cids).to_a.compact.map { |id| - CUSTOMER_REPO.find(id).catch_only(CustomerRepo::NotFound) { nil } + customer_repo.find(id).catch_only(CustomerRepo::NotFound) { nil } } ).then do |customers| Hash[customers.compact.map { |c| [c.customer_id, JobCustomer.new(c)] }] @@ -237,16 +230,23 @@ def fetch_customers(cids) end SIM_QUERY = "SELECT iccid, customer_id FROM sims WHERE iccid = ANY ($1)" -def load_customers!(sims) - DB.query_defer(SIM_QUERY, [sims.keys]).then { |rows| - fetch_customers(rows.map { |row| row["customer_id"] }).then do |customers| - rows.each do |row| - sims[row["iccid"]]&.customer = customers[row["customer_id"]] - end +def load_customers!(sims, db: DB, customer_repo: CUSTOMER_REPO) + db.query_defer(SIM_QUERY, [sims.keys]).then { |rows| + load_customer_rows!(sims, rows, customer_repo) + } +end - sims +def load_customer_rows!(sims, rows, customer_repo) + fetch_customers( + rows.map { |row| row["customer_id"] }, + customer_repo: customer_repo + ).then do |customers| + rows.each do |row| + sims[row["iccid"]]&.customer = customers[row["customer_id"]] end - } + + sims + end end def decide_sim_actions(sims) @@ -261,31 +261,109 @@ def decide_sim_actions(sims) }.compact end -EM.run do - REDIS = EM::Hiredis.connect - DB = Postgres.connect(dbname: "jmp") - - BlatherNotify.start( - CONFIG[:xmpp][:jid], - CONFIG[:xmpp][:password] - ).then { SIM_REPO.all }.then { |sims| - load_customers!(decide_sim_actions(sims)) - }.then { |items| - items = items.values.select { |item| item.customer&.currency } - EMPromise.all(items.map(&:call)) +def report_sim_job_error(e) + LOG.error e + + case e + when Exception then Sentry.capture_exception(e) + when EM::HttpClient + LOG.error e.response_header + Sentry.capture_message("HTTP Error: #{e.response_header.http_status} " \ + "@ #{e.last_effective_url}") + else + Sentry.capture_message(e.to_s) + end +end + +def process_sim_actions(actions, customer_repo: CUSTOMER_REPO) + actions.reduce(EMPromise.resolve(nil)) do |promise, action| + promise.then { process_sim_action(action, customer_repo) } + end +end + +def process_sim_action(action, customer_repo) + EMPromise.resolve(nil).then { + reload_action_customer(action, customer_repo).then { |customer| + action.customer = customer + action.call + } }.catch { |e| - LOG.error e - - case e - when Exception - Sentry.capture_exception(e) - when EM::HttpClient - LOG.error e.response_header - Sentry.capture_message( - "HTTP Error: #{e.response_header.http_status} @ #{e.last_effective_url}" + report_sim_job_error(e) + nil + } +end + +def reload_action_customer(action, customer_repo) + customer_repo.find(action.customer.customer_id).then do |customer| + customer.billing_customer(customer_repo).then do |billing_customer| + JobCustomer.new( + customer.with(balance: billing_customer.balance), + billing_customer: billing_customer ) - else - Sentry.capture_message(e.to_s) end - }.then { EM.stop } + end +end + +if $0 == __FILE__ + require "pg/em/connection_pool" + require "eventmachine" + require "em-hiredis" + require "dhall" + require "ougai" + require "sentry-ruby" + + $stdout.sync = true + LOG = Ougai::Logger.new($stdout) + LOG.level = ENV.fetch("LOG_LEVEL", "info") + LOG.formatter = Ougai::Formatters::Readable.new( + nil, + nil, + plain: !$stdout.isatty + ) + + def log + LOG + end + + Sentry.init do |config| + config.background_worker_threads = 0 + end + + raise "Need a Dhall config" unless ARGV[0] + + CONFIG = Dhall::Coder + .new(safe: Dhall::Coder::JSON_LIKE + [Symbol, Proc]) + .load("#{ARGV.first} : #{SCHEMA}", transform_keys: :to_sym) + + CONFIG[:keep_area_codes_in] = {} + CONFIG[:creds] = {} + + require_relative "../lib/async_braintree" + require_relative "../lib/blather_notify" + require_relative "../lib/customer_repo" + require_relative "../lib/expiring_lock" + require_relative "../lib/low_balance" + require_relative "../lib/postgres" + require_relative "../lib/sim_repo" + require_relative "../lib/sim_pricing" + require_relative "../lib/transaction" + + BRAINTREE = AsyncBraintree.new(**CONFIG[:braintree]) + CUSTOMER_REPO = CustomerRepo.new + SIM_REPO = SIMRepo.new + + EM.run do + REDIS = EM::Hiredis.connect + DB = Postgres.connect(dbname: "jmp") + + BlatherNotify.start( + CONFIG[:xmpp][:jid], + CONFIG[:xmpp][:password] + ).then { SIM_REPO.all }.then { |sims| + load_customers!(decide_sim_actions(sims)) + }.then { |items| + items = items.values.select { |item| item.customer&.currency } + process_sim_actions(items) + }.catch(&method(:report_sim_job_error)).then { EM.stop } + end end diff --git a/test/test_sim_job.rb b/test/test_sim_job.rb new file mode 100644 index 0000000000000000000000000000000000000000..1842226549d7d08831d3ef4ee8b3a7e7e57386dd --- /dev/null +++ b/test/test_sim_job.rb @@ -0,0 +1,234 @@ +# frozen_string_literal: true + +require "bigdecimal" +require "em-http" +require "test_helper" +require "customer" +require "low_balance" +require "sim" + +load File.expand_path("../bin/sim_job", __dir__) + +class SimJobTest < Minitest::Test + class Ledger + attr_reader :charges + + def initialize(balances) + @balances = balances + @charges = [] + end + + def balance_for(customer_id) + @balances.fetch(customer_id, BigDecimal("0")) + end + + def charge(customer_id, amount) + @balances[customer_id] = balance_for(customer_id) - amount + @charges << [customer_id, amount] + end + end + + class CustomerRepo + def initialize(customers, ledger) + @customers = customers + @ledger = ledger + end + + def find(customer_id) + customer = @customers.fetch(customer_id) + EMPromise.resolve(customer.with(balance: balance_for(customer))) + end + + def balance_for(customer) + @ledger.balance_for(customer.billing_customer_id) + end + end + + class RecordingAction + attr_accessor :customer + attr_reader :balances_seen, :iccid + + def initialize(iccid, customer, ledger:, price: BigDecimal("5"), error: nil) + @iccid = iccid + @customer = customer + @ledger = ledger + @price = price + @error = error + @balances_seen = [] + end + + def call + @balances_seen << customer.balance + return EMPromise.reject(@error) if @error + + if customer.balance >= @price + @ledger.charge(customer.billing_customer_id, @price) + end + EMPromise.resolve(nil) + end + end + + class SuccessfulLowBalance + def initialize(amount) + @amount = amount + end + + def notify! + EMPromise.resolve(@amount) + end + end + + def test_process_sim_actions_reloads_balance_between_actions + ledger = Ledger.new("parent" => BigDecimal("5")) + repo = CustomerRepo.new( + { + "parent" => customer("parent"), + "first" => customer("first", parent_customer_id: "parent"), + "second" => customer("second", parent_customer_id: "parent") + }, + ledger + ) + first = RecordingAction.new( + "first", customer("first", parent_customer_id: "parent"), ledger: ledger + ) + second = RecordingAction.new( + "second", customer("second", parent_customer_id: "parent"), ledger: ledger + ) + + process_sim_actions([first, second], customer_repo: repo).sync + + assert_equal [BigDecimal("5")], first.balances_seen + assert_equal [BigDecimal("0")], second.balances_seen + assert_equal [["parent", BigDecimal("5")]], ledger.charges + end + em :test_process_sim_actions_reloads_balance_between_actions + + def test_process_sim_actions_continues_after_failure + ledger = Ledger.new("first" => BigDecimal("5"), "second" => BigDecimal("5")) + repo = CustomerRepo.new( + { + "first" => customer("first"), + "second" => customer("second") + }, + ledger + ) + error = RuntimeError.new("boom") + first = RecordingAction.new( + "first", customer("first"), ledger: ledger, error: error + ) + second = RecordingAction.new("second", customer("second"), ledger: ledger) + + process_sim_actions([first, second], customer_repo: repo).sync + + assert_equal [BigDecimal("5")], first.balances_seen + assert_equal [BigDecimal("5")], second.balances_seen + assert_equal [["second", BigDecimal("5")]], ledger.charges + end + em :test_process_sim_actions_continues_after_failure + + def test_job_customer_preserves_billing_customer_when_balance_changes + parent = customer( + "parent", + balance: BigDecimal("45"), + auto_top_up_amount: 15 + ) + child = customer( + "child", + parent_customer_id: "parent", + balance: BigDecimal("0") + ) + + updated = JobCustomer.new(child, billing_customer: parent) + .with(balance: BigDecimal("40")) + billing_customer = updated.billing_customer.sync + + assert_kind_of JobCustomer, updated + assert_equal "child", updated.customer_id + assert_equal BigDecimal("40"), updated.balance + assert_kind_of JobCustomer, billing_customer + assert_equal "parent", billing_customer.customer_id + assert_equal BigDecimal("40"), billing_customer.balance + end + em :test_job_customer_preserves_billing_customer_when_balance_changes + + def test_sim_top_up_refills_when_refill_exactly_fits_monthly_limit + top_up = SimTopUp.new(sim, customer: customer(balance: BigDecimal("100"))) + refill = nil + + top_up.stub(:amount_spent, EMPromise.resolve(BigDecimal("1"))) do + top_up.stub(:monthly_limit, EMPromise.resolve(BigDecimal("46"))) do + top_up.stub(:refill_price, BigDecimal("45")) do + top_up.stub(:refill_and_bill, lambda { |*args| + refill = args + EMPromise.resolve(nil) + }) do + top_up.call.sync + end + end + end + end + + assert_equal( + [5120, BigDecimal("45"), "5GB Data Topup for 123"], + refill + ) + end + em :test_sim_top_up_refills_when_refill_exactly_fits_monthly_limit + + def test_sim_top_up_warns_when_refill_would_exceed_monthly_limit + top_up = SimTopUp.new(sim, customer: customer(balance: BigDecimal("100"))) + warned = false + refilled = false + + top_up.stub(:amount_spent, EMPromise.resolve(BigDecimal("45"))) do + top_up.stub(:monthly_limit, EMPromise.resolve(BigDecimal("46"))) do + top_up.stub(:refill_price, BigDecimal("45")) do + top_up.stub(:refill_and_bill, lambda { |*| + refilled = true + EMPromise.resolve(nil) + }) do + top_up.stub(:warn, lambda { + warned = true + EMPromise.resolve(nil) + }) do + top_up.call.sync + end + end + end + end + end + + assert warned + refute refilled + end + em :test_sim_top_up_warns_when_refill_would_exceed_monthly_limit + + def test_sim_annual_retries_with_low_balance_result_added_to_balance + annual = SimAnnual.new(sim, customer: customer(balance: BigDecimal("0"))) + balances_billed = [] + low_balance = SuccessfulLowBalance.new(BigDecimal("10")) + + LowBalance.stub(:for, EMPromise.resolve(low_balance)) do + annual.stub(:annual_price, BigDecimal("10")) do + annual.stub(:refill_and_bill, lambda { |*| + balances_billed << annual.customer.balance + EMPromise.resolve(nil) + }) do + annual.call.sync + end + end + end + + assert_equal [BigDecimal("10")], balances_billed + end + em :test_sim_annual_retries_with_low_balance_result_added_to_balance + + def sim + SIM.new( + iccid: "123", + remaining_usage_kb: 1, + remaining_days: 365, + notes: "" + ) + end +end