/* ** task_queue.c - Task::Queue implementation */ #include #include #include #include #include #include #include "task.h" typedef struct mrb_task_queue { uint8_t closed; } mrb_task_queue; static void mrb_task_queue_free(mrb_state *mrb, void *ptr) { mrb_free(mrb, ptr); } static const struct mrb_data_type mrb_task_queue_type = { "Task::Queue", mrb_task_queue_free, }; static mrb_value wait_retry_; static struct RClass *task_error_class_; /* Wake the highest-priority task waiting on this queue */ static void queue_wake_one_waiter(mrb_state *mrb, mrb_task_queue *q) { mrb_task_disable_irq(); mrb_task *curr = q_waiting_; while (curr) { mrb_task *next = curr->next; if (curr->reason == MRB_TASK_REASON_QUEUE && curr->wait.queue == q) { q_delete_task(mrb, curr); curr->status = MRB_TASK_STATUS_READY; curr->reason = MRB_TASK_REASON_NONE; curr->wait.queue = NULL; q_insert_task(mrb, curr); switching_ = TRUE; break; } curr = next; } mrb_task_enable_irq(); } /* Wake all tasks waiting on this queue (used by close) */ static void queue_wake_all_waiters(mrb_state *mrb, mrb_task_queue *q) { mrb_bool woke_any = FALSE; mrb_task_disable_irq(); mrb_task *curr = q_waiting_; while (curr) { mrb_task *next = curr->next; if (curr->reason == MRB_TASK_REASON_QUEUE && curr->wait.queue == q) { q_delete_task(mrb, curr); curr->status = MRB_TASK_STATUS_READY; curr->reason = MRB_TASK_REASON_NONE; curr->wait.queue = NULL; q_insert_task(mrb, curr); woke_any = TRUE; } curr = next; } if (woke_any) { switching_ = TRUE; } mrb_task_enable_irq(); } static mrb_value queue_initialize(mrb_state *mrb, mrb_value self) { mrb_task_queue *q = (mrb_task_queue*)mrb_malloc(mrb, sizeof(mrb_task_queue)); q->closed = 0; mrb_data_init(self, q, &mrb_task_queue_type); mrb_iv_set(mrb, self, mrb_intern_lit(mrb, "@items"), mrb_ary_new(mrb)); return self; } static mrb_value queue_push(mrb_state *mrb, mrb_value self) { mrb_value obj; mrb_get_args(mrb, "o", &obj); mrb_task_queue *q = (mrb_task_queue*)mrb_data_get_ptr(mrb, self, &mrb_task_queue_type); if (!q) mrb_raise(mrb, E_ARGUMENT_ERROR, "invalid queue"); if (q->closed) mrb_raise(mrb, task_error_class_, "queue closed"); mrb_value items = mrb_iv_get(mrb, self, mrb_intern_lit(mrb, "@items")); mrb_ary_push(mrb, items, obj); queue_wake_one_waiter(mrb, q); return self; } /* * __pop_try: try to pop one item. Returns: * - the item if available * - nil if closed and empty * - raises Task::Error if non_block and empty * - Task::Queue::WAIT_RETRY sentinel if the current task was put to WAITING * * Ruby-level pop loops on WAIT_RETRY. */ static mrb_value queue_pop_try(mrb_state *mrb, mrb_value self) { mrb_bool non_block = FALSE; mrb_get_args(mrb, "|b", &non_block); mrb_task_queue *q = (mrb_task_queue*)mrb_data_get_ptr(mrb, self, &mrb_task_queue_type); if (!q) mrb_raise(mrb, E_ARGUMENT_ERROR, "invalid queue"); mrb_value items = mrb_iv_get(mrb, self, mrb_intern_lit(mrb, "@items")); /* Item available - return it */ if (RARRAY_LEN(items) > 0) { return mrb_ary_shift(mrb, items); } /* Closed and empty */ if (q->closed) { return mrb_nil_value(); } /* Non-blocking and empty */ if (non_block) { mrb_raise(mrb, task_error_class_, "queue empty"); } /* Blocking pop only works inside a task */ if (mrb->c == mrb->root_c) { mrb_raise(mrb, E_RUNTIME_ERROR, "blocking pop can only be called from within a task"); } /* Blocking pop requires the scheduler to be running */ task_check_scheduler_lock(mrb); /* Guard against yielding from inside a C function boundary */ mrb_callinfo *ci; for (ci = mrb->c->ci; ci >= mrb->c->cibase; ci--) { if (ci->cci > 0) { mrb_raise(mrb, E_RUNTIME_ERROR, "blocking pop cannot be called from within a C function boundary"); } } /* Move current task to WAITING */ mrb_task *current = MRB2TASK(mrb); mrb_task_disable_irq(); q_delete_task(mrb, current); current->status = MRB_TASK_STATUS_WAITING; current->reason = MRB_TASK_REASON_QUEUE; current->wait.queue = q; q_insert_task(mrb, current); mrb_task_enable_irq(); switching_ = TRUE; /* Return sentinel; the Ruby pop loop will retry after wakeup */ return wait_retry_; } static mrb_value queue_size(mrb_state *mrb, mrb_value self) { mrb_value items = mrb_iv_get(mrb, self, mrb_intern_lit(mrb, "@items")); return mrb_int_value(mrb, RARRAY_LEN(items)); } static mrb_value queue_empty_p(mrb_state *mrb, mrb_value self) { mrb_value items = mrb_iv_get(mrb, self, mrb_intern_lit(mrb, "@items")); return mrb_bool_value(RARRAY_LEN(items) == 0); } static mrb_value queue_clear(mrb_state *mrb, mrb_value self) { mrb_value items = mrb_iv_get(mrb, self, mrb_intern_lit(mrb, "@items")); mrb_ary_clear(mrb, items); return self; } static mrb_value queue_close(mrb_state *mrb, mrb_value self) { mrb_task_queue *q = (mrb_task_queue*)mrb_data_get_ptr(mrb, self, &mrb_task_queue_type); if (!q) mrb_raise(mrb, E_ARGUMENT_ERROR, "invalid queue"); if (!q->closed) { q->closed = 1; queue_wake_all_waiters(mrb, q); } return self; } static mrb_value queue_closed_p(mrb_state *mrb, mrb_value self) { mrb_task_queue *q = (mrb_task_queue*)mrb_data_get_ptr(mrb, self, &mrb_task_queue_type); if (!q) mrb_raise(mrb, E_ARGUMENT_ERROR, "invalid queue"); return mrb_bool_value(q->closed); } static mrb_value queue_num_waiting(mrb_state *mrb, mrb_value self) { mrb_task_queue *q = (mrb_task_queue*)mrb_data_get_ptr(mrb, self, &mrb_task_queue_type); if (!q) mrb_raise(mrb, E_ARGUMENT_ERROR, "invalid queue"); uint32_t count = 0; mrb_task_disable_irq(); mrb_task *curr = q_waiting_; while (curr) { if (curr->reason == MRB_TASK_REASON_QUEUE && curr->wait.queue == q) { count++; } curr = curr->next; } mrb_task_enable_irq(); return mrb_int_value(mrb, (mrb_int)count); } void mrb_init_task_queue(mrb_state *mrb, struct RClass *task_class) { struct RClass *queue_class; queue_class = mrb_define_class_under_id(mrb, task_class, MRB_SYM(Queue), mrb->object_class); MRB_SET_INSTANCE_TT(queue_class, MRB_TT_DATA); task_error_class_ = mrb_class_get_under_id(mrb, task_class, MRB_SYM(Error)); /* Allocate and store WAIT_RETRY sentinel (rooted by the class constant table) */ wait_retry_ = mrb_obj_new(mrb, mrb->object_class, 0, NULL); mrb_define_const_id(mrb, queue_class, MRB_SYM(WAIT_RETRY), wait_retry_); mrb_define_method_id(mrb, queue_class, MRB_SYM(initialize), queue_initialize, MRB_ARGS_NONE()); mrb_define_method_id(mrb, queue_class, MRB_SYM(__push), queue_push, MRB_ARGS_REQ(1)); mrb_define_method_id(mrb, queue_class, MRB_SYM(__pop_try), queue_pop_try, MRB_ARGS_OPT(1)); mrb_define_method_id(mrb, queue_class, MRB_SYM(size), queue_size, MRB_ARGS_NONE()); mrb_define_method_id(mrb, queue_class, MRB_SYM(length), queue_size, MRB_ARGS_NONE()); mrb_define_method_id(mrb, queue_class, MRB_SYM_Q(empty), queue_empty_p, MRB_ARGS_NONE()); mrb_define_method_id(mrb, queue_class, MRB_SYM(clear), queue_clear, MRB_ARGS_NONE()); mrb_define_method_id(mrb, queue_class, MRB_SYM(close), queue_close, MRB_ARGS_NONE()); mrb_define_method_id(mrb, queue_class, MRB_SYM_Q(closed), queue_closed_p, MRB_ARGS_NONE()); mrb_define_method_id(mrb, queue_class, MRB_SYM(num_waiting), queue_num_waiting, MRB_ARGS_NONE()); }