#include <open62541/client_highlevel.h>
#include <open62541/client_highlevel_async.h>
#include "ua_client_internal.h"
struct UA_Client_MonitoredItem_ForDelete {
UA_Client *client;
UA_Client_Subscription *sub;
UA_UInt32 *monitoredItemId;
};
static enum ZIP_CMP
UA_ClientHandle_cmp(const void *a, const void *b) {
const UA_Client_MonitoredItem *aa = (const UA_Client_MonitoredItem *)a;
const UA_Client_MonitoredItem *bb = (const UA_Client_MonitoredItem *)b;
if(aa->parameters.clientHandle < bb->parameters.clientHandle)
return ZIP_CMP_LESS;
if(aa->parameters.clientHandle > bb->parameters.clientHandle)
return ZIP_CMP_MORE;
return ZIP_CMP_EQ;
}
ZIP_FUNCTIONS(MonitorItemsTree, UA_Client_MonitoredItem, zipfields,
UA_Client_MonitoredItem, zipfields, UA_ClientHandle_cmp)
static UA_Client_MonitoredItem *
findMonitoredItemByHandle(UA_Client_Subscription *sub, UA_UInt32 clientHandle) {
UA_Client_MonitoredItem dummy;
memset(&dummy, 0, sizeof(dummy));
dummy.parameters.clientHandle = clientHandle;
return ZIP_FIND(MonitorItemsTree, &sub->monitoredItems, &dummy);
}
static void *
MonitoredItem_findPendingByHandle(void *data, UA_Client_MonitoredItem *mon) {
UA_UInt32 clientHandle = *(UA_UInt32*)data;
if(mon->pendingParameters.clientHandle == clientHandle)
return mon;
return NULL;
}
static void
MonitoredItem_delete(UA_Client *client, UA_Client_Subscription *sub,
UA_Client_MonitoredItem *mon);
static void
Subscription_create(UA_Client *client, UA_Client_Subscription *newSub,
UA_CreateSubscriptionResponse *response) {
UA_LOCK_ASSERT(&client->clientMutex);
UA_EventLoop *el = client->config.eventLoop;
newSub->subscriptionId = response->subscriptionId;
newSub->sequenceNumber = 0;
newSub->lastActivity = el->dateTime_nowMonotonic(el);
newSub->publishingInterval = response->revisedPublishingInterval;
newSub->maxKeepAliveCount = response->revisedMaxKeepAliveCount;
ZIP_INIT(&newSub->monitoredItems);
LIST_INSERT_HEAD(&client->subscriptions, newSub, listEntry);
__Client_Subscriptions_backgroundPublish(client);
}
static void
Subscriptions_create_handler(UA_Client *client, void *data,
UA_UInt32 requestId, void *r) {
UA_LOCK_ASSERT(&client->clientMutex);
UA_CreateSubscriptionResponse *response = (UA_CreateSubscriptionResponse *)r;
CustomCallback *cc = (CustomCallback *)data;
UA_Client_Subscription *newSub = (UA_Client_Subscription *)cc->clientData;
if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD) {
UA_free(newSub);
goto cleanup;
}
Subscription_create(client, newSub, response);
cleanup:
if(cc->userCallback)
cc->userCallback(client, cc->userData, requestId, response);
UA_free(cc);
}
UA_CreateSubscriptionResponse
UA_Client_Subscriptions_create(UA_Client *client,
const UA_CreateSubscriptionRequest request,
void *subscriptionContext,
UA_Client_StatusChangeNotificationCallback statusChangeCallback,
UA_Client_DeleteSubscriptionCallback deleteCallback) {
lockClient(client);
UA_CreateSubscriptionResponse response;
UA_Client_Subscription *sub = (UA_Client_Subscription *)
UA_malloc(sizeof(UA_Client_Subscription));
if(!sub) {
UA_CreateSubscriptionResponse_init(&response);
response.responseHeader.serviceResult = UA_STATUSCODE_BADOUTOFMEMORY;
unlockClient(client);
return response;
}
sub->context = subscriptionContext;
sub->statusChangeCallback = statusChangeCallback;
sub->deleteCallback = deleteCallback;
__Client_Service(client, &request, &UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONREQUEST],
&response, &UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONRESPONSE]);
if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD) {
UA_free(sub);
unlockClient(client);
return response;
}
Subscription_create(client, sub, &response);
unlockClient(client);
return response;
}
UA_StatusCode
UA_Client_Subscriptions_create_async(UA_Client *client, const UA_CreateSubscriptionRequest request,
void *subscriptionContext,
UA_Client_StatusChangeNotificationCallback statusChangeCallback,
UA_Client_DeleteSubscriptionCallback deleteCallback,
UA_ClientAsyncCreateSubscriptionCallback createCallback,
void *userdata, UA_UInt32 *requestId) {
CustomCallback *cc = (CustomCallback *)UA_calloc(1, sizeof(CustomCallback));
if(!cc)
return UA_STATUSCODE_BADOUTOFMEMORY;
UA_Client_Subscription *sub = (UA_Client_Subscription *)
UA_malloc(sizeof(UA_Client_Subscription));
if(!sub) {
UA_free(cc);
return UA_STATUSCODE_BADOUTOFMEMORY;
}
sub->context = subscriptionContext;
sub->statusChangeCallback = statusChangeCallback;
sub->deleteCallback = deleteCallback;
cc->userCallback = (UA_ClientAsyncServiceCallback)createCallback;
cc->userData = userdata;
cc->clientData = sub;
UA_StatusCode res =
__UA_Client_AsyncService(client, &request,
&UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONREQUEST],
Subscriptions_create_handler,
&UA_TYPES[UA_TYPES_CREATESUBSCRIPTIONRESPONSE],
cc, requestId);
if(res != UA_STATUSCODE_GOOD) {
UA_free(cc);
UA_free(sub);
}
return res;
}
static UA_Client_Subscription *
findSubscriptionById(const UA_Client *client, UA_UInt32 subscriptionId) {
UA_Client_Subscription *sub = NULL;
LIST_FOREACH(sub, &client->subscriptions, listEntry) {
if(sub->subscriptionId == subscriptionId)
break;
}
return sub;
}
static void
Subscription_modify(UA_Client *client, UA_Client_Subscription *sub,
const UA_ModifySubscriptionResponse *response) {
sub->publishingInterval = response->revisedPublishingInterval;
sub->maxKeepAliveCount = response->revisedMaxKeepAliveCount;
}
static void
Subscription_modify_handler(UA_Client *client, void *data,
UA_UInt32 requestId, void *r) {
UA_LOCK_ASSERT(&client->clientMutex);
UA_ModifySubscriptionResponse *response = (UA_ModifySubscriptionResponse *)r;
CustomCallback *cc = (CustomCallback *)data;
UA_Client_Subscription *sub =
findSubscriptionById(client, (UA_UInt32)(uintptr_t)cc->clientData);
if(sub) {
Subscription_modify(client, sub, response);
} else {
UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
"No internal representation of subscription %" PRIu32,
(UA_UInt32)(uintptr_t)cc->clientData);
}
if(cc->userCallback)
cc->userCallback(client, cc->userData, requestId, response);
UA_free(cc);
}
UA_StatusCode
UA_Client_Subscriptions_getContext(UA_Client *client, UA_UInt32 subscriptionId,
void **subContext) {
if(!client || !subContext)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
if(!sub) {
unlockClient(client);
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
}
*subContext = sub->context;
unlockClient(client);
return UA_STATUSCODE_GOOD;
}
UA_StatusCode
UA_Client_Subscriptions_setContext(UA_Client *client, UA_UInt32 subscriptionId,
void *subContext) {
if(!client)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
if(!sub) {
unlockClient(client);
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
}
sub->context = subContext;
unlockClient(client);
return UA_STATUSCODE_GOOD;
}
UA_ModifySubscriptionResponse
UA_Client_Subscriptions_modify(UA_Client *client,
const UA_ModifySubscriptionRequest request) {
UA_ModifySubscriptionResponse response;
UA_ModifySubscriptionResponse_init(&response);
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
if(!sub) {
unlockClient(client);
response.responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
return response;
}
__Client_Service(client,
&request, &UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONREQUEST],
&response, &UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONRESPONSE]);
sub = findSubscriptionById(client, request.subscriptionId);
if(!sub) {
response.responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
unlockClient(client);
return response;
}
Subscription_modify(client, sub, &response);
unlockClient(client);
return response;
}
UA_StatusCode
UA_Client_Subscriptions_modify_async(UA_Client *client,
const UA_ModifySubscriptionRequest request,
UA_ClientAsyncModifySubscriptionCallback callback,
void *userdata, UA_UInt32 *requestId) {
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
if(!sub) {
unlockClient(client);
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
}
CustomCallback *cc = (CustomCallback *)UA_calloc(1, sizeof(CustomCallback));
if(!cc) {
unlockClient(client);
return UA_STATUSCODE_BADOUTOFMEMORY;
}
cc->clientData = (void *)(uintptr_t)request.subscriptionId;
cc->userData = userdata;
cc->userCallback = (UA_ClientAsyncServiceCallback)callback;
UA_StatusCode res =
__Client_AsyncService(client, &request,
&UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONREQUEST],
Subscription_modify_handler,
&UA_TYPES[UA_TYPES_MODIFYSUBSCRIPTIONRESPONSE],
cc, requestId);
unlockClient(client);
return res;
}
static void *
MonitoredItem_delete_wrapper(void *data, UA_Client_MonitoredItem *mon) {
struct UA_Client_MonitoredItem_ForDelete *deleteMonitoredItem =
(struct UA_Client_MonitoredItem_ForDelete *)data;
if(deleteMonitoredItem != NULL) {
if(deleteMonitoredItem->monitoredItemId != NULL &&
(mon->monitoredItemId != *deleteMonitoredItem->monitoredItemId)) {
return NULL;
}
MonitoredItem_delete(deleteMonitoredItem->client, deleteMonitoredItem->sub, mon);
}
return NULL;
}
static void
__Client_Subscription_deleteInternal(UA_Client *client,
UA_Client_Subscription *sub) {
struct UA_Client_MonitoredItem_ForDelete deleteMonitoredItem;
memset(&deleteMonitoredItem, 0, sizeof(struct UA_Client_MonitoredItem_ForDelete));
deleteMonitoredItem.client = client;
deleteMonitoredItem.sub = sub;
ZIP_ITER(MonitorItemsTree, &sub->monitoredItems,
MonitoredItem_delete_wrapper, &deleteMonitoredItem);
if(sub->deleteCallback) {
void *subC = sub->context;
UA_UInt32 subId = sub->subscriptionId;
sub->deleteCallback(client, subId, subC);
}
LIST_REMOVE(sub, listEntry);
UA_free(sub);
}
static void
__Client_Subscription_processDelete(UA_Client *client,
const UA_DeleteSubscriptionsRequest *request,
const UA_DeleteSubscriptionsResponse *response) {
if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD)
return;
if(request->subscriptionIdsSize != response->resultsSize)
return;
for(size_t i = 0; i < request->subscriptionIdsSize; i++) {
if(response->results[i] != UA_STATUSCODE_GOOD &&
response->results[i] != UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID)
continue;
UA_Client_Subscription *sub =
findSubscriptionById(client, request->subscriptionIds[i]);
if(!sub) {
UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
"No internal representation of subscription %" PRIu32,
request->subscriptionIds[i]);
continue;
}
__Client_Subscription_deleteInternal(client, sub);
}
}
typedef struct {
UA_DeleteSubscriptionsRequest request;
UA_ClientAsyncServiceCallback userCallback;
void *userData;
} DeleteSubscriptionCallback;
static void
Subscriptions_delete_handler(UA_Client *client, void *data,
UA_UInt32 requestId, void *r) {
UA_DeleteSubscriptionsResponse *response =
(UA_DeleteSubscriptionsResponse *)r;
DeleteSubscriptionCallback *dsc =
(DeleteSubscriptionCallback*)data;
lockClient(client);
__Client_Subscription_processDelete(client, &dsc->request, response);
if(dsc->userCallback)
dsc->userCallback(client, dsc->userData, requestId, response);
UA_DeleteSubscriptionsRequest_clear(&dsc->request);
UA_free(dsc);
unlockClient(client);
}
UA_StatusCode
UA_Client_Subscriptions_delete_async(UA_Client *client,
const UA_DeleteSubscriptionsRequest request,
UA_ClientAsyncDeleteSubscriptionsCallback callback,
void *userdata, UA_UInt32 *requestId) {
DeleteSubscriptionCallback *dsc = (DeleteSubscriptionCallback*)
UA_malloc(sizeof(DeleteSubscriptionCallback));
if(!dsc)
return UA_STATUSCODE_BADOUTOFMEMORY;
dsc->userCallback = (UA_ClientAsyncServiceCallback)callback;
dsc->userData = userdata;
UA_StatusCode res = UA_DeleteSubscriptionsRequest_copy(&request, &dsc->request);
if(res != UA_STATUSCODE_GOOD) {
UA_free(dsc);
return res;
}
res = __UA_Client_AsyncService(client, &request,
&UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSREQUEST],
Subscriptions_delete_handler,
&UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSRESPONSE],
dsc, requestId);
if(res != UA_STATUSCODE_GOOD) {
UA_DeleteSubscriptionsRequest_clear(&dsc->request);
UA_free(dsc);
}
return res;
}
UA_DeleteSubscriptionsResponse
UA_Client_Subscriptions_delete(UA_Client *client,
const UA_DeleteSubscriptionsRequest request) {
lockClient(client);
UA_DeleteSubscriptionsResponse response;
__Client_Service(client, &request,
&UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSREQUEST],
&response, &UA_TYPES[UA_TYPES_DELETESUBSCRIPTIONSRESPONSE]);
__Client_Subscription_processDelete(client, &request, &response);
unlockClient(client);
return response;
}
UA_StatusCode
UA_Client_Subscriptions_deleteSingle(UA_Client *client, UA_UInt32 subscriptionId) {
UA_DeleteSubscriptionsRequest request;
UA_DeleteSubscriptionsRequest_init(&request);
request.subscriptionIds = &subscriptionId;
request.subscriptionIdsSize = 1;
UA_DeleteSubscriptionsResponse response =
UA_Client_Subscriptions_delete(client, request);
UA_StatusCode retval = response.responseHeader.serviceResult;
if(retval != UA_STATUSCODE_GOOD) {
UA_DeleteSubscriptionsResponse_clear(&response);
return retval;
}
if(response.resultsSize != 1) {
UA_DeleteSubscriptionsResponse_clear(&response);
return UA_STATUSCODE_BADINTERNALERROR;
}
retval = response.results[0];
UA_DeleteSubscriptionsResponse_clear(&response);
return retval;
}
static void
EventFields_clear(UA_KeyValueMap *eventFields) {
for(size_t i = 0; i < eventFields->mapSize; i++)
UA_Variant_init(&eventFields->map[i].value);
UA_KeyValueMap_clear(eventFields);
}
static void
MonitoredItem_delete(UA_Client *client, UA_Client_Subscription *sub,
UA_Client_MonitoredItem *mon) {
UA_LOCK_ASSERT(&client->clientMutex);
ZIP_REMOVE(MonitorItemsTree, &sub->monitoredItems, mon);
if(mon->deleteCallback)
mon->deleteCallback(client, sub->subscriptionId, sub->context,
mon->monitoredItemId, mon->context);
EventFields_clear(&mon->eventFields);
UA_MonitoringParameters_clear(&mon->parameters);
UA_MonitoringParameters_clear(&mon->pendingParameters);
UA_free(mon);
}
static UA_StatusCode
prepareEventFieldsMap(UA_KeyValueMap *eventFields,
UA_MonitoringParameters *params) {
UA_ExtensionObject *eo = ¶ms->filter;
if(eo->content.decoded.type != &UA_TYPES[UA_TYPES_EVENTFILTER])
return UA_STATUSCODE_GOOD;
UA_EventFilter *ef = (UA_EventFilter*)eo->content.decoded.data;
UA_StatusCode res = UA_STATUSCODE_GOOD;
if(ef->selectClausesSize == 0)
return UA_STATUSCODE_GOOD;
eventFields->map = (UA_KeyValuePair*)
UA_calloc(ef->selectClausesSize, sizeof(UA_KeyValuePair));
if(!eventFields->map)
return UA_STATUSCODE_BADOUTOFMEMORY;
eventFields->mapSize = ef->selectClausesSize;
for(size_t i = 0; i < eventFields->mapSize; i++) {
res |= UA_SimpleAttributeOperand_print(&ef->selectClauses[i],
&eventFields->map[i].key.name);
}
return res;
}
static UA_StatusCode
MonitoredItem_createBegin(UA_Client *client, UA_Client_Subscription *sub,
UA_MonitoredItemCreateRequest *item,
UA_Client_DeleteMonitoredItemCallback deleteCallback,
void *context, void *handlingCallback,
UA_Client_MonitoredItem **outMon) {
UA_Client_MonitoredItem *mon = (UA_Client_MonitoredItem *)
UA_calloc(1, sizeof(UA_Client_MonitoredItem));
if(!mon)
return UA_STATUSCODE_BADOUTOFMEMORY;
item->requestedParameters.clientHandle = ++client->monitoredItemHandles;
UA_StatusCode res =
UA_MonitoringParameters_copy(&item->requestedParameters,
&mon->parameters);
if(res != UA_STATUSCODE_GOOD) {
UA_free(mon);
return res;
}
mon->isEventMonitoredItem =
(item->itemToMonitor.attributeId == UA_ATTRIBUTEID_EVENTNOTIFIER);
mon->context = context;
mon->deleteCallback = deleteCallback;
mon->handler.dataChangeCallback =
(UA_Client_DataChangeNotificationCallback)(uintptr_t)handlingCallback;
ZIP_INSERT(MonitorItemsTree, &sub->monitoredItems, mon);
UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
"Subscription %" PRIu32 " | Added a MonitoredItem with handle %" PRIu32,
sub->subscriptionId, mon->parameters.clientHandle);
*outMon = mon;
return UA_STATUSCODE_GOOD;
}
static void
MonitoredItem_createFinish(UA_Client *client, UA_Client_Subscription *sub,
UA_Client_MonitoredItem *mon,
UA_MonitoredItemCreateResult *result) {
UA_assert(result->statusCode == UA_STATUSCODE_GOOD);
mon->monitoredItemId = result->monitoredItemId;
}
static void
Client_MonitoredItems_create(UA_Client *client,
const UA_CreateMonitoredItemsRequest *constRequest,
void **contexts, void **handlingCallbacks,
UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
UA_CreateMonitoredItemsResponse *response) {
UA_LOCK_ASSERT(&client->clientMutex);
UA_CreateMonitoredItemsResponse_init(response);
if(constRequest->itemsToCreateSize == 0) {
response->responseHeader.serviceResult = UA_STATUSCODE_BADNOTHINGTODO;
return;
}
UA_Client_Subscription *sub =
findSubscriptionById(client, constRequest->subscriptionId);
if(!sub) {
response->responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
return;
}
UA_CreateMonitoredItemsRequest request;
UA_StatusCode res =
UA_CreateMonitoredItemsRequest_copy(constRequest, &request);
if(res != UA_STATUSCODE_GOOD) {
response->responseHeader.serviceResult = res;
return;
}
UA_STACKARRAY(UA_Client_MonitoredItem*, mons, request.itemsToCreateSize);
memset(mons, 0, sizeof(UA_Client_MonitoredItem*) * request.itemsToCreateSize);
for(size_t i = 0; i < request.itemsToCreateSize; i++) {
void *context = (contexts) ? contexts[i] : NULL;
void *handlingCallback = (handlingCallbacks) ? handlingCallbacks[i] : NULL;
UA_Client_DeleteMonitoredItemCallback deleteCallback =
(deleteCallbacks) ? deleteCallbacks[i] : NULL;
res |= MonitoredItem_createBegin(client, sub, &request.itemsToCreate[i],
deleteCallback, context,
handlingCallback, &mons[i]);
}
if(res != UA_STATUSCODE_GOOD) {
for(size_t i = 0; i < request.itemsToCreateSize; i++) {
if(mons[i])
MonitoredItem_delete(client, sub, mons[i]);
else if(deleteCallbacks && contexts)
deleteCallbacks[i](client, request.subscriptionId,
sub->context, 0, contexts[i]);
}
UA_CreateMonitoredItemsRequest_clear(&request);
response->responseHeader.serviceResult = res;
return;
}
__Client_Service(client, &request,
&UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSREQUEST],
response, &UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSRESPONSE]);
if(response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
response->resultsSize != request.itemsToCreateSize)
response->responseHeader.serviceResult = UA_STATUSCODE_BADINTERNALERROR;
for(size_t i = 0; i < request.itemsToCreateSize; i++) {
UA_assert(mons[i]);
UA_MonitoredItemCreateResult *item = &response->results[i];
if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD ||
item->statusCode != UA_STATUSCODE_GOOD) {
MonitoredItem_delete(client, sub, mons[i]);
continue;
}
MonitoredItem_createFinish(client, sub, mons[i], item);
}
UA_CreateMonitoredItemsRequest_clear(&request);
}
UA_CreateMonitoredItemsResponse
UA_Client_MonitoredItems_createDataChanges(UA_Client *client,
const UA_CreateMonitoredItemsRequest request,
void **contexts,
UA_Client_DataChangeNotificationCallback *callbacks,
UA_Client_DeleteMonitoredItemCallback *deleteCallbacks) {
UA_CreateMonitoredItemsResponse response;
lockClient(client);
Client_MonitoredItems_create(client, &request, contexts, (void **)callbacks,
deleteCallbacks, &response);
unlockClient(client);
return response;
}
UA_MonitoredItemCreateResult
UA_Client_MonitoredItems_createDataChange(UA_Client *client, UA_UInt32 subscriptionId,
UA_TimestampsToReturn timestampsToReturn,
const UA_MonitoredItemCreateRequest item,
void *context,
UA_Client_DataChangeNotificationCallback callback,
UA_Client_DeleteMonitoredItemCallback deleteCallback) {
UA_CreateMonitoredItemsRequest request;
UA_CreateMonitoredItemsRequest_init(&request);
request.subscriptionId = subscriptionId;
request.timestampsToReturn = timestampsToReturn;
request.itemsToCreate = (UA_MonitoredItemCreateRequest*)(uintptr_t)&item;
request.itemsToCreateSize = 1;
UA_CreateMonitoredItemsResponse response =
UA_Client_MonitoredItems_createDataChanges(client, request, &context,
&callback, &deleteCallback);
UA_MonitoredItemCreateResult result;
UA_MonitoredItemCreateResult_init(&result);
if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD)
result.statusCode = response.responseHeader.serviceResult;
if(result.statusCode == UA_STATUSCODE_GOOD &&
response.resultsSize != 1)
result.statusCode = UA_STATUSCODE_BADINTERNALERROR;
if(result.statusCode == UA_STATUSCODE_GOOD) {
result = response.results[0];
UA_MonitoredItemCreateResult_init(&response.results[0]);
}
UA_CreateMonitoredItemsResponse_clear(&response);
return result;
}
UA_CreateMonitoredItemsResponse
UA_Client_MonitoredItems_createEvents(UA_Client *client,
const UA_CreateMonitoredItemsRequest request,
void **contexts,
UA_Client_EventNotificationCallback *callback,
UA_Client_DeleteMonitoredItemCallback *deleteCallback) {
UA_CreateMonitoredItemsResponse response;
lockClient(client);
Client_MonitoredItems_create(client, &request, contexts, (void **)callback,
deleteCallback, &response);
unlockClient(client);
return response;
}
UA_MonitoredItemCreateResult
UA_Client_MonitoredItems_createEvent(UA_Client *client, UA_UInt32 subscriptionId,
UA_TimestampsToReturn timestampsToReturn,
const UA_MonitoredItemCreateRequest item, void *context,
UA_Client_EventNotificationCallback callback,
UA_Client_DeleteMonitoredItemCallback deleteCallback) {
UA_CreateMonitoredItemsRequest request;
UA_CreateMonitoredItemsRequest_init(&request);
request.subscriptionId = subscriptionId;
request.timestampsToReturn = timestampsToReturn;
request.itemsToCreate = (UA_MonitoredItemCreateRequest*)(uintptr_t)&item;
request.itemsToCreateSize = 1;
UA_CreateMonitoredItemsResponse response =
UA_Client_MonitoredItems_createEvents(client, request, &context,
&callback, &deleteCallback);
UA_MonitoredItemCreateResult result;
UA_MonitoredItemCreateResult_init(&result);
if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD)
result.statusCode = response.responseHeader.serviceResult;
if(result.statusCode == UA_STATUSCODE_GOOD &&
response.resultsSize != 1)
result.statusCode = UA_STATUSCODE_BADINTERNALERROR;
if(result.statusCode == UA_STATUSCODE_GOOD) {
result = response.results[0];
UA_MonitoredItemCreateResult_init(&response.results[0]);
}
UA_CreateMonitoredItemsResponse_clear(&response);
return result;
}
static void
MonitoredItems_create_async_handler(UA_Client *client, void *data,
UA_UInt32 requestId, void *resp) {
CustomCallback *cc = (CustomCallback*)data;
UA_UInt32 *handles = (UA_UInt32*)cc->clientData;
UA_CreateMonitoredItemsResponse *response =
(UA_CreateMonitoredItemsResponse *)resp;
lockClient(client);
UA_UInt32 subId = handles[0];
UA_UInt32 monSize = handles[1];
UA_UInt32 *monHandles = handles + 2;
if(response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
response->resultsSize != monSize)
response->responseHeader.serviceResult = UA_STATUSCODE_BADINTERNALERROR;
UA_Client_Subscription *sub = findSubscriptionById(client, subId);
if(!sub)
response->responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
for(size_t i = 0; sub && i < monSize; i++) {
UA_Client_MonitoredItem *mon =
findMonitoredItemByHandle(sub, monHandles[i]);
if(!mon)
continue;
UA_MonitoredItemCreateResult *item = &response->results[i];
if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD ||
item->statusCode != UA_STATUSCODE_GOOD) {
MonitoredItem_delete(client, sub, mon);
continue;
}
MonitoredItem_createFinish(client, sub, mon, item);
}
if(cc->userCallback) {
UA_ClientAsyncCreateMonitoredItemsCallback cb =
(UA_ClientAsyncCreateMonitoredItemsCallback)cc->userCallback;
cb(client, cc->userData, requestId, response);
}
UA_free(handles);
UA_free(cc);
unlockClient(client);
}
static UA_StatusCode
Client_MonitoredItems_createAsync(UA_Client *client,
const UA_CreateMonitoredItemsRequest *constRequest,
void **contexts, void **handlingCallbacks,
UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
UA_ClientAsyncServiceCallback createCallback,
void *userdata, UA_UInt32 *requestId) {
UA_LOCK_ASSERT(&client->clientMutex);
UA_Client_Subscription *sub =
findSubscriptionById(client, constRequest->subscriptionId);
if(!sub)
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
if(constRequest->itemsToCreateSize == 0)
return UA_STATUSCODE_BADNOTHINGTODO;
UA_CreateMonitoredItemsRequest request;
UA_StatusCode res =
UA_CreateMonitoredItemsRequest_copy(constRequest, &request);
if(res != UA_STATUSCODE_GOOD)
return res;
CustomCallback *cc = (CustomCallback*)UA_calloc(1, sizeof(CustomCallback));
if(!cc) {
UA_CreateMonitoredItemsRequest_clear(&request);
return UA_STATUSCODE_BADOUTOFMEMORY;
}
UA_UInt32 *handles = (UA_UInt32*)
UA_malloc(sizeof(UA_UInt32) * (request.itemsToCreateSize + 2));
if(!handles) {
UA_free(cc);
UA_CreateMonitoredItemsRequest_clear(&request);
return UA_STATUSCODE_BADOUTOFMEMORY;
}
UA_STACKARRAY(UA_Client_MonitoredItem*, mons, request.itemsToCreateSize);
memset(mons, 0, sizeof(UA_Client_MonitoredItem*) * request.itemsToCreateSize);
for(size_t i = 0; i < request.itemsToCreateSize; i++) {
void *context = (contexts) ? contexts[i] : NULL;
void *handlingCallback = (handlingCallbacks) ? handlingCallbacks[i] : NULL;
UA_Client_DeleteMonitoredItemCallback deleteCallback =
(deleteCallbacks) ? deleteCallbacks[i] : NULL;
res |= MonitoredItem_createBegin(client, sub, &request.itemsToCreate[i],
deleteCallback, context,
handlingCallback, &mons[i]);
}
if(res != UA_STATUSCODE_GOOD) {
for(size_t i = 0; i < request.itemsToCreateSize; i++) {
if(mons[i])
MonitoredItem_delete(client, sub, mons[i]);
else if(deleteCallbacks && contexts)
deleteCallbacks[i](client, request.subscriptionId,
sub->context, 0, contexts[i]);
}
UA_free(handles);
UA_free(cc);
UA_CreateMonitoredItemsRequest_clear(&request);
return res;
}
handles[0] = sub->subscriptionId;
handles[1] = (UA_UInt32)request.itemsToCreateSize;
for(size_t i = 0; i < request.itemsToCreateSize; i++) {
handles[i+2] = mons[i]->parameters.clientHandle;
}
cc->clientData = handles;
cc->userCallback = createCallback;
cc->userData = userdata;
res = __Client_AsyncService(client, &request,
&UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSREQUEST],
MonitoredItems_create_async_handler,
&UA_TYPES[UA_TYPES_CREATEMONITOREDITEMSRESPONSE],
cc, requestId);
if(res != UA_STATUSCODE_GOOD) {
for(size_t i = 0; i < request.itemsToCreateSize; i++) {
UA_assert(mons[i]);
MonitoredItem_delete(client, sub, mons[i]);
}
UA_free(handles);
UA_free(cc);
}
UA_CreateMonitoredItemsRequest_clear(&request);
return res;
}
UA_StatusCode
UA_Client_MonitoredItems_createDataChanges_async(UA_Client *client,
const UA_CreateMonitoredItemsRequest request,
void **contexts,
UA_Client_DataChangeNotificationCallback *callbacks,
UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
UA_ClientAsyncCreateMonitoredItemsCallback createCallback,
void *userdata, UA_UInt32 *requestId) {
lockClient(client);
UA_StatusCode res =
Client_MonitoredItems_createAsync(client, &request, contexts,
(void **)callbacks, deleteCallbacks,
(UA_ClientAsyncServiceCallback)createCallback,
userdata, requestId);
unlockClient(client);
return res;
}
UA_StatusCode
UA_Client_MonitoredItems_createEvents_async(UA_Client *client,
const UA_CreateMonitoredItemsRequest request,
void **contexts,
UA_Client_EventNotificationCallback *callbacks,
UA_Client_DeleteMonitoredItemCallback *deleteCallbacks,
UA_ClientAsyncCreateMonitoredItemsCallback createCallback,
void *userdata, UA_UInt32 *requestId) {
lockClient(client);
UA_StatusCode res =
Client_MonitoredItems_createAsync(client, &request, contexts,
(void **)callbacks, deleteCallbacks,
(UA_ClientAsyncServiceCallback)createCallback,
userdata, requestId);
unlockClient(client);
return res;
}
static void
MonitoredItems_delete(UA_Client *client, UA_Client_Subscription *sub,
const UA_DeleteMonitoredItemsRequest *request,
const UA_DeleteMonitoredItemsResponse *response) {
#ifdef __clang_analyzer__
return;
#endif
struct UA_Client_MonitoredItem_ForDelete deleteMonitoredItem;
memset(&deleteMonitoredItem, 0, sizeof(struct UA_Client_MonitoredItem_ForDelete));
deleteMonitoredItem.client = client;
deleteMonitoredItem.sub = sub;
for(size_t i = 0; i < response->resultsSize; i++) {
if(response->results[i] != UA_STATUSCODE_GOOD &&
response->results[i] != UA_STATUSCODE_BADMONITOREDITEMIDINVALID) {
continue;
}
deleteMonitoredItem.monitoredItemId = &request->monitoredItemIds[i];
ZIP_ITER(MonitorItemsTree,&sub->monitoredItems,
MonitoredItem_delete_wrapper, &deleteMonitoredItem);
}
}
static void
MonitoredItems_delete_handler(UA_Client *client, void *d, UA_UInt32 requestId, void *r) {
UA_Client_Subscription *sub = NULL;
CustomCallback *cc = (CustomCallback *)d;
UA_DeleteMonitoredItemsResponse *response = (UA_DeleteMonitoredItemsResponse *)r;
UA_DeleteMonitoredItemsRequest *request =
(UA_DeleteMonitoredItemsRequest *)cc->clientData;
lockClient(client);
if(response->responseHeader.serviceResult != UA_STATUSCODE_GOOD)
goto cleanup;
sub = findSubscriptionById(client, request->subscriptionId);
if(!sub) {
UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
"No internal representation of subscription %" PRIu32,
request->subscriptionId);
goto cleanup;
}
MonitoredItems_delete(client, sub, request, response);
cleanup:
if(cc->userCallback)
cc->userCallback(client, cc->userData, requestId, response);
UA_DeleteMonitoredItemsRequest_delete(request);
UA_free(cc);
unlockClient(client);
}
UA_DeleteMonitoredItemsResponse
UA_Client_MonitoredItems_delete(UA_Client *client,
const UA_DeleteMonitoredItemsRequest request) {
UA_DeleteMonitoredItemsResponse response;
__UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSREQUEST],
&response, &UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSRESPONSE]);
if(response.responseHeader.serviceResult != UA_STATUSCODE_GOOD)
return response;
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
if(!sub) {
UA_LOG_INFO(client->config.logging, UA_LOGCATEGORY_CLIENT,
"No internal representation of subscription %" PRIu32,
request.subscriptionId);
unlockClient(client);
return response;
}
MonitoredItems_delete(client, sub, &request, &response);
unlockClient(client);
return response;
}
UA_StatusCode
UA_Client_MonitoredItems_delete_async(UA_Client *client,
const UA_DeleteMonitoredItemsRequest request,
UA_ClientAsyncDeleteMonitoredItemsCallback callback,
void *userdata, UA_UInt32 *requestId) {
CustomCallback *cc = (CustomCallback *)UA_calloc(1, sizeof(CustomCallback));
if(!cc)
return UA_STATUSCODE_BADOUTOFMEMORY;
UA_DeleteMonitoredItemsRequest *req_copy = UA_DeleteMonitoredItemsRequest_new();
if(!req_copy) {
UA_free(cc);
return UA_STATUSCODE_BADOUTOFMEMORY;
}
UA_DeleteMonitoredItemsRequest_copy(&request, req_copy);
cc->clientData = req_copy;
cc->userCallback = (UA_ClientAsyncServiceCallback)callback;
cc->userData = userdata;
UA_StatusCode res =
__UA_Client_AsyncService(client, &request,
&UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSREQUEST],
MonitoredItems_delete_handler,
&UA_TYPES[UA_TYPES_DELETEMONITOREDITEMSRESPONSE],
cc, requestId);
if(res != UA_STATUSCODE_GOOD) {
UA_DeleteMonitoredItemsRequest_delete(req_copy);
UA_free(cc);
}
return res;
}
UA_StatusCode
UA_Client_MonitoredItems_deleteSingle(UA_Client *client, UA_UInt32 subscriptionId,
UA_UInt32 monitoredItemId) {
UA_DeleteMonitoredItemsRequest request;
UA_DeleteMonitoredItemsRequest_init(&request);
request.subscriptionId = subscriptionId;
request.monitoredItemIds = &monitoredItemId;
request.monitoredItemIdsSize = 1;
UA_DeleteMonitoredItemsResponse response =
UA_Client_MonitoredItems_delete(client, request);
UA_StatusCode retval = response.responseHeader.serviceResult;
if(retval != UA_STATUSCODE_GOOD) {
UA_DeleteMonitoredItemsResponse_clear(&response);
return retval;
}
if(response.resultsSize != 1) {
UA_DeleteMonitoredItemsResponse_clear(&response);
return UA_STATUSCODE_BADINTERNALERROR;
}
retval = response.results[0];
UA_DeleteMonitoredItemsResponse_clear(&response);
return retval;
}
static void *
MonitoredItem_findByID(void *data, UA_Client_MonitoredItem *mon) {
UA_UInt32 monitorId = *(UA_UInt32*)data;
if(monitorId && (mon->monitoredItemId == monitorId))
return mon;
return NULL;
}
static UA_Client_MonitoredItem *
findMonitoredItemById(UA_Client_Subscription *sub, UA_UInt32 monitoredItemId) {
return (UA_Client_MonitoredItem *)
ZIP_ITER(MonitorItemsTree, &sub->monitoredItems,
MonitoredItem_findByID, &monitoredItemId);
}
static UA_StatusCode
MonitoredItems_prepareModify(UA_Client *client, UA_Client_Subscription *sub,
UA_ModifyMonitoredItemsRequest *request) {
UA_STACKARRAY(UA_Client_MonitoredItem*, mons, request->itemsToModifySize);
UA_STACKARRAY(UA_MonitoringParameters, preparedParameters,
request->itemsToModifySize);
memset(mons, 0, sizeof(UA_Client_MonitoredItem*) *
request->itemsToModifySize);
memset(preparedParameters, 0, sizeof(UA_MonitoringParameters) *
request->itemsToModifySize);
for(size_t i = 0; i < request->itemsToModifySize; ++i) {
UA_MonitoredItemModifyRequest *mimr = &request->itemsToModify[i];
mimr->requestedParameters.clientHandle = 0;
mons[i] = findMonitoredItemById(sub, mimr->monitoredItemId);
if(!mons[i])
continue;
mimr->requestedParameters.clientHandle = ++client->monitoredItemHandles;
UA_StatusCode res =
UA_MonitoringParameters_copy(&mimr->requestedParameters,
&preparedParameters[i]);
if(res != UA_STATUSCODE_GOOD) {
for(size_t j = 0; j <= i; j++)
UA_MonitoringParameters_clear(&preparedParameters[j]);
return res;
}
}
for(size_t i = 0; i < request->itemsToModifySize; ++i) {
UA_Client_MonitoredItem *mon = mons[i];
if(!mon)
continue;
UA_MonitoringParameters_clear(&mon->pendingParameters);
mon->pendingParameters = preparedParameters[i];
memset(&preparedParameters[i], 0, sizeof(UA_MonitoringParameters));
}
return UA_STATUSCODE_GOOD;
}
static void
MonitoredItems_reconcileModify(UA_Client_Subscription *sub,
const UA_ModifyMonitoredItemsRequest *request,
UA_ModifyMonitoredItemsResponse *response) {
UA_Boolean validResponse = response &&
response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
response->resultsSize == request->itemsToModifySize;
if(response && response->responseHeader.serviceResult == UA_STATUSCODE_GOOD &&
!validResponse)
response->responseHeader.serviceResult = UA_STATUSCODE_BADINTERNALERROR;
for(size_t i = 0; i < request->itemsToModifySize; ++i) {
if(validResponse && response->results[i].statusCode == UA_STATUSCODE_GOOD)
continue;
const UA_MonitoredItemModifyRequest *mimr = &request->itemsToModify[i];
UA_Client_MonitoredItem *mon =
findMonitoredItemById(sub, mimr->monitoredItemId);
if(mon && mon->pendingParameters.clientHandle ==
mimr->requestedParameters.clientHandle)
UA_MonitoringParameters_clear(&mon->pendingParameters);
}
}
UA_ModifyMonitoredItemsResponse
UA_Client_MonitoredItems_modify(UA_Client *client,
const UA_ModifyMonitoredItemsRequest request) {
UA_ModifyMonitoredItemsResponse response;
UA_ModifyMonitoredItemsResponse_init(&response);
UA_ModifyMonitoredItemsRequest modifiedRequest;
UA_StatusCode res = UA_ModifyMonitoredItemsRequest_copy(&request, &modifiedRequest);
if(res != UA_STATUSCODE_GOOD) {
response.responseHeader.serviceResult = res;
return response;
}
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
if(!sub) {
unlockClient(client);
UA_ModifyMonitoredItemsRequest_clear(&modifiedRequest);
response.responseHeader.serviceResult = UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
return response;
}
res = MonitoredItems_prepareModify(client, sub, &modifiedRequest);
if(res != UA_STATUSCODE_GOOD) {
response.responseHeader.serviceResult = res;
unlockClient(client);
UA_ModifyMonitoredItemsRequest_clear(&modifiedRequest);
return response;
}
__Client_Service(client, &modifiedRequest,
&UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSREQUEST], &response,
&UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSRESPONSE]);
MonitoredItems_reconcileModify(sub, &modifiedRequest, &response);
unlockClient(client);
UA_ModifyMonitoredItemsRequest_clear(&modifiedRequest);
return response;
}
static void
MonitoredItems_modify_async_handler(UA_Client *client, void *data,
UA_UInt32 requestId, void *resp) {
CustomCallback *cc = (CustomCallback*)data;
UA_ModifyMonitoredItemsRequest *request =
(UA_ModifyMonitoredItemsRequest*)cc->clientData;
UA_ModifyMonitoredItemsResponse *response =
(UA_ModifyMonitoredItemsResponse*)resp;
lockClient(client);
UA_Client_Subscription *sub =
findSubscriptionById(client, request->subscriptionId);
if(sub)
MonitoredItems_reconcileModify(sub, request, response);
if(cc->userCallback) {
UA_ClientAsyncModifyMonitoredItemsCallback cb =
(UA_ClientAsyncModifyMonitoredItemsCallback)cc->userCallback;
cb(client, cc->userData, requestId, response);
}
UA_ModifyMonitoredItemsRequest_delete(request);
UA_free(cc);
unlockClient(client);
}
UA_StatusCode
UA_Client_MonitoredItems_modify_async(UA_Client *client,
const UA_ModifyMonitoredItemsRequest request,
UA_ClientAsyncModifyMonitoredItemsCallback callback,
void *userdata, UA_UInt32 *requestId) {
UA_ModifyMonitoredItemsRequest *requestCopy =
UA_ModifyMonitoredItemsRequest_new();
if(!requestCopy)
return UA_STATUSCODE_BADOUTOFMEMORY;
UA_StatusCode res =
UA_ModifyMonitoredItemsRequest_copy(&request, requestCopy);
if(res != UA_STATUSCODE_GOOD) {
UA_ModifyMonitoredItemsRequest_delete(requestCopy);
return res;
}
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, request.subscriptionId);
if(!sub) {
unlockClient(client);
UA_ModifyMonitoredItemsRequest_delete(requestCopy);
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
}
res = MonitoredItems_prepareModify(client, sub, requestCopy);
if(res != UA_STATUSCODE_GOOD) {
unlockClient(client);
UA_ModifyMonitoredItemsRequest_delete(requestCopy);
return res;
}
CustomCallback *cc = (CustomCallback*)UA_calloc(1, sizeof(CustomCallback));
if(!cc) {
MonitoredItems_reconcileModify(sub, requestCopy, NULL);
unlockClient(client);
UA_ModifyMonitoredItemsRequest_delete(requestCopy);
return UA_STATUSCODE_BADOUTOFMEMORY;
}
cc->clientData = requestCopy;
cc->userCallback = (UA_ClientAsyncServiceCallback)callback;
cc->userData = userdata;
UA_StatusCode statusCode = __Client_AsyncService(
client, requestCopy, &UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSREQUEST],
MonitoredItems_modify_async_handler,
&UA_TYPES[UA_TYPES_MODIFYMONITOREDITEMSRESPONSE], cc, requestId);
if(statusCode != UA_STATUSCODE_GOOD) {
MonitoredItems_reconcileModify(sub, requestCopy, NULL);
UA_ModifyMonitoredItemsRequest_delete(requestCopy);
UA_free(cc);
}
unlockClient(client);
return statusCode;
}
UA_StatusCode
UA_Client_MonitoredItem_getContext(UA_Client *client, UA_UInt32 subscriptionId,
UA_UInt32 monitoredItemId, void **monContext) {
if(!client || !monContext)
return UA_STATUSCODE_BADINVALIDARGUMENT;
*monContext = NULL;
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
if(!sub) {
unlockClient(client);
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
}
UA_StatusCode itemstatus = UA_STATUSCODE_BADMONITOREDITEMIDINVALID;
UA_Client_MonitoredItem *monItem = findMonitoredItemById(sub, monitoredItemId);
if(monItem) {
*monContext = monItem->context;
itemstatus = UA_STATUSCODE_GOOD;
}
unlockClient(client);
return itemstatus;
}
UA_StatusCode
UA_Client_MonitoredItem_setContext(UA_Client *client, UA_UInt32 subscriptionId,
UA_UInt32 monitoredItemId, void *monContext) {
if(!client)
return UA_STATUSCODE_BADINVALIDARGUMENT;
lockClient(client);
UA_Client_Subscription *sub = findSubscriptionById(client, subscriptionId);
if(!sub) {
unlockClient(client);
return UA_STATUSCODE_BADSUBSCRIPTIONIDINVALID;
}
UA_StatusCode itemstatus = UA_STATUSCODE_BADMONITOREDITEMIDINVALID;
UA_Client_MonitoredItem *monItem = findMonitoredItemById(sub, monitoredItemId);
if(monItem) {
monItem->context = monContext;
itemstatus = UA_STATUSCODE_GOOD;
}
unlockClient(client);
return itemstatus;
}
UA_StatusCode
__Client_preparePublishRequest(UA_Client *client, UA_PublishRequest *request) {
UA_LOCK_ASSERT(&client->clientMutex);
UA_Client_NotificationsAckNumber *ack;
LIST_FOREACH(ack, &client->pendingNotificationsAcks, listEntry)
++request->subscriptionAcknowledgementsSize;
request->subscriptionAcknowledgements = (UA_SubscriptionAcknowledgement*)
UA_Array_new(request->subscriptionAcknowledgementsSize,
&UA_TYPES[UA_TYPES_SUBSCRIPTIONACKNOWLEDGEMENT]);
if(!request->subscriptionAcknowledgements) {
request->subscriptionAcknowledgementsSize = 0;
return UA_STATUSCODE_BADOUTOFMEMORY;
}
size_t i = 0;
UA_Client_NotificationsAckNumber *ack_tmp;
LIST_FOREACH_SAFE(ack, &client->pendingNotificationsAcks, listEntry, ack_tmp) {
LIST_REMOVE(ack, listEntry);
UA_SubscriptionAcknowledgement *reqAck = &request->subscriptionAcknowledgements[i];
reqAck->sequenceNumber = ack->subAck.sequenceNumber;
reqAck->subscriptionId = ack->subAck.subscriptionId;
UA_free(ack);
i++;
}
return UA_STATUSCODE_GOOD;
}
static UA_UInt32
__nextSequenceNumber(UA_UInt32 sequenceNumber) {
UA_UInt32 nextSequenceNumber = sequenceNumber + 1;
if(nextSequenceNumber == 0)
nextSequenceNumber = 1;
return nextSequenceNumber;
}
static UA_Client_MonitoredItem *
findMonitoredItemForNotification(UA_Client_Subscription *sub,
UA_UInt32 clientHandle) {
UA_Client_MonitoredItem *mon =
findMonitoredItemByHandle(sub, clientHandle);
if(mon)
return mon;
mon = (UA_Client_MonitoredItem*)
ZIP_ITER(MonitorItemsTree, &sub->monitoredItems,
MonitoredItem_findPendingByHandle, &clientHandle);
if(!mon)
return NULL;
ZIP_REMOVE(MonitorItemsTree, &sub->monitoredItems, mon);
UA_MonitoringParameters_clear(&mon->parameters);
mon->parameters = mon->pendingParameters;
UA_MonitoringParameters_init(&mon->pendingParameters);
EventFields_clear(&mon->eventFields);
ZIP_INSERT(MonitorItemsTree, &sub->monitoredItems, mon);
return mon;
}
static void
processDataChangeNotification(UA_Client *client, UA_Client_Subscription *sub,
UA_DataChangeNotification *dataChangeNotification) {
UA_LOCK_ASSERT(&client->clientMutex);
for(size_t j = 0; j < dataChangeNotification->monitoredItemsSize; ++j) {
UA_MonitoredItemNotification *min = &dataChangeNotification->monitoredItems[j];
UA_Client_MonitoredItem *mon =
findMonitoredItemForNotification(sub, min->clientHandle);
if(!mon) {
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"Could not process a notification with clienthandle %" PRIu32
" on subscription %" PRIu32, min->clientHandle, sub->subscriptionId);
continue;
}
if(mon->isEventMonitoredItem) {
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"MonitoredItem is configured for Events. But received a "
"DataChangeNotification.");
continue;
}
if(mon->handler.dataChangeCallback) {
void *subC = sub->context;
void *monC = mon->context;
UA_UInt32 subId = sub->subscriptionId;
UA_UInt32 monId = mon->monitoredItemId;
mon->handler.dataChangeCallback(client, subId, subC, monId, monC, &min->value);
}
}
}
static void
processEventNotification(UA_Client *client, UA_Client_Subscription *sub,
UA_EventNotificationList *eventNotificationList) {
UA_LOCK_ASSERT(&client->clientMutex);
for(size_t j = 0; j < eventNotificationList->eventsSize; ++j) {
UA_EventFieldList *efl = &eventNotificationList->events[j];
UA_Client_MonitoredItem *mon =
findMonitoredItemForNotification(sub, efl->clientHandle);
if(!mon) {
UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
"Could not process a notification with clienthandle %" PRIu32
" on subscription %" PRIu32, efl->clientHandle,
sub->subscriptionId);
continue;
}
if(!mon->isEventMonitoredItem) {
UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
"MonitoredItem is configured for DataChanges. But received a "
"EventNotification");
continue;
}
if(mon->eventFields.mapSize == 0) {
UA_StatusCode res =
prepareEventFieldsMap(&mon->eventFields, &mon->parameters);
if(res != UA_STATUSCODE_GOOD) {
EventFields_clear(&mon->eventFields);
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"Could not prepare event fields: %s",
UA_StatusCode_name(res));
continue;
}
}
if(mon->eventFields.mapSize != efl->eventFieldsSize) {
UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
"MonitoredItem received a EventNotification with the "
"wrong number of event fields");
continue;
}
for(size_t i = 0; i < mon->eventFields.mapSize; i++)
mon->eventFields.map[i].value = efl->eventFields[i];
mon->handler.eventCallback(client, sub->subscriptionId, sub->context,
mon->monitoredItemId, mon->context,
mon->eventFields);
}
}
static void
processNotificationMessage(UA_Client *client, UA_Client_Subscription *sub,
UA_ExtensionObject *msg) {
UA_LOCK_ASSERT(&client->clientMutex);
if(msg->encoding != UA_EXTENSIONOBJECT_DECODED)
return;
if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_DATACHANGENOTIFICATION]) {
UA_DataChangeNotification *dataChangeNotification =
(UA_DataChangeNotification *)msg->content.decoded.data;
processDataChangeNotification(client, sub, dataChangeNotification);
return;
}
if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_EVENTNOTIFICATIONLIST]) {
UA_EventNotificationList *eventNotificationList =
(UA_EventNotificationList *)msg->content.decoded.data;
processEventNotification(client, sub, eventNotificationList);
return;
}
if(msg->content.decoded.type == &UA_TYPES[UA_TYPES_STATUSCHANGENOTIFICATION]) {
if(sub->statusChangeCallback) {
void *subC = sub->context;
UA_UInt32 subId = sub->subscriptionId;
sub->statusChangeCallback(client, subId, subC,
(UA_StatusChangeNotification*)msg->content.decoded.data);
} else {
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"Dropped a StatusChangeNotification since no "
"callback is registered");
}
return;
}
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"Unknown notification message type");
}
void
__Client_Subscriptions_processPublishResponse(UA_Client *client, UA_PublishRequest *request,
UA_PublishResponse *response) {
UA_LOCK_ASSERT(&client->clientMutex);
client->currentlyOutStandingPublishRequests--;
switch(response->responseHeader.serviceResult) {
case UA_STATUSCODE_BADTOOMANYPUBLISHREQUESTS:
if(client->config.outStandingPublishRequests > 1) {
client->config.outStandingPublishRequests--;
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Too many PublishRequest, reduce "
"outStandingPublishRequests to %" PRId16,
client->config.outStandingPublishRequests);
} else {
UA_LOG_ERROR(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Too many PublishRequests when "
"outStandingPublishRequests = 1");
UA_Client_Subscriptions_deleteSingle(client, response->subscriptionId);
}
return;
case UA_STATUSCODE_BADSHUTDOWN:
UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Received BadShutdown status");
return;
case UA_STATUSCODE_BADNOSUBSCRIPTION:
UA_LOG_DEBUG(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Received BadNoSubscription status");
return;
default:
break;
}
UA_Client_Subscription *sub = findSubscriptionById(client, response->subscriptionId);
if(!sub) {
response->responseHeader.serviceResult = UA_STATUSCODE_BADNOSUBSCRIPTION;
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Received response for an unknown Subscription");
return;
}
switch(response->responseHeader.serviceResult) {
case UA_STATUSCODE_BADSESSIONCLOSED:
__Client_Subscription_deleteInternal(client, sub);
return;
case UA_STATUSCODE_BADTIMEOUT:
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Aborted with BadTimeout status");
if(client->config.subscriptionInactivityCallback) {
void *subC = sub->context;
UA_UInt32 subId = sub->subscriptionId;
client->config.subscriptionInactivityCallback(client, subId, subC);
}
return;
case UA_STATUSCODE_GOOD:
break;
default:
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Received %s status",
UA_StatusCode_name(response->responseHeader.serviceResult));
return;
}
UA_EventLoop *el = client->config.eventLoop;
sub->lastActivity = el->dateTime_nowMonotonic(el);
UA_NotificationMessage *msg = &response->notificationMessage;
if(__nextSequenceNumber(sub->sequenceNumber) != msg->sequenceNumber) {
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Invalid subscription sequence number: "
"Expected %" PRIu32 " but got %" PRIu32,
__nextSequenceNumber(sub->sequenceNumber), msg->sequenceNumber);
}
if(msg->notificationDataSize)
sub->sequenceNumber = msg->sequenceNumber;
for(size_t k = 0; k < msg->notificationDataSize; ++k)
processNotificationMessage(client, sub, &msg->notificationData[k]);
for(size_t i = 0; i < response->availableSequenceNumbersSize; i++) {
if(response->availableSequenceNumbers[i] != msg->sequenceNumber)
continue;
UA_Client_NotificationsAckNumber *tmpAck = (UA_Client_NotificationsAckNumber*)
UA_malloc(sizeof(UA_Client_NotificationsAckNumber));
if(!tmpAck) {
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"PublishResponse: Not enough memory to store the pending "
"acknowledgement for Subscription %" PRIu32, sub->subscriptionId);
break;
}
tmpAck->subAck.sequenceNumber = msg->sequenceNumber;
tmpAck->subAck.subscriptionId = sub->subscriptionId;
LIST_INSERT_HEAD(&client->pendingNotificationsAcks, tmpAck, listEntry);
break;
}
}
static void
processPublishResponseAsync(UA_Client *client, void *userdata,
UA_UInt32 requestId, void *response) {
UA_PublishRequest *req = (UA_PublishRequest*)userdata;
UA_PublishResponse *res = (UA_PublishResponse*)response;
lockClient(client);
__Client_Subscriptions_processPublishResponse(client, req, res);
UA_PublishRequest_delete(req);
__Client_Subscriptions_backgroundPublish(client);
unlockClient(client);
}
void
__Client_Subscriptions_clear(UA_Client *client) {
UA_Client_NotificationsAckNumber *n;
UA_Client_NotificationsAckNumber *tmp;
LIST_FOREACH_SAFE(n, &client->pendingNotificationsAcks, listEntry, tmp) {
LIST_REMOVE(n, listEntry);
UA_free(n);
}
UA_Client_Subscription *sub;
UA_Client_Subscription *tmps;
LIST_FOREACH_SAFE(sub, &client->subscriptions, listEntry, tmps)
__Client_Subscription_deleteInternal(client, sub);
client->monitoredItemHandles = 0;
}
void
__Client_Subscriptions_backgroundPublishInactivityCheck(UA_Client *client) {
UA_LOCK_ASSERT(&client->clientMutex);
UA_EventLoop *el = client->config.eventLoop;
UA_DateTime nowm = el->dateTime_nowMonotonic(el);
UA_Client_Subscription *sub;
LIST_FOREACH(sub, &client->subscriptions, listEntry) {
UA_DateTime maxSilence = (UA_DateTime)
((sub->publishingInterval * sub->maxKeepAliveCount) +
client->config.timeout) * UA_DATETIME_MSEC;
if(maxSilence + sub->lastActivity < nowm) {
sub->lastActivity = nowm;
if(client->config.subscriptionInactivityCallback) {
void *subC = sub->context;
UA_UInt32 subId = sub->subscriptionId;
client->config.subscriptionInactivityCallback(client, subId, subC);
}
UA_LOG_WARNING(client->config.logging, UA_LOGCATEGORY_CLIENT,
"Inactivity for Subscription %" PRIu32 ".",
sub->subscriptionId);
}
}
}
void
__Client_Subscriptions_backgroundPublish(UA_Client *client) {
UA_LOCK_ASSERT(&client->clientMutex);
if(client->sessionState != UA_SESSIONSTATE_ACTIVATED)
return;
if(!LIST_FIRST(&client->subscriptions))
return;
while(client->currentlyOutStandingPublishRequests < client->config.outStandingPublishRequests) {
UA_PublishRequest *request = UA_PublishRequest_new();
if(!request)
return;
request->requestHeader.timeoutHint = 10 * 60 * 1000;
UA_StatusCode retval = __Client_preparePublishRequest(client, request);
if(retval != UA_STATUSCODE_GOOD) {
UA_PublishRequest_delete(request);
return;
}
retval = __Client_AsyncService(client, request,
&UA_TYPES[UA_TYPES_PUBLISHREQUEST],
processPublishResponseAsync,
&UA_TYPES[UA_TYPES_PUBLISHRESPONSE],
(void*)request, NULL);
if(retval != UA_STATUSCODE_GOOD) {
UA_PublishRequest_delete(request);
return;
}
client->currentlyOutStandingPublishRequests++;
}
}
UA_SetPublishingModeResponse
UA_Client_Subscriptions_setPublishingMode(UA_Client *client,
const UA_SetPublishingModeRequest request) {
UA_SetPublishingModeResponse response;
__UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_SETPUBLISHINGMODEREQUEST],
&response, &UA_TYPES[UA_TYPES_SETPUBLISHINGMODERESPONSE]);
return response;
}
UA_SetMonitoringModeResponse
UA_Client_MonitoredItems_setMonitoringMode(UA_Client *client,
const UA_SetMonitoringModeRequest request) {
UA_SetMonitoringModeResponse response;
__UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_SETMONITORINGMODEREQUEST],
&response, &UA_TYPES[UA_TYPES_SETMONITORINGMODERESPONSE]);
return response;
}
UA_StatusCode
UA_Client_MonitoredItems_setMonitoringMode_async(UA_Client *client,
const UA_SetMonitoringModeRequest request,
UA_ClientAsyncSetMonitoringModeCallback callback,
void *userdata, UA_UInt32 *requestId) {
return __UA_Client_AsyncService(client, &request,
&UA_TYPES[UA_TYPES_SETMONITORINGMODEREQUEST],
(UA_ClientAsyncServiceCallback)callback,
&UA_TYPES[UA_TYPES_SETMONITORINGMODERESPONSE],
userdata, requestId);
}
UA_SetTriggeringResponse
UA_Client_MonitoredItems_setTriggering(UA_Client *client,
const UA_SetTriggeringRequest request) {
UA_SetTriggeringResponse response;
__UA_Client_Service(client, &request, &UA_TYPES[UA_TYPES_SETTRIGGERINGREQUEST],
&response, &UA_TYPES[UA_TYPES_SETTRIGGERINGRESPONSE]);
return response;
}
UA_StatusCode
UA_Client_MonitoredItems_setTriggering_async(UA_Client *client,
const UA_SetTriggeringRequest request,
UA_ClientAsyncSetTriggeringCallback callback,
void *userdata, UA_UInt32 *requestId) {
return __UA_Client_AsyncService(client, &request,
&UA_TYPES[UA_TYPES_SETTRIGGERINGREQUEST],
(UA_ClientAsyncServiceCallback)callback,
&UA_TYPES[UA_TYPES_SETTRIGGERINGRESPONSE],
userdata, requestId);
}