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
6 changes: 3 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -95,13 +95,13 @@ Make sure that your `.env` configurations exists.
Add new section to the `docker-compose.{dev,prod}.yml` files.

```
hawk-worker-telegram:
hawk-worker-sender:
image: "codexteamuser/hawk-workers:prod"
env_file:
- .env
- workers/telegram/.env
- workers/sender/.env
restart: unless-stopped
entrypoint: /usr/local/bin/node runner.js hawk-worker-telegram
entrypoint: /usr/local/bin/node runner.js hawk-worker-sender
```

## Error handling
Expand Down
42 changes: 3 additions & 39 deletions docker-compose.dev.yml
Original file line number Diff line number Diff line change
Expand Up @@ -121,51 +121,15 @@ services:
- ./:/usr/src/app
- workers-deps:/usr/src/app/node_modules

hawk-worker-email:
hawk-worker-sender:
build:
dockerfile: "dev.Dockerfile"
context: .
env_file:
- .env
- workers/email/.env
- workers/sender/.env
restart: unless-stopped
entrypoint: yarn run-email
volumes:
- ./:/usr/src/app
- workers-deps:/usr/src/app/node_modules

hawk-worker-telegram:
build:
dockerfile: "dev.Dockerfile"
context: .
env_file:
- .env
restart: unless-stopped
entrypoint: yarn run-telegram
volumes:
- ./:/usr/src/app
- workers-deps:/usr/src/app/node_modules

hawk-worker-slack:
build:
dockerfile: "dev.Dockerfile"
context: .
env_file:
- .env
restart: unless-stopped
entrypoint: yarn run-slack
volumes:
- ./:/usr/src/app
- workers-deps:/usr/src/app/node_modules

hawk-worker-webhook:
build:
dockerfile: "dev.Dockerfile"
context: .
env_file:
- .env
restart: unless-stopped
entrypoint: yarn run-webhook
entrypoint: yarn run-sender
volumes:
- ./:/usr/src/app
- workers-deps:/usr/src/app/node_modules
Expand Down
15 changes: 3 additions & 12 deletions docker-compose.prod.yml
Original file line number Diff line number Diff line change
Expand Up @@ -67,20 +67,11 @@ services:
restart: unless-stopped
entrypoint: /usr/local/bin/node runner.js hawk-worker-notifier

hawk-worker-email:
hawk-worker-sender:
image: "codexteamuser/hawk-workers:prod"
network_mode: host
env_file:
- .env
- workers/email/.env
- workers/sender/.env
restart: unless-stopped
entrypoint: /usr/local/bin/node runner.js hawk-worker-email

hawk-worker-telegram:
image: "codexteamuser/hawk-workers:prod"
network_mode: host
env_file:
- .env
- workers/telegram/.env
restart: unless-stopped
entrypoint: /usr/local/bin/node runner.js hawk-worker-telegram
entrypoint: /usr/local/bin/node runner.js hawk-worker-sender
39 changes: 32 additions & 7 deletions lib/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,11 @@ export abstract class Worker {
*/
private registryConnection: amqp.Connection;

/**
* Connection to Registry owned by another module: not opened and not closed by the worker
*/
private sharedRegistryConnection: amqp.Connection;

/**
* Channel is a "transport-way" between Consumer and Registry inside the connection
* One connection can has several channels.
Expand Down Expand Up @@ -121,6 +126,16 @@ export abstract class Worker {
// return [ this.metricSuccessfullyProcessedMessages ];
// }

/**
* Share an existing Registry connection instead of opening a new one.
* Call before start(); the worker still opens its own channel, but does not close the connection
*
* @param connection - connection to Registry to use
*/
public useRegistryConnection(connection: amqp.Connection): void {
this.sharedRegistryConnection = connection;
}

/**
* Start consuming messages
*/
Expand Down Expand Up @@ -214,11 +229,6 @@ export abstract class Worker {
* Connect to RabbitMQ server
*/
private async connect(): Promise<void> {
/**
* Connect to RabbitMQ
*/
this.registryConnection = await amqp.connect(this.registryUrl);

const errorHandler = (error: Error): void => {
this.logger.error('Error in RabbitMQ has been occurred', error);
HawkCatcher.send(error, {
Expand All @@ -231,7 +241,19 @@ export abstract class Worker {
process.exit(1);
};

this.registryConnection.on('error', errorHandler);
if (this.sharedRegistryConnection) {
/**
* Shared connection is opened and observed by its owner
*/
this.registryConnection = this.sharedRegistryConnection;
} else {
/**
* Connect to RabbitMQ
*/
this.registryConnection = await amqp.connect(this.registryUrl);

this.registryConnection.on('error', errorHandler);
}

/**
* Open channel inside the connection
Expand Down Expand Up @@ -397,7 +419,10 @@ export abstract class Worker {
await this.channelWithRegistry.close();
}

if (this.registryConnection) {
/**
* Shared connection is closed by its owner
*/
if (this.registryConnection && !this.sharedRegistryConnection) {
await this.registryConnection.close();
}

Expand Down
7 changes: 7 additions & 0 deletions lib/workerNames.js
Original file line number Diff line number Diff line change
Expand Up @@ -37,4 +37,11 @@ workersDir.forEach(file => {
workers[file.name.toUpperCase()] = pkg.workerType;
});

/**
* Queues of the sender worker channels (one queue per channel, see workers/sender)
*/
['email', 'telegram', 'slack', 'webhook', 'loop'].forEach(channel => {
workers[channel.toUpperCase()] = `sender/${channel}`;
});

module.exports = workers;
12 changes: 3 additions & 9 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -25,33 +25,27 @@
"test:sentry": "jest workers/sentry --config workers/sentry/jest.config.js",
"test:javascript": "jest workers/javascript",
"test:release": "jest workers/release",
"test:slack": "jest workers/slack",
"test:loop": "jest workers/loop",
"test:sender": "jest workers/sender",
"test:limiter": "jest workers/limiter --runInBand",
"test:grouper": "jest workers/grouper",
"test:diff": "jest ./workers/grouper/tests/diff.test.ts",
"test:paymaster": "jest workers/paymaster",
"test:notifier": "jest workers/notifier",
"test:js": "jest workers/javascript",
"test:task-manager": "jest workers/task-manager",
"test:webhook": "jest workers/webhook",
"test:clear": "jest --clearCache",
"run-default": "yarn worker hawk-worker-default",
"run-sentry": "yarn worker hawk-worker-sentry",
"run-js": "yarn worker hawk-worker-javascript",
"run-slack": "yarn worker hawk-worker-slack",
"run-loop": "yarn worker hawk-worker-loop",
"run-sender": "yarn worker hawk-worker-sender",
"run-grouper": "yarn worker hawk-worker-grouper",
"run-archiver": "yarn worker hawk-worker-archiver",
"run-accountant": "yarn worker hawk-worker-accountant",
"run-paymaster": "yarn worker hawk-worker-paymaster",
"run-notifier": "yarn worker hawk-worker-notifier",
"run-release": "yarn worker hawk-worker-release",
"run-email": "yarn worker hawk-worker-email",
"run-telegram": "yarn worker hawk-worker-telegram",
"run-limiter": "yarn worker hawk-worker-limiter",
"run-task-manager": "yarn worker hawk-worker-task-manager",
"run-webhook": "yarn worker hawk-worker-webhook"
"run-task-manager": "yarn worker hawk-worker-task-manager"
},
"dependencies": {
"@babel/parser": "^7.26.9",
Expand Down
25 changes: 0 additions & 25 deletions workers/email/.env.sample

This file was deleted.

1 change: 0 additions & 1 deletion workers/email/.gitignore

This file was deleted.

24 changes: 0 additions & 24 deletions workers/email/README.md

This file was deleted.

22 changes: 0 additions & 22 deletions workers/email/package.json

This file was deleted.

9 changes: 0 additions & 9 deletions workers/email/src/env.ts

This file was deleted.

25 changes: 0 additions & 25 deletions workers/email/src/index.ts

This file was deleted.

Loading
Loading