Integrations data (#770)

* feat: OSC settings

* feat: HTTP settings
This commit is contained in:
Carlos Valente
2024-02-11 21:25:23 +01:00
committed by GitHub
parent 5355e45b80
commit 53963a9ad7
52 changed files with 742 additions and 1620 deletions
@@ -1,23 +1,22 @@
import got from 'got';
import { HttpSettings, HttpSubscription, HttpSubscriptionOptions, LogOrigin } from 'ontime-types';
import { HttpSettings, HttpSubscription, LogOrigin } from 'ontime-types';
import IIntegration, { TimerLifeCycleKey } from './IIntegration.js';
import { parseTemplateNested } from './integrationUtils.js';
import { dbModel } from '../../models/dataModel.js';
import { logger } from '../../classes/Logger.js';
import { validateHttpSubscriptionObject } from '../../utils/parserFunctions.js';
type Action = TimerLifeCycleKey | string;
/**
* @description Class contains logic towards outgoing HTTP communications
* @class
*/
export class HttpIntegration implements IIntegration<HttpSubscriptionOptions> {
subscriptions: HttpSubscription;
export class HttpIntegration implements IIntegration<HttpSubscription> {
subscriptions: HttpSubscription[];
enabled: boolean;
constructor() {
this.subscriptions = dbModel.http.subscriptions;
this.subscriptions = [];
this.enabled = false;
}
/**
@@ -25,64 +24,40 @@ export class HttpIntegration implements IIntegration<HttpSubscriptionOptions> {
*/
init(config: HttpSettings) {
const { subscriptions, enabledOut } = config;
if (!enabledOut) {
return {
success: false,
message: 'HTTP output disabled',
};
}
this.initSubscriptions(subscriptions);
return {
success: true,
message: 'HTTP integration client ready',
};
this.enabled = enabledOut;
}
initSubscriptions(subscriptionOptions: HttpSubscription) {
if (validateHttpSubscriptionObject(subscriptionOptions)) {
this.subscriptions = { ...subscriptionOptions };
}
initSubscriptions(subscriptions: HttpSubscription[]) {
this.subscriptions = subscriptions;
}
dispatch(action: Action, state?: object) {
if (!action) {
return {
success: false,
message: 'HTTP called with no action',
};
dispatch(action: TimerLifeCycleKey, state?: object) {
// noop
if (!this.enabled || !action) {
return;
}
// check subscriptions for action
const eventSubscriptions = this.subscriptions?.[action] || [];
eventSubscriptions.forEach((sub) => {
const { enabled, message } = sub;
if (enabled && message) {
const parsedMessage = parseTemplateNested(message, state || {});
try {
const parsedUrl = new URL(parsedMessage);
this.emit(parsedUrl);
} catch (err) {
logger.error(LogOrigin.Tx, `HTTP Integration: ${err}`);
return {
success: false,
message: `${err}`,
};
}
for (let i = 0; i < this.subscriptions.length; i++) {
const { cycle, message, enabled } = this.subscriptions[i];
if (cycle !== action || !enabled || !message) {
continue;
}
});
const parsedMessage = parseTemplateNested(message, state || {});
try {
const parsedUrl = new URL(parsedMessage);
this.emit(parsedUrl);
} catch (error) {
logger.error(LogOrigin.Tx, `HTTP Integration: ${error}`);
}
}
}
async emit(path: URL) {
try {
await got.get(path, {
retry: { limit: 0 },
});
} catch (err) {
logger.error(LogOrigin.Tx, `HTTP integration: ${err}`);
}
await got.get(path, {
retry: { limit: 0 },
});
}
shutdown() {}
@@ -1,24 +1,11 @@
import { TimerLifeCycle, Subscription } from 'ontime-types';
import { TimerLifeCycle } from 'ontime-types';
export type TimerLifeCycleKey = keyof typeof TimerLifeCycle;
export default interface IIntegration<T> {
subscriptions: Subscription<T>;
init: (config: unknown) => OperationReturn;
dispatch: (action: TimerLifeCycleKey, state?: object) => OperationReturn;
subscriptions: T[];
init: (config: unknown) => void;
dispatch: (action: TimerLifeCycleKey, state?: object) => void;
emit: (...args: unknown[]) => unknown;
shutdown: () => void;
}
// either went well, or explain what failed
type OperationReturn = ReturnOnSuccess | ReturnOnError;
type ReturnOnSuccess = {
success: true;
message?: string;
};
type ReturnOnError = {
success: false;
message: string;
};
@@ -1,5 +1,8 @@
import { LogOrigin } from 'ontime-types';
import IIntegration, { TimerLifeCycleKey } from './IIntegration.js';
import { eventStore } from '../../stores/EventStore.js';
import { logger } from '../../classes/Logger.js';
class IntegrationService {
private integrations: IIntegration<unknown>[];
@@ -24,7 +27,7 @@ class IntegrationService {
}
shutdown() {
console.log('Shutdown integrations');
logger.info(LogOrigin.Tx, `Shutdown Integrations`);
this.integrations.forEach((integration) => {
integration.shutdown();
});
@@ -1,25 +1,28 @@
import { ArgumentType, Client, Message } from 'node-osc';
import { OSCSettings, OscSubscription, OscSubscriptionOptions } from 'ontime-types';
import { LogOrigin, MaybeNumber, MaybeString, OSCSettings, OscSubscription } from 'ontime-types';
import IIntegration, { TimerLifeCycleKey } from './IIntegration.js';
import { parseTemplateNested } from './integrationUtils.js';
import { isObject } from '../../utils/varUtils.js';
import { dbModel } from '../../models/dataModel.js';
import { validateOscSubscriptionObject } from '../../utils/parserFunctions.js';
type Action = TimerLifeCycleKey | string;
import { logger } from '../../classes/Logger.js';
/**
* @description Class contains logic towards outgoing OSC communications
* @class
*/
export class OscIntegration implements IIntegration<OscSubscriptionOptions> {
export class OscIntegration implements IIntegration<OscSubscription> {
protected oscClient: null | Client;
subscriptions: OscSubscription;
subscriptions: OscSubscription[];
targetIP: MaybeString;
portOut: MaybeNumber;
enabledOut: boolean;
constructor() {
this.oscClient = null;
this.subscriptions = dbModel.osc.subscriptions;
this.subscriptions = [];
this.targetIP = null;
this.portOut = null;
this.enabledOut = false;
}
/**
@@ -27,75 +30,58 @@ export class OscIntegration implements IIntegration<OscSubscriptionOptions> {
*/
init(config: OSCSettings) {
const { targetIP, portOut, subscriptions, enabledOut } = config;
if (!enabledOut) {
this.oscClient?.close();
return {
success: false,
message: 'OSC output disabled',
};
}
this.initSubscriptions(subscriptions);
// runtime validation
const validateType = typeof targetIP !== 'string' || typeof portOut !== 'number';
const validateNull = !targetIP || !portOut;
if (validateType || validateNull) {
return {
success: false,
message: 'Config options incorrect',
};
if (!enabledOut && this.enabledOut) {
this.targetIP = targetIP;
this.portOut = portOut;
this.enabledOut = enabledOut;
this.shutdown();
return;
}
if (this.oscClient && targetIP === this.targetIP && portOut === this.portOut) {
// nothing changed that would mean we need a new client
return;
}
this.targetIP = targetIP;
this.portOut = portOut;
this.enabledOut = enabledOut;
try {
// this allows re-calling the init function during runtime
this.oscClient?.close();
logger.info(LogOrigin.Tx, 'Initialising OSC integration...');
this.oscClient = new Client(targetIP, portOut);
return {
success: true,
message: `OSC integration client connected to ${targetIP}:${portOut}`,
};
} catch (error) {
this.oscClient = null;
return {
success: false,
message: `Failed initialising OSC Client: ${error}`,
};
throw new Error(`Failed initialising OSC client: ${error}`);
}
return `OSC integration client connected to ${targetIP}:${portOut}`;
}
initSubscriptions(subscriptionOptions: OscSubscription) {
if (validateOscSubscriptionObject(subscriptionOptions)) {
this.subscriptions = { ...subscriptionOptions };
}
initSubscriptions(subscriptions: OscSubscription[]) {
this.subscriptions = subscriptions;
}
dispatch(action: Action, state?: object) {
if (!this.oscClient) {
return {
success: false,
message: 'Client not initialised',
};
dispatch(action: TimerLifeCycleKey, state?: object) {
// noop
if (!this.oscClient || !action) {
return;
}
if (!action) {
return {
success: false,
message: 'OSC called with no action',
};
}
// check subscriptions for action
const eventSubscriptions = this.subscriptions?.[action] || [];
eventSubscriptions.forEach((sub) => {
const { enabled, message } = sub;
if (enabled && message) {
const parsedMessage = parseTemplateNested(message, state || {});
this.emit(parsedMessage);
for (let i = 0; i < this.subscriptions.length; i++) {
const { cycle, message, enabled } = this.subscriptions[i];
if (cycle !== action || !enabled || !message) {
continue;
}
});
const parsedMessage = parseTemplateNested(message, state || {});
try {
this.emit(parsedMessage);
} catch (error) {
logger.error(LogOrigin.Tx, `OSC Integration: ${error}`);
}
}
}
emit(path: string, payload?: ArgumentType) {
@@ -105,33 +91,18 @@ export class OscIntegration implements IIntegration<OscSubscriptionOptions> {
const message = new Message(path);
if (payload) {
try {
if (isObject(payload)) {
message.append(JSON.stringify(payload));
} else {
message.append(payload);
}
} catch (error) {
console.log('OSC ERROR', error, payload);
if (isObject(payload)) {
message.append(JSON.stringify(payload));
} else {
message.append(payload);
}
}
this.oscClient.send(message, (error) => {
if (error) {
return {
success: false,
message: `Error sending message: ${JSON.stringify(error)}`,
};
}
return {
success: true,
message: 'OSC Message sent',
};
});
this.oscClient.send(message);
}
shutdown() {
console.log('Shutting down OSC integration');
logger.info(LogOrigin.Tx, 'Shutting down OSC integration');
if (this.oscClient) {
this.oscClient?.close();
this.oscClient = null;