summaryrefslogtreecommitdiff
path: root/qpid/extras/dispatch/src/router_node.c
diff options
context:
space:
mode:
authorTed Ross <tross@apache.org>2013-04-26 16:34:33 +0000
committerTed Ross <tross@apache.org>2013-04-26 16:34:33 +0000
commit64ab8be9c34528ef71ca5c58ff075ed57a48c9e0 (patch)
tree82c66d38e2d710d780b22c71d5afd6e00297082e /qpid/extras/dispatch/src/router_node.c
parentcf474845241b0c206710dbd9c87fe2e752c512a0 (diff)
downloadqpid-python-64ab8be9c34528ef71ca5c58ff075ed57a48c9e0.tar.gz
NO-JIRA - Development update to Dispatch Router
- Began refactoring of the routing table to support in-process destinations and multi-hop paths. - Added API for the internal management agent. Began integrating the agent with the router module for communication. - Added field parsing to handle topological addresses. - Added tests. git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1476280 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/extras/dispatch/src/router_node.c')
-rw-r--r--qpid/extras/dispatch/src/router_node.c190
1 files changed, 165 insertions, 25 deletions
diff --git a/qpid/extras/dispatch/src/router_node.c b/qpid/extras/dispatch/src/router_node.c
index 0513b08a6b..65756be215 100644
--- a/qpid/extras/dispatch/src/router_node.c
+++ b/qpid/extras/dispatch/src/router_node.c
@@ -18,13 +18,34 @@
*/
#include <stdio.h>
+#include <string.h>
#include <qpid/dispatch.h>
#include "dispatch_private.h"
-static char *module="ROUTER_NODE";
+static char *module = "ROUTER";
+
+//static char *local_prefix = "_local/";
+//static char *topo_prefix = "_topo/";
+
+/**
+ * Address Types and Processing:
+ *
+ * Address Hash Compare onReceive onEmit
+ * =============================================================================
+ * _local/<local> L<local> handler forward
+ * _topo/<area>/<router>/<local> A<area> forward forward
+ * _topo/<my-area>/<router>/<local> R<router> forward forward
+ * _topo/<my-area>/<my-router>/<local> L<local> forward+handler forward
+ * _topo/<area>/all/<local> A<area> forward forward
+ * _topo/<my-area>/all/<local> L<local> forward+handler forward
+ * _topo/all/all/<local> L<local> forward+handler forward
+ * <mobile> M<mobile> forward+handler forward
+ */
struct dx_router_t {
dx_dispatch_t *dx;
+ const char *router_area;
+ const char *router_id;
dx_node_t *node;
dx_link_list_t in_links;
dx_link_list_t out_links;
@@ -41,11 +62,32 @@ typedef struct {
dx_message_list_t out_fifo;
} dx_router_link_t;
-
ALLOC_DECLARE(dx_router_link_t);
ALLOC_DEFINE(dx_router_link_t);
+typedef struct {
+ const char *id;
+ dx_router_link_t *next_hop;
+ // list of valid origins (pointers to router_node) - (bit masks?)
+} dx_router_node_t;
+
+ALLOC_DECLARE(dx_router_node_t);
+ALLOC_DEFINE(dx_router_node_t);
+
+
+struct dx_address_t {
+ int is_local;
+ dx_router_message_cb handler; // In-Process Consumer
+ void *handler_context;
+ dx_router_link_t *rlink; // Locally-Connected Consumer - TODO: Make this a list
+ dx_router_node_t *rnode; // Remotely-Connected Consumer - TODO: Make this a list
+};
+
+ALLOC_DECLARE(dx_address_t);
+ALLOC_DEFINE(dx_address_t);
+
+
/**
* Outbound Delivery Handler
*/
@@ -119,22 +161,42 @@ static void router_rx_handler(void* context, dx_link_t *link, pn_delivery_t *del
if (valid_message) {
dx_field_iterator_t *iter = dx_message_field_iterator(msg, DX_FIELD_TO);
- dx_router_link_t *rlink;
+ dx_address_t *addr;
if (iter) {
- dx_field_iterator_reset(iter, ITER_VIEW_NO_HOST);
+ dx_field_iterator_reset_view(iter, ITER_VIEW_ADDRESS_HASH);
sys_mutex_lock(router->lock);
- int result = hash_retrieve(router->out_hash, iter, (void*) &rlink);
+ hash_retrieve(router->out_hash, iter, (void*) &addr);
dx_field_iterator_free(iter);
- if (result == 0) {
+ if (addr) {
+ //
+ // To field is valid and contains a known destination. Handle the various
+ // cases for forwarding.
+ //
+ // Forward to the in-process handler for this message if there is one.
+ //
+ if (addr->handler)
+ addr->handler(addr->handler_context, msg);
+
//
- // To field is valid and contains a known destination. Enqueue on
- // the output fifo for the next-hop-to-destination.
+ // Forward to the local link for the locally-connected consumer, if present.
+ // TODO - Don't forward if this is a "_local" address.
//
- pn_link_t* pn_outlink = dx_link_pn(rlink->link);
- DEQ_INSERT_TAIL(rlink->out_fifo, msg);
- pn_link_offered(pn_outlink, DEQ_SIZE(rlink->out_fifo));
- dx_link_activate(rlink->link);
+ if (addr->rlink) {
+ pn_link_t* pn_outlink = dx_link_pn(addr->rlink->link);
+ DEQ_INSERT_TAIL(addr->rlink->out_fifo, msg);
+ pn_link_offered(pn_outlink, DEQ_SIZE(addr->rlink->out_fifo));
+ dx_link_activate(addr->rlink->link);
+ }
+
+ //
+ // Forward to the next-hop for a remotely-connected consumer, if present.
+ // Don't forward if this is a "_local" address.
+ //
+ if (addr->rnode) {
+ // TODO
+ }
+
} else {
//
// To field contains an unknown address. Release the message.
@@ -242,17 +304,30 @@ static int router_outgoing_link_handler(void* context, dx_link_t *link)
pn_link_t *pn_link = dx_link_pn(link);
const char *r_tgt = pn_terminus_get_address(pn_link_remote_target(pn_link));
- sys_mutex_lock(router->lock);
dx_router_link_t *rlink = new_dx_router_link_t();
rlink->link = link;
DEQ_INIT(rlink->out_fifo);
dx_link_set_context(link, rlink);
- dx_field_iterator_t *iter = dx_field_iterator_string(r_tgt, ITER_VIEW_NO_HOST);
- int result = hash_insert(router->out_hash, iter, rlink);
+ dx_address_t *addr;
+
+ dx_field_iterator_t *iter = dx_field_iterator_string(r_tgt, ITER_VIEW_ADDRESS_HASH);
+
+ sys_mutex_lock(router->lock);
+ hash_retrieve(router->out_hash, iter, (void**) &addr);
+ if (!addr) {
+ addr = new_dx_address_t();
+ addr->is_local = 0;
+ addr->handler = 0;
+ addr->handler_context = 0;
+ addr->rlink = 0;
+ addr->rnode = 0;
+ hash_insert(router->out_hash, iter, addr);
+ }
dx_field_iterator_free(iter);
- if (result == 0) {
+ if (addr->rlink == 0) {
+ addr->rlink = rlink;
pn_terminus_copy(pn_link_source(pn_link), pn_link_remote_source(pn_link));
pn_terminus_copy(pn_link_target(pn_link), pn_link_remote_target(pn_link));
pn_link_open(pn_link);
@@ -314,14 +389,14 @@ static int router_link_detach_handler(void* context, dx_link_t *link, int closed
if (pn_link_is_sender(pn_link)) {
item = DEQ_HEAD(router->out_links);
- dx_field_iterator_t *iter = dx_field_iterator_string(r_tgt, ITER_VIEW_NO_HOST);
- dx_router_link_t *rlink;
+ dx_field_iterator_t *iter = dx_field_iterator_string(r_tgt, ITER_VIEW_ADDRESS_HASH);
+ dx_address_t *addr;
if (iter) {
- int result = hash_retrieve(router->out_hash, iter, (void*) &rlink);
- if (result == 0) {
- dx_field_iterator_reset(iter, ITER_VIEW_NO_HOST);
+ hash_retrieve(router->out_hash, iter, (void**) &addr);
+ if (addr) {
hash_remove(router->out_hash, iter);
- free_dx_router_link_t(rlink);
+ free_dx_router_link_t(addr->rlink);
+ free_dx_address_t(addr);
dx_log(module, LOG_TRACE, "Removed local address: %s", r_tgt);
}
dx_field_iterator_free(iter);
@@ -383,7 +458,7 @@ static dx_node_type_t router_node = {"router", 0, 0,
static int type_registered = 0;
-dx_router_t *dx_router(dx_dispatch_t *dx)
+dx_router_t *dx_router(dx_dispatch_t *dx, const char *area, const char *id)
{
if (!type_registered) {
type_registered = 1;
@@ -397,8 +472,10 @@ dx_router_t *dx_router(dx_dispatch_t *dx)
DEQ_INIT(router->out_links);
DEQ_INIT(router->in_fifo);
- router->dx = dx;
- router->lock = sys_mutex();
+ router->dx = dx;
+ router->lock = sys_mutex();
+ router->router_area = area;
+ router->router_id = id;
router->timer = dx_timer(dx, dx_router_timer_handler, (void*) router);
dx_timer_schedule(router->timer, 0); // Immediate
@@ -406,10 +483,22 @@ dx_router_t *dx_router(dx_dispatch_t *dx)
router->out_hash = hash(10, 32, 0);
router->dtag = 1;
+ //
+ // Inform the field iterator module of this router's id and area. The field iterator
+ // uses this to offload some of the address-processing load from the router.
+ //
+ dx_field_iterator_set_address(area, id);
+
return router;
}
+void dx_router_setup_agent(dx_dispatch_t *dx)
+{
+ // TODO
+}
+
+
void dx_router_free(dx_router_t *router)
{
dx_container_set_default_node_type(router->dx, 0, 0, DX_DIST_BOTH);
@@ -417,3 +506,54 @@ void dx_router_free(dx_router_t *router)
free(router);
}
+
+dx_address_t *dx_router_register_address(dx_dispatch_t *dx,
+ bool is_local,
+ const char *address,
+ dx_router_message_cb handler,
+ void *context)
+{
+ char addr[1000];
+ dx_address_t *ad = new_dx_address_t();
+ dx_field_iterator_t *iter;
+ int result;
+
+ if (!ad)
+ return 0;
+
+ ad->is_local = is_local;
+ ad->handler = handler;
+ ad->handler_context = context;
+ ad->rlink = 0;
+
+ if (ad->is_local)
+ strcpy(addr, "L"); // Local Hash-Key Space
+ else
+ strcpy(addr, "M"); // Mobile Hash-Key Space
+
+ strcat(addr, address);
+ iter = dx_field_iterator_string(addr, ITER_VIEW_NO_HOST);
+ result = hash_insert(dx->router->out_hash, iter, ad);
+ dx_field_iterator_free(iter);
+ if (result != 0) {
+ free_dx_address_t(ad);
+ return 0;
+ }
+
+ dx_log(module, LOG_TRACE, "In-Process Address Registered: %s", address);
+ return ad;
+}
+
+
+void dx_router_unregister_address(dx_address_t *ad)
+{
+ free_dx_address_t(ad);
+}
+
+
+void dx_router_send(dx_dispatch_t *dx,
+ const char *address,
+ dx_message_t *msg)
+{
+}
+