diff options
Diffstat (limited to 'main.c')
| -rw-r--r-- | main.c | 309 |
1 files changed, 309 insertions, 0 deletions
@@ -0,0 +1,309 @@ +#include <curl/curl.h> +#include <sqlite3.h> +#include <stddef.h> +#include <stdio.h> +#include <stdlib.h> +#include <string.h> +#include <regex.h> +#include <pthread.h> +#include <unistd.h> + +#include "db_init.h" + +#define N_THREADS 20 +#define QUEUE_MAX 1000 + + +typedef struct { + char *urls[QUEUE_MAX]; + int front; + int rear; + int size; +} queue; + +typedef struct { + sqlite3 *db; + char *url; +} curl_ctx; + +regex_t regex; + +pthread_mutex_t q_lock = PTHREAD_MUTEX_INITIALIZER; +pthread_cond_t q_cond = PTHREAD_COND_INITIALIZER; +queue *q; + +int shutdown_flag = 0; + +void queue_init(queue *q) { + q->front = 0; + q->rear = 0; + q->size = 0; +} + +int enqueue(queue *q, char *c) { + pthread_mutex_lock(&q_lock); + if (q->size == QUEUE_MAX) { + pthread_mutex_unlock(&q_lock); + return 1; + } + q->urls[q->rear] = c; + q->rear = (q->rear + 1) % QUEUE_MAX; + q->size++; + pthread_mutex_unlock(&q_lock); + return 0; +} + +char* dequeue(queue *q) { + pthread_mutex_lock(&q_lock); + while (q->size == 0 && !shutdown_flag) { + pthread_cond_wait(&q_cond, &q_lock); + } + if (q->size == 0 && shutdown_flag) { + pthread_mutex_unlock(&q_lock); + return NULL; + } + char *dqed = q->urls[q->front]; + q->front = (q->front + 1) % QUEUE_MAX; + q->size--; + pthread_mutex_unlock(&q_lock); + return dqed; +} + +void queue_print(queue *q) { + pthread_mutex_lock(&q_lock); + for (size_t i = 0; i < q->size; i++) { + size_t iter = (q->front + i) % QUEUE_MAX; + printf("%s\n", q->urls[iter]); + } + pthread_mutex_unlock(&q_lock); +} + +int db_insert_url(sqlite3 *db, const char *url) { + const char *sql = "INSERT OR IGNORE INTO urls (url, status) VALUES (?, 'pending')"; + sqlite3_stmt *stmt; + + if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) { + fprintf(stderr, "insert url prepare failed: %s\n", sqlite3_errmsg(db)); + return -1; + } + + sqlite3_bind_text(stmt, 1, url, -1, SQLITE_STATIC); + + int rc = sqlite3_step(stmt); + if (rc != SQLITE_DONE) { + fprintf(stderr, "insert failed: %s\n", sqlite3_errmsg(db)); + sqlite3_finalize(stmt); + return -1; + } + sqlite3_finalize(stmt); + + return sqlite3_changes(db); // 1 if inserted, 0 if it already existed +} + +int db_set_status(sqlite3 *db, const char *url, const char *status) { + const char *sql = "UPDATE urls SET status = ? WHERE url = ?"; + sqlite3_stmt *stmt; + + if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) { + fprintf(stderr, "set status prepare failed: %s\n", sqlite3_errmsg(db)); + return -1; + } + + sqlite3_bind_text(stmt, 1, status, -1, SQLITE_STATIC); + sqlite3_bind_text(stmt, 2, url, -1, SQLITE_STATIC); + + int rc = sqlite3_step(stmt); + sqlite3_finalize(stmt); + + if (rc != SQLITE_DONE) { + fprintf(stderr, "update failed: %s\n", sqlite3_errmsg(db)); + return -1; + } + return sqlite3_changes(db); // 1 if a row was updated, 0 if url not found +} + +sqlite3_int64 db_get_url_id(sqlite3 *db, const char *url) { + sqlite3_stmt *stmt; + sqlite3_int64 id = -1; + + if (sqlite3_prepare_v2(db, "SELECT id FROM urls WHERE url = ?", -1, &stmt, NULL) != SQLITE_OK) { + fprintf(stderr, "get url id prepare failed: %s\n", sqlite3_errmsg(db)); + return -1; + } + sqlite3_bind_text(stmt, 1, url, -1, SQLITE_STATIC); + if (sqlite3_step(stmt) == SQLITE_ROW) { + id = sqlite3_column_int64(stmt, 0); + } + sqlite3_finalize(stmt); + return id; +} + +int db_insert_link(sqlite3 *db, const char *from_url, const char *to_url) { + sqlite3_int64 from_id = db_get_url_id(db, from_url); + sqlite3_int64 to_id = db_get_url_id(db, to_url); + + if (from_id < 0 || to_id < 0) { + fprintf(stderr, "db_insert_link: url not found (from=%lld to=%lld)\n", + (long long)from_id, (long long)to_id); + return -1; + } + + sqlite3_stmt *stmt; + if (sqlite3_prepare_v2(db, "INSERT OR IGNORE INTO links (from_id, to_id) VALUES (?, ?)", -1, &stmt, NULL) != SQLITE_OK) { + fprintf(stderr, "insert link prepare failed: %s\n", sqlite3_errmsg(db)); + return -1; + } + sqlite3_bind_int64(stmt, 1, from_id); + sqlite3_bind_int64(stmt, 2, to_id); + + int rc = sqlite3_step(stmt); + sqlite3_finalize(stmt); + + if (rc != SQLITE_DONE) { + fprintf(stderr, "db_insert_link failed: %s\n", sqlite3_errmsg(db)); + return -1; + } + return sqlite3_changes(db); +} + +static int fill_queue_callback(void *q, int argc, char **argv, char **col_name) { + char *enq = strdup(argv[1]); + printf("enqueueing %s\n", enq); + enqueue(q, enq); + return 0; +} + +void fill_queue(sqlite3 *DB) { + pthread_mutex_lock(&q_lock); + int space = QUEUE_MAX - q->size; + pthread_mutex_unlock(&q_lock); + + char sql[128]; + snprintf(sql, sizeof(sql), + "SELECT id, url FROM urls WHERE status = 'pending' LIMIT %d", space); + int rc = sqlite3_exec(DB, sql, fill_queue_callback, q, NULL); + if (rc != SQLITE_OK) { + fprintf(stderr, "enqueue from db failed"); + } + +} + +static size_t curl_callback(char *str, size_t size, size_t nmemb, void *ctx) { + + sqlite3 *db = ((curl_ctx *)ctx)->db; + char *orig_url = ((curl_ctx *)ctx)->url; + size_t n = size * nmemb; + + char *buf = malloc(n + 1); + if (!buf) return n; + memcpy(buf, str, n); + buf[n] = '\0'; + char *cursor = buf; + regmatch_t match; + + // look for urls + while (regexec(®ex, cursor, 1, &match, 0) == 0) { + size_t len = match.rm_eo - match.rm_so; + char *url = malloc(len + 1); + memcpy(url, cursor + match.rm_so, len); + url[len] = '\0'; + + //printf("found url: %s\n", url); + + //insert url into db + int inserted = db_insert_url(db, url); + db_insert_link(db, orig_url, url); + if (inserted) { + if (enqueue(q, url)) free(url); // queue was full, free + } else { + free(url); // was already in db, free + } + + cursor += match.rm_eo; + } + free(buf); + return n; +} + + +void *worker(void *arg) { + sqlite3 *db = db_open_worker("crawl.db"); + printf("worker go\n"); + + while (!shutdown_flag) { + char *url = dequeue(q); + if (!url) break; // shutdown + if (strstr(url, "wikipedia")) { + db_set_status(db, url, "skipped"); + free(url); + continue; //i like wikipedia but not that much + } + + db_set_status(db, url, "in_progress"); + + curl_ctx ctx; + ctx.db = db; + ctx.url = url; + + CURL *curl = curl_easy_init(); + curl_easy_setopt(curl, CURLOPT_URL, url); + curl_easy_setopt(curl, CURLOPT_WRITEDATA, &ctx); + curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, curl_callback); + curl_easy_setopt(curl, CURLOPT_FOLLOWLOCATION, 1L); + curl_easy_setopt(curl, CURLOPT_USERAGENT, "Mozilla/5.0"); + curl_easy_setopt(curl, CURLOPT_TIMEOUT, 30L); + + CURLcode res = curl_easy_perform(curl); + char *status_str = res == CURLE_OK ? "done" : "failed"; + db_set_status(db, url, status_str); + printf("%s: %s\n", status_str, url); + curl_easy_cleanup(curl); + free(url); + } + sqlite3_close(db); + return NULL; +} + +int main(void) { + //sqlite init + sqlite3 *DB = db_init("crawl.db"); + + //curl init + curl_global_init(CURL_GLOBAL_ALL); + + //init queue + q = malloc(sizeof(queue)); + queue_init(q); + fill_queue(DB); + + //init regex + const char *pattern = "https?://[^[:space:]\"<>#]+"; + if (regcomp(®ex, pattern, REG_EXTENDED) != 0) { + printf("Failed to compile regex\n"); + return 1; + } + + //init worker threads + pthread_t threads[N_THREADS]; + for (int i = 0; i < N_THREADS; i++) + pthread_create(&threads[i], NULL, worker, NULL); + + while (1) { + sleep(1); + printf("queue: %d\n", q->size); + if (q->size*5 < QUEUE_MAX) { + fill_queue(DB); + } + } + + //collect workers + pthread_mutex_lock(&q_lock); + shutdown_flag = 1; + pthread_cond_broadcast(&q_cond); + pthread_mutex_unlock(&q_lock); + + sqlite3_close(DB); + curl_global_cleanup(); + regfree(®ex); + return 0; +} |
