Skip to content

Commit 969fecc

Browse files
Merge pull request #273 from Jayking40/main
feat: implement isolated worker queues for Soroban subscription events
2 parents 492f89d + 89e7b2b commit 969fecc

9 files changed

Lines changed: 1691 additions & 2 deletions

.env.example

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -209,6 +209,50 @@ INACTIVE_RETENTION_YEARS=3
209209
PII_SCRUBBING_ENABLED=true
210210
PII_SCRUBBING_CRON_SCHEDULE=0 2 * * 0 # Weekly on Sunday at 2 AM
211211

212+
# Isolated Soroban Subscription Worker Queues (Issue #230)
213+
# These BullMQ queues run in isolation — each on-chain event type gets its own
214+
# worker so that a billing spike cannot starve payment-failure dunning, etc.
215+
SOROBAN_QUEUE_MAX_ATTEMPTS=5
216+
SOROBAN_QUEUE_BACKOFF_DELAY_MS=2000
217+
SOROBAN_QUEUE_DEFAULT_CONCURRENCY=5
218+
SOROBAN_QUEUE_RETAIN_COMPLETED=200
219+
SOROBAN_QUEUE_RETAIN_FAILED=100
220+
221+
# SubscriptionBilled queue
222+
SOROBAN_QUEUE_BILLED_CONCURRENCY=10
223+
SOROBAN_QUEUE_BILLED_MAX_ATTEMPTS=5
224+
SOROBAN_QUEUE_BILLED_BACKOFF_MS=2000
225+
SOROBAN_QUEUE_BILLED_RATE_MAX=100
226+
SOROBAN_QUEUE_BILLED_RATE_WINDOW_MS=10000
227+
228+
# TrialStarted queue
229+
SOROBAN_QUEUE_TRIAL_CONCURRENCY=8
230+
SOROBAN_QUEUE_TRIAL_MAX_ATTEMPTS=5
231+
SOROBAN_QUEUE_TRIAL_BACKOFF_MS=2000
232+
SOROBAN_QUEUE_TRIAL_RATE_MAX=50
233+
SOROBAN_QUEUE_TRIAL_RATE_WINDOW_MS=10000
234+
235+
# PaymentFailed queue (higher retries, slower backoff for dunning safety)
236+
SOROBAN_QUEUE_PAYMENT_FAILED_CONCURRENCY=4
237+
SOROBAN_QUEUE_PAYMENT_FAILED_MAX_ATTEMPTS=7
238+
SOROBAN_QUEUE_PAYMENT_FAILED_BACKOFF_MS=5000
239+
SOROBAN_QUEUE_PAYMENT_FAILED_RATE_MAX=30
240+
SOROBAN_QUEUE_PAYMENT_FAILED_RATE_WINDOW_MS=10000
241+
242+
# PaymentFailedGracePeriodStarted queue
243+
SOROBAN_QUEUE_GRACE_PERIOD_CONCURRENCY=4
244+
SOROBAN_QUEUE_GRACE_PERIOD_MAX_ATTEMPTS=7
245+
SOROBAN_QUEUE_GRACE_PERIOD_BACKOFF_MS=5000
246+
SOROBAN_QUEUE_GRACE_PERIOD_RATE_MAX=30
247+
SOROBAN_QUEUE_GRACE_PERIOD_RATE_WINDOW_MS=10000
248+
249+
# How often (ms) the subscription worker logs a queue stats snapshot
250+
SOROBAN_WORKER_STATS_INTERVAL_MS=60000
251+
252+
# Set to 'false' to disable the subscription worker pool when running via worker.js
253+
# (useful when running workers/sorobanSubscriptionWorker.js as a separate process)
254+
SOROBAN_SUBSCRIPTION_WORKERS_ENABLED=true
255+
212256
# Sandbox Environment Configuration
213257
SANDBOX_ENABLED=false
214258
SANDBOX_MODE=testnet # testnet or mainnet

package-lock.json

Lines changed: 1 addition & 1 deletion
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

package.json

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,9 @@
1919
"reconciliation:health": "node worker.js --reconciliation --health",
2020
"pii-scrub": "node workers/piiScrubbingWorker.js",
2121
"pii-scrub:dry-run": "node workers/piiScrubbingWorker.js 3 --dry-run",
22+
"subscription-worker": "node workers/sorobanSubscriptionWorker.js",
23+
"subscription-worker:dev": "nodemon workers/sorobanSubscriptionWorker.js",
24+
"subscription-worker:health": "node workers/sorobanSubscriptionWorker.js --health",
2225
"test:soroban": "jest --testPathPattern=soroban",
2326
"test:pii": "jest piiScrubbing.test.js",
2427
"migrate": "node migrations/runMigrations.js",

src/config.js

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -207,6 +207,51 @@ async function loadConfig(env = process.env, vaultService = null) {
207207
emailQueue: env.RABBITMQ_EMAIL_QUEUE || 'substream_emails_queue',
208208
leaderboardQueue: env.RABBITMQ_LEADERBOARD_QUEUE || 'substream_leaderboard_queue',
209209
},
210+
// Isolated BullMQ queue configuration for on-chain Soroban subscription events.
211+
// Each event type gets its own queue with independently-tunable settings.
212+
sorobanQueues: {
213+
defaultMaxAttempts: Number(env.SOROBAN_QUEUE_MAX_ATTEMPTS || 5),
214+
defaultBackoffDelay: Number(env.SOROBAN_QUEUE_BACKOFF_DELAY_MS || 2000),
215+
defaultConcurrency: Number(env.SOROBAN_QUEUE_DEFAULT_CONCURRENCY || 5),
216+
defaultRetainCompleted: Number(env.SOROBAN_QUEUE_RETAIN_COMPLETED || 200),
217+
defaultRetainFailed: Number(env.SOROBAN_QUEUE_RETAIN_FAILED || 100),
218+
SubscriptionBilled: {
219+
concurrency: Number(env.SOROBAN_QUEUE_BILLED_CONCURRENCY || 10),
220+
maxAttempts: Number(env.SOROBAN_QUEUE_BILLED_MAX_ATTEMPTS || 5),
221+
backoffDelay: Number(env.SOROBAN_QUEUE_BILLED_BACKOFF_MS || 2000),
222+
rateLimiter: {
223+
max: Number(env.SOROBAN_QUEUE_BILLED_RATE_MAX || 100),
224+
duration: Number(env.SOROBAN_QUEUE_BILLED_RATE_WINDOW_MS || 10000),
225+
},
226+
},
227+
TrialStarted: {
228+
concurrency: Number(env.SOROBAN_QUEUE_TRIAL_CONCURRENCY || 8),
229+
maxAttempts: Number(env.SOROBAN_QUEUE_TRIAL_MAX_ATTEMPTS || 5),
230+
backoffDelay: Number(env.SOROBAN_QUEUE_TRIAL_BACKOFF_MS || 2000),
231+
rateLimiter: {
232+
max: Number(env.SOROBAN_QUEUE_TRIAL_RATE_MAX || 50),
233+
duration: Number(env.SOROBAN_QUEUE_TRIAL_RATE_WINDOW_MS || 10000),
234+
},
235+
},
236+
PaymentFailed: {
237+
concurrency: Number(env.SOROBAN_QUEUE_PAYMENT_FAILED_CONCURRENCY || 4),
238+
maxAttempts: Number(env.SOROBAN_QUEUE_PAYMENT_FAILED_MAX_ATTEMPTS || 7),
239+
backoffDelay: Number(env.SOROBAN_QUEUE_PAYMENT_FAILED_BACKOFF_MS || 5000),
240+
rateLimiter: {
241+
max: Number(env.SOROBAN_QUEUE_PAYMENT_FAILED_RATE_MAX || 30),
242+
duration: Number(env.SOROBAN_QUEUE_PAYMENT_FAILED_RATE_WINDOW_MS || 10000),
243+
},
244+
},
245+
PaymentFailedGracePeriodStarted: {
246+
concurrency: Number(env.SOROBAN_QUEUE_GRACE_PERIOD_CONCURRENCY || 4),
247+
maxAttempts: Number(env.SOROBAN_QUEUE_GRACE_PERIOD_MAX_ATTEMPTS || 7),
248+
backoffDelay: Number(env.SOROBAN_QUEUE_GRACE_PERIOD_BACKOFF_MS || 5000),
249+
rateLimiter: {
250+
max: Number(env.SOROBAN_QUEUE_GRACE_PERIOD_RATE_MAX || 30),
251+
duration: Number(env.SOROBAN_QUEUE_GRACE_PERIOD_RATE_WINDOW_MS || 10000),
252+
},
253+
},
254+
},
210255
substream: {
211256
baseDomain: env.SUBSTREAM_BASE_DOMAIN || 'substream.app',
212257
backendUrl: env.SUBSTREAM_BACKEND_URL || 'http://localhost:3000',

0 commit comments

Comments
 (0)