#include "ua_server_internal.h"
#include "ua_subscription.h"
#include "../ua_types_encoding_binary.h"
#ifdef UA_ENABLE_SUBSCRIPTIONS
#define UA_DETECT_DEADBAND(TYPE) do { \
TYPE v1 = *(const TYPE*)data1; \
TYPE v2 = *(const TYPE*)data2; \
UA_UInt64 mag = (v1 > v2) ? (UA_UInt64)v1 - (UA_UInt64)v2 \
: (UA_UInt64)v2 - (UA_UInt64)v1; \
UA_Double diff = (UA_Double)mag; \
return (diff > deadband); \
} while(false)
#define UA_DETECT_DEADBAND_FLOAT(TYPE) do { \
TYPE v1 = *(const TYPE*)data1; \
TYPE v2 = *(const TYPE*)data2; \
UA_Double diff = (v1 > v2) ? (UA_Double)v1 - (UA_Double)v2 \
: (UA_Double)v2 - (UA_Double)v1; \
return (diff > deadband); \
} while(false)
static UA_Boolean
detectScalarDeadBand(const void *data1, const void *data2,
const UA_DataType *type, const UA_Double deadband) {
if(type->typeKind == UA_DATATYPEKIND_SBYTE) {
UA_DETECT_DEADBAND(UA_SByte);
} else if(type->typeKind == UA_DATATYPEKIND_BYTE) {
UA_DETECT_DEADBAND(UA_Byte);
} else if(type->typeKind == UA_DATATYPEKIND_INT16) {
UA_DETECT_DEADBAND(UA_Int16);
} else if(type->typeKind == UA_DATATYPEKIND_UINT16) {
UA_DETECT_DEADBAND(UA_UInt16);
} else if(type->typeKind == UA_DATATYPEKIND_INT32) {
UA_DETECT_DEADBAND(UA_Int32);
} else if(type->typeKind == UA_DATATYPEKIND_UINT32) {
UA_DETECT_DEADBAND(UA_UInt32);
} else if(type->typeKind == UA_DATATYPEKIND_INT64) {
UA_DETECT_DEADBAND(UA_Int64);
} else if(type->typeKind == UA_DATATYPEKIND_UINT64) {
UA_DETECT_DEADBAND(UA_UInt64);
} else if(type->typeKind == UA_DATATYPEKIND_FLOAT) {
UA_DETECT_DEADBAND_FLOAT(UA_Float);
} else if(type->typeKind == UA_DATATYPEKIND_DOUBLE) {
UA_DETECT_DEADBAND_FLOAT(UA_Double);
} else {
return false;
}
}
static UA_Boolean
detectVariantDeadband(const UA_Variant *value, const UA_Variant *oldValue,
const UA_Double deadbandValue) {
if(value->arrayLength != oldValue->arrayLength)
return true;
if(value->type != oldValue->type)
return true;
if(UA_Variant_isScalar(value) != UA_Variant_isScalar(oldValue))
return true;
size_t length = 1;
if(!UA_Variant_isScalar(value))
length = value->arrayLength;
uintptr_t data = (uintptr_t)value->data;
uintptr_t oldData = (uintptr_t)oldValue->data;
UA_UInt32 memSize = value->type->memSize;
for(size_t i = 0; i < length; ++i) {
if(detectScalarDeadBand((const void*)data, (const void*)oldData,
value->type, deadbandValue))
return true;
data += memSize;
oldData += memSize;
}
return false;
}
static UA_Boolean
detectValueChange(UA_Server *server, UA_MonitoredItem *mon, const UA_DataValue *dv) {
UA_LOCK_ASSERT(&server->serviceMutex);
if(dv->hasStatus != mon->lastValue.hasStatus ||
dv->status != mon->lastValue.status) {
return true;
}
UA_DataChangeTrigger trigger = UA_DATACHANGETRIGGER_STATUSVALUE;
const UA_DataChangeFilter *dcf = NULL;
const UA_ExtensionObject *filter = &mon->parameters.filter;
if(filter->content.decoded.type == &UA_TYPES[UA_TYPES_DATACHANGEFILTER]) {
dcf = (UA_DataChangeFilter*)filter->content.decoded.data;
trigger = dcf->trigger;
}
if(trigger == UA_DATACHANGETRIGGER_STATUS)
return false;
UA_assert(trigger == UA_DATACHANGETRIGGER_STATUSVALUE ||
trigger == UA_DATACHANGETRIGGER_STATUSVALUETIMESTAMP);
if(dv->hasValue != mon->lastValue.hasValue)
return true;
if(dcf && dcf->deadbandType == UA_DEADBANDTYPE_ABSOLUTE &&
dv->value.type != NULL && UA_DataType_isNumeric(dv->value.type))
return detectVariantDeadband(&dv->value, &mon->lastValue.value,
dcf->deadbandValue);
if(trigger == UA_DATACHANGETRIGGER_STATUSVALUETIMESTAMP) {
if(dv->hasSourceTimestamp != mon->lastValue.hasSourceTimestamp)
return true;
if(dv->hasSourceTimestamp &&
dv->sourceTimestamp != mon->lastValue.sourceTimestamp)
return true;
}
return !UA_equal(&dv->value, &mon->lastValue.value,
&UA_TYPES[UA_TYPES_VARIANT]);
}
UA_StatusCode
UA_MonitoredItem_createDataChangeNotification(UA_Server *server, UA_MonitoredItem *mon,
const UA_DataValue *dv) {
UA_DataValue valueCopy;
UA_StatusCode retval = UA_DataValue_copy(dv, &valueCopy);
if(retval != UA_STATUSCODE_GOOD)
return retval;
UA_Notification *n = UA_Notification_new();
if(!n) {
UA_DataValue_clear(&valueCopy);
return UA_STATUSCODE_BADOUTOFMEMORY;
}
n->mon = mon;
n->data.dataChange.value = valueCopy;
n->data.dataChange.clientHandle = mon->parameters.clientHandle;
UA_Notification_enqueueAndTrigger(server, n);
return UA_STATUSCODE_GOOD;
}
void
UA_MonitoredItem_processSampledValue(UA_Server *server, UA_MonitoredItem *mon,
UA_DataValue *value) {
UA_assert(mon->itemToMonitor.attributeId != UA_ATTRIBUTEID_EVENTNOTIFIER);
UA_LOCK_ASSERT(&server->serviceMutex);
UA_Boolean changed = detectValueChange(server, mon, value);
if(!changed) {
UA_LOG_DEBUG_SUBSCRIPTION(server->config.logging, mon->subscription,
"MonitoredItem %" PRIi32 " | "
"The value has not changed", mon->monitoredItemId);
UA_DataValue_clear(value);
return;
}
UA_StatusCode res =
UA_MonitoredItem_createDataChangeNotification(server, mon, value);
if(res != UA_STATUSCODE_GOOD) {
UA_LOG_WARNING_SUBSCRIPTION(server->config.logging, mon->subscription,
"MonitoredItem %" PRIi32 " | "
"Processing the sample returned the statuscode %s",
mon->monitoredItemId, UA_StatusCode_name(res));
UA_DataValue_clear(value);
return;
}
UA_DataValue_clear(&mon->lastValue);
mon->lastValue = *value;
if(!mon->subscription) {
UA_LocalMonitoredItem *localMon = (UA_LocalMonitoredItem*) mon;
void *nodeContext = NULL;
getNodeContext(server, mon->itemToMonitor.nodeId, &nodeContext);
localMon->callback.dataChangeCallback(server,
mon->monitoredItemId, localMon->context,
&mon->itemToMonitor.nodeId, nodeContext,
mon->itemToMonitor.attributeId, value);
}
}
static void
processMonitoredItemAsyncRead(UA_Server *server, UA_MonitoredItem *mon,
const UA_DataValue *result) {
mon->outstandingAsyncReads--;
UA_DataValue *mut_result = (UA_DataValue*)(uintptr_t)result;
if(mut_result->status == UA_STATUSCODE_BADREQUESTCANCELLEDBYREQUEST)
return;
UA_MonitoredItem_processSampledValue(server, mon, mut_result);
UA_DataValue_init(mut_result);
}
void
UA_MonitoredItem_sample(UA_Server *server, UA_MonitoredItem *mon) {
UA_LOCK_ASSERT(&server->serviceMutex);
UA_assert(mon->itemToMonitor.attributeId != UA_ATTRIBUTEID_EVENTNOTIFIER);
UA_Subscription *sub = mon->subscription;
UA_LOG_DEBUG_SUBSCRIPTION(server->config.logging, sub, "MonitoredItem %" PRIi32
" | Sample callback called", mon->monitoredItemId);
UA_Session *session = (sub) ? sub->session : &server->adminSession;
UA_StatusCode res = UA_STATUSCODE_BADTOOMANYOPERATIONS;
if(UA_LIKELY(mon->outstandingAsyncReads < UA_MONITOREDITEM_ASYNC_MAX)) {
res = read_async(server, session, &mon->itemToMonitor, mon->timestampsToReturn,
(UA_ServerAsyncReadResultCallback)processMonitoredItemAsyncRead, mon, 0);
}
if(res == UA_STATUSCODE_GOOD) {
mon->outstandingAsyncReads++;
} else {
UA_DataValue dv;
UA_DataValue_init(&dv);
dv.hasStatus = true;
dv.status = res;
UA_MonitoredItem_processSampledValue(server, mon, &dv);
}
}
#endif