bulk-update-migration-status.mjs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358
  1. #!/usr/bin/env node
  2. /**
  3. * This script validates subscriptions that are canceled or expired and have migration metadata set to "in_progress",
  4. * then updates the metadata to "cancelled" if validation passes.
  5. *
  6. * TODO: This script can be deleted after being run in production.
  7. *
  8. * Usage:
  9. * node scripts/stripe/bulk-update-migration-status.mjs [OPTS] [INPUT-FILE]
  10. *
  11. * Options:
  12. * --output PATH Output file path (default: /tmp/bulk_update_output_<timestamp>.csv)
  13. * Use '-' to write to stdout
  14. * --commit Apply changes (without this flag, runs in dry-run mode)
  15. * --concurrency N Number of subscriptions to process concurrently (default: 10)
  16. * --stripe-rate-limit N Requests per second for Stripe (default: 50)
  17. * --stripe-api-retries N Number of retries on Stripe 429s (default: 5)
  18. * --stripe-retry-delay-ms N Delay between Stripe retries in ms (default: 1000)
  19. * --help Show a help message
  20. *
  21. * CSV Input Format:
  22. * The CSV must have the following columns:
  23. * - subscription_id: Stripe subscription id
  24. * - target_stripe_account: Either 'stripe-uk' or 'stripe-us'
  25. *
  26. * Output:
  27. * Writes a CSV with columns:
  28. * - subscription_id: The subscription id processed
  29. * - target_stripe_account: The Stripe account
  30. * - status: Result status (validated, updated, invalid-status, invalid-metadata, or error)
  31. * - note: Additional information about the status
  32. */
  33. import fs from 'node:fs'
  34. import path from 'node:path'
  35. import * as csv from 'csv'
  36. import minimist from 'minimist'
  37. import PQueue from 'p-queue'
  38. import { z } from '../../app/src/infrastructure/Validation.mjs'
  39. import { scriptRunner } from '../lib/ScriptRunner.mjs'
  40. import { getRegionClient } from '../../modules/subscriptions/app/src/StripeClient.mjs'
  41. import { ReportError } from './helpers.mjs'
  42. import {
  43. createRateLimitedApiWrappers,
  44. DEFAULT_STRIPE_RATE_LIMIT,
  45. DEFAULT_STRIPE_API_RETRIES,
  46. DEFAULT_STRIPE_RETRY_DELAY_MS,
  47. } from './RateLimiter.mjs'
  48. const DEFAULT_CONCURRENCY = 10
  49. // rate limiters - initialized in main()
  50. let rateLimiters
  51. function usage() {
  52. console.error(`Usage: node scripts/stripe/bulk-update-migration-status.mjs [OPTS] [INPUT-FILE]
  53. Options:
  54. --output PATH Output file path (default: /tmp/bulk_update_output_<timestamp>.csv)
  55. Use '-' to write to stdout
  56. --commit Apply changes (without this, runs in dry-run mode)
  57. --concurrency N Number of subscriptions to process concurrently (default: ${DEFAULT_CONCURRENCY})
  58. --stripe-rate-limit N Requests per second for Stripe (default: ${DEFAULT_STRIPE_RATE_LIMIT})
  59. --stripe-api-retries N Number of retries on Stripe 429s (default: ${DEFAULT_STRIPE_API_RETRIES})
  60. --stripe-retry-delay-ms N Delay between Stripe retries in ms (default: ${DEFAULT_STRIPE_RETRY_DELAY_MS})
  61. --help Show this help message
  62. `)
  63. }
  64. async function main(trackProgress) {
  65. const opts = parseArgs()
  66. const timestamp = new Date().toISOString().replace(/[:.]/g, '-')
  67. const outputFile = opts.output ?? `/tmp/bulk_update_output_${timestamp}.csv`
  68. // initialize rate limiters
  69. rateLimiters = createRateLimitedApiWrappers({
  70. stripeRateLimit: opts.stripeRateLimit,
  71. stripeApiRetries: opts.stripeApiRetries,
  72. stripeRetryDelayMs: opts.stripeRetryDelayMs,
  73. })
  74. await trackProgress(
  75. 'Starting bulk validation and update of subscription metadata'
  76. )
  77. await trackProgress(`Run mode: ${opts.commit ? 'COMMIT' : 'DRY RUN'}`)
  78. await trackProgress(`Rate limit: Stripe ${opts.stripeRateLimit}/s`)
  79. await trackProgress(`Concurrency: ${opts.concurrency}`)
  80. const inputStream = opts.inputFile
  81. ? fs.createReadStream(opts.inputFile)
  82. : process.stdin
  83. const csvReader = getCsvReader(inputStream)
  84. const csvWriter = getCsvWriter(outputFile)
  85. await trackProgress(`Output: ${outputFile === '-' ? 'stdout' : outputFile}`)
  86. let processedCount = 0
  87. let successCount = 0
  88. let errorCount = 0
  89. const queue = new PQueue({ concurrency: opts.concurrency })
  90. const maxQueueSize = opts.concurrency
  91. try {
  92. for await (const input of csvReader) {
  93. if (queue.size >= maxQueueSize) {
  94. await queue.onSizeLessThan(maxQueueSize)
  95. }
  96. queue.add(async () => {
  97. try {
  98. const result = await processValidation(input, opts.commit)
  99. csvWriter.write({
  100. subscription_id: input.subscription_id,
  101. target_stripe_account: input.target_stripe_account,
  102. status: result.status,
  103. note:
  104. result.note ||
  105. (opts.commit ? '' : 'dry run - no changes applied'),
  106. })
  107. if (result.status === 'updated' || result.status === 'validated') {
  108. successCount++
  109. } else {
  110. errorCount++
  111. }
  112. } catch (err) {
  113. errorCount++
  114. if (err instanceof ReportError) {
  115. csvWriter.write({
  116. subscription_id: input.subscription_id,
  117. target_stripe_account: input.target_stripe_account,
  118. status: err.status,
  119. note: err.message,
  120. })
  121. } else {
  122. csvWriter.write({
  123. subscription_id: input.subscription_id,
  124. target_stripe_account: input.target_stripe_account,
  125. status: 'error',
  126. note: err.message,
  127. })
  128. await trackProgress(
  129. `Error processing ${input.subscription_id}: ${err.message}`
  130. )
  131. }
  132. }
  133. processedCount++
  134. if (processedCount % 10 === 0) {
  135. await trackProgress(
  136. `Processed ${processedCount} subscriptions (${successCount} ${opts.commit ? 'updated' : 'validated'}, ${errorCount} errors)`
  137. )
  138. }
  139. })
  140. }
  141. } finally {
  142. await queue.onIdle()
  143. }
  144. await trackProgress(`✅ Total processed: ${processedCount}`)
  145. if (opts.commit) {
  146. await trackProgress(`✅ Successfully updated: ${successCount}`)
  147. } else {
  148. await trackProgress(`✅ Successfully validated: ${successCount}`)
  149. await trackProgress('ℹ️ DRY RUN: No changes were applied')
  150. }
  151. await trackProgress(`❌ Errors: ${errorCount}`)
  152. await trackProgress('🎉 Script completed!')
  153. csvWriter.end()
  154. }
  155. function parseArgs() {
  156. const args = minimist(process.argv.slice(2), {
  157. string: [
  158. 'output',
  159. 'concurrency',
  160. 'stripe-rate-limit',
  161. 'stripe-api-retries',
  162. 'stripe-retry-delay-ms',
  163. ],
  164. boolean: ['commit', 'help'],
  165. default: {
  166. commit: false,
  167. concurrency: DEFAULT_CONCURRENCY,
  168. 'stripe-rate-limit': DEFAULT_STRIPE_RATE_LIMIT,
  169. 'stripe-api-retries': DEFAULT_STRIPE_API_RETRIES,
  170. 'stripe-retry-delay-ms': DEFAULT_STRIPE_RETRY_DELAY_MS,
  171. },
  172. unknown: arg => {
  173. if (arg.startsWith('-')) {
  174. console.error(`Unknown option: ${arg}`)
  175. usage()
  176. process.exit(1)
  177. }
  178. return true
  179. },
  180. })
  181. if (args.help) {
  182. usage()
  183. process.exit(0)
  184. }
  185. const inputFile = args._[0]
  186. const paramsSchema = z.object({
  187. output: z.string().optional(),
  188. commit: z.boolean(),
  189. concurrency: z.number().int().positive(),
  190. stripeRateLimit: z.number().positive(),
  191. stripeApiRetries: z.number().int().nonnegative(),
  192. stripeRetryDelayMs: z.number().int().nonnegative(),
  193. inputFile: z.string().optional(),
  194. })
  195. try {
  196. return paramsSchema.parse({
  197. output: args.output,
  198. commit: args.commit,
  199. concurrency: Number(args.concurrency),
  200. stripeRateLimit: Number(args['stripe-rate-limit']),
  201. stripeApiRetries: Number(args['stripe-api-retries']),
  202. stripeRetryDelayMs: Number(args['stripe-retry-delay-ms']),
  203. inputFile,
  204. })
  205. } catch (err) {
  206. console.error('Invalid arguments:', err.message)
  207. usage()
  208. process.exit(1)
  209. }
  210. }
  211. function getCsvReader(inputStream) {
  212. const parser = csv.parse({ columns: true })
  213. inputStream.pipe(parser)
  214. return parser
  215. }
  216. function getCsvWriter(outputFile) {
  217. if (outputFile === '-') {
  218. const writer = csv.stringify({
  219. columns: ['subscription_id', 'target_stripe_account', 'status', 'note'],
  220. header: true,
  221. })
  222. writer.on('error', err => {
  223. console.error(err)
  224. process.exit(1)
  225. })
  226. writer.pipe(process.stdout)
  227. return writer
  228. }
  229. fs.mkdirSync(path.dirname(outputFile), { recursive: true })
  230. const outputStream = fs.createWriteStream(outputFile)
  231. const writer = csv.stringify({
  232. columns: ['subscription_id', 'target_stripe_account', 'status', 'note'],
  233. header: true,
  234. })
  235. writer.on('error', err => {
  236. console.error(err)
  237. process.exit(1)
  238. })
  239. writer.pipe(outputStream)
  240. return writer
  241. }
  242. async function processValidation(input, commit) {
  243. const {
  244. subscription_id: subscriptionId,
  245. target_stripe_account: targetStripeAccount,
  246. } = input
  247. // get Stripe client for the target account (strip 'stripe-' prefix if present)
  248. const region = targetStripeAccount.replace(/^stripe-/, '')
  249. const stripeClient = getRegionClient(region)
  250. // fetch subscription
  251. let subscription
  252. try {
  253. subscription = await rateLimiters.requestWithRetries(
  254. stripeClient.serviceName,
  255. () => stripeClient.stripe.subscriptions.retrieve(subscriptionId),
  256. {
  257. operation: 'subscriptions.retrieve',
  258. subscriptionId,
  259. region: stripeClient.serviceName,
  260. }
  261. )
  262. } catch (err) {
  263. throw new ReportError(
  264. 'subscription-not-found',
  265. `Subscription not found: ${err.message}`
  266. )
  267. }
  268. const validStatuses = ['canceled']
  269. if (!validStatuses.includes(subscription.status)) {
  270. throw new ReportError(
  271. 'invalid-status',
  272. `Subscription status is ${subscription.status}, expected canceled`
  273. )
  274. }
  275. if (
  276. subscription.metadata?.recurly_to_stripe_migration_status !== 'in_progress'
  277. ) {
  278. throw new ReportError(
  279. 'invalid-metadata',
  280. `Migration status is ${subscription.metadata?.recurly_to_stripe_migration_status}, expected in_progress`
  281. )
  282. }
  283. if (!commit) {
  284. return {
  285. status: 'validated',
  286. note: 'Subscription is valid for update',
  287. }
  288. }
  289. try {
  290. await rateLimiters.requestWithRetries(
  291. stripeClient.serviceName,
  292. () =>
  293. stripeClient.updateSubscriptionMetadata(subscriptionId, {
  294. recurly_to_stripe_migration_status: 'cancelled',
  295. }),
  296. {
  297. operation: 'updateSubscriptionMetadata',
  298. subscriptionId,
  299. region: stripeClient.serviceName,
  300. }
  301. )
  302. return {
  303. status: 'updated',
  304. note: `Updated metadata for subscription ${subscriptionId}`,
  305. }
  306. } catch (err) {
  307. throw new ReportError(
  308. 'update-failed',
  309. `Failed to update metadata: ${err.message}`
  310. )
  311. }
  312. }
  313. try {
  314. await scriptRunner(main)
  315. process.exit(0)
  316. } catch (error) {
  317. console.error(error)
  318. process.exit(1)
  319. }