All checks were successful
farmcontrol/farmcontrol-ws/pipeline/head This commit looks good
230 lines
6.3 KiB
JavaScript
230 lines
6.3 KiB
JavaScript
import log4js from 'log4js';
|
|
import _ from 'lodash';
|
|
import { loadConfig } from '../config.js';
|
|
import { natsServer } from '../database/nats.js';
|
|
import { expandObjectIds } from '../database/utils.js';
|
|
import { formatTraceData } from '../utils.js';
|
|
const config = loadConfig();
|
|
|
|
// Setup logger
|
|
const logger = log4js.getLogger('Update Manager');
|
|
logger.level = config.server.logLevel;
|
|
|
|
const normalizeFilter = filter =>
|
|
filter && typeof filter === 'object' && !Array.isArray(filter) ? filter : {};
|
|
|
|
const getFilterValue = (object, key) => {
|
|
if (key.endsWith('._id')) {
|
|
const refPath = key.slice(0, -4);
|
|
const ref = _.get(object, refPath);
|
|
if (ref && typeof ref === 'object' && ref._id) {
|
|
return ref._id;
|
|
}
|
|
return ref;
|
|
}
|
|
|
|
return _.get(object, key);
|
|
};
|
|
|
|
const valuesMatch = (actual, expected) => {
|
|
if (actual == expected) {
|
|
return true;
|
|
}
|
|
|
|
if (actual != null && expected != null) {
|
|
return String(actual) === String(expected);
|
|
}
|
|
|
|
return false;
|
|
};
|
|
|
|
const matchesFilter = (object, filter) => {
|
|
if (!filter || Object.keys(filter).length === 0) {
|
|
return true;
|
|
}
|
|
|
|
if (object == null) {
|
|
return false;
|
|
}
|
|
|
|
const normalizedObject =
|
|
typeof object === 'object' && !Array.isArray(object)
|
|
? object
|
|
: { _id: object };
|
|
|
|
for (const [key, expectedValue] of Object.entries(filter)) {
|
|
if (!valuesMatch(getFilterValue(normalizedObject, key), expectedValue)) {
|
|
return false;
|
|
}
|
|
}
|
|
|
|
return true;
|
|
};
|
|
|
|
const stableStringify = value => {
|
|
if (Array.isArray(value)) {
|
|
return `[${value.map(stableStringify).join(',')}]`;
|
|
}
|
|
|
|
if (value && typeof value === 'object') {
|
|
return `{${Object.keys(value)
|
|
.sort()
|
|
.map(key => `${JSON.stringify(key)}:${stableStringify(value[key])}`)
|
|
.join(',')}}`;
|
|
}
|
|
|
|
return JSON.stringify(value);
|
|
};
|
|
|
|
const getSubscriptionOwner = (socketId, filter) =>
|
|
`${socketId}:${stableStringify(normalizeFilter(filter))}`;
|
|
|
|
const getSubscriptionKey = (subject, owner) => `${subject}:${owner}`;
|
|
|
|
/**
|
|
* UpdateManager handles tracking object updates and broadcasts update events via websockets.
|
|
*/
|
|
export class UpdateManager {
|
|
constructor(socketClient) {
|
|
this.socketClient = socketClient;
|
|
this.subscriptions = new Set();
|
|
}
|
|
|
|
matchesObjectTypeFilter(objectType, filter, value) {
|
|
return matchesFilter(value, normalizeFilter(filter));
|
|
}
|
|
|
|
emitObjectTypeEvent(eventName, objectType, filter, value) {
|
|
const normalizedFilter = normalizeFilter(filter);
|
|
const matches = this.matchesObjectTypeFilter(
|
|
objectType,
|
|
normalizedFilter,
|
|
value
|
|
);
|
|
|
|
if (!matches) {
|
|
logger.trace(
|
|
`Filtered ${eventName} event: ${formatTraceData({
|
|
objectType,
|
|
filter: normalizedFilter,
|
|
value
|
|
})}`
|
|
);
|
|
return;
|
|
}
|
|
|
|
this.socketClient.socket.emit(eventName, {
|
|
object: value,
|
|
objectType: objectType,
|
|
filter: normalizedFilter
|
|
});
|
|
}
|
|
|
|
async subscribeToObjectNew(objectType, filter = {}) {
|
|
const normalizedFilter = normalizeFilter(filter);
|
|
const subject = `${objectType}s.new`;
|
|
const owner = getSubscriptionOwner(
|
|
this.socketClient.socketId,
|
|
normalizedFilter
|
|
);
|
|
|
|
await natsServer.subscribe(subject, owner, async (key, value) => {
|
|
logger.trace(`Object new event: ${formatTraceData(value)}`);
|
|
this.emitObjectTypeEvent(
|
|
'objectNew',
|
|
objectType,
|
|
normalizedFilter,
|
|
value
|
|
);
|
|
});
|
|
|
|
this.subscriptions.add(getSubscriptionKey(subject, owner));
|
|
return { success: true };
|
|
}
|
|
|
|
async subscribeToObjectDelete(objectType, filter = {}) {
|
|
const normalizedFilter = normalizeFilter(filter);
|
|
const subject = `${objectType}s.delete`;
|
|
const owner = getSubscriptionOwner(
|
|
this.socketClient.socketId,
|
|
normalizedFilter
|
|
);
|
|
|
|
await natsServer.subscribe(subject, owner, async (key, value) => {
|
|
logger.trace(`Object delete event: ${formatTraceData(value)}`);
|
|
this.emitObjectTypeEvent(
|
|
'objectDelete',
|
|
objectType,
|
|
normalizedFilter,
|
|
value
|
|
);
|
|
});
|
|
|
|
this.subscriptions.add(getSubscriptionKey(subject, owner));
|
|
return { success: true };
|
|
}
|
|
|
|
async subscribeToObjectUpdate(id, objectType) {
|
|
logger.debug('Subscribing to object update...', id, objectType);
|
|
const subject = `${objectType}s.${id}.object`;
|
|
const owner = this.socketClient.socketId;
|
|
|
|
await natsServer.subscribe(subject, owner, (key, value) => {
|
|
const expandedValue = expandObjectIds(value);
|
|
logger.trace('Object update event:', id);
|
|
this.socketClient.socket.emit('objectUpdate', {
|
|
_id: id,
|
|
objectType: objectType,
|
|
object: { ...expandedValue }
|
|
});
|
|
});
|
|
|
|
this.subscriptions.add(getSubscriptionKey(subject, owner));
|
|
return { success: true };
|
|
}
|
|
|
|
async removeObjectNewListener(objectType, filter = {}) {
|
|
const subject = `${objectType}s.new`;
|
|
const owner = getSubscriptionOwner(this.socketClient.socketId, filter);
|
|
|
|
await natsServer.removeSubscription(subject, owner);
|
|
this.subscriptions.delete(getSubscriptionKey(subject, owner));
|
|
return { success: true };
|
|
}
|
|
|
|
async removeObjectDeleteListener(objectType, filter = {}) {
|
|
const subject = `${objectType}s.delete`;
|
|
const owner = getSubscriptionOwner(this.socketClient.socketId, filter);
|
|
|
|
await natsServer.removeSubscription(subject, owner);
|
|
this.subscriptions.delete(getSubscriptionKey(subject, owner));
|
|
return { success: true };
|
|
}
|
|
|
|
async removeObjectUpdateListener(id, objectType) {
|
|
const subject = `${objectType}s.${id}.object`;
|
|
const owner = this.socketClient.socketId;
|
|
|
|
await natsServer.removeSubscription(subject, owner);
|
|
this.subscriptions.delete(getSubscriptionKey(subject, owner));
|
|
return { success: true };
|
|
}
|
|
|
|
async removeAllListeners() {
|
|
logger.debug('Removing all update listeners...');
|
|
const removePromises = Array.from(this.subscriptions).map(
|
|
subscriptionKey => {
|
|
const separatorIndex = subscriptionKey.indexOf(':');
|
|
const subject = subscriptionKey.slice(0, separatorIndex);
|
|
const owner = subscriptionKey.slice(separatorIndex + 1);
|
|
return natsServer.removeSubscription(subject, owner);
|
|
}
|
|
);
|
|
|
|
await Promise.all(removePromises);
|
|
this.subscriptions.clear();
|
|
logger.debug(`Removed ${removePromises.length} update listener(s)`);
|
|
return { success: true };
|
|
}
|
|
}
|