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
7 changes: 7 additions & 0 deletions .claude/settings.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
{
"attribution": {
"commit": "",
"pr": "",
"sessionUrl": false
}
}
4 changes: 3 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -51,4 +51,6 @@ codecov
# Don't commit MacOS screenshots
*-darwin.png
/.run/All Tests.run.xml
.claude/
# Personal Claude Code files stay local; the shared project settings are committed.
.claude/*
!.claude/settings.json
124 changes: 110 additions & 14 deletions src/api/telemetry/TelemetryAPI.js
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ const SUBSCRIBE_STRATEGY = {
export default class TelemetryAPI {
#isGreedyLAD;
#subscribeCache;
#subscriptionObservers;
#hasReturnedFirstData;

get SUBSCRIBE_STRATEGY() {
Expand Down Expand Up @@ -110,6 +111,7 @@ export default class TelemetryAPI {
this.#isGreedyLAD = true;
this.BatchingWebSocket = BatchingWebSocket;
this.#subscribeCache = {};
this.#subscriptionObservers = new Set();
this.#hasReturnedFirstData = false;
}

Expand Down Expand Up @@ -385,6 +387,68 @@ export default class TelemetryAPI {
this.requestInterceptorRegistry.addInterceptor(requestInterceptorDef);
}

/**
* Observe every telemetry subscription in the application as it is created.
*
* The observer is invoked once per subscription, and NOT once per datum. This
* allows an observer to resolve whatever it needs from the domain object
* (metadata, value formatters) a single time and close over the result,
* instead of resolving it again for every sample that arrives.
*
* Observers are invoked for subscriptions that already exist at the time of
* registration, as well as for any created subsequently. An observer declines
* a subscription by returning undefined.
*
* Only telemetry delivered via subscriptions is observed. The results of
* historical requests are not.
*
* @param {function(import('openmct').DomainObject): (function(object): void | undefined)} observeSubscription
* invoked with each subscribed domain object. Returns a function to be
* called with each datum received for that object, or undefined to
* ignore the subscription.
* @returns {Function} a function which removes this observer from all
* subscriptions, both current and future
* @method addSubscriptionObserver
*/
addSubscriptionObserver(observeSubscription) {
this.#subscriptionObservers.add(observeSubscription);

Object.values(this.#subscribeCache).forEach((subscriber) => {
this.#attachObserverToSubscriber(observeSubscription, subscriber);
});

return () => this.#removeSubscriptionObserver(observeSubscription);
}

/**
* Offer a single subscription to a single observer, recording the datum
* handler it returns. The registration function is recorded alongside the
* handler so that the observer can later be removed without a second lookup
* structure.
*
* @private
*/
#attachObserverToSubscriber(observeSubscription, subscriber) {
const observeDatum = observeSubscription(subscriber.domainObject);

if (observeDatum !== undefined) {
subscriber.datumObservers.push({ observeSubscription, observeDatum });
}
}

/**
* @private
*/
#removeSubscriptionObserver(observeSubscription) {
this.#subscriptionObservers.delete(observeSubscription);

Object.values(this.#subscribeCache).forEach((subscriber) => {
subscriber.datumObservers = subscriber.datumObservers.filter(
(datumObserver) => datumObserver.observeSubscription !== observeSubscription
);
});
}

/**
* Retrieve the request interceptors for a given domain object.
* @private
Expand Down Expand Up @@ -554,8 +618,19 @@ export default class TelemetryAPI {
if (!subscriber) {
subscriber = this.#subscribeCache[cacheKey] = {
latestCallbacks: [],
batchCallbacks: []
batchCallbacks: [],
datumObservers: [],
// Retained so that observers registered after this subscription was
// created can still be offered it. Where several callers share a cache
// key, this is the first caller's instance. That is harmless: metadata
// is cached per instance in a WeakMap and does not change at runtime.
domainObject
};

this.#subscriptionObservers.forEach((observeSubscription) => {
this.#attachObserverToSubscriber(observeSubscription, subscriber);
});

if (provider) {
subscriber.unsubscribe = provider.subscribe(
domainObject,
Expand All @@ -574,9 +649,40 @@ export default class TelemetryAPI {
}

// Guarantees that view receive telemetry in the expected form
//
// Subscription observers are notified BEFORE the view callbacks. A
// telemetry driven clock ticks from an observer, and a tick moves the time
// bounds that TelemetryCollection filters incoming data against. Notifying
// views first would mean the very first datum of every stream is discarded
// as out of bounds, because a tick trims data rather than re-requesting it.
function invokeCallbackWithRequestedStrategy(data) {
const latestDatum = getLatestDatum(data);

notifySubscriptionObservers(latestDatum, subscriber.datumObservers);
invokeCallbacksWithArray(data, subscriber.batchCallbacks);
invokeCallbacksWithSingleValue(data, subscriber.latestCallbacks);
invokeCallbacksWithSingleValue(latestDatum, subscriber.latestCallbacks);
}

function getLatestDatum(data) {
const latestDatum = Array.isArray(data) ? data[data.length - 1] : data;

if (latestDatum === undefined || latestDatum === null) {
throw new Error(
'Attempt to invoke telemetry subscription callback with no telemetry datum'
);
}

return latestDatum;
}

function notifySubscriptionObservers(latestDatum, datumObservers) {
if (datumObservers.length === 0) {
return;
}

datumObservers.forEach((datumObserver) => {
datumObserver.observeDatum(latestDatum);
});
}

function invokeCallbacksWithArray(data, batchCallbacks) {
Expand All @@ -596,19 +702,9 @@ export default class TelemetryAPI {
});
}

function invokeCallbacksWithSingleValue(data, latestCallbacks) {
if (Array.isArray(data)) {
data = data[data.length - 1];
}

if (data === undefined || data === null) {
throw new Error(
'Attempt to invoke telemetry subscription callback with no telemetry datum'
);
}

function invokeCallbacksWithSingleValue(latestDatum, latestCallbacks) {
latestCallbacks.forEach((cb) => {
cb(data);
cb(latestDatum);
});
}

Expand Down
123 changes: 123 additions & 0 deletions src/api/telemetry/TelemetryAPISpec.js
Original file line number Diff line number Diff line change
Expand Up @@ -453,6 +453,129 @@ describe('Telemetry API', () => {
expect(telemetryProvider.subscribe.calls.mostRecent().args[2].strategy).toBe('latest');
});
});

describe('subscription observers', () => {
let providerCallbacks;
let observeDatum;
let observeSubscription;

function emitTelemetry(data) {
providerCallbacks.forEach((providerCallback) => {
providerCallback(data);
});
}

beforeEach(() => {
providerCallbacks = [];
telemetryProvider.supportsSubscribe.and.returnValue(true);
telemetryProvider.subscribe.and.callFake((obj, providerCallback) => {
providerCallbacks.push(providerCallback);

return jasmine.createSpy('unsubscribe');
});
telemetryAPI.addProvider(telemetryProvider);

observeDatum = jasmine.createSpy('observeDatum');
observeSubscription = jasmine
.createSpy('observeSubscription')
.and.returnValue(observeDatum);
});

it('invokes the observer once with the domain object, when a subscription is created', () => {
telemetryAPI.addSubscriptionObserver(observeSubscription);
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callback'));

expect(observeSubscription).toHaveBeenCalledOnceWith(domainObject);
});

it('invokes the observer for subscriptions that already exist', () => {
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callback'));
telemetryAPI.addSubscriptionObserver(observeSubscription);

expect(observeSubscription).toHaveBeenCalledOnceWith(domainObject);

emitTelemetry({ value: 1 });

expect(observeDatum).toHaveBeenCalledWith({ value: 1 });
});

it('does not invoke the observer again when a second callback joins a subscription', () => {
telemetryAPI.addSubscriptionObserver(observeSubscription);
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callbackOne'));
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callbackTwo'));

expect(observeSubscription).toHaveBeenCalledTimes(1);
});

it('invokes the datum observer with each datum received', () => {
telemetryAPI.addSubscriptionObserver(observeSubscription);
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callback'));

emitTelemetry({ value: 1 });
emitTelemetry({ value: 2 });

expect(observeDatum).toHaveBeenCalledTimes(2);
expect(observeDatum).toHaveBeenCalledWith({ value: 2 });
});

it('invokes the datum observer with only the last element of a batch', () => {
telemetryAPI.addSubscriptionObserver(observeSubscription);
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callback'));

emitTelemetry([{ value: 1 }, { value: 2 }, { value: 3 }]);

expect(observeDatum).toHaveBeenCalledOnceWith({ value: 3 });
});

it('never invokes an observer that declined the subscription', () => {
const decliningObserver = jasmine.createSpy('decliningObserver').and.returnValue(undefined);

telemetryAPI.addSubscriptionObserver(decliningObserver);
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callback'));

emitTelemetry({ value: 1 });

expect(decliningObserver).toHaveBeenCalledTimes(1);
});

it('stops delivering to a removed observer, leaving others intact', () => {
const otherObserveDatum = jasmine.createSpy('otherObserveDatum');
const viewCallback = jasmine.createSpy('viewCallback');

const removeObserver = telemetryAPI.addSubscriptionObserver(observeSubscription);
telemetryAPI.addSubscriptionObserver(() => otherObserveDatum);
telemetryAPI.subscribe(domainObject, viewCallback);

removeObserver();
emitTelemetry({ value: 1 });

expect(observeDatum).not.toHaveBeenCalled();
expect(otherObserveDatum).toHaveBeenCalledWith({ value: 1 });
expect(viewCallback).toHaveBeenCalledWith({ value: 1 });
});

it('delivers to view callbacks normally when no observers are registered', () => {
const latestCallback = jasmine.createSpy('latestCallback');

telemetryAPI.subscribe(domainObject, latestCallback);

emitTelemetry([{ value: 1 }, { value: 2 }]);

expect(latestCallback).toHaveBeenCalledOnceWith({ value: 2 });
});

it('throws when telemetry is received with no datum', () => {
telemetryAPI.addSubscriptionObserver(observeSubscription);
telemetryAPI.subscribe(domainObject, jasmine.createSpy('callback'));

expect(() => emitTelemetry(undefined)).toThrowError(
'Attempt to invoke telemetry subscription callback with no telemetry datum'
);
expect(() => emitTelemetry([])).toThrowError(
'Attempt to invoke telemetry subscription callback with no telemetry datum'
);
});
});
});

describe('metadata', () => {
Expand Down
43 changes: 0 additions & 43 deletions src/plugins/latestDataClock/LADClock.js

This file was deleted.

2 changes: 2 additions & 0 deletions src/plugins/plugins.js
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ import RemoteClock from './remoteClock/plugin.js';
import StaticRootPlugin from './staticRootPlugin/plugin.js';
import SummaryWidget from './summaryWidget/plugin.js';
import Tabs from './tabs/plugin.js';
import TelemetryClock from './telemetryClock/plugin.js';
import TelemetryMean from './telemetryMean/plugin.js';
import TelemetryTablePlugin from './telemetryTable/plugin.js';
import DarkMatter from './themes/darkmatter.js';
Expand Down Expand Up @@ -109,6 +110,7 @@ plugins.example.ExampleStaleness = ExampleStaleness;
plugins.UTCTimeSystem = UTCTimeSystem;
plugins.LocalTimeSystem = LocalTimeSystem;
plugins.RemoteClock = RemoteClock;
plugins.TelemetryClock = TelemetryClock;

plugins.MyItems = MyItems;

Expand Down
Loading
Loading