Skip to content
Open
Show file tree
Hide file tree
Changes from 8 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
1 change: 1 addition & 0 deletions services/pub-sub/.npmignore
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
src
119 changes: 119 additions & 0 deletions services/pub-sub/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
# `@byndyusoft-ui/pub-sub`

> A performant Pub/Sub interface with controlled instance management

### Installation

```bash
npm i @byndyusoft-ui/pub-sub
```

## Usage

#### Import the class

```ts
import PubSub from '@byndyusoft-ui/pub-sub';
```

#### Define your channels
Create a type that defines the channels and their corresponding callback signatures.

```ts
type ChannelsType = {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Вот этот тип мне кажется проблемой. С ним у нас есть место, которое должно знать о всех событиях, которые надо обрабатывать. Как будто бы появляется лишняя связь между разными частями приложения.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Согласен, с глобальным экземпляром есть такая проблема. Тут можно использовать pub-sub только внутри модуля, если это возможно. Или не типизировать глобальный экземпляр и делать адаптер в каждом модуле.

addTodo: (data: TodoType) => void;
removeTodo: (todoId: number) => void;
removeAll: () => void;
// For async callbacks:
asyncMessage: (data: string) => Promise<void>;
};
```

#### Create an instance
```ts
const pubSubInstance = new PubSub<ChannelsType>();
```

#### Subscribe & Unsubscribe
Basic Subscription
```ts
const addTodoCallback = (data: TodoType) => {
console.log('Added new todo:', data);
};

// subscribe
pubSubInstance.subscribe('addTodo', addTodoCallback);

// unsubscribe
pubSubInstance.unsubscribe('addTodo', addTodoCallback);
```

#### One-Time Subscription
Use `subscribeOnce` to subscribe to an event that should be handled only once:

```ts
pubSubInstance.subscribeOnce('addTodo', (data) => {
console.log('This callback will only be executed once:', data);
});
```

#### Unsubscribe All
Remove all callbacks from a specific channel or from all channels:

```ts
// Unsubscribe all from a specific channel
pubSubInstance.unsubscribeAll('addTodo');

// Unsubscribe all from all channels
pubSubInstance.unsubscribeAll();
```

#### Publish Events

Synchronous Publish

```ts
pubSubInstance.publish('addTodo', { id: 1, text: 'Some todo'});
```

Asynchronous Publish
Use `publishAsync` to publish data and wait for asynchronous subscribers:

```ts
pubSubInstance.subscribe('asyncMessage', async (data) => {
await new Promise((resolve) => setTimeout(resolve, 1000));
console.log(`Async received: ${data}`);
});
```

#### Publish asynchronously
Use publishAsync to publish data and handle asynchronous subscribers.

```ts

pubSubInstance.subscribe('asyncMessage', async (data) => {
await new Promise((resolve) => setTimeout(resolve, 1000));
console.log(`Async received: ${data}`);
});

await pubSubInstance.publishAsync('asyncMessage', 'This is asynchronous!');
```


#### Get All Subscriptions
For debugging or monitoring, you can retrieve current subscriptions:

```ts
const subscriptions = pubSubInstance.allSubscribes();
console.log(subscriptions);
// Output example:
// [ { channel: 'addTodo', subscribers: 2 }, { channel: 'asyncMessage', subscribers: 1 } ]
```

#### Reset Subscriptions
Clear all channels and their subscribers:

```ts
pubSubInstance.reset();
```

34 changes: 34 additions & 0 deletions services/pub-sub/package.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
{
"name": "@byndyusoft-ui/pub-sub",
"version": "0.0.1",
"description": "Byndyusoft UI Service",
"keywords": [
"byndyusoft",
"byndyusoft-ui",
"channels",
"publish",
"subscribe",
"Pub/Sub"
],
"author": "Gleb Fomin <gleb.fom28@gmail.com>",
"homepage": "https://github.com/Byndyusoft/ui/tree/master/services/pub-sub#readme",
"license": "Apache-2.0",
"main": "dist/index.js",
"types": "dist/index.d.ts",
"repository": {
"type": "git",
"url": "git+https://github.com/Byndyusoft/ui.git"
},
"scripts": {
"build": "tsc --project tsconfig.build.json",
"clean": "rimraf dist",
"lint": "eslint src --config ../../eslint.config.js",
"test": "jest --config ../../jest.config.js --roots services/pub-sub/src"
},
"bugs": {
"url": "https://github.com/Byndyusoft/ui/issues"
},
"publishConfig": {
"access": "public"
}
}
1 change: 1 addition & 0 deletions services/pub-sub/src/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
export { default } from './pubSub';
126 changes: 126 additions & 0 deletions services/pub-sub/src/pubSub.tests.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
import PubSub from './pubSub';

interface IChannels {
testChannel: (data?: string) => void;
asyncChannel: (data?: string) => Promise<void>;
}

describe('services/pub-sub', () => {
const pubSub = new PubSub<IChannels>();

afterEach(() => {
pubSub.reset();
});

test('should subscribe and publish to a channel', () => {
const callback = jest.fn();
pubSub.subscribe('testChannel', callback);

pubSub.publish('testChannel', 'Hello, World!');

expect(callback).toHaveBeenCalledTimes(1);
expect(callback).toHaveBeenCalledWith('Hello, World!');
});

test('should not call callback if no subscribers', () => {
const callback = jest.fn();
pubSub.publish('testChannel');

expect(callback).not.toHaveBeenCalled();
});

test('should unsubscribe from a channel', () => {
const callback = jest.fn();
pubSub.subscribe('testChannel', callback);
pubSub.unsubscribe('testChannel', callback);

pubSub.publish('testChannel');

expect(callback).not.toHaveBeenCalled();
});

test('should warn if no subscribers are present for a channel', () => {
console.warn = jest.fn();

pubSub.publish('testChannel', 'No one is listening');

expect(console.warn).toHaveBeenCalledWith('No subscribers for channel: testChannel');
});

test('should handle async subscribe callbacks', async () => {
const asyncCallback = jest.fn().mockResolvedValue(undefined);
pubSub.subscribe('asyncChannel', asyncCallback);

await pubSub.publishAsync('asyncChannel', 'Async data');

expect(asyncCallback).toHaveBeenCalledTimes(1);
expect(asyncCallback).toHaveBeenCalledWith('Async data');
});

test('should reset all subscriptions', () => {
const callback1 = jest.fn();
const callback2 = jest.fn();

pubSub.subscribe('testChannel', callback1);
pubSub.subscribe('testChannel', callback2);

pubSub.reset();

pubSub.publish('testChannel');

expect(callback1).not.toHaveBeenCalled();
expect(callback2).not.toHaveBeenCalled();
});

test('should unsubscribe all callbacks for all channels using unsubscribeAll', () => {
const callback1 = jest.fn();
const callback2 = jest.fn();

pubSub.subscribe('testChannel', callback1);
pubSub.subscribe('asyncChannel', callback2);

pubSub.unsubscribeAll();

pubSub.publish('testChannel', 'Test data');
pubSub.publish('asyncChannel', 'Test data');

expect(callback1).not.toHaveBeenCalled();
expect(callback2).not.toHaveBeenCalled();
});

test('should call subscribeOnce callback only once', () => {
const callback = jest.fn();
pubSub.subscribeOnce('testChannel', callback);

// First publish should trigger the callback.
pubSub.publish('testChannel', 'Test message 1');

// Subsequent publish should not trigger the callback.
pubSub.publish('testChannel', 'Test message 2');

expect(callback).toHaveBeenCalledTimes(1);
expect(callback).toHaveBeenCalledWith('Test message 1');
});

test('should return all subscriptions info', () => {
const callback1 = jest.fn();
const callback2 = jest.fn();

pubSub.subscribe('testChannel', callback1);
pubSub.subscribe('testChannel', callback2);
pubSub.subscribe('asyncChannel', callback1);

const result = pubSub.allSubscribes();

const testChannelInfo = result.find(item => item.channel === 'testChannel');
const asyncChannelInfo = result.find(item => item.channel === 'asyncChannel');

expect(testChannelInfo).toBeDefined();

expect(testChannelInfo?.subscribers).toBe(2);

expect(asyncChannelInfo).toBeDefined();

expect(asyncChannelInfo?.subscribers).toBe(1);
});
});
112 changes: 112 additions & 0 deletions services/pub-sub/src/pubSub.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
import { TAllSubscribesResult, TChannelData, TChannelMap, TDefaultChannels } from './pubSub.types';

class PubSub<ChannelsRecord extends TDefaultChannels<ChannelsRecord>> {
private channels: TChannelMap<ChannelsRecord> = new Map();

/**
* Subscribe to the channel.
*/
subscribe<ChannelKey extends keyof ChannelsRecord>(
channel: ChannelKey,
callback: ChannelsRecord[ChannelKey]
): void {
if (!this.channels.has(channel)) {
this.channels.set(channel, new Set());
}
(this.channels.get(channel) as Set<ChannelsRecord[ChannelKey]>).add(callback);
}

/**
* Unsubscribe from the channel.
*/
unsubscribe<ChannelKey extends keyof ChannelsRecord>(
channel: ChannelKey,
callback: ChannelsRecord[ChannelKey]
): void {
const channelSet = this.channels.get(channel);
if (channelSet) {
channelSet.delete(callback);
if (channelSet.size === 0) {
this.channels.delete(channel);
}
}
}

/**
* After the first execution, the callback is automatically unsubscribed.
*/
subscribeOnce<ChannelKey extends keyof ChannelsRecord>(
channel: ChannelKey,
callback: ChannelsRecord[ChannelKey]
): void {
const onceCallback: ChannelsRecord[ChannelKey] = ((data?: TChannelData<ChannelsRecord, ChannelKey>) => {
this.unsubscribe(channel, onceCallback);
return callback(data);
}) as ChannelsRecord[ChannelKey];

this.subscribe(channel, onceCallback);
}

/**
* Unsubscribe all callbacks for a specific channel or all channels.
*/
unsubscribeAll = <ChannelKey extends keyof ChannelsRecord>(channel?: ChannelKey): void => {
if (channel) {
this.channels.delete(channel);
} else {
this.channels.clear();
}
};

/**
* Publishing to the channel.
*/
publish<ChannelKey extends keyof ChannelsRecord>(
channel: ChannelKey,
data?: TChannelData<ChannelsRecord, ChannelKey>
): void {
const channelSet = this.channels.get(channel);
if (channelSet) {
for (const callback of channelSet) {
callback(data);
}
} else {
console.warn(`No subscribers for channel: ${String(channel)}`);
}
}

async publishAsync<ChannelKey extends keyof ChannelsRecord>(
channel: ChannelKey,
data?: TChannelData<ChannelsRecord, ChannelKey>
): Promise<void> {
const channelSet = this.channels.get(channel);

if (!channelSet) {
console.warn(`No subscribers for channel: ${String(channel)}`);
return;
}

const promises = Array.from(channelSet).map(callback => Promise.resolve(callback(data)));

await Promise.all(promises);
}

/**
* Returns an array containing information about all current subscriptions.
*/
allSubscribes(): TAllSubscribesResult<ChannelsRecord> {
return Array.from(this.channels.entries()).map(([channel, subscribers]) => ({
channel,
subscribers: subscribers.size
}));
}

/**
* Reset all subscriptions.
*/
reset(): void {
this.channels.clear();
}
}

export default PubSub;
Loading