add support for timeouts on synchronous requests
authorFelix Fietkau <nbd@openwrt.org>
Fri, 11 Feb 2011 01:40:39 +0000 (02:40 +0100)
committerFelix Fietkau <nbd@openwrt.org>
Fri, 11 Feb 2011 01:40:39 +0000 (02:40 +0100)
cli.c
libubus.c
libubus.h
ubusmsg.h

diff --git a/cli.c b/cli.c
index d020b8e652cbaea953dcb92ea5d6c0b682c5ccd3..76fb0c424703bd66ed4616938a42bb9cda4aa3db 100644 (file)
--- a/cli.c
+++ b/cli.c
@@ -4,6 +4,7 @@
 #include "libubus.h"
 
 static struct blob_buf b;
+static int timeout = 30;
 
 static const char *format_type(void *priv, struct blob_attr *attr)
 {
@@ -96,7 +97,7 @@ static int ubus_cli_call(struct ubus_context *ctx, int argc, char **argv)
        if (ret)
                return ret;
 
-       return ubus_invoke(ctx, id, argv[1], b.head, receive_call_result_data, NULL);
+       return ubus_invoke(ctx, id, argv[1], b.head, receive_call_result_data, NULL, timeout * 1000);
 }
 
 static int ubus_cli_listen(struct ubus_context *ctx, int argc, char **argv)
@@ -194,11 +195,14 @@ int main(int argc, char **argv)
 
        progname = argv[0];
 
-       while ((ch = getopt(argc, argv, "s:")) != -1) {
+       while ((ch = getopt(argc, argv, "s:t:")) != -1) {
                switch (ch) {
                case 's':
                        ubus_socket = optarg;
                        break;
+               case 't':
+                       timeout = atoi(optarg);
+                       break;
                default:
                        return usage(progname);
                }
index 0a0df973b6276000f06a2950b0f38af0339c2758..d84e970c39022658ec6082658b517697088ac1b7 100644 (file)
--- a/libubus.c
+++ b/libubus.c
@@ -21,6 +21,7 @@ const char *__ubus_strerror[__UBUS_STATUS_LAST] = {
        [UBUS_STATUS_NOT_FOUND] = "Not found",
        [UBUS_STATUS_NO_DATA] = "No response",
        [UBUS_STATUS_PERMISSION_DENIED] = "Permission denied",
+       [UBUS_STATUS_TIMEOUT] = "Request timed out",
 };
 
 static struct blob_buf b;
@@ -157,10 +158,8 @@ static bool recv_retry(int fd, struct iovec *iov, bool wait)
                        if (errno == EINTR)
                                continue;
 
-                       if (errno != EAGAIN) {
-                               perror("read");
+                       if (errno != EAGAIN)
                                return false;
-                       }
                }
                if (!wait && !bytes)
                        return false;
@@ -187,13 +186,13 @@ static bool ubus_validate_hdr(struct ubus_msghdr *hdr)
        return true;
 }
 
-static bool get_next_msg(struct ubus_context *ctx, bool wait)
+static bool get_next_msg(struct ubus_context *ctx)
 {
        struct iovec iov = STATIC_IOV(ctx->msgbuf.hdr);
 
        /* receive header + start attribute */
        iov.iov_len += sizeof(struct blob_attr);
-       if (!recv_retry(ctx->sock.fd, &iov, wait))
+       if (!recv_retry(ctx->sock.fd, &iov, false))
                return false;
 
        iov.iov_len = blob_len(ctx->msgbuf.hdr.data);
@@ -409,47 +408,77 @@ static void ubus_handle_data(struct uloop_fd *u, unsigned int events)
        struct ubus_context *ctx = container_of(u, struct ubus_context, sock);
        struct ubus_msghdr *hdr = &ctx->msgbuf.hdr;
 
-       while (get_next_msg(ctx, false))
+       while (get_next_msg(ctx)) {
                ubus_process_msg(ctx, hdr);
+               if (uloop_cancelled)
+                       break;
+       }
 
        if (u->eof)
                ctx->connection_lost(ctx);
 }
 
-int ubus_complete_request(struct ubus_context *ctx, struct ubus_request *req)
+static void ubus_sync_req_cb(struct ubus_request *req, int ret)
 {
-       struct ubus_msghdr *hdr = &ctx->msgbuf.hdr;
-
-       if (!list_empty(&req->list))
-               list_del(&req->list);
+       req->status_msg = true;
+       req->status_code = ret;
+       uloop_end();
+}
 
-       while (1) {
-               if (req->status_msg)
-                       return req->status_code;
+struct ubus_sync_req_cb {
+       struct uloop_timeout timeout;
+       struct ubus_request *req;
+};
 
-               if (req->cancelled)
-                       return UBUS_STATUS_NO_DATA;
+static void ubus_sync_req_timeout_cb(struct uloop_timeout *timeout)
+{
+       struct ubus_sync_req_cb *cb;
 
-               if (!get_next_msg(ctx, true))
-                       return UBUS_STATUS_NO_DATA;
+       cb = container_of(timeout, struct ubus_sync_req_cb, timeout);
+       ubus_sync_req_cb(cb->req, UBUS_STATUS_TIMEOUT);
+}
 
-               if (hdr->seq != req->seq || hdr->peer != req->peer)
-                       goto skip;
+int ubus_complete_request(struct ubus_context *ctx, struct ubus_request *req,
+                         int timeout)
+{
+       struct ubus_sync_req_cb cb;
+       ubus_complete_handler_t complete_cb = req->complete_cb;
+       bool registered = ctx->sock.registered;
+       bool cancelled = uloop_cancelled;
+       int status = UBUS_STATUS_NO_DATA;
 
-               switch(hdr->type) {
-               case UBUS_MSG_STATUS:
-                       return ubus_process_req_status(req, hdr);
-               case UBUS_MSG_DATA:
-                       if (req->data_cb || req->raw_data_cb)
-                               ubus_req_data(req, hdr);
-                       continue;
-               default:
-                       goto skip;
-               }
+       if (!registered) {
+               uloop_init();
+               ubus_add_uloop(ctx);
+       }
 
-skip:
-               ubus_process_msg(ctx, hdr);
+       if (timeout) {
+               memset(&cb, 0, sizeof(cb));
+               cb.req = req;
+               cb.timeout.cb = ubus_sync_req_timeout_cb;
+               uloop_timeout_set(&cb.timeout, timeout);
        }
+
+       ubus_complete_request_async(ctx, req);
+       req->complete_cb = ubus_sync_req_cb;
+
+       uloop_run();
+
+       if (timeout)
+               uloop_timeout_cancel(&cb.timeout);
+
+       if (req->status_msg)
+               status = req->status_code;
+
+       req->complete_cb = complete_cb;
+       if (req->complete_cb)
+               req->complete_cb(req, status);
+
+       uloop_cancelled = cancelled;
+       if (!registered)
+               uloop_fd_delete(&ctx->sock);
+
+       return status;
 }
 
 struct ubus_lookup_request {
@@ -493,7 +522,7 @@ int ubus_lookup(struct ubus_context *ctx, const char *path,
        lookup.req.raw_data_cb = ubus_lookup_cb;
        lookup.req.priv = priv;
        lookup.cb = cb;
-       return ubus_complete_request(ctx, &lookup.req);
+       return ubus_complete_request(ctx, &lookup.req, 0);
 }
 
 static void ubus_lookup_id_cb(struct ubus_request *req, int type, struct blob_attr *msg)
@@ -523,7 +552,7 @@ int ubus_lookup_id(struct ubus_context *ctx, const char *path, uint32_t *id)
        req.raw_data_cb = ubus_lookup_id_cb;
        req.priv = id;
 
-       return ubus_complete_request(ctx, &req);
+       return ubus_complete_request(ctx, &req, 0);
 }
 
 int ubus_send_reply(struct ubus_context *ctx, struct ubus_request_data *req,
@@ -557,14 +586,15 @@ int ubus_invoke_async(struct ubus_context *ctx, uint32_t obj, const char *method
 }
 
 int ubus_invoke(struct ubus_context *ctx, uint32_t obj, const char *method,
-                struct blob_attr *msg, ubus_data_handler_t cb, void *priv)
+                struct blob_attr *msg, ubus_data_handler_t cb, void *priv,
+               int timeout)
 {
        struct ubus_request req;
 
        ubus_invoke_async(ctx, obj, method, msg, &req);
        req.data_cb = cb;
        req.priv = priv;
-       return ubus_complete_request(ctx, &req);
+       return ubus_complete_request(ctx, &req, timeout);
 }
 
 static void ubus_add_object_cb(struct ubus_request *req, int type, struct blob_attr *msg)
@@ -664,7 +694,7 @@ static int __ubus_add_object(struct ubus_context *ctx, struct ubus_object *obj)
 
        req.raw_data_cb = ubus_add_object_cb;
        req.priv = obj;
-       ret = ubus_complete_request(ctx, &req);
+       ret = ubus_complete_request(ctx, &req, 0);
        if (ret)
                return ret;
 
@@ -712,7 +742,7 @@ int ubus_remove_object(struct ubus_context *ctx, struct ubus_object *obj)
 
        req.raw_data_cb = ubus_remove_object_cb;
        req.priv = obj;
-       ret = ubus_complete_request(ctx, &req);
+       ret = ubus_complete_request(ctx, &req, 0);
        if (ret)
                return ret;
 
@@ -766,7 +796,7 @@ int ubus_register_event_handler(struct ubus_context *ctx,
                blobmsg_add_string(&b2, "pattern", pattern);
 
        return ubus_invoke(ctx, UBUS_SYSTEM_OBJECT_EVENT, "register", b2.head,
-                         NULL, NULL);
+                         NULL, NULL, 0);
 }
 
 int ubus_send_event(struct ubus_context *ctx, const char *id,
@@ -786,7 +816,7 @@ int ubus_send_event(struct ubus_context *ctx, const char *id,
        if (ubus_start_request(ctx, &req, b.head, UBUS_MSG_INVOKE, UBUS_SYSTEM_OBJECT_EVENT) < 0)
                return UBUS_STATUS_INVALID_ARGUMENT;
 
-       return ubus_complete_request(ctx, &req);
+       return ubus_complete_request(ctx, &req, 0);
 }
 
 static void ubus_default_connection_lost(struct ubus_context *ctx)
index 9e78ad43ebaf03efa52b5bbcaa86de981578d644..8b5085af788a76d08960692059e4f805e2ea4bb0 100644 (file)
--- a/libubus.h
+++ b/libubus.h
@@ -154,7 +154,8 @@ static inline void ubus_handle_event(struct ubus_context *ctx)
 /* ----------- raw request handling ----------- */
 
 /* wait for a request to complete and return its status */
-int ubus_complete_request(struct ubus_context *ctx, struct ubus_request *req);
+int ubus_complete_request(struct ubus_context *ctx, struct ubus_request *req,
+                         int timeout);
 
 /* complete a request asynchronously */
 void ubus_complete_request_async(struct ubus_context *ctx,
@@ -180,11 +181,12 @@ int ubus_remove_object(struct ubus_context *ctx, struct ubus_object *obj);
 
 /* invoke a method on a specific object */
 int ubus_invoke(struct ubus_context *ctx, uint32_t obj, const char *method,
-                struct blob_attr *msg, ubus_data_handler_t cb, void *priv);
+               struct blob_attr *msg, ubus_data_handler_t cb, void *priv,
+               int timeout);
 
 /* asynchronous version of ubus_invoke() */
 int ubus_invoke_async(struct ubus_context *ctx, uint32_t obj, const char *method,
-                      struct blob_attr *msg, struct ubus_request *req);
+                     struct blob_attr *msg, struct ubus_request *req);
 
 /* send a reply to an incoming object method call */
 int ubus_send_reply(struct ubus_context *ctx, struct ubus_request_data *req,
index a34b341c3534aab33df7f25cedbcbbabddf42b07..00bf125c544ff3363aafb5423b4d5abb8144e847 100644 (file)
--- a/ubusmsg.h
+++ b/ubusmsg.h
@@ -71,6 +71,7 @@ enum ubus_msg_status {
        UBUS_STATUS_NOT_FOUND,
        UBUS_STATUS_NO_DATA,
        UBUS_STATUS_PERMISSION_DENIED,
+       UBUS_STATUS_TIMEOUT,
        __UBUS_STATUS_LAST
 };