+
+
+
+
-
-
-
-
+
+
+ {buttons.map((button, index) => (
+
+ handleBtnClicked({
+ featureId: button.feature_id,
+ value: button.value,
+ })
+ }
+ variant="gradientPrimary"
+ >
+ {button.text}
+
+ ))}
+
+
Billing
+
+ {customer?.name ? customer?.name : customerId}, you have access to:
+
+
+
Pricing
+
+
+
+
+ Make a test purchase to see how Autumn handles it. Use the Stripe
+ test card 4242 4242 4242 4242 {" "}
+ with any expiration date, CVC and cardholder details.
+
+
+
);
}
-
-const APIPlayground = ({
- title,
- endpoint,
- request,
- response,
- loading,
-}: {
- title: string;
- endpoint: string;
- request: any;
- response: any;
- loading: boolean;
-}) => {
- return (
-
-
-
{title}
-
- {endpoint}
-
-
-
-
-
Request
-
- {JSON.stringify(request, null, 2)}
-
-
-
-
-
Response
-
- {loading ? (
-
- ) : response === null ? (
- "No response"
- ) : (
-
- )}
-
-
-
- );
-};
diff --git a/frontend/src/views/demo/autumnBackend.tsx b/frontend/src/views/demo/autumnBackend.tsx
new file mode 100644
index 000000000..ef31c02b4
--- /dev/null
+++ b/frontend/src/views/demo/autumnBackend.tsx
@@ -0,0 +1,49 @@
+import axios, { AxiosInstance } from "axios";
+
+export const createAxiosInstance = (secretKey: string, endpoint: string) => {
+ return axios.create({
+ baseURL: endpoint + "/v1",
+ headers: {
+ Authorization: `Bearer ${secretKey}`,
+ },
+ });
+};
+//Check access to Pro features and email balance
+export const checkAccess = async ({
+ axiosInstance,
+ customerId,
+ featureId,
+}: {
+ axiosInstance: AxiosInstance;
+ customerId: string;
+ featureId: string;
+}) => {
+ const { data } = await axiosInstance.post("/entitled", {
+ customer_id: customerId,
+ feature_id: featureId,
+ });
+ return data;
+};
+
+//Send usage event for
+export const sendUsage = async ({
+ axiosInstance,
+ customerId,
+ featureId,
+ value,
+}: {
+ axiosInstance: AxiosInstance;
+ customerId: string;
+ featureId: string;
+ value: number;
+}) => {
+ const { data } = await axiosInstance.post("/events", {
+ customer_id: customerId,
+ event_name: featureId,
+ properties: {
+ value: value,
+ },
+ });
+
+ return data;
+};
diff --git a/package-lock.json b/package-lock.json
index 744a708cc..92c85de36 100644
--- a/package-lock.json
+++ b/package-lock.json
@@ -9307,6 +9307,12 @@
"typescript": "^5.7.3"
}
},
+ "node_modules/@ioredis/commands": {
+ "version": "1.2.0",
+ "resolved": "https://registry.npmjs.org/@ioredis/commands/-/commands-1.2.0.tgz",
+ "integrity": "sha512-Sx1pU8EM64o2BrqNpEO1CNLtKQwyhuXuqyfH7oGKCk+1a33d2r5saW8zNwm3j6BTExtjrv2BxTgzzkMwts6vGg==",
+ "license": "MIT"
+ },
"node_modules/@isaacs/cliui": {
"version": "8.0.2",
"resolved": "https://registry.npmjs.org/@isaacs/cliui/-/cliui-8.0.2.tgz",
@@ -11112,9 +11118,9 @@
"license": "MIT"
},
"node_modules/@types/node": {
- "version": "22.13.0",
- "resolved": "https://registry.npmjs.org/@types/node/-/node-22.13.0.tgz",
- "integrity": "sha512-ClIbNe36lawluuvq3+YYhnIN2CELi+6q8NpnM7PYp4hBn/TatfboPgVSm2rwKRfnV2M+Ty9GWDFI64KEe+kysA==",
+ "version": "22.13.4",
+ "resolved": "https://registry.npmjs.org/@types/node/-/node-22.13.4.tgz",
+ "integrity": "sha512-ywP2X0DYtX3y08eFVx5fNIw7/uIv8hYUKgXoK8oayJlLnKcRfEYCxWMVE1XagUdVtCJlZT1AU4LXEABW+L1Peg==",
"license": "MIT",
"dependencies": {
"undici-types": "~6.20.0"
@@ -12035,6 +12041,15 @@
"node": ">=6"
}
},
+ "node_modules/cluster-key-slot": {
+ "version": "1.1.2",
+ "resolved": "https://registry.npmjs.org/cluster-key-slot/-/cluster-key-slot-1.1.2.tgz",
+ "integrity": "sha512-RMr0FhtfXemyinomL4hrWcYJxmX6deFdCxpJzhDttxgO1+bcCnkk+9drydLVDmAMG7NE6aN/fl4F7ucU/90gAA==",
+ "license": "Apache-2.0",
+ "engines": {
+ "node": ">=0.10.0"
+ }
+ },
"node_modules/cmdk": {
"version": "1.0.0",
"resolved": "https://registry.npmjs.org/cmdk/-/cmdk-1.0.0.tgz",
@@ -12792,6 +12807,15 @@
"node": ">=0.4.0"
}
},
+ "node_modules/denque": {
+ "version": "2.1.0",
+ "resolved": "https://registry.npmjs.org/denque/-/denque-2.1.0.tgz",
+ "integrity": "sha512-HVQE3AAb/pxF8fQAoiqpvg9i3evqug3hoiwakOyZAwJm+6vZehbkYXZ0l4JxS+I3QxM97v5aaRNhj8v5oBhekw==",
+ "license": "Apache-2.0",
+ "engines": {
+ "node": ">=0.10"
+ }
+ },
"node_modules/dequal": {
"version": "2.0.3",
"resolved": "https://registry.npmjs.org/dequal/-/dequal-2.0.3.tgz",
@@ -14278,6 +14302,30 @@
"url": "https://github.com/sponsors/colinhacks"
}
},
+ "node_modules/ioredis": {
+ "version": "5.5.0",
+ "resolved": "https://registry.npmjs.org/ioredis/-/ioredis-5.5.0.tgz",
+ "integrity": "sha512-7CutT89g23FfSa8MDoIFs2GYYa0PaNiW/OrT+nRyjRXHDZd17HmIgy+reOQ/yhh72NznNjGuS8kbCAcA4Ro4mw==",
+ "license": "MIT",
+ "dependencies": {
+ "@ioredis/commands": "^1.1.1",
+ "cluster-key-slot": "^1.1.0",
+ "debug": "^4.3.4",
+ "denque": "^2.1.0",
+ "lodash.defaults": "^4.2.0",
+ "lodash.isarguments": "^3.1.0",
+ "redis-errors": "^1.2.0",
+ "redis-parser": "^3.0.0",
+ "standard-as-callback": "^2.1.0"
+ },
+ "engines": {
+ "node": ">=12.22.0"
+ },
+ "funding": {
+ "type": "opencollective",
+ "url": "https://opencollective.com/ioredis"
+ }
+ },
"node_modules/ip-address": {
"version": "9.0.5",
"resolved": "https://registry.npmjs.org/ip-address/-/ip-address-9.0.5.tgz",
@@ -15113,6 +15161,18 @@
"integrity": "sha512-TwuEnCnxbc3rAvhf/LbG7tJUDzhqXyFnv3dtzLOPgCG/hODL7WFnsbwktkD7yUV0RrreP/l1PALq/YSg6VvjlA==",
"license": "MIT"
},
+ "node_modules/lodash.defaults": {
+ "version": "4.2.0",
+ "resolved": "https://registry.npmjs.org/lodash.defaults/-/lodash.defaults-4.2.0.tgz",
+ "integrity": "sha512-qjxPLHd3r5DnsdGacqOMU6pb/avJzdh9tFX2ymgoZE27BmjXrNy/y4LoaiTeAb+O3gL8AfpJGtqfX/ae2leYYQ==",
+ "license": "MIT"
+ },
+ "node_modules/lodash.isarguments": {
+ "version": "3.1.0",
+ "resolved": "https://registry.npmjs.org/lodash.isarguments/-/lodash.isarguments-3.1.0.tgz",
+ "integrity": "sha512-chi4NHZlZqZD18a0imDHnZPrDeBbTtVN7GXMwuGdRH9qotxAjYs3aVLKc7zNOG9eddR5Ksd8rvFEBc9SsggPpg==",
+ "license": "MIT"
+ },
"node_modules/lodash.merge": {
"version": "4.6.2",
"resolved": "https://registry.npmjs.org/lodash.merge/-/lodash.merge-4.6.2.tgz",
@@ -17713,6 +17773,27 @@
"node": ">= 12.13.0"
}
},
+ "node_modules/redis-errors": {
+ "version": "1.2.0",
+ "resolved": "https://registry.npmjs.org/redis-errors/-/redis-errors-1.2.0.tgz",
+ "integrity": "sha512-1qny3OExCf0UvUV/5wpYKf2YwPcOqXzkwKKSmKHiE6ZMQs5heeE/c8eXK+PNllPvmjgAbfnsbpkGZWy8cBpn9w==",
+ "license": "MIT",
+ "engines": {
+ "node": ">=4"
+ }
+ },
+ "node_modules/redis-parser": {
+ "version": "3.0.0",
+ "resolved": "https://registry.npmjs.org/redis-parser/-/redis-parser-3.0.0.tgz",
+ "integrity": "sha512-DJnGAeenTdpMEH6uAJRK/uiyEIH9WVsUmoLwzudwGJUwZPp80PDBWPHXSAGNPwNvIXAbe7MSUB1zQFugFml66A==",
+ "license": "MIT",
+ "dependencies": {
+ "redis-errors": "^1.0.0"
+ },
+ "engines": {
+ "node": ">=4"
+ }
+ },
"node_modules/refractor": {
"version": "3.6.0",
"resolved": "https://registry.npmjs.org/refractor/-/refractor-3.6.0.tgz",
@@ -18437,6 +18518,12 @@
"dev": true,
"license": "BSD-3-Clause"
},
+ "node_modules/standard-as-callback": {
+ "version": "2.1.0",
+ "resolved": "https://registry.npmjs.org/standard-as-callback/-/standard-as-callback-2.1.0.tgz",
+ "integrity": "sha512-qoRRSyROncaz1z0mvYqIE4lCd9p2R90i6GxW3uZv5ucSu8tU7B5HXUP1gG8pVZsYNVaXjk8ClXHPttLyxAL48A==",
+ "license": "MIT"
+ },
"node_modules/std-env": {
"version": "3.8.0",
"resolved": "https://registry.npmjs.org/std-env/-/std-env-3.8.0.tgz",
@@ -19592,6 +19679,7 @@
"express": "^4.21.1",
"http-status-codes": "^2.3.0",
"inngest": "^3.30.0",
+ "ioredis": "^5.5.0",
"ksuid": "^3.0.0",
"nodemon": "^3.1.7",
"openai": "^4.85.2",
@@ -19611,7 +19699,7 @@
"@types/chai": "^5.0.1",
"@types/chai-http": "^3.0.5",
"@types/mocha": "^10.0.10",
- "@types/node": "^22.13.0",
+ "@types/node": "^22.13.4",
"@types/pg": "^8.11.10",
"mocha": "^11.1.0",
"puppeteer": "^24.2.0",
@@ -19873,10 +19961,6 @@
"@types/node": ">=18"
}
},
- "server/node_modules/@ioredis/commands": {
- "version": "1.2.0",
- "license": "MIT"
- },
"server/node_modules/@msgpackr-extract/msgpackr-extract-darwin-arm64": {
"version": "3.0.3",
"cpu": [
@@ -20338,13 +20422,6 @@
"node": ">=8"
}
},
- "server/node_modules/cluster-key-slot": {
- "version": "1.1.2",
- "license": "Apache-2.0",
- "engines": {
- "node": ">=0.10.0"
- }
- },
"server/node_modules/commander": {
"version": "12.1.0",
"license": "MIT",
@@ -20412,13 +20489,6 @@
"node": ">=0.10.0"
}
},
- "server/node_modules/denque": {
- "version": "2.1.0",
- "license": "Apache-2.0",
- "engines": {
- "node": ">=0.10"
- }
- },
"server/node_modules/depd": {
"version": "2.0.0",
"license": "MIT",
@@ -20711,47 +20781,6 @@
"@types/node": ">=18"
}
},
- "server/node_modules/ioredis": {
- "version": "5.4.1",
- "license": "MIT",
- "dependencies": {
- "@ioredis/commands": "^1.1.1",
- "cluster-key-slot": "^1.1.0",
- "debug": "^4.3.4",
- "denque": "^2.1.0",
- "lodash.defaults": "^4.2.0",
- "lodash.isarguments": "^3.1.0",
- "redis-errors": "^1.2.0",
- "redis-parser": "^3.0.0",
- "standard-as-callback": "^2.1.0"
- },
- "engines": {
- "node": ">=12.22.0"
- },
- "funding": {
- "type": "opencollective",
- "url": "https://opencollective.com/ioredis"
- }
- },
- "server/node_modules/ioredis/node_modules/debug": {
- "version": "4.4.0",
- "license": "MIT",
- "dependencies": {
- "ms": "^2.1.3"
- },
- "engines": {
- "node": ">=6.0"
- },
- "peerDependenciesMeta": {
- "supports-color": {
- "optional": true
- }
- }
- },
- "server/node_modules/ioredis/node_modules/ms": {
- "version": "2.1.3",
- "license": "MIT"
- },
"server/node_modules/ipaddr.js": {
"version": "1.9.1",
"license": "MIT",
@@ -20803,14 +20832,6 @@
"node": ">=8"
}
},
- "server/node_modules/lodash.defaults": {
- "version": "4.2.0",
- "license": "MIT"
- },
- "server/node_modules/lodash.isarguments": {
- "version": "3.1.0",
- "license": "MIT"
- },
"server/node_modules/lodash.isstring": {
"version": "4.0.1",
"license": "MIT"
@@ -21309,23 +21330,6 @@
"recase": "dist/cli.js"
}
},
- "server/node_modules/redis-errors": {
- "version": "1.2.0",
- "license": "MIT",
- "engines": {
- "node": ">=4"
- }
- },
- "server/node_modules/redis-parser": {
- "version": "3.0.0",
- "license": "MIT",
- "dependencies": {
- "redis-errors": "^1.0.0"
- },
- "engines": {
- "node": ">=4"
- }
- },
"server/node_modules/require-main-filename": {
"version": "2.0.0",
"license": "ISC"
@@ -21469,10 +21473,6 @@
"node": ">=8"
}
},
- "server/node_modules/standard-as-callback": {
- "version": "2.1.0",
- "license": "MIT"
- },
"server/node_modules/statuses": {
"version": "2.0.1",
"license": "MIT",
diff --git a/server/package.json b/server/package.json
index 21a74e691..2016489a3 100644
--- a/server/package.json
+++ b/server/package.json
@@ -5,7 +5,7 @@
"main": "index.js",
"type": "module",
"scripts": {
- "dev": "nodemon --no-deprecation --exec tsx src/index.ts",
+ "dev": "nodemon --no-deprecation --exec tsx src/index.ts --ignore scripts",
"start": "tsx src/index.ts",
"queue:dev": "tsx watch src/queue.ts",
"build": "tsc -b",
@@ -51,6 +51,7 @@
"express": "^4.21.1",
"http-status-codes": "^2.3.0",
"inngest": "^3.30.0",
+ "ioredis": "^5.5.0",
"ksuid": "^3.0.0",
"nodemon": "^3.1.7",
"openai": "^4.85.2",
@@ -70,7 +71,7 @@
"@types/chai": "^5.0.1",
"@types/chai-http": "^3.0.5",
"@types/mocha": "^10.0.10",
- "@types/node": "^22.13.0",
+ "@types/node": "^22.13.4",
"@types/pg": "^8.11.10",
"mocha": "^11.1.0",
"puppeteer": "^24.2.0",
diff --git a/server/src/errors/logger.ts b/server/src/errors/logger.ts
index 76f5c955c..e2ac66ab0 100644
--- a/server/src/errors/logger.ts
+++ b/server/src/errors/logger.ts
@@ -31,7 +31,8 @@ export const initLogger = () => {
},
});
- logger.info("Logger initialized");
+ // logger.info("Logger initialized");
+ console.log("Logger initialized");
return logger;
};
diff --git a/server/src/index.ts b/server/src/index.ts
index 855a5a966..35d4b689d 100644
--- a/server/src/index.ts
+++ b/server/src/index.ts
@@ -10,28 +10,27 @@ import webhooksRouter from "./external/webhooks/webhooksRouter.js";
import { envMiddleware } from "./middleware/envMiddleware.js";
import pg from "pg";
-import { initQueue, initWorkers } from "./queue/queue.js";
+import { initWorkers } from "./queue/queue.js";
import http from "http";
import { initWs } from "./websockets/initWs.js";
import { publicRouter } from "./internal/public/publicRouter.js";
import { initLogger } from "./errors/logger.js";
-import RecaseError, { handleRequestError } from "./utils/errorUtils.js";
+import { QueueManager } from "./queue/QueueManager.js";
const init = async () => {
const app = express();
- const server = http.createServer(app);
- const wss = initWs(server);
+
const logger = initLogger();
+ const server = http.createServer(app);
const pgClient = new pg.Client(process.env.SUPABASE_CONNECTION_STRING || "");
await pgClient.connect();
- const queue = initQueue();
- const workers = initWorkers(queue);
+ await QueueManager.getInstance(); // initialize the queue manager
+ await initWorkers();
app.use((req: any, res, next) => {
req.pg = pgClient;
- req.queue = queue;
req.logger = logger;
next();
});
diff --git a/server/src/internal/api/events/eventRouter.ts b/server/src/internal/api/events/eventRouter.ts
index c9dd4856a..09f9b1f4d 100644
--- a/server/src/internal/api/events/eventRouter.ts
+++ b/server/src/internal/api/events/eventRouter.ts
@@ -6,7 +6,7 @@ import {
Event,
Feature,
} from "@autumn/shared";
-import { handleRequestError } from "@/utils/errorUtils.js";
+import RecaseError, { handleRequestError } from "@/utils/errorUtils.js";
import { generateId } from "@/utils/genUtils.js";
import { EventService } from "./EventService.js";
@@ -16,6 +16,7 @@ import { Queue } from "bullmq";
import { createNewCustomer } from "../customers/cusUtils.js";
import { SupabaseClient } from "@supabase/supabase-js";
import { OrgService } from "@/internal/orgs/OrgService.js";
+import { QueueManager } from "@/queue/QueueManager.js";
export const eventsRouter = Router();
@@ -147,15 +148,39 @@ export const handleEventSent = async ({
});
if (affectedFeatures.length > 0) {
- let queue: Queue = req.queue;
- queue.add("update-balance", {
+ // let queue: Queue = req.queue;
+
+ const payload = {
customerId: customer.internal_id,
customer,
features: affectedFeatures,
event,
org,
env,
- });
+ };
+
+ const queue = await QueueManager.getQueue({ useBackup: false });
+
+ try {
+ // Add timeout to queue operation
+ await queue.add("update-balance", payload);
+ // console.log("Added update-balance to queue");
+ } catch (error: any) {
+ try {
+ console.log("Adding update-balance to backup queue");
+ const backupQueue = await QueueManager.getQueue({ useBackup: true });
+ await backupQueue.add("update-balance", payload);
+ } catch (error: any) {
+ throw new RecaseError({
+ message: "Failed to add update-balance to queue (backup)",
+ code: "EVENT_QUEUE_ERROR",
+ statusCode: 500,
+ data: {
+ message: error.message,
+ },
+ });
+ }
+ }
}
};
diff --git a/server/src/internal/dev/devRouter.ts b/server/src/internal/dev/devRouter.ts
index 8cceacc2e..8d776f51f 100644
--- a/server/src/internal/dev/devRouter.ts
+++ b/server/src/internal/dev/devRouter.ts
@@ -5,9 +5,43 @@ import { Router } from "express";
import { ApiKeyService } from "./ApiKeyService.js";
import { OrgService } from "../orgs/OrgService.js";
import { createKey } from "./api-keys/apiKeyUtils.js";
+import { SupabaseClient } from "@supabase/supabase-js";
export const devRouter = Router();
+export const handleCreateApiKey = async ({
+ sb,
+ env,
+ name,
+ orgId,
+ orgSlug,
+}: {
+ sb: SupabaseClient;
+ env: AppEnv;
+ name: string;
+ orgId: string;
+ orgSlug: string;
+}) => {
+ // 1. Create API key
+ let prefix = "am_sk_test";
+ if (env === AppEnv.Live) {
+ prefix = "am_sk_live";
+ }
+
+ const apiKey = await createKey({
+ sb,
+ env,
+ name,
+ orgId,
+ prefix,
+ meta: {
+ org_slug: orgSlug,
+ },
+ });
+
+ return apiKey;
+};
+
devRouter.get("/data", withOrgAuth, async (req: any, res) => {
const apiKeys = await ApiKeyService.getByOrg(req.sb, req.orgId, req.env);
const org = await OrgService.getFullOrg({ sb: req.sb, orgId: req.orgId });
diff --git a/server/src/internal/orgs/OrgService.ts b/server/src/internal/orgs/OrgService.ts
index 10f265708..8621cff16 100644
--- a/server/src/internal/orgs/OrgService.ts
+++ b/server/src/internal/orgs/OrgService.ts
@@ -86,6 +86,10 @@ export class OrgService {
.single();
if (error) {
+ if (error.code === "PGRST116") {
+ return null;
+ }
+
throw new RecaseError({
message: "Failed to get org from supabase",
code: ErrCode.OrgNotFound,
diff --git a/server/src/internal/prices/priceUtils.ts b/server/src/internal/prices/priceUtils.ts
index baadbd342..368ed7a7c 100644
--- a/server/src/internal/prices/priceUtils.ts
+++ b/server/src/internal/prices/priceUtils.ts
@@ -304,7 +304,7 @@ export const roundPriceAmounts = (price: Price) => {
const config = price.config as UsagePriceConfig;
for (let i = 0; i < config.usage_tiers.length; i++) {
config.usage_tiers[i].amount = Number(
- config.usage_tiers[i].amount.toFixed(2)
+ config.usage_tiers[i].amount.toFixed(10)
);
}
diff --git a/server/src/middleware/apiMiddleware.ts b/server/src/middleware/apiMiddleware.ts
index cf880eb89..922a424ce 100644
--- a/server/src/middleware/apiMiddleware.ts
+++ b/server/src/middleware/apiMiddleware.ts
@@ -44,19 +44,20 @@ export const apiAuthMiddleware = async (req: any, res: any, next: any) => {
console.log(`Autumn API verification failed`);
}
} catch (error) {
- console.log("Failed to fetch key from Autumn");
+ console.log("Error: Failed to fetch key from Autumn");
+ console.log("Error: ", error);
}
// Fallback: Verify via Unkey
try {
const result = await validateApiKey(apiKey);
- await migrateKey({
- sb: req.sb,
- keyId: result.keyId ?? "",
- meta: { org_slug: result.meta?.org_slug },
- apiKey,
- });
+ // await migrateKey({
+ // sb: req.sb,
+ // keyId: result.keyId ?? "",
+ // meta: { org_slug: result.meta?.org_slug },
+ // apiKey,
+ // });
// console.log(`Unkey verification successul for ${result.meta?.org_slug}`);
req.orgId = result.ownerId;
@@ -68,7 +69,8 @@ export const apiAuthMiddleware = async (req: any, res: any, next: any) => {
next();
} catch (error) {
- console.log("Unkey API verification failed");
+ console.log("WARNING: Unkey API verification failed");
+ console.log(error);
withOrgAuth(req, res, next);
return;
}
diff --git a/server/src/queue/QueueManager.ts b/server/src/queue/QueueManager.ts
new file mode 100644
index 000000000..e93644bf3
--- /dev/null
+++ b/server/src/queue/QueueManager.ts
@@ -0,0 +1,161 @@
+import { Queue } from "bullmq";
+import { getRedisConnection } from "./queue.js";
+import { Redis } from "ioredis";
+
+export class QueueManager {
+ private static instance: QueueManager;
+ private queue: Queue | null = null;
+ private backupQueue: Queue | null = null;
+
+ private mainConnection: Redis | null = null;
+ private backupConnection: Redis | null = null;
+
+ private constructor() {
+ this.initializePromise = this.initQueue();
+ }
+
+ private initializePromise: Promise
;
+ public static async getInstance(): Promise {
+ if (!QueueManager.instance) {
+ QueueManager.instance = new QueueManager();
+ }
+ // Wait for initialization to complete
+ await QueueManager.instance.initializePromise;
+ return QueueManager.instance;
+ }
+
+ // 1. Create main redis connection
+ private async pingRedis({
+ useBackup,
+ keepConnection = false,
+ }: {
+ useBackup: boolean;
+ keepConnection?: boolean;
+ }) {
+ // 1. Connect to redis
+
+ const redisUrl = useBackup
+ ? process.env.REDIS_BACKUP_URL
+ : process.env.REDIS_URL;
+
+ const connection = new Redis(redisUrl!, {
+ retryStrategy: (times) => {
+ return 5000;
+ },
+ });
+
+ connection.on("error", (error) => {
+ console.log(
+ `Redis connection error (${useBackup ? "backup" : "main"}): ${
+ error.message
+ }`
+ );
+
+ if (!keepConnection) {
+ process.exit(1);
+ }
+ });
+
+ // Check if connection is live...
+ await connection.ping();
+
+ if (!keepConnection) {
+ await connection.quit();
+ }
+ return connection;
+ }
+
+ private async createConnections() {
+ console.log("2. Creating redis connections (for workers...)");
+
+ this.mainConnection = await this.pingRedis({
+ useBackup: false,
+ keepConnection: true,
+ });
+ this.backupConnection = await this.pingRedis({
+ useBackup: true,
+ keepConnection: true,
+ });
+ }
+
+ private async initQueue() {
+ console.log("Initializing Queue Manager...");
+ console.group();
+ // 1. Create redis connections
+ console.log("1. Pinging main & backup redis");
+ this.mainConnection = await this.pingRedis({ useBackup: false });
+ this.backupConnection = await this.pingRedis({ useBackup: true });
+
+ await this.createConnections();
+ // 2. Initialize main and backup queues
+ console.log("2. Initializing main & backup queues");
+ const mainQueue = new Queue("autumn", {
+ connection: {
+ url: process.env.REDIS_URL,
+ enableOfflineQueue: false,
+ retryStrategy: (times) => {
+ return 5000;
+ },
+ },
+ });
+
+ const backupQueue = new Queue("autumn", {
+ connection: {
+ url: process.env.REDIS_BACKUP_URL,
+ enableOfflineQueue: false,
+ },
+ });
+
+ // Set up error handling for the queue
+ mainQueue.on("error", async (error: any) => {
+ console.error("QUEUE ERROR:", error.message);
+ if (error.code !== "ECONNREFUSED") {
+ }
+ });
+
+ backupQueue.on("error", async (error: any) => {
+ console.error("BACKUP QUEUE ERROR:", error.message);
+ if (error.code !== "ECONNREFUSED") {
+ }
+ });
+
+ this.queue = mainQueue;
+ this.backupQueue = backupQueue;
+ console.groupEnd();
+ }
+
+ // Create workers
+
+ public static async getQueue({
+ useBackup,
+ }: {
+ useBackup: boolean;
+ }): Promise {
+ const queueManager = await QueueManager.getInstance();
+ if (!queueManager.queue || !queueManager.backupQueue) {
+ throw new Error("Queue not initialized");
+ }
+ return useBackup ? queueManager.backupQueue : queueManager.queue;
+ }
+
+ public static async getConnection({
+ useBackup,
+ }: {
+ useBackup: boolean;
+ }): Promise {
+ const queueManager = await QueueManager.getInstance();
+ if (!queueManager.mainConnection || !queueManager.backupConnection) {
+ throw new Error("Connection not initialized");
+ }
+ return useBackup
+ ? queueManager.backupConnection
+ : queueManager.mainConnection;
+ }
+
+ public getBackupConnection(): Redis {
+ if (!this.backupConnection) {
+ throw new Error("Backup connection not initialized");
+ }
+ return this.backupConnection;
+ }
+}
diff --git a/server/src/queue/queue.ts b/server/src/queue/queue.ts
index 960cd2ac5..15b32cfd6 100644
--- a/server/src/queue/queue.ts
+++ b/server/src/queue/queue.ts
@@ -1,54 +1,74 @@
import { Job, Queue, Worker } from "bullmq";
import { runUpdateBalanceTask } from "@/trigger/updateBalanceTask.js";
-import { Redis } from "ioredis";
+import { QueueManager } from "./QueueManager.js";
-const redis = new Redis(process.env.REDIS_URL || "redis://localhost:6379");
+const NUM_WORKERS = 5;
+
+export const getRedisConnection = ({
+ useBackup = false,
+}: {
+ useBackup?: boolean;
+}) => {
+ let redisUrl = process.env.REDIS_URL || "redis://localhost:6379";
+
+ if (useBackup) {
+ redisUrl = process.env.REDIS_BACKUP_URL || "redis://localhost:6379";
+ }
+
+ return {
+ connection: {
+ url: redisUrl,
+ // enableOfflineQueue: false,
+ },
+ };
+};
+
+async function acquireLock({
+ customerId,
+ timeout = 30000,
+ useBackup = false,
+}: {
+ customerId: string;
+ timeout?: number;
+ useBackup?: boolean;
+}): Promise {
+ // const redis = getRedisClient({ useBackup });
+ const redis = await QueueManager.getConnection({ useBackup });
-async function acquireLock(
- customerId: string,
- timeout = 30000
-): Promise {
const lockKey = `lock:customer:${customerId}`;
const acquired = await redis.set(lockKey, "1", "PX", timeout, "NX");
return acquired === "OK";
}
-async function releaseLock(customerId: string): Promise {
+async function releaseLock({
+ customerId,
+ useBackup,
+}: {
+ customerId: string;
+ useBackup: boolean;
+}): Promise {
+ const redis = await QueueManager.getConnection({ useBackup });
const lockKey = `lock:customer:${customerId}`;
await redis.del(lockKey);
}
-const getRedisConnection = () => {
- let redisUrl = process.env.REDIS_URL || "redis://localhost:6379";
-
- return {
- connection: {
- url: redisUrl,
- },
- };
-};
-
-export const initQueue = () => {
- try {
- return new Queue("autumn", getRedisConnection());
- } catch (error) {
- console.error("Error initialising queue:\n", error);
- process.exit(1);
- }
-};
-
-const numWorkers = 5;
-
-const initWorker = (id: number, queue: Queue) => {
- // Create supabase client
-
+const initWorker = ({
+ id,
+ queue,
+ useBackup,
+}: {
+ id: number;
+ queue: Queue;
+ useBackup: boolean;
+}) => {
let worker = new Worker(
"autumn",
async (job: Job) => {
+ // console.log("JOB ID:", job.id, `(${useBackup ? "BACKUP" : "MAIN"})`);
+ // console.log("EVENT ID:", job.data.event.id);
const { customerId } = job.data;
- while (!(await acquireLock(customerId, 10000))) {
- // console.log(`Customer ${customer.id} locked by another worker`);
+ while (!(await acquireLock({ customerId, timeout: 10000, useBackup }))) {
await queue.add(job.name, job.data, {
delay: 50,
});
@@ -60,37 +80,57 @@ const initWorker = (id: number, queue: Queue) => {
} catch (error) {
console.error("Error updating balance:", error);
} finally {
- await releaseLock(customerId);
+ await releaseLock({ customerId, useBackup });
}
},
-
{
- ...getRedisConnection(),
- concurrency: 3,
+ ...getRedisConnection({ useBackup }),
+ concurrency: 1,
+ removeOnComplete: {
+ count: 0,
+ },
+ removeOnFail: {
+ count: 0,
+ },
+ drainDelay: 1000,
+ maxStalledCount: 0,
}
);
worker.on("ready", () => {
- console.log(`Worker ${id} ready`);
+ console.log(`Worker ${id} ready (${useBackup ? "BACKUP" : "MAIN"})`);
});
- worker.on("error", async (error) => {
- console.log("WORKER ERROR:\n");
- console.log(error);
- await new Promise((resolve) => setTimeout(resolve, 5000));
+ worker.on("stalled", (jobId: string) => {
+ console.log(`Worker ${id} stalled (${useBackup ? "BACKUP" : "MAIN"})`);
+ console.log("JOB ID:", jobId);
+ });
+
+ // Check jobs left in queue
+
+ worker.on("error", async (error: any) => {
+ if (error.code !== "ECONNREFUSED") {
+ console.log("WORKER ERROR:", error.message);
+ }
});
worker.on("failed", (job, error) => {
- console.log("WORKER FAILED:\n");
- console.log(error);
+ console.log("WORKER FAILED:", error.message);
});
};
-export const initWorkers = (queue: Queue) => {
+export const initWorkers = async () => {
const workers = [];
- for (let i = 0; i < numWorkers; i++) {
- workers.push(initWorker(i, queue));
+
+ const mainQueue = await QueueManager.getQueue({ useBackup: false });
+ const backupQueue = await QueueManager.getQueue({ useBackup: true });
+
+ for (let i = 0; i < NUM_WORKERS; i++) {
+ workers.push(initWorker({ id: i, queue: mainQueue, useBackup: false }));
+ workers.push(initWorker({ id: i, queue: backupQueue, useBackup: true }));
}
+ // Get stalled jobs
+
return workers;
};
diff --git a/server/src/trigger/updateBalanceTask.ts b/server/src/trigger/updateBalanceTask.ts
index 0fb4c79ba..4b61e56ea 100644
--- a/server/src/trigger/updateBalanceTask.ts
+++ b/server/src/trigger/updateBalanceTask.ts
@@ -256,7 +256,9 @@ export const updateCustomerBalance = async ({
},
});
} else {
- console.log("No usage-based entitlement found");
+ console.log(
+ ` - Remaining deduction: ${toDeduct}, no usage-based entitlement found`
+ );
}
}
@@ -272,7 +274,9 @@ export const runUpdateBalanceTask = async (payload: any) => {
const { customer, features, event, org, env } = payload;
console.log("--------------------------------");
- console.log("Inside updateBalanceTask...");
+ console.log(
+ `UPDATING BALANCE FOR CUSTOMER (${customer.id}), ORG: ${org.slug}`
+ );
console.log("1. Updating customer balance...");
const cusEnts: any = await updateCustomerBalance({
diff --git a/server/src/utils/genUtils.ts b/server/src/utils/genUtils.ts
index 7a56b4fb2..39649f2b8 100644
--- a/server/src/utils/genUtils.ts
+++ b/server/src/utils/genUtils.ts
@@ -18,7 +18,9 @@ export const compareObjects = (obj1: any, obj2: any) => {
};
export const keyToTitle = (key: string) => {
- return key.replace(/-/g, " ").replace(/\b\w/g, (char) => char.toUpperCase());
+ return key
+ .replace(/[-_]/g, " ")
+ .replace(/\b\w/g, (char) => char.toUpperCase());
};
export const notNullOrUndefined = (value: any) => {