diff --git a/bun.lock b/bun.lock index 48c8775cc..7da7f7b01 100644 --- a/bun.lock +++ b/bun.lock @@ -26,7 +26,7 @@ "@autumn/shared": "workspace:*", "chalk": "^5.3.0", "dotenv": "^16.5.0", - "drizzle-orm": "^0.44.7", + "drizzle-orm": "catalog:", "inquirer": "^12.6.3", "ora": "^9.0.0", "p-limit": "^7.2.0", @@ -3187,7 +3187,7 @@ "@anthropic-ai/sdk/@types/node": ["@types/node@18.19.130", "", { "dependencies": { "undici-types": "~5.26.4" } }, "sha512-GRaXQx6jGfL8sKfaIDD6OupbIHBr9jv7Jnaml9tB7l4v068PAOXqfcujMMo5PhbIs6ggR1XODELqahT2R8v0fg=="], - "@autumn/scripts/drizzle-orm": ["drizzle-orm@0.44.7", "", { "peerDependencies": { "@aws-sdk/client-rds-data": ">=3", "@cloudflare/workers-types": ">=4", "@electric-sql/pglite": ">=0.2.0", "@libsql/client": ">=0.10.0", "@libsql/client-wasm": ">=0.10.0", "@neondatabase/serverless": ">=0.10.0", "@op-engineering/op-sqlite": ">=2", "@opentelemetry/api": "^1.4.1", "@planetscale/database": ">=1.13", "@prisma/client": "*", "@tidbcloud/serverless": "*", "@types/better-sqlite3": "*", "@types/pg": "*", "@types/sql.js": "*", "@upstash/redis": ">=1.34.7", "@vercel/postgres": ">=0.8.0", "@xata.io/client": "*", "better-sqlite3": ">=7", "bun-types": "*", "expo-sqlite": ">=14.0.0", "gel": ">=2", "knex": "*", "kysely": "*", "mysql2": ">=2", "pg": ">=8", "postgres": ">=3", "sql.js": ">=1", "sqlite3": ">=5" }, "optionalPeers": ["@aws-sdk/client-rds-data", "@cloudflare/workers-types", "@electric-sql/pglite", "@libsql/client", "@libsql/client-wasm", "@neondatabase/serverless", "@op-engineering/op-sqlite", "@opentelemetry/api", "@planetscale/database", "@prisma/client", "@tidbcloud/serverless", "@types/better-sqlite3", "@types/pg", "@types/sql.js", "@upstash/redis", "@vercel/postgres", "@xata.io/client", "better-sqlite3", "bun-types", "expo-sqlite", "gel", "knex", "kysely", "mysql2", "pg", "postgres", "sql.js", "sqlite3"] }, "sha512-quIpnYznjU9lHshEOAYLoZ9s3jweleHlZIAWR/jX9gAWNg/JhQ1wj0KGRf7/Zm+obRrYd9GjPVJg790QY9N5AQ=="], + "@autumn/shared/@types/bun": ["@types/bun@1.3.4", "", { "dependencies": { "bun-types": "1.3.4" } }, "sha512-EEPTKXHP+zKGPkhRLv+HI0UEX8/o+65hqARxLy8Ov5rIxMBPNTjeZww00CIihrIQGEQBYg+0roO5qOnS/7boGA=="], "@autumn/vite/@types/node": ["@types/node@22.19.1", "", { "dependencies": { "undici-types": "~6.21.0" } }, "sha512-LCCV0HdSZZZb34qifBsyWlUmok6W7ouER+oQIGBScS8EsZsQbrtFTUrDX4hOl+CS6p7cnNC4td+qrSVGSCTUfQ=="], @@ -3961,6 +3961,8 @@ "@anthropic-ai/sdk/@types/node/undici-types": ["undici-types@5.26.5", "", {}, "sha512-JlCMO+ehdEIKqlFxk6IfVoAUVmgz7cU7zD/h9XZ0qzeosSHmUJVOzSQvvYSYWXkFXC+IfLKSIffhv0sVZup6pA=="], + "@autumn/shared/@types/bun/bun-types": ["bun-types@1.3.4", "", { "dependencies": { "@types/node": "*" } }, "sha512-5ua817+BZPZOlNaRgGBpZJOSAQ9RQ17pkwPD0yR7CfJg+r8DgIILByFifDTa+IPDDxzf5VNhtNlcKqFzDgJvlQ=="], + "@autumn/vite/@types/node/undici-types": ["undici-types@6.21.0", "", {}, "sha512-iwDZqg0QAGrg9Rav5H4n0M64c3mkR59cJ6wQp+7C4nI0gsmExaedaYLNO44eT4AtBBwjbTiGPMlt2Md0T9H9JQ=="], "@aws-crypto/sha256-browser/@smithy/util-utf8/@smithy/util-buffer-from": ["@smithy/util-buffer-from@2.2.0", "", { "dependencies": { "@smithy/is-array-buffer": "^2.2.0", "tslib": "^2.6.2" } }, "sha512-IJdWBbTcMQ6DA0gdNhh/BwrLkDR+ADW5Kr1aZmd4k3DIF6ezMV4R2NIAmT08wQJ3yUK82thHWmC/TnK/wpMMIA=="], diff --git a/scripts/package.json b/scripts/package.json index 43ce88d63..63238826f 100644 --- a/scripts/package.json +++ b/scripts/package.json @@ -12,7 +12,7 @@ "@autumn/shared": "workspace:*", "chalk": "^5.3.0", "dotenv": "^16.5.0", - "drizzle-orm": "^0.44.7", + "drizzle-orm": "catalog:", "inquirer": "^12.6.3", "ora": "^9.0.0", "p-limit": "^7.2.0" diff --git a/server/src/db/initDrizzle.ts b/server/src/db/initDrizzle.ts index c361e2afd..607c29b3c 100644 --- a/server/src/db/initDrizzle.ts +++ b/server/src/db/initDrizzle.ts @@ -9,9 +9,10 @@ import postgres from "postgres"; export const client = postgres(process.env.DATABASE_URL!); export const db = drizzle(client, { schema }); -export const initDrizzle = (params?: { maxConnections?: number }) => { +export const initDrizzle = (params?: { maxConnections?: number, replica?: boolean }) => { const maxConnections = params?.maxConnections || 10; - const client = postgres(process.env.DATABASE_URL!, { + const dbUrl = (params?.replica ? process.env.DATABASE_REPLICA_URL : process.env.DATABASE_URL) ?? ""; + const client = postgres(dbUrl, { max: maxConnections, }); diff --git a/server/src/utils/checkUtils/checkCustomerCorrect.ts b/server/src/utils/checkUtils/checkCustomerCorrect.ts index adc043e0b..05220d01c 100644 --- a/server/src/utils/checkUtils/checkCustomerCorrect.ts +++ b/server/src/utils/checkUtils/checkCustomerCorrect.ts @@ -34,8 +34,7 @@ import { isOneOffPrice, } from "@server/internal/products/prices/priceUtils/usagePriceUtils/classifyUsagePrice"; import { - isFreeProduct, - isOneOff, + isFreeProduct } from "@server/internal/products/productUtils"; import type Stripe from "stripe"; import { formatUnixToDateTime, nullish } from "../genUtils"; @@ -244,38 +243,38 @@ const compareActualItems = async ({ } }; -// If all cus products are free, then should have no sub -const checkAllFreeProducts = async ({ - db, - fullCus, - subs, -}: { - db: DrizzleCli; - fullCus: FullCustomer; - subs: Stripe.Subscription[]; -}) => { - const cusProducts = fullCus.customer_products; - const allFreeOrOneOff = cusProducts.every((cp) => { - const product = cusProductToProduct({ cusProduct: cp }); - return isFreeProduct(product.prices) || isOneOff(product.prices); - }); +// // If all cus products are free, then should have no sub +// const checkAllFreeProducts = async ({ +// db, +// fullCus, +// subs, +// }: { +// db: DrizzleCli; +// fullCus: FullCustomer; +// subs: Stripe.Subscription[]; +// }) => { +// const cusProducts = fullCus.customer_products; +// const allFreeOrOneOff = cusProducts.every((cp) => { +// const product = cusProductToProduct({ cusProduct: cp }); +// return isFreeProduct(product.prices) || isOneOff(product.prices); +// }); - if (allFreeOrOneOff) { - // Make sure no subs exist for this customer - const sub = subs.find( - (sub) => - sub.customer === fullCus.processor?.id && - (sub.status === "active" || sub.status === "past_due"), - ); +// if (allFreeOrOneOff) { +// // Make sure no subs exist for this customer +// const sub = subs.find( +// (sub) => +// sub.customer === fullCus.processor?.id && +// (sub.status === "active" || sub.status === "past_due"), +// ); - if (fullCus.org_id === "6bWdIqEuRHBrReXbTb30l9beMFVZ3Ts3") return true; +// if (fullCus.org_id === "6bWdIqEuRHBrReXbTb30l9beMFVZ3Ts3") return true; - assert(!sub, `no sub should exist for this customer`); - return true; - } +// assert(!sub, `no sub should exist for this customer`); +// return true; +// } - return false; -}; +// return false; +// }; export const checkCusSubCorrect = async ({ db, @@ -292,12 +291,12 @@ export const checkCusSubCorrect = async ({ org: Organization; env: AppEnv; }) => { - const allFree = await checkAllFreeProducts({ - db, - fullCus, - subs, - }); - if (allFree) return; + // const allFree = await checkAllFreeProducts({ + // db, + // fullCus, + // subs, + // }); + // if (allFree) return; // 1. Only 1 sub ID available const cusProducts = fullCus.customer_products; diff --git a/server/src/utils/checkUtils/checkCustomerState.ts b/server/src/utils/checkUtils/checkCustomerState.ts index bf4fda1ed..7a39f7e88 100644 --- a/server/src/utils/checkUtils/checkCustomerState.ts +++ b/server/src/utils/checkUtils/checkCustomerState.ts @@ -1,8 +1,12 @@ import type { AppEnv, FullCustomer, Organization } from "@autumn/shared"; import type { DrizzleCli } from "@server/db/initDrizzle.js"; import type Stripe from "stripe"; +import type { AutumnContext } from "../../honoUtils/HonoEnv.js"; import { checkCusProducts } from "./checkCusProducts.js"; -import { checkCusSubCorrect, SubItemMismatchError } from "./checkCustomerCorrect.js"; +import { + checkCusSubCorrect, + SubItemMismatchError, +} from "./checkCustomerCorrect.js"; import { checkSubCountMatch } from "./checkSubCountMatch.js"; import { saveCheckState } from "./saveCheckState.js"; import type { StateCheckResult } from "./stateCheckTypes.js"; @@ -15,20 +19,21 @@ export type { RedisChecksState } from "./stateCheckTypes.js"; * Does not throw - captures all errors in the result object. */ export const runCustomerStateChecks = async ({ - db, + // db, + ctx, fullCus, subs, schedules, - org, - env, + // org, + // env, }: { - db: DrizzleCli; + // db: DrizzleCli; + ctx: AutumnContext; fullCus: FullCustomer; subs: Stripe.Subscription[]; schedules: Stripe.SubscriptionSchedule[]; - org: Organization; - env: AppEnv; }): Promise => { + const { db, org, env } = ctx; const result: StateCheckResult = { passed: true, errors: [], @@ -36,7 +41,6 @@ export const runCustomerStateChecks = async ({ checks: [], }; - // Check: Subscription correctness (from checkCusSubCorrect) await testSubscriptionCorrectness({ db, @@ -56,6 +60,7 @@ export const runCustomerStateChecks = async ({ // Check: Subscription IDs match Stripe await checkSubCountMatch({ + ctx, fullCus, subs, result, diff --git a/server/src/utils/checkUtils/checkSubCountMatch.ts b/server/src/utils/checkUtils/checkSubCountMatch.ts index 6e41800ac..378c464f2 100644 --- a/server/src/utils/checkUtils/checkSubCountMatch.ts +++ b/server/src/utils/checkUtils/checkSubCountMatch.ts @@ -1,12 +1,17 @@ import type { FullCustomer } from "@autumn/shared"; +import { metadata } from "@autumn/shared"; +import { count, inArray } from "drizzle-orm"; import type Stripe from "stripe"; +import type { AutumnContext } from "@/honoUtils/HonoEnv"; import type { StateCheckResult } from "./stateCheckTypes"; export const checkSubCountMatch = async ({ + ctx, fullCus, subs, result, }: { + ctx: AutumnContext; fullCus: FullCustomer; subs: Stripe.Subscription[]; result: StateCheckResult; @@ -15,16 +20,41 @@ export const checkSubCountMatch = async ({ const subIds = [...new Set(cusProducts.flatMap((cp) => cp.subscription_ids))]; const stripeSubs = subs.filter((sub) => { + if (sub.status === "incomplete") return false; const subCustomerId = typeof sub.customer === "string" ? sub.customer : sub.customer?.id; return subCustomerId === fullCus.processor?.id; }); if (stripeSubs.length !== subIds.length) { - result.passed = false; - result.errors.push( - `Expected ${subIds.length} subs in total, found ${stripeSubs.length} in Stripe`, - ); + // 1. Check if there are any invoice metadata's waiting to be processed... + // const invoiceMetadata = await MetadataService.getByStripeInvoiceId({ + // db: ctx.db, + // stripeInvoiceId: stripeSubs.map((sub) => sub.latest_invoice), + // type: MetadataType.InvoiceCheckout, + // }); + const invoiceIds = await stripeSubs.map((sub) => sub.latest_invoice); + // Count of rows in metadata table with stripe_invoice_id in invoiceIds + const metadataCount = await ctx.db + .select({ count: count() }) + .from(metadata) + .where(inArray(metadata.stripe_invoice_id, invoiceIds as string[])); + + if (stripeSubs.length - (metadataCount?.[0]?.count || 0) !== subIds.length) { + result.passed = false; + const errorMsg = `Expected ${subIds.length} subs in total, found ${stripeSubs.length} in Stripe`; + result.errors.push(errorMsg); + result.checks.push({ + name: "Subscription Count Match", + type: "sub_count_match", + passed: false, + message: errorMsg, + }); + } + // result.passed = false; + // result.errors.push( + // `Expected ${subIds.length} subs in total, found ${stripeSubs.length} in Stripe`, + // ); } else { result.checks.push({ name: `Subscription Match`, diff --git a/server/src/utils/checkUtils/saveCheckState.ts b/server/src/utils/checkUtils/saveCheckState.ts index 7609850b1..f86304bfd 100644 --- a/server/src/utils/checkUtils/saveCheckState.ts +++ b/server/src/utils/checkUtils/saveCheckState.ts @@ -24,6 +24,7 @@ export const saveCheckState = async ({ if (newFailedChecks.length === 0) return; const stateKey = `state:${org.id}:${env}:${fullCus.internal_id}`; + const existingState = (await upstash.get(stateKey)) as RedisChecksState | null; if (existingState) { @@ -57,17 +58,14 @@ export const saveCheckState = async ({ } } - // Status transition: "new" → "ongoing" if no changes but still failing - const noChanges = checksToRemove.size === 0 && checksToAdd.size === 0; - const newStatus = - noChanges && existingState.status === "new" ? "ongoing" : existingState.status; + // Keep existing status - all status changes should be done through the dashboard const updatedState: RedisChecksState = { ...existingState, - status: newStatus, checks: updatedChecks, }; await upstash.set(stateKey, JSON.stringify(updatedState)); + } else { // Create new state const newState: RedisChecksState = { @@ -89,6 +87,7 @@ export const saveCheckState = async ({ })), }; await upstash.set(stateKey, JSON.stringify(newState)); + } }; diff --git a/server/src/utils/checkUtils/stateCheckTypes.ts b/server/src/utils/checkUtils/stateCheckTypes.ts index 0a7134093..9da007cd7 100644 --- a/server/src/utils/checkUtils/stateCheckTypes.ts +++ b/server/src/utils/checkUtils/stateCheckTypes.ts @@ -17,6 +17,7 @@ export type StateCheckResult = { | "subscription_correctness" | "customer_product_correctness" | "sub_id_matching" + | "sub_count_match" | "group_uniqueness" | "entitlement_price_correctness" | "overall_status";