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
80 changes: 80 additions & 0 deletions src/glib/glib_ssl.c
Original file line number Diff line number Diff line change
Expand Up @@ -67,3 +67,83 @@
return NATS_UPDATE_ERR_STACK(s);
}

static bool hashNoErrorOnNoSSL = false;

void
nats_hashNoErrorOnNoSSL(bool noError)
{
hashNoErrorOnNoSSL = noError;
}

natsStatus
nats_hashNew(nats_hash **new_hash)
{
#if defined(NATS_HAS_TLS)
EVP_MD_CTX *h = EVP_MD_CTX_new();
if (h == NULL)
return nats_setError(NATS_SSL_ERROR, "unable to create hash: %s", NATS_SSL_ERR_REASON_STRING);

Check warning on line 84 in src/glib/glib_ssl.c

View check run for this annotation

Codecov / codecov/patch

src/glib/glib_ssl.c#L84

Added line #L84 was not covered by tests

if (!EVP_DigestInit_ex(h, EVP_sha256(), NULL))
{
EVP_MD_CTX_free(h);
return nats_setError(NATS_SSL_ERROR, "unable to create hash: %s", NATS_SSL_ERR_REASON_STRING);

Check warning on line 89 in src/glib/glib_ssl.c

View check run for this annotation

Codecov / codecov/patch

src/glib/glib_ssl.c#L88-L89

Added lines #L88 - L89 were not covered by tests
}
*new_hash = (nats_hash*) h;
return NATS_OK;
#else
if (hashNoErrorOnNoSSL)
{
*new_hash = NULL;
return NATS_OK;
}
return nats_setError(NATS_ILLEGAL_STATE, "%s", NO_SSL_ERR);
#endif
}

natsStatus
nats_hashWrite(nats_hash *hash, const void *data, int dataLen)
{
#if defined(NATS_HAS_TLS)
if (!EVP_DigestUpdate((EVP_MD_CTX*) hash, data, (size_t) dataLen))
return nats_setError(NATS_SSL_ERROR, "error writing into hash: %s", NATS_SSL_ERR_REASON_STRING);

Check warning on line 108 in src/glib/glib_ssl.c

View check run for this annotation

Codecov / codecov/patch

src/glib/glib_ssl.c#L108

Added line #L108 was not covered by tests
return NATS_OK;
#else
if (hashNoErrorOnNoSSL)
return NATS_OK;
return nats_setError(NATS_ILLEGAL_STATE, "%s", NO_SSL_ERR);

Check warning on line 113 in src/glib/glib_ssl.c

View check run for this annotation

Codecov / codecov/patch

src/glib/glib_ssl.c#L113

Added line #L113 was not covered by tests
#endif
}

natsStatus
nats_hashSum(nats_hash *hash, unsigned char *digest, unsigned int *len)
{
#if defined(NATS_HAS_TLS)
if (!EVP_DigestFinal_ex((EVP_MD_CTX*) hash, digest, len))
return nats_setError(NATS_SSL_ERROR, "error finalizing hash: %s", NATS_SSL_ERR_REASON_STRING);

Check warning on line 122 in src/glib/glib_ssl.c

View check run for this annotation

Codecov / codecov/patch

src/glib/glib_ssl.c#L122

Added line #L122 was not covered by tests

return NATS_OK;

#else
if (hashNoErrorOnNoSSL)
{
const char *nss = "not supported";
const unsigned int slen = (unsigned int) strlen(nss);
memcpy(digest, nss, (size_t) slen);
*len = slen;
return NATS_OK;
}
return nats_setError(NATS_ILLEGAL_STATE, "%s", NO_SSL_ERR);

Check warning on line 135 in src/glib/glib_ssl.c

View check run for this annotation

Codecov / codecov/patch

src/glib/glib_ssl.c#L135

Added line #L135 was not covered by tests
#endif
}

void
nats_hashDestroy(nats_hash *hash)
{
if (hash == NULL)
return;

#if defined(NATS_HAS_TLS)
EVP_MD_CTX_free((EVP_MD_CTX*) hash);
#endif
}

18 changes: 17 additions & 1 deletion src/js.c
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,9 @@ _destroyOptions(jsOptions *o)
static void
_freeContext(jsCtx *js)
{
natsConnection *nc = NULL;
natsConnection *nc = NULL;
void *arg = NULL;
js_onReleaseCb cb = NULL;

natsStrHash_Destroy(js->pm);
natsSubscription_Destroy(js->rsub);
Expand All @@ -90,8 +92,13 @@ _freeContext(jsCtx *js)
natsMutex_Destroy(js->mu);
natsTimer_Destroy(js->pmtmr);
nc = js->nc;
cb = js->onReleaseCb;
arg = js->onReleaseCbArg;
NATS_FREE(js);

if (cb != NULL)
cb(arg);

natsConn_release(nc);
}

Expand Down Expand Up @@ -3651,3 +3658,12 @@ jsSub_checkOrderedMsg(natsSubscription *sub, natsMsg *msg, bool *reset)
}
return NATS_UPDATE_ERR_STACK(s);
}

void
js_setOnReleasedCb(jsCtx *js, js_onReleaseCb cb, void *arg)
{
js_lock(js);
js->onReleaseCb = cb;
js->onReleaseCbArg = arg;
js_unlock(js);
}
5 changes: 4 additions & 1 deletion src/js.h
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ void js_unlock(jsCtx *js);
// We know what we are doing :-)

#define js_lock(js) (natsMutex_Lock((js)->mu))
#define js_unlock(c) (natsMutex_Unlock((js)->mu))
#define js_unlock(js) (natsMutex_Unlock((js)->mu))

#endif // DEV_MODE

Expand Down Expand Up @@ -291,3 +291,6 @@ js_checkFetchedMsg(natsSubscription *sub, natsMsg *msg, uint64_t fetchID, bool c

natsStatus
js_maybeFetchMore(natsSubscription *sub, jsFetch *fetch);

void
js_setOnReleasedCb(jsCtx *js, js_onReleaseCb cb, void *arg);
106 changes: 60 additions & 46 deletions src/jsm.c
Original file line number Diff line number Diff line change
Expand Up @@ -32,28 +32,6 @@

} apiPaged;

static natsStatus
_marshalTimeUTC(natsBuffer *buf, bool sep, const char *fieldName, int64_t timeUTC)
{
natsStatus s = NATS_OK;
char dbuf[36] = {'\0'};

s = nats_EncodeTimeUTC(dbuf, sizeof(dbuf), timeUTC);
if (s != NATS_OK)
return nats_setError(NATS_ERR, "unable to encode data for field '%s' value %" PRId64, fieldName, timeUTC);

if (sep)
s = natsBuf_AppendByte(buf, ',');

IFOK(s, natsBuf_AppendByte(buf, '"'));
IFOK(s, natsBuf_Append(buf, fieldName, -1));
IFOK(s, natsBuf_Append(buf, "\":\"", -1));
IFOK(s, natsBuf_Append(buf, dbuf, -1));
IFOK(s, natsBuf_AppendByte(buf, '"'));

return NATS_UPDATE_ERR_STACK(s);
}

//
// Stream related functions
//
Expand Down Expand Up @@ -356,7 +334,7 @@
if ((s == NATS_OK) && (source->OptStartSeq > 0))
s = nats_marshalLong(buf, true, "opt_start_seq", source->OptStartSeq);
if ((s == NATS_OK) && (source->OptStartTime > 0))
IFOK(s, _marshalTimeUTC(buf, true, "opt_start_time", source->OptStartTime));
IFOK(s, nats_marshalTimeUTC(buf, true, "opt_start_time", source->OptStartTime));
if (source->FilterSubject != NULL)
{
IFOK(s, natsBuf_Append(buf, ",\"filter_subject\":\"", -1));
Expand Down Expand Up @@ -887,7 +865,7 @@
if ((s == NATS_OK) && cfg->DiscardNewPerSubject)
IFOK(s, natsBuf_Append(buf, ",\"discard_new_per_subject\":true", -1));

IFOK(s, nats_marshalMetadata(buf, true, "metadata", cfg->Metadata));
IFOK(s, nats_marshalMetadata(buf, true, "metadata", &(cfg->Metadata)));
IFOK(s, _marshalStorageCompression(cfg->Compression, buf));
IFOK(s, nats_marshalULong(buf, true, "first_seq", cfg->FirstSeq));
IFOK(s, _marshalSubjectTransformConfig(&cfg->SubjectTransform, buf));
Expand Down Expand Up @@ -1779,14 +1757,34 @@
natsStatus
js_PurgeStream(jsCtx *js, const char *stream, jsOptions *opts, jsErrCode *errCode)
{
natsStatus s = _purgeOrDelete(true, js, stream, opts, errCode);
natsStatus s = NATS_OK;
jsErrCode jsErr = 0;
jsErrCode *pErr = (errCode == NULL ? &jsErr : errCode);

s = _purgeOrDelete(true, js, stream, opts, pErr);
// On "not found", clear the error stack and return the error.
if ((s == NATS_NOT_FOUND) && ((*pErr) == JSStreamNotFoundErr))
{
nats_clearLastError();
return s;
}
return NATS_UPDATE_ERR_STACK(s);
}

natsStatus
js_DeleteStream(jsCtx *js, const char *stream, jsOptions *opts, jsErrCode *errCode)
{
natsStatus s = _purgeOrDelete(false, js, stream, opts, errCode);
natsStatus s = NATS_OK;
jsErrCode jsErr = 0;
jsErrCode *pErr = (errCode == NULL ? &jsErr : errCode);

s = _purgeOrDelete(false, js, stream, opts, pErr);
// On "not found", clear the error stack and return the error.
if ((s == NATS_NOT_FOUND) && ((*pErr) == JSStreamNotFoundErr))
{
nats_clearLastError();
return s;
}
return NATS_UPDATE_ERR_STACK(s);
}

Expand Down Expand Up @@ -2927,7 +2925,7 @@
if ((s == NATS_OK) && (cfg->OptStartSeq > 0))
s = nats_marshalLong(buf, true, "opt_start_seq", cfg->OptStartSeq);
if ((s == NATS_OK) && (cfg->OptStartTime > 0))
s = _marshalTimeUTC(buf, true, "opt_start_time", cfg->OptStartTime);
s = nats_marshalTimeUTC(buf, true, "opt_start_time", cfg->OptStartTime);
IFOK(s, _marshalAckPolicy(buf, cfg->AckPolicy));
if ((s == NATS_OK) && (cfg->AckWait > 0))
s = nats_marshalLong(buf, true, "ack_wait", cfg->AckWait);
Expand All @@ -2941,9 +2939,9 @@
}
if ((s == NATS_OK) && (cfg->FilterSubjectsLen > 0))
s = nats_marshalStringArray(buf, true, "filter_subjects", cfg->FilterSubjects, cfg->FilterSubjectsLen);
IFOK(s, nats_marshalMetadata(buf, true, "metadata", cfg->Metadata));
IFOK(s, nats_marshalMetadata(buf, true, "metadata", &(cfg->Metadata)));
if ((s == NATS_OK) && (cfg->PauseUntil > 0))
s = _marshalTimeUTC(buf, true, "pause_until", cfg->PauseUntil);
s = nats_marshalTimeUTC(buf, true, "pause_until", cfg->PauseUntil);
if ((s == NATS_OK) && !nats_IsStringEmpty(cfg->PriorityPolicy))
{
s = natsBuf_Append(buf, ",\"priority_policy\":\"", -1);
Expand Down Expand Up @@ -3390,6 +3388,8 @@
bool freePfx = false;
natsConnection *nc = NULL;
natsMsg *resp = NULL;
jsErrCode jsErr = 0;
jsErrCode *pErr = (errCode == NULL ? &jsErr : errCode);
jsOptions o;

if (errCode != NULL)
Expand Down Expand Up @@ -3420,17 +3420,16 @@
IFOK_JSR(s, natsConnection_Request(&resp, nc, subj, NULL, 0, o.Wait));

// If we got a response, check for error or return the consumer info result.
IFOK(s, _unmarshalConsumerCreateOrGetResp(new_ci, resp, errCode));
IFOK(s, _unmarshalConsumerCreateOrGetResp(new_ci, resp, pErr));

Check warning on line 3423 in src/jsm.c

View check run for this annotation

Codecov / codecov/patch

src/jsm.c#L3423

Added line #L3423 was not covered by tests

NATS_FREE(subj);
natsMsg_Destroy(resp);

if (s == NATS_NOT_FOUND)
if ((s == NATS_NOT_FOUND) && ((*pErr) == JSConsumerNotFoundErr))
{
nats_clearLastError();
return s;
}

return NATS_UPDATE_ERR_STACK(s);
}

Expand All @@ -3444,6 +3443,8 @@
natsConnection *nc = NULL;
natsMsg *resp = NULL;
bool success = false;
jsErrCode jsErr = 0;
jsErrCode *pErr = (errCode == NULL ? &jsErr : errCode);
jsOptions o;

if (errCode != NULL)
Expand Down Expand Up @@ -3474,13 +3475,18 @@
IFOK_JSR(s, natsConnection_Request(&resp, nc, subj, NULL, 0, o.Wait));

// If we got a response, check for error and success result.
IFOK(s, _unmarshalSuccessResp(&success, resp, errCode));
IFOK(s, _unmarshalSuccessResp(&success, resp, pErr));

Check warning on line 3478 in src/jsm.c

View check run for this annotation

Codecov / codecov/patch

src/jsm.c#L3478

Added line #L3478 was not covered by tests
if ((s == NATS_OK) && !success)
s = nats_setError(s, "failed to delete consumer '%s'", consumer);

NATS_FREE(subj);
natsMsg_Destroy(resp);

if ((s == NATS_NOT_FOUND) && ((*pErr) == JSConsumerNotFoundErr))
{
nats_clearLastError();
return s;
}
return NATS_UPDATE_ERR_STACK(s);
}

Expand Down Expand Up @@ -3548,7 +3554,7 @@
s = natsBuf_Create(&buf, 256);
IFOK(s, natsBuf_AppendByte(buf, '{'));
if ((s == NATS_OK) && (pauseUntil > 0)) {
s = _marshalTimeUTC(buf, false, "pause_until", pauseUntil);
s = nats_marshalTimeUTC(buf, false, "pause_until", pauseUntil);
}
IFOK(s, natsBuf_AppendByte(buf, '}'));

Expand All @@ -3575,13 +3581,15 @@
const char *stream, const char *consumer,
uint64_t pauseUntil, jsOptions *opts, jsErrCode *errCode)
{
natsStatus s = NATS_OK;
char *subj = NULL;
bool freePfx = false;
natsConnection *nc = NULL;
natsBuffer *buf = NULL;
natsMsg *resp = NULL;
jsOptions o;
natsStatus s = NATS_OK;
char *subj = NULL;
bool freePfx = false;
natsConnection *nc = NULL;
natsBuffer *buf = NULL;
natsMsg *resp = NULL;
jsErrCode jsErr = 0;
jsErrCode *pErr = (errCode == NULL ? &jsErr : errCode);
jsOptions o;

if (errCode != NULL)
*errCode = 0;
Expand Down Expand Up @@ -3613,18 +3621,17 @@
IFOK_JSR(s, natsConnection_Request(&resp, nc, subj, natsBuf_Data(buf), natsBuf_Len(buf), o.Wait));

// If we got a response, check for error or return the consumer info result.
IFOK(s, _unmarshalConsumerPauseResp(new_cpr, resp, errCode));
IFOK(s, _unmarshalConsumerPauseResp(new_cpr, resp, pErr));

Check warning on line 3624 in src/jsm.c

View check run for this annotation

Codecov / codecov/patch

src/jsm.c#L3624

Added line #L3624 was not covered by tests

NATS_FREE(subj);
natsMsg_Destroy(resp);
natsBuf_Destroy(buf);

if (s == NATS_NOT_FOUND)
if ((s == NATS_NOT_FOUND) && ((*pErr) == JSConsumerNotFoundErr))
{
nats_clearLastError();
return s;
}

return NATS_UPDATE_ERR_STACK(s);
}

Expand All @@ -3638,6 +3645,8 @@
natsConnection *nc = NULL;
natsMsg *resp = NULL;
bool success = false;
jsErrCode jsErr = 0;
jsErrCode *pErr = (errCode == NULL ? &jsErr : errCode);
jsOptions o;
char jsonBuf[64];

Expand Down Expand Up @@ -3675,13 +3684,18 @@
IFOK_JSR(s, natsConnection_RequestString(&resp, nc, subj, jsonBuf, o.Wait));

// If we got a response, check for error and success result.
IFOK(s, _unmarshalSuccessResp(&success, resp, errCode));
IFOK(s, _unmarshalSuccessResp(&success, resp, pErr));

Check warning on line 3687 in src/jsm.c

View check run for this annotation

Codecov / codecov/patch

src/jsm.c#L3687

Added line #L3687 was not covered by tests
if ((s == NATS_OK) && !success)
s = nats_setError(s, "failed to unpin group '%s' at consumer '%s'", group, consumer);

NATS_FREE(subj);
natsMsg_Destroy(resp);

if ((s == NATS_NOT_FOUND) && ((*pErr) == JSConsumerNotFoundErr))
{
nats_clearLastError();
return s;

Check warning on line 3697 in src/jsm.c

View check run for this annotation

Codecov / codecov/patch

src/jsm.c#L3696-L3697

Added lines #L3696 - L3697 were not covered by tests
}
return NATS_UPDATE_ERR_STACK(s);
}

Expand Down Expand Up @@ -4121,7 +4135,7 @@
c->FilterSubjectsLen++;
}
}
IFOK(s, nats_cloneMetadata(&(c->Metadata), org->Metadata));
IFOK(s, nats_cloneMetadata(&(c->Metadata), &(org->Metadata)));
if (s == NATS_OK)
*clone = c;
else
Expand Down
Loading