mirror of https://github.com/tstack/lnav
[curl] add a curl looper to handle url requests
parent
248dd78f63
commit
f286950854
@ -0,0 +1,258 @@
|
|||||||
|
/**
|
||||||
|
* Copyright (c) 2015, Timothy Stack
|
||||||
|
*
|
||||||
|
* All rights reserved.
|
||||||
|
*
|
||||||
|
* Redistribution and use in source and binary forms, with or without
|
||||||
|
* modification, are permitted provided that the following conditions are met:
|
||||||
|
*
|
||||||
|
* * Redistributions of source code must retain the above copyright notice, this
|
||||||
|
* list of conditions and the following disclaimer.
|
||||||
|
* * Redistributions in binary form must reproduce the above copyright notice,
|
||||||
|
* this list of conditions and the following disclaimer in the documentation
|
||||||
|
* and/or other materials provided with the distribution.
|
||||||
|
* * Neither the name of Timothy Stack nor the names of its contributors
|
||||||
|
* may be used to endorse or promote products derived from this software
|
||||||
|
* without specific prior written permission.
|
||||||
|
*
|
||||||
|
* THIS SOFTWARE IS PROVIDED BY THE REGENTS AND CONTRIBUTORS ''AS IS'' AND ANY
|
||||||
|
* EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||||
|
* WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
|
||||||
|
* DISCLAIMED. IN NO EVENT SHALL THE REGENTS OR CONTRIBUTORS BE LIABLE FOR ANY
|
||||||
|
* DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
|
||||||
|
* (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
|
||||||
|
* LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
|
||||||
|
* ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
|
||||||
|
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
|
||||||
|
* SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
||||||
|
*
|
||||||
|
* @file curl_looper.cc
|
||||||
|
*/
|
||||||
|
|
||||||
|
#include "config.h"
|
||||||
|
|
||||||
|
#ifdef HAVE_LIBCURL
|
||||||
|
#include <curl/multi.h>
|
||||||
|
|
||||||
|
#include "curl_looper.hh"
|
||||||
|
|
||||||
|
using namespace std;
|
||||||
|
|
||||||
|
struct curl_request_eq {
|
||||||
|
curl_request_eq(const std::string &name) : cre_name(name) {
|
||||||
|
};
|
||||||
|
|
||||||
|
bool operator()(const curl_request *cr) const {
|
||||||
|
return this->cre_name == cr->get_name();
|
||||||
|
};
|
||||||
|
|
||||||
|
bool operator()(const pair<mstime_t, curl_request *> &pair) const {
|
||||||
|
return this->cre_name == pair.second->get_name();
|
||||||
|
};
|
||||||
|
|
||||||
|
const std::string &cre_name;
|
||||||
|
};
|
||||||
|
|
||||||
|
int curl_request::debug_cb(CURL *handle,
|
||||||
|
curl_infotype type,
|
||||||
|
char *data,
|
||||||
|
size_t size,
|
||||||
|
void *userp) {
|
||||||
|
curl_request *cr = (curl_request *) userp;
|
||||||
|
|
||||||
|
if (type == CURLINFO_TEXT) {
|
||||||
|
while (size > 0 && isspace(data[size - 1])) {
|
||||||
|
size -= 1;
|
||||||
|
}
|
||||||
|
log_debug("%s:%.*s", cr->get_name().c_str(), size, data);
|
||||||
|
}
|
||||||
|
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
void *curl_looper::trampoline(void *arg)
|
||||||
|
{
|
||||||
|
curl_looper *cl = (curl_looper *) arg;
|
||||||
|
|
||||||
|
return cl->run();
|
||||||
|
}
|
||||||
|
|
||||||
|
void *curl_looper::run()
|
||||||
|
{
|
||||||
|
log_info("curl looper thread started");
|
||||||
|
while (this->cl_looping) {
|
||||||
|
this->loop_body();
|
||||||
|
}
|
||||||
|
log_info("curl looper thread exiting");
|
||||||
|
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
|
||||||
|
void curl_looper::loop_body()
|
||||||
|
{
|
||||||
|
mstime_t current_time = getmstime();
|
||||||
|
int timeout = this->compute_timeout(current_time);
|
||||||
|
|
||||||
|
if (this->cl_handle_to_request.empty()) {
|
||||||
|
mutex_guard mg(this->cl_mutex);
|
||||||
|
|
||||||
|
if (this->cl_new_requests.empty() && this->cl_close_requests.empty()) {
|
||||||
|
mstime_t deadline = current_time + timeout;
|
||||||
|
struct timespec ts;
|
||||||
|
|
||||||
|
ts.tv_sec = deadline / 1000ULL;
|
||||||
|
ts.tv_nsec = (deadline % 1000ULL) * 1000 * 1000;
|
||||||
|
log_trace("no requests in progress, waiting %d ms for new ones",
|
||||||
|
timeout);
|
||||||
|
pthread_cond_timedwait(&this->cl_cond, &this->cl_mutex, &ts);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
this->perform_io();
|
||||||
|
|
||||||
|
this->check_for_finished_requests();
|
||||||
|
|
||||||
|
this->check_for_new_requests();
|
||||||
|
|
||||||
|
this->requeue_requests(current_time + 5);
|
||||||
|
}
|
||||||
|
|
||||||
|
void curl_looper::perform_io()
|
||||||
|
{
|
||||||
|
if (this->cl_handle_to_request.empty()) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
mstime_t current_time = getmstime();
|
||||||
|
int timeout = this->compute_timeout(current_time);
|
||||||
|
int running_handles;
|
||||||
|
|
||||||
|
curl_multi_wait(this->cl_curl_multi,
|
||||||
|
NULL,
|
||||||
|
0,
|
||||||
|
timeout,
|
||||||
|
NULL);
|
||||||
|
curl_multi_perform(this->cl_curl_multi, &running_handles);
|
||||||
|
}
|
||||||
|
|
||||||
|
void curl_looper::requeue_requests(mstime_t up_to_time)
|
||||||
|
{
|
||||||
|
while (!this->cl_poll_queue.empty() &&
|
||||||
|
this->cl_poll_queue.front().first <= up_to_time) {
|
||||||
|
curl_request *cr = this->cl_poll_queue.front().second;
|
||||||
|
|
||||||
|
log_debug("%s:polling request is ready again -- %p",
|
||||||
|
cr->get_name().c_str(), cr);
|
||||||
|
this->cl_handle_to_request[cr->get_handle()] = cr;
|
||||||
|
curl_multi_add_handle(this->cl_curl_multi, cr->get_handle());
|
||||||
|
this->cl_poll_queue.erase(this->cl_poll_queue.begin());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
void curl_looper::check_for_new_requests() {
|
||||||
|
mutex_guard mg(this->cl_mutex);
|
||||||
|
|
||||||
|
while (!this->cl_new_requests.empty()) {
|
||||||
|
curl_request *cr = this->cl_new_requests.back();
|
||||||
|
|
||||||
|
log_info("%s:new curl request %p",
|
||||||
|
cr->get_name().c_str(),
|
||||||
|
cr);
|
||||||
|
this->cl_handle_to_request[cr->get_handle()] = cr;
|
||||||
|
curl_multi_add_handle(this->cl_curl_multi, cr->get_handle());
|
||||||
|
this->cl_new_requests.pop_back();
|
||||||
|
}
|
||||||
|
while (!this->cl_close_requests.empty()) {
|
||||||
|
const std::string &name = this->cl_close_requests.back();
|
||||||
|
vector<curl_request *>::iterator all_iter = find_if(
|
||||||
|
this->cl_all_requests.begin(),
|
||||||
|
this->cl_all_requests.end(),
|
||||||
|
curl_request_eq(name));
|
||||||
|
|
||||||
|
log_info("attempting to close request -- %s", name.c_str());
|
||||||
|
if (all_iter != this->cl_all_requests.end()) {
|
||||||
|
map<CURL *, curl_request *>::iterator act_iter;
|
||||||
|
vector<pair<mstime_t, curl_request *> >::iterator poll_iter;
|
||||||
|
curl_request *cr = *all_iter;
|
||||||
|
|
||||||
|
log_info("%s:closing request -- %p",
|
||||||
|
cr->get_name().c_str(), cr);
|
||||||
|
(*all_iter)->close();
|
||||||
|
act_iter = this->cl_handle_to_request.find(cr);
|
||||||
|
if (act_iter != this->cl_handle_to_request.end()) {
|
||||||
|
this->cl_handle_to_request.erase(act_iter);
|
||||||
|
curl_multi_remove_handle(this->cl_curl_multi,
|
||||||
|
cr->get_handle());
|
||||||
|
delete cr;
|
||||||
|
}
|
||||||
|
poll_iter = find_if(this->cl_poll_queue.begin(),
|
||||||
|
this->cl_poll_queue.end(),
|
||||||
|
curl_request_eq(name));
|
||||||
|
if (poll_iter != this->cl_poll_queue.end()) {
|
||||||
|
this->cl_poll_queue.erase(poll_iter);
|
||||||
|
}
|
||||||
|
this->cl_all_requests.erase(all_iter);
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
log_error("Unable to find request with the name -- %s",
|
||||||
|
name.c_str());
|
||||||
|
}
|
||||||
|
|
||||||
|
this->cl_close_requests.pop_back();
|
||||||
|
|
||||||
|
pthread_cond_broadcast(&this->cl_cond);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
void curl_looper::check_for_finished_requests()
|
||||||
|
{
|
||||||
|
CURLMsg *msg;
|
||||||
|
int msgs_left;
|
||||||
|
|
||||||
|
while ((msg = curl_multi_info_read(this->cl_curl_multi, &msgs_left)) != NULL) {
|
||||||
|
if (msg->msg != CURLMSG_DONE) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
CURL *easy = msg->easy_handle;
|
||||||
|
map<CURL *, curl_request *>::iterator iter = this->cl_handle_to_request.find(easy);
|
||||||
|
|
||||||
|
curl_multi_remove_handle(this->cl_curl_multi, easy);
|
||||||
|
if (iter != this->cl_handle_to_request.end()) {
|
||||||
|
curl_request *cr = iter->second;
|
||||||
|
long delay_ms;
|
||||||
|
|
||||||
|
this->cl_handle_to_request.erase(iter);
|
||||||
|
delay_ms = cr->complete(msg->data.result);
|
||||||
|
if (delay_ms < 0) {
|
||||||
|
vector<curl_request *>::iterator all_iter;
|
||||||
|
|
||||||
|
log_info("%s:curl_request %p finished, deleting...",
|
||||||
|
cr->get_name().c_str(), cr);
|
||||||
|
{
|
||||||
|
mutex_guard mg(this->cl_mutex);
|
||||||
|
|
||||||
|
all_iter = find(this->cl_all_requests.begin(),
|
||||||
|
this->cl_all_requests.end(),
|
||||||
|
cr);
|
||||||
|
if (all_iter != this->cl_all_requests.end()) {
|
||||||
|
this->cl_all_requests.erase(all_iter);
|
||||||
|
}
|
||||||
|
pthread_cond_broadcast(&cl_cond);
|
||||||
|
}
|
||||||
|
delete cr;
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
log_debug("%s:curl_request %p is polling, requeueing in %d",
|
||||||
|
cr->get_name().c_str(),
|
||||||
|
cr,
|
||||||
|
delay_ms);
|
||||||
|
this->cl_poll_queue.push_back(
|
||||||
|
make_pair(getmstime() + delay_ms, cr));
|
||||||
|
sort(this->cl_poll_queue.begin(), this->cl_poll_queue.end());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#endif
|
@ -0,0 +1,240 @@
|
|||||||
|
/**
|
||||||
|
* Copyright (c) 2015, Timothy Stack
|
||||||
|
*
|
||||||
|
* All rights reserved.
|
||||||
|
*
|
||||||
|
* Redistribution and use in source and binary forms, with or without
|
||||||
|
* modification, are permitted provided that the following conditions are met:
|
||||||
|
*
|
||||||
|
* * Redistributions of source code must retain the above copyright notice, this
|
||||||
|
* list of conditions and the following disclaimer.
|
||||||
|
* * Redistributions in binary form must reproduce the above copyright notice,
|
||||||
|
* this list of conditions and the following disclaimer in the documentation
|
||||||
|
* and/or other materials provided with the distribution.
|
||||||
|
* * Neither the name of Timothy Stack nor the names of its contributors
|
||||||
|
* may be used to endorse or promote products derived from this software
|
||||||
|
* without specific prior written permission.
|
||||||
|
*
|
||||||
|
* THIS SOFTWARE IS PROVIDED BY THE REGENTS AND CONTRIBUTORS ''AS IS'' AND ANY
|
||||||
|
* EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||||
|
* WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
|
||||||
|
* DISCLAIMED. IN NO EVENT SHALL THE REGENTS OR CONTRIBUTORS BE LIABLE FOR ANY
|
||||||
|
* DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
|
||||||
|
* (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
|
||||||
|
* LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
|
||||||
|
* ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
|
||||||
|
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
|
||||||
|
* SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
||||||
|
*
|
||||||
|
* @file curl_looper.hh
|
||||||
|
*/
|
||||||
|
|
||||||
|
#ifndef curl_looper_hh
|
||||||
|
#define curl_looper_hh
|
||||||
|
|
||||||
|
#include <map>
|
||||||
|
#include <string>
|
||||||
|
#include <vector>
|
||||||
|
|
||||||
|
#ifndef HAVE_LIBCURL
|
||||||
|
|
||||||
|
typedef int CURLcode;
|
||||||
|
|
||||||
|
class curl_request {
|
||||||
|
public:
|
||||||
|
curl_request(const std::string &name) {
|
||||||
|
};
|
||||||
|
};
|
||||||
|
|
||||||
|
class curl_looper {
|
||||||
|
public:
|
||||||
|
void start() { };
|
||||||
|
void add_request(curl_request *cr) { };
|
||||||
|
void close_request(const std::string &name) { };
|
||||||
|
void process_all() { };
|
||||||
|
};
|
||||||
|
|
||||||
|
#else
|
||||||
|
#include <curl/curl.h>
|
||||||
|
|
||||||
|
#include "auto_mem.hh"
|
||||||
|
#include "lnav_log.hh"
|
||||||
|
#include "lnav_util.hh"
|
||||||
|
#include "pthreadpp.hh"
|
||||||
|
|
||||||
|
class curl_request {
|
||||||
|
public:
|
||||||
|
curl_request(const std::string &name)
|
||||||
|
: cr_name(name),
|
||||||
|
cr_open(true),
|
||||||
|
cr_handle(curl_easy_cleanup),
|
||||||
|
cr_completions(0) {
|
||||||
|
this->cr_handle.reset(curl_easy_init());
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_NOSIGNAL, 1);
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_ERRORBUFFER, this->cr_error_buffer);
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_DEBUGFUNCTION, debug_cb);
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_DEBUGDATA, this);
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_VERBOSE, 1);
|
||||||
|
};
|
||||||
|
|
||||||
|
virtual ~curl_request() {
|
||||||
|
|
||||||
|
};
|
||||||
|
|
||||||
|
const std::string &get_name() const {
|
||||||
|
return this->cr_name;
|
||||||
|
};
|
||||||
|
|
||||||
|
virtual void close() {
|
||||||
|
this->cr_open = false;
|
||||||
|
};
|
||||||
|
|
||||||
|
bool is_open() {
|
||||||
|
return this->cr_open;
|
||||||
|
};
|
||||||
|
|
||||||
|
CURL *get_handle() const {
|
||||||
|
return this->cr_handle;
|
||||||
|
};
|
||||||
|
|
||||||
|
int get_completions() const {
|
||||||
|
return this->cr_completions;
|
||||||
|
};
|
||||||
|
|
||||||
|
virtual long complete(CURLcode result) {
|
||||||
|
double total_time = 0, download_size = 0, download_speed = 0;
|
||||||
|
|
||||||
|
this->cr_completions += 1;
|
||||||
|
curl_easy_getinfo(this->cr_handle, CURLINFO_TOTAL_TIME, &total_time);
|
||||||
|
log_debug("%s: total_time=%f", this->cr_name.c_str(), total_time);
|
||||||
|
curl_easy_getinfo(this->cr_handle, CURLINFO_SIZE_DOWNLOAD, &download_size);
|
||||||
|
log_debug("%s: download_size=%f", this->cr_name.c_str(), download_size);
|
||||||
|
curl_easy_getinfo(this->cr_handle, CURLINFO_SPEED_DOWNLOAD, &download_speed);
|
||||||
|
log_debug("%s: download_speed=%f", this->cr_name.c_str(), download_speed);
|
||||||
|
|
||||||
|
return -1;
|
||||||
|
};
|
||||||
|
|
||||||
|
protected:
|
||||||
|
|
||||||
|
static int debug_cb(CURL *handle,
|
||||||
|
curl_infotype type,
|
||||||
|
char *data,
|
||||||
|
size_t size,
|
||||||
|
void *userp);
|
||||||
|
|
||||||
|
const std::string cr_name;
|
||||||
|
bool cr_open;
|
||||||
|
auto_mem<CURL> cr_handle;
|
||||||
|
char cr_error_buffer[CURL_ERROR_SIZE];
|
||||||
|
int cr_completions;
|
||||||
|
};
|
||||||
|
|
||||||
|
class curl_looper {
|
||||||
|
public:
|
||||||
|
curl_looper()
|
||||||
|
: cl_started(false),
|
||||||
|
cl_looping(true),
|
||||||
|
cl_curl_multi(curl_multi_cleanup) {
|
||||||
|
this->cl_curl_multi.reset(curl_multi_init());
|
||||||
|
pthread_mutex_init(&this->cl_mutex, NULL);
|
||||||
|
pthread_cond_init(&this->cl_cond, NULL);
|
||||||
|
};
|
||||||
|
|
||||||
|
~curl_looper() {
|
||||||
|
this->stop();
|
||||||
|
pthread_cond_destroy(&this->cl_cond);
|
||||||
|
pthread_mutex_destroy(&this->cl_mutex);
|
||||||
|
}
|
||||||
|
|
||||||
|
void start() {
|
||||||
|
if (pthread_create(&this->cl_thread, NULL, trampoline, this) == 0) {
|
||||||
|
this->cl_started = true;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
void stop() {
|
||||||
|
if (this->cl_started) {
|
||||||
|
void *result;
|
||||||
|
|
||||||
|
this->cl_looping = false;
|
||||||
|
{
|
||||||
|
mutex_guard mg(this->cl_mutex);
|
||||||
|
|
||||||
|
pthread_cond_broadcast(&this->cl_cond);
|
||||||
|
}
|
||||||
|
log_debug("waiting for curl_looper thread");
|
||||||
|
pthread_join(this->cl_thread, &result);
|
||||||
|
log_debug("curl_looper thread joined");
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
void process_all() {
|
||||||
|
this->check_for_new_requests();
|
||||||
|
|
||||||
|
this->requeue_requests(LONG_MAX);
|
||||||
|
|
||||||
|
while (!this->cl_handle_to_request.empty()) {
|
||||||
|
this->perform_io();
|
||||||
|
|
||||||
|
this->check_for_finished_requests();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
void add_request(curl_request *cr) {
|
||||||
|
mutex_guard mg(this->cl_mutex);
|
||||||
|
|
||||||
|
require(cr != NULL);
|
||||||
|
|
||||||
|
this->cl_all_requests.push_back(cr);
|
||||||
|
this->cl_new_requests.push_back(cr);
|
||||||
|
pthread_cond_broadcast(&this->cl_cond);
|
||||||
|
};
|
||||||
|
|
||||||
|
void close_request(const std::string &name) {
|
||||||
|
mutex_guard mg(this->cl_mutex);
|
||||||
|
|
||||||
|
this->cl_close_requests.push_back(name);
|
||||||
|
pthread_cond_broadcast(&this->cl_cond);
|
||||||
|
};
|
||||||
|
|
||||||
|
private:
|
||||||
|
|
||||||
|
void *run();
|
||||||
|
void loop_body();
|
||||||
|
void perform_io();
|
||||||
|
void check_for_new_requests();
|
||||||
|
void check_for_finished_requests();
|
||||||
|
void requeue_requests(mstime_t up_to_time);
|
||||||
|
|
||||||
|
int compute_timeout(mstime_t current_time) const {
|
||||||
|
int retval = 1000;
|
||||||
|
|
||||||
|
if (!this->cl_poll_queue.empty()) {
|
||||||
|
retval = std::max(
|
||||||
|
1LL, this->cl_poll_queue.front().first - current_time);
|
||||||
|
}
|
||||||
|
|
||||||
|
ensure(retval > 0);
|
||||||
|
|
||||||
|
return retval;
|
||||||
|
};
|
||||||
|
|
||||||
|
static void *trampoline(void *arg);
|
||||||
|
|
||||||
|
bool cl_started;
|
||||||
|
pthread_t cl_thread;
|
||||||
|
volatile bool cl_looping;
|
||||||
|
auto_mem<CURLM> cl_curl_multi;
|
||||||
|
pthread_mutex_t cl_mutex;
|
||||||
|
pthread_cond_t cl_cond;
|
||||||
|
std::vector<curl_request *> cl_all_requests;
|
||||||
|
std::vector<curl_request *> cl_new_requests;
|
||||||
|
std::vector<std::string> cl_close_requests;
|
||||||
|
std::map<CURL *, curl_request *> cl_handle_to_request;
|
||||||
|
std::vector<std::pair<mstime_t, curl_request *> > cl_poll_queue;
|
||||||
|
|
||||||
|
};
|
||||||
|
#endif
|
||||||
|
|
||||||
|
#endif
|
@ -0,0 +1,51 @@
|
|||||||
|
/**
|
||||||
|
* Copyright (c) 2015, Timothy Stack
|
||||||
|
*
|
||||||
|
* All rights reserved.
|
||||||
|
*
|
||||||
|
* Redistribution and use in source and binary forms, with or without
|
||||||
|
* modification, are permitted provided that the following conditions are met:
|
||||||
|
*
|
||||||
|
* * Redistributions of source code must retain the above copyright notice, this
|
||||||
|
* list of conditions and the following disclaimer.
|
||||||
|
* * Redistributions in binary form must reproduce the above copyright notice,
|
||||||
|
* this list of conditions and the following disclaimer in the documentation
|
||||||
|
* and/or other materials provided with the distribution.
|
||||||
|
* * Neither the name of Timothy Stack nor the names of its contributors
|
||||||
|
* may be used to endorse or promote products derived from this software
|
||||||
|
* without specific prior written permission.
|
||||||
|
*
|
||||||
|
* THIS SOFTWARE IS PROVIDED BY THE REGENTS AND CONTRIBUTORS ''AS IS'' AND ANY
|
||||||
|
* EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||||
|
* WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
|
||||||
|
* DISCLAIMED. IN NO EVENT SHALL THE REGENTS OR CONTRIBUTORS BE LIABLE FOR ANY
|
||||||
|
* DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
|
||||||
|
* (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
|
||||||
|
* LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
|
||||||
|
* ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
|
||||||
|
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
|
||||||
|
* SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
||||||
|
*
|
||||||
|
* @file pthreadpp.hh
|
||||||
|
*/
|
||||||
|
|
||||||
|
#ifndef pthreadpp_hh
|
||||||
|
#define pthreadpp_hh
|
||||||
|
|
||||||
|
#include <pthread.h>
|
||||||
|
|
||||||
|
class mutex_guard {
|
||||||
|
public:
|
||||||
|
mutex_guard(pthread_mutex_t &mutex) : mg_mutex(mutex) {
|
||||||
|
pthread_mutex_lock(&mutex);
|
||||||
|
};
|
||||||
|
|
||||||
|
~mutex_guard() {
|
||||||
|
pthread_mutex_unlock(&this->mg_mutex);
|
||||||
|
};
|
||||||
|
|
||||||
|
private:
|
||||||
|
pthread_mutex_t &mg_mutex;
|
||||||
|
};
|
||||||
|
|
||||||
|
#endif
|
@ -0,0 +1,143 @@
|
|||||||
|
/**
|
||||||
|
* Copyright (c) 2015, Timothy Stack
|
||||||
|
*
|
||||||
|
* All rights reserved.
|
||||||
|
*
|
||||||
|
* Redistribution and use in source and binary forms, with or without
|
||||||
|
* modification, are permitted provided that the following conditions are met:
|
||||||
|
*
|
||||||
|
* * Redistributions of source code must retain the above copyright notice, this
|
||||||
|
* list of conditions and the following disclaimer.
|
||||||
|
* * Redistributions in binary form must reproduce the above copyright notice,
|
||||||
|
* this list of conditions and the following disclaimer in the documentation
|
||||||
|
* and/or other materials provided with the distribution.
|
||||||
|
* * Neither the name of Timothy Stack nor the names of its contributors
|
||||||
|
* may be used to endorse or promote products derived from this software
|
||||||
|
* without specific prior written permission.
|
||||||
|
*
|
||||||
|
* THIS SOFTWARE IS PROVIDED BY THE REGENTS AND CONTRIBUTORS ''AS IS'' AND ANY
|
||||||
|
* EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
|
||||||
|
* WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
|
||||||
|
* DISCLAIMED. IN NO EVENT SHALL THE REGENTS OR CONTRIBUTORS BE LIABLE FOR ANY
|
||||||
|
* DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
|
||||||
|
* (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
|
||||||
|
* LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
|
||||||
|
* ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
|
||||||
|
* (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
|
||||||
|
* SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
||||||
|
*/
|
||||||
|
|
||||||
|
#ifndef url_loader_hh
|
||||||
|
#define url_loader_hh
|
||||||
|
|
||||||
|
#ifdef HAVE_LIBCURL
|
||||||
|
#include <curl/curl.h>
|
||||||
|
|
||||||
|
class url_loader : public curl_request {
|
||||||
|
public:
|
||||||
|
url_loader(const std::string &url) : curl_request(url), ul_resume_offset(0) {
|
||||||
|
char piper_tmpname[PATH_MAX];
|
||||||
|
const char *tmpdir;
|
||||||
|
|
||||||
|
if ((tmpdir = getenv("TMPDIR")) == NULL) {
|
||||||
|
tmpdir = _PATH_VARTMP;
|
||||||
|
}
|
||||||
|
snprintf(piper_tmpname, sizeof(piper_tmpname),
|
||||||
|
"%s/lnav.url.XXXXXX",
|
||||||
|
tmpdir);
|
||||||
|
if ((this->ul_fd = mkstemp(piper_tmpname)) == -1) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
unlink(piper_tmpname);
|
||||||
|
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_URL, this->cr_name.c_str());
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_WRITEFUNCTION, write_cb);
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_WRITEDATA, this);
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_FILETIME, 1);
|
||||||
|
};
|
||||||
|
|
||||||
|
int get_fd() const {
|
||||||
|
return this->ul_fd.get();
|
||||||
|
};
|
||||||
|
|
||||||
|
auto_fd copy_fd() const {
|
||||||
|
return this->ul_fd;
|
||||||
|
};
|
||||||
|
|
||||||
|
long complete(CURLcode result) {
|
||||||
|
curl_request::complete(result);
|
||||||
|
|
||||||
|
switch (result) {
|
||||||
|
case CURLE_OK:
|
||||||
|
break;
|
||||||
|
case CURLE_BAD_DOWNLOAD_RESUME:
|
||||||
|
break;
|
||||||
|
default:
|
||||||
|
log_error("%s:curl failure -- %ld %s",
|
||||||
|
this->cr_name.c_str(), result, curl_easy_strerror(result));
|
||||||
|
write(this->ul_fd, this->cr_error_buffer, strlen(this->cr_error_buffer));
|
||||||
|
return -1;
|
||||||
|
}
|
||||||
|
|
||||||
|
long file_time;
|
||||||
|
CURLcode rc;
|
||||||
|
|
||||||
|
rc = curl_easy_getinfo(this->cr_handle, CURLINFO_FILETIME, &file_time);
|
||||||
|
if (rc == CURLE_OK) {
|
||||||
|
time_t current_time;
|
||||||
|
|
||||||
|
time(¤t_time);
|
||||||
|
if (file_time == -1 ||
|
||||||
|
(current_time - file_time) < FOLLOW_IF_MODIFIED_SINCE) {
|
||||||
|
char range[64];
|
||||||
|
struct stat st;
|
||||||
|
off_t start;
|
||||||
|
|
||||||
|
fstat(this->ul_fd, &st);
|
||||||
|
if (st.st_size > 0) {
|
||||||
|
start = st.st_size - 1;
|
||||||
|
this->ul_resume_offset = 1;
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
start = 0;
|
||||||
|
this->ul_resume_offset = 0;
|
||||||
|
}
|
||||||
|
snprintf(range, sizeof(range), "%ld-", (long) start);
|
||||||
|
curl_easy_setopt(this->cr_handle, CURLOPT_RANGE, range);
|
||||||
|
return 2000;
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
log_debug("URL was not recently modified, not tailing: %s",
|
||||||
|
this->cr_name.c_str());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
log_error("Could not get file time for URL: %s -- %s",
|
||||||
|
this->cr_name.c_str(), curl_easy_strerror(rc));
|
||||||
|
}
|
||||||
|
|
||||||
|
return -1;
|
||||||
|
};
|
||||||
|
|
||||||
|
private:
|
||||||
|
static const long FOLLOW_IF_MODIFIED_SINCE = 60 * 60;
|
||||||
|
|
||||||
|
static ssize_t write_cb(void *contents, size_t size, size_t nmemb, void *userp) {
|
||||||
|
url_loader *ul = (url_loader *) userp;
|
||||||
|
char *c_contents = (char *) contents;
|
||||||
|
ssize_t retval;
|
||||||
|
|
||||||
|
c_contents += ul->ul_resume_offset;
|
||||||
|
retval = write(ul->ul_fd, c_contents, (size * nmemb) - ul->ul_resume_offset);
|
||||||
|
retval += ul->ul_resume_offset;
|
||||||
|
ul->ul_resume_offset = 0;
|
||||||
|
return retval;
|
||||||
|
}
|
||||||
|
|
||||||
|
auto_fd ul_fd;
|
||||||
|
off_t ul_resume_offset;
|
||||||
|
};
|
||||||
|
#endif
|
||||||
|
|
||||||
|
#endif
|
@ -0,0 +1,48 @@
|
|||||||
|
#! /bin/bash
|
||||||
|
|
||||||
|
if test x"$SFTP_TEST_URL" == x""; then
|
||||||
|
exit 0
|
||||||
|
fi
|
||||||
|
|
||||||
|
run_test ${lnav_test} -n \
|
||||||
|
file://${test_dir}/logfile_access_log.0
|
||||||
|
|
||||||
|
check_output "file URL is not working" <<EOF
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:26 +0000] "GET /vmw/cgi/tramp HTTP/1.0" 200 134 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkboot.gz HTTP/1.0" 404 46210 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkernel.gz HTTP/1.0" 200 78929 "-" "gPXE/0.9.7"
|
||||||
|
EOF
|
||||||
|
|
||||||
|
cp ${test_dir}/logfile_access_log.0 curl_access_log.0
|
||||||
|
|
||||||
|
run_test ${lnav_test} -n \
|
||||||
|
$SFTP_TEST_URL/`pwd`/curl_access_log.0
|
||||||
|
|
||||||
|
check_output "sftp URL is not working" <<EOF
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:26 +0000] "GET /vmw/cgi/tramp HTTP/1.0" 200 134 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkboot.gz HTTP/1.0" 404 46210 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkernel.gz HTTP/1.0" 200 78929 "-" "gPXE/0.9.7"
|
||||||
|
EOF
|
||||||
|
|
||||||
|
run_test ${lnav_test} -n \
|
||||||
|
-c ":poll-now" \
|
||||||
|
-c ":shexec echo foo >> curl_access_log.0" \
|
||||||
|
$SFTP_TEST_URL/`pwd`/curl_access_log.0
|
||||||
|
|
||||||
|
check_output "sftp URL is not working" <<EOF
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:26 +0000] "GET /vmw/cgi/tramp HTTP/1.0" 200 134 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkboot.gz HTTP/1.0" 404 46210 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkernel.gz HTTP/1.0" 200 78929 "-" "gPXE/0.9.7"
|
||||||
|
foo
|
||||||
|
EOF
|
||||||
|
|
||||||
|
run_test ${lnav_test} -n \
|
||||||
|
-c ":open $SFTP_TEST_URL/`pwd`/curl_access_log.0" \
|
||||||
|
${test_dir}/logfile_empty.0
|
||||||
|
|
||||||
|
check_output "sftp URL is not working" <<EOF
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:26 +0000] "GET /vmw/cgi/tramp HTTP/1.0" 200 134 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkboot.gz HTTP/1.0" 404 46210 "-" "gPXE/0.9.7"
|
||||||
|
192.168.202.254 - - [20/Jul/2009:22:59:29 +0000] "GET /vmw/vSphere/default/vmkernel.gz HTTP/1.0" 200 78929 "-" "gPXE/0.9.7"
|
||||||
|
foo
|
||||||
|
EOF
|
Loading…
Reference in New Issue