Skip to content
Open
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
65 changes: 54 additions & 11 deletions src/datum_stratum_dupes.c
Original file line number Diff line number Diff line change
Expand Up @@ -59,15 +59,26 @@ void datum_stratum_dupes_init(void *sdata_v) {
}

dupes = sdata->dupes;

dupes->ptr = calloc((datum_config.stratum_v1_max_clients_per_thread * datum_config.stratum_v1_vardiff_target_shares_min * (datum_config.stratum_v1_share_stale_seconds/60) * 16), sizeof(T_DATUM_STRATUM_DUPE_ITEM) );

// Sized once, then allocated, so the array and max_items cannot disagree.
//
// The floor is because stratum.max_clients_per_thread is range checked for an upper
// bound and not a lower one, and this is the product of three configured values. A zero
// or a negative sizes the table at nothing, and nothing is not a table that merely
// overflows quickly: the expand grows it by 25%, 25% of zero is zero, and the gateway
// takes a share it then has nowhere to put. Sixteen is also the point below which the
// same rounding stops the table growing at all.
int max_items = datum_config.stratum_v1_max_clients_per_thread * datum_config.stratum_v1_vardiff_target_shares_min * (datum_config.stratum_v1_share_stale_seconds/60) * 16;
if (max_items < 16) max_items = 16;

dupes->ptr = calloc(max_items, sizeof(T_DATUM_STRATUM_DUPE_ITEM) );
if (!dupes->ptr) {
DLOG_FATAL("Could not allocate RAM for dupe struct (big one, %lu bytes)",(unsigned long)(datum_config.stratum_v1_max_clients_per_thread * datum_config.stratum_v1_vardiff_target_shares_min * (datum_config.stratum_v1_share_stale_seconds/60) * 16) * sizeof(T_DATUM_STRATUM_DUPE_ITEM));
DLOG_FATAL("Could not allocate RAM for dupe struct (big one, %lu bytes)",(unsigned long)max_items * sizeof(T_DATUM_STRATUM_DUPE_ITEM));
panic_from_thread(__LINE__);
return;
}
dupes->max_items = (datum_config.stratum_v1_max_clients_per_thread * datum_config.stratum_v1_vardiff_target_shares_min * (datum_config.stratum_v1_share_stale_seconds/60) * 16);

dupes->max_items = max_items;
dupes->current_items = 0;

DLOG_DEBUG("Initialized dupe check thread data. %"PRIu64" bytes of RAM used for %d max entries @ %p for %p", (uint64_t)dupes->max_items * (uint64_t)sizeof(T_DATUM_STRATUM_DUPE_ITEM), dupes->max_items, dupes, sdata);
Expand Down Expand Up @@ -207,6 +218,12 @@ void datum_stratum_dupes_cleanup(T_DATUM_STRATUM_DUPES *dupes, bool full_wipe) {
if (full_wipe) {
// we're just cleaning up after a new block or whatever
memset(dupes->ptr, 0, sizeof(T_DATUM_STRATUM_DUPE_ITEM) * dupes->max_items);
// The buckets have to go with the items they point at. Leaving them meant every
// bucket still named a slot that had just been zeroed and was about to be handed
// out again to a new entry, so the first share on such a bucket could link a slot
// to itself and the next walk of that chain would never terminate. Nothing calls
// this with full_wipe today, which is the only reason that has not been seen.
memset(dupes->index, 0, sizeof(dupes->index));
dupes->current_items = 0;
return;
}
Expand Down Expand Up @@ -249,6 +266,20 @@ void datum_stratum_dupes_cleanup(T_DATUM_STRATUM_DUPES *dupes, bool full_wipe) {
T_DATUM_STRATUM_DUPE_ITEM *datum_stratum_add_new_dupe(T_DATUM_STRATUM_DUPES *dupes, uint64_t nonce, unsigned short job_index, uint64_t ntime_val, unsigned int version_bits, unsigned char *extranonce_bin, T_DATUM_STRATUM_DUPE_ITEM *insert_after) {
T_DATUM_STRATUM_DUPE_ITEM *i;

// The caller makes room before it walks the list, so this is a bug rather than a
// full table. Refusing the entry loses one share's dupe protection; writing past
// the end of the array corrupts the heap.
if (dupes->current_items >= dupes->max_items) {
// Once per run, not once per share. This fires on the share path, so a bug that
// made it reachable would otherwise write a line per share submitted.
static bool reported = false;
if (!reported) {
reported = true;
DLOG_ERROR("Dupe table full at insert (%d/%d); dropping the entry rather than writing past it. This should not be reachable; please report it.", dupes->current_items, dupes->max_items);
}
return NULL;
}

i = &dupes->ptr[dupes->current_items];
if (!i) {
DLOG_FATAL("Could not add entry to dupe table!");
Expand All @@ -269,11 +300,14 @@ T_DATUM_STRATUM_DUPE_ITEM *datum_stratum_add_new_dupe(T_DATUM_STRATUM_DUPES *dup
insert_after->next = i;
}
dupes->current_items++;

if (dupes->current_items >= dupes->max_items) {
datum_stratum_dupes_cleanup(dupes, false);
}


// The cleanup that used to be here ran between taking this pointer and returning it,
// and both of the things it can do invalidate it: the sort moves every item, and the
// expand reallocates the array. The caller stores what it gets back into the bucket
// index, so the index ended up holding a pointer into the freed array and the next
// share on that nonce read it. It has moved to the top of datum_stratum_check_for_dupe,
// which is the only place there is no insertion point to invalidate.

return i;
}

Expand All @@ -293,7 +327,14 @@ bool datum_stratum_check_for_dupe(T_DATUM_STRATUM_THREADPOOL_DATA *t, uint64_t n
}

dupes = t->dupes;


// Make room before reading anything out of the table. Everything below this line
// either holds a pointer into the array or an insertion point in a bucket, and a
// cleanup invalidates both, so this is the last moment it can safely run.
if (dupes->current_items >= dupes->max_items) {
datum_stratum_dupes_cleanup(dupes, false);
}

if (dupes->index[nonce_index] == NULL) {
// first nonce of its kind!
// not a duplicate
Expand All @@ -314,6 +355,8 @@ bool datum_stratum_check_for_dupe(T_DATUM_STRATUM_THREADPOOL_DATA *t, uint64_t n
} else {
// we need to replace the first item in a list, so... let's make a new entry
p = datum_stratum_add_new_dupe(dupes, nonce, job_index, ntime_val, version_bits, extranonce_bin, NULL);
// A refused entry leaves the bucket as it was rather than unlinking it
if (!p) return false;
dupes->index[nonce_index] = p;
p->next = i;
}
Expand Down
225 changes: 225 additions & 0 deletions src/datum_stratum_dupes_tests.c
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@

#include <stdlib.h>

#include "datum_conf.h"
#include "datum_stratum.h"
#include "datum_stratum_dupes.h"
#include "datum_utils.h"
Expand Down Expand Up @@ -69,6 +70,230 @@ static void datum_pow_dupe_tests(void) {
free(thread_data);
}

// Fill the table past max_items, which is the only way the cleanup/expand path runs.
// Nothing above reaches it: that test has 8 slots and inserts 2.
static void datum_dupe_table_fill_tests(void) {
const int saved_clients = datum_config.stratum_v1_max_clients_per_thread;
const int saved_shares = datum_config.stratum_v1_vardiff_target_shares_min;
const int saved_stale = datum_config.stratum_v1_share_stale_seconds;

// max_items = clients * shares_min * (stale/60) * 16, so this is a 16 slot table
datum_config.stratum_v1_max_clients_per_thread = 1;
datum_config.stratum_v1_vardiff_target_shares_min = 1;
datum_config.stratum_v1_share_stale_seconds = 60;

T_DATUM_STRATUM_THREADPOOL_DATA * const thread_data = calloc(1, sizeof(*thread_data));
datum_test(thread_data != NULL);
if (!thread_data) return;
datum_stratum_dupes_init(thread_data);

// One live job. Nothing ages out, so the cleanup has nothing to prune and takes the
// expand path, which is what a busy thread does between block changes.
T_DATUM_STRATUM_JOB * const job = calloc(1, sizeof(*job));
T_DATUM_STRATUM_JOB * const saved_job = global_cur_stratum_jobs[1];
datum_test(job != NULL);
if (!job) { free(thread_data); return; }
job->tsms = current_time_millis();
global_cur_stratum_jobs[1] = job;

unsigned char extranonce[12] = {0};

// Distinct low 16 bits, so every share lands in a bucket of its own and each insert
// is the "first nonce of its kind" case.
for (int i = 0; i < 64; ++i) {
const uint64_t nonce = ((uint64_t)i << 32) | (uint64_t)((i * 7) + 1);
datum_test(!datum_stratum_check_for_dupe(thread_data, nonce, 1, 1000 + i, 0, extranonce));
}

// The table has to have actually grown, or nothing above went through the path this
// test exists for and the assertions below prove nothing.
T_DATUM_STRATUM_DUPES * const dupes = thread_data->dupes;
datum_test(dupes->max_items > 16);
datum_test(dupes->current_items <= dupes->max_items);

// Every one of those is a duplicate now. Re-probing walks each bucket from its index
// entry, which is where a pointer left over from before a reallocation is read.
for (int i = 0; i < 64; ++i) {
const uint64_t nonce = ((uint64_t)i << 32) | (uint64_t)((i * 7) + 1);
datum_test(datum_stratum_check_for_dupe(thread_data, nonce, 1, 1000 + i, 0, extranonce));
}

global_cur_stratum_jobs[1] = saved_job;
free(dupes->ptr);
free(thread_data->dupes);
free(thread_data);
free(job);

datum_config.stratum_v1_max_clients_per_thread = saved_clients;
datum_config.stratum_v1_vardiff_target_shares_min = saved_shares;
datum_config.stratum_v1_share_stale_seconds = saved_stale;
}

// Every pointer reachable from the bucket index must name a live entry, and no chain may
// loop. A chain that loops never terminates in datum_stratum_check_for_dupe, which runs on
// the stratum thread with that thread's clients waiting on it.
static void datum_dupe_index_is_sound(const T_DATUM_STRATUM_DUPES *dupes) {
for (int b = 0; b < 65536; ++b) {
const T_DATUM_STRATUM_DUPE_ITEM *i = dupes->index[b];
int walked = 0;
while (i) {
const ptrdiff_t slot = i - dupes->ptr;
// Named rather than asserted inline, because datum_test reports the expression
// it was given and these are what the reader wants to see in a failure
const bool bucket_points_at_a_live_entry =
slot >= 0 && slot < dupes->current_items;
if (!datum_test(bucket_points_at_a_live_entry)) return;
const bool bucket_chain_terminates = ++walked <= dupes->current_items;
if (!datum_test(bucket_chain_terminates)) return;
i = i->next;
}
}
}

// The other half of the cleanup, and the half a gateway actually reaches: entries old
// enough to age out are pruned rather than the array being grown. The prune sorts the
// array, which moves every entry, so an insertion point taken before it is stale after it
// in the same way a reallocation makes one stale.
static void datum_dupe_table_prune_tests(void) {
const int saved_clients = datum_config.stratum_v1_max_clients_per_thread;
const int saved_shares = datum_config.stratum_v1_vardiff_target_shares_min;
const int saved_stale = datum_config.stratum_v1_share_stale_seconds;

datum_config.stratum_v1_max_clients_per_thread = 1;
datum_config.stratum_v1_vardiff_target_shares_min = 1;
datum_config.stratum_v1_share_stale_seconds = 60;

T_DATUM_STRATUM_THREADPOOL_DATA * const thread_data = calloc(1, sizeof(*thread_data));
datum_test(thread_data != NULL);
if (!thread_data) return;
datum_stratum_dupes_init(thread_data);
T_DATUM_STRATUM_DUPES * const dupes = thread_data->dupes;

// Job 1 is current, job 2 is old enough for its shares to age out. A real gateway has
// both at once: the jobs it is handing out now, and the ones from a few minutes ago.
T_DATUM_STRATUM_JOB * const fresh = calloc(1, sizeof(*fresh));
T_DATUM_STRATUM_JOB * const old = calloc(1, sizeof(*old));
T_DATUM_STRATUM_JOB * const saved_fresh = global_cur_stratum_jobs[1];
T_DATUM_STRATUM_JOB * const saved_old = global_cur_stratum_jobs[2];
datum_test(fresh != NULL && old != NULL);
if (!fresh || !old) { free(fresh); free(old); free(thread_data); return; }
fresh->tsms = current_time_millis();
old->tsms = current_time_millis() - 600000;
global_cur_stratum_jobs[1] = fresh;
global_cur_stratum_jobs[2] = old;

unsigned char extranonce[12] = {0};

// Mostly stale, so the cleanup frees well over its 5% and takes the prune path
for (int i = 0; i < 256; ++i) {
const uint64_t nonce = ((uint64_t)i << 32) | (uint64_t)((i * 11) + 3);
const unsigned short job = (i % 8) ? 2 : 1;
datum_stratum_check_for_dupe(thread_data, nonce, job, 2000 + i, 0, extranonce);
datum_dupe_index_is_sound(dupes);
}

global_cur_stratum_jobs[1] = saved_fresh;
global_cur_stratum_jobs[2] = saved_old;
free(dupes->ptr);
free(thread_data->dupes);
free(thread_data);
free(fresh);
free(old);

datum_config.stratum_v1_max_clients_per_thread = saved_clients;
datum_config.stratum_v1_vardiff_target_shares_min = saved_shares;
datum_config.stratum_v1_share_stale_seconds = saved_stale;
}

// Many cleanup cycles of both kinds, against the invariant the insert now relies on:
// datum_stratum_check_for_dupe must leave room for the next entry, every time. If it ever
// does not, datum_stratum_add_new_dupe refuses a share's dupe record and says so in the
// log, which is a thing an operator should never see.
//
// Both kinds matter because they fail differently. A prune sorts the array in place and an
// expand reallocates it, and before the fix each left a different flavour of stale pointer
// in the bucket index.
static void datum_dupe_table_cycle_tests(void) {
const int saved_clients = datum_config.stratum_v1_max_clients_per_thread;
const int saved_shares = datum_config.stratum_v1_vardiff_target_shares_min;
const int saved_stale = datum_config.stratum_v1_share_stale_seconds;

datum_config.stratum_v1_max_clients_per_thread = 2;
datum_config.stratum_v1_vardiff_target_shares_min = 2;
datum_config.stratum_v1_share_stale_seconds = 60;

T_DATUM_STRATUM_THREADPOOL_DATA * const thread_data = calloc(1, sizeof(*thread_data));
datum_test(thread_data != NULL);
if (!thread_data) return;
datum_stratum_dupes_init(thread_data);
T_DATUM_STRATUM_DUPES * const dupes = thread_data->dupes;
const int initial_max = dupes->max_items;

// Job 0 included deliberately. It is a real job slot on a running gateway, and a zeroed
// table entry reads as job_index 0, so the sort consults job 0 for entries that are not
// entries. Keeping it live here is the arrangement that would expose that if it bit.
T_DATUM_STRATUM_JOB * const jobs = calloc(4, sizeof(*jobs));
T_DATUM_STRATUM_JOB *saved[4];
datum_test(jobs != NULL);
if (!jobs) { free(thread_data); return; }
for (int j = 0; j < 4; ++j) {
saved[j] = global_cur_stratum_jobs[j];
global_cur_stratum_jobs[j] = &jobs[j];
}

unsigned char extranonce[12] = {0};
bool saw_expand = false;

for (int i = 0; i < 20000; ++i) {
// Jobs age as the run goes on, so cleanups alternate between having plenty to
// prune and having nothing to prune and needing to grow.
const uint64_t now = current_time_millis();
jobs[0].tsms = now;
jobs[1].tsms = now;
jobs[2].tsms = (i % 3) ? now - 600000 : now;
jobs[3].tsms = (i % 7) ? now - 600000 : now;

const uint64_t nonce = ((uint64_t)i * 2654435761u) ^ ((uint64_t)i << 24);
extranonce[0] = (unsigned char)i;
extranonce[11] = (unsigned char)(i >> 8);

// A table that is not yet full cannot clean up during the call, so a share that is
// not a duplicate has to become exactly one new entry. Anything else means the
// insert refused it, which is the only way the fix can go wrong quietly: the share
// is still accepted, but nothing remembers it and a real resubmission slips past.
const int before = dupes->current_items;
const bool had_room = before < dupes->max_items;
const bool dupe = datum_stratum_check_for_dupe(thread_data, nonce,
(unsigned short)(i & 3), 3000 + (i & 0xff), (unsigned int)i, extranonce);
if (had_room && !dupe) {
const bool the_new_share_became_an_entry =
dupes->current_items == before + 1;
datum_test(the_new_share_became_an_entry);
}

// Never past the end of the array, cleanup or no cleanup
const bool entries_fit_the_array = dupes->current_items <= dupes->max_items;
datum_test(entries_fit_the_array);
if (dupes->max_items > initial_max) saw_expand = true;
}

datum_dupe_index_is_sound(dupes);
datum_test(saw_expand);

for (int j = 0; j < 4; ++j) global_cur_stratum_jobs[j] = saved[j];
free(dupes->ptr);
free(thread_data->dupes);
free(thread_data);
free(jobs);

datum_config.stratum_v1_max_clients_per_thread = saved_clients;
datum_config.stratum_v1_vardiff_target_shares_min = saved_shares;
datum_config.stratum_v1_share_stale_seconds = saved_stale;
}

void datum_stratum_dupes_tests(void) {
datum_pow_dupe_tests();
datum_dupe_table_fill_tests();
datum_dupe_table_prune_tests();
datum_dupe_table_cycle_tests();
}