// Wire // ==== // TCP on bytes: a List of octets in and out, where Base's socket // effects decode and encode UTF-8. After bend-kit-wire (MIT-0). #ifndef GROUNDS_WIRE #define GROUNDS_WIRE static Term gw_list(Env e, const char* p, u64 n) { Term xs = term_pak(CID_NIL, 0); for (u64 i = n; i > 0; i -= 1) { xs = io_node(e, CID_CON, ((uint8_t*)p)[i - 1], xs); } return xs; } // A value past 255 is not an octet: bad is set, and the send fails. static char* gw_octets(Env e, Term xs, u64* len, bool* bad) { u64 cap = 64; u64 n = 0; char* buf = io_mem(malloc(cap)); *bad = false; while (term_aux(xs) == CID_CON) { Term fb[2]; spare_free(e, cls_fit(2), ctr_take(e, xs, 2, fb)); if (n + 1 > cap) { cap *= 2; buf = io_mem(realloc(buf, cap)); } *bad = *bad || (u64)fb[0] > 255; buf[n++] = (char)(fb[0] & 0xFF); xs = fb[1]; } *len = n; return buf; } #endif #ifdef CID_WIRE_RECV // The loop parked the request until the socket was readable; a recv that // still finds nothing (the socket is non-blocking) parks again. static Term gw_recv_more(Env e, IoWork* w) { int fd = (int)w->hand; ssize_t n = io_sys_end(w, recv(fd, w->data, (size_t)w->made, 0)); if (w->code == EAGAIN) { return io_wait_on(w, fd, POLLIN, 0, gw_recv_more); } Term r = w->code ? io_fail(e, w->code, NULL) : io_done(e, gw_list(e, w->data, (u64)n)); free(w->data); return io_tup(e, io_hand(w->hand), r); } Term gw_recv_run(Env e, Term* f, IoWork* w) { w->hand = (intptr_t)io_hand_v(f[0]); w->made = f[1] < INT32_MAX ? (intptr_t)f[1] : INT32_MAX; w->data = io_mem(malloc((size_t)w->made + 1)); return gw_recv_more(e, w); } static void __attribute__((constructor)) gw_recv_use(void) { io_eff(CID_WIRE_RECV, gw_recv_run, IO_READ); } #endif #ifdef CID_WIRE_RECV_TIMEOUT // wire_recv with a deadline, as Base's TCP.poll: parks on the socket and // the clock, whichever fires first; None{} past the deadline. static Term gw_poll_end(Env e, IoWork* w, Term r) { free(w->data); return io_tup(e, io_hand(w->hand), r); } static Term gw_poll_more(Env e, IoWork* w) { int fd = (int)w->hand; u64 at = io_wait_time(w); ssize_t n = io_sys_end(w, recv(fd, w->data, (size_t)w->made, 0)); if (w->code == EAGAIN) { return io_tick() < at ? io_wait_on(w, fd, POLLIN, at, gw_poll_more) : gw_poll_end(e, w, io_done(e, term_pak(CID_NONE, 0))); } return gw_poll_end(e, w, w->code ? io_fail(e, w->code, NULL) : io_done(e, io_box(e, CID_SOME, gw_list(e, w->data, (u64)n)))); } Term gw_poll_run(Env e, Term* f, IoWork* w) { w->hand = (intptr_t)io_hand_v(f[0]); w->made = f[1] < INT32_MAX ? (intptr_t)f[1] : INT32_MAX; w->data = io_mem(malloc((size_t)w->made + 1)); return io_wait_on(w, (int)w->hand, POLLIN, io_tick() + (u64)f[2] * 1000000ull, gw_poll_more); } static void __attribute__((constructor)) gw_poll_use(void) { io_eff(CID_WIRE_RECV_TIMEOUT, gw_poll_run, 0); } #endif #ifdef CID_WIRE_SEND // Sends what is left; a full socket parks until it is writable. static Term gw_send_more(Env e, IoWork* w) { int fd = (int)w->hand; while (w->code == 0 && (u64)w->made < w->size) { ssize_t n = send(fd, w->data + w->made, w->size - (u64)w->made, 0); if (n < 0 && errno == EAGAIN) { return io_wait_on(w, fd, POLLOUT, 0, gw_send_more); } w->made += io_sys_end(w, n); } Term r = w->code != 0 ? io_fail(e, w->code, NULL) : io_done(e, term_pak(CID_UNIT, 0)); free(w->data); return io_tup(e, io_hand(w->hand), r); } Term gw_send_run(Env e, Term* f, IoWork* w) { bool bad; w->hand = (intptr_t)io_hand_v(f[0]); w->data = gw_octets(e, f[1], &w->size, &bad); w->made = 0; w->code = bad ? EINVAL : 0; return gw_send_more(e, w); } static void __attribute__((constructor)) gw_send_use(void) { io_eff(CID_WIRE_SEND, gw_send_run, 0); } #endif