класс SizedQueue
Этот класс представляет очереди заданной ёмкости. Операция добавления может быть заблокирована, если ёмкость заполнена.
См. Queue, чтобы увидеть пример того, как работает SizedQueue.
Публичные методы класса
static VALUE
rb_szqueue_initialize(VALUE self, VALUE vmax)
{
long max;
struct rb_szqueue *sq = szqueue_ptr(self);
max = NUM2LONG(vmax);
if (max <= 0) {
rb_raise(rb_eArgError, "queue size must be positive");
}
RB_OBJ_WRITE(self, &sq->q.que, ary_buf_new());
list_head_init(szqueue_waitq(sq));
list_head_init(szqueue_pushq(sq));
sq->max = max;
return self;
} Создаёт очередь фиксированной длины максимальной ёмкостью max.
Публичные методы экземпляра
Добавляет object в очередь.
Если в очереди нет места, ожидается освобождение места, если non_block не равно true. Если non_block равно true, поток не приостанавливается и поднимается исключение ThreadError.
static VALUE
rb_szqueue_clear(VALUE self)
{
struct rb_szqueue *sq = szqueue_ptr(self);
rb_ary_clear(check_array(self, sq->q.que));
wakeup_all(szqueue_pushq(sq));
return self;
} Удаляет все объекты из очереди.
static VALUE
rb_szqueue_close(VALUE self)
{
if (!queue_closed_p(self)) {
struct rb_szqueue *sq = szqueue_ptr(self);
FL_SET(self, QUEUE_CLOSED);
wakeup_all(szqueue_waitq(sq));
wakeup_all(szqueue_pushq(sq));
}
return self;
} Аналогично Queue#close.
Разница заключается в поведении с потоками, ожидающими добавления в очередь.
Если есть ожидающие потоки добавления, они прерываются путём поднятия исключения ClosedQueueError('очередь закрыта').
Извлекает данные из очереди.
Если очередь пуста, вызывающий поток приостанавливается до тех пор, пока данные не будут добавлены в очередь. Если non_block равно true, поток не приостанавливается, и поднимается исключение ThreadError.
static VALUE
rb_szqueue_empty_p(VALUE self)
{
struct rb_szqueue *sq = szqueue_ptr(self);
return queue_length(self, &sq->q) == 0 ? Qtrue : Qfalse;
} Возвращает true если очередь пуста.
Добавляет object в очередь.
Если в очереди нет места, ожидается освобождение места, если non_block не равно true. Если non_block равно true, поток не приостанавливается, и поднимается исключение ThreadError.
static VALUE
rb_szqueue_length(VALUE self)
{
struct rb_szqueue *sq = szqueue_ptr(self);
return LONG2NUM(queue_length(self, &sq->q));
} Возвращает длину очереди.
static VALUE
rb_szqueue_max_get(VALUE self)
{
return LONG2NUM(szqueue_ptr(self)->max);
} Возвращает максимальную ёмкость очереди.
static VALUE
rb_szqueue_max_set(VALUE self, VALUE vmax)
{
long max = NUM2LONG(vmax);
long diff = 0;
struct rb_szqueue *sq = szqueue_ptr(self);
if (max <= 0) {
rb_raise(rb_eArgError, "queue size must be positive");
}
if (max > sq->max) {
diff = max - sq->max;
}
sq->max = max;
sync_wakeup(szqueue_pushq(sq), diff);
return vmax;
} Устанавливает максимальную ёмкость очереди на заданное number.
static VALUE
rb_szqueue_num_waiting(VALUE self)
{
struct rb_szqueue *sq = szqueue_ptr(self);
return INT2NUM(sq->q.num_waiting + sq->num_waiting_push);
} Возвращает количество потоков, ожидающих в очереди.
static VALUE
rb_szqueue_pop(int argc, VALUE *argv, VALUE self)
{
int should_block = queue_pop_should_block(argc, argv);
return szqueue_do_pop(self, should_block);
} Извлекает данные из очереди.
Если очередь пуста, вызывающий поток приостанавливается до тех пор, пока данные не будут добавлены в очередь. Если non_block равно true, поток не приостанавливается, и поднимается исключение ThreadError.
static VALUE
rb_szqueue_push(int argc, VALUE *argv, VALUE self)
{
struct rb_szqueue *sq = szqueue_ptr(self);
int should_block = szqueue_push_should_block(argc, argv);
while (queue_length(self, &sq->q) >= sq->max) {
if (!should_block) {
rb_raise(rb_eThreadError, "queue full");
}
else if (queue_closed_p(self)) {
break;
}
else {
rb_execution_context_t *ec = GET_EC();
COROUTINE_STACK_LOCAL(struct queue_waiter, qw);
struct list_head *pushq = szqueue_pushq(sq);
qw->w.self = self;
qw->w.th = ec->thread_ptr;
qw->w.fiber = ec->fiber_ptr;
qw->as.sq = sq;
list_add_tail(pushq, &qw->w.node);
sq->num_waiting_push++;
rb_ensure(queue_sleep, self, szqueue_sleep_done, (VALUE)qw);
}
}
if (queue_closed_p(self)) {
raise_closed_queue_error(self);
}
return queue_do_push(self, &sq->q, argv[0]);
} Добавляет object в очередь.
Если в очереди нет места, ожидается освобождение места, если non_block не равно true. Если non_block равно true, поток не приостанавливается, и поднимается исключение ThreadError.
Извлекает данные из очереди.
Если очередь пуста, вызывающий поток приостанавливается до тех пор, пока данные не будут добавлены в очередь. Если non_block равно true, поток не приостанавливается, и поднимается исключение ThreadError.
Ruby Core © 1993–2020 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.