diff --git a/errors/errno.fugu b/errors/errno.fugu index 127ece6..f32e761 100755 --- a/errors/errno.fugu +++ b/errors/errno.fugu @@ -372,6 +372,12 @@ errors libaos LIB_ERR_ { failure MSGBUF_CANNOT_GROW "Failed to grow message buffer while marshalling", failure RCK_NOTIFY "Failure in rck_notify()", failure IPI_NOTIFY "Failure in ipi_notify()", + failure URPC_BAD_CAPTYPE "This type of capability cannot be transferred to another core", + failure UMP_NOT_REGISTERED "The requested server is not yet registered", + failure UMP_ALREADY_REGISTERED "Failed to register server, it is already registered", + failure UMP_REGISTER_GLOBAL "Failure in register_service_global()", + failure UMP_GET_GLOBAL "Failure in get_global()", + failure UMP_REGISTER "Failure in ump_register()", // IDC binding/export and Monitor client interface failure MONITOR_CLIENT_BIND "Error in monitor_client_lmp_bind()", diff --git a/hake/menu.lst.armv8_a57_qemu b/hake/menu.lst.armv8_a57_qemu index 8a733bb..fff9c2c 100644 --- a/hake/menu.lst.armv8_a57_qemu +++ b/hake/menu.lst.armv8_a57_qemu @@ -5,9 +5,10 @@ bootdriver /armv8/sbin/boot_armv8_generic cpudriver /armv8/sbin/cpu_a57_qemu loglevel=3 serial=0x9000000 logmask=128 module /armv8/sbin/init -module /armv8/sbin/hello spawn_remote +module /armv8/sbin/hello echoclient module /armv8/sbin/memeater module /armv8/sbin/mallocator module /armv8/sbin/stackoverflow +module /armv8/sbin/echoserver # End of file, this needs to have a certain length... diff --git a/hake/menu.lst.armv8_imx8x b/hake/menu.lst.armv8_imx8x index 927f20b..3cd2895 100644 --- a/hake/menu.lst.armv8_imx8x +++ b/hake/menu.lst.armv8_imx8x @@ -5,7 +5,8 @@ bootdriver /armv8/sbin/boot_armv8_generic cpudriver /armv8/sbin/cpu_imx8x module /armv8/sbin/init -module /armv8/sbin/hello spawn_remote +module /armv8/sbin/hello echoclient module /armv8/sbin/memeater module /armv8/sbin/mallocator module /armv8/sbin/stackoverflow +module /armv8/sbin/echoserver diff --git a/include/aos/aos_rpc.h b/include/aos/aos_rpc.h index 9caf240..abf60e2 100644 --- a/include/aos/aos_rpc.h +++ b/include/aos/aos_rpc.h @@ -16,6 +16,7 @@ #define _LIB_BARRELFISH_AOS_MESSAGES_H #include +#include #define RPC_SHARED_SIZE PAGE_SIZE @@ -35,6 +36,10 @@ enum rpc_mtype { RPC_MTYPE_ALLOCATE_PID, RPC_MTYPE_GET_BOOTINFO, RPC_MTYPE_NOP, + RPC_MTYPE_UMP_REGISTER, + RPC_MTYPE_UMP_CONNECT, + RPC_MTYPE_UMP_REGISTER_GLOBAL, + RPC_MTYPE_UMP_GET_GLOBAL, RPC_MTYPE_COUNT // How many message types exist }; @@ -136,6 +141,20 @@ errval_t aos_rpc_process_get_name(struct aos_rpc *chan, domainid_t pid, errval_t aos_rpc_process_get_all_pids(struct aos_rpc *chan, domainid_t **pids, size_t *pid_count); +/** + * \brief Register an UMP server + * \arg server_id The server id to register. + * \arg accept_ep_cap An LMP endpoint which will receive incoming connections. + */ +errval_t aos_rpc_ump_register(struct aos_rpc *rpc, enum ump_server_id server_id, struct capref accept_ep_cap, struct capref *reply_ep_cap); + +/** + * \brief Connect to an UMP server + * \arg server_id The server id to connect to. + * \arg retcap Will be filled with the UMP frame cap. + */ +errval_t aos_rpc_ump_connect(struct aos_rpc *rpc, enum ump_server_id server_id, struct capref *retcap); + /** * \brief Returns the RPC channel to init. */ diff --git a/include/aos/aos_urpc.h b/include/aos/aos_urpc.h index bf354e1..15f4ad3 100644 --- a/include/aos/aos_urpc.h +++ b/include/aos/aos_urpc.h @@ -32,6 +32,7 @@ struct aos_urpc { struct aos_urpc_server { struct generic_rpc_server g; struct aos_urpc_meta *meta; + struct capref ret_cap; struct thread_sem sem; }; @@ -51,4 +52,6 @@ int urpc_client_loop(void *arg); int urpc_server(void *arg); +void urpc_server_async_reply(struct generic_rpc_server *g_rpc, errval_t ret_err); + #endif // _LIB_BARRELFISH_AOS_URPC_H diff --git a/include/aos/ump_binding.h b/include/aos/ump_binding.h new file mode 100644 index 0000000..1f89119 --- /dev/null +++ b/include/aos/ump_binding.h @@ -0,0 +1,29 @@ +#ifndef _LIB_BARRELFISH_UMP_BINDUNG_H +#define _LIB_BARRELFISH_UMP_BINDUNG_H + +#include + +#define UMP_FRAME_SIZE BASE_PAGE_SIZE + +enum ump_server_id { + UMP_SERVER_ECHO, + UMP_SERVER_COUNT // How many servers exist +}; + +typedef errval_t (*ump_binding_connect_callback_t)(void *arg, struct capref cap); + +// fields are private +struct ump_binding_server { + struct lmp_chan chan; + ump_binding_connect_callback_t connect_callback; + void *connect_callback_arg; + errval_t ret_err; +}; + +errval_t ump_binding_register ( + struct ump_binding_server *server, + enum ump_server_id server_id, + ump_binding_connect_callback_t cb, void *cb_arg +); + +#endif // _LIB_BARRELFISH_UMP_BINDUNG_H diff --git a/include/spawn/ump_binding_server.h b/include/spawn/ump_binding_server.h new file mode 100644 index 0000000..73adebb --- /dev/null +++ b/include/spawn/ump_binding_server.h @@ -0,0 +1,18 @@ +#ifndef _UMP_BINDING_SERVER_H_ +#define _UMP_BINDING_SERVER_H_ + +#include + +errval_t ump_binding_register_service_global (enum ump_server_id server_id, coreid_t coreid); +errval_t ump_binding_get_global (enum ump_server_id server_id, coreid_t *coreid); + +errval_t ump_binding_register_service ( + enum ump_server_id server_id, struct capref accept_ep_cap, struct capref *reply_ep_cap +); + +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 +); + +#endif /* _UMP_BINDING_SERVER_H_ */ diff --git a/lib/aos/Hakefile b/lib/aos/Hakefile index 0f9e870..fe91d60 100644 --- a/lib/aos/Hakefile +++ b/lib/aos/Hakefile @@ -25,6 +25,7 @@ "slot_alloc/twolevel_slot_alloc.c", "aos_rpc.c", "aos_urpc.c", + "ump_binding.c", "performance.c", "capabilities.c", "coreset.c", diff --git a/lib/aos/aos_rpc.c b/lib/aos/aos_rpc.c index b4f2371..12ac19e 100644 --- a/lib/aos/aos_rpc.c +++ b/lib/aos/aos_rpc.c @@ -14,6 +14,7 @@ #include #include +#include #include @@ -362,6 +363,24 @@ aos_rpc_process_get_all_pids(struct aos_rpc *rpc, domainid_t **pids, return err; } + +errval_t aos_rpc_ump_register(struct aos_rpc *rpc, enum ump_server_id server_id, struct capref accept_ep_cap, struct capref *reply_ep_cap) { + return do_aos_rpc( + rpc, RPC_MTYPE_UMP_REGISTER, + accept_ep_cap, 0, server_id, 0, + reply_ep_cap, NULL, NULL, NULL + ); +} + +errval_t aos_rpc_ump_connect(struct aos_rpc *rpc, enum ump_server_id server_id, struct capref *ret_cap) { + return do_aos_rpc( + rpc, RPC_MTYPE_UMP_CONNECT, + NULL_CAP, 0, server_id, 0, + ret_cap, NULL, NULL, NULL + ); +} + + // We are allowed to change the signature for this one errval_t aos_rpc_init(struct aos_rpc *rpc, struct lmp_chan *chan, void *shared_mem) { diff --git a/lib/aos/aos_urpc.c b/lib/aos/aos_urpc.c index 0dcfa3d..f64291c 100644 --- a/lib/aos/aos_urpc.c +++ b/lib/aos/aos_urpc.c @@ -75,8 +75,13 @@ errval_t do_aos_urpc( if(ret_cap != NULL) { if(rpc->meta->cap.type != ObjType_Null) { - assert(rpc->meta->cap.type == ObjType_RAM); - err = ram_forge(*ret_cap, rpc->meta->cap.u.ram.base, rpc->meta->cap.u.ram.bytes, my_core_id); + if (rpc->meta->cap.type == ObjType_RAM) { + err = ram_forge(*ret_cap, rpc->meta->cap.u.ram.base, rpc->meta->cap.u.ram.bytes, my_core_id); + } else if (rpc->meta->cap.type == ObjType_Frame) { + err = frame_forge(*ret_cap, rpc->meta->cap.u.frame.base, rpc->meta->cap.u.frame.bytes, my_core_id); + } else { + err = LIB_ERR_URPC_BAD_CAPTYPE; + } if (err_is_fail(err)) return err; } else { *ret_cap = NULL_CAP; @@ -99,6 +104,7 @@ int urpc_client_loop(void *arg) { } } +static void urpc_server_send_reply(struct aos_urpc_server *urpc, errval_t err); static void urpc_server_handler(void *arg) { struct aos_urpc_server *urpc = arg; @@ -109,7 +115,7 @@ static void urpc_server_handler(void *arg) { uintptr_t arg0 = urpc->meta->a2; uintptr_t arg1 = urpc->meta->a3; - struct capref ret_cap = NULL_CAP; + urpc->ret_cap = NULL_CAP; urpc->meta->a1 = 0; urpc->meta->a2 = 0; urpc->meta->a3 = 0; @@ -125,13 +131,19 @@ static void urpc_server_handler(void *arg) { err = handler( &urpc->g, NULL_CAP, arg_size, arg0, arg1, - &ret_cap, &urpc->meta->a1, &urpc->meta->a2, &urpc->meta->a3 + &urpc->ret_cap, &urpc->meta->a1, &urpc->meta->a2, &urpc->meta->a3 ); } } + + if (err == AOS_ERR_RPC_ASYNC_REPLY) return; + urpc_server_send_reply(urpc, err); +} + +static void urpc_server_send_reply(struct aos_urpc_server *urpc, errval_t err) { urpc->meta->a0 = err; - if (!capref_is_null(ret_cap)) { - err = cap_direct_identify(ret_cap, &urpc->meta->cap); + if (!capref_is_null(urpc->ret_cap)) { + err = cap_direct_identify(urpc->ret_cap, &urpc->meta->cap); if (err_is_fail(err)) { DEBUG_ERR(err, "in cap_direct_identify while handling URPC"); abort(); @@ -158,6 +170,12 @@ static void urpc_server_handler(void *arg) { thread_sem_post(&urpc->sem); } +void urpc_server_async_reply(struct generic_rpc_server *g_rpc, errval_t ret_err) { + struct aos_urpc_server *urpc = (struct aos_urpc_server *)g_rpc; + urpc_server_send_reply(urpc, ret_err); +} + + errval_t aos_urpc_get_bootinfo(struct aos_urpc * rpc, struct bootinfo_serialized ** ret) { size_t ret_len; @@ -214,7 +232,7 @@ int urpc_server(void *arg) { perf_add_now(&p, "triggered_closure"); } #endif - + // wait until the rpc is handled thread_sem_wait(&urpc->sem); diff --git a/lib/aos/ump_binding.c b/lib/aos/ump_binding.c new file mode 100644 index 0000000..949ac41 --- /dev/null +++ b/lib/aos/ump_binding.c @@ -0,0 +1,85 @@ +#include +#include +#include + + +static void ump_server_register_connect (struct ump_binding_server *server); +static void ump_server_handle_connect (void *arg); +static void ump_server_send_reply(void *arg); + + +errval_t ump_binding_register ( + struct ump_binding_server *server, + enum ump_server_id server_id, + ump_binding_connect_callback_t cb, void *cb_arg +) { + errval_t err; + + server->connect_callback = cb; + server->connect_callback_arg = cb_arg; + + lmp_chan_init(&server->chan); + err = endpoint_create(DEFAULT_LMP_BUF_WORDS, &server->chan.local_cap, &server->chan.endpoint); + if (err_is_fail(err)) return err_push(err, LIB_ERR_ENDPOINT_CREATE); + + err = lmp_chan_alloc_recv_slot(&server->chan); + if (err_is_fail(err)) return err_push(err, LIB_ERR_LMP_ALLOC_RECV_SLOT); + + struct aos_rpc *rpc = aos_rpc_get_init_channel(); + err = aos_rpc_ump_register(rpc, server_id, server->chan.local_cap, &server->chan.remote_cap); + if (err_is_fail(err)) return err_push(err, LIB_ERR_UMP_REGISTER); + + ump_server_register_connect(server); + + return SYS_ERR_OK; +} + +static void ump_server_register_connect (struct ump_binding_server *server) { + errval_t err = lmp_chan_register_recv(&server->chan, get_default_waitset(), MKCLOSURE(ump_server_handle_connect, server)); + if (err_is_fail(err)) { + DEBUG_ERR(err, "Could not register receive handler"); + } +} + +static void ump_server_handle_connect (void *arg) { + errval_t err; + struct ump_binding_server *server = (struct ump_binding_server *)arg; + + struct lmp_recv_msg msg = LMP_RECV_MSG_INIT; + struct capref arg_cap; + + err = lmp_chan_recv(&server->chan, &msg, &arg_cap); + assert(err_is_ok(err)); + + err = lmp_chan_alloc_recv_slot(&server->chan); + if (err_is_fail(err)) DEBUG_ERR(err, "Failed to allocate recv slot"); + + server->ret_err = server->connect_callback(server->connect_callback_arg, arg_cap); + + ump_server_send_reply(server); +} + +static void ump_server_send_reply(void *arg) { + struct ump_binding_server *server = (struct ump_binding_server *)arg; + errval_t err; + + err = lmp_chan_send4( + &server->chan, LMP_FLAG_YIELD | LMP_FLAG_SYNC, + NULL_CAP, server->ret_err, 0, 0, 0 + ); + if (err_is_fail(err)) { + if (!lmp_err_is_transient(err)) { + DEBUG_ERR(err, "Could not send connect reply"); + return; + } + // Cannot send right now, try again later + err = lmp_chan_register_send(&server->chan, get_default_waitset(), MKCLOSURE(ump_server_send_reply, arg)); + if (err_is_fail(err)) { + DEBUG_ERR(err, "Could not register send handler"); + return; + } + } + + // Sent reply, wait for next connect call + ump_server_register_connect(server); +} diff --git a/lib/grading/rpc.c b/lib/grading/rpc.c index b869090..f6cd588 100644 --- a/lib/grading/rpc.c +++ b/lib/grading/rpc.c @@ -28,7 +28,7 @@ void grading_rpc_handler_serial_putchar(char c) void grading_rpc_handler_ram_cap(size_t bytes, size_t alignment) { - debug_printf("grading_rpc_handler_ram_cap(0x%"PRIxPTR", 0x%"PRIxPTR")\n", bytes, alignment); + // debug_printf("grading_rpc_handler_ram_cap(0x%"PRIxPTR", 0x%"PRIxPTR")\n", bytes, alignment); } void grading_rpc_handler_process_spawn(char* cmdline, coreid_t core) diff --git a/lib/spawn/Hakefile b/lib/spawn/Hakefile index a550f90..357889c 100644 --- a/lib/spawn/Hakefile +++ b/lib/spawn/Hakefile @@ -13,7 +13,7 @@ [ build library { target = "spawn", - cFiles = [ "spawn.c", "rpc_server.c" ], + cFiles = [ "spawn.c", "rpc_server.c", "ump_binding_server.c" ], addLibraries = [ "elf", "argv", "multiboot" ] }, build library { diff --git a/lib/spawn/rpc_server.c b/lib/spawn/rpc_server.c index 9bed700..5cfd525 100644 --- a/lib/spawn/rpc_server.c +++ b/lib/spawn/rpc_server.c @@ -3,6 +3,7 @@ #include #include #include +#include #include #include @@ -41,6 +42,13 @@ static void rpc_server_handle_recv(void *arg) { err = lmp_chan_recv(rpc->chan, &msg, &arg_cap); assert(err_is_ok(err)); + if (!capref_is_null(arg_cap)) { + err = lmp_chan_alloc_recv_slot(rpc->chan); + if (err_is_fail(err)) { + DEBUG_ERR(err, "Failed to allocate recv slot"); + } + } + rpc->ret_cap = NULL_CAP; rpc->ret_size = 0; rpc->ret0 = 0; @@ -54,8 +62,12 @@ static void rpc_server_handle_recv(void *arg) { if (handler == NULL) { err = AOS_ERR_RPC_UNKNOWN_MSG_TYPE; } else { - if (my_core_id == 0 || - (msg.words[0] == RPC_MTYPE_PROCESS_SPAWN && msg.words[1] == my_core_id)/* dont forward spawn rpcs */) { + if ( + my_core_id == 0 || + (msg.words[0] == RPC_MTYPE_PROCESS_SPAWN && msg.words[1] == my_core_id) || // dont forward spawn rpcs + msg.words[0] == RPC_MTYPE_UMP_REGISTER || + msg.words[0] == RPC_MTYPE_UMP_CONNECT + ) { err = handler( &rpc->g, arg_cap, msg.words[1], msg.words[2], msg.words[3], @@ -105,8 +117,8 @@ static void rpc_server_send_reply(void *arg) { err = lmp_chan_register_send(rpc->chan, get_default_waitset(), MKCLOSURE(rpc_server_send_reply, arg)); if (err_is_fail(err)) { DEBUG_ERR(err, "Could not register send handler"); - return; } + return; } // Sent reply, wait for next RPC call @@ -440,6 +452,52 @@ static errval_t handle_rpc_get_bootinfo( return SYS_ERR_OK; } +static errval_t handle_rpc_ump_register( + struct generic_rpc_server *rpc, + struct capref arg_cap, size_t arg_size, uintptr_t arg0, uintptr_t arg1, + struct capref *ret_cap, size_t *ret_size, uintptr_t *ret0, uintptr_t *ret1 +) { + enum ump_server_id server_id = arg0; + + return ump_binding_register_service(server_id, arg_cap, ret_cap); +} + +static errval_t handle_rpc_ump_connect( + struct generic_rpc_server *rpc, + struct capref arg_cap, size_t arg_size, uintptr_t arg0, uintptr_t arg1, + struct capref *ret_cap, size_t *ret_size, uintptr_t *ret0, uintptr_t *ret1 +) { + enum ump_server_id server_id = arg0; + + return ump_binding_connect(server_id, ret_cap, (void (*)(void *, errval_t))rpc->async_reply, rpc); +} + +static errval_t handle_rpc_ump_register_global( + struct generic_rpc_server *rpc, + struct capref arg_cap, size_t arg_size, uintptr_t arg0, uintptr_t arg1, + struct capref *ret_cap, size_t *ret_size, uintptr_t *ret0, uintptr_t *ret1 +) { + enum ump_server_id server_id = arg0; + coreid_t coreid = arg1; + + return ump_binding_register_service_global(server_id, coreid); +} + +static errval_t handle_rpc_ump_get_global( + struct generic_rpc_server *rpc, + struct capref arg_cap, size_t arg_size, uintptr_t arg0, uintptr_t arg1, + struct capref *ret_cap, size_t *ret_size, uintptr_t *ret0, uintptr_t *ret1 +) { + errval_t err; + enum ump_server_id server_id = arg0; + + coreid_t coreid; + err = ump_binding_get_global(server_id, &coreid); + if (err_is_fail(err)) return err; + *ret0 = coreid; + return SYS_ERR_OK; +} + rpc_handler_t rpc_handlers[RPC_MTYPE_COUNT] = { [RPC_MTYPE_SEND_NUMBER] = handle_rpc_send_number, [RPC_MTYPE_SEND_STRING] = handle_rpc_send_string, @@ -454,4 +512,8 @@ rpc_handler_t rpc_handlers[RPC_MTYPE_COUNT] = { [RPC_MTYPE_ALLOCATE_PID] = handle_rpc_allocate_pid, [RPC_MTYPE_NOP] = handle_rpc_nop, [RPC_MTYPE_GET_BOOTINFO] = handle_rpc_get_bootinfo, + [RPC_MTYPE_UMP_REGISTER] = handle_rpc_ump_register, + [RPC_MTYPE_UMP_CONNECT] = handle_rpc_ump_connect, + [RPC_MTYPE_UMP_REGISTER_GLOBAL] = handle_rpc_ump_register_global, + [RPC_MTYPE_UMP_GET_GLOBAL] = handle_rpc_ump_get_global, }; diff --git a/lib/spawn/spawn.c b/lib/spawn/spawn.c index 95499e0..1c8d0db 100644 --- a/lib/spawn/spawn.c +++ b/lib/spawn/spawn.c @@ -97,6 +97,10 @@ static void handle_child_recv(void *arg) { case RPC_MTYPE_CHILD_ENDPOINT: debug_printf("SPAWN: Received child init_chan endpoint\n"); si->init_chan.remote_cap = cap; + err = lmp_chan_alloc_recv_slot(&si->init_chan); + if (err_is_fail(err)) { + USER_PANIC_ERR(err, "Failed to allocate recv slot"); + } err = lmp_chan_send1(&si->init_chan, LMP_FLAG_YIELD | LMP_FLAG_SYNC, NULL_CAP, SYS_ERR_OK); if (err_is_fail(err)) { USER_PANIC_ERR(err, "Failed to send reply"); @@ -675,7 +679,8 @@ errval_t spawn_load_argv(int argc, char *argv[], struct spawninfo *si, rpc_server_init(&si->rpc_server, &si->init_chan, rpc_shared_memory); // Receive the child's endpoint for the init channel - lmp_chan_alloc_recv_slot(&si->init_chan); + err = lmp_chan_alloc_recv_slot(&si->init_chan); + if (err_is_fail(err)) return err_push(err, LIB_ERR_LMP_ALLOC_RECV_SLOT); err = lmp_chan_register_recv(&si->init_chan, get_default_waitset(), MKCLOSURE(handle_child_recv, si)); if (err_is_fail(err)) return err; diff --git a/lib/spawn/ump_binding_server.c b/lib/spawn/ump_binding_server.c new file mode 100644 index 0000000..73c8729 --- /dev/null +++ b/lib/spawn/ump_binding_server.c @@ -0,0 +1,229 @@ +#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); + + 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); +} diff --git a/platforms/Hakefile b/platforms/Hakefile index 9df7578..92f54a6 100644 --- a/platforms/Hakefile +++ b/platforms/Hakefile @@ -12,7 +12,7 @@ let -- Default list of modules to build/install - modules_common = [ "/sbin/" ++ f | f <- [ "init", "hello", "memeater", "mallocator", "stackoverflow" + modules_common = [ "/sbin/" ++ f | f <- [ "init", "hello", "memeater", "mallocator", "stackoverflow", "echoserver" ] ] in [ diff --git a/usr/echoserver/Hakefile b/usr/echoserver/Hakefile new file mode 100644 index 0000000..be49d16 --- /dev/null +++ b/usr/echoserver/Hakefile @@ -0,0 +1,18 @@ +-------------------------------------------------------------------------- +-- Copyright (c) 2007-2010, ETH Zurich. +-- All rights reserved. +-- +-- This file is distributed under the terms in the attached LICENSE file. +-- If you do not find this file, copies can be found by writing to: +-- ETH Zurich D-INFK, Haldeneggsteig 4, CH-8092 Zurich. Attn: Systems Group. +-- +-- Hakefile for /usr/init +-- +-------------------------------------------------------------------------- + +[ build application + { + target = "echoserver", + cFiles = [ "main.c" ] + } +] diff --git a/usr/echoserver/main.c b/usr/echoserver/main.c new file mode 100644 index 0000000..ce1bf32 --- /dev/null +++ b/usr/echoserver/main.c @@ -0,0 +1,40 @@ +#include +#include +#include +#include +#include + +static errval_t connect_server (void *arg, struct capref cap) { + debug_printf("Echo server: incoming connection\n"); + + struct paging_state *pstate = get_current_paging_state(); + uint8_t *ump_data; + errval_t err = paging_map_frame(pstate, (void **)&ump_data, BASE_PAGE_SIZE, cap); + if (err_is_fail(err)) USER_PANIC_ERR(err, "Failed to map ump frame"); + + // TODO: Create an UMP channel + ump_data[0] = 123; + + return SYS_ERR_OK; +} + +int main (int argc, char *argv[]) { + errval_t err; + struct ump_binding_server server; + + debug_printf("Echo server registering\n"); + + err = ump_binding_register(&server, UMP_SERVER_ECHO, connect_server, NULL); + if (err_is_fail(err)) USER_PANIC_ERR(err, "Failed to register UMP server"); + + debug_printf("Echo server listening\n"); + + struct waitset *default_ws = get_default_waitset(); + while (true) { + err = event_dispatch(default_ws); + if (err_is_fail(err)) { + DEBUG_ERR(err, "in event_dispatch"); + abort(); + } + } +} diff --git a/usr/hello/hello.c b/usr/hello/hello.c index 5c16eb0..30877f8 100644 --- a/usr/hello/hello.c +++ b/usr/hello/hello.c @@ -21,9 +21,10 @@ #include #include -#define HELLO_CATCH_COMMAND "catch" -#define HELLO_SPAWN_COMMAND "spawn" -#define HELLO_SPAWN_COMMAND_REMOTE "spawn_remote" +#define HELLO_CMD_CATCH "catch" +#define HELLO_CMD_SPAWN "spawn" +#define HELLO_CMD_SPAWN_REMOTE "spawn_remote" +#define HELLO_CMD_ECHOCLIENT "echoclient" #define HELLO_CMDLINE_READ_LEN 100 @@ -57,7 +58,7 @@ int main(int argc, char *argv[]) if (err_is_fail(err)) USER_PANIC_ERR(err, "Failed to send RPC"); // try to spawn a child using rpc - if (argc > 1 && !strncmp(argv[1], HELLO_SPAWN_COMMAND, sizeof(HELLO_SPAWN_COMMAND))) { + if (argc > 1 && !strncmp(argv[1], HELLO_CMD_SPAWN, sizeof(HELLO_CMD_SPAWN))) { debug_printf("Waiting for 1s before showing shell because kernel getchar blocks everything!\n"); barrelfish_usleep(1000000); @@ -112,7 +113,7 @@ int main(int argc, char *argv[]) } // if we receive the catch command, RUN! - if (argc > 1 && !strncmp(argv[1], HELLO_CATCH_COMMAND, sizeof(HELLO_CATCH_COMMAND))) { + if (argc > 1 && !strncmp(argv[1], HELLO_CMD_CATCH, sizeof(HELLO_CMD_CATCH))) { int iter = 0; while(1) { printf("Catch me if you can! %d\n", iter); @@ -122,7 +123,7 @@ int main(int argc, char *argv[]) } } - if (argc > 1 && !strncmp(argv[1], HELLO_SPAWN_COMMAND_REMOTE, sizeof(HELLO_SPAWN_COMMAND_REMOTE))) { + if (argc > 1 && !strncmp(argv[1], HELLO_CMD_SPAWN_REMOTE, sizeof(HELLO_CMD_SPAWN_REMOTE))) { printf("Spawning process on core 1...\n"); domainid_t pid; rpc = aos_rpc_get_process_channel(); @@ -131,5 +132,23 @@ int main(int argc, char *argv[]) printf("Spawn succeeded, pid=%"PRIuDOMAINID"\n", pid); } + if (argc > 1 && !strncmp(argv[1], HELLO_CMD_ECHOCLIENT, sizeof(HELLO_CMD_ECHOCLIENT))) { + debug_printf("Connecting to echo server...\n"); + struct capref echo_cap; + for (int attempts = 0; attempts < 50; attempts++) { + err = aos_rpc_ump_connect(rpc, UMP_SERVER_ECHO, &echo_cap); + if (err != LIB_ERR_UMP_NOT_REGISTERED) break; + barrelfish_usleep(100000); + } + if (err_is_fail(err)) USER_PANIC_ERR(err, "Failed to connect to echo server"); + debug_printf("Connected to echo server.\n"); + + struct paging_state *pstate = get_current_paging_state(); + uint8_t *ump_data; + err = paging_map_frame(pstate, (void **)&ump_data, BASE_PAGE_SIZE, echo_cap); + if (err_is_fail(err)) USER_PANIC_ERR(err, "Failed to map ump frame"); + debug_printf("Read byte: %d\n", ump_data[0]); + } + return EXIT_SUCCESS; } diff --git a/usr/init/main.c b/usr/init/main.c index 43bccd3..cf1ae87 100644 --- a/usr/init/main.c +++ b/usr/init/main.c @@ -22,6 +22,7 @@ #include #include #include +#include #include #include "mem_alloc.h" @@ -136,7 +137,7 @@ bsp_main(int argc, char *argv[]) { struct aos_urpc_server *urpc_to_bsp_server = malloc(sizeof(struct aos_urpc_server)); urpc_to_bsp_server->g.shared_mem = &urpc->shared_mem_to_bsp; - urpc_to_bsp_server->g.async_reply = NULL; + urpc_to_bsp_server->g.async_reply = urpc_server_async_reply; urpc_to_bsp_server->meta = &urpc->meta_to_bsp; thread_create(urpc_server, urpc_to_bsp_server); @@ -243,7 +244,7 @@ app_main(int argc, char *argv[]) { urpc_to_bsp.meta = &urpc->meta_to_bsp; urpc_to_app_server.g.shared_mem = &urpc->shared_mem_to_app; - urpc_to_app_server.g.async_reply = NULL; + urpc_to_app_server.g.async_reply = urpc_server_async_reply; urpc_to_app_server.meta = &urpc->meta_to_app; ram_alloc_set(ram_alloc_remote_core); @@ -301,6 +302,14 @@ app_main(int argc, char *argv[]) { // Grading grading_test_early(); + // Spawn echoserver + struct spawninfo echoserver_si; + domainid_t echoserver_pid; + err = spawn_load_by_name("echoserver", &echoserver_si, &echoserver_pid); + if (err_is_fail(err)) { + DEBUG_ERR(err, "when spawning echoserver"); + } + // Grading grading_test_late();