Skip to content
Open
Show file tree
Hide file tree
Changes from 3 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
1 change: 1 addition & 0 deletions components/stratum/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ SRCS
"mining.c"
"mining_job.c"
"stratum_api.c"
"stratum_timing.c"
"stratum_socket.c"
"coinbase_decoder.c"
"segwit_addr.c"
Expand Down
7 changes: 0 additions & 7 deletions components/stratum/include/stratum_api.h
Original file line number Diff line number Diff line change
Expand Up @@ -87,11 +87,6 @@ typedef struct sv1_conn {
int active_job_ids_count;
} sv1_conn_t;

typedef struct RequestTiming
{
int64_t timestamp_us;
bool tracking;
} RequestTiming;

esp_transport_handle_t STRATUM_V1_transport_init(tls_mode tls, const char * cert);

Expand Down Expand Up @@ -121,6 +116,4 @@ int STRATUM_V1_submit_share(esp_transport_handle_t transport, int send_uid, cons
const char *extranonce_2, const uint32_t ntime, const uint32_t nonce,
const uint32_t version_bits, uint64_t *out_sent_time_us);

float STRATUM_V1_get_response_time_ms(int request_id, int64_t receive_time_us);

#endif // STRATUM_API_H
39 changes: 39 additions & 0 deletions components/stratum/include/stratum_timing.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
#ifndef STRATUM_TIMING_H
#define STRATUM_TIMING_H

#include <stdint.h>
#include <stdbool.h>

#define STRATUM_TIMING_SLOTS 64
Comment thread
mutatrum marked this conversation as resolved.

typedef struct {
uint64_t submit_time_us[STRATUM_TIMING_SLOTS];
} stratum_timing_tracker_t;

/**
* @brief Record a submission timestamp for a given request / share ID.
*
* @param tracker Pointer to the timing tracker.
* @param id The request ID / share ID / sequence number.
* @param submit_time_us The submission timestamp in microseconds.
*/
void stratum_timing_record(stratum_timing_tracker_t *tracker, uint32_t id, uint64_t submit_time_us);

/**
* @brief Calculate round-trip response time in milliseconds for an ID, and clear the slot.
*
* @param tracker Pointer to the timing tracker.
* @param id The response ID / share ID / sequence number.
* @param now_us The current timestamp in microseconds.
* @return Latency in milliseconds, or -1.0f if the ID was not tracked or invalid.
*/
float stratum_timing_calculate_ms(stratum_timing_tracker_t *tracker, uint32_t id, uint64_t now_us);

/**
* @brief Reset/clear all tracking slots (e.g. on disconnect or pool switch).
*
* @param tracker Pointer to the timing tracker.
*/
void stratum_timing_reset(stratum_timing_tracker_t *tracker);

#endif /* STRATUM_TIMING_H */
61 changes: 1 addition & 60 deletions components/stratum/stratum_api.c
Original file line number Diff line number Diff line change
Expand Up @@ -38,28 +38,6 @@ static char * json_rpc_buffer = NULL;
static size_t json_rpc_buffer_size = 0;
static size_t json_rpc_buffer_len = 0;

static RequestTiming *request_timings = NULL;

static RequestTiming* get_request_timing(int request_id) {
if (request_id < 0) return NULL;
int index = request_id % MAX_REQUEST_IDS;
return &request_timings[index];
}

float STRATUM_V1_get_response_time_ms(int request_id, int64_t receive_time_us)
{
if (request_id < 0) return -1.0;

RequestTiming *timing = get_request_timing(request_id);
if (!timing || !timing->tracking) {
return -1.0;
}

float response_time = (receive_time_us - timing->timestamp_us) / 1000.0f;
timing->tracking = false;
return response_time;
}

esp_transport_handle_t STRATUM_V1_transport_init(tls_mode tls, const char * cert)
{
esp_transport_handle_t transport;
Expand Down Expand Up @@ -116,25 +94,6 @@ bool STRATUM_V1_initialize_buffer(void)
json_rpc_buffer_size = BUFFER_SIZE;
json_rpc_buffer[0] = '\0';

if (request_timings == NULL) {
request_timings = heap_caps_malloc(sizeof(RequestTiming) * MAX_REQUEST_IDS, MALLOC_CAP_SPIRAM);
if (request_timings == NULL) {
request_timings = malloc(sizeof(RequestTiming) * MAX_REQUEST_IDS);
}
if (request_timings == NULL) {
ESP_LOGE(TAG, "Failed to allocate memory for request_timings");
free(json_rpc_buffer);
json_rpc_buffer = NULL;
json_rpc_buffer_size = 0;
return false;
}
}

for (int i = 0; i < MAX_REQUEST_IDS; i++) {
request_timings[i].timestamp_us = 0;
request_timings[i].tracking = false;
}

return true;
}

Expand Down Expand Up @@ -773,19 +732,6 @@ bool STRATUM_V1_parse(StratumApiV1Message *message, const char *stratum_json, mi
return result;
}



static void stamp_tx(int request_id, uint64_t timestamp_us)
{
if (request_id >= 1) {
RequestTiming *timing = get_request_timing(request_id);
if (timing) {
timing->timestamp_us = timestamp_us;
timing->tracking = true;
}
}
}

static void debug_stratum_tx(const char * msg)
{
char *newline = strchr(msg, '\n');
Expand Down Expand Up @@ -887,14 +833,9 @@ int STRATUM_V1_submit_share(esp_transport_handle_t transport, int send_uid, cons

int ret = esp_transport_write(transport, submit_msg, strlen(submit_msg), TRANSPORT_TIMEOUT_MS);

uint64_t now = esp_timer_get_time();
if (out_sent_time_us) {
*out_sent_time_us = now;
}
if (out_sent_time_us) *out_sent_time_us = esp_timer_get_time();

debug_stratum_tx(submit_msg);

stamp_tx(send_uid, now);

return ret;
}
Expand Down
32 changes: 32 additions & 0 deletions components/stratum/stratum_timing.c
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
#include "stratum_timing.h"
#include <string.h>

void stratum_timing_record(stratum_timing_tracker_t *tracker, uint32_t id, uint64_t submit_time_us)
{
if (!tracker) {
return;
}
tracker->submit_time_us[id % STRATUM_TIMING_SLOTS] = submit_time_us;
}

float stratum_timing_calculate_ms(stratum_timing_tracker_t *tracker, uint32_t id, uint64_t now_us)
{
if (!tracker) {
return -1.0f;
}
uint32_t slot = id % STRATUM_TIMING_SLOTS;
uint64_t submit_time_us = tracker->submit_time_us[slot];
if (submit_time_us == 0 || now_us < submit_time_us) {
return -1.0f;
}
tracker->submit_time_us[slot] = 0;
return (float)(now_us - submit_time_us) / 1000.0f;
}

void stratum_timing_reset(stratum_timing_tracker_t *tracker)
{
if (!tracker) {
return;
}
memset(tracker->submit_time_us, 0, sizeof(tracker->submit_time_us));
}
77 changes: 77 additions & 0 deletions components/stratum/test/test_stratum_timing.c
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
#include "unity.h"
#include "stratum_timing.h"
#include <string.h>

TEST_CASE("stratum_timing record and calculate latency", "[stratum_timing]")
{
stratum_timing_tracker_t tracker;
stratum_timing_reset(&tracker);

// Record request 1 at t = 1,000,000 us (1.0s)
stratum_timing_record(&tracker, 1, 1000000ULL);

// Calculate response at t = 1,025,000 us (25ms later)
float latency_ms = stratum_timing_calculate_ms(&tracker, 1, 1025000ULL);
TEST_ASSERT_FLOAT_WITHIN(0.01f, 25.0f, latency_ms);

// Second call should return -1.0f as the slot is cleared
latency_ms = stratum_timing_calculate_ms(&tracker, 1, 1030000ULL);
TEST_ASSERT_EQUAL_FLOAT(-1.0f, latency_ms);
}

TEST_CASE("stratum_timing wraparound and slot reuse", "[stratum_timing]")
{
stratum_timing_tracker_t tracker;
stratum_timing_reset(&tracker);

uint32_t id_first = 5;
uint32_t id_second = 5 + STRATUM_TIMING_SLOTS; // Wraps around to the exact same slot

// Record first ID
stratum_timing_record(&tracker, id_first, 2000000ULL);

// Second ID wraps around and overwrites the slot
stratum_timing_record(&tracker, id_second, 2100000ULL);

// Calculate for the second ID
float latency_ms = stratum_timing_calculate_ms(&tracker, id_second, 2115500ULL);
TEST_ASSERT_FLOAT_WITHIN(0.01f, 15.5f, latency_ms);

// Slot is now cleared, checking first or second ID yields -1.0f
TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(&tracker, id_first, 2120000ULL));
TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(&tracker, id_second, 2120000ULL));
}

TEST_CASE("stratum_timing reset", "[stratum_timing]")
{
stratum_timing_tracker_t tracker;
stratum_timing_reset(&tracker);

stratum_timing_record(&tracker, 10, 5000000ULL);
stratum_timing_record(&tracker, 11, 5000100ULL);
stratum_timing_record(&tracker, 12, 5000200ULL);

stratum_timing_reset(&tracker);

TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(&tracker, 10, 5050000ULL));
TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(&tracker, 11, 5050000ULL));
TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(&tracker, 12, 5050000ULL));
}

TEST_CASE("stratum_timing invalid timestamps and null safety", "[stratum_timing]")
{
stratum_timing_tracker_t tracker;
stratum_timing_reset(&tracker);

// NULL tracker safety
stratum_timing_record(NULL, 1, 1000ULL);
stratum_timing_reset(NULL);
TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(NULL, 1, 2000ULL));

// Backward clock jump (now_us < submit_time_us)
stratum_timing_record(&tracker, 7, 5000000ULL);
TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(&tracker, 7, 4999999ULL));

// Non-existent ID
TEST_ASSERT_EQUAL_FLOAT(-1.0f, stratum_timing_calculate_ms(&tracker, 999, 6000000ULL));
}
12 changes: 10 additions & 2 deletions main/tasks/stratum_v1_client.c
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
#include <esp_heap_caps.h>
#include "esp_transport_ssl.h"
#include "freertos/task.h"
#include "stratum_timing.h"

#define MAX_EXTRANONCE_2_LEN 32
#define TRANSPORT_TIMEOUT_MS 5000
Expand All @@ -27,6 +28,7 @@ static const char *TAG = "stratum_v1";

static StratumApiV1Message *s_v1_msg = NULL;
static sv1_conn_t *s_v1_conn = NULL;
static stratum_timing_tracker_t s_v1_timing = {0};

static bool add_active_job_id(char active_job_ids[][MAX_JOB_ID_LEN], int *count, const char *job_id)
{
Expand Down Expand Up @@ -81,6 +83,7 @@ int stratum_v1_submit_share(GlobalState *GLOBAL_STATE, const asic_job_t *active_
}

int uid = s_v1_conn->send_uid++;
uint64_t now = 0;
int ret = STRATUM_V1_submit_share(
transport,
uid,
Expand All @@ -90,9 +93,13 @@ int stratum_v1_submit_share(GlobalState *GLOBAL_STATE, const asic_job_t *active_
active_job->ntime,
nonce,
version_bits,
sent_time_us);
&now);

if (ret >= 0) {
stratum_timing_record(&s_v1_timing, (uint32_t)uid, now);
if (sent_time_us) {
*sent_time_us = now;
}
if (GLOBAL_STATE->SYSTEM_MODULE.shares_pending < UINT16_MAX) {
GLOBAL_STATE->SYSTEM_MODULE.shares_pending++;
}
Expand Down Expand Up @@ -120,6 +127,7 @@ void stratum_v1_close_connection(GlobalState *GLOBAL_STATE)
pthread_mutex_unlock(&GLOBAL_STATE->transport_mutex);

SYSTEM_reset_pool_session(GLOBAL_STATE);
stratum_timing_reset(&s_v1_timing);
}

esp_err_t stratum_v1_run(GlobalState *GLOBAL_STATE, uint16_t pool_idx)
Expand Down Expand Up @@ -379,7 +387,7 @@ esp_err_t stratum_v1_run(GlobalState *GLOBAL_STATE, uint16_t pool_idx)
break;

case STRATUM_RESULT: {
float response_time_ms = STRATUM_V1_get_response_time_ms(s_v1_msg->message_id, receive_time_us);
float response_time_ms = stratum_timing_calculate_ms(&s_v1_timing, (uint32_t)s_v1_msg->message_id, (uint64_t)receive_time_us);
Comment thread
mutatrum marked this conversation as resolved.
Outdated
if (response_time_ms >= 0) {
if (GLOBAL_STATE->SYSTEM_MODULE.shares_pending > 0) {
GLOBAL_STATE->SYSTEM_MODULE.shares_pending--;
Expand Down
16 changes: 6 additions & 10 deletions main/tasks/stratum_v2_client.c
Original file line number Diff line number Diff line change
Expand Up @@ -16,20 +16,19 @@
#include "libbase58.h"
#include "device_config.h"
#include "esp_heap_caps.h"
#include "stratum_timing.h"

#include <string.h>
#include <stdlib.h>
#include <math.h>

#define TRANSPORT_TIMEOUT_MS 5000
#define SV2_MAX_FRAME_SIZE 8192
#define SV2_SUBMIT_TIMING_SLOTS 32

static const char *TAG = "stratum_v2";

static sv2_conn_t *s_v2_conn = NULL;

static int64_t stratum_v2_submit_time_us[SV2_SUBMIT_TIMING_SLOTS] = {0};
static stratum_timing_tracker_t s_v2_timing = {0};

static bool add_active_job_id(uint32_t *active_job_ids, int *count, uint32_t job_id)
{
Expand Down Expand Up @@ -109,8 +108,8 @@ void stratum_v2_close_connection(GlobalState *GLOBAL_STATE)
}
pthread_mutex_unlock(&GLOBAL_STATE->transport_mutex);

memset(stratum_v2_submit_time_us, 0, sizeof(stratum_v2_submit_time_us));
SYSTEM_reset_pool_session(GLOBAL_STATE);
stratum_timing_reset(&s_v2_timing);
}

static void stratum_v2_update_pending_shares(GlobalState *GLOBAL_STATE)
Expand All @@ -127,7 +126,7 @@ static void stratum_v2_update_pending_shares(GlobalState *GLOBAL_STATE)

static void stratum_v2_track_submit(GlobalState *GLOBAL_STATE, uint32_t sequence_number)
{
stratum_v2_submit_time_us[sequence_number % SV2_SUBMIT_TIMING_SLOTS] = esp_timer_get_time();
stratum_timing_record(&s_v2_timing, sequence_number, esp_timer_get_time());
stratum_v2_update_pending_shares(GLOBAL_STATE);
}

Expand Down Expand Up @@ -795,14 +794,11 @@ esp_err_t stratum_v2_run(GlobalState *GLOBAL_STATE, uint16_t pool_idx)
(unsigned long)accepted_count, (unsigned long)pending);
accepted_count = pending;
}
int slot = last_sequence_number % SV2_SUBMIT_TIMING_SLOTS;
int64_t submit_time_us = stratum_v2_submit_time_us[slot];
if (submit_time_us > 0) {
float response_time_ms = (float)(esp_timer_get_time() - submit_time_us) / 1000.0f;
float response_time_ms = stratum_timing_calculate_ms(&s_v2_timing, last_sequence_number, esp_timer_get_time());
if (response_time_ms >= 0) {
ESP_LOGI(TAG, "Shares accepted: %lu (%.1f ms)", accepted_count, response_time_ms);
GLOBAL_STATE->SYSTEM_MODULE.response_time = response_time_ms;
GLOBAL_STATE->SYSTEM_MODULE.response_share_batch = (uint16_t)accepted_count;
stratum_v2_submit_time_us[slot] = 0;
} else {
ESP_LOGI(TAG, "Shares accepted: %lu", accepted_count);
}
Expand Down
Loading