#include #include #include #include #include #include extern coreid_t my_core_id; extern struct aos_urpc urpc_to_bsp; extern struct aos_urpc urpc_to_app; extern struct waitset urpc_to_app_ws; // This is only valid on the BSP core. static coreid_t server_coreid[UMP_SERVER_COUNT]; // This only contains channels for servers on the same core. static struct lmp_chan server_chan[UMP_SERVER_COUNT]; struct waitlist_entry { struct capref cap; void (*callback)(void *arg, errval_t err); void *callback_arg; struct waitlist_entry *next; }; static struct waitlist_entry *server_waitlist_head[UMP_SERVER_COUNT]; static struct waitlist_entry *server_waitlist_tail[UMP_SERVER_COUNT]; errval_t ump_binding_register_service_global (enum ump_server_id server_id, coreid_t coreid) { assert(my_core_id == 0); if (server_id >= UMP_SERVER_COUNT) return ERR_INVALID_ARGS; if (server_coreid[server_id] != 0) return LIB_ERR_UMP_ALREADY_REGISTERED; server_coreid[server_id] = 1 + coreid; return SYS_ERR_OK; } errval_t ump_binding_get_global (enum ump_server_id server_id, coreid_t *coreid) { assert(my_core_id == 0); if (server_id >= UMP_SERVER_COUNT) return ERR_INVALID_ARGS; if (server_coreid[server_id] == 0) return LIB_ERR_UMP_NOT_REGISTERED; *coreid = server_coreid[server_id] - 1; return SYS_ERR_OK; } errval_t ump_binding_register_service ( enum ump_server_id server_id, struct capref accept_ep_cap, struct capref *reply_ep_cap ) { errval_t err; // Create the LMP channel for accepting connections if (server_id >= UMP_SERVER_COUNT) return ERR_INVALID_ARGS; struct lmp_chan *lc = &server_chan[server_id]; if (lc->endpoint != NULL) return LIB_ERR_UMP_ALREADY_REGISTERED; lmp_chan_init(lc); err = endpoint_create(DEFAULT_LMP_BUF_WORDS, &lc->local_cap, &lc->endpoint); if (err_is_fail(err)) return err_push(err, LIB_ERR_ENDPOINT_CREATE); lc->remote_cap = accept_ep_cap; *reply_ep_cap = lc->local_cap; // Register the service in the global table if (my_core_id == 0) { err = ump_binding_register_service_global(server_id, 0); } else { err = do_aos_urpc( &urpc_to_bsp, RPC_MTYPE_UMP_REGISTER_GLOBAL, NULL_CAP, 0, server_id, my_core_id, NULL, NULL, NULL, NULL ); } if (err_is_fail(err)) return err_push(err, LIB_ERR_UMP_REGISTER_GLOBAL); return SYS_ERR_OK; } struct connect_urpc_arg { struct waitset_chanstate chan; void (*callback)(void *arg, errval_t err); void *callback_arg; struct capref *ret_cap; }; static void do_connect_urpc(void *arg) { errval_t err; struct connect_urpc_arg *urpc_arg = arg; err = do_aos_urpc( &urpc_to_app, RPC_MTYPE_UMP_CONNECT, NULL_CAP, 0, 0, 0, urpc_arg->ret_cap, NULL, NULL, NULL ); urpc_arg->callback(urpc_arg->callback_arg, err); free(urpc_arg); } static void ump_connect_send (void *arg); static void ump_connect_recv (void *arg); errval_t ump_binding_connect ( enum ump_server_id server_id, struct capref *ret_cap, void (*callback)(void *arg, errval_t err), void *callback_arg ) { errval_t err; // Get the service from the global table coreid_t coreid = 0; if (my_core_id == 0) { err = ump_binding_get_global(server_id, &coreid); } else { uintptr_t coreid_tmp; err = do_aos_urpc( &urpc_to_bsp, RPC_MTYPE_UMP_GET_GLOBAL, NULL_CAP, 0, server_id, 0, NULL, NULL, &coreid_tmp, NULL ); coreid = coreid_tmp; } if (err_is_fail(err)) { if (err == LIB_ERR_UMP_NOT_REGISTERED) return err; return err_push(err, LIB_ERR_UMP_GET_GLOBAL); } if (coreid == my_core_id) { // Allocate ump frame err = frame_alloc(ret_cap, UMP_FRAME_SIZE, NULL); if (err_is_fail(err)) return err_push(err, LIB_ERR_FRAME_ALLOC); // Initialize ump frame uint8_t *frame_data; err = paging_map_frame( get_current_paging_state(), (void **)&frame_data, UMP_FRAME_SIZE, *ret_cap ); if (err_is_fail(err)) return err_push(err, LIB_ERR_PMAP_MAP); // NOTE rueegges: the ring buffer MUST be initialized to all zero so // all ring buffer entries are initially owned by the // respective senders memset(frame_data, 0, UMP_FRAME_SIZE); err = paging_unmap(get_current_paging_state(), (void *)frame_data); if (err_is_fail(err)) return err_push(err, LIB_ERR_PMAP_UNMAP); // Create a waitlist entry for sending the connect call to the server struct waitlist_entry *entry = malloc(sizeof(struct waitlist_entry)); if (entry == NULL) return LIB_ERR_MALLOC_FAIL; entry->cap = *ret_cap; entry->callback = callback; entry->callback_arg = callback_arg; entry->next = NULL; if (server_waitlist_tail[server_id] == NULL) { server_waitlist_head[server_id] = entry; ump_connect_send((void *)server_id); } else { server_waitlist_tail[server_id]->next = entry; } server_waitlist_tail[server_id] = entry; return AOS_ERR_RPC_ASYNC_REPLY; } else if (my_core_id == 0) { struct connect_urpc_arg *urpc_arg = malloc(sizeof(struct connect_urpc_arg)); if (urpc_arg == NULL) return LIB_ERR_MALLOC_FAIL; urpc_arg->ret_cap = ret_cap; urpc_arg->callback = callback; urpc_arg->callback_arg = callback_arg; waitset_chanstate_init(&urpc_arg->chan, CHANTYPE_OTHER); waitset_chan_trigger_closure(&urpc_to_app_ws, &urpc_arg->chan, MKCLOSURE(do_connect_urpc, urpc_arg)); return AOS_ERR_RPC_ASYNC_REPLY; } else { return do_aos_urpc( &urpc_to_bsp, RPC_MTYPE_UMP_CONNECT, NULL_CAP, 0, server_id, 0, ret_cap, NULL, NULL, NULL ); } } static void ump_connect_send (void *arg) { enum ump_server_id server_id = (enum ump_server_id)arg; errval_t err; struct waitlist_entry *entry = server_waitlist_head[server_id]; assert(entry != NULL); err = lmp_chan_send4( &server_chan[server_id], LMP_FLAG_YIELD | LMP_FLAG_SYNC, entry->cap, 0, 0, 0, 0 ); if (err_is_fail(err)) { if (!lmp_err_is_transient(err)) { DEBUG_ERR(err, "Could not send connect message"); return; } // Cannot send right now, try again later err = lmp_chan_register_send(&server_chan[server_id], get_default_waitset(), MKCLOSURE(ump_connect_send, arg)); if (err_is_fail(err)) { DEBUG_ERR(err, "Could not register send handler"); } return; } // Sent connect, wait for reply err = lmp_chan_register_recv(&server_chan[server_id], get_default_waitset(), MKCLOSURE(ump_connect_recv, arg)); if (err_is_fail(err)) { DEBUG_ERR(err, "Could not register receive handler"); } } static void ump_connect_recv (void *arg) { enum ump_server_id server_id = (enum ump_server_id)arg; errval_t err; struct lmp_recv_msg msg = LMP_RECV_MSG_INIT; err = lmp_chan_recv(&server_chan[server_id], &msg, NULL); assert(err_is_ok(err)); struct waitlist_entry *entry = server_waitlist_head[server_id]; assert(entry->callback != NULL); entry->callback(entry->callback_arg, msg.words[0]); server_waitlist_head[server_id] = entry->next; if (entry->next == NULL) { server_waitlist_tail[server_id] = NULL; } else { ump_connect_send(arg); } free(entry); }