Change Log 0.8.0
Version 0.8.0 introduces three major additions. First, a complete WLDT Monitoring System — a new it.wldt.monitoring package providing a pluggable, handler-based metric infrastructure that is automatically injected into every DT component and tracks a rich set of framework-native metrics out of the box. Second, Augmentation Function refinements: a unified result-handling type (AugmentationFunctionResultList), explicit state tracking for stateful functions, and full monitoring integration at the augmentation layer. Third, a WldtEventBus thread-safety fix that eliminates a ConcurrentModificationException under concurrent publish/subscribe.
Key updates:
- New
MonitoringInterfacesystem — one perDigitalTwin, auto-injected into Model, Adapters, Augmentation Functions, and Storage; opt-in per component; developer extension point viaMonitoringInterfaceHandler - 5 metric types —
WldtCounter,WldtUpDownCounter,WldtGauge,WldtTimer,WldtHistogram— all with fluent mutation API and snapshot semantics for callbacks - Framework-native metrics — automatic tracking of latency, success/error counts, and state transitions for DT_MODEL, Physical Adapter, Digital Adapter, and lifecycle without any developer action
- Augmentation monitoring integration —
AugmentationFunctionbase class gains built-in monitoring support;handleMetricsRegistration()hook for custom function-level metrics - Storage monitoring integration —
StorageManager,WldtStorage, andDefaultQueryManagerinstrumented automatically; query and write metrics available under theSTORAGEcomponent AugmentationFunctionResultList— unified return type for augmentation function executions, carrying either a list of results or an error (mutually exclusive)- Stateful function state tracking —
isRunning()onStatefulAugmentationFunction; guard checks inDigitalTwinModel; newisStatefulAugmentationFunctionRunning()query methods
New: WLDT Monitoring System
A new it.wldt.monitoring package provides the infrastructure for collecting, routing, and reacting to metrics emitted by any DT component.
Each DigitalTwin instance owns one MonitoringInterface which is automatically injected into all components at startup (DT Model, Physical Adapters, Digital Adapters, Augmentation Functions, and Storage). Monitoring is entirely opt-in: no metrics are collected and no callbacks are fired until both a configuration and a handler are attached.
Setup
MonitoringInterfaceConfiguration config = new MonitoringInterfaceConfiguration.Builder()
.withDtModelMonitoring()
.withPhysicalAdapterMonitoring()
.withDigitalAdapterMonitoring()
.withAugmentationMonitoring()
.withStorageMonitoring()
// .withAllMonitoring() ← enables all of the above at once
.withCustomNamespace("myapp.dt") // namespace prefix for custom metrics (default: "custom")
.build();
digitalTwin.getMonitoringInterface().setConfiguration(config);
digitalTwin.getMonitoringInterface().setHandler(new MyMonitoringHandler());
Handler
Extend MonitoringInterfaceHandler and override whichever callbacks are needed:
public class MyMonitoringHandler extends MonitoringInterfaceHandler {
@Override
public void onMetricRegistered(WldtMetricComponent component, WldtMetric metric) {
// metric is an empty snapshot — identity fields only, no value yet
// Use this to set up external counters (e.g., create a Prometheus counter)
}
@Override
public void onMetricUpdated(WldtMetricComponent component, WldtMetric metric) {
// metric is an independent snapshot of the live instance at callback time
if (component == WldtMetricComponent.DT_MODEL && metric instanceof WldtTimer) {
WldtTimer t = (WldtTimer) metric;
myHistogram.observe(t.getDurationMs());
}
}
}
Key behavior:
onMetricRegisteredfires once per metric name, passing anemptySnapshot()— uninitialized, carrying only identity fieldsonMetricUpdatedfires on every update (including on first registration if the metric was constructed with an initial value)- Both callbacks receive independent snapshot copies; mutations to the live metric do not affect them
- Handler exceptions are caught and logged — monitoring failures never affect DT processing
Additional MonitoringInterface methods:
| Method | Purpose |
|---|---|
registerMetric(WldtMetric) | Pre-register bypassing component flag gating; fires onMetricRegistered; use at DT startup |
deregisterMetric(String fullName) | Remove metric; next push starts fresh |
getMetric(String fullName) | Get live mutable instance by namespace.name |
getAllMetrics() | Unmodifiable snapshot of all registered metrics |
isMetricRegistered(String) | Check presence |
isActive() | true if both config and handler are set |
isActive(WldtMetricComponent) | true if active and the given component’s flag is enabled |
Metric Types
All metric constructors accept digitalTwinId as their first parameter, followed by namespace, name, and WldtMetricComponent. An optional long / double / long, double, double, double argument (depending on type) sets an initial value; without it the metric is created uninitialized (isInitialized() == false). An additional Map<String, Object> metadata overload is available on every constructor for tagging metrics with arbitrary key-value context.
The namespace for all framework-internal metrics is "core", returned by CoreMonitoringUtils.buildCoreNamespace(). Use any string namespace for custom metrics.
| Type | Use case | Key supporting fields |
|---|---|---|
WldtCounter | Monotonically increasing count (events processed, errors) | getValue(), getDelta(), getTotalIncrements() |
WldtUpDownCounter | Bidirectional count (active connections, running functions) | getValue(), getDelta(), getPeakValue(), getTroughValue() |
WldtGauge | Point-in-time observed value (temperature, queue depth) | getValue(), getPreviousValue(), getDelta(), getMinObserved(), getMaxObserved() |
WldtTimer | Operation latency in milliseconds | getDurationMs(), getDurationSeconds(), getMeanDurationMs(), getMinDurationMs(), getMaxDurationMs(), getObservationCount() |
WldtHistogram | Statistical distribution; supports single observations and pre-aggregated windows | getCount(), getMean(), getTotalCount(), getGlobalMin(), getGlobalMax(), getGlobalMean() |
All mutation methods return this for fluent chaining. WldtTimer also exposes updateSince(long startMs) which computes now − startMs automatically.
Quick examples
String dtId = "my-dt";
String ns = CoreMonitoringUtils.buildCoreNamespace(); // "core"
// Counter — register, then increment
monitoringInterface.registerMetric(
new WldtCounter(dtId, ns, "events_processed", WldtMetricComponent.DT_MODEL, 0L));
monitoringInterface.increaseCounter(ns, "events_processed");
monitoringInterface.increaseCounter(ns, "events_processed", 5L);
// UpDownCounter — register, then go up and down
monitoringInterface.registerMetric(
new WldtUpDownCounter(dtId, ns, "active_connections", WldtMetricComponent.PHYSICAL_ADAPTER, 0L));
monitoringInterface.increaseCounter(ns, "active_connections");
monitoringInterface.decreaseCounter(ns, "active_connections");
// Gauge — custom namespace
monitoringInterface.registerMetric(
new WldtGauge(dtId, "myapp.sensor", "temperature", WldtMetricComponent.CUSTOM, 20.0));
monitoringInterface.updateGauge("myapp.sensor", "temperature", 23.5);
// Timer — compute elapsed time from start
long start = System.currentTimeMillis();
// ... operation ...
monitoringInterface.registerMetric(
new WldtTimer(dtId, ns, "processing_latency_ms", WldtMetricComponent.DT_MODEL, 0L));
monitoringInterface.updateTimerSince(ns, "processing_latency_ms", start);
// Histogram — single observations
monitoringInterface.registerMetric(
new WldtHistogram(dtId, ns, "msg_size_bytes", WldtMetricComponent.PHYSICAL_ADAPTER,
1L, 120.0, 120.0, 120.0));
monitoringInterface.histogramObservation(ns, "msg_size_bytes", 95.0);
monitoringInterface.histogramObservation(ns, "msg_size_bytes", 10L, 1150.0, 80.0, 160.0); // pre-aggregated
Full documentation in monitoring.md.
Framework-Native Metrics
Enabling a component’s monitoring flag is sufficient — the framework registers and updates all metrics listed below automatically. All use namespace "core" and constants from CoreMonitoringUtils.
DT_MODEL
Enabled by withDtModelMonitoring().
| Metric name | Type | Tracks |
|---|---|---|
pt_property_variation_exec_time | WldtTimer | PA property variation processing time |
pt_property_variation_exec_success_count | WldtCounter | Successful PA property variation handlers |
pt_property_variation_exec_error_count | WldtCounter | Failed PA property variation handlers |
pt_event_notification_exec_time | WldtTimer | PA event notification processing time |
pt_event_notification_exec_success_count | WldtCounter | Successful PA event notification handlers |
pt_event_notification_exec_error_count | WldtCounter | Failed PA event notification handlers |
pt_rel_instance_created_exec_time | WldtTimer | Relationship-created handler time |
pt_rel_instance_created_exec_success_count | WldtCounter | Successful relationship-created handlers |
pt_rel_instance_created_exec_error_count | WldtCounter | Failed relationship-created handlers |
pt_rel_instance_deleted_exec_time | WldtTimer | Relationship-deleted handler time |
pt_rel_instance_deleted_exec_success_count | WldtCounter | Successful relationship-deleted handlers |
pt_rel_instance_deleted_exec_error_count | WldtCounter | Failed relationship-deleted handlers |
digital_action_exec_time | WldtTimer | Digital action request processing time |
digital_action_exec_success_count | WldtCounter | Successful digital action handlers |
digital_action_exec_error_count | WldtCounter | Failed digital action handlers |
dt_state_computation_exec_time | WldtTimer | Full DT state computation time |
dt_state_computation_exec_success_count | WldtCounter | Successful state computations |
dt_state_computation_exec_error_count | WldtCounter | Failed state computations |
af_stateless_exec_success_count | WldtCounter | Successful stateless AF invocations (DTM level) |
af_stateless_exec_error_count | WldtCounter | Failed stateless AF invocations (DTM level) |
af_stateful_start_success_count | WldtCounter | Successful stateful AF start requests (DTM level) |
af_stateful_start_error_count | WldtCounter | Failed stateful AF start requests (DTM level) |
af_stateful_stop_success_count | WldtCounter | Successful stateful AF stop requests (DTM level) |
af_stateful_stop_error_count | WldtCounter | Failed stateful AF stop requests (DTM level) |
PHYSICAL_ADAPTER
Enabled by withPhysicalAdapterMonitoring().
| Metric name | Type | Tracks |
|---|---|---|
pa_property_event_pub_success_count | WldtCounter | Property variation events published successfully |
pa_property_event_pub_error_count | WldtCounter | Property variation event publication failures |
pa_property_event_notification_pub_success_count | WldtCounter | Event notification messages published successfully |
pa_property_event_notification_pub_error_count | WldtCounter | Event notification publication failures |
pa_property_rel_created_event_pub_success_count | WldtCounter | Relationship-created events published successfully |
pa_property_rel_created_event_pub_error_count | WldtCounter | Relationship-created event publication failures |
pa_property_rel_deleted_event_pub_success_count | WldtCounter | Relationship-deleted events published successfully |
pa_property_rel_deleted_event_pub_error_count | WldtCounter | Relationship-deleted event publication failures |
pa_action_computation_exec_time | WldtTimer | Physical action computation time |
pa_action_computation_exec_success_count | WldtCounter | Successful physical action computations |
pa_action_computation_exec_error_count | WldtCounter | Failed physical action computations |
DIGITAL_ADAPTER
Enabled by withDigitalAdapterMonitoring().
| Metric name | Type | Tracks |
|---|---|---|
da_action_event_pub_success_count | WldtCounter | Digital action events published successfully |
da_action_event_pub_error_count | WldtCounter | Digital action event publication failures |
da_state_update_processing_exec_time | WldtTimer | DT state update processing time |
da_state_update_processing_exec_success_count | WldtCounter | Successful state update processing |
da_state_update_processing_exec_error_count | WldtCounter | Failed state update processing |
da_event_notification_processing_exec_time | WldtTimer | DT event notification processing time |
da_event_notification_processing_exec_success_count | WldtCounter | Successful event notification processing |
da_event_notification_processing_exec_error_count | WldtCounter | Failed event notification processing |
Lifecycle
dt_lifecycle_value (WldtUpDownCounter) tracks the current LifeCycleState as its numeric ordinal. Updated automatically by DigitalTwin on every lifecycle transition.
Augmentation Function Monitoring Integration
AugmentationFunction now participates in the monitoring system through three new protected fields injected by AugmentationFunctionHandler when a function is registered:
| Field | Type | Value after injection |
|---|---|---|
monitoringInterface | MonitoringInterface | The owning DT’s shared monitoring interface |
digitalTwinId | String | The owning DT’s unique id |
metricsNamespace | String | CoreMonitoringUtils.buildCoreNamespace() → "core" |
The injection is triggered automatically via setMonitoringInterface(MonitoringInterface, String digitalTwinId). Developers do not call this method directly.
handleMetricsRegistration() hook
Override this protected method to register custom function-level metrics at injection time:
public class MyStatelessFunction extends StatelessAugmentationFunction {
private static final String EXEC_TIME = "my_function_exec_time";
public MyStatelessFunction(String id) {
super(id, "MyFunction", AugmentationFunctionType.STATELESS);
}
@Override
protected void handleMetricsRegistration() {
if (monitoringInterface != null && monitoringInterface.isActive(WldtMetricComponent.AUGMENTATION)) {
Map<String, Object> meta = new HashMap<>();
meta.put(AugmentationFunction.METRIC_METADATA_AF_FUNCTION_ID_KEY, getId());
monitoringInterface.registerMetric(
new WldtTimer(digitalTwinId, metricsNamespace, EXEC_TIME,
WldtMetricComponent.AUGMENTATION, meta));
}
}
@Override
protected AugmentationFunctionResultList run(AugmentationFunctionRequest request)
throws AugmentationFunctionException {
long start = System.currentTimeMillis();
// ... function logic ...
monitoringInterface.updateTimerSince(metricsNamespace, EXEC_TIME, start);
return new AugmentationFunctionResultList(/* result */);
}
}
AugmentationFunction.METRIC_METADATA_AF_FUNCTION_ID_KEY ("af_function_id") is the standard metadata key for tagging function-level metrics with the function’s unique id, enabling per-function routing in the handler.
Framework-tracked AUGMENTATION metrics
Enabled by withAugmentationMonitoring(). All use namespace "core".
| Metric name | Type | Tracks |
|---|---|---|
af_handler_count | WldtUpDownCounter | Number of registered handlers (manager level) |
af_stateful_running_count | WldtUpDownCounter | Stateful functions currently running (manager level) |
af_handler_registered_stateless_count | WldtUpDownCounter | Stateless functions registered per handler |
af_handler_registered_stateful_count | WldtUpDownCounter | Stateful functions registered per handler |
af_handler_stateful_running_count | WldtUpDownCounter | Stateful functions currently running per handler |
af_function_stateless_exec_time | WldtTimer | Stateless function execution time |
af_function_stateless_exec_success_count | WldtCounter | Successful stateless executions |
af_function_stateless_exec_error_count | WldtCounter | Failed stateless executions |
af_function_stateful_start_exec_time | WldtTimer | Time to start a stateful function |
af_function_stateful_start_success_count | WldtCounter | Successful stateful starts |
af_function_stateful_start_error_count | WldtCounter | Failed stateful starts |
af_function_stateful_stop_exec_time | WldtTimer | Time to stop a stateful function |
af_function_stateful_stop_success_count | WldtCounter | Successful stateful stops |
af_function_stateful_stop_error_count | WldtCounter | Failed stateful stops |
af_function_query_exec_time | WldtTimer | Storage query time from a function |
af_function_query_exec_success_count | WldtCounter | Successful storage queries from functions |
af_function_query_exec_error_count | WldtCounter | Failed storage queries from functions |
af_function_state_update_exec_time | WldtTimer | Time to dispatch a DT state update to a stateful function |
af_function_state_update_success_count | WldtCounter | Successful state update dispatches |
af_function_state_update_error_count | WldtCounter | Failed state update dispatches |
Full list including result/error/registration dispatch metrics: monitoring.md — Framework-Native Metrics.
Storage Monitoring Integration
StorageManager, WldtStorage, and DefaultQueryManager are instrumented automatically when withStorageMonitoring() is set in the configuration. No developer action is required. WldtStorage gains a concrete (non-abstract) setMonitoringInterface(MonitoringInterface) method — existing implementors require no changes.
Framework-tracked STORAGE metrics
All use namespace "core".
| Metric name | Type | Tracks |
|---|---|---|
af_storage_query_success_count | WldtCounter | Successful storage queries |
af_storage_query_error_count | WldtCounter | Failed storage queries |
af_storage_query_exec_time | WldtTimer | Query execution time |
af_storage_write_pa_description_success_count | WldtCounter | Successful PA description writes |
af_storage_write_pa_description_error_count | WldtCounter | Failed PA description writes |
af_storage_write_pa_description_exec_time | WldtTimer | PA description write time |
af_storage_write_dt_state_success_count | WldtCounter | Successful DT state writes |
af_storage_write_dt_state_error_count | WldtCounter | Failed DT state writes |
af_storage_write_dt_state_exec_time | WldtTimer | DT state write time |
af_storage_write_af_success_count | WldtCounter | Successful augmentation function result writes |
af_storage_write_af_error_count | WldtCounter | Failed augmentation function result writes |
af_storage_write_af_exec_time | WldtTimer | Augmentation function result write time |
af_storage_write_de_success_count | WldtCounter | Successful digital event writes |
af_storage_write_de_error_count | WldtCounter | Failed digital event writes |
af_storage_write_de_exec_time | WldtTimer | Digital event write time |
af_storage_write_pe_success_count | WldtCounter | Successful physical event writes |
af_storage_write_pe_error_count | WldtCounter | Failed physical event writes |
af_storage_write_pe_exec_time | WldtTimer | Physical event write time |
af_storage_write_lifecycle_event_success_count | WldtCounter | Successful lifecycle event writes |
af_storage_write_lifecycle_event_error_count | WldtCounter | Failed lifecycle event writes |
af_storage_write_lifecycle_event_exec_time | WldtTimer | Lifecycle event write time |
Augmentation: Result Handling Refactoring
New: AugmentationFunctionResultList
A new AugmentationFunctionResultList class has been added to it.wldt.augmentation.result. It extends ArrayList<AugmentationFunctionResult<?>> and is the unified return type for all augmentation function executions.
A list instance carries either results or an error — not both:
- Once an error is set via
setAugmentationFunctionError(), all mutating operations (add,set,remove,clear) throwIllegalStateException. Read operations remain available. hasError()— indicates error state.getAugmentationFunctionError()— retrieves the embedded error, ornullif none.
Constructors:
// Empty list — append results with the standard List API
new AugmentationFunctionResultList()
// Single-result convenience constructor
new AugmentationFunctionResultList(AugmentationFunctionResult<?> result)
// Error-state constructor — list is immediately immutable
new AugmentationFunctionResultList(AugmentationFunctionError error)
Removed: StatelessAugmentationListener
The StatelessAugmentationListener interface has been removed. AugmentationFunctionHandler no longer implements it. The notifyError() method is no longer available in StatelessAugmentationFunction. Stateless functions must now report errors by returning a new AugmentationFunctionResultList(error) directly from run().
Migration:
// Before (0.7.0)
protected List<AugmentationFunctionResult<?>> run(AugmentationFunctionRequest request)
throws AugmentationFunctionException {
try {
// ...
return Collections.singletonList(result);
} catch (Exception e) {
notifyError(new AugmentationFunctionError(AugmentationFunctionErrorType.ERROR, e.getMessage()));
return Collections.emptyList();
}
}
// After (0.8.0)
protected AugmentationFunctionResultList run(AugmentationFunctionRequest request)
throws AugmentationFunctionException {
try {
// ...
return new AugmentationFunctionResultList(result);
} catch (Exception e) {
return new AugmentationFunctionResultList(
new AugmentationFunctionError(AugmentationFunctionErrorType.ERROR, e.getMessage())
);
}
}
Stateful functions continue to have notifyError() available for asynchronous error reporting via the handler. onStatefulAugmentationFunctionError() has been removed from StatefulAugmentationListener — errors from stateful functions are now embedded in the AugmentationFunctionResultList passed to onStatefulAugmentationFunctionResult().
API Signature Changes
All methods that previously accepted or returned List<AugmentationFunctionResult<?>> now use AugmentationFunctionResultList:
| Component | Method / Callback | Change |
|---|---|---|
StatelessAugmentationFunction | run() return type | List<...> → AugmentationFunctionResultList |
AugmentationFunctionHandler | handleAugmentationFunctionExecution() return type | List<...> → AugmentationFunctionResultList |
StatefulAugmentationListener | onStatefulAugmentationFunctionResult() parameter | List<...> → AugmentationFunctionResultList |
DigitalTwinModel | onAugmentationFunctionResultEvent() parameter | List<...> → AugmentationFunctionResultList |
AugmentationFunctionResultWldtEvent | payload type | List<...> → AugmentationFunctionResultList |
StatefulAugmentationFunction | notifyResult() parameter | List<...> → AugmentationFunctionResultList |
Augmentation: Stateful Function State Management
StatefulAugmentationFunction.isRunning()
A new public method isRunning() has been added to StatefulAugmentationFunction. It returns true if the function has been started and has not yet been stopped. The flag is managed automatically by the framework; subclasses do not set it directly.
public final boolean isRunning()
Guard checks in DigitalTwinModel
startAugmentationFunction() is now a no-op if the target function is already running. stopAugmentationFunction() is a no-op if the target function is already stopped. Both methods log an informational message and return early without publishing an event.
New: isStatefulAugmentationFunctionRunning()
Two new protected methods allow the DTM to query the running state of a stateful function:
// Search across all registered handlers
protected boolean isStatefulAugmentationFunctionRunning(String augmentationFunctionId)
throws AugmentationFunctionException
// Target a specific handler
protected boolean isStatefulAugmentationFunctionRunning(String augmentationFunctionHandlerId,
String augmentationFunctionId)
throws AugmentationFunctionException
Both throw AugmentationFunctionException if the function is not found or is not of type STATEFUL.
Request ID attached to errors automatically
Both StatelessAugmentationFunction and StatefulAugmentationFunction now automatically attach the current request id to any AugmentationFunctionError passed to notifyError(). The augmentationFunctionRequestId field on the error object is set by the framework — no developer action required.
Fix: WldtEventBus Thread Safety
The internal event publishing loop in WldtEventBus now iterates over a copy of the subscription event type set rather than the live key set. This eliminates a ConcurrentModificationException that could occur when a QueryExecutor callback modified the subscription map while publishing was in progress.
This is an internal fix with no API surface change.
Breaking Changes
1. StatelessAugmentationListener removed
The StatelessAugmentationListener interface (it.wldt.augmentation.listener) has been removed. Any code that implements or references this interface must be updated. Errors from stateless functions must now be returned inside an AugmentationFunctionResultList from run() — see the migration example above.
2. API signatures: List<AugmentationFunctionResult<?>> → AugmentationFunctionResultList
All method signatures listed in the API Signature Changes table above must be updated. In most cases the change is mechanical: replace List<AugmentationFunctionResult<?>> with AugmentationFunctionResultList and update return Collections.emptyList() / return Collections.singletonList(r) to use the new constructors.
3. onStatefulAugmentationFunctionError() removed from StatefulAugmentationListener
Implementations of StatefulAugmentationListener that override onStatefulAugmentationFunctionError() must remove the override. Stateful function errors are now delivered via onStatefulAugmentationFunctionResult() with the error embedded in the AugmentationFunctionResultList.
The monitoring system is not a breaking change. It is entirely opt-in and no existing DigitalTwinModel subclass, PhysicalAdapter, DigitalAdapter, AugmentationFunctionHandler, or WldtStorage implementation requires modification.