[PATCH 15/31] Connection: simplify in-connection handeling (kill 1 function, and break shoutcast)

Niv Sardi <[email protected]>
Newsgroups gmane.comp.audio.icecast.devel
Message-ID <[email protected]>
Signed-off-by: Niv Sardi <[email protected]>
---
 src/connection.c |  395 ++++++++++++++++++------------------------------------
 1 files changed, 131 insertions(+), 264 deletions(-)

diff --git a/src/connection.c b/src/connection.c
index ed8bea9..5c9e96e 100644
--- a/src/connection.c
+++ b/src/connection.c
@@ -108,7 +108,6 @@ static spin_t _connection_lock;
 static volatile unsigned long _current_id = 0;
 static int _initialized = 0;
 
-static volatile client_queue_t *_req_queue = NULL, **_req_queue_tail = &_req_queue;
 static volatile client_queue_t *_con_queue = NULL, **_con_queue_tail = &_con_queue;
 static int ssl_ok;
 #ifdef HAVE_OPENSSL
@@ -120,7 +119,8 @@ cache_file_contents banned_ip, allowed_ip;
 
 rwlock_t _source_shutdown_rwlock;
 
-static void _handle_connection(void);
+static void _handle_shoutcast_compatible (int shoutcast, char *shoutcast_mount);
+static int _handle_client (client_t *client);
 
 static int compare_ip (void *arg, void *a, void *b)
 {
@@ -146,8 +146,6 @@ void connection_initialize(void)
     thread_mutex_create(&move_clients_mutex);
     thread_rwlock_create(&_source_shutdown_rwlock);
     thread_cond_init(&global.shutdown_cond);
-    _req_queue = NULL;
-    _req_queue_tail = &_req_queue;
     _con_queue = NULL;
     _con_queue_tail = &_con_queue;
 
@@ -553,100 +551,6 @@ static connection_t *_accept_connection(int duration)
     return NULL;
 }
 
-
-/* add client to connection queue. At this point some header information
- * has been collected, so we now pass it onto the connection thread for
- * further processing
- */
-static void _add_connection (client_queue_t *node)
-{
-    *_con_queue_tail = node;
-    _con_queue_tail = (volatile client_queue_t **)&node->next;
-}
-
-
-/* this returns queued clients for the connection thread. headers are
- * already provided, but need to be parsed.
- */
-static client_queue_t *_get_connection(void)
-{
-    client_queue_t *node = NULL;
-
-    /* common case, no new connections so don't bother taking locks */
-    if (_con_queue)
-    {
-        node = (client_queue_t *)_con_queue;
-        _con_queue = node->next;
-        if (_con_queue == NULL)
-            _con_queue_tail = &_con_queue;
-        node->next = NULL;
-    }
-    return node;
-}
-
-
-/* run along queue checking for any data that has come in or a timeout */
-static void process_request_queue (int timeout)
-{
-    client_queue_t **node_ref = (client_queue_t **)&_req_queue;
-
-    while (*node_ref)
-    {
-        client_queue_t *node = *node_ref;
-        client_t *client = node->client;
-        int len = PER_CLIENT_REFBUF_SIZE - 1 - node->offset;
-        char *buf = client->refbuf->data + node->offset;
-
-        if (len > 0)
-        {
-            if (client->con->con_time + timeout <= time(NULL))
-                len = 0;
-            else
-                len = client_read_bytes (client, buf, len);
-        }
-
-        if (len > 0)
-        {
-            node->offset += len;
-            if ((len = util_find_eos_delim(client->refbuf, -node->offset,
-			    node->shoutcast?HEADER_READ_LINE:HEADER_READ_ENTIRE)) < 0) {
-                node->stream_offset = len;
-                if ((client_queue_t **)_req_queue_tail == &(node->next))
-                    _req_queue_tail = (volatile client_queue_t **)node_ref;
-                *node_ref = node->next;
-                node->next = NULL;
-                _add_connection (node);
-                continue;
-            }
-        }
-        else
-        {
-            if (len == 0 || client->con->error)
-            {
-                if ((client_queue_t **)_req_queue_tail == &node->next)
-                    _req_queue_tail = (volatile client_queue_t **)node_ref;
-                *node_ref = node->next;
-                client_destroy (client);
-                free (node);
-                continue;
-            }
-        }
-        node_ref = &node->next;
-    }
-    _handle_connection();
-}
-
-
-/* add node to the queue of requests. This is where the clients are when
- * initial http details are read.
- */
-static void _add_request_queue (client_queue_t *node)
-{
-    *_req_queue_tail = node;
-    _req_queue_tail = (volatile client_queue_t **)&node->next;
-}
-
-
 void connection_accept_loop (void)
 {
     connection_t *con;
@@ -663,68 +567,99 @@ void connection_accept_loop (void)
     {
         con = _accept_connection (duration);
 
-        if (con)
-        {
-            client_queue_t *node;
-            ice_config_t *config;
-            client_t *client = NULL;
-            listener_t *listener;
+        if (!con) {
+            duration = 300; /* use longer timeouts when nothing waiting */
+            continue;
+        }
 
-            global_lock();
-            if (client_create (&client, con, NULL) < 0)
-            {
-                global_unlock();
-                client_send_403 (client, "Icecast connection limit reached");
-                /* don't be too eager as this is an imposed hard limit */
-                thread_sleep (400000);
-                continue;
-            }
+        ice_config_t *config;
+        client_t *client = NULL;
+        listener_t *listener;
+        refbuf_t *header = NULL;
+        http_parser_t *parser = NULL;
+        int hdrsize = 0;
+        int shoutcast = 0;
+        char *shoutcast_mount = NULL;
+
+        header = refbuf_new (PER_CLIENT_REFBUF_SIZE);
+        hdrsize = util_read_header (con, header, HEADER_READ_ENTIRE);
+        if (hdrsize < 0)
+        {
+            global_unlock();
+            ERROR ("Header read failed");
+            thread_sleep (400000);
+            continue;
+        }
 
-            /* setup client for reading incoming http */
-            client->refbuf->data [PER_CLIENT_REFBUF_SIZE-1] = '\000';
+        /* process normal HTTP headers */
+        parser = httpp_create_parser();
+        httpp_initialize(parser, NULL);
+        if (!httpp_parse (parser, header->data, hdrsize))
+        {
+            ERROR0("HTTP request parsing failed");
+            client_destroy (client);
+            continue;
+        }
 
-            if (sock_set_blocking (client->con->sock, 0) || sock_set_nodelay (client->con->sock))
-            {
-                global_unlock();
-                WARN0 ("failed to set tcp options on client connection, dropping");
-                client_destroy (client);
-                continue;
-            }
+        if (httpp_getvar (parser, HTTPP_VAR_ERROR_MESSAGE))
+        {
+            ERROR("Error(%s)", httpp_getvar(parser, HTTPP_VAR_ERROR_MESSAGE));
+            break;
+        }
 
-            node = calloc (1, sizeof (client_queue_t));
-            if (node == NULL)
-            {
-                global_unlock();
-                client_destroy (client);
-                continue;
-            }
-            node->client = client;
+        global_lock();
+        if (client_create (&client, con, parser) < 0)
+        {
+            global_unlock();
+            client_send_403 (client, "Icecast connection limit reached");
+            /* don't be too eager as this is an imposed hard limit */
+            thread_sleep (400000);
+            continue;
+        }
 
-            config = config_get_config();
-            listener = config_get_listen_sock (config, client->con);
+        client_set_queue (client, header);
 
-            if (listener)
-            {
-                if (listener->shoutcast_compat)
-                    node->shoutcast = 1;
-                if (listener->ssl && ssl_ok)
-                    connection_uses_ssl (client->con);
-                if (listener->shoutcast_mount)
-                    node->shoutcast_mount = strdup (listener->shoutcast_mount);
-            }
+        if (sock_set_blocking (client->con->sock, 0) || sock_set_nodelay (client->con->sock))
+        {
             global_unlock();
-            config_release_config();
-
-            _add_request_queue (node);
-            stats_event_inc (NULL, "connections");
-            duration = 5;
+            WARN0 ("failed to set tcp options on client connection, dropping");
+            client_destroy (client);
+            continue;
         }
-        else
+
+        header->len -= hdrsize;
+        memmove(header->data, header->data + hdrsize, header->len);
+        client_set_queue (client, header);
+        refbuf_release(header);
+
+//        client->pos = hdrsize;
+
+        config = config_get_config();
+        listener = config_get_listen_sock (config, client->con);
+
+        if (listener)
         {
-            if (_req_queue == NULL)
-                duration = 300; /* use longer timeouts when nothing waiting */
+            if (listener->shoutcast_compat)
+                shoutcast = 1;
+            if (listener->ssl && ssl_ok)
+                connection_uses_ssl (client->con);
+            if (listener->shoutcast_mount)
+                shoutcast_mount = strdup (listener->shoutcast_mount);
+        }
+        global_unlock();
+        config_release_config();
+
+        stats_event_inc (NULL, "connections");
+        duration = 5;
+
+        if (client->con->con_time + timeout <= time(NULL))
+            continue;
+
+        if (shoutcast) {
+            _handle_shoutcast_compatible (shoutcast, shoutcast_mount);
+        } else {
+            _handle_client (client);
         }
-        process_request_queue (timeout);
     }
 
     /* Give all the other threads notification to shut down */
@@ -1113,22 +1048,22 @@ static void _handle_get_request (client_t *client, char *passed_uri)
     if (uri != passed_uri) free (uri);
 }
 
-static void _handle_shoutcast_compatible (client_queue_t *node)
+static void _handle_shoutcast_compatible (int shoutcast, char *shoutcast_mount)
 {
+/*
+  SHOUTCAST IS BROKEN
+*/
     char *http_compliant;
     int http_compliant_len = 0;
     http_parser_t *parser;
     ice_config_t *config = config_get_config ();
-    char *shoutcast_mount;
-    client_t *client = node->client;
+    client_t *client = NULL; //node->client;
     refbuf_t *refbuf = client->refbuf;
 
-    if (node->shoutcast_mount)
-        shoutcast_mount = node->shoutcast_mount;
-    else
+    if (!shoutcast_mount)
         shoutcast_mount = config->shoutcast_mount;
 
-    if (node->shoutcast == 1)
+    if (shoutcast == 1)
     {
         char *source_password;
         mount_proxy *mountinfo = config_find_mount (config, shoutcast_mount);
@@ -1143,8 +1078,6 @@ static void _handle_shoutcast_compatible (client_queue_t *node)
         if ((hdrlen = util_find_eos_delim (client->refbuf, 0, HEADER_READ_LINE)) < 0) {
             client_destroy (client);
             free (source_password);
-            free (node->shoutcast_mount);
-            free (node);
             return;
         }
 
@@ -1153,20 +1086,17 @@ static void _handle_shoutcast_compatible (client_queue_t *node)
             /* send this non-blocking but if there is only a partial write
              * then leave to header timeout */
             sock_write (client->con->sock, "OK2\r\nicy-caps:11\r\n\r\n");
-            node->offset -= hdrlen;
-            memmove (refbuf->data, refbuf->data + hdrlen, node->offset+1);
-            node->shoutcast = 2;
+            shoutcast = 2;
             /* we've checked the password, now send it back for reading headers */
-            _add_request_queue (node);
+//            _add_request_queue (node);
             free (source_password);
             return;
         }
 
-	INFO1 ("password does not match \"%s\"", client->refbuf->data);
+        INFO1 ("password does not match \"%s\"", refbuf->data);
+
         client_destroy (client);
         free (source_password);
-        free (node->shoutcast_mount);
-        free (node);
         return;
     }
     /* actually make a copy as we are dropping the config lock */
@@ -1175,7 +1105,7 @@ static void _handle_shoutcast_compatible (client_queue_t *node)
     /* Here we create a valid HTTP request based of the information
        that was passed in via the non-HTTP style protocol above. This
        means we can use some of our existing code to handle this case */
-    http_compliant_len = 20 + strlen (shoutcast_mount) + node->offset;
+//    http_compliant_len = 20 + strlen (shoutcast_mount) + node->offset;
     http_compliant = (char *)calloc(1, http_compliant_len);
     snprintf (http_compliant, http_compliant_len,
             "SOURCE %s HTTP/1.0\r\n%s", shoutcast_mount, refbuf->data);
@@ -1183,15 +1113,6 @@ static void _handle_shoutcast_compatible (client_queue_t *node)
     httpp_initialize(parser, NULL);
     if (httpp_parse (parser, http_compliant, strlen(http_compliant)))
     {
-        /* we may have more than just headers, so prepare for it */
-        if (node->stream_offset == node->offset)
-            refbuf->len = 0;
-        else
-        {
-            char *ptr = refbuf->data;
-            refbuf->len = node->offset - node->stream_offset;
-            memmove (ptr, ptr + node->stream_offset, refbuf->len);
-        }
         client->parser = parser;
         source_startup (client, shoutcast_mount, SHOUTCAST_SOURCE_AUTH);
     }
@@ -1201,106 +1122,52 @@ static void _handle_shoutcast_compatible (client_queue_t *node)
     }
     free (http_compliant);
     free (shoutcast_mount);
-    free (node->shoutcast_mount);
-    free (node);
     return;
 }
 
-
-/* Connection thread. Here we take clients off the connection queue and check
- * the contents provided. We set up the parser then hand off to the specific
- * request handler.
- */
-static void _handle_connection(void)
+static int _handle_client (client_t *client)
 {
-    http_parser_t *parser;
     const char *rawuri;
-    client_queue_t *node;
-
-    while (1)
-    {
-        node = _get_connection();
-        if (! node)
-            break;
-
-        client_t *client = node->client;
-        refbuf_t *refbuf = client->refbuf;
-        char *uri;
-
-
-        /* Check for special shoutcast compatability processing */
-        if (node->shoutcast)
-        {
-            _handle_shoutcast_compatible (node);
-            continue;
-        }
-
-        /* process normal HTTP headers */
-        parser = httpp_create_parser();
-        httpp_initialize(parser, NULL);
-        client->parser = parser;
-        if (!httpp_parse (parser, refbuf->data, node->offset))
-        {
-            free (node);
-            ERROR0("HTTP request parsing failed");
-            client_destroy (client);
-            continue;
-        }
-
-        /* we may have more than just headers, so prepare for it */
-        if (node->stream_offset == node->offset) {
-            refbuf->len = 0;
-        } else {
-            char *ptr = refbuf->data;
-            refbuf->len = node->offset - node->stream_offset;
-            memmove (ptr, ptr + node->stream_offset, refbuf->len);
-        }
-
-        rawuri = httpp_getvar(parser, HTTPP_VAR_URI);
+    http_parser_t *parser = client->parser;
+    char *uri;
 
-        /* assign a port-based shoutcast mountpoint if required */
-        if (node->shoutcast_mount && strcmp (rawuri, "/admin.cgi") == 0)
-            httpp_set_query_param (client->parser, "mount", node->shoutcast_mount);
+    rawuri = httpp_getvar(parser, HTTPP_VAR_URI);
 
-        free (node->shoutcast_mount);
-        free (node);
-
-        if (strcmp("ICE",  httpp_getvar(parser, HTTPP_VAR_PROTOCOL)) &&
-            strcmp("HTTP", httpp_getvar(parser, HTTPP_VAR_PROTOCOL))) {
-            ERROR0("Bad HTTP protocol detected");
-            client_destroy (client);
-            continue;
-        }
+    if (strcmp("ICE",  httpp_getvar(parser, HTTPP_VAR_PROTOCOL)) &&
+        strcmp("HTTP", httpp_getvar(parser, HTTPP_VAR_PROTOCOL))) {
+        ERROR0("Bad HTTP protocol detected");
+        client_destroy (client);
+        return 0;
+    }
 
-        uri = util_normalise_uri(rawuri);
+    uri = util_normalise_uri(rawuri);
 
-        if (uri == NULL)
-        {
-            client_destroy (client);
-            continue;
-        }
-
-        if (parser->req_type == httpp_req_source) {
-            _handle_source_request (client, uri);
-        }
-        else if (parser->req_type == httpp_req_post) {
-            _handle_post_request (client, uri);
-        }
-        else if (parser->req_type == httpp_req_stats) {
-            _handle_stats_request (client, uri);
-        }
-        else if (parser->req_type == httpp_req_get) {
-            _handle_get_request (client, uri);
-        }
-        else {
-            ERROR0("Wrong request type from client");
-            client_send_400 (client, "unknown request");
-        }
+    if (uri == NULL)
+    {
+        client_destroy (client);
+        return 0;
+    }
 
-        free(uri);
+    if (parser->req_type == httpp_req_source) {
+        _handle_source_request (client, uri);
+    }
+    else if (parser->req_type == httpp_req_post) {
+        _handle_post_request (client, uri);
+    }
+    else if (parser->req_type == httpp_req_stats) {
+        _handle_stats_request (client, uri);
+    }
+    else if (parser->req_type == httpp_req_get) {
+        _handle_get_request (client, uri);
+    }
+    else {
+        ERROR0("Wrong request type from client");
+        client_send_400 (client, "unknown request");
     }
-}
 
+    free(uri);
+    return 1;
+}
 
 /* called when listening thread is not checking for incoming connections */
 int connection_setup_sockets (ice_config_t *config)
-- 
1.7.1
lmpx.com only provides a reader for public news (NNTP) servers. It is not affiliated with the servers or forums shown here and is not responsible for the content of articles, which is written by their respective authors.