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
2 changes: 1 addition & 1 deletion tpu_sync/fault_injection/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ cc_library(
"@com_google_absl//absl/base:core_headers",
"@com_google_absl//absl/base:no_destructor",
"@com_google_absl//absl/container:flat_hash_map",
"@com_google_absl//absl/log",
"@com_google_absl//absl/log:absl_check",
"@com_google_absl//absl/random",
"@com_google_absl//absl/status",
Expand All @@ -50,7 +51,6 @@ cc_test(
features = ["-use_header_modules"],
deps = [
":hooks",
"@com_google_absl//absl/container:flat_hash_set",
"@com_google_absl//absl/status",
"@com_google_absl//absl/strings",
"@com_google_absl//absl/synchronization",
Expand Down
144 changes: 144 additions & 0 deletions tpu_sync/fault_injection/fault_injector.cc
Original file line number Diff line number Diff line change
Expand Up @@ -16,14 +16,20 @@

#include <poll.h>
#include <sys/socket.h>
#include <unistd.h>

#include <algorithm>
#include <cerrno>
#include <cmath>
#include <cstdint>
#include <cstdio>
#include <cstdlib>
#include <fstream>
#include <iterator>
#include <stdexcept>
#include <string>
#include <string_view>
#include <thread>
#include <utility>
#include <vector>

Expand All @@ -33,8 +39,11 @@
#include "absl/log/absl_check.h"
#include "absl/random/random.h"
#include "absl/status/status.h"
#include "absl/strings/ascii.h"
#include "absl/strings/match.h"
#include "absl/strings/numbers.h"
#include "absl/strings/str_cat.h"
#include "absl/strings/str_split.h"
#include "absl/strings/strip.h"
#include "absl/synchronization/mutex.h"
#include "absl/time/clock.h"
Expand All @@ -59,8 +68,85 @@ absl::Span<const std::string_view> HooksFor(FaultInjectionType action) {
return {};
}

constexpr std::string_view kFaultInjectionFileEnvVar =
"RAIDEN_FAULT_INJECTION_FILE";

__attribute__((constructor)) void InitFaultInjectorFromEnvAtLoad() {
if (const char* p = std::getenv(kFaultInjectionFileEnvVar.data());
p != nullptr && p[0] != '\0') {
(void)GetFaultInjector();
}
}

// Plain-text format: one rule per line (`#` comments and blank lines ignored):
// <hook|pattern*> <fail|delay> <probability> [min_delay_ms [max_delay_ms]]
absl::Status LoadRulesFromText(std::string_view content,
FaultInjector& injector) {
FaultInjectionRules rules;
for (std::string_view line :
absl::StrSplit(content, '\n', absl::SkipEmpty())) {
line = absl::StripAsciiWhitespace(line);
if (line.empty() || line[0] == '#') continue;
std::vector<std::string_view> cols =
absl::StrSplit(line, absl::ByAnyChar(" \t,"), absl::SkipEmpty());
if (cols.size() < 3 || cols.size() > 5) {
return absl::InvalidArgumentError(
absl::StrCat("invalid rule line: ", line));
}
FaultInjectionRule r;
r.hook = std::string(cols[0]);
if (cols[1] == "delay") {
r.action = FaultInjectionType::kDelay;
} else if (cols[1] == "fail") {
r.action = FaultInjectionType::kFail;
} else {
return absl::InvalidArgumentError(
absl::StrCat("unknown action: ", cols[1]));
}
if (!absl::SimpleAtod(cols[2], &r.probability)) {
return absl::InvalidArgumentError(
absl::StrCat("invalid probability: ", cols[2]));
}
if (cols.size() >= 4 && !absl::SimpleAtoi(cols[3], &r.min_delay_ms)) {
return absl::InvalidArgumentError(
absl::StrCat("invalid min_delay_ms: ", cols[3]));
}
if (cols.size() == 5 && !absl::SimpleAtoi(cols[4], &r.max_delay_ms)) {
return absl::InvalidArgumentError(
absl::StrCat("invalid max_delay_ms: ", cols[4]));
}
rules.push_back(std::move(r));
}
return injector.Install(rules);
}

void WriteStatusFile(std::string_view status_path, bool armed,
uint64_t total_hits,
const absl::flat_hash_map<std::string, uint64_t>& hits) {
std::string out_str =
absl::StrCat("armed=", armed ? 1 : 0, "\ntotal_hits=", total_hits, "\n");
for (const auto& [hook, count] : hits) {
absl::StrAppend(&out_str, hook, "=", count, "\n");
}
std::string tmp_path = absl::StrCat(status_path, ".tmp");
if (std::ofstream out(tmp_path); out) {
out << out_str;
out.close();
(void)std::rename(tmp_path.c_str(), std::string(status_path).c_str());
}
}

} // namespace

FaultInjector::FaultInjector() {
if (const char* p = std::getenv(kFaultInjectionFileEnvVar.data());
p != nullptr && p[0] != '\0') {
StartFileWatcher(p);
}
}

FaultInjector::~FaultInjector() { StopFileWatcher(); }

bool FaultInjector::IsHookActive(std::string_view hook) const noexcept {
if (!HasActiveInjections()) return false;
absl::ReaderMutexLock lock(mu_);
Expand Down Expand Up @@ -181,6 +267,64 @@ absl::flat_hash_map<std::string, uint64_t> FaultInjector::GetHitCounts() const {
return rule_hits_;
}

void FaultInjector::StartFileWatcher(std::string_view file_path,
absl::Duration poll_interval) {
StopFileWatcher();
if (file_path.empty()) return;
{
absl::MutexLock lock(watcher_mu_);
watcher_stopping_ = false;
}
watcher_thread_ = std::thread(&FaultInjector::WatcherLoop, this,
std::string(file_path), poll_interval);
}

void FaultInjector::StopFileWatcher() {
if (!watcher_thread_.joinable()) return;
{
absl::MutexLock lock(watcher_mu_);
watcher_stopping_ = true;
}
watcher_thread_.join();
}

void FaultInjector::WatcherLoop(std::string file_path,
absl::Duration poll_interval) {
const std::string status_path = absl::StrCat(file_path, ".status.", getpid());
std::string last_content;
bool status_written = false;
bool last_armed = false;
uint64_t last_hits = 0;

while (true) {
if (std::ifstream in(file_path); in) {
std::string content((std::istreambuf_iterator<char>(in)), {});
if (content != last_content) {
last_content = std::move(content);
LoadRulesFromText(last_content, *this).IgnoreError();
}
} else if (!last_content.empty()) {
Install({}).IgnoreError();
last_content.clear();
}

bool armed = HasActiveInjections();
uint64_t total_hits = GetHitCount();
if (!status_written || armed != last_armed || total_hits != last_hits) {
WriteStatusFile(status_path, armed, total_hits, GetHitCounts());
status_written = true;
last_armed = armed;
last_hits = total_hits;
}

absl::MutexLock lock(watcher_mu_);
if (watcher_mu_.AwaitWithTimeout(absl::Condition(&watcher_stopping_),
poll_interval)) {
break;
}
}
}

void FaultInjector::ExecuteDelay(std::string_view hook) {
EvaluateAndSleep(hook);
}
Expand Down
20 changes: 18 additions & 2 deletions tpu_sync/fault_injection/fault_injector.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
#include <cstdint>
#include <string>
#include <string_view>
#include <thread>
#include <vector>

#include "absl/base/optimization.h"
Expand All @@ -28,6 +29,7 @@
#include "absl/random/random.h"
#include "absl/status/status.h"
#include "absl/synchronization/mutex.h"
#include "absl/time/time.h"
#include "tpu_sync/fault_injection/hooks.h" // IWYU pragma: export

namespace tpu_raiden {
Expand Down Expand Up @@ -55,10 +57,10 @@ using FaultInjectionRules = std::vector<FaultInjectionRule>;

class FaultInjector final {
public:
FaultInjector() = default;
FaultInjector();
FaultInjector(const FaultInjector&) = delete;
FaultInjector& operator=(const FaultInjector&) = delete;
~FaultInjector() = default;
~FaultInjector();

// Returns true if any fault injection rules are currently
// active.
Expand Down Expand Up @@ -92,6 +94,13 @@ class FaultInjector final {
absl::flat_hash_map<std::string, uint64_t> GetHitCounts() const
ABSL_LOCKS_EXCLUDED(mu_);

// Polls `file_path` every `poll_interval`, arming rules when the file exists,
// disarming when deleted, and writing status to `<file_path>.status.<pid>`.
void StartFileWatcher(std::string_view file_path,
absl::Duration poll_interval = absl::Milliseconds(200))
ABSL_LOCKS_EXCLUDED(watcher_mu_, mu_);
void StopFileWatcher() ABSL_LOCKS_EXCLUDED(watcher_mu_);

// Slow-path execution helpers invoked only when HasActiveInjections() is
// true. A delay always runs to completion.
void ExecuteDelay(std::string_view hook) ABSL_LOCKS_EXCLUDED(mu_);
Expand All @@ -112,8 +121,15 @@ class FaultInjector final {
void ExecuteSocket(std::string_view hook, int fd) ABSL_LOCKS_EXCLUDED(mu_);

private:
void WatcherLoop(std::string file_path, absl::Duration poll_interval)
ABSL_LOCKS_EXCLUDED(watcher_mu_, mu_);

static inline std::atomic<bool> has_active_injections_{false};

std::thread watcher_thread_;
absl::Mutex watcher_mu_;
bool watcher_stopping_ ABSL_GUARDED_BY(watcher_mu_) = false;

// Evaluates the hook and sleeps for the delay if a delay action is chosen.
FaultInjectionAction EvaluateAndSleep(std::string_view hook)
ABSL_LOCKS_EXCLUDED(mu_);
Expand Down
97 changes: 95 additions & 2 deletions tpu_sync/fault_injection/fault_injector_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,10 @@

#include <cerrno>
#include <cstdint>
#include <cstdio>
#include <fstream>
#include <functional>
#include <iterator>
#include <stdexcept>
#include <string>
#include <string_view>
Expand All @@ -27,6 +31,8 @@
#include <gmock/gmock.h>
#include <gtest/gtest.h>
#include "absl/status/status.h"
#include "absl/strings/match.h"
#include "absl/strings/str_cat.h"
#include "absl/synchronization/notification.h"
#include "absl/time/clock.h"
#include "absl/time/time.h"
Expand Down Expand Up @@ -63,8 +69,14 @@ void WaitForHit(std::string_view hook) {

class FaultInjectorTest : public ::testing::Test {
protected:
void SetUp() override { GetFaultInjector().Reset(); }
void TearDown() override { GetFaultInjector().Reset(); }
void SetUp() override {
GetFaultInjector().StopFileWatcher();
GetFaultInjector().Reset();
}
void TearDown() override {
GetFaultInjector().StopFileWatcher();
GetFaultInjector().Reset();
}
};

TEST_F(FaultInjectorTest, DefaultStateIsInactive) {
Expand Down Expand Up @@ -562,5 +574,86 @@ TEST_F(FaultInjectorTest, FastPathOverheadIsSubNanosecond) {
EXPECT_LT(ns_per_check, 5.0) << "ns_per_check=" << ns_per_check;
}

bool WaitUntil(const std::function<bool()>& predicate) {
absl::Time deadline = absl::Now() + absl::Seconds(5);
while (absl::Now() < deadline) {
if (predicate()) return true;
absl::SleepFor(absl::Milliseconds(10));
}
return predicate();
}

std::string ReadTextFile(const std::string& path) {
std::ifstream in(path);
if (!in) return "";
return std::string((std::istreambuf_iterator<char>(in)), {});
}

TEST_F(FaultInjectorTest, FileWatcherArmsDisarmsAndPreservesStatusMetrics) {
const std::string rules_path =
absl::StrCat(::testing::TempDir(), "/raiden_test_faults.txt");
const std::string status_path =
absl::StrCat(rules_path, ".status.", getpid());
(void)std::remove(rules_path.c_str());
(void)std::remove(status_path.c_str());

GetFaultInjector().StartFileWatcher(rules_path, absl::Milliseconds(20));

// Create rules file to arm kTestHookAlpha.
{
std::ofstream out(rules_path);
out << "transfer_recv_session.pull.request fail 1.0\n";
}
ASSERT_TRUE(WaitUntil([]() { return FaultInjector::HasActiveInjections(); }));
EXPECT_THAT(FaultInjectStatus(kTestHookAlpha),
StatusIs(absl::StatusCode::kInternal));
EXPECT_THAT(FaultInjectStatus(kTestHookAlpha),
StatusIs(absl::StatusCode::kInternal));

ASSERT_TRUE(WaitUntil([&]() {
std::string s = ReadTextFile(status_path);
return absl::StrContains(s, "armed=1\n") &&
absl::StrContains(s, "total_hits=2\n");
}));
EXPECT_THAT(ReadTextFile(status_path),
HasSubstr("transfer_recv_session.pull.request=2\n"));

// Delete rules file to disarm and verify cumulative metrics survive.
ASSERT_EQ(std::remove(rules_path.c_str()), 0);
ASSERT_TRUE(
WaitUntil([]() { return !FaultInjector::HasActiveInjections(); }));
ABSL_EXPECT_OK(FaultInjectStatus(kTestHookAlpha));

ASSERT_TRUE(WaitUntil([&]() {
return absl::StrContains(ReadTextFile(status_path), "armed=0\n");
}));
std::string status_disarmed = ReadTextFile(status_path);
EXPECT_THAT(status_disarmed, HasSubstr("total_hits=2\n"));
EXPECT_THAT(status_disarmed,
HasSubstr("transfer_recv_session.pull.request=2\n"));

// Recreate rules file and verify metrics accumulate across disarm -> re-arm.
{
std::ofstream out(rules_path);
out << absl::StrCat(kTestHookBeta, " delay 1.0 1 1\n");
}
ASSERT_TRUE(WaitUntil(
[]() { return GetFaultInjector().IsHookActive(kTestHookBeta); }));
FaultInjectDelay(kTestHookBeta);

ASSERT_TRUE(WaitUntil([&]() {
std::string s = ReadTextFile(status_path);
return absl::StrContains(s, "armed=1\n") &&
absl::StrContains(s, "total_hits=3\n");
}));
std::string status_rearmed = ReadTextFile(status_path);
EXPECT_THAT(status_rearmed, HasSubstr(absl::StrCat(kTestHookAlpha, "=2\n")));
EXPECT_THAT(status_rearmed, HasSubstr(absl::StrCat(kTestHookBeta, "=1\n")));

GetFaultInjector().StopFileWatcher();
(void)std::remove(rules_path.c_str());
(void)std::remove(status_path.c_str());
}

} // namespace
} // namespace tpu_raiden
Loading