diff --git a/.claude/settings.json b/.claude/settings.json new file mode 100644 index 0000000000..b7f9b0a629 --- /dev/null +++ b/.claude/settings.json @@ -0,0 +1,7 @@ +{ + "attribution": { + "commit": "", + "pr": "", + "sessionUrl": false + } +} diff --git a/.gitignore b/.gitignore index c19cbfa2b4..d78cd80ef2 100644 --- a/.gitignore +++ b/.gitignore @@ -51,4 +51,6 @@ codecov # Don't commit MacOS screenshots *-darwin.png /.run/All Tests.run.xml -.claude/ \ No newline at end of file +# Personal Claude Code files stay local; the shared project settings are committed. +.claude/* +!.claude/settings.json diff --git a/src/api/telemetry/TelemetryAPI.js b/src/api/telemetry/TelemetryAPI.js index 944e467fc3..7c97cf3d2a 100644 --- a/src/api/telemetry/TelemetryAPI.js +++ b/src/api/telemetry/TelemetryAPI.js @@ -82,6 +82,7 @@ const SUBSCRIBE_STRATEGY = { export default class TelemetryAPI { #isGreedyLAD; #subscribeCache; + #subscriptionObservers; #hasReturnedFirstData; get SUBSCRIBE_STRATEGY() { @@ -110,6 +111,7 @@ export default class TelemetryAPI { this.#isGreedyLAD = true; this.BatchingWebSocket = BatchingWebSocket; this.#subscribeCache = {}; + this.#subscriptionObservers = new Set(); this.#hasReturnedFirstData = false; } @@ -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 @@ -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, @@ -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) { @@ -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); }); } diff --git a/src/api/telemetry/TelemetryAPISpec.js b/src/api/telemetry/TelemetryAPISpec.js index 648b1172ca..3ccc547599 100644 --- a/src/api/telemetry/TelemetryAPISpec.js +++ b/src/api/telemetry/TelemetryAPISpec.js @@ -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', () => { diff --git a/src/plugins/latestDataClock/LADClock.js b/src/plugins/latestDataClock/LADClock.js deleted file mode 100644 index 7490d3f901..0000000000 --- a/src/plugins/latestDataClock/LADClock.js +++ /dev/null @@ -1,43 +0,0 @@ -/***************************************************************************** - * Open MCT Web, Copyright (c) 2014-2024, United States Government - * as represented by the Administrator of the National Aeronautics and Space - * Administration. All rights reserved. - * - * Open MCT Web is licensed under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * http://www.apache.org/licenses/LICENSE-2.0. - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT - * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the - * License for the specific language governing permissions and limitations - * under the License. - * - * Open MCT Web includes source code licensed under additional open source - * licenses. See the Open Source Licenses file (LICENSES.md) included with - * this source code distribution or the Licensing information page available - * at runtime from the About dialog for additional information. - *****************************************************************************/ - -import LocalClock from '../../../src/plugins/utcTimeSystem/LocalClock.js'; - -class LADClock extends LocalClock { - /** - * A {@link Clock} that mocks a "latest available data" type tick source. - * This is for testing purposes only, and behaves identically to a local clock. - * It DOES NOT tick on receipt of data. - * @constructor - */ - constructor(period) { - super(period); - - this.key = 'test-lad'; - this.mode = 'lad'; - this.cssClass = 'icon-suitcase'; - this.name = 'Latest available data'; - this.description = 'Updates when new data is available'; - } -} - -export default LADClock; diff --git a/src/plugins/plugins.js b/src/plugins/plugins.js index 938bc51c09..d4ffdf774f 100644 --- a/src/plugins/plugins.js +++ b/src/plugins/plugins.js @@ -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'; @@ -109,6 +110,7 @@ plugins.example.ExampleStaleness = ExampleStaleness; plugins.UTCTimeSystem = UTCTimeSystem; plugins.LocalTimeSystem = LocalTimeSystem; plugins.RemoteClock = RemoteClock; +plugins.TelemetryClock = TelemetryClock; plugins.MyItems = MyItems; diff --git a/src/plugins/remoteClock/RemoteClock.js b/src/plugins/remoteClock/RemoteClock.js index b7c6407855..315447345f 100644 --- a/src/plugins/remoteClock/RemoteClock.js +++ b/src/plugins/remoteClock/RemoteClock.js @@ -19,8 +19,8 @@ * this source code distribution or the Licensing information page available * at runtime from the About dialog for additional information. *****************************************************************************/ +import clockReadyRequestInterceptor from '../../utils/clock/clockReadyRequestInterceptor.js'; import DefaultClock from '../../utils/clock/DefaultClock.js'; -import remoteClockRequestInterceptor from './requestInterceptor.js'; /** * A {@link openmct.TimeAPI.Clock} that updates the temporal bounds of the @@ -49,11 +49,11 @@ export default class RemoteClock extends DefaultClock { this.formatTime = undefined; this.metadata = undefined; + // A value of zero means "has not ticked yet". The request interceptor + // relies on this to know when the clock has a real time to report. this.lastTick = 0; - this.openmct.telemetry.addRequestInterceptor( - remoteClockRequestInterceptor(this.openmct, this.identifier, this.#waitForReady.bind(this)) - ); + this.openmct.telemetry.addRequestInterceptor(clockReadyRequestInterceptor(this.openmct, this)); this._processDatum = this._processDatum.bind(this); } @@ -146,26 +146,4 @@ export default class RemoteClock extends DefaultClock { return timeFormatter.format(datum); }; } - - /** - * Waits for the clock to have a non-default tick value. - */ - #waitForReady() { - const waitForInitialTick = (resolve) => { - const tickListener = () => { - if (this.lastTick > 0) { - const offsets = this.openmct.time.getClockOffsets(); - this.openmct.time.off('tick', tickListener); // Unregister the tick listener - resolve({ - start: this.lastTick + offsets.start, - end: this.lastTick + offsets.end - }); - } - }; - - this.openmct.time.on('tick', tickListener); - }; - - return new Promise(waitForInitialTick); - } } diff --git a/src/plugins/remoteClock/requestInterceptor.js b/src/plugins/remoteClock/requestInterceptor.js deleted file mode 100644 index f4a075970f..0000000000 --- a/src/plugins/remoteClock/requestInterceptor.js +++ /dev/null @@ -1,77 +0,0 @@ -/***************************************************************************** - * Open MCT, Copyright (c) 2014-2024, United States Government - * as represented by the Administrator of the National Aeronautics and Space - * Administration. All rights reserved. - * - * Open MCT is licensed under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * http://www.apache.org/licenses/LICENSE-2.0. - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT - * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the - * License for the specific language governing permissions and limitations - * under the License. - * - * Open MCT includes source code licensed under additional open source - * licenses. See the Open Source Licenses file (LICENSES.md) included with - * this source code distribution or the Licensing information page available - * at runtime from the About dialog for additional information. - *****************************************************************************/ - -/** - * Intercepts requests to ensure the remote clock is ready. - * - * @param {import('../../openmct').OpenMCT} openmct - The OpenMCT instance. - * @param {import('../../openmct').Identifier} _remoteClockIdentifier - The identifier for the remote clock. - * @param {Function} waitForBounds - A function that returns a promise resolving to the initial bounds. - * @returns {Object} The request interceptor. - */ -function remoteClockRequestInterceptor(openmct, _remoteClockIdentifier, waitForBounds) { - let remoteClockLoaded = false; - - return { - /** - * Determines if the interceptor applies to the given request. - * - * @param {Object} _ - Unused parameter. - * @param {import('../../api/telemetry/TelemetryAPI').TelemetryRequestOptions} request - The request object. - * @returns {boolean} True if the interceptor applies, false otherwise. - */ - appliesTo: (_, request) => { - // Get the activeClock from the Global Time Context - /** @type {import("../../api/time/TimeContext").default} */ - const { activeClock } = openmct.time; - - // this type of request does not rely on clock having bounds - if (request.strategy === 'latest' && request.timeContext.isRealTime()) { - return false; - } - - return activeClock?.key === 'remote-clock' && !remoteClockLoaded; - }, - /** - * Invokes the interceptor to modify the request. - * - * @param {Object} request - The request object. - * @returns {Promise} The modified request object. - */ - invoke: async (request) => { - const timeContext = request?.timeContext ?? openmct.time; - - // Wait for initial bounds if the request is for real-time data. - // Otherwise, use the bounds provided by the request. - if (timeContext.isRealTime()) { - const { start, end } = await waitForBounds(); - remoteClockLoaded = true; - request.start = start; - request.end = end; - } - - return request; - } - }; -} - -export default remoteClockRequestInterceptor; diff --git a/src/plugins/telemetryClock/TelemetryClock.js b/src/plugins/telemetryClock/TelemetryClock.js new file mode 100644 index 0000000000..941cecdac5 --- /dev/null +++ b/src/plugins/telemetryClock/TelemetryClock.js @@ -0,0 +1,170 @@ +/***************************************************************************** + * Open MCT, Copyright (c) 2014-2024, United States Government + * as represented by the Administrator of the National Aeronautics and Space + * Administration. All rights reserved. + * + * Open MCT is licensed under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0. + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations + * under the License. + * + * Open MCT includes source code licensed under additional open source + * licenses. See the Open Source Licenses file (LICENSES.md) included with + * this source code distribution or the Licensing information page available + * at runtime from the About dialog for additional information. + *****************************************************************************/ + +import throttle from 'lodash/throttle'; + +import { TIME_CONTEXT_EVENTS } from '../../api/time/constants.js'; +import clockReadyRequestInterceptor from '../../utils/clock/clockReadyRequestInterceptor.js'; +import DefaultClock from '../../utils/clock/DefaultClock.js'; + +const DEFAULT_TICK_PERIOD_MILLISECONDS = 100; +const NO_TICK_YET = 0; + +/** + * A {@link openmct.TimeAPI.Clock} that derives the current time from the + * timestamps of all telemetry arriving anywhere in the application, rather + * than from the local system clock. + * + * This allows Open MCT to follow a time source that is not wall clock aligned, + * such as a simulation, a replay, or a test bed running faster or slower than + * real time. + * + * The clock is monotonic. Only a timestamp strictly newer than the last one + * emitted will advance it, so out of order telemetry cannot drag time + * backwards. It does not tick at all until the first datum arrives, and never + * falls back to the local system clock. + * + * @param {import('../../openmct').OpenMCT} openmct + * @param {number} [tickPeriod] the minimum interval, in milliseconds, between + * ticks. Telemetry arriving within a single period is coalesced into + * one tick carrying the newest timestamp seen. + * @constructor + */ +export default class TelemetryClock extends DefaultClock { + #openmct; + #newestTimestampSeen; + #tickWithNewestTimestamp; + #stopObservingSubscriptions; + #restartOnTimeSystemChange; + + constructor(openmct, tickPeriod = DEFAULT_TICK_PERIOD_MILLISECONDS) { + super(); + + this.key = 'telemetry-clock'; + this.name = 'Telemetry Clock'; + this.cssClass = 'icon-telemetry'; + this.description = + 'Ticks from the timestamps of all incoming telemetry. Does not tick until telemetry arrives.'; + + this.#openmct = openmct; + + // A value of zero means "has not ticked yet". Both the request interceptor + // and TimeContext.now() rely on this. + this.lastTick = NO_TICK_YET; + this.#newestTimestampSeen = NO_TICK_YET; + + this.#tickWithNewestTimestamp = throttle(this.#emitNewestTimestamp.bind(this), tickPeriod); + this.#restartOnTimeSystemChange = this.#restartForNewTimeSystem.bind(this); + + this.#openmct.telemetry.addRequestInterceptor( + clockReadyRequestInterceptor(this.#openmct, this) + ); + } + + start() { + this.#observeAllTelemetrySubscriptions(); + this.#openmct.time.on(TIME_CONTEXT_EVENTS.timeSystem, this.#restartOnTimeSystemChange); + } + + stop() { + this.#openmct.time.off(TIME_CONTEXT_EVENTS.timeSystem, this.#restartOnTimeSystemChange); + this.#tickWithNewestTimestamp.cancel(); + this.#stopObservingSubscriptions(); + this.#stopObservingSubscriptions = undefined; + + // lastTick is deliberately left alone. Monotonicity should survive + // switching away to another clock and back again. + } + + /** + * Watch every telemetry subscription in the application, present and future, + * for timestamps. + */ + #observeAllTelemetrySubscriptions() { + this.#stopObservingSubscriptions = this.#openmct.telemetry.addSubscriptionObserver( + (domainObject) => this.#observeTimestampsFrom(domainObject) + ); + } + + /** + * Resolve the value formatter for the active time system once per + * subscription and close over it, so that no lookup is required for each + * datum that arrives. + * + * Returns undefined for objects that carry no timestamp in the active time + * system. Those objects are reconsidered if the time system changes. + */ + #observeTimestampsFrom(domainObject) { + const metadata = this.#openmct.telemetry.getMetadata(domainObject); + const timeSystemKey = this.#openmct.time.getTimeSystem().key; + const timestampMetadata = metadata.value(timeSystemKey); + + if (timestampMetadata === undefined) { + return undefined; + } + + const timestampFormatter = this.#openmct.telemetry.getValueFormatter(timestampMetadata); + + return (datum) => this.#observeTimestamp(timestampFormatter.parse(datum)); + } + + /** + * Timestamps in a new time system are on a different scale and epoch + * entirely, so the time we have already reported is meaningless and the + * monotonic gate would block forever. Return to the un-ticked state and wait + * for data exactly as at startup, which also re-arms the request interceptor. + * + * Rebuilding every observer, rather than refreshing a list of them, means + * objects that carried no timestamp in the previous time system are + * reconsidered under the new one. + */ + #restartForNewTimeSystem() { + this.#tickWithNewestTimestamp.cancel(); + + this.lastTick = NO_TICK_YET; + this.#newestTimestampSeen = NO_TICK_YET; + + this.#stopObservingSubscriptions(); + this.#observeAllTelemetrySubscriptions(); + } + + /** + * The monotonic gate. A datum carrying no value for the active time system + * parses to NaN, and every comparison against NaN is false, so it is ignored + * without needing a guard. + */ + #observeTimestamp(timestamp) { + if (timestamp > this.#newestTimestampSeen) { + this.#newestTimestampSeen = timestamp; + this.#tickWithNewestTimestamp(); + } + } + + /** + * Reads the accumulated timestamp rather than taking it as an argument, so + * that a coalesced tick always carries the newest value seen rather than + * whichever one happened to trigger it. + */ + #emitNewestTimestamp() { + this.tick(this.#newestTimestampSeen); + } +} diff --git a/src/plugins/telemetryClock/TelemetryClockSpec.js b/src/plugins/telemetryClock/TelemetryClockSpec.js new file mode 100644 index 0000000000..045c5f5b37 --- /dev/null +++ b/src/plugins/telemetryClock/TelemetryClockSpec.js @@ -0,0 +1,304 @@ +/***************************************************************************** + * Open MCT, Copyright (c) 2014-2024, United States Government + * as represented by the Administrator of the National Aeronautics and Space + * Administration. All rights reserved. + * + * Open MCT is licensed under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0. + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations + * under the License. + * + * Open MCT includes source code licensed under additional open source + * licenses. See the Open Source Licenses file (LICENSES.md) included with + * this source code distribution or the Licensing information page available + * at runtime from the About dialog for additional information. + *****************************************************************************/ + +import { createOpenMct, resetApplicationState } from 'utils/testing'; + +const TELEMETRY_CLOCK_KEY = 'telemetry-clock'; +const TICK_PERIOD = 100; +const NO_TICK_YET = 0; + +describe('the TelemetryClock plugin', () => { + let openmct; + let telemetryClock; + let tickCallback; + let providerCallbacks; + let domainObject; + let timestampMetadataByTimeSystem; + + /** + * Delivers a datum to every active subscription, as a telemetry provider + * would. + */ + function receiveTelemetry(timestamp) { + providerCallbacks.forEach((providerCallback) => { + providerCallback({ timestamp }); + }); + } + + function subscribeToTelemetry() { + return openmct.telemetry.subscribe(domainObject, () => {}); + } + + function activateTelemetryClock() { + openmct.time.setClock(TELEMETRY_CLOCK_KEY); + } + + beforeEach((done) => { + providerCallbacks = []; + domainObject = { + name: 'some-telemetry', + type: 'sample-type', + identifier: { namespace: 'test', key: 'telemetry' } + }; + + // Objects report a timestamp in UTC only, so that the time system change + // behaviour can be exercised by switching to a system they do not support. + timestampMetadataByTimeSystem = { utc: { key: 'utc', source: 'timestamp' } }; + + openmct = createOpenMct(); + openmct.install(openmct.plugins.TelemetryClock(TICK_PERIOD)); + openmct.on('start', done); + openmct.startHeadless(); + }); + + afterEach(() => { + return resetApplicationState(openmct); + }); + + describe('once activated', () => { + beforeEach(() => { + openmct.telemetry.addProvider({ + supportsSubscribe: () => true, + subscribe: (object, providerCallback) => { + providerCallbacks.push(providerCallback); + + return () => { + providerCallbacks = providerCallbacks.filter( + (candidate) => candidate !== providerCallback + ); + }; + } + }); + + spyOn(openmct.telemetry, 'getMetadata').and.returnValue({ + value: (timeSystemKey) => timestampMetadataByTimeSystem[timeSystemKey] + }); + spyOn(openmct.telemetry, 'getValueFormatter').and.returnValue({ + parse: (datum) => datum.timestamp + }); + + telemetryClock = openmct.time + .getAllClocks() + .filter((clock) => clock.key === TELEMETRY_CLOCK_KEY)[0]; + + tickCallback = jasmine.createSpy('tickCallback'); + openmct.time.on('tick', tickCallback); + + jasmine.clock().install(); + jasmine.clock().mockDate(); + }); + + afterEach(() => { + jasmine.clock().uninstall(); + openmct.time.setClock('local'); + }); + + it('is registered with the expected identity', () => { + expect(telemetryClock.key).toEqual(TELEMETRY_CLOCK_KEY); + expect(telemetryClock.name).toEqual('Telemetry Clock'); + }); + + it('does not tick before any telemetry has arrived', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + jasmine.clock().tick(TICK_PERIOD * 10); + + expect(tickCallback).not.toHaveBeenCalled(); + expect(telemetryClock.currentValue()).toEqual(NO_TICK_YET); + }); + + it('ticks immediately on the first datum, without waiting for the period', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + + expect(tickCallback).toHaveBeenCalledOnceWith(1000); + expect(telemetryClock.currentValue()).toEqual(1000); + }); + + it('coalesces a burst into a single tick carrying the newest timestamp', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + receiveTelemetry(1010); + receiveTelemetry(1050); + receiveTelemetry(1030); + + expect(tickCallback).toHaveBeenCalledTimes(1); + + jasmine.clock().tick(TICK_PERIOD); + + expect(tickCallback).toHaveBeenCalledTimes(2); + expect(tickCallback).toHaveBeenCalledWith(1050); + }); + + it('never goes backwards', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD); + + receiveTelemetry(500); + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD * 5); + + expect(tickCallback).toHaveBeenCalledTimes(1); + expect(telemetryClock.currentValue()).toEqual(1000); + }); + + it('stops ticking, and leaves no timer armed, when telemetry stops', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD * 10); + + const tickCountWhenTelemetryStopped = tickCallback.calls.count(); + + jasmine.clock().tick(TICK_PERIOD * 100); + + expect(tickCallback).toHaveBeenCalledTimes(tickCountWhenTelemetryStopped); + expect(telemetryClock.currentValue()).toEqual(1000); + }); + + it('ticks immediately again after an idle gap', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD * 10); + tickCallback.calls.reset(); + + receiveTelemetry(2000); + + expect(tickCallback).toHaveBeenCalledOnceWith(2000); + }); + + it('observes subscriptions that already existed when it was activated', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + + expect(tickCallback).toHaveBeenCalledOnceWith(1000); + }); + + it('observes subscriptions created after it was activated', () => { + activateTelemetryClock(); + subscribeToTelemetry(); + + receiveTelemetry(1000); + + expect(tickCallback).toHaveBeenCalledOnceWith(1000); + }); + + it('ignores objects that carry no timestamp in the active time system', () => { + timestampMetadataByTimeSystem = {}; + + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD * 5); + + expect(tickCallback).not.toHaveBeenCalled(); + }); + + it('stops observing telemetry once another clock is selected', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD); + openmct.time.setClock('local'); + tickCallback.calls.reset(); + + receiveTelemetry(5000); + jasmine.clock().tick(TICK_PERIOD); + + expect(tickCallback).not.toHaveBeenCalledWith(5000); + }); + + it('retains the time it reached when switched away from and back', () => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD); + + openmct.time.setClock('local'); + activateTelemetryClock(); + + expect(telemetryClock.currentValue()).toEqual(1000); + + receiveTelemetry(500); + jasmine.clock().tick(TICK_PERIOD * 5); + + expect(telemetryClock.currentValue()).toEqual(1000); + }); + + describe('when the time system changes', () => { + beforeEach(() => { + subscribeToTelemetry(); + activateTelemetryClock(); + + receiveTelemetry(1000); + jasmine.clock().tick(TICK_PERIOD); + tickCallback.calls.reset(); + }); + + it('returns to its un-ticked state, because timestamps are not comparable across time systems', () => { + timestampMetadataByTimeSystem = { other: { key: 'other', source: 'timestamp' } }; + openmct.time.addTimeSystem({ + key: 'other', + name: 'Other', + timeFormat: 'utc', + durationFormat: 'duration', + isUTCBased: false + }); + openmct.time.setTimeSystem('other', { start: 0, end: 1 }); + + expect(telemetryClock.currentValue()).toEqual(NO_TICK_YET); + }); + + it('accepts a timestamp that would have been rejected under the previous time system', () => { + timestampMetadataByTimeSystem = { other: { key: 'other', source: 'timestamp' } }; + openmct.time.addTimeSystem({ + key: 'other', + name: 'Other', + timeFormat: 'utc', + durationFormat: 'duration', + isUTCBased: false + }); + openmct.time.setTimeSystem('other', { start: 0, end: 1 }); + + receiveTelemetry(10); + + expect(tickCallback).toHaveBeenCalledWith(10); + }); + }); + }); +}); diff --git a/src/plugins/latestDataClock/plugin.js b/src/plugins/telemetryClock/plugin.js similarity index 58% rename from src/plugins/latestDataClock/plugin.js rename to src/plugins/telemetryClock/plugin.js index 2219d6b3d4..918442120b 100644 --- a/src/plugins/latestDataClock/plugin.js +++ b/src/plugins/telemetryClock/plugin.js @@ -1,9 +1,9 @@ /***************************************************************************** - * Open MCT Web, Copyright (c) 2014-2024, United States Government + * Open MCT, Copyright (c) 2014-2024, United States Government * as represented by the Administrator of the National Aeronautics and Space * Administration. All rights reserved. * - * Open MCT Web is licensed under the Apache License, Version 2.0 (the + * Open MCT is licensed under the Apache License, Version 2.0 (the * "License"); you may not use this file except in compliance with the License. * You may obtain a copy of the License at * http://www.apache.org/licenses/LICENSE-2.0. @@ -14,16 +14,26 @@ * License for the specific language governing permissions and limitations * under the License. * - * Open MCT Web includes source code licensed under additional open source + * Open MCT includes source code licensed under additional open source * licenses. See the Open Source Licenses file (LICENSES.md) included with * this source code distribution or the Licensing information page available * at runtime from the About dialog for additional information. *****************************************************************************/ -import LADClock from './LADClock.js'; +import TelemetryClock from './TelemetryClock.js'; -export default function () { - return function (openmct) { - openmct.time.addClock(new LADClock()); +/** + * Installs a clock that ticks from the timestamps of all incoming telemetry, + * instead of from the local system clock. + * + * The clock must also be named by a time conductor menu option before it can + * be selected. + * + * @param {number} [tickPeriod] the minimum interval, in milliseconds, between + * ticks + */ +export default function (tickPeriod) { + return function install(openmct) { + openmct.time.addClock(new TelemetryClock(openmct, tickPeriod)); }; } diff --git a/src/utils/clock/clockReadyRequestInterceptor.js b/src/utils/clock/clockReadyRequestInterceptor.js new file mode 100644 index 0000000000..92d6a43e5f --- /dev/null +++ b/src/utils/clock/clockReadyRequestInterceptor.js @@ -0,0 +1,97 @@ +/***************************************************************************** + * Open MCT, Copyright (c) 2014-2024, United States Government + * as represented by the Administrator of the National Aeronautics and Space + * Administration. All rights reserved. + * + * Open MCT is licensed under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * http://www.apache.org/licenses/LICENSE-2.0. + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations + * under the License. + * + * Open MCT includes source code licensed under additional open source + * licenses. See the Open Source Licenses file (LICENSES.md) included with + * this source code distribution or the Licensing information page available + * at runtime from the About dialog for additional information. + *****************************************************************************/ + +const NO_TICK_YET = 0; + +/** + * Defers realtime requests until a telemetry driven clock has produced its + * first tick. + * + * Such clocks derive their time from incoming telemetry rather than the wall + * clock, so before any telemetry has arrived they have no time to report and + * the temporal bounds are meaningless. Rather than let views request data for + * an arbitrary window, this interceptor parks realtime requests until the + * clock knows what time it is, then rewrites them with real bounds. + * + * This does not deadlock. TelemetryCollection.load() starts its historical + * request without awaiting it, and then establishes its subscriptions + * synchronously. Subscriptions are therefore live while a request is parked + * here, so telemetry can still arrive and produce the tick that releases it. + * + * @param {import('../../openmct').OpenMCT} openmct + * @param {import('./DefaultClock').default} clock the clock whose readiness + * gates these requests + * @returns {import('../../api/telemetry/TelemetryAPI').RequestInterceptorDef} + */ +export default function clockReadyRequestInterceptor(openmct, clock) { + return { + appliesTo: (_identifier, request) => { + const { activeClock } = openmct.time; + + // This type of request does not rely on the clock having bounds. + if (request.strategy === 'latest' && request.timeContext.isRealTime()) { + return false; + } + + return activeClock?.key === clock.key && clock.currentValue() === NO_TICK_YET; + }, + invoke: async (request) => { + const timeContext = request?.timeContext ?? openmct.time; + + // Requests for fixed bounds already carry the bounds they want. + if (timeContext.isRealTime()) { + const firstTimestamp = await waitForFirstTick(openmct, clock); + const offsets = timeContext.getClockOffsets(); + + request.start = firstTimestamp + offsets.start; + request.end = firstTimestamp + offsets.end; + } + + return request; + } + }; +} + +/** + * Resolves with the clock's first tick value, or immediately with its current + * value if it has already ticked. + * + * Listens on the time context rather than on the clock itself. DefaultClock + * starts and stops based on its own listener count, so attaching a listener + * directly to the clock would interfere with its lifecycle. + */ +function waitForFirstTick(openmct, clock) { + return new Promise((resolve) => { + if (clock.currentValue() > NO_TICK_YET) { + resolve(clock.currentValue()); + + return; + } + + function resolveOnFirstTick(timestamp) { + openmct.time.off('tick', resolveOnFirstTick); + resolve(timestamp); + } + + openmct.time.on('tick', resolveOnFirstTick); + }); +}