Skip to content
Open
Show file tree
Hide file tree
Changes from 11 commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
b06243c
[C Driver] Add resolved variables to driver context to separate real …
mikeb01 Aug 11, 2026
05260d0
[C Driver] Validate that affinity configuration options do not overlap.
mikeb01 Aug 11, 2026
8182965
[C Driver] Rearrange cpuset and affinity validation to incorporate ne…
mikeb01 Aug 11, 2026
5a65c1a
[Driver/C] Initial pairwise implementation of L3 cache affinity valid…
EmilJohn24 Aug 12, 2026
47ef46a
[Driver/C] Transfer peer table building logic to aeron_topology.c. Im…
EmilJohn24 Aug 13, 2026
13fb595
[Driver/C] Refactor to use struct holding almost all information need…
EmilJohn24 Aug 14, 2026
7cce802
[Driver/C] Minor test reorganization in TopologyTest.
EmilJohn24 Aug 14, 2026
99ad211
[Driver/C] Add a TODO regarding potential peer table redesign.
EmilJohn24 Aug 14, 2026
506bcce
[Driver/C] Implement die locality checking using group table approach…
EmilJohn24 Aug 16, 2026
35a5aed
[Driver/C] Add peer table generation inside group table call for L3 l…
EmilJohn24 Aug 16, 2026
29d76b7
[Driver/C] Add extra checks for group table functions following curso…
EmilJohn24 Aug 17, 2026
f18460c
[Driver/C] Change error handling in validation to skip specific check…
EmilJohn24 Aug 17, 2026
52bf6e6
[Driver/C] Set resolved CPU affinities inside validation logic as well.
EmilJohn24 Aug 17, 2026
68c0f91
[Driver/C] Add teardown for output stream in DriverTest
EmilJohn24 Aug 17, 2026
3f8b57c
[Driver/C] Additional null checks for group table generation.
EmilJohn24 Aug 17, 2026
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
150 changes: 144 additions & 6 deletions aeron-driver/src/main/c/aeron_driver.c
Original file line number Diff line number Diff line change
Expand Up @@ -578,10 +578,14 @@ void aeron_driver_context_print_configuration(aeron_driver_context_t *context)
(uint64_t)context->network_publication_max_messages_per_send);
AERON_FPRINTF(fpout, "\n resource_free_limit=%" PRIu32, context->resource_free_limit);
AERON_FPRINTF(fpout, "\n conductor_cpu_affinity_no=%" PRId32, context->conductor_cpu_affinity_no);
AERON_FPRINTF(fpout, "\n conductor_cpu_affinity_resolved=%" PRId32, context->conductor_cpu_affinity_resolved);
AERON_FPRINTF(fpout, "\n receiver_cpu_affinity_no=%" PRId32, context->receiver_cpu_affinity_no);
AERON_FPRINTF(fpout, "\n receiver_cpu_affinity_resolved=%" PRId32, context->receiver_cpu_affinity_resolved);
AERON_FPRINTF(fpout, "\n sender_cpu_affinity_no=%" PRId32, context->sender_cpu_affinity_no);
AERON_FPRINTF(fpout, "\n sender_cpu_affinity_resolved=%" PRId32, context->sender_cpu_affinity_resolved);
AERON_FPRINTF(fpout, "\n native_resource_agent_cpu_affinity_no=%" PRId32, context->native_resource_agent_cpu_affinity_no);
AERON_FPRINTF(fpout, "\n cpuset_affinity=%" PRId32, context->cpuset_affinity);
AERON_FPRINTF(fpout, "\n native_resource_agent_cpu_affinity_resolved=%" PRId32, context->native_resource_agent_cpu_affinity_resolved);
AERON_FPRINTF(fpout, "\n cpuset_affinity=%s", context->cpuset_affinity ? "true" : "false");
AERON_FPRINTF(fpout, "\n cpuset_warnings_as_errors=%" PRId32, context->cpuset_warnings_as_errors);

AERON_FPRINTF(fpout, "\n epoch_clock=%s",
Expand Down Expand Up @@ -1183,23 +1187,157 @@ int aeron_driver_apply_cpuset_affinity(aeron_driver_context_t *context)
const int total_warnings_count =
alignment_warnings_count + cluster_locality_warnings_count + l3_locality_warnings_count;

if (context->cpuset_warnings_as_errors && 0 < total_warnings_count)
if (aeron_driver_context_apply_cpuset_affinity(context, cpus, cpu_count) < 0)
{
AERON_SET_ERR(EINVAL, "cpuset warnings as errors, %d warnings", total_warnings_count);
AERON_APPEND_ERR("%s", "failed to apply cpuset affinity");
goto error;
}

if (aeron_driver_context_apply_cpuset_affinity(context, cpus, cpu_count) < 0)
aeron_free(cpus);
return total_warnings_count;

error:
aeron_free(cpus);
return -1;
}

static int aeron_driver_validate_affinity_pair(
const int32_t a,
const int32_t b,
const char *a_name,
const char *b_name,
FILE *output)
{
if (AERON_NULL_VALUE != a && AERON_NULL_VALUE != b && a == b)
{
fprintf(output, "WARN: %s and %s are sharing cpu affinity=%" PRId32 "\n", a_name, b_name, a);
return 1;
}

return 0;
}

int aeron_driver_validate_unshared_affinity(aeron_driver_context_t* context, FILE *output)
{
int warnings = 0;

warnings += aeron_driver_validate_affinity_pair(
context->conductor_cpu_affinity_no, context->sender_cpu_affinity_no, "conductor", "sender", output);
warnings += aeron_driver_validate_affinity_pair(
context->conductor_cpu_affinity_no, context->receiver_cpu_affinity_no, "conductor", "receiver", output);
warnings += aeron_driver_validate_affinity_pair(
context->conductor_cpu_affinity_no, context->native_resource_agent_cpu_affinity_no, "conductor", "native_resource_agent", output);
warnings += aeron_driver_validate_affinity_pair(
context->sender_cpu_affinity_no, context->receiver_cpu_affinity_no, "sender", "receiver", output);
warnings += aeron_driver_validate_affinity_pair(
context->sender_cpu_affinity_no, context->native_resource_agent_cpu_affinity_no, "sender", "native_resource_agent", output);
warnings += aeron_driver_validate_affinity_pair(
context->receiver_cpu_affinity_no, context->native_resource_agent_cpu_affinity_no, "receiver", "native_resource_agent", output);

return warnings;
}

int aeron_driver_validate_group_locality(
const aeron_topology_cpu_info_t *cpu_info,
const char *sysfs_prop_descriptor,
FILE *output)
{
if (cpu_info->group_count <= 1)
{
return 0;
}

fprintf(output, "WARN: cpu affinities span %d %s:\n", cpu_info->group_count, sysfs_prop_descriptor);
// TODO: Consider alternative, which is just to print all group IDs individually in a list
// This will allow us to remove the group_ids array processing
// Another possibility is to pre-sort the cpu list by group ID instead of having a list for them
for (int g = 0; g < cpu_info->group_count; g++)
{
// If group IDs were not listed, just use the index (e.g., in L3 cache case)
const int current_group = NULL != cpu_info->group_ids ? cpu_info->group_ids[g] : g;
fprintf(output, " group %d:", current_group);
for (int i = 0; i < cpu_info->cpu_count; i++)
{
const aeron_topology_cpu_group_t *entry = &cpu_info->cpus[i];
if (current_group == entry->group_id)
{
fprintf(
output, " %s (cpu=%d [configured=%d])",
entry->extra_info->name, entry->cpu, entry->extra_info->original_cpu);
}
}
fprintf(output, "\n");
}

return 1;
}

int aeron_driver_validate_and_apply_affinity_configuration(aeron_driver_context_t *context)
{
int unshared_affinity_warnings;
if ((unshared_affinity_warnings = aeron_driver_validate_unshared_affinity(context, stderr)) < 0)
{
AERON_APPEND_ERR("%s", "failed to validate unshared affinity");
goto error;
}

int cpuset_warnings = 0;
if ((cpuset_warnings = aeron_driver_apply_cpuset_affinity(context)) < 0)
{
AERON_APPEND_ERR("%s", "failed to apply cpuset affinity");
goto error;
}

aeron_free(cpus);
int l3_locality_warnings = 0;
int die_locality_warnings = 0;
#ifdef __linux__
aeron_topology_extra_info_t extra_info[4] = {
{ "conductor", context->conductor_cpu_affinity_no },
{ "sender", context->sender_cpu_affinity_no },
{ "receiver", context->receiver_cpu_affinity_no },
{ "native_resource_agent", context->native_resource_agent_cpu_affinity_no }
};
aeron_topology_cpu_group_t cpu_groups[4] = {
{ context->conductor_cpu_affinity_resolved, AERON_NULL_VALUE, &extra_info[0] },
{ context->sender_cpu_affinity_resolved, AERON_NULL_VALUE, &extra_info[1] },
{ context->receiver_cpu_affinity_resolved, AERON_NULL_VALUE, &extra_info[2] },
{ context->native_resource_agent_cpu_affinity_resolved, AERON_NULL_VALUE, &extra_info[3] }
};
int *peers[4] = { NULL, NULL, NULL, NULL };
int peer_count[4] = { 0, 0, 0, 0 };
aeron_topology_cpu_info_t cpu_info = { cpu_groups, 4, peers, peer_count, NULL, 0 };

if (aeron_topology_build_l3_group_table(AERON_TOPOLOGY_SYS_CPU_PATH, &cpu_info) < 0)
{
AERON_APPEND_ERR("%s", "failed to build l3 group table");
aeron_topology_cpu_info_free(&cpu_info);
goto error;
}
l3_locality_warnings = aeron_driver_validate_group_locality(&cpu_info, "L3 cache domain", stderr);

if (aeron_topology_build_die_locality_group_table(AERON_TOPOLOGY_SYS_CPU_PATH, &cpu_info) < 0)
{
AERON_APPEND_ERR("%s", "failed to build die locality group table");
aeron_topology_cpu_info_free(&cpu_info);
goto error;
}
die_locality_warnings = aeron_driver_validate_group_locality(&cpu_info, "dies", stderr);

aeron_topology_cpu_info_free(&cpu_info);
#endif

const int total_warnings_count =
unshared_affinity_warnings + cpuset_warnings + l3_locality_warnings + die_locality_warnings;

if (context->cpuset_warnings_as_errors && 0 < total_warnings_count)
{
AERON_SET_ERR(EINVAL, "cpuset warnings as errors, %d warnings", total_warnings_count);
goto error;
}

return 0;

error:
aeron_free(cpus);
return -1;
}

Expand Down
2 changes: 2 additions & 0 deletions aeron-driver/src/main/c/aeron_driver.h
Original file line number Diff line number Diff line change
Expand Up @@ -47,4 +47,6 @@ bool aeron_is_driver_active_with_cnc(

int aeron_driver_apply_cpuset_affinity(aeron_driver_context_t *context);

int aeron_driver_validate_and_apply_affinity_configuration(aeron_driver_context_t *context);

#endif //AERON_DRIVER_H
40 changes: 25 additions & 15 deletions aeron-driver/src/main/c/aeron_driver_context.c
Original file line number Diff line number Diff line change
Expand Up @@ -778,24 +778,31 @@ int aeron_driver_context_init(aeron_driver_context_t **context)
_context->conductor_cpu_affinity_no,
-1,
255);
_context->conductor_cpu_affinity_resolved = _context->conductor_cpu_affinity_no;

_context->receiver_cpu_affinity_no = aeron_config_parse_int32(
AERON_RECEIVER_CPU_AFFINITY_ENV_VAR,
getenv(AERON_RECEIVER_CPU_AFFINITY_ENV_VAR),
_context->receiver_cpu_affinity_no,
-1,
255);
_context->receiver_cpu_affinity_resolved = _context->receiver_cpu_affinity_no;

_context->sender_cpu_affinity_no = aeron_config_parse_int32(
AERON_SENDER_CPU_AFFINITY_ENV_VAR,
getenv(AERON_SENDER_CPU_AFFINITY_ENV_VAR),
_context->sender_cpu_affinity_no,
-1,
255);
_context->sender_cpu_affinity_resolved = _context->sender_cpu_affinity_no;

_context->native_resource_agent_cpu_affinity_no = aeron_config_parse_int32(
AERON_DRIVER_NATIVE_RESOURCE_AGENT_CPU_AFFINITY_ENV_VAR,
getenv(AERON_DRIVER_NATIVE_RESOURCE_AGENT_CPU_AFFINITY_ENV_VAR),
_context->native_resource_agent_cpu_affinity_no,
-1,
255);
_context->native_resource_agent_cpu_affinity_resolved = _context->native_resource_agent_cpu_affinity_no;

_context->cpuset_affinity = aeron_parse_bool(
getenv(AERON_DRIVER_CPUSET_AFFINITY_ENV_VAR), AERON_DRIVER_CPUSET_AFFINITY_DEFAULT);
Expand Down Expand Up @@ -3407,32 +3414,32 @@ void aeron_set_thread_affinity_on_start(void *state, const char *role_name)
{
aeron_driver_context_t *context = (aeron_driver_context_t *)state;
int result = 0;
if (0 <= context->conductor_cpu_affinity_no &&
if (0 <= context->conductor_cpu_affinity_resolved &&
(0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_CONDUCTOR, role_name) ||
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_CONDUCTOR_NEW, role_name) ||
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_SHARED, role_name) ||
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_SHARED_NEW, role_name)))
{
result = aeron_thread_set_affinity(role_name, (uint8_t)context->conductor_cpu_affinity_no);
result = aeron_thread_set_affinity(role_name, (uint8_t)context->conductor_cpu_affinity_resolved);
}
else if (0 <= context->sender_cpu_affinity_no &&
else if (0 <= context->sender_cpu_affinity_resolved &&
(0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_SENDER, role_name) ||
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_SENDER_NEW, role_name) ||
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_SHARED_NETWORK, role_name) ||
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_SHARED_NETWORK_NEW, role_name)))
{
result = aeron_thread_set_affinity(role_name, (uint8_t)context->sender_cpu_affinity_no);
result = aeron_thread_set_affinity(role_name, (uint8_t)context->sender_cpu_affinity_resolved);
}
else if (0 <= context->receiver_cpu_affinity_no &&
else if (0 <= context->receiver_cpu_affinity_resolved &&
(0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_RECEIVER, role_name) ||
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_RECEIVER_NEW, role_name)))
{
result = aeron_thread_set_affinity(role_name, (uint8_t)context->receiver_cpu_affinity_no);
result = aeron_thread_set_affinity(role_name, (uint8_t)context->receiver_cpu_affinity_resolved);
}
else if (0 <= context->native_resource_agent_cpu_affinity_no &&
else if (0 <= context->native_resource_agent_cpu_affinity_resolved &&
0 == strcmp(AERON_DRIVER_AGENT_ROLE_NAME_NATIVE_RESOURCE_AGENT, role_name))
{
result = aeron_thread_set_affinity(role_name, (uint8_t)context->native_resource_agent_cpu_affinity_no);
result = aeron_thread_set_affinity(role_name, (uint8_t)context->native_resource_agent_cpu_affinity_resolved);
Comment thread
cursor[bot] marked this conversation as resolved.
}

if (result < 0)
Expand Down Expand Up @@ -3538,9 +3545,8 @@ bool aeron_driver_context_get_cpuset_warnings_as_errors(aeron_driver_context_t *
}

static int aeron_driver_context_apply_cpuset_affinity_per_cpu(
const int *cpus, int cpu_count, const char *name, int32_t *affinity_ptr)
const int *cpus, int cpu_count, const char *name, int32_t affinity, int32_t *affinity_out)
{
const int32_t affinity = *affinity_ptr;
if (cpu_count <= affinity)
{
AERON_SET_ERR(
Expand All @@ -3551,7 +3557,7 @@ static int aeron_driver_context_apply_cpuset_affinity_per_cpu(
if (-1 < affinity)
{
const int32_t cpuset_conductor_affinity = cpus[(int)affinity];
*affinity_ptr = cpuset_conductor_affinity;
*affinity_out = cpuset_conductor_affinity;
}

return 0;
Expand All @@ -3560,28 +3566,32 @@ static int aeron_driver_context_apply_cpuset_affinity_per_cpu(
int aeron_driver_context_apply_cpuset_affinity(aeron_driver_context_t *context, const int *cpus, int cpu_count)
{
if (aeron_driver_context_apply_cpuset_affinity_per_cpu(
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_CONDUCTOR, &context->conductor_cpu_affinity_no))
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_CONDUCTOR,
context->conductor_cpu_affinity_no, &context->conductor_cpu_affinity_resolved))
{
AERON_APPEND_ERR("%s", "");
return -1;
}

if (aeron_driver_context_apply_cpuset_affinity_per_cpu(
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_RECEIVER, &context->receiver_cpu_affinity_no))
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_RECEIVER,
context->receiver_cpu_affinity_no, &context->receiver_cpu_affinity_resolved))
{
AERON_APPEND_ERR("%s", "");
return -1;
}

if (aeron_driver_context_apply_cpuset_affinity_per_cpu(
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_SENDER, &context->sender_cpu_affinity_no))
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_SENDER,
context->sender_cpu_affinity_no, &context->sender_cpu_affinity_resolved))
{
AERON_APPEND_ERR("%s", "");
return -1;
}

if (aeron_driver_context_apply_cpuset_affinity_per_cpu(
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_NATIVE_RESOURCE_AGENT, &context->native_resource_agent_cpu_affinity_no))
cpus, cpu_count, AERON_DRIVER_AGENT_ROLE_NAME_NATIVE_RESOURCE_AGENT,
context->native_resource_agent_cpu_affinity_no, &context->native_resource_agent_cpu_affinity_resolved))
{
AERON_APPEND_ERR("%s", "");
return -1;
Expand Down
4 changes: 4 additions & 0 deletions aeron-driver/src/main/c/aeron_driver_context.h
Original file line number Diff line number Diff line change
Expand Up @@ -232,9 +232,13 @@ typedef struct aeron_driver_context_stct
uint32_t max_resend; /* aeron.max.resend = 16 */

int32_t conductor_cpu_affinity_no; /* aeron.conductor.cpu.affinity = -1 */
int32_t conductor_cpu_affinity_resolved;
int32_t receiver_cpu_affinity_no; /* aeron.receiver.cpu.affinity = -1 */
int32_t receiver_cpu_affinity_resolved;
int32_t sender_cpu_affinity_no; /* aeron.sender.cpu.affinity = -1 */
int32_t sender_cpu_affinity_resolved;
int32_t native_resource_agent_cpu_affinity_no; /* aeron.driver.native.resource.agent.cpu.affinity = -1 */
int32_t native_resource_agent_cpu_affinity_resolved;
int32_t stream_session_limit; /* aeron.driver.stream.session.limit = INT32_MAX */
bool cpuset_affinity; /* aeron.driver.cpuset.affinity = false */
bool cpuset_warnings_as_errors; /* aeron.driver.cpuset.warnings_as_errors = false */
Expand Down
Loading
Loading