MCPcopy Create free account
hub / github.com/Snapchat/KeyDB / RM_StreamIteratorStart

Function RM_StreamIteratorStart

src/module.cpp:3560–3597  ·  view source on GitHub ↗

Sets up a stream iterator. * * - `key`: The stream key opened for reading using RedisModule_OpenKey(). * - `flags`: * - `REDISMODULE_STREAM_ITERATOR_EXCLUSIVE`: Don't include `start` and `end` * in the iterated range. * - `REDISMODULE_STREAM_ITERATOR_REVERSE`: Iterate in reverse order, starting * from the `end` of the range. * - `start`: The lower bound of the range. Use NULL f

Source from the content-addressed store, hash-verified

3558 * RedisModule_StreamIteratorStop(key);
3559 */
3560int RM_StreamIteratorStart(RedisModuleKey *key, int flags, RedisModuleStreamID *start, RedisModuleStreamID *end) {
3561 /* check args */
3562 if (!key ||
3563 (flags & ~(REDISMODULE_STREAM_ITERATOR_EXCLUSIVE |
3564 REDISMODULE_STREAM_ITERATOR_REVERSE))) {
3565 errno = EINVAL; /* key missing or invalid flags */
3566 return REDISMODULE_ERR;
3567 } else if (!key->value || key->value->type != OBJ_STREAM) {
3568 errno = ENOTSUP;
3569 return REDISMODULE_ERR; /* not a stream */
3570 } else if (key->iter) {
3571 errno = EBADF; /* iterator already started */
3572 return REDISMODULE_ERR;
3573 }
3574
3575 /* define range for streamIteratorStart() */
3576 streamID lower, upper;
3577 if (start) lower = {start->ms, start->seq};
3578 if (end) upper = {end->ms, end->seq};
3579 if (flags & REDISMODULE_STREAM_ITERATOR_EXCLUSIVE) {
3580 if ((start && streamIncrID(&lower) != C_OK) ||
3581 (end && streamDecrID(&upper) != C_OK)) {
3582 errno = EDOM; /* end is 0-0 or start is MAX-MAX? */
3583 return REDISMODULE_ERR;
3584 }
3585 }
3586
3587 /* create iterator */
3588 stream *s = (stream*)ptrFromObj(key->value);
3589 int rev = flags & REDISMODULE_STREAM_ITERATOR_REVERSE;
3590 streamIterator *si = (streamIterator*)zmalloc(sizeof(*si));
3591 streamIteratorStart(si, s, start ? &lower : NULL, end ? &upper : NULL, rev);
3592 key->iter = si;
3593 key->u.stream.currentid.ms = 0; /* for RM_StreamIteratorDelete() */
3594 key->u.stream.currentid.seq = 0;
3595 key->u.stream.numfieldsleft = 0; /* for RM_StreamIteratorNextField() */
3596 return REDISMODULE_OK;
3597}
3598
3599/* Stops a stream iterator created using RedisModule_StreamIteratorStart() and
3600 * reclaims its memory.

Callers

nothing calls this directly

Calls 5

streamIncrIDFunction · 0.85
streamDecrIDFunction · 0.85
ptrFromObjFunction · 0.85
zmallocFunction · 0.85
streamIteratorStartFunction · 0.85

Tested by

no test coverage detected