channel.c 43 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196
  1. /**
  2. * Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
  3. * SPDX-License-Identifier: Apache-2.0.
  4. */
  5. #include <aws/io/channel.h>
  6. #include <aws/common/atomics.h>
  7. #include <aws/common/clock.h>
  8. #include <aws/common/mutex.h>
  9. #include <aws/io/event_loop.h>
  10. #include <aws/io/logging.h>
  11. #include <aws/io/message_pool.h>
  12. #include <aws/io/statistics.h>
  13. #ifdef _MSC_VER
  14. # pragma warning(disable : 4204) /* non-constant aggregate initializer */
  15. #endif
  16. static size_t s_message_pool_key = 0; /* Address of variable serves as key in hash table */
  17. enum {
  18. KB_16 = 16 * 1024,
  19. };
  20. size_t g_aws_channel_max_fragment_size = KB_16;
  21. #define INITIAL_STATISTIC_LIST_SIZE 5
  22. enum aws_channel_state {
  23. AWS_CHANNEL_SETTING_UP,
  24. AWS_CHANNEL_ACTIVE,
  25. AWS_CHANNEL_SHUTTING_DOWN,
  26. AWS_CHANNEL_SHUT_DOWN,
  27. };
  28. struct aws_shutdown_notification_task {
  29. struct aws_task task;
  30. int error_code;
  31. struct aws_channel_slot *slot;
  32. bool shutdown_immediately;
  33. };
  34. struct shutdown_task {
  35. struct aws_channel_task task;
  36. struct aws_channel *channel;
  37. int error_code;
  38. bool shutdown_immediately;
  39. };
  40. struct aws_channel {
  41. struct aws_allocator *alloc;
  42. struct aws_event_loop *loop;
  43. struct aws_channel_slot *first;
  44. struct aws_message_pool *msg_pool;
  45. enum aws_channel_state channel_state;
  46. struct aws_shutdown_notification_task shutdown_notify_task;
  47. aws_channel_on_shutdown_completed_fn *on_shutdown_completed;
  48. void *shutdown_user_data;
  49. struct aws_atomic_var refcount;
  50. struct aws_task deletion_task;
  51. struct aws_task statistics_task;
  52. struct aws_crt_statistics_handler *statistics_handler;
  53. uint64_t statistics_interval_start_time_ms;
  54. struct aws_array_list statistic_list;
  55. struct {
  56. struct aws_linked_list list;
  57. } channel_thread_tasks;
  58. struct {
  59. struct aws_mutex lock;
  60. struct aws_linked_list list;
  61. struct aws_task scheduling_task;
  62. struct shutdown_task shutdown_task;
  63. bool is_channel_shut_down;
  64. } cross_thread_tasks;
  65. size_t window_update_batch_emit_threshold;
  66. struct aws_channel_task window_update_task;
  67. bool read_back_pressure_enabled;
  68. bool window_update_scheduled;
  69. };
  70. struct channel_setup_args {
  71. struct aws_allocator *alloc;
  72. struct aws_channel *channel;
  73. aws_channel_on_setup_completed_fn *on_setup_completed;
  74. void *user_data;
  75. struct aws_task task;
  76. };
  77. static void s_on_msg_pool_removed(struct aws_event_loop_local_object *object) {
  78. struct aws_message_pool *msg_pool = object->object;
  79. AWS_LOGF_TRACE(
  80. AWS_LS_IO_CHANNEL,
  81. "static: message pool %p has been purged "
  82. "from the event-loop: likely because of shutdown",
  83. (void *)msg_pool);
  84. struct aws_allocator *alloc = msg_pool->alloc;
  85. aws_message_pool_clean_up(msg_pool);
  86. aws_mem_release(alloc, msg_pool);
  87. aws_mem_release(alloc, object);
  88. }
  89. static void s_on_channel_setup_complete(struct aws_task *task, void *arg, enum aws_task_status task_status) {
  90. (void)task;
  91. struct channel_setup_args *setup_args = arg;
  92. struct aws_message_pool *message_pool = NULL;
  93. struct aws_event_loop_local_object *local_object = NULL;
  94. AWS_LOGF_DEBUG(AWS_LS_IO_CHANNEL, "id=%p: setup complete, notifying caller.", (void *)setup_args->channel);
  95. if (task_status == AWS_TASK_STATUS_RUN_READY) {
  96. struct aws_event_loop_local_object stack_obj;
  97. AWS_ZERO_STRUCT(stack_obj);
  98. local_object = &stack_obj;
  99. if (aws_event_loop_fetch_local_object(setup_args->channel->loop, &s_message_pool_key, local_object)) {
  100. local_object = aws_mem_calloc(setup_args->alloc, 1, sizeof(struct aws_event_loop_local_object));
  101. if (!local_object) {
  102. goto cleanup_setup_args;
  103. }
  104. message_pool = aws_mem_acquire(setup_args->alloc, sizeof(struct aws_message_pool));
  105. if (!message_pool) {
  106. goto cleanup_local_obj;
  107. }
  108. AWS_LOGF_DEBUG(
  109. AWS_LS_IO_CHANNEL,
  110. "id=%p: no message pool is currently stored in the event-loop "
  111. "local storage, adding %p with max message size %zu, "
  112. "message count 4, with 4 small blocks of 128 bytes.",
  113. (void *)setup_args->channel,
  114. (void *)message_pool,
  115. g_aws_channel_max_fragment_size);
  116. struct aws_message_pool_creation_args creation_args = {
  117. .application_data_msg_data_size = g_aws_channel_max_fragment_size,
  118. .application_data_msg_count = 4,
  119. .small_block_msg_count = 4,
  120. .small_block_msg_data_size = 128,
  121. };
  122. if (aws_message_pool_init(message_pool, setup_args->alloc, &creation_args)) {
  123. goto cleanup_msg_pool_mem;
  124. }
  125. local_object->key = &s_message_pool_key;
  126. local_object->object = message_pool;
  127. local_object->on_object_removed = s_on_msg_pool_removed;
  128. if (aws_event_loop_put_local_object(setup_args->channel->loop, local_object)) {
  129. goto cleanup_msg_pool;
  130. }
  131. } else {
  132. message_pool = local_object->object;
  133. AWS_LOGF_DEBUG(
  134. AWS_LS_IO_CHANNEL,
  135. "id=%p: message pool %p found in event-loop local storage: using it.",
  136. (void *)setup_args->channel,
  137. (void *)message_pool);
  138. }
  139. setup_args->channel->msg_pool = message_pool;
  140. setup_args->channel->channel_state = AWS_CHANNEL_ACTIVE;
  141. setup_args->on_setup_completed(setup_args->channel, AWS_OP_SUCCESS, setup_args->user_data);
  142. aws_channel_release_hold(setup_args->channel);
  143. aws_mem_release(setup_args->alloc, setup_args);
  144. return;
  145. }
  146. goto cleanup_setup_args;
  147. cleanup_msg_pool:
  148. aws_message_pool_clean_up(message_pool);
  149. cleanup_msg_pool_mem:
  150. aws_mem_release(setup_args->alloc, message_pool);
  151. cleanup_local_obj:
  152. aws_mem_release(setup_args->alloc, local_object);
  153. cleanup_setup_args:
  154. setup_args->on_setup_completed(setup_args->channel, AWS_OP_ERR, setup_args->user_data);
  155. aws_channel_release_hold(setup_args->channel);
  156. aws_mem_release(setup_args->alloc, setup_args);
  157. }
  158. static void s_schedule_cross_thread_tasks(struct aws_task *task, void *arg, enum aws_task_status status);
  159. static void s_destroy_partially_constructed_channel(struct aws_channel *channel) {
  160. if (channel == NULL) {
  161. return;
  162. }
  163. aws_array_list_clean_up(&channel->statistic_list);
  164. aws_mem_release(channel->alloc, channel);
  165. }
  166. struct aws_channel *aws_channel_new(struct aws_allocator *alloc, const struct aws_channel_options *creation_args) {
  167. AWS_PRECONDITION(creation_args);
  168. AWS_PRECONDITION(creation_args->event_loop);
  169. AWS_PRECONDITION(creation_args->on_setup_completed);
  170. struct aws_channel *channel = aws_mem_calloc(alloc, 1, sizeof(struct aws_channel));
  171. if (!channel) {
  172. return NULL;
  173. }
  174. AWS_LOGF_DEBUG(AWS_LS_IO_CHANNEL, "id=%p: Beginning creation and setup of new channel.", (void *)channel);
  175. channel->alloc = alloc;
  176. channel->loop = creation_args->event_loop;
  177. channel->on_shutdown_completed = creation_args->on_shutdown_completed;
  178. channel->shutdown_user_data = creation_args->shutdown_user_data;
  179. if (aws_array_list_init_dynamic(
  180. &channel->statistic_list, alloc, INITIAL_STATISTIC_LIST_SIZE, sizeof(struct aws_crt_statistics_base *))) {
  181. goto on_error;
  182. }
  183. /* Start refcount at 2:
  184. * 1 for self-reference, released from aws_channel_destroy()
  185. * 1 for the setup task, released when task executes */
  186. aws_atomic_init_int(&channel->refcount, 2);
  187. struct channel_setup_args *setup_args = aws_mem_calloc(alloc, 1, sizeof(struct channel_setup_args));
  188. if (!setup_args) {
  189. goto on_error;
  190. }
  191. channel->channel_state = AWS_CHANNEL_SETTING_UP;
  192. aws_linked_list_init(&channel->channel_thread_tasks.list);
  193. aws_linked_list_init(&channel->cross_thread_tasks.list);
  194. channel->cross_thread_tasks.lock = (struct aws_mutex)AWS_MUTEX_INIT;
  195. if (creation_args->enable_read_back_pressure) {
  196. channel->read_back_pressure_enabled = true;
  197. /* we probably only need room for one fragment, but let's avoid potential deadlocks
  198. * on things like tls that need extra head-room. */
  199. channel->window_update_batch_emit_threshold = g_aws_channel_max_fragment_size * 2;
  200. }
  201. aws_task_init(
  202. &channel->cross_thread_tasks.scheduling_task,
  203. s_schedule_cross_thread_tasks,
  204. channel,
  205. "schedule_cross_thread_tasks");
  206. setup_args->alloc = alloc;
  207. setup_args->channel = channel;
  208. setup_args->on_setup_completed = creation_args->on_setup_completed;
  209. setup_args->user_data = creation_args->setup_user_data;
  210. aws_task_init(&setup_args->task, s_on_channel_setup_complete, setup_args, "on_channel_setup_complete");
  211. aws_event_loop_schedule_task_now(creation_args->event_loop, &setup_args->task);
  212. return channel;
  213. on_error:
  214. s_destroy_partially_constructed_channel(channel);
  215. return NULL;
  216. }
  217. static void s_cleanup_slot(struct aws_channel_slot *slot) {
  218. if (slot) {
  219. if (slot->handler) {
  220. aws_channel_handler_destroy(slot->handler);
  221. }
  222. aws_mem_release(slot->alloc, slot);
  223. }
  224. }
  225. void aws_channel_destroy(struct aws_channel *channel) {
  226. AWS_LOGF_DEBUG(AWS_LS_IO_CHANNEL, "id=%p: destroying channel.", (void *)channel);
  227. aws_channel_release_hold(channel);
  228. }
  229. static void s_final_channel_deletion_task(struct aws_task *task, void *arg, enum aws_task_status status) {
  230. (void)task;
  231. (void)status;
  232. struct aws_channel *channel = arg;
  233. struct aws_channel_slot *current = channel->first;
  234. if (!current || !current->handler) {
  235. /* Allow channels with no valid slots to skip shutdown process */
  236. channel->channel_state = AWS_CHANNEL_SHUT_DOWN;
  237. }
  238. AWS_ASSERT(channel->channel_state == AWS_CHANNEL_SHUT_DOWN);
  239. while (current) {
  240. struct aws_channel_slot *tmp = current->adj_right;
  241. s_cleanup_slot(current);
  242. current = tmp;
  243. }
  244. aws_array_list_clean_up(&channel->statistic_list);
  245. aws_channel_set_statistics_handler(channel, NULL);
  246. aws_mem_release(channel->alloc, channel);
  247. }
  248. void aws_channel_acquire_hold(struct aws_channel *channel) {
  249. size_t prev_refcount = aws_atomic_fetch_add(&channel->refcount, 1);
  250. AWS_ASSERT(prev_refcount != 0);
  251. (void)prev_refcount;
  252. }
  253. void aws_channel_release_hold(struct aws_channel *channel) {
  254. size_t prev_refcount = aws_atomic_fetch_sub(&channel->refcount, 1);
  255. AWS_ASSERT(prev_refcount != 0);
  256. if (prev_refcount == 1) {
  257. /* Refcount is now 0, finish cleaning up channel memory. */
  258. if (aws_channel_thread_is_callers_thread(channel)) {
  259. s_final_channel_deletion_task(NULL, channel, AWS_TASK_STATUS_RUN_READY);
  260. } else {
  261. aws_task_init(&channel->deletion_task, s_final_channel_deletion_task, channel, "final_channel_deletion");
  262. aws_event_loop_schedule_task_now(channel->loop, &channel->deletion_task);
  263. }
  264. }
  265. }
  266. struct channel_shutdown_task_args {
  267. struct aws_channel *channel;
  268. struct aws_allocator *alloc;
  269. int error_code;
  270. struct aws_task task;
  271. };
  272. static int s_channel_shutdown(struct aws_channel *channel, int error_code, bool shutdown_immediately);
  273. static void s_on_shutdown_completion_task(struct aws_task *task, void *arg, enum aws_task_status status);
  274. static void s_shutdown_task(struct aws_channel_task *task, void *arg, enum aws_task_status status) {
  275. (void)task;
  276. (void)status;
  277. struct shutdown_task *shutdown_task = arg;
  278. struct aws_channel *channel = shutdown_task->channel;
  279. int error_code = shutdown_task->error_code;
  280. bool shutdown_immediately = shutdown_task->shutdown_immediately;
  281. if (channel->channel_state < AWS_CHANNEL_SHUTTING_DOWN) {
  282. AWS_LOGF_DEBUG(AWS_LS_IO_CHANNEL, "id=%p: beginning shutdown process", (void *)channel);
  283. struct aws_channel_slot *slot = channel->first;
  284. channel->channel_state = AWS_CHANNEL_SHUTTING_DOWN;
  285. if (slot) {
  286. AWS_LOGF_TRACE(
  287. AWS_LS_IO_CHANNEL,
  288. "id=%p: shutting down slot %p (the first one) in the read direction",
  289. (void *)channel,
  290. (void *)slot);
  291. aws_channel_slot_shutdown(slot, AWS_CHANNEL_DIR_READ, error_code, shutdown_immediately);
  292. return;
  293. }
  294. channel->channel_state = AWS_CHANNEL_SHUT_DOWN;
  295. AWS_LOGF_TRACE(AWS_LS_IO_CHANNEL, "id=%p: shutdown completed", (void *)channel);
  296. aws_mutex_lock(&channel->cross_thread_tasks.lock);
  297. channel->cross_thread_tasks.is_channel_shut_down = true;
  298. aws_mutex_unlock(&channel->cross_thread_tasks.lock);
  299. if (channel->on_shutdown_completed) {
  300. channel->shutdown_notify_task.task.fn = s_on_shutdown_completion_task;
  301. channel->shutdown_notify_task.task.arg = channel;
  302. channel->shutdown_notify_task.error_code = error_code;
  303. aws_event_loop_schedule_task_now(channel->loop, &channel->shutdown_notify_task.task);
  304. }
  305. }
  306. }
  307. static int s_channel_shutdown(struct aws_channel *channel, int error_code, bool shutdown_immediately) {
  308. bool need_to_schedule = true;
  309. aws_mutex_lock(&channel->cross_thread_tasks.lock);
  310. if (channel->cross_thread_tasks.shutdown_task.task.task_fn) {
  311. need_to_schedule = false;
  312. AWS_LOGF_DEBUG(
  313. AWS_LS_IO_CHANNEL, "id=%p: Channel shutdown is already pending, not scheduling another.", (void *)channel);
  314. } else {
  315. aws_channel_task_init(
  316. &channel->cross_thread_tasks.shutdown_task.task,
  317. s_shutdown_task,
  318. &channel->cross_thread_tasks.shutdown_task,
  319. "channel_shutdown");
  320. channel->cross_thread_tasks.shutdown_task.shutdown_immediately = shutdown_immediately;
  321. channel->cross_thread_tasks.shutdown_task.channel = channel;
  322. channel->cross_thread_tasks.shutdown_task.error_code = error_code;
  323. }
  324. aws_mutex_unlock(&channel->cross_thread_tasks.lock);
  325. if (need_to_schedule) {
  326. AWS_LOGF_TRACE(AWS_LS_IO_CHANNEL, "id=%p: channel shutdown task is scheduled", (void *)channel);
  327. aws_channel_schedule_task_now(channel, &channel->cross_thread_tasks.shutdown_task.task);
  328. }
  329. return AWS_OP_SUCCESS;
  330. }
  331. int aws_channel_shutdown(struct aws_channel *channel, int error_code) {
  332. return s_channel_shutdown(channel, error_code, false);
  333. }
  334. struct aws_io_message *aws_channel_acquire_message_from_pool(
  335. struct aws_channel *channel,
  336. enum aws_io_message_type message_type,
  337. size_t size_hint) {
  338. struct aws_io_message *message = aws_message_pool_acquire(channel->msg_pool, message_type, size_hint);
  339. if (AWS_LIKELY(message)) {
  340. message->owning_channel = channel;
  341. AWS_LOGF_TRACE(
  342. AWS_LS_IO_CHANNEL,
  343. "id=%p: acquired message %p of capacity %zu from pool %p. Requested size was %zu",
  344. (void *)channel,
  345. (void *)message,
  346. message->message_data.capacity,
  347. (void *)channel->msg_pool,
  348. size_hint);
  349. }
  350. return message;
  351. }
  352. struct aws_channel_slot *aws_channel_slot_new(struct aws_channel *channel) {
  353. struct aws_channel_slot *new_slot = aws_mem_calloc(channel->alloc, 1, sizeof(struct aws_channel_slot));
  354. if (!new_slot) {
  355. return NULL;
  356. }
  357. AWS_LOGF_TRACE(AWS_LS_IO_CHANNEL, "id=%p: creating new slot %p.", (void *)channel, (void *)new_slot);
  358. new_slot->alloc = channel->alloc;
  359. new_slot->channel = channel;
  360. if (!channel->first) {
  361. channel->first = new_slot;
  362. }
  363. return new_slot;
  364. }
  365. int aws_channel_current_clock_time(struct aws_channel *channel, uint64_t *time_nanos) {
  366. return aws_event_loop_current_clock_time(channel->loop, time_nanos);
  367. }
  368. int aws_channel_fetch_local_object(
  369. struct aws_channel *channel,
  370. const void *key,
  371. struct aws_event_loop_local_object *obj) {
  372. return aws_event_loop_fetch_local_object(channel->loop, (void *)key, obj);
  373. }
  374. int aws_channel_put_local_object(
  375. struct aws_channel *channel,
  376. const void *key,
  377. const struct aws_event_loop_local_object *obj) {
  378. (void)key;
  379. return aws_event_loop_put_local_object(channel->loop, (struct aws_event_loop_local_object *)obj);
  380. }
  381. int aws_channel_remove_local_object(
  382. struct aws_channel *channel,
  383. const void *key,
  384. struct aws_event_loop_local_object *removed_obj) {
  385. return aws_event_loop_remove_local_object(channel->loop, (void *)key, removed_obj);
  386. }
  387. static void s_channel_task_run(struct aws_task *task, void *arg, enum aws_task_status status) {
  388. struct aws_channel_task *channel_task = AWS_CONTAINER_OF(task, struct aws_channel_task, wrapper_task);
  389. struct aws_channel *channel = arg;
  390. /* Any task that runs after shutdown completes is considered canceled */
  391. if (channel->channel_state == AWS_CHANNEL_SHUT_DOWN) {
  392. status = AWS_TASK_STATUS_CANCELED;
  393. }
  394. aws_linked_list_remove(&channel_task->node);
  395. channel_task->task_fn(channel_task, channel_task->arg, status);
  396. }
  397. static void s_schedule_cross_thread_tasks(struct aws_task *task, void *arg, enum aws_task_status status) {
  398. (void)task;
  399. struct aws_channel *channel = arg;
  400. struct aws_linked_list cross_thread_task_list;
  401. aws_linked_list_init(&cross_thread_task_list);
  402. /* Grab contents of cross-thread task list while we have the lock */
  403. aws_mutex_lock(&channel->cross_thread_tasks.lock);
  404. aws_linked_list_swap_contents(&channel->cross_thread_tasks.list, &cross_thread_task_list);
  405. aws_mutex_unlock(&channel->cross_thread_tasks.lock);
  406. /* If the channel has shut down since the cross-thread tasks were scheduled, run tasks immediately as canceled */
  407. if (channel->channel_state == AWS_CHANNEL_SHUT_DOWN) {
  408. status = AWS_TASK_STATUS_CANCELED;
  409. }
  410. while (!aws_linked_list_empty(&cross_thread_task_list)) {
  411. struct aws_linked_list_node *node = aws_linked_list_pop_front(&cross_thread_task_list);
  412. struct aws_channel_task *channel_task = AWS_CONTAINER_OF(node, struct aws_channel_task, node);
  413. if ((channel_task->wrapper_task.timestamp == 0) || (status == AWS_TASK_STATUS_CANCELED)) {
  414. /* Run "now" tasks, and canceled tasks, immediately */
  415. channel_task->task_fn(channel_task, channel_task->arg, status);
  416. } else {
  417. /* "Future" tasks are scheduled with the event-loop. */
  418. aws_linked_list_push_back(&channel->channel_thread_tasks.list, &channel_task->node);
  419. aws_event_loop_schedule_task_future(
  420. channel->loop, &channel_task->wrapper_task, channel_task->wrapper_task.timestamp);
  421. }
  422. }
  423. }
  424. void aws_channel_task_init(
  425. struct aws_channel_task *channel_task,
  426. aws_channel_task_fn *task_fn,
  427. void *arg,
  428. const char *type_tag) {
  429. AWS_ZERO_STRUCT(*channel_task);
  430. channel_task->task_fn = task_fn;
  431. channel_task->arg = arg;
  432. channel_task->type_tag = type_tag;
  433. }
  434. static void s_register_pending_task_in_event_loop(
  435. struct aws_channel *channel,
  436. struct aws_channel_task *channel_task,
  437. uint64_t run_at_nanos) {
  438. AWS_LOGF_TRACE(
  439. AWS_LS_IO_CHANNEL,
  440. "id=%p: scheduling task with wrapper task id %p.",
  441. (void *)channel,
  442. (void *)&channel_task->wrapper_task);
  443. /* If channel is shut down, run task immediately as canceled */
  444. if (channel->channel_state == AWS_CHANNEL_SHUT_DOWN) {
  445. AWS_LOGF_DEBUG(
  446. AWS_LS_IO_CHANNEL,
  447. "id=%p: Running %s channel task immediately as canceled due to shut down channel",
  448. (void *)channel,
  449. channel_task->type_tag);
  450. channel_task->task_fn(channel_task, channel_task->arg, AWS_TASK_STATUS_CANCELED);
  451. return;
  452. }
  453. aws_linked_list_push_back(&channel->channel_thread_tasks.list, &channel_task->node);
  454. if (run_at_nanos == 0) {
  455. aws_event_loop_schedule_task_now(channel->loop, &channel_task->wrapper_task);
  456. } else {
  457. aws_event_loop_schedule_task_future(
  458. channel->loop, &channel_task->wrapper_task, channel_task->wrapper_task.timestamp);
  459. }
  460. }
  461. static void s_register_pending_task_cross_thread(struct aws_channel *channel, struct aws_channel_task *channel_task) {
  462. AWS_LOGF_TRACE(
  463. AWS_LS_IO_CHANNEL,
  464. "id=%p: scheduling task with wrapper task id %p from "
  465. "outside the event-loop thread.",
  466. (void *)channel,
  467. (void *)&channel_task->wrapper_task);
  468. /* Outside event-loop thread... */
  469. bool should_cancel_task = false;
  470. /* Begin Critical Section */
  471. aws_mutex_lock(&channel->cross_thread_tasks.lock);
  472. if (channel->cross_thread_tasks.is_channel_shut_down) {
  473. should_cancel_task = true; /* run task outside critical section to avoid deadlock */
  474. } else {
  475. bool list_was_empty = aws_linked_list_empty(&channel->cross_thread_tasks.list);
  476. aws_linked_list_push_back(&channel->cross_thread_tasks.list, &channel_task->node);
  477. if (list_was_empty) {
  478. aws_event_loop_schedule_task_now(channel->loop, &channel->cross_thread_tasks.scheduling_task);
  479. }
  480. }
  481. aws_mutex_unlock(&channel->cross_thread_tasks.lock);
  482. /* End Critical Section */
  483. if (should_cancel_task) {
  484. channel_task->task_fn(channel_task, channel_task->arg, AWS_TASK_STATUS_CANCELED);
  485. }
  486. }
  487. static void s_reset_pending_channel_task(
  488. struct aws_channel *channel,
  489. struct aws_channel_task *channel_task,
  490. uint64_t run_at_nanos) {
  491. /* Reset every property on channel task other than user's fn & arg.*/
  492. aws_task_init(&channel_task->wrapper_task, s_channel_task_run, channel, channel_task->type_tag);
  493. channel_task->wrapper_task.timestamp = run_at_nanos;
  494. aws_linked_list_node_reset(&channel_task->node);
  495. }
  496. /* Common functionality for scheduling "now" and "future" tasks.
  497. * For "now" tasks, pass 0 for `run_at_nanos` */
  498. static void s_register_pending_task(
  499. struct aws_channel *channel,
  500. struct aws_channel_task *channel_task,
  501. uint64_t run_at_nanos) {
  502. s_reset_pending_channel_task(channel, channel_task, run_at_nanos);
  503. if (aws_channel_thread_is_callers_thread(channel)) {
  504. s_register_pending_task_in_event_loop(channel, channel_task, run_at_nanos);
  505. } else {
  506. s_register_pending_task_cross_thread(channel, channel_task);
  507. }
  508. }
  509. void aws_channel_schedule_task_now(struct aws_channel *channel, struct aws_channel_task *task) {
  510. s_register_pending_task(channel, task, 0);
  511. }
  512. void aws_channel_schedule_task_now_serialized(struct aws_channel *channel, struct aws_channel_task *task) {
  513. s_reset_pending_channel_task(channel, task, 0);
  514. s_register_pending_task_cross_thread(channel, task);
  515. }
  516. void aws_channel_schedule_task_future(
  517. struct aws_channel *channel,
  518. struct aws_channel_task *task,
  519. uint64_t run_at_nanos) {
  520. s_register_pending_task(channel, task, run_at_nanos);
  521. }
  522. bool aws_channel_thread_is_callers_thread(struct aws_channel *channel) {
  523. return aws_event_loop_thread_is_callers_thread(channel->loop);
  524. }
  525. static void s_update_channel_slot_message_overheads(struct aws_channel *channel) {
  526. size_t overhead = 0;
  527. struct aws_channel_slot *slot_iter = channel->first;
  528. while (slot_iter) {
  529. slot_iter->upstream_message_overhead = overhead;
  530. if (slot_iter->handler) {
  531. overhead += slot_iter->handler->vtable->message_overhead(slot_iter->handler);
  532. }
  533. slot_iter = slot_iter->adj_right;
  534. }
  535. }
  536. int aws_channel_slot_set_handler(struct aws_channel_slot *slot, struct aws_channel_handler *handler) {
  537. slot->handler = handler;
  538. slot->handler->slot = slot;
  539. s_update_channel_slot_message_overheads(slot->channel);
  540. return aws_channel_slot_increment_read_window(slot, slot->handler->vtable->initial_window_size(handler));
  541. }
  542. int aws_channel_slot_remove(struct aws_channel_slot *slot) {
  543. if (slot->adj_right) {
  544. slot->adj_right->adj_left = slot->adj_left;
  545. if (slot == slot->channel->first) {
  546. slot->channel->first = slot->adj_right;
  547. }
  548. }
  549. if (slot->adj_left) {
  550. slot->adj_left->adj_right = slot->adj_right;
  551. }
  552. if (slot == slot->channel->first) {
  553. slot->channel->first = NULL;
  554. }
  555. s_update_channel_slot_message_overheads(slot->channel);
  556. s_cleanup_slot(slot);
  557. return AWS_OP_SUCCESS;
  558. }
  559. int aws_channel_slot_replace(struct aws_channel_slot *remove, struct aws_channel_slot *new_slot) {
  560. new_slot->adj_left = remove->adj_left;
  561. if (remove->adj_left) {
  562. remove->adj_left->adj_right = new_slot;
  563. }
  564. new_slot->adj_right = remove->adj_right;
  565. if (remove->adj_right) {
  566. remove->adj_right->adj_left = new_slot;
  567. }
  568. if (remove == remove->channel->first) {
  569. remove->channel->first = new_slot;
  570. }
  571. s_update_channel_slot_message_overheads(remove->channel);
  572. s_cleanup_slot(remove);
  573. return AWS_OP_SUCCESS;
  574. }
  575. int aws_channel_slot_insert_right(struct aws_channel_slot *slot, struct aws_channel_slot *to_add) {
  576. to_add->adj_right = slot->adj_right;
  577. if (slot->adj_right) {
  578. slot->adj_right->adj_left = to_add;
  579. }
  580. slot->adj_right = to_add;
  581. to_add->adj_left = slot;
  582. return AWS_OP_SUCCESS;
  583. }
  584. int aws_channel_slot_insert_end(struct aws_channel *channel, struct aws_channel_slot *to_add) {
  585. /* It's actually impossible there's not a first if the user went through the aws_channel_slot_new() function.
  586. * But also check that a user didn't call insert_end if it's the first slot in the channel since first would already
  587. * have been set. */
  588. if (AWS_LIKELY(channel->first && channel->first != to_add)) {
  589. struct aws_channel_slot *cur = channel->first;
  590. while (cur->adj_right) {
  591. cur = cur->adj_right;
  592. }
  593. return aws_channel_slot_insert_right(cur, to_add);
  594. }
  595. AWS_ASSERT(0);
  596. return AWS_OP_ERR;
  597. }
  598. int aws_channel_slot_insert_left(struct aws_channel_slot *slot, struct aws_channel_slot *to_add) {
  599. to_add->adj_left = slot->adj_left;
  600. if (slot->adj_left) {
  601. slot->adj_left->adj_right = to_add;
  602. }
  603. slot->adj_left = to_add;
  604. to_add->adj_right = slot;
  605. if (slot == slot->channel->first) {
  606. slot->channel->first = to_add;
  607. }
  608. return AWS_OP_SUCCESS;
  609. }
  610. int aws_channel_slot_send_message(
  611. struct aws_channel_slot *slot,
  612. struct aws_io_message *message,
  613. enum aws_channel_direction dir) {
  614. if (dir == AWS_CHANNEL_DIR_READ) {
  615. AWS_ASSERT(slot->adj_right);
  616. AWS_ASSERT(slot->adj_right->handler);
  617. if (!slot->channel->read_back_pressure_enabled || slot->adj_right->window_size >= message->message_data.len) {
  618. AWS_LOGF_TRACE(
  619. AWS_LS_IO_CHANNEL,
  620. "id=%p: sending read message of size %zu, "
  621. "from slot %p to slot %p with handler %p.",
  622. (void *)slot->channel,
  623. message->message_data.len,
  624. (void *)slot,
  625. (void *)slot->adj_right,
  626. (void *)slot->adj_right->handler);
  627. slot->adj_right->window_size -= message->message_data.len;
  628. return aws_channel_handler_process_read_message(slot->adj_right->handler, slot->adj_right, message);
  629. }
  630. AWS_LOGF_ERROR(
  631. AWS_LS_IO_CHANNEL,
  632. "id=%p: sending message of size %zu, "
  633. "from slot %p to slot %p with handler %p, but this would exceed the channel's "
  634. "read window, this is always a programming error.",
  635. (void *)slot->channel,
  636. message->message_data.len,
  637. (void *)slot,
  638. (void *)slot->adj_right,
  639. (void *)slot->adj_right->handler);
  640. return aws_raise_error(AWS_IO_CHANNEL_READ_WOULD_EXCEED_WINDOW);
  641. }
  642. AWS_ASSERT(slot->adj_left);
  643. AWS_ASSERT(slot->adj_left->handler);
  644. AWS_LOGF_TRACE(
  645. AWS_LS_IO_CHANNEL,
  646. "id=%p: sending write message of size %zu, "
  647. "from slot %p to slot %p with handler %p.",
  648. (void *)slot->channel,
  649. message->message_data.len,
  650. (void *)slot,
  651. (void *)slot->adj_left,
  652. (void *)slot->adj_left->handler);
  653. return aws_channel_handler_process_write_message(slot->adj_left->handler, slot->adj_left, message);
  654. }
  655. struct aws_io_message *aws_channel_slot_acquire_max_message_for_write(struct aws_channel_slot *slot) {
  656. AWS_PRECONDITION(slot);
  657. AWS_PRECONDITION(slot->channel);
  658. AWS_PRECONDITION(aws_channel_thread_is_callers_thread(slot->channel));
  659. const size_t overhead = aws_channel_slot_upstream_message_overhead(slot);
  660. if (overhead >= g_aws_channel_max_fragment_size) {
  661. AWS_LOGF_ERROR(
  662. AWS_LS_IO_CHANNEL, "id=%p: Upstream overhead exceeds channel's max message size.", (void *)slot->channel);
  663. aws_raise_error(AWS_ERROR_INVALID_STATE);
  664. return NULL;
  665. }
  666. const size_t size_hint = g_aws_channel_max_fragment_size - overhead;
  667. return aws_channel_acquire_message_from_pool(slot->channel, AWS_IO_MESSAGE_APPLICATION_DATA, size_hint);
  668. }
  669. static void s_window_update_task(struct aws_channel_task *channel_task, void *arg, enum aws_task_status status) {
  670. (void)channel_task;
  671. struct aws_channel *channel = arg;
  672. channel->window_update_scheduled = false;
  673. if (status == AWS_TASK_STATUS_RUN_READY && channel->channel_state < AWS_CHANNEL_SHUTTING_DOWN) {
  674. /* get the right-most slot to start the updates. */
  675. struct aws_channel_slot *slot = channel->first;
  676. while (slot->adj_right) {
  677. slot = slot->adj_right;
  678. }
  679. while (slot->adj_left) {
  680. struct aws_channel_slot *upstream_slot = slot->adj_left;
  681. if (upstream_slot->handler) {
  682. slot->window_size = aws_add_size_saturating(slot->window_size, slot->current_window_update_batch_size);
  683. size_t update_size = slot->current_window_update_batch_size;
  684. slot->current_window_update_batch_size = 0;
  685. if (aws_channel_handler_increment_read_window(upstream_slot->handler, upstream_slot, update_size)) {
  686. AWS_LOGF_ERROR(
  687. AWS_LS_IO_CHANNEL,
  688. "channel %p: channel update task failed with status %d",
  689. (void *)slot->channel,
  690. aws_last_error());
  691. aws_channel_shutdown(channel, aws_last_error());
  692. return;
  693. }
  694. }
  695. slot = slot->adj_left;
  696. }
  697. }
  698. }
  699. int aws_channel_slot_increment_read_window(struct aws_channel_slot *slot, size_t window) {
  700. if (slot->channel->read_back_pressure_enabled && slot->channel->channel_state < AWS_CHANNEL_SHUTTING_DOWN) {
  701. slot->current_window_update_batch_size =
  702. aws_add_size_saturating(slot->current_window_update_batch_size, window);
  703. if (!slot->channel->window_update_scheduled &&
  704. slot->window_size <= slot->channel->window_update_batch_emit_threshold) {
  705. slot->channel->window_update_scheduled = true;
  706. aws_channel_task_init(
  707. &slot->channel->window_update_task, s_window_update_task, slot->channel, "window update task");
  708. aws_channel_schedule_task_now(slot->channel, &slot->channel->window_update_task);
  709. }
  710. }
  711. return AWS_OP_SUCCESS;
  712. }
  713. int aws_channel_slot_shutdown(
  714. struct aws_channel_slot *slot,
  715. enum aws_channel_direction dir,
  716. int err_code,
  717. bool free_scarce_resources_immediately) {
  718. AWS_ASSERT(slot->handler);
  719. AWS_LOGF_TRACE(
  720. AWS_LS_IO_CHANNEL,
  721. "id=%p: shutting down slot %p, with handler %p "
  722. "in %s direction with error code %d",
  723. (void *)slot->channel,
  724. (void *)slot,
  725. (void *)slot->handler,
  726. (dir == AWS_CHANNEL_DIR_READ) ? "read" : "write",
  727. err_code);
  728. return aws_channel_handler_shutdown(slot->handler, slot, dir, err_code, free_scarce_resources_immediately);
  729. }
  730. static void s_on_shutdown_completion_task(struct aws_task *task, void *arg, enum aws_task_status status) {
  731. (void)status;
  732. struct aws_shutdown_notification_task *shutdown_notify = (struct aws_shutdown_notification_task *)task;
  733. struct aws_channel *channel = arg;
  734. AWS_ASSERT(channel->channel_state == AWS_CHANNEL_SHUT_DOWN);
  735. /* Cancel tasks that have been scheduled with the event loop */
  736. while (!aws_linked_list_empty(&channel->channel_thread_tasks.list)) {
  737. struct aws_linked_list_node *node = aws_linked_list_front(&channel->channel_thread_tasks.list);
  738. struct aws_channel_task *channel_task = AWS_CONTAINER_OF(node, struct aws_channel_task, node);
  739. AWS_LOGF_DEBUG(
  740. AWS_LS_IO_CHANNEL,
  741. "id=%p: during shutdown, canceling task %p",
  742. (void *)channel,
  743. (void *)&channel_task->wrapper_task);
  744. /* The task will remove itself from the list when it's canceled */
  745. aws_event_loop_cancel_task(channel->loop, &channel_task->wrapper_task);
  746. }
  747. /* Cancel off-thread tasks, which haven't made it to the event-loop thread yet */
  748. aws_mutex_lock(&channel->cross_thread_tasks.lock);
  749. bool cancel_cross_thread_tasks = !aws_linked_list_empty(&channel->cross_thread_tasks.list);
  750. aws_mutex_unlock(&channel->cross_thread_tasks.lock);
  751. if (cancel_cross_thread_tasks) {
  752. aws_event_loop_cancel_task(channel->loop, &channel->cross_thread_tasks.scheduling_task);
  753. }
  754. AWS_ASSERT(aws_linked_list_empty(&channel->channel_thread_tasks.list));
  755. AWS_ASSERT(aws_linked_list_empty(&channel->cross_thread_tasks.list));
  756. channel->on_shutdown_completed(channel, shutdown_notify->error_code, channel->shutdown_user_data);
  757. }
  758. static void s_run_shutdown_write_direction(struct aws_task *task, void *arg, enum aws_task_status status) {
  759. (void)arg;
  760. (void)status;
  761. struct aws_shutdown_notification_task *shutdown_notify = (struct aws_shutdown_notification_task *)task;
  762. task->fn = NULL;
  763. task->arg = NULL;
  764. struct aws_channel_slot *slot = shutdown_notify->slot;
  765. aws_channel_handler_shutdown(
  766. slot->handler, slot, AWS_CHANNEL_DIR_WRITE, shutdown_notify->error_code, shutdown_notify->shutdown_immediately);
  767. }
  768. int aws_channel_slot_on_handler_shutdown_complete(
  769. struct aws_channel_slot *slot,
  770. enum aws_channel_direction dir,
  771. int err_code,
  772. bool free_scarce_resources_immediately) {
  773. AWS_LOGF_DEBUG(
  774. AWS_LS_IO_CHANNEL,
  775. "id=%p: handler %p shutdown in %s dir completed.",
  776. (void *)slot->channel,
  777. (void *)slot->handler,
  778. (dir == AWS_CHANNEL_DIR_READ) ? "read" : "write");
  779. if (slot->channel->channel_state == AWS_CHANNEL_SHUT_DOWN) {
  780. return AWS_OP_SUCCESS;
  781. }
  782. if (dir == AWS_CHANNEL_DIR_READ) {
  783. if (slot->adj_right && slot->adj_right->handler) {
  784. return aws_channel_handler_shutdown(
  785. slot->adj_right->handler, slot->adj_right, dir, err_code, free_scarce_resources_immediately);
  786. }
  787. /* break the shutdown sequence so we don't have handlers having to deal with their memory disappearing out from
  788. * under them during a shutdown process. */
  789. slot->channel->shutdown_notify_task.slot = slot;
  790. slot->channel->shutdown_notify_task.shutdown_immediately = free_scarce_resources_immediately;
  791. slot->channel->shutdown_notify_task.error_code = err_code;
  792. slot->channel->shutdown_notify_task.task.fn = s_run_shutdown_write_direction;
  793. slot->channel->shutdown_notify_task.task.arg = NULL;
  794. aws_event_loop_schedule_task_now(slot->channel->loop, &slot->channel->shutdown_notify_task.task);
  795. return AWS_OP_SUCCESS;
  796. }
  797. if (slot->adj_left && slot->adj_left->handler) {
  798. return aws_channel_handler_shutdown(
  799. slot->adj_left->handler, slot->adj_left, dir, err_code, free_scarce_resources_immediately);
  800. }
  801. if (slot->channel->first == slot) {
  802. slot->channel->channel_state = AWS_CHANNEL_SHUT_DOWN;
  803. aws_mutex_lock(&slot->channel->cross_thread_tasks.lock);
  804. slot->channel->cross_thread_tasks.is_channel_shut_down = true;
  805. aws_mutex_unlock(&slot->channel->cross_thread_tasks.lock);
  806. if (slot->channel->on_shutdown_completed) {
  807. slot->channel->shutdown_notify_task.task.fn = s_on_shutdown_completion_task;
  808. slot->channel->shutdown_notify_task.task.arg = slot->channel;
  809. slot->channel->shutdown_notify_task.error_code = err_code;
  810. aws_event_loop_schedule_task_now(slot->channel->loop, &slot->channel->shutdown_notify_task.task);
  811. }
  812. }
  813. return AWS_OP_SUCCESS;
  814. }
  815. size_t aws_channel_slot_downstream_read_window(struct aws_channel_slot *slot) {
  816. AWS_ASSERT(slot->adj_right);
  817. return slot->channel->read_back_pressure_enabled ? slot->adj_right->window_size : SIZE_MAX;
  818. }
  819. size_t aws_channel_slot_upstream_message_overhead(struct aws_channel_slot *slot) {
  820. return slot->upstream_message_overhead;
  821. }
  822. void aws_channel_handler_destroy(struct aws_channel_handler *handler) {
  823. AWS_ASSERT(handler->vtable && handler->vtable->destroy);
  824. handler->vtable->destroy(handler);
  825. }
  826. int aws_channel_handler_process_read_message(
  827. struct aws_channel_handler *handler,
  828. struct aws_channel_slot *slot,
  829. struct aws_io_message *message) {
  830. AWS_ASSERT(handler->vtable && handler->vtable->process_read_message);
  831. return handler->vtable->process_read_message(handler, slot, message);
  832. }
  833. int aws_channel_handler_process_write_message(
  834. struct aws_channel_handler *handler,
  835. struct aws_channel_slot *slot,
  836. struct aws_io_message *message) {
  837. AWS_ASSERT(handler->vtable && handler->vtable->process_write_message);
  838. return handler->vtable->process_write_message(handler, slot, message);
  839. }
  840. int aws_channel_handler_increment_read_window(
  841. struct aws_channel_handler *handler,
  842. struct aws_channel_slot *slot,
  843. size_t size) {
  844. AWS_ASSERT(handler->vtable && handler->vtable->increment_read_window);
  845. return handler->vtable->increment_read_window(handler, slot, size);
  846. }
  847. int aws_channel_handler_shutdown(
  848. struct aws_channel_handler *handler,
  849. struct aws_channel_slot *slot,
  850. enum aws_channel_direction dir,
  851. int error_code,
  852. bool free_scarce_resources_immediately) {
  853. AWS_ASSERT(handler->vtable && handler->vtable->shutdown);
  854. return handler->vtable->shutdown(handler, slot, dir, error_code, free_scarce_resources_immediately);
  855. }
  856. size_t aws_channel_handler_initial_window_size(struct aws_channel_handler *handler) {
  857. AWS_ASSERT(handler->vtable && handler->vtable->initial_window_size);
  858. return handler->vtable->initial_window_size(handler);
  859. }
  860. struct aws_channel_slot *aws_channel_get_first_slot(struct aws_channel *channel) {
  861. return channel->first;
  862. }
  863. static void s_reset_statistics(struct aws_channel *channel) {
  864. AWS_FATAL_ASSERT(aws_channel_thread_is_callers_thread(channel));
  865. struct aws_channel_slot *current_slot = channel->first;
  866. while (current_slot) {
  867. struct aws_channel_handler *handler = current_slot->handler;
  868. if (handler != NULL && handler->vtable->reset_statistics != NULL) {
  869. handler->vtable->reset_statistics(handler);
  870. }
  871. current_slot = current_slot->adj_right;
  872. }
  873. }
  874. static void s_channel_gather_statistics_task(struct aws_task *task, void *arg, enum aws_task_status status) {
  875. if (status != AWS_TASK_STATUS_RUN_READY) {
  876. return;
  877. }
  878. struct aws_channel *channel = arg;
  879. if (channel->statistics_handler == NULL) {
  880. return;
  881. }
  882. if (channel->channel_state == AWS_CHANNEL_SHUTTING_DOWN || channel->channel_state == AWS_CHANNEL_SHUT_DOWN) {
  883. return;
  884. }
  885. uint64_t now_ns = 0;
  886. if (aws_channel_current_clock_time(channel, &now_ns)) {
  887. return;
  888. }
  889. uint64_t now_ms = aws_timestamp_convert(now_ns, AWS_TIMESTAMP_NANOS, AWS_TIMESTAMP_MILLIS, NULL);
  890. struct aws_array_list *statistics_list = &channel->statistic_list;
  891. aws_array_list_clear(statistics_list);
  892. struct aws_channel_slot *current_slot = channel->first;
  893. while (current_slot) {
  894. struct aws_channel_handler *handler = current_slot->handler;
  895. if (handler != NULL && handler->vtable->gather_statistics != NULL) {
  896. handler->vtable->gather_statistics(handler, statistics_list);
  897. }
  898. current_slot = current_slot->adj_right;
  899. }
  900. struct aws_crt_statistics_sample_interval sample_interval = {
  901. .begin_time_ms = channel->statistics_interval_start_time_ms, .end_time_ms = now_ms};
  902. aws_crt_statistics_handler_process_statistics(
  903. channel->statistics_handler, &sample_interval, statistics_list, channel);
  904. s_reset_statistics(channel);
  905. uint64_t reschedule_interval_ns = aws_timestamp_convert(
  906. aws_crt_statistics_handler_get_report_interval_ms(channel->statistics_handler),
  907. AWS_TIMESTAMP_MILLIS,
  908. AWS_TIMESTAMP_NANOS,
  909. NULL);
  910. aws_event_loop_schedule_task_future(channel->loop, task, now_ns + reschedule_interval_ns);
  911. channel->statistics_interval_start_time_ms = now_ms;
  912. }
  913. int aws_channel_set_statistics_handler(struct aws_channel *channel, struct aws_crt_statistics_handler *handler) {
  914. AWS_FATAL_ASSERT(aws_channel_thread_is_callers_thread(channel));
  915. if (channel->statistics_handler) {
  916. aws_crt_statistics_handler_destroy(channel->statistics_handler);
  917. aws_event_loop_cancel_task(channel->loop, &channel->statistics_task);
  918. channel->statistics_handler = NULL;
  919. }
  920. if (handler != NULL) {
  921. aws_task_init(&channel->statistics_task, s_channel_gather_statistics_task, channel, "gather_statistics");
  922. uint64_t now_ns = 0;
  923. if (aws_channel_current_clock_time(channel, &now_ns)) {
  924. return AWS_OP_ERR;
  925. }
  926. uint64_t report_time_ns = now_ns + aws_timestamp_convert(
  927. aws_crt_statistics_handler_get_report_interval_ms(handler),
  928. AWS_TIMESTAMP_MILLIS,
  929. AWS_TIMESTAMP_NANOS,
  930. NULL);
  931. channel->statistics_interval_start_time_ms =
  932. aws_timestamp_convert(now_ns, AWS_TIMESTAMP_NANOS, AWS_TIMESTAMP_MILLIS, NULL);
  933. s_reset_statistics(channel);
  934. aws_event_loop_schedule_task_future(channel->loop, &channel->statistics_task, report_time_ns);
  935. }
  936. channel->statistics_handler = handler;
  937. return AWS_OP_SUCCESS;
  938. }
  939. struct aws_event_loop *aws_channel_get_event_loop(struct aws_channel *channel) {
  940. return channel->loop;
  941. }
  942. int aws_channel_trigger_read(struct aws_channel *channel) {
  943. if (channel == NULL) {
  944. return aws_raise_error(AWS_ERROR_INVALID_ARGUMENT);
  945. }
  946. if (!aws_channel_thread_is_callers_thread(channel)) {
  947. return aws_raise_error(AWS_ERROR_INVALID_STATE);
  948. }
  949. struct aws_channel_slot *slot = channel->first;
  950. if (slot == NULL) {
  951. return aws_raise_error(AWS_ERROR_INVALID_STATE);
  952. }
  953. struct aws_channel_handler *handler = slot->handler;
  954. if (handler == NULL) {
  955. return aws_raise_error(AWS_ERROR_INVALID_STATE);
  956. }
  957. if (handler->vtable->trigger_read != NULL) {
  958. handler->vtable->trigger_read(handler);
  959. }
  960. return AWS_OP_SUCCESS;
  961. }