Skip to content
Merged
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
23 changes: 15 additions & 8 deletions tests/ut/cpp/hierarchical/test_scheduler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@

#include "call_config.h"
#include "common/host_span.h"
#include "common/host_span_names.h"
#include "common/host_span_scope.h"
#include "orchestrator.h"
#include "ring.h"
Expand Down Expand Up @@ -1210,13 +1211,13 @@ TEST(WorkerManagerTest, DispatchAndCompletionEmitHostSpans) {

worker.dispatch(WorkerDispatch{slot, 0});
ASSERT_TRUE(endpoint_ptr->wait_submitted(1));
EXPECT_TRUE(captured_host_span("node.dispatch"));
EXPECT_TRUE(captured_host_span(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::Dispatch)));

const WorkerDispatch submitted = endpoint_ptr->submitted().front();
endpoint_ptr->emit(WorkerProgressKind::COMPLETED, submitted);
worker.progress();
ASSERT_EQ(completed.size(), 1u);
EXPECT_TRUE(captured_host_span("node.complete"));
EXPECT_TRUE(captured_host_span(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::Complete)));

worker.stop();
allocator.shutdown();
Expand Down Expand Up @@ -1474,8 +1475,8 @@ TEST(WorkerManagerTest, TwoFrameLeaseSlotsDoNotDefineFifoOrAcceptance) {
ASSERT_TRUE(endpoint.poll_progress(progress));
EXPECT_EQ(progress.kind, WorkerProgressKind::ACCEPTED);
EXPECT_EQ(progress.dispatch.dispatch_id, 42u);
EXPECT_TRUE(captured_host_span("node.frame_submit"));
EXPECT_TRUE(captured_host_span("node.activate"));
EXPECT_TRUE(captured_host_span(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::FrameSubmit)));
EXPECT_TRUE(captured_host_span(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::Activate)));
allocator.shutdown();
}

Expand Down Expand Up @@ -2042,10 +2043,13 @@ TEST_F(ProgressSchedulerFixture, GroupSubmitReportsNoSingleWorkerAndNoSingleInde
);
orchestrator.close_run_submission(run);

ASSERT_TRUE(captured_host_span("node.submit"));
ASSERT_TRUE(captured_host_span(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::Submit)));
std::ostringstream expected;
expected << "run_id=" << run << " task_slot=" << group.task_slot << " group_index=-1 group_size=2 role=facade";
EXPECT_EQ(captured_host_span_attrs("node.submit"), expected.str());
EXPECT_EQ(
captured_host_span_attrs(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::Submit)),
expected.str()
);
}

TEST_F(ProgressSchedulerFixture, SuccessorStagesButActivatesOnlyAfterFifoPromotion) {
Expand Down Expand Up @@ -2515,11 +2519,14 @@ TEST_F(SchedulerFixture, NextLevelSubmitEmitsAPairableHostSpan) {

auto submitted = orch.submit_next_level(C(0x42), single_tensor_args(0xCAFE, TensorArgType::OUTPUT), cfg, 0);

ASSERT_TRUE(captured_host_span("node.submit"));
ASSERT_TRUE(captured_host_span(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::Submit)));
std::ostringstream expected;
expected << "run_id=" << run_id << " task_slot=" << submitted.task_slot
<< " group_index=0 group_size=1 worker_id=0 role=facade";
EXPECT_EQ(captured_host_span_attrs("node.submit"), expected.str());
EXPECT_EQ(
captured_host_span_attrs(simpler::host_trace::host_span_name(simpler::host_trace::HostSpan::Submit)),
expected.str()
);

mock_worker.wait_running();
mock_worker.complete();
Expand Down
27 changes: 26 additions & 1 deletion tests/ut/py/test_worker/test_host_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import inspect
import multiprocessing.shared_memory as shared_memory_mod
import struct
import subprocess
import sys
import threading
import time
Expand Down Expand Up @@ -65,6 +66,7 @@
_mailbox_store_i32,
_pack_py_callable_payload,
)
from simpler.worker_level import WorkerLevel

from ._harness import chip_callable, fake_chip_l3, requires_sim_binaries

Expand Down Expand Up @@ -2755,7 +2757,8 @@ def bad_graph(*_args):
with pytest.raises(ValueError, match="bad graph"):
worker._submit_l3_locked(bad_graph, None, cast(Any, object()))

assert emitted == [("node.graph_build", 1, 0, 0, 100, 175, "run_id=1 role=facade")]
expected_name = f"{worker._host_span_prefix}.graph_build"
assert emitted == [(expected_name, 1, 0, 0, 100, 175, "run_id=1 role=facade")]

def test_unsettled_graph_cancellation_abandons_the_handle_before_close(self):
worker, events = self._submission_failure_worker(failures=2)
Expand Down Expand Up @@ -7655,3 +7658,25 @@ def test_unregister_removes_only_after_last_digest_ref(self):
shm.unlink()
payload_shm.close()
payload_shm.unlink()


def test_the_cpp_pre_bind_level_word_is_the_ladder_word_for_l3():
"""`host_span_names.h` hand-writes the L3 word a third time, as its pre-bind default.

`WorkerLevel` is the source of truth and `strace_timing.py`'s copy is pinned to
it by its own test, but the C++ default is pinned to nothing — so a level word
could be renamed in Python while C++ keeps emitting the old one until a Worker
binds the prefix. Every span emitted before that binding would carry the stale
word, and no test would notice.

Read in a child process on purpose: the prefix freezes on first bind, so any
test in this process that constructed a Worker would leave the *bound* word
here instead of the default. Passing an empty word binds nothing (the setter
returns early) while the binding still reports what is currently in effect.
"""
source = "from _task_interface import _set_host_span_level_prefix as bind; print(bind(''), end='')"
completed = subprocess.run([sys.executable, "-c", source], capture_output=True, text=True, check=True, timeout=120)

assert completed.stdout == WorkerLevel.node.name, (
f"C++ pre-bind level word is {completed.stdout!r}, ladder says {WorkerLevel.node.name!r}"
)
Loading