Skip to content
Closed
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
25 changes: 21 additions & 4 deletions libraries/hermes/src/Socket/Socket.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
EventResponse,
isInstanceOfEventResponse,
JoinRoomResponse,
LeaveRoomResponse,
SocketPush,
} from '../interfaces/index.js';
import ObjectHash from 'object-hash';
Expand Down Expand Up @@ -184,11 +185,13 @@ export class SocketController extends ConduitRouter {
}

async handleSocketPush(push: SocketPush) {
const localOnly = push.localOnly === true;
if (push.event === 'join-room') {
if (push.rooms.length === 0) return;
const filteredSockets = await this.findAndFilterSockets(
push.receivers,
push.namespace,
localOnly,
);
for (const socket of filteredSockets) {
ConduitGrpcSdk.Logger.info(
Expand All @@ -203,6 +206,7 @@ export class SocketController extends ConduitRouter {
const filteredSockets = await this.findAndFilterSockets(
push.receivers,
push.namespace,
localOnly,
);
for (const socket of filteredSockets) {
for (const room of push.rooms) {
Expand All @@ -221,20 +225,31 @@ export class SocketController extends ConduitRouter {
ConduitGrpcSdk.Logger.info(
`Emitting event: ${push.event} to all sockets in namespace: ${push.namespace}`,
);
this.io.of(push.namespace).emit(push.event, push.data);
const nsp = this.io.of(push.namespace);
if (localOnly) {
nsp.local.emit(push.event, push.data);
} else {
nsp.emit(push.event, push.data);
}
} else {
if (push.rooms.length !== 0) {
ConduitGrpcSdk.Logger.info(
`Emitting event: ${push.event} to rooms: ${push.rooms.join(
', ',
)} in namespace: ${push.namespace}`,
);
this.io.of(push.namespace).to(push.rooms).emit(push.event, push.data);
const target = this.io.of(push.namespace).to(push.rooms);
if (localOnly) {
target.local.emit(push.event, push.data);
} else {
target.emit(push.event, push.data);
}
}
if (push.receivers.length !== 0) {
const filteredSockets = await this.findAndFilterSockets(
push.receivers,
push.namespace,
localOnly,
);
for (const socket of filteredSockets) {
ConduitGrpcSdk.Logger.info(
Expand All @@ -248,7 +263,7 @@ export class SocketController extends ConduitRouter {
}

private async handleResponse(
res: EventResponse | JoinRoomResponse,
res: EventResponse | JoinRoomResponse | LeaveRoomResponse,
socket: Socket,
namespace: string,
) {
Expand Down Expand Up @@ -307,8 +322,10 @@ export class SocketController extends ConduitRouter {
async findAndFilterSockets(
userIds: string[],
namespace: string,
localOnly: boolean = false,
): Promise<RemoteSocket<any, any>[]> {
const sockets = await this.io.of(namespace).fetchSockets();
const nsp = this.io.of(namespace);
const sockets = localOnly ? await nsp.local.fetchSockets() : await nsp.fetchSockets();
const userIdSet = new Set(userIds);
return sockets.filter(socket => {
if (socket.data && socket.data.user) {
Expand Down
9 changes: 8 additions & 1 deletion libraries/hermes/src/interfaces/Socket.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,14 @@ export type JoinRoomResponse = {
rooms: string[];
};

export type ConduitSocketHandlerResponse = Promise<EventResponse | JoinRoomResponse>;
export type LeaveRoomResponse = {
event: 'leave-room';
rooms: string[];
};

export type ConduitSocketHandlerResponse = Promise<
EventResponse | JoinRoomResponse | LeaveRoomResponse
>;

export type ConduitSocketEventHandler = (
request: ConduitSocketParameters,
Expand Down
1 change: 1 addition & 0 deletions libraries/hermes/src/interfaces/SocketPush.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,4 +4,5 @@ export interface SocketPush {
receivers: string[];
rooms: string[];
namespace: string;
localOnly?: boolean;
}
40 changes: 40 additions & 0 deletions modules/router/README.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -77,3 +77,43 @@ If you sent any other string that those provided conduit will not utilize cachin
The number that you specify also sets the expiry time of conduit's cache. All caching takes place,
AFTER middleware execution
- cacheControl?: string;

## Event Relays

Event Relays map an **exact** Redis bus channel to a Socket.io event on the dedicated `/events/` namespace.
They are configured through the Admin API (`/router/event-relays`) and are **ephemeral**: Redis pub/sub is not replayed, and missed events are not stored.

Each relay stores a ReBAC check (`permission` on `resourceType:resourceId`). Clients never choose a raw room name. Optional `notes` describe the relay for operators.

### Client contract

```javascript
import { io } from 'socket.io-client';

const socket = io(`${SOCKET_BASE_URL}/events/`, {
path: '/realtime',
extraHeaders: {
authorization: `Bearer ${accessToken}`,
},
});

socket.emit('subscribe', relayId, resourceId);
socket.on('order-updated', payload => {
// payload is the rendered JSON template
});
socket.emit('unsubscribe', relayId, resourceId);
```

- Subscribe is allowed only when `User:<id>` has the relay `permission` on `<resourceType>:<resourceId>`.
- The outbound event name is the relay `socketEvent`.
- Placeholders in `messageTemplate` use `{{payload.path}}` against the bus JSON payload.

### Admin API

| Method | Path | Description |
|:-------|:-----|:------------|
| `GET` | `/router/event-relays` | Paginated list (`skip`, `limit`, `search`) |
| `GET` | `/router/event-relays/:id` | Single relay |
| `POST` | `/router/event-relays` | Create |
| `PATCH` | `/router/event-relays/:id` | Update |
| `DELETE` | `/router/event-relays/:id` | Delete |
1 change: 1 addition & 0 deletions modules/router/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
"prebuild:bundle": "pnpm --filter @conduitplatform/service-bundle run build",
"build:bundle": "rimraf bundle && node ../../libraries/service-bundle/dist/cli.js generate-manifest && tsup && node ../../libraries/service-bundle/dist/cli.js copy-assets && node ../../libraries/service-bundle/dist/cli.js generate-lockfile",
"generateTypes": "sh build.sh",
"test": "npx tsc -p tsconfig.test.json && node --test dist-test/event-relays/*.test.js",
"build:docker": "docker build -t ghcr.io/conduitplatform/router:latest -f ./Dockerfile ../../ && docker push ghcr.io/conduitplatform/router:latest"
},
"license": "ISC",
Expand Down
28 changes: 27 additions & 1 deletion modules/router/src/Router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,11 @@ import * as adminRoutes from './admin/routes/index.js';
import metricsSchema from './metrics/index.js';
import { ConfigController, ManagedModule } from '@conduitplatform/module-tools';
import { fileURLToPath } from 'node:url';
import {
createEventRelayPusher,
createEventsSocket,
EventRelayManager,
} from './event-relays/index.js';

const __filename = fileURLToPath(import.meta.url);
const __dirname = path.dirname(__filename);
Expand Down Expand Up @@ -69,6 +74,8 @@ export default class ConduitDefaultRouter extends ManagedModule<Config> {
private hasAppliedMiddleware: string[] = [];
private _refreshTimeout: NodeJS.Timeout | null = null;
private _haInitialized = false;
private eventRelayManager: EventRelayManager;
private eventsSocket?: ConduitSocket;

constructor(peerManifestRoot?: string) {
super('router', peerManifestRoot);
Expand Down Expand Up @@ -102,7 +109,16 @@ export default class ConduitDefaultRouter extends ManagedModule<Config> {
}

async onRegister() {
this.adminRouter = new AdminHandlers(this.grpcServer, this.grpcSdk, this);
this.eventRelayManager = new EventRelayManager(
this.grpcSdk,
createEventRelayPusher(data => this._internalRouter.socketPush(data)),
);
this.adminRouter = new AdminHandlers(
this.grpcServer,
this.grpcSdk,
this,
this.eventRelayManager,
);
this._security = new SecurityModule(this.grpcSdk, this);
}

Expand Down Expand Up @@ -132,8 +148,11 @@ export default class ConduitDefaultRouter extends ManagedModule<Config> {
}
if (config.transports.sockets) {
this._internalRouter.initSockets();
this.registerEventsNamespace();
await this.eventRelayManager.start();
atLeastOne = true;
} else {
await this.eventRelayManager?.stop();
this._internalRouter.stopSockets();
}

Expand Down Expand Up @@ -336,6 +355,13 @@ export default class ConduitDefaultRouter extends ManagedModule<Config> {
this._internalRouter.registerConduitRoute(route);
}

private registerEventsNamespace() {
if (!this.eventsSocket) {
this.eventsSocket = createEventsSocket(this.grpcSdk, this.eventRelayManager);
}
this._internalRouter.registerConduitSocket(this.eventsSocket);
}

protected registerSchemas(): Promise<unknown> {
const promises = Object.values(models).map(model => {
const modelInstance = model.getInstance(this.grpcSdk.database!);
Expand Down
128 changes: 128 additions & 0 deletions modules/router/src/admin/event-relays.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
import { status } from '@grpc/grpc-js';
import {
GrpcError,
ParsedRouterRequest,
Query,
UnparsedRouterResponse,
} from '@conduitplatform/grpc-sdk';
import { isNil } from 'lodash-es';
import { EventRelay } from '../models/index.js';
import { EventRelayManager } from '../event-relays/EventRelayManager.js';
import { EventRelayInput, validateEventRelayInput } from '../event-relays/validation.js';
import { EventRelayValidationError } from '../event-relays/validationError.js';
import { buildSearchQuery, parsePagination } from '../event-relays/search.js';

export class EventRelayAdmin {
constructor(private readonly manager: EventRelayManager) {}

async listEventRelays(call: ParsedRouterRequest): Promise<UnparsedRouterResponse> {
const { skip, limit } = parsePagination(
call.request.params.skip,
call.request.params.limit,
);
const { search } = call.request.params;
const query = buildSearchQuery(search) as Query<EventRelay>;
const relays = await EventRelay.getInstance().findMany(query, {
skip,
limit,
sort: { updatedAt: -1 },
});
const count = await EventRelay.getInstance().countDocuments(query);
return { relays, count };
}

async getEventRelay(call: ParsedRouterRequest): Promise<UnparsedRouterResponse> {
const relay = await EventRelay.getInstance().findOne({
_id: call.request.params.id,
});
if (isNil(relay)) {
throw new GrpcError(status.NOT_FOUND, 'Event relay not found');
}
return relay;
}

async createEventRelay(call: ParsedRouterRequest): Promise<UnparsedRouterResponse> {
const input = parseInput(call.request.params as EventRelayInput);
await assertUniqueName(input.name);
const relay = await EventRelay.getInstance().create({
...input,
messageTemplate: input.messageTemplate as EventRelay['messageTemplate'],
});
await this.manager.notifyChanged();
return relay;
}

async patchEventRelay(call: ParsedRouterRequest): Promise<UnparsedRouterResponse> {
const existing = await EventRelay.getInstance().findOne({
_id: call.request.params.id,
});
if (isNil(existing)) {
throw new GrpcError(status.NOT_FOUND, 'Event relay not found');
}

const merged: EventRelayInput = {
name: call.request.params.name ?? existing.name,
notes:
call.request.params.notes === undefined
? existing.notes
: call.request.params.notes,
active:
call.request.params.active === undefined
? existing.active
: call.request.params.active,
busEvent: call.request.params.busEvent ?? existing.busEvent,
socketEvent: call.request.params.socketEvent ?? existing.socketEvent,
resourceType: call.request.params.resourceType ?? existing.resourceType,
resourceIdPath: call.request.params.resourceIdPath ?? existing.resourceIdPath,
permission: call.request.params.permission ?? existing.permission,
messageTemplate:
call.request.params.messageTemplate === undefined
? existing.messageTemplate
: call.request.params.messageTemplate,
};
const input = parseInput(merged);
if (input.name !== existing.name) {
await assertUniqueName(input.name, existing._id);
}

const updated = await EventRelay.getInstance().findByIdAndUpdate(existing._id, {
...input,
messageTemplate: input.messageTemplate as EventRelay['messageTemplate'],
});
await this.manager.notifyChanged();
return updated!;
}

async deleteEventRelay(call: ParsedRouterRequest): Promise<UnparsedRouterResponse> {
const existing = await EventRelay.getInstance().findOne({
_id: call.request.params.id,
});
if (isNil(existing)) {
throw new GrpcError(status.NOT_FOUND, 'Event relay not found');
}
await EventRelay.getInstance().deleteOne({ _id: existing._id });
await this.manager.notifyChanged();
return { message: 'Event relay deleted' };
}
}

function parseInput(params: EventRelayInput): EventRelayInput {
try {
return validateEventRelayInput(params);
} catch (err) {
if (err instanceof EventRelayValidationError) {
throw new GrpcError(status.INVALID_ARGUMENT, err.message);
}
throw err;
}
}

async function assertUniqueName(name: string, excludeId?: string): Promise<void> {
const existing = await EventRelay.getInstance().findOne({ name });
if (existing && existing._id !== excludeId) {
throw new GrpcError(
status.ALREADY_EXISTS,
`An event relay named '${name}' already exists`,
);
}
}
Loading
Loading