
#include <stdio.h>
#include <malloc.h>
#include <pthread.h>
#include <curl/curl.h>


#define MAX_TEST_URL 2

#define SIMULATE_FULL_BUFFER 1
#define TEST_SEQUENCED_DATA  0
#define TEST_BUFFER_SIZE     16384



const char * const urls[]= {
  "http://192.168.131.65/test_dat/seq_data_1.bin",
  "http://192.168.131.65/test_dat/seq_data_2.bin",
  "http://192.168.131.65/test_dat/seq_data_3.bin",
  "http://192.168.131.65/test_dat/seq_data_4.bin"
};


/***************************************************************************/
static size_t callback_content(void *ptr, size_t size, size_t nmemb, void *userp);
static int callback_debug(CURL *handle, curl_infotype type, char *data, size_t size, void *userp);


struct http_request_context_t
{
  CURL *curl;

  struct http_context_t *httpContext;

  unsigned char expectNextByte;
  unsigned long bytesReceived;

  unsigned char buffer[CURL_MAX_WRITE_SIZE];

  int shallRelease;
  int shallResume;
  int paused;
};

struct http_context_t
{
  CURLSH *curlShare;
  CURLM *curlMulti;

  int quit;
  pthread_mutex_t shareLock;

  struct http_request_context_t  *httpRequestContexts[MAX_TEST_URL];

};

/***************************************************************************/

static struct http_context_t      gHttpContext;



/***************************************************************************/

static void curl_shared_lock(CURL *handle, curl_lock_data data, curl_lock_access access, void *userptr)
{
  struct http_context_t *httpContext = (struct http_context_t *)userptr;

  pthread_mutex_lock(&httpContext->shareLock);    
}

static void curl_shared_unlock(CURL *handle, curl_lock_data data, void *userptr)
{
  struct http_context_t *httpContext = (struct http_context_t *)userptr;

  pthread_mutex_unlock(&httpContext->shareLock);
}
/***************************************************************************/

static int add_request(struct http_context_t *httpContext, struct http_request_context_t *request)
{
  int n;
  int success = 0;

  for(n = 0; (n < MAX_TEST_URL) && !success; n++) 
  {
    if (!httpContext->httpRequestContexts[n]) 
    {
      httpContext->httpRequestContexts[n] = request;
      success = 1;
    }
  }

  return success;
}

static int remove_request(struct http_context_t *httpContext, struct http_request_context_t *request)
{
  int n;
  int success = 0;

  for(n = 0; (n < MAX_TEST_URL) && !success; n++) 
  {
    if (httpContext->httpRequestContexts[n] == request) 
    {
      httpContext->httpRequestContexts[n] = NULL;
      success = 1;
    }
  }

  return success;
}

static struct http_request_context_t *find_request(struct http_context_t *httpContext, CURL *curlHandle)
{
  int n;

  for(n = 0; (n < MAX_TEST_URL); n++) 
  {
    if (!httpContext->httpRequestContexts[n]) 
    {
      if (httpContext->httpRequestContexts[n]->curl == curlHandle) 
      {
        return httpContext->httpRequestContexts[n];
      }
    }
  }

  return NULL;
}


/***************************************************************************/

static int send_request(struct http_context_t *httpContext, char *url)
{
  struct http_request_context_t *tempRequest;

  tempRequest = malloc(sizeof(struct http_request_context_t));

  if(!tempRequest) 
  {
    return 0;
  }

  memset(tempRequest, 0, (sizeof(struct http_request_context_t)));

  printf("creating curl handle for: GET %s\n", url);
  tempRequest->curl = curl_easy_init();

  if (!tempRequest->curl)
  {
    free(tempRequest);
    return 0;
  }

  tempRequest->httpContext = httpContext;

  curl_easy_setopt(tempRequest->curl, CURLOPT_URL, url);
  curl_easy_setopt(tempRequest->curl, CURLOPT_WRITEFUNCTION, callback_content);
  curl_easy_setopt(tempRequest->curl, CURLOPT_WRITEDATA, tempRequest);

#if 1
  curl_easy_setopt(tempRequest->curl, CURLOPT_VERBOSE, 1);
  curl_easy_setopt(tempRequest->curl, CURLOPT_DEBUGFUNCTION, callback_debug);
  curl_easy_setopt(tempRequest->curl, CURLOPT_DEBUGDATA, tempRequest);
#endif

  curl_easy_setopt(tempRequest->curl, CURLOPT_TIMEOUT, 30);

  printf("send_request: created request[%p] (curl %p)\n", tempRequest, tempRequest->curl);

  pthread_mutex_lock(&tempRequest->httpContext->shareLock);
  add_request(tempRequest->httpContext, tempRequest);
  pthread_mutex_unlock(&tempRequest->httpContext->shareLock);

  curl_multi_add_handle(tempRequest->httpContext->curlMulti, tempRequest->curl);
}

/***************************************************************************/

static int callback_debug(CURL *handle, curl_infotype type,
             char *data, size_t size, void *userp)
{
  struct data *config = (struct data *)userp;
  const char *text;

  switch (type) {
  case CURLINFO_TEXT:
    printf("== Info: %s", data);
  default:
    return 0;

  case CURLINFO_HEADER_OUT:
    text = "=> Send header";
    break;
  case CURLINFO_DATA_OUT:
    text = "=> Send data";
    break;
  case CURLINFO_HEADER_IN:
    text = "<= Recv header";
    break;
  case CURLINFO_DATA_IN:
    text = "<= Recv data";
    break;
  }

  printf("%s (%d bytes)\n", text, size);

  return 0;
}


static size_t callback_content(void *ptr, size_t size, size_t nmemb, void *userp)
{
  unsigned long sizeToWrite = (size*nmemb);
  int n;
  unsigned char *dataptr = (unsigned char *)ptr;

  struct http_request_context_t *request = (struct http_request_context_t *)userp;

  printf("callback_content: entered for request[%p] - size %d\n", request, sizeToWrite);

  for(n = 0; n < 16; n++)
  {
    printf( "%02x ", dataptr[n]);
  }
  printf( "....\n");

  pthread_mutex_lock(&request->httpContext->shareLock);

#if TEST_SEQUENCED_DATA
  if (request->expectNextByte != dataptr[0])
  {
    fprintf(stderr, "***** data error[request %p]: expected %d, got %d\n", request, request->expectNextByte, dataptr[0]);
  }
#endif

#if SIMULATE_FULL_BUFFER
  /* simulate a full buffer */
  if((request->bytesReceived + sizeToWrite) > TEST_BUFFER_SIZE) 
  {
    printf("callback_content: buffer full for request[%p] - pausing...\n", request);
    request->paused = 1;

    /* signal other thread to read from buffer here*/
    pthread_mutex_unlock(&request->httpContext->shareLock);

    /* tell curl that we're full */
    return CURL_WRITEFUNC_PAUSE;
  }
#endif 
  /* else try to add to buffer here */
  printf("callback_content: data added to buffer...\n", request);

  request->bytesReceived += sizeToWrite;

#if TEST_SEQUENCED_DATA
  request->expectNextByte = (unsigned char)(((unsigned long)(request->expectNextByte) + sizeToWrite) & 0xff);
#endif

  pthread_mutex_unlock(&request->httpContext->shareLock);

  /* signal other thread to read from buffer here*/

  return size*nmemb;
}


static void curl_thread_checkCompletions(struct http_context_t *httpContext)
{
  CURLMsg *pstCurlMsg;
  int numCurlMsgs;
  struct http_request_context_t *request = NULL;
  CURLcode code;

  pstCurlMsg = curl_multi_info_read(httpContext->curlMulti, &numCurlMsgs);

  while (pstCurlMsg)
  {
    switch (pstCurlMsg->msg)
    {
      case CURLMSG_DONE:
      {
        pthread_mutex_lock(&httpContext->shareLock);

        request = find_request(httpContext, pstCurlMsg->easy_handle);

        if (request)
        {
          if (CURLE_OK == pstCurlMsg->data.result) 
          {
            printf("Request [%p (%p)] marked complete\n", request, request->curl);
          }
          else
          {
            printf("Request [%p (%p)]: HTTP request failed\n", request, request->curl);
            printf(" -- The error was #%d: %s\n",pstCurlMsg->data.result,curl_easy_strerror( pstCurlMsg->data.result ) );
          }

          request->shallRelease = 1;
        }

        pthread_mutex_unlock(&httpContext->shareLock);
        break;
      }

      default:
        printf("Unknown CURLMSG_* [%d] for curl handle %p\n", pstCurlMsg->msg, pstCurlMsg->easy_handle);
    }

    pstCurlMsg = curl_multi_info_read(httpContext->curlMulti, &numCurlMsgs);
  }
}


static void curl_thread_checkForReleaseOrResume(struct http_context_t *httpContext)
{
  int n;

  for(n = 0; n < MAX_TEST_URL; n++) 
  {
    if(httpContext->httpRequestContexts[n]) 
    {
      struct http_request_context_t *request = httpContext->httpRequestContexts[n];

      if (request->shallResume) 
      {
        printf("curl_thread_checkForReleaseOrResume: unpausing request %p...\n", request);

        pthread_mutex_lock(&request->httpContext->shareLock);
        request->shallResume = 0;
        request->paused = 0;
        pthread_mutex_unlock(&request->httpContext->shareLock);

        curl_easy_pause(request->curl, CURLPAUSE_CONT);
      }
      else if (request->shallRelease) 
      {
        printf("curl_thread_checkForReleaseOrResume: removing request %p...\n", request);

        pthread_mutex_lock(&request->httpContext->shareLock);
        remove_request(httpContext, request);
        pthread_mutex_unlock(&request->httpContext->shareLock);

        curl_multi_remove_handle(httpContext->curlMulti, request->curl);
        curl_easy_cleanup(request->curl);

        free(request);
      }
    }
  }
}

static void *curl_thread(void *data)
{
  struct http_context_t *httpContext = (struct http_context_t *)data;
  CURLMcode mcode;
  CURLcode code;
  int stillRunning = 0;

  httpContext->curlShare = curl_share_init();
  if (!httpContext->curlShare)
  {
    fprintf(stderr,"couldn't create share handle\n");
    return;
  }

  curl_share_setopt(httpContext->curlShare,CURLSHOPT_USERDATA, httpContext);
  curl_share_setopt(httpContext->curlShare,CURLSHOPT_SHARE, CURL_LOCK_DATA_COOKIE);
  curl_share_setopt(httpContext->curlShare,CURLSHOPT_LOCKFUNC, curl_shared_lock);
  curl_share_setopt(httpContext->curlShare,CURLSHOPT_UNLOCKFUNC, curl_shared_unlock);

  /* init a multi stack */
  httpContext->curlMulti = curl_multi_init();

  if (!httpContext->curlMulti)
  {
    fprintf(stderr,"couldn't create multi handle\n");
    return;
  }

  curl_multi_setopt(httpContext->curlMulti, CURLMOPT_PIPELINING, 1);
 
  while(!httpContext->quit)
  {
    mcode = curl_multi_perform(httpContext->curlMulti, &stillRunning);

    curl_thread_checkCompletions(httpContext);

    if (mcode != CURLM_CALL_MULTI_PERFORM)
    {
      curl_thread_checkForReleaseOrResume(httpContext);

      if (stillRunning)
      {
        struct timeval timeout;
        int rc; /* select() return code */
        fd_set fdread;
        fd_set fdwrite;
        fd_set fdexcep;
        int maxfd;

        FD_ZERO(&fdread);
        FD_ZERO(&fdwrite);
        FD_ZERO(&fdexcep);

        timeout.tv_sec = 0;
        timeout.tv_usec = 50000;
    
        /* get file descriptors from the transfers */
        curl_multi_fdset(httpContext->curlMulti, &fdread, &fdwrite, &fdexcep, &maxfd);

        //printf("debug: curl-thread: still running = %d -- selecting...\n", stillRunning);
        
        rc = select(maxfd+1, &fdread, &fdwrite, &fdexcep, &timeout);
    
        switch(rc) {
        case -1:
          /* select error */
          printf("curl-thread: select error\n");
          break;

        case 0:
          //printf("curl-thread: select TIMEOUT\n");
          break;

        default:
          //printf("curl-thread: select says stuff to do [%d]\n", rc);
          break;
        }
      }
    }
  }

  curl_multi_cleanup(httpContext->curlMulti);
  curl_share_cleanup(httpContext->curlShare);
}


int main(int argc, char **argv)
{
  pthread_t curlThreadId;
  int i;
  int error;

  memset(&gHttpContext, 0, sizeof(struct http_context_t));
  pthread_mutex_init(&gHttpContext.shareLock, NULL);


  /* Must initialize libcurl before any threads are started */
  curl_global_init(CURL_GLOBAL_ALL);


  error = pthread_create(&curlThreadId,
                           NULL, /* default attributes please */
                           curl_thread,
                           &gHttpContext);

  if(0 == error)
  {
    for (i = 0; i < MAX_TEST_URL; i++)
    {
      printf("adding %s\n", urls[i]);
      send_request(&gHttpContext, urls[i]);
    }
  
#if SIMULATE_FULL_BUFFER
    while(!gHttpContext.quit) 
    {
      sleep(1);

      pthread_mutex_lock(&gHttpContext.shareLock);
      for (i = 0; i < MAX_TEST_URL; i++)
      {
        if(gHttpContext.httpRequestContexts[i]) 
        {
          if (gHttpContext.httpRequestContexts[i]->paused) 
          {
            gHttpContext.httpRequestContexts[i]->bytesReceived = 0;

            printf("buffer now empty, shall resume request %p\n", gHttpContext.httpRequestContexts[i]);
            gHttpContext.httpRequestContexts[i]->shallResume = 1;
          }
        }
      }
      pthread_mutex_unlock(&gHttpContext.shareLock);
    }
#endif
  
    error = pthread_join(curlThreadId, NULL);
    fprintf(stderr, "curl thread terminated\n");
  }
  else
  {
    fprintf(stderr, "Couldn't run curl thread errno %d\n", error);
  }

  curl_global_cleanup();

  pthread_mutex_destroy(&gHttpContext.shareLock);

  return 0;
}


