Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/datum_gateway.c
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ void datum_pow_tests(void);
void datum_protocol_tests(void);
void datum_stratum_dupes_tests(void);
void datum_utils_tests(void);
void datum_logger_tests(void);
void datum_submitblock_tests(void);

static error_t parse_opt(int key, char *arg, struct argp_state *state) {
Expand All @@ -111,6 +112,7 @@ static error_t parse_opt(int key, char *arg, struct argp_state *state) {
break;
case 0x101: // test
datum_utils_tests();
datum_logger_tests();
datum_conf_tests();
datum_blocktemplates_tests();
datum_coinbaser_tests();
Expand Down
126 changes: 114 additions & 12 deletions src/datum_logger.c
Original file line number Diff line number Diff line change
Expand Up @@ -51,13 +51,15 @@
#include <stdbool.h>
#include <pthread.h>
#include <errno.h>
#include <fcntl.h>
#include <stdatomic.h>

#include "datum_logger.h"
#include "datum_utils.h"

const char *level_text[] = { " ALL", "DEBUG", " INFO", " WARN", "ERROR", "FATAL" };

volatile bool datum_logger_initialized = false;
static atomic_bool datum_logger_initialized = false;
static FILE *log_handle = NULL;
volatile bool log_reopen_signal = false;

Expand Down Expand Up @@ -427,12 +429,18 @@ void datum_logger_hup_signal(int signum) {
log_reopen_signal = true;
}

int datum_logger_init(void) {
const struct sigaction hup_sigaction = { .sa_handler = datum_logger_hup_signal, };
if (0 != sigaction(SIGHUP, &hup_sigaction, NULL)) {
DLOG_ERROR("Failed to setup SIGHUP handler: %s", strerror(errno));
}

static void datum_logger_init_cleanup(void) {
if (log_handle) fclose(log_handle);
log_handle = NULL;
free(dlog_queue[0]);
dlog_queue[0] = dlog_queue[1] = NULL;
free(msg_buffer[0]);
msg_buffer[0] = msg_buffer[1] = NULL;
dlog_queue_max_entries = 0;
}

static int datum_logger_init_with_thread_create(
int (*thread_create)(pthread_t *, const pthread_attr_t *, void *(*)(void *), void *)) {
pthread_t pthread_datum_logger_thread;

// Set the queue and the log file up here, before the writer thread
Expand All @@ -449,21 +457,115 @@ int datum_logger_init(void) {
dlog_queue[0] = calloc(dlog_queue_max_entries * 2 * sizeof(DLOG_MSG),1);
if (!dlog_queue[0]) {
DLOG(DLOG_LEVEL_FATAL, "Could not allocate memory for logger queue list!");
return -1;
goto fail;
}
dlog_queue[1] = &dlog_queue[0][dlog_queue_max_entries];

if ((log_to_file) && (log_file[0] != 0)) {
log_handle = fopen(log_file,"a");
if (!log_handle) {
DLOG(DLOG_LEVEL_FATAL, "Could not open log file (%s): %s!", log_file, strerror(errno));
return -1;
goto fail;
}
}

const int err = thread_create(&pthread_datum_logger_thread, NULL, datum_logger_thread, NULL);
if (err) {
DLOG_FATAL("Could not start logging thread: %s!", strerror(err));
goto fail;
}
// Publish only once a writer exists. Until then, callers log synchronously
// and cannot access queues that a failed startup needs to free.
datum_logger_initialized = true;

pthread_create(&pthread_datum_logger_thread, NULL, datum_logger_thread, NULL);

return 0;

fail:
datum_logger_init_cleanup();
return -1;
}

int datum_logger_init(void) {
const struct sigaction hup_sigaction = { .sa_handler = datum_logger_hup_signal, };
if (0 != sigaction(SIGHUP, &hup_sigaction, NULL)) {
DLOG_ERROR("Failed to setup SIGHUP handler: %s", strerror(errno));
}
return datum_logger_init_with_thread_create(pthread_create);
}

static int logger_test_create_result;
static int logger_test_file_fd;

static int logger_test_thread_create(pthread_t *thread, const pthread_attr_t *attr,
void *(*start)(void *), void *arg) {
(void)thread;
(void)attr;
(void)arg;
datum_test(start == datum_logger_thread);
datum_test(!datum_logger_initialized);
datum_test(msg_buffer[0] && msg_buffer[1] && dlog_queue[0] && dlog_queue[1]);
datum_test(log_handle != NULL);
logger_test_file_fd = log_handle ? fileno(log_handle) : -1;
return logger_test_create_result;
}

void datum_logger_tests(void) {
// Run before the production logger starts; injection never launches a writer.
datum_test(!datum_logger_initialized);
if (datum_logger_initialized) return;
char test_path[] = "/tmp/datum-logger-test-XXXXXX";
const int test_fd = mkstemp(test_path);
FILE *capture = tmpfile();
const int saved_stdout = dup(STDOUT_FILENO);
datum_test(test_fd >= 0 && capture && saved_stdout >= 0);
if (test_fd < 0 || !capture || saved_stdout < 0) {
if (test_fd >= 0) { close(test_fd); unlink(test_path); }
if (capture) fclose(capture);
if (saved_stdout >= 0) close(saved_stdout);
return;
}
close(test_fd);
const bool saved_file = log_to_file, saved_console = log_to_console, saved_stderr = log_to_stderr;
const int saved_level = log_level_console;
char saved_path[sizeof(log_file)];
memcpy(saved_path, log_file, sizeof(saved_path));
log_to_file = log_to_console = true;
log_to_stderr = false;
log_level_console = DLOG_LEVEL_ALL;
strcpy(log_file, test_path);
fflush(stdout);
datum_test(dup2(fileno(capture), STDOUT_FILENO) >= 0);
logger_test_create_result = EAGAIN;
datum_test(datum_logger_init_with_thread_create(logger_test_thread_create) == -1);
datum_test(!datum_logger_initialized && !log_handle);
datum_test(!msg_buffer[0] && !msg_buffer[1] && !dlog_queue[0] && !dlog_queue[1]);
datum_test(dlog_queue_max_entries == 0);
datum_test(fcntl(logger_test_file_fd, F_GETFD) == -1 && errno == EBADF);
DLOG_INFO("logger failure fallback test");
datum_test(dlog_queue_next[0] == 0 && dlog_queue_next[1] == 0);
fflush(stdout);
rewind(capture);
char output[1024] = {0};
fread(output, 1, sizeof(output) - 1, capture);
datum_test(strstr(output, "Could not start logging thread:") != NULL);
datum_test(strstr(output, "logger failure fallback test") != NULL);
datum_test(dup2(saved_stdout, STDOUT_FILENO) >= 0);
close(saved_stdout);
fclose(capture);

// Successful init has usable queues immediately, before the writer runs.
logger_test_create_result = 0;
datum_test(datum_logger_init_with_thread_create(logger_test_thread_create) == 0);
datum_test(datum_logger_initialized);
DLOG_INFO("logger immediate queue test");
datum_test(dlog_queue_next[0] == 1);
datum_test(!strcmp(dlog_queue[0][0].msg, "logger immediate queue test"));
datum_logger_initialized = false;
dlog_queue_next[0] = msg_buf_idx[0] = 0;
datum_logger_init_cleanup();
log_to_file = saved_file;
log_to_console = saved_console;
log_to_stderr = saved_stderr;
log_level_console = saved_level;
memcpy(log_file, saved_path, sizeof(log_file));
unlink(test_path);
}
5 changes: 5 additions & 0 deletions src/datum_stratum.c
Original file line number Diff line number Diff line change
Expand Up @@ -1126,6 +1126,11 @@ int client_mining_submit(T_DATUM_CLIENT_DATA *c, uint64_t id, json_t *params_obj

// need to build the full coinbase txn
coinbase_index = job_id_bin[7];
// Quick difficulty jobs use Q instead of N; the empty coinbase index
// still identifies subsidy-only work while the coinbaser is pending.
if (quickdiff && coinbase_index == DATUM_COINBASE_ID_EMPTY) {
empty_work = true;
}
if (coinbase_index >= MAX_COINBASE_TYPES) {
if (!(empty_work && coinbase_index == DATUM_COINBASE_ID_EMPTY)) {
send_unknown_work_error(c, id);
Expand Down
83 changes: 83 additions & 0 deletions src/datum_stratum_tests.c
Original file line number Diff line number Diff line change
Expand Up @@ -298,6 +298,88 @@ static void datum_blake2b_coinbase_selection_tests(void) {
free(sdata);
}

static void datum_blake2b_quickdiff_coinbase_tests(void) {
T_DATUM_THREAD_DATA *thread = calloc(1, sizeof(*thread));
T_DATUM_STRATUM_THREADPOOL_DATA *sdata = calloc(1, sizeof(*sdata));
T_DATUM_CLIENT_DATA client = {0};
T_DATUM_MINER_DATA miner = {0};
T_DATUM_STRATUM_JOB job = {0};
T_DATUM_TEMPLATE_DATA tdata = {0};
T_DATUM_STRATUM_JOB *saved_job = global_cur_stratum_jobs[0];
const int saved_active = atomic_load(&datum_protocol_client_active);
static const struct {
bool ready, quickdiff;
const char *job_id;
} cases[] = {
{false, false, "N0000000000c0deff"},
{false, true, "Q0000000000c0deff"},
{true, true, "Q0000000000c0de04"},
};
static const char expected[] =
"{\"error\":[23,\"H-not-zero\",null],\"id\":42,\"result\":null}\n";

if (!datum_test(thread && sdata)) {
free(thread);
free(sdata);
return;
}
thread->app_thread_data = sdata;
client.datum_thread = thread;
client.app_client_data = &miner;
miner.sdata = sdata;
sdata->cur_stratum_job = &job;
job.block_template = &tdata;
job.job_state = JOB_STATE_FULL_PRIORITY_WAIT_COINBASER;
strcpy(job.job_id, "0000000000c0de");
strcpy(job.version, "20000000");
job.subsidy_only_coinbase.coinb1_len = 1;
job.subsidy_only_coinbase.coinb1_bin[0] = 0xff;
job.coinbase[COINBASE_TYPE_YUGE] = job.subsidy_only_coinbase;
tdata.curtime = 1000;
tdata.version = 0x20000000;
tdata.bits_uint = 0x1d00ffff;
datum_stratum_job_refresh_blake2b(&job);
global_cur_stratum_jobs[0] = &job;
atomic_store(&datum_protocol_client_active, 3);
if (!datum_protocol_is_active()) {
unsigned char notice[36] = {DATUM_ABW_DRAFT_REVISION, DATUM_ABW_ASSIGNMENT_ACTIVE, 0};
memset(notice + 3, 0x5a, 32);
notice[35] = 0xfe;
datum_test(datum_protocol_abw_assignment_notice(sizeof(notice), notice) == 1);
}
datum_test(datum_protocol_is_active());

for (size_t i = 0; i < sizeof(cases) / sizeof(cases[0]); ++i) {
sdata->full_coinbase_ready = cases[i].ready;
miner.current_diff = miner.last_sent_diff = 2ULL << i;
miner.stratum_job_diffs[0] = 1;
client.out_buf = 0;
datum_test(send_mining_notify(&client, true, cases[i].quickdiff, false) == 0);
json_t *notify = json_loadb(client.w_buffer, client.out_buf, 0, NULL);
const char *job_id = json_string_value(json_array_get(json_object_get(notify, "params"), 0));
if (datum_test(job_id != NULL)) {
datum_test(!strcmp(job_id, cases[i].job_id));
datum_test(miner.quickdiff_active == cases[i].quickdiff);
datum_test(miner.stratum_job_diffs[0] == (cases[i].quickdiff ? 1 : miner.last_sent_diff));
if (cases[i].quickdiff) datum_test(miner.quickdiff_value == miner.last_sent_diff);
json_t *params = json_pack("[sssss]", "miner", job_id,
"0000000000000000", "00000000", "00000000");
client.out_buf = 0;
datum_test(client_mining_submit(&client, 42, params) == 0);
// The advertised job must reach proof validation, not unknown-work.
datum_test(client.out_buf == (int)strlen(expected));
datum_test(!memcmp(client.w_buffer, expected, strlen(expected)));
json_decref(params);
}
json_decref(notify);
}
datum_protocol_abw_reset();
atomic_store(&datum_protocol_client_active, saved_active);
global_cur_stratum_jobs[0] = saved_job;
free(thread);
free(sdata);
}

static void datum_stratum_abw_block_request_tests(void) {
unsigned char xor_key[16];
unsigned char raw_hash[32], masked_hash[32];
Expand Down Expand Up @@ -588,6 +670,7 @@ void datum_stratum_tests(void) {
datum_stratum_minimum_difficulty_configure_tests();
datum_stratum_string_request_id_tests();
datum_blake2b_coinbase_selection_tests();
datum_blake2b_quickdiff_coinbase_tests();
datum_blake2b_h_not_zero_tests();
datum_blake2b_malformed_submit_job_tests();
datum_blake2b_client_pot_commitment_tests();
Expand Down
3 changes: 2 additions & 1 deletion src/datum_utils.c
Original file line number Diff line number Diff line change
Expand Up @@ -369,8 +369,9 @@ bool hex_to_bin_exact(const char *hex, unsigned char *bin, const size_t bin_len)
if (!hex || !bin) return false;
for (size_t i = 0; i < bin_len; i++) {
high = hex_value(hex[i<<1]);
if (high < 0) return false;
low = hex_value(hex[(i<<1)+1]);
if (high < 0 || low < 0) return false;
if (low < 0) return false;
bin[i] = (unsigned char)((high << 4) | low);
}
return hex[bin_len<<1] == 0;
Expand Down
13 changes: 13 additions & 0 deletions src/datum_utils_tests.c
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
#include <stddef.h>
#include <stdint.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <strings.h>

Expand Down Expand Up @@ -116,6 +117,18 @@ void datum_utils_tests_hex(void) {
datum_test(!hex_to_u32("123456789", &value));
datum_test(!hex_to_u32(NULL, &value));
datum_test(!hex_to_u32("00000000", NULL));

/* Exact allocations expose reads past the terminator under ASan. */
for (size_t len = 0; len < 8; ++len) {
unsigned char bin[4];
char *truncated = malloc(len + 1);
if (!datum_test(truncated != NULL)) break;
memcpy(truncated, "1234aBcD", len);
truncated[len] = '\0';
datum_test(!hex_to_bin_exact(truncated, bin, sizeof(bin)));
datum_test(!hex_to_u32(truncated, &value));
free(truncated);
}
}

void datum_utils_tests_secure_strequals(void) {
Expand Down