Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
291 changes: 93 additions & 198 deletions asset-transfer-basic/rest-api-typescript/package-lock.json

Large diffs are not rendered by default.

9 changes: 4 additions & 5 deletions asset-transfer-basic/rest-api-typescript/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,10 @@
"description": "Asset Transfer Basic REST API implemented in TypeScript",
"main": "dist/index.js",
"engines": {
"node": ">=12"
"node": ">=20"
},
"dependencies": {
"bullmq": "^1.47.2",
"bullmq": "^6.3.4",
"cors": "^2.8.5",
"dotenv": "^10.0.0",
"env-var": "^7.0.1",
Expand All @@ -16,7 +16,7 @@
"fabric-network": "^2.2.20",
"helmet": "^4.6.0",
"http-status-codes": "^2.1.4",
"ioredis": "^4.27.8",
"ioredis": "^5.9.0",
"long": "^5.2.3",
"passport": "^0.6.0",
"passport-headerapikey": "^1.2.2",
Expand All @@ -28,7 +28,6 @@
"@tsconfig/node12": "^12.1.0",
"@types/cors": "^2.8.12",
"@types/express": "^5.0.3",
"@types/ioredis": "^4.26.4",
"@types/jest": "^27.4.1",
"@types/node": "^12.20.55",
"@types/passport": "^1.0.7",
Expand All @@ -40,7 +39,7 @@
"eslint": "^8.49.0",
"eslint-config-prettier": "^8.3.0",
"eslint-plugin-prettier": "^3.4.0",
"ioredis-mock": "^5.6.0",
"ioredis-mock": "^8.13.1",
"jest": "^29.7.0",
"jest-mock-extended": "^3.0.5",
"pino-pretty": "^5.0.2",
Expand Down
22 changes: 0 additions & 22 deletions asset-transfer-basic/rest-api-typescript/src/config.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -192,28 +192,6 @@ describe('Config values', () => {
});
});

describe('submitJobQueueScheduler', () => {
it('defaults to "true"', () => {
const config = require('./config');
expect(config.submitJobQueueScheduler).toBe(true);
});

it('can be configured using the "SUBMIT_JOB_QUEUE_SCHEDULER" environment variable', () => {
process.env.SUBMIT_JOB_QUEUE_SCHEDULER = 'false';
const config = require('./config');
expect(config.submitJobQueueScheduler).toBe(false);
});

it('throws an error when the "SUBMIT_JOB_QUEUE_SCHEDULER" environment variable has an invalid boolean value', () => {
process.env.SUBMIT_JOB_QUEUE_SCHEDULER = '11';
expect(() => {
require('./config');
}).toThrow(
'env-var: "SUBMIT_JOB_QUEUE_SCHEDULER" should be either "true", "false", "TRUE", or "FALSE". An example of a valid value would be: true'
);
});
});

describe('asLocalhost', () => {
it('defaults to "true"', () => {
const config = require('./config');
Expand Down
11 changes: 0 additions & 11 deletions asset-transfer-basic/rest-api-typescript/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,17 +90,6 @@ export const maxFailedSubmitJobs = env
.example('1000')
.asIntPositive();

/**
* Whether to initialise a scheduler for the submit job queue
* There must be at least on queue scheduler to handle retries and you may want
* more than one for redundancy
*/
export const submitJobQueueScheduler = env
.get('SUBMIT_JOB_QUEUE_SCHEDULER')
.default('true')
.example('true')
.asBoolStrict();

/**
* Whether to convert discovered host addresses to be 'localhost'
* This should be set to 'true' when running a docker composed fabric network on the
Expand Down
18 changes: 2 additions & 16 deletions asset-transfer-basic/rest-api-typescript/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,19 +13,14 @@ import {
getContracts,
getNetwork,
} from './fabric';
import {
initJobQueue,
initJobQueueScheduler,
initJobQueueWorker,
} from './jobs';
import { initJobQueue, initJobQueueWorker } from './jobs';
import { logger } from './logger';
import { createServer } from './server';
import { isMaxmemoryPolicyNoeviction } from './redis';
import { Queue, QueueScheduler, Worker } from 'bullmq';
import { Queue, Worker } from 'bullmq';

let jobQueue: Queue | undefined;
let jobQueueWorker: Worker | undefined;
let jobQueueScheduler: QueueScheduler | undefined;

async function main() {
logger.info('Checking Redis config');
Expand Down Expand Up @@ -65,10 +60,6 @@ async function main() {
logger.info('Initialising submit job queue');
jobQueue = initJobQueue();
jobQueueWorker = initJobQueueWorker(app);
if (config.submitJobQueueScheduler === true) {
logger.info('Initialising submit job queue scheduler');
jobQueueScheduler = initJobQueueScheduler();
}
app.locals.jobq = jobQueue;

logger.info('Starting REST server');
Expand All @@ -80,11 +71,6 @@ async function main() {
main().catch(async (err) => {
logger.error({ err }, 'Unxepected error');

if (jobQueueScheduler != undefined) {
logger.debug('Closing job queue scheduler');
await jobQueueScheduler.close();
}

if (jobQueueWorker != undefined) {
logger.debug('Closing job queue worker');
await jobQueueWorker.close();
Expand Down
8 changes: 4 additions & 4 deletions asset-transfer-basic/rest-api-typescript/src/jobs.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -176,8 +176,8 @@ describe('updateJobData', () => {

await updateJobData(mockJob, mockTransaction);

expect(mockJob.update).toBeCalledTimes(1);
expect(mockJob.update).toBeCalledWith({
expect(mockJob.updateData).toBeCalledTimes(1);
expect(mockJob.updateData).toBeCalledWith({
transactionIds: ['txn1', 'txn2'],
transactionState: mockSavedState,
});
Expand All @@ -186,8 +186,8 @@ describe('updateJobData', () => {
it('removes the serialized state from the job data if a transaction is not specified', async () => {
await updateJobData(mockJob, undefined);

expect(mockJob.update).toBeCalledTimes(1);
expect(mockJob.update).toBeCalledWith({
expect(mockJob.updateData).toBeCalledTimes(1);
expect(mockJob.updateData).toBeCalledWith({
transactionIds: ['txn1'],
transactionState: undefined,
});
Expand Down
26 changes: 7 additions & 19 deletions asset-transfer-basic/rest-api-typescript/src/jobs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
* retry support for failing jobs
*/

import { ConnectionOptions, Job, Queue, QueueScheduler, Worker } from 'bullmq';
import { ConnectionOptions, Job, Queue, Worker } from 'bullmq';
import { Application } from 'express';
import { Contract, Transaction } from 'fabric-network';
import * as config from './config';
Expand Down Expand Up @@ -75,6 +75,10 @@ export const initJobQueue = (): Queue => {
/**
* Set up a worker to process submit jobs on the queue, using the
* processSubmitTransactionJob function below
*
* The worker also manages stalled and delayed jobs, which is what makes
* retries with backoff work. Before BullMQ v2 that was the job of a separate
* QueueScheduler, which no longer exists.
*/
export const initJobQueueWorker = (app: Application): Worker => {
const worker = new Worker<JobData, JobResult>(
Expand Down Expand Up @@ -209,23 +213,6 @@ export const processSubmitTransactionJob = async (
}
};

/**
* Set up a scheduler for the submit job queue
*
* This manages stalled and delayed jobs and is required for retries with backoff
*/
export const initJobQueueScheduler = (): QueueScheduler => {
const queueScheduler = new QueueScheduler(config.JOB_QUEUE_NAME, {
connection,
});

queueScheduler.on('failed', (jobId, failedReason) => {
logger.error({ jobId, failedReason }, 'Queue sceduler failure');
});

return queueScheduler;
};

/**
* Helper to add a new submit transaction job to the queue
*/
Expand Down Expand Up @@ -271,7 +258,8 @@ export const updateJobData = async (
newData.transactionState = undefined;
}

await job.update(newData);
// BullMQ v5 renamed Job.update to Job.updateData
await job.updateData(newData);
};

/**
Expand Down
6 changes: 4 additions & 2 deletions asset-transfer-basic/rest-api-typescript/src/redis.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,10 +28,12 @@ export const isMaxmemoryPolicyNoeviction = async (): Promise<boolean> => {
try {
redis = new IORedis(redisOptions);

const maxmemoryPolicyConfig = await (redis as Redis).config(
// ioredis 5 types the CONFIG GET reply as unknown. Over RESP2 it is
// a flat array of alternating parameter names and values.
const maxmemoryPolicyConfig = (await redis.config(
'GET',
'maxmemory-policy'
);
)) as string[];
logger.debug({ maxmemoryPolicyConfig }, 'Got maxmemory-policy config');

if (
Expand Down
3 changes: 2 additions & 1 deletion asset-transfer-basic/rest-api-typescript/tsconfig.json
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,8 @@
"extends": "@tsconfig/node12/tsconfig.json",
"compilerOptions": {
"sourceMap": true,
"outDir": "./dist"
"outDir": "./dist",
"types": ["node", "jest"]
},
"include": [
"src/"
Expand Down
Loading