Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 53 additions & 24 deletions include/neuron/msg.h
Original file line number Diff line number Diff line change
Expand Up @@ -888,10 +888,11 @@ typedef struct neu_resp_tag_value {
typedef neu_resp_tag_value_t neu_tag_value_t;

typedef struct neu_resp_tag_value_meta {
char tag[NEU_TAG_NAME_LEN];
neu_dvalue_t value;
neu_tag_meta_t metas[NEU_TAG_META_SIZE];
neu_datatag_t datatag;
char tag[NEU_TAG_NAME_LEN];
neu_dvalue_t value;
neu_tag_meta_t *metas;
int n_meta;
neu_datatag_t datatag;
} neu_resp_tag_value_meta_t;

static inline UT_icd *neu_resp_tag_value_meta_icd()
Expand Down Expand Up @@ -962,10 +963,11 @@ typedef struct {
} neu_resp_read_group_t;

typedef struct neu_resp_tag_value_meta_paginate {
char tag[NEU_TAG_NAME_LEN];
neu_dvalue_t value;
neu_tag_meta_t metas[NEU_TAG_META_SIZE];
neu_datatag_t datatag;
char tag[NEU_TAG_NAME_LEN];
neu_dvalue_t value;
neu_tag_meta_t *metas;
int n_meta;
neu_datatag_t datatag;
} neu_resp_tag_value_meta_paginate_t;

static inline UT_icd *neu_resp_tag_value_meta_paginate_icd()
Expand Down Expand Up @@ -996,6 +998,13 @@ static inline void neu_resp_read_free(neu_resp_read_group_t *resp)
tag_value->value.type < NEU_TYPE_ARRAY_STRING) {
free(tag_value->value.value.bools.bools);
}

if (tag_value->metas != NULL) {
for (int i = 0; i < tag_value->n_meta; i++) {
neu_free_dvalue(&tag_value->metas[i].value);
}
free(tag_value->metas);
}
}
free(resp->driver);
free(resp->group);
Expand All @@ -1013,6 +1022,13 @@ neu_resp_read_paginate_free(neu_resp_read_group_paginate_t *resp)
tag_value->value.type < NEU_TYPE_ARRAY_STRING) {
free(tag_value->value.value.bools.bools);
}

if (tag_value->metas != NULL) {
for (int i = 0; i < tag_value->n_meta; i++) {
neu_free_dvalue(&tag_value->metas[i].value);
}
free(tag_value->metas);
}
}
free(resp->driver);
free(resp->group);
Expand Down Expand Up @@ -1090,6 +1106,13 @@ static inline void neu_trans_data_free(neu_reqresp_trans_data_t *data)
tag_value->value.type < NEU_TYPE_ARRAY_STRING) {
free(tag_value->value.value.bools.bools);
}

if (tag_value->metas != NULL) {
for (int i = 0; i < tag_value->n_meta; i++) {
neu_free_dvalue(&tag_value->metas[i].value);
}
free(tag_value->metas);
}
}
utarray_free(data->tags);
free(data->group);
Expand All @@ -1108,18 +1131,21 @@ static inline void neu_tag_value_to_json(neu_resp_tag_value_meta_t *tag_value,
tag_json->name = tag_value->tag;
tag_json->error = 0;

for (int k = 0; k < NEU_TAG_META_SIZE; k++) {
if (strlen(tag_value->metas[k].name) > 0) {
tag_json->n_meta++;
} else {
break;
}
}
// for (int k = 0; k < NEU_TAG_META_SIZE; k++) {
// if (strlen(tag_value->metas[k].name) > 0) {
// tag_json->n_meta++;
// } else {
// break;
// }
// }

tag_json->n_meta = tag_value->n_meta;

if (tag_json->n_meta > 0) {
tag_json->metas = (neu_json_tag_meta_t *) calloc(
tag_json->n_meta, sizeof(neu_json_tag_meta_t));
}
neu_json_metas_to_json(tag_value->metas, NEU_TAG_META_SIZE, tag_json);
neu_json_metas_to_json(tag_value->metas, tag_value->n_meta, tag_json);

tag_json->datatag.bias = tag_value->datatag.bias;

Expand Down Expand Up @@ -1317,18 +1343,21 @@ neu_tag_value_to_json_paginate(neu_resp_tag_value_meta_paginate_t *tag_value,
memcpy(tag_json->datatag.meta, tag_value->datatag.meta,
NEU_TAG_META_LENGTH);

for (int k = 0; k < NEU_TAG_META_SIZE; k++) {
if (strlen(tag_value->metas[k].name) > 0) {
tag_json->n_meta++;
} else {
break;
}
}
// for (int k = 0; k < NEU_TAG_META_SIZE; k++) {
// if (strlen(tag_value->metas[k].name) > 0) {
// tag_json->n_meta++;
// } else {
// break;
// }
// }

tag_json->n_meta = tag_value->n_meta;

if (tag_json->n_meta > 0) {
tag_json->metas = (neu_json_tag_meta_t *) calloc(
tag_json->n_meta, sizeof(neu_json_tag_meta_t));
}
neu_json_metas_to_json_paginate(tag_value->metas, NEU_TAG_META_SIZE,
neu_json_metas_to_json_paginate(tag_value->metas, tag_value->n_meta,
tag_json);

switch (tag_value->value.type) {
Expand Down
4 changes: 2 additions & 2 deletions plugins/mqtt/mqtt_handle.c
Original file line number Diff line number Diff line change
Expand Up @@ -1186,7 +1186,7 @@ int handle_read_response(neu_plugin_t *plugin, neu_json_mqtt_t *mqtt_json,
break;
}

for (int i = 0; i < NEU_TAG_META_SIZE; i++) {
for (int i = 0; i < tag_value->n_meta; i++) {
if (strlen(tag_value->metas[i].name) > 0) {
if (strncmp(tag_value->metas[i].name, "q", 1) == 0) {
tag->has_q = true;
Expand Down Expand Up @@ -1416,7 +1416,7 @@ int handle_trans_data(neu_plugin_t * plugin,
break;
}

for (int i = 0; i < NEU_TAG_META_SIZE; i++) {
for (int i = 0; i < tag_value->n_meta; i++) {
if (strlen(tag_value->metas[i].name) > 0) {
if (strncmp(tag_value->metas[i].name, "q", 1) == 0) {
tag->has_q = true;
Expand Down
67 changes: 43 additions & 24 deletions src/adapter/driver/cache.c
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,8 @@ struct elem {
neu_dvalue_t value;
neu_dvalue_t value_old;

neu_tag_meta_t metas[NEU_TAG_META_SIZE];
neu_tag_meta_t *metas;
int n_meta;

tkey_t key;
UT_hash_handle hh;
Expand Down Expand Up @@ -115,6 +116,13 @@ void neu_driver_cache_destroy(neu_driver_cache_t *cache)
free(elem->value.value.bools.bools);
}

if (elem->metas != NULL) {
for (int i = 0; i < elem->n_meta; i++) {
neu_free_dvalue(&elem->metas[i].value);
}
free(elem->metas);
}

free(elem);
}

Expand Down Expand Up @@ -764,9 +772,16 @@ void neu_driver_cache_update_change(neu_driver_cache_t *cache,
}
elem->value.type = value.type;

memset(elem->metas, 0, sizeof(neu_tag_meta_t) * NEU_TAG_META_SIZE);
for (int i = 0; i < n_meta; i++) {
memcpy(&elem->metas[i], &metas[i], sizeof(neu_tag_meta_t));
if (metas != NULL) {
if (elem->metas != NULL) {
free(elem->metas);
}
elem->metas = calloc(n_meta, sizeof(neu_tag_meta_t));
for (int i = 0; i < n_meta; i++) {
memcpy(&elem->metas[i], &metas[i], sizeof(neu_tag_meta_t));
}

elem->n_meta = n_meta;
}
}

Expand All @@ -784,7 +799,7 @@ void neu_driver_cache_update(neu_driver_cache_t *cache, const char *group,

int neu_driver_cache_meta_get(neu_driver_cache_t *cache, const char *group,
const char *tag, neu_driver_cache_value_t *value,
neu_tag_meta_t *metas, int n_meta)
neu_tag_meta_t **metas, int *n_meta)
{
struct elem *elem = NULL;
int ret = -1;
Expand All @@ -798,8 +813,14 @@ int neu_driver_cache_meta_get(neu_driver_cache_t *cache, const char *group,
value->value.type = elem->value.type;
value->value.precision = elem->value.precision;

assert(n_meta <= NEU_TAG_META_SIZE);
memcpy(metas, elem->metas, sizeof(neu_tag_meta_t) * NEU_TAG_META_SIZE);
// assert(n_meta <= NEU_TAG_META_SIZE);
if (elem->metas) {
*metas = calloc(elem->n_meta, sizeof(neu_tag_meta_t));
memcpy(*metas, elem->metas, sizeof(neu_tag_meta_t) * elem->n_meta);
} else {
*metas = NULL;
}
*n_meta = elem->n_meta;

switch (elem->value.type) {
case NEU_TYPE_INT8:
Expand Down Expand Up @@ -941,13 +962,6 @@ int neu_driver_cache_meta_get(neu_driver_cache_t *cache, const char *group,
break;
}

for (int i = 0; i < NEU_TAG_META_SIZE; i++) {
if (strlen(elem->metas[i].name) > 0) {
memcpy(&value->metas[i], &elem->metas[i],
sizeof(neu_tag_meta_t));
}
}

ret = 0;
}

Expand All @@ -959,7 +973,7 @@ int neu_driver_cache_meta_get(neu_driver_cache_t *cache, const char *group,
int neu_driver_cache_meta_get_changed(neu_driver_cache_t *cache,
const char *group, const char *tag,
neu_driver_cache_value_t *value,
neu_tag_meta_t *metas, int n_meta)
neu_tag_meta_t **metas, int *n_meta)
{
struct elem *elem = NULL;
int ret = -1;
Expand All @@ -973,8 +987,14 @@ int neu_driver_cache_meta_get_changed(neu_driver_cache_t *cache,
value->value.type = elem->value.type;
value->value.precision = elem->value.precision;

assert(n_meta <= NEU_TAG_META_SIZE);
memcpy(metas, elem->metas, sizeof(neu_tag_meta_t) * NEU_TAG_META_SIZE);
// assert(n_meta <= NEU_TAG_META_SIZE);
if (elem->metas) {
*metas = calloc(elem->n_meta, sizeof(neu_tag_meta_t));
memcpy(*metas, elem->metas, sizeof(neu_tag_meta_t) * elem->n_meta);
} else {
*metas = NULL;
}
*n_meta = elem->n_meta;

switch (elem->value.type) {
case NEU_TYPE_INT8:
Expand Down Expand Up @@ -1116,13 +1136,6 @@ int neu_driver_cache_meta_get_changed(neu_driver_cache_t *cache,
break;
}

for (int i = 0; i < NEU_TAG_META_SIZE; i++) {
if (strlen(elem->metas[i].name) > 0) {
memcpy(&value->metas[i], &elem->metas[i],
sizeof(neu_tag_meta_t));
}
}

if (elem->value.type != NEU_TYPE_ERROR) {
elem->changed = false;
}
Expand Down Expand Up @@ -1165,6 +1178,12 @@ void neu_driver_cache_del(neu_driver_cache_t *cache, const char *group,
elem->value.type < NEU_TYPE_ARRAY_STRING) {
free(elem->value.value.bools.bools);
}
if (elem->metas != NULL) {
for (int i = 0; i < elem->n_meta; i++) {
neu_free_dvalue(&elem->metas[i].value);
}
free(elem->metas);
}
free(elem);
}

Expand Down
9 changes: 4 additions & 5 deletions src/adapter/driver/cache.h
Original file line number Diff line number Diff line change
Expand Up @@ -50,17 +50,16 @@ void neu_driver_cache_update_trace(neu_driver_cache_t *cache, const char *group,
void *neu_driver_cache_get_trace(neu_driver_cache_t *cache, const char *group);

typedef struct {
neu_dvalue_t value;
int64_t timestamp;
neu_tag_meta_t metas[NEU_TAG_META_SIZE];
neu_dvalue_t value;
int64_t timestamp;
} neu_driver_cache_value_t;

int neu_driver_cache_meta_get(neu_driver_cache_t *cache, const char *group,
const char *tag, neu_driver_cache_value_t *value,
neu_tag_meta_t *metas, int n_meta);
neu_tag_meta_t **metas, int *n_meta);
int neu_driver_cache_meta_get_changed(neu_driver_cache_t *cache,
const char *group, const char *tag,
neu_driver_cache_value_t *value,
neu_tag_meta_t *metas, int n_meta);
neu_tag_meta_t **metas, int *n_meta);

#endif
23 changes: 15 additions & 8 deletions src/adapter/driver/driver.c
Original file line number Diff line number Diff line change
Expand Up @@ -2338,6 +2338,13 @@ static int report_callback(void *usr_data)
} else {
neu_free_dvalue(&tag_value->value);
}

if (tag_value->metas != NULL) {
for (int i = 0; i < tag_value->n_meta; i++) {
neu_free_dvalue(&tag_value->metas[i].value);
}
free(tag_value->metas);
}
}
utarray_free(data->tags);
free(data->group);
Expand Down Expand Up @@ -2542,15 +2549,15 @@ static void read_report_group(int64_t timestamp, int64_t timeout,

if (neu_tag_attribute_test(tag, NEU_ATTRIBUTE_SUBSCRIBE)) {
if (neu_driver_cache_meta_get_changed(cache, group, tag->name,
&value, tag_value.metas,
NEU_TAG_META_SIZE) != 0) {
&value, &tag_value.metas,
&tag_value.n_meta) != 0) {
nlog_debug("tag: %s not changed", tag->name);
continue;
}
} else {
if (neu_driver_cache_meta_get(cache, group, tag->name, &value,
tag_value.metas,
NEU_TAG_META_SIZE) != 0) {
&tag_value.metas,
&tag_value.n_meta) != 0) {
strcpy(tag_value.tag, tag->name);
tag_value.value.type = NEU_TYPE_ERROR;
tag_value.value.value.i32 = NEU_ERR_PLUGIN_TAG_NOT_READY;
Expand Down Expand Up @@ -2737,8 +2744,8 @@ static void read_group(int64_t timestamp, int64_t timeout,
tag_value.datatag.bias = tag->bias;

if (neu_driver_cache_meta_get(cache, group, tag->name, &value,
tag_value.metas,
NEU_TAG_META_SIZE) != 0) {
&tag_value.metas,
&tag_value.n_meta) != 0) {
tag_value.value.type = NEU_TYPE_ERROR;
tag_value.value.value.i32 = NEU_ERR_PLUGIN_TAG_NOT_READY;

Expand Down Expand Up @@ -2925,8 +2932,8 @@ static void read_group_paginate(int64_t timestamp, int64_t timeout,
}

if (neu_driver_cache_meta_get(cache, group, tag->name, &value,
tag_value.metas,
NEU_TAG_META_SIZE) != 0) {
&tag_value.metas,
&tag_value.n_meta) != 0) {
tag_value.value.type = NEU_TYPE_ERROR;
tag_value.value.value.i32 = NEU_ERR_PLUGIN_TAG_NOT_READY;

Expand Down