Let runner implement the engine ABI for the engine worker.

This consolidates the duplicated controller interaction in runner and engine worker.
There is some unique interaction left in the runner:

  - Dumping serailized FuzzTest config.

    Needed so that one can pass the fuzzing options to the test binary command when directly using the controller. But this uses C++-based struct. We should either drop this or make it dump controller options as strings, as part of the engine ABI.

  - Notifying the absense of custom mutator to use the controller's builtin mutator.

    Technically the runner has its own mutator for LLVMFuzzerMutate, but it's using the legacy ByteArrayMutator for minimal dependency, while the controller has the FuzzTest-based mutator. This should not be a problem once we remove the legacy support in runner (and let legacy fuzzers use the FuzzTest LLVM fuzzer wrapper instead)

Also, move fork server to the engine worker because we need to make sure that fork server starts before handling persistent mode.

PiperOrigin-RevId: 995125597
diff --git a/centipede/BUILD b/centipede/BUILD
index e546515..5e1690e 100644
--- a/centipede/BUILD
+++ b/centipede/BUILD
@@ -985,52 +985,6 @@
     # This has no dependency and only uses string, pair and vector.
 )
 
-cc_library(
-    name = "engine_worker",
-    srcs = [
-        "engine_worker.cc",
-    ],
-    hdrs = ["engine_worker_abi.h"],
-    deps = [
-        ":engine_abi",
-        ":execution_metadata",
-        ":feature",
-        ":runner_fork_server",
-        ":runner_request",
-        ":runner_result",
-        ":runner_utils",
-        ":shared_memory_blob_sequence",
-        "@abseil-cpp//absl/base:nullability",
-        "@com_google_fuzztest//common:defs",
-    ],
-)
-
-bool_flag(
-    name = "bundle_controller_binary",
-    build_setting_default = True,
-)
-
-config_setting(
-    name = "is_controller_binary_bundled",
-    flag_values = {":bundle_controller_binary": "True"},
-)
-
-cc_library(
-    name = "engine_controller_with_subprocess",
-    srcs = ["engine_controller_with_subprocess.cc"],
-    hdrs = ["engine_controller_abi.h"],
-    data = select({
-        ":is_controller_binary_bundled": [
-            ":centipede_uninstrumented",
-        ],
-        "//conditions:default": [],
-    }),
-    deps = [
-        "@com_google_fuzztest//centipede:engine_abi",
-        "@com_google_fuzztest//fuzztest/internal:escaping",
-    ],
-)
-
 # The runner library is special:
 #   * It must not be instrumented with asan, sancov, etc.
 #   * It must not have heavy dependencies, and ideally not at all.
@@ -1051,8 +1005,6 @@
     "runner_dl_info.cc",
     "runner_dl_info.h",
     "runner_interface.h",
-    "runner_utils.cc",
-    "runner_utils.h",
     "sancov_callbacks.cc",
     "sancov_interceptors.cc",
     "sancov_object_array.cc",
@@ -1099,11 +1051,12 @@
     ":knobs",
     ":mutation_data",
     ":rolling_hash",
-    ":engine_abi",
     ":runner_cmp_trace",
-    ":runner_fork_server",
     ":runner_request",
     ":runner_result",
+    ":runner_utils",
+    ":engine_abi",
+    ":engine_worker",
     ":shared_memory_blob_sequence",
     "@com_google_fuzztest//common:defs",
     "@abseil-cpp//absl/base:core_headers",
@@ -1112,6 +1065,53 @@
     "@abseil-cpp//absl/types:span",
 ]
 
+cc_library(
+    name = "engine_worker",
+    srcs = [
+        "engine_worker.cc",
+    ],
+    hdrs = ["engine_worker_abi.h"],
+    copts = RUNNER_COPTS,
+    deps = [
+        ":engine_abi",
+        ":execution_metadata",
+        ":feature",
+        ":runner_fork_server",
+        ":runner_request",
+        ":runner_result",
+        ":runner_utils",
+        ":shared_memory_blob_sequence",
+        "@abseil-cpp//absl/base:nullability",
+        "@com_google_fuzztest//common:defs",
+    ],
+)
+
+bool_flag(
+    name = "bundle_controller_binary",
+    build_setting_default = True,
+)
+
+config_setting(
+    name = "is_controller_binary_bundled",
+    flag_values = {":bundle_controller_binary": "True"},
+)
+
+cc_library(
+    name = "engine_controller_with_subprocess",
+    srcs = ["engine_controller_with_subprocess.cc"],
+    hdrs = ["engine_controller_abi.h"],
+    data = select({
+        ":is_controller_binary_bundled": [
+            ":centipede_uninstrumented",
+        ],
+        "//conditions:default": [],
+    }),
+    deps = [
+        "@com_google_fuzztest//centipede:engine_abi",
+        "@com_google_fuzztest//fuzztest/internal:escaping",
+    ],
+)
+
 # A fuzz target needs to link with this library in order to run with Centipede.
 # The fuzz target must provide its own main().
 #
diff --git a/centipede/runner.cc b/centipede/runner.cc
index c2d3eef..e4d8384 100644
--- a/centipede/runner.cc
+++ b/centipede/runner.cc
@@ -24,6 +24,7 @@
 
 #include <fcntl.h>
 #include <pthread.h>  // NOLINT: use pthread to avoid extra dependencies.
+#include <sched.h>
 #include <sys/resource.h>
 #include <sys/socket.h>
 #include <sys/time.h>
@@ -51,6 +52,8 @@
 #include "absl/base/optimization.h"
 #include "absl/types/span.h"
 #include "./centipede/byte_array_mutator.h"
+#include "./centipede/engine_abi.h"
+#include "./centipede/engine_worker_abi.h"
 #include "./centipede/execution_metadata.h"
 #include "./centipede/feature.h"
 #include "./centipede/mutation_data.h"
@@ -77,6 +80,42 @@
 
 GlobalRunnerStateManager state_manager __attribute__((init_priority(200)));
 
+class SpinlockGuard {
+ public:
+  SpinlockGuard(const SpinlockGuard&) = delete;
+  SpinlockGuard& operator=(const SpinlockGuard&) = delete;
+  SpinlockGuard(SpinlockGuard&& other) : lock_(other.lock_) {
+    other.lock_ = nullptr;
+  }
+  SpinlockGuard& operator=(SpinlockGuard&& other) = delete;
+
+  static SpinlockGuard Acquire(std::atomic<bool>& lock) {
+    while (lock.exchange(true, std::memory_order_acquire)) {
+      sched_yield();
+    }
+    return SpinlockGuard{&lock};
+  }
+
+  static SpinlockGuard TryAcquire(std::atomic<bool>& lock) {
+    if (lock.exchange(true, std::memory_order_acquire)) {
+      return SpinlockGuard{nullptr};
+    }
+    return SpinlockGuard{&lock};
+  }
+
+  operator bool() const { return lock_ != nullptr; }
+
+  ~SpinlockGuard() {
+    if (lock_ != nullptr) {
+      lock_->store(false, std::memory_order_release);
+    }
+  }
+
+ private:
+  explicit SpinlockGuard(std::atomic<bool>* lock) : lock_(lock) {}
+  std::atomic<bool>* lock_;
+};
+
 }  // namespace
 
 static size_t GetPeakRSSMb() {
@@ -248,6 +287,16 @@
   state->input_start_time = curr_time;
 }
 
+// Returns a random seed. No need for a more sophisticated seed.
+static unsigned GetRandomSeed() { return time(nullptr); }
+
+static void EnsureBuiltinMutator() {
+  if (state->byte_array_mutator == nullptr) {
+    static auto* mutator = new ByteArrayMutator(state->knobs, GetRandomSeed());
+    state->byte_array_mutator = mutator;
+  }
+}
+
 // Byte array mutation fallback for a custom mutator, as defined here:
 // https://github.com/google/fuzzing/blob/master/docs/structure-aware-fuzzing.md
 extern "C" __attribute__((weak)) size_t
@@ -265,6 +314,7 @@
   }
 
   ByteArray array(data, data + size);
+  EnsureBuiltinMutator();
   state->byte_array_mutator->set_max_len(max_size);
   state->byte_array_mutator->Mutate(array);
   if (array.size() > max_size) {
@@ -492,90 +542,6 @@
   return true;
 }
 
-// Handles an ExecutionRequest, see RequestExecution(). Reads inputs from
-// `inputs_blobseq`, runs them, saves coverage features to `outputs_blobseq`.
-// Returns EXIT_SUCCESS on success and EXIT_FAILURE otherwise.
-static int ExecuteInputsFromShmem(BlobSequence& inputs_blobseq,
-                                  BlobSequence& outputs_blobseq,
-                                  RunnerCallbacks& callbacks) {
-  size_t num_inputs = 0;
-  if (!IsExecutionRequest(inputs_blobseq.Read())) return EXIT_FAILURE;
-  if (!IsNumInputs(inputs_blobseq.Read(), num_inputs)) return EXIT_FAILURE;
-
-  std::vector<void*> inputs;
-  inputs.reserve(num_inputs);
-  for (size_t i = 0; i < num_inputs; i++) {
-    auto blob = inputs_blobseq.Read();
-    // TODO(kcc): distinguish bad input from end of stream.
-    if (!blob.IsValid()) break;  // no more blobs to read.
-    if (!IsDataInput(blob)) return EXIT_FAILURE;
-
-    // TODO(kcc): [impl] handle sizes larger than kMaxDataSize.
-    size_t size = std::min(kMaxDataSize, blob.size);
-    inputs.push_back(callbacks.DeserializeInput({blob.data, size}));
-  }
-
-  CentipedeBeginExecutionBatch();
-
-  for (void* input : inputs) {
-    if (!StartSendingOutputsToEngine(outputs_blobseq)) break;
-
-    RunOneInput(input, callbacks);
-
-    if (state->has_failure_description.load()) break;
-
-    if (!FinishSendingOutputsToEngine(outputs_blobseq)) break;
-  }
-
-  CentipedeEndExecutionBatch();
-
-  for (void* input : inputs) {
-    callbacks.FreeInput(input);
-  }
-
-  return state->has_failure_description.load() ? EXIT_FAILURE : EXIT_SUCCESS;
-}
-
-// Dumps seed inputs to `output_dir`. Also see `GetSeedsViaExternalBinary()`.
-static void DumpSeedsToDir(RunnerCallbacks& callbacks, const char* output_dir) {
-  size_t seed_index = 0;
-  // Declare it on the outer scope to save allocations.
-  ByteArray serialized;
-  auto dump_seed_callback = [&](void* seed) {
-    serialized.clear();
-    callbacks.SerializeInput(seed, [&](ByteSpan bytes) {
-      // Cannot use `vector::insert` due to potential conflict of this file
-      // without sanitizers when linking other files with sanitizers that uses
-      // the same symbol. Other symbols used here seem safe. This a dirty hack
-      // that is expected to go away soon.
-      const size_t cur = serialized.size();
-      serialized.resize(cur + bytes.size());
-      std::memcpy(serialized.data() + cur, bytes.data(), bytes.size());
-    });
-    callbacks.FreeInput(seed);
-    // Cap seed index within 9 digits. If this was triggered, the dumping would
-    // take forever..
-    if (seed_index >= 1000000000) return;
-    char seed_path_buf[PATH_MAX];
-    const size_t num_path_chars =
-        snprintf(seed_path_buf, PATH_MAX, "%s/%09lu", output_dir, seed_index);
-    PrintErrorAndExitIf(num_path_chars >= PATH_MAX,
-                        "seed path reaches PATH_MAX");
-    FILE* output_file = fopen(seed_path_buf, "w");
-    const size_t num_bytes_written =
-        fwrite(serialized.data(), 1, serialized.size(), output_file);
-    PrintErrorAndExitIf(num_bytes_written != serialized.size(),
-                        "wrong number of bytes written for cf table");
-    fclose(output_file);
-    ++seed_index;
-  };
-  callbacks.GetPresetSeedInputs(dump_seed_callback);
-  const size_t min_seeds = state->flag_helper.GetIntFlag("min_seeds=", 32);
-  while (seed_index < min_seeds) {
-    dump_seed_callback(callbacks.GetRandomSeedInput());
-  }
-}
-
 // Dumps serialized target config to `output_file_path`. Also see
 // `GetSerializedTargetConfigViaExternalBinary()`.
 static void DumpSerializedTargetConfigToFile(RunnerCallbacks& callbacks,
@@ -590,127 +556,6 @@
   fclose(output_file);
 }
 
-// Returns a random seed. No need for a more sophisticated seed.
-// TODO(kcc): [as-needed] optionally pass an external seed.
-static unsigned GetRandomSeed() { return time(nullptr); }
-
-void MutateInputs(RunnerCallbacks& callbacks,
-                  absl::Span<const MutationInputRef> inputs, size_t num_mutants,
-                  std::function<void(MutantRef)> new_mutant_callback) {
-  static unsigned int seed = GetRandomSeed();
-  std::vector<void*> input_objects;
-  input_objects.resize(inputs.size(), nullptr);
-  for (size_t i = 0; i < inputs.size(); ++i) {
-    input_objects[i] = callbacks.DeserializeInput(inputs[i].data);
-  }
-  ByteArray mutant_data;
-  ExecutionMetadata empty_metadata;
-  for (size_t i = 0; i < num_mutants; ++i) {
-    const size_t origin_index = rand_r(&seed) % inputs.size();
-    void* mutant = nullptr;
-    if (rand_r(&seed) % 100 < state->run_time_flags.crossover_level) {
-      // Perform crossover `crossover_level`% of the time.
-      const size_t other_index = rand_r(&seed) % inputs.size();
-      mutant = callbacks.CrossOver(input_objects[origin_index],
-                                   inputs[origin_index].metadata != nullptr
-                                       ? *inputs[origin_index].metadata
-                                       : empty_metadata,
-                                   input_objects[other_index],
-                                   inputs[other_index].metadata != nullptr
-                                       ? *inputs[other_index].metadata
-                                       : empty_metadata);
-    } else {
-      mutant = callbacks.Mutate(input_objects[origin_index],
-                                inputs[origin_index].metadata != nullptr
-                                    ? *inputs[origin_index].metadata
-                                    : empty_metadata);
-    }
-    mutant_data.clear();
-    callbacks.SerializeInput(mutant, [&mutant_data](ByteSpan bytes) {
-      // Cannot use `vector::insert` due to potential conflict of this file
-      // without sanitizers when linking other files with sanitizers that uses
-      // the same symbol. Other symbols used here seem safe.
-      // This a dirty hack that is expected to go away soon.
-      const size_t cur = mutant_data.size();
-      mutant_data.resize(cur + bytes.size());
-      std::memcpy(mutant_data.data() + cur, bytes.data(), bytes.size());
-    });
-    new_mutant_callback(
-        MutantRef{{(unsigned char*)mutant_data.data(), mutant_data.size()},
-                  origin_index});
-    callbacks.FreeInput(mutant);
-  }
-  for (void* input_object : input_objects) {
-    callbacks.FreeInput(input_object);
-  }
-}
-
-// Handles a Mutation Request, see RequestMutation().
-// Mutates inputs read from `inputs_blobseq`,
-// writes the mutants to `outputs_blobseq`
-// Returns EXIT_SUCCESS on success and EXIT_FAILURE on failure
-// so that main() can return its result.
-// If both `custom_mutator_cb` and `custom_crossover_cb` are nullptr,
-// returns EXIT_FAILURE.
-//
-// TODO(kcc): [impl] make use of custom_crossover_cb, if available.
-static int MutateInputsFromShmem(BlobSequence& inputs_blobseq,
-                                 BlobSequence& outputs_blobseq,
-                                 RunnerCallbacks& callbacks) {
-  // Read max_num_mutants.
-  size_t num_mutants = 0;
-  size_t num_inputs = 0;
-  if (!IsMutationRequest(inputs_blobseq.Read())) return EXIT_FAILURE;
-  if (!IsNumMutants(inputs_blobseq.Read(), num_mutants)) return EXIT_FAILURE;
-  if (!IsNumInputs(inputs_blobseq.Read(), num_inputs)) return EXIT_FAILURE;
-
-  // Mutation input with ownership.
-  struct MutationInput {
-    ByteArray data;
-    ExecutionMetadata metadata;
-  };
-  // TODO(kcc): unclear if we can continue using std::vector (or other STL)
-  // in the runner. But for now use std::vector.
-  // Collect the inputs into a vector. We copy them instead of using pointers
-  // into shared memory so that the user code doesn't touch the shared memory.
-  std::vector<MutationInput> inputs;
-  inputs.reserve(num_inputs);
-  std::vector<MutationInputRef> input_refs;
-  input_refs.reserve(num_inputs);
-  for (size_t i = 0; i < num_inputs; ++i) {
-    // If inputs_blobseq have overflown in the engine, we still want to
-    // handle the first few inputs.
-    ExecutionMetadata metadata;
-    if (!IsExecutionMetadata(inputs_blobseq.Read(), metadata)) {
-      break;
-    }
-    auto blob = inputs_blobseq.Read();
-    if (!IsDataInput(blob)) break;
-    inputs.push_back(
-        MutationInput{/*data=*/ByteArray{blob.data, blob.data + blob.size},
-                      /*metadata=*/std::move(metadata)});
-    input_refs.push_back(
-        MutationInputRef{/*data=*/inputs.back().data,
-                         /*metadata=*/&inputs.back().metadata});
-  }
-
-  if (!inputs.empty()) {
-    state->byte_array_mutator->SetMetadata(inputs[0].metadata);
-  }
-
-  if (!MutationResult::WriteHasCustomMutator(callbacks.HasCustomMutator(),
-                                             outputs_blobseq)) {
-    return EXIT_FAILURE;
-  }
-  if (!callbacks.HasCustomMutator()) return EXIT_SUCCESS;
-
-  MutateInputs(callbacks, input_refs, num_mutants, [&](MutantRef mutant) {
-    (void)MutationResult::WriteMutant(mutant, outputs_blobseq);
-  });
-
-  return EXIT_SUCCESS;
-}
-
 void LegacyRunnerCallbacks::SerializeInput(
     void* input, const std::function<void(ByteSpan)>& bytes_sink) {
   const auto* input_object = reinterpret_cast<const ByteArray*>(input);
@@ -729,6 +574,7 @@
 void* LegacyRunnerCallbacks::Mutate(void* origin,
                                     const ExecutionMetadata& origin_metadata) {
   const auto* origin_ba = reinterpret_cast<const ByteArray*>(origin);
+  EnsureBuiltinMutator();
   state->byte_array_mutator->SetMetadata(origin_metadata);
   const size_t max_mutant_size = state->run_time_flags.max_len;
   const size_t size = std::min(max_mutant_size, origin_ba->size());
@@ -759,6 +605,7 @@
   const auto* origin_ba = reinterpret_cast<const ByteArray*>(origin);
   const auto* other_ba = reinterpret_cast<const ByteArray*>(other);
   static unsigned int seed = GetRandomSeed();
+  EnsureBuiltinMutator();
   state->byte_array_mutator->SetMetadata((rand_r(&seed) & 1) ? origin_metadata
                                                              : other_metadata);
   const size_t max_mutant_size = state->run_time_flags.max_len;
@@ -824,70 +671,8 @@
   }
 }
 
-// Create a fake reference to ForkServerCallMeVeryEarly() here so that the
-// fork server module is not dropped during linking.
-// Alternatives are
-//  * Use -Wl,--whole-archive when linking with the runner archive.
-//  * Use -Wl,-u,ForkServerCallMeVeryEarly when linking with the runner archive.
-//    (requires ForkServerCallMeVeryEarly to be extern "C").
-// These alternatives require extra flags and are thus more fragile.
-// We declare ForkServerCallMeVeryEarly() here instead of doing it in some
-// header file, because we want to keep the fork server header-free.
-extern void ForkServerCallMeVeryEarly();
-[[maybe_unused]] auto fake_reference_for_fork_server =
-    &ForkServerCallMeVeryEarly;
-
-void MaybeConnectToPersistentMode() {
-  if (state->persistent_mode_socket_path == nullptr) {
-    return;
-  }
-  state->persistent_mode_socket = socket(AF_UNIX, SOCK_STREAM, 0);
-  if (state->persistent_mode_socket < 0) {
-    fprintf(stderr, "Failed to create persistent mode socket\n");
-  }
-
-  struct sockaddr_un addr{};
-  addr.sun_family = AF_UNIX;
-  const size_t socket_path_len = strlen(state->persistent_mode_socket_path);
-  RunnerCheck(
-      socket_path_len < sizeof(addr.sun_path),
-      "persistent mode socket path string must be fit in sockaddr_un.sun_path");
-  std::memcpy(addr.sun_path, state->persistent_mode_socket_path,
-              socket_path_len);
-
-  int connect_ret = 0;
-  do {
-    connect_ret = connect(state->persistent_mode_socket,
-                          (struct sockaddr*)&addr, sizeof(addr));
-  } while (connect_ret == -1 && errno == EINTR);
-  if (connect_ret == -1) {
-    fprintf(stderr, "Failed to connect the persistent mode socket to %s\n",
-            state->persistent_mode_socket_path);
-    (void)close(state->persistent_mode_socket);
-    state->persistent_mode_socket = -1;
-  }
-
-  int flags = fcntl(state->persistent_mode_socket, F_GETFD);
-  if (flags == -1) {
-    fprintf(stderr, "fcntl(F_GETFD) failed\n");
-    (void)close(state->persistent_mode_socket);
-    state->persistent_mode_socket = -1;
-  }
-  flags |= FD_CLOEXEC;
-  if (fcntl(state->persistent_mode_socket, F_SETFD, flags) == -1) {
-    fprintf(stderr, "fcntl(F_SETFD) failed\n");
-    (void)close(state->persistent_mode_socket);
-    state->persistent_mode_socket = -1;
-  }
-}
-
 GlobalRunnerState::GlobalRunnerState() {
-  // Make sure fork server is started if needed.
-  ForkServerCallMeVeryEarly();
-
-  // Connecting to the persistent mode socket should be immediately after.
-  MaybeConnectToPersistentMode();
-
+  FuzzTestWorkerInitEarly();
   SancovRuntimeInitialize();
 
   // TODO(kcc): move some code from CentipedeRunnerMain() here so that it works
@@ -918,84 +703,223 @@
   }
 }
 
-static int HandleSharedMemoryRequest(RunnerCallbacks& callbacks,
-                                     BlobSequence& inputs_blobseq,
-                                     BlobSequence& outputs_blobseq) {
-  state->has_failure_description = false;
-  // Read the first blob. It indicates what further actions to take.
-  auto request_type_blob = inputs_blobseq.Read();
-  if (IsMutationRequest(request_type_blob)) {
-    // Mutation request.
-    inputs_blobseq.Reset();
-    static auto mutator = new ByteArrayMutator(state->knobs, GetRandomSeed());
-    state->byte_array_mutator = mutator;
-    // Since we are mutating, no need to spend time collecting the coverage.
-    // We still pay for executing the coverage callbacks, but those will
-    // return immediately.
-    const int old_traced = CentipedeSetCurrentThreadTraced(/*traced=*/0);
-    const int result =
-        MutateInputsFromShmem(inputs_blobseq, outputs_blobseq, callbacks);
-    CentipedeSetCurrentThreadTraced(old_traced);
-    return result;
-  }
-  if (IsExecutionRequest(request_type_blob)) {
-    // Execution request.
-    inputs_blobseq.Reset();
-    return ExecuteInputsFromShmem(inputs_blobseq, outputs_blobseq, callbacks);
-  }
-  return EXIT_FAILURE;
+namespace {
+
+struct Input {
+  void* content;
+  ExecutionMetadata metadata;
+};
+
+struct RunnerAdapterCtx {
+  bool need_full_cleanup;
+  RunnerCallbacks& callbacks;
+};
+
+void SetUpCoverageDomains(FuzzTestAdapterCtx* ctx,
+                          const FuzzTestCoverageDomainRegistry* registry) {
+  SanCovRuntimeSetUpCoverageDomains(registry);
 }
 
-static int HandlePersistentMode(RunnerCallbacks& callbacks,
-                                BlobSequence& inputs_blobseq,
-                                BlobSequence& outputs_blobseq) {
-  bool first = true;
-  while (true) {
-    PersistentModeRequest req;
-    if (!ReadAll(state->persistent_mode_socket, reinterpret_cast<char*>(&req),
-                 1)) {
-      perror("Failed to read request from persistent mode socket");
-      return EXIT_FAILURE;
-    }
-    if (first) {
-      first = false;
-    } else {
-      // Reset stdout/stderr.
-      for (int fd = 1; fd <= 2; fd++) {
-        lseek(fd, 0, SEEK_SET);
-        // NOTE: Allow ftruncate() to fail by ignoring its return; that's okay
-        // to happen when the stdout/stderr are not redirected to a file.
-        (void)ftruncate(fd, 0);
-      }
-      fprintf(
-          stderr, "Centipede fuzz target runner (%s); flags: %s\n",
-          req == PersistentModeRequest::kExit ? "exiting persistent mode"
-                                              : "persistent mode batch",
-          CentipedeGetRunnerFlags() ? CentipedeGetRunnerFlags() : "(unset)");
-    }
-    if (req == PersistentModeRequest::kExit) break;
-    RunnerCheck(req == PersistentModeRequest::kRunBatch,
-                "Unknown persistent mode request");
-    const int result =
-        HandleSharedMemoryRequest(callbacks, inputs_blobseq, outputs_blobseq);
-    inputs_blobseq.Reset();
-    outputs_blobseq.Reset();
-    if (!WriteAll(state->persistent_mode_socket,
-                  reinterpret_cast<const char*>(&result), sizeof(result))) {
-      perror("Failed to write response to the persistent mode socket");
-      return EXIT_FAILURE;
-    }
-  }
-  return EXIT_SUCCESS;
+void GetPresetSeedInputs(FuzzTestAdapterCtx* ctx,
+                         const FuzzTestInputSink* sink) {
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  runner_ctx->callbacks.GetPresetSeedInputs([&](void* input) {
+    sink->Emit(sink->ctx, reinterpret_cast<FuzzTestInputHandle>(
+                              new Input{/*content=*/input,
+                                        /*metadata=*/{}}));
+  });
 }
 
-// If state->run_time_flags.shmem_size_mb is non-zero, state->arg1 and
-// state->arg2 are the names of in/out shared memory locations. Read inputs and
-// write outputs via shared memory.
+void GetRandomSeedInput(FuzzTestAdapterCtx* ctx,
+                        const FuzzTestInputSink* sink) {
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  sink->Emit(sink->ctx,
+             reinterpret_cast<FuzzTestInputHandle>(new Input{
+                 /*content=*/runner_ctx->callbacks.GetRandomSeedInput(),
+                 /*metadata=*/{}}));
+}
+
+void NotifyControllerAndExitIfNoCustomMutator(RunnerCallbacks& callbacks) {
+  if (!callbacks.HasCustomMutator()) {
+    // This is a ugly hack to let Centipede use the builtin mutator
+    // needed by the LLVM fuzzer runners (used by some engine tests). We
+    // cannot remove it until the LLVM fuzzers can be migrated to use
+    // the FuzzTest LLVM fuzzer wrapper.
+    SharedMemoryBlobSequence outputs_blobseq(
+        sancov_state->arg2,
+        sancov_state->flag_helper.GetIntFlag("shmem_size_mb=", 0) << 20);
+    RunnerCheck(MutationResult::WriteHasCustomMutator(false, outputs_blobseq),
+                "Failed to write the indicator for no custom mutator");
+    std::_Exit(0);
+  }
+}
+
+void Mutate(FuzzTestAdapterCtx* ctx, FuzzTestInputHandle origin_handle,
+            int shrink, const FuzzTestInputSink* sink) {
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  runner_ctx->need_full_cleanup = true;
+  NotifyControllerAndExitIfNoCustomMutator(runner_ctx->callbacks);
+  const auto* origin = reinterpret_cast<Input*>(origin_handle);
+  auto mutant = new Input{
+      /*content=*/runner_ctx->callbacks.Mutate(origin->content,
+                                               origin->metadata),
+      /*metadata=*/{},
+  };
+  sink->Emit(sink->ctx, reinterpret_cast<FuzzTestInputHandle>(mutant));
+}
+
+void CrossOver(FuzzTestAdapterCtx* ctx, FuzzTestInputHandle origin_handle,
+               FuzzTestInputHandle other_handle,
+               const FuzzTestInputSink* sink) {
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  runner_ctx->need_full_cleanup = true;
+  NotifyControllerAndExitIfNoCustomMutator(runner_ctx->callbacks);
+  const auto* origin = reinterpret_cast<Input*>(origin_handle);
+  const auto* other = reinterpret_cast<Input*>(other_handle);
+  auto mutant = new Input{
+      /*content=*/runner_ctx->callbacks.CrossOver(
+          origin->content, origin->metadata, other->content, other->metadata),
+      /*metadata=*/{},
+  };
+  sink->Emit(sink->ctx, reinterpret_cast<FuzzTestInputHandle>(mutant));
+}
+
+void Execute(FuzzTestAdapterCtx* ctx, FuzzTestInputHandle handle,
+             const FuzzTestFeedbackSink* sink) {
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  auto* input = reinterpret_cast<Input*>(handle);
+  SanCovRuntimeClearCoverage(runner_ctx->need_full_cleanup);
+  if (runner_ctx->need_full_cleanup) {
+    runner_ctx->need_full_cleanup = false;
+  }
+  const int old_traced = CentipedeSetCurrentThreadTraced(/*traced=*/1);
+  RunOneInput(input->content, runner_ctx->callbacks);
+  CentipedeSetCurrentThreadTraced(old_traced);
+  {
+    LockGuard lock(state->execution_result_override_mu);
+    bool has_overridden_execution_result = false;
+    if (state->execution_result_override != nullptr) {
+      RunnerCheck(state->execution_result_override->results().size() <= 1,
+                  "unexpected number of overridden execution results");
+      has_overridden_execution_result =
+          state->execution_result_override->results().size() == 1;
+    }
+    if (has_overridden_execution_result) {
+      auto& result = state->execution_result_override->results()[0];
+      SanCovRuntimeConvertToEngineFeatures(result.mutable_features().data(),
+                                           result.mutable_features().size());
+      const FuzzTestUint64sView features = {
+          result.features().data(),
+          result.features().size(),
+      };
+      sink->EmitCoverageFeatures(sink->ctx, &features);
+      input->metadata = result.metadata();
+      return;
+    }
+  }
+  SanCovRuntimeEmitFeatures(sink);
+  input->metadata = SanCovRuntimeGetExecutionMetadata();
+}
+
+void DeserializeInputContent(FuzzTestAdapterCtx* ctx,
+                             const FuzzTestBytesView* view,
+                             const FuzzTestInputSink* sink) {
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  auto* input = new Input{
+      /*content=*/runner_ctx->callbacks.DeserializeInput(
+          {view->data, view->size}),
+      /*metadata=*/{},
+  };
+  sink->Emit(sink->ctx, reinterpret_cast<FuzzTestInputHandle>(input));
+}
+
+void UpdateInputMetadata(FuzzTestAdapterCtx* ctx, const FuzzTestBytesView* view,
+                         FuzzTestInputHandle handle) {
+  auto* input = reinterpret_cast<Input*>(handle);
+  input->metadata.cmp_data = {view->data, view->data + view->size};
+}
+
+void SerializeInputContent(FuzzTestAdapterCtx* ctx, FuzzTestInputHandle handle,
+                           const FuzzTestBytesSink* sink) {
+  auto* input = reinterpret_cast<Input*>(handle);
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  runner_ctx->callbacks.SerializeInput(input->content, [&](ByteSpan bytes) {
+    const auto input_bytes = FuzzTestBytesView{
+        reinterpret_cast<const uint8_t*>(bytes.data()),
+        bytes.size(),
+    };
+    sink->Emit(sink->ctx, &input_bytes);
+  });
+}
+
+void SerializeInputMetadata(FuzzTestAdapterCtx* ctx, FuzzTestInputHandle handle,
+                            const FuzzTestBytesSink* sink) {
+  auto* input = reinterpret_cast<Input*>(handle);
+  const auto input_bytes = FuzzTestBytesView{
+      reinterpret_cast<const uint8_t*>(input->metadata.cmp_data.data()),
+      input->metadata.cmp_data.size(),
+  };
+  sink->Emit(sink->ctx, &input_bytes);
+}
+
+void FreeInput(FuzzTestAdapterCtx* ctx, FuzzTestInputHandle handle) {
+  auto* input = reinterpret_cast<Input*>(handle);
+  auto* runner_ctx = reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  runner_ctx->callbacks.FreeInput(input->content);
+  delete input;
+}
+
+void FreeCtx(FuzzTestAdapterCtx* ctx) {
+  delete reinterpret_cast<RunnerAdapterCtx*>(ctx);
+  auto guard = SpinlockGuard::Acquire(state->diagnostic_sink_spinlock);
+  state->diagnostic_sink = nullptr;
+}
+
+void GetTestName(FuzzTestAdapterManagerCtx* ctx,
+                 const FuzzTestBytesSink* sink) {
+  // Provide the test name exactly specified from the flag. This is hacky
+  // but should work as the user of runner should call RunnerMain at most
+  // once.
+  static const char* test_name = state->flag_helper.GetStringFlag("test=");
+  if (test_name == nullptr) return;
+  static size_t len = strlen(test_name);
+  const auto bytes = FuzzTestBytesView{
+      reinterpret_cast<const uint8_t*>(test_name),
+      len,
+  };
+  sink->Emit(sink->ctx, &bytes);
+}
+
+void ConstructAdapter(FuzzTestAdapterManagerCtx* ctx,
+                      const FuzzTestDiagnosticSink* diagnostic_sink,
+                      FuzzTestAdapter* adapter_out) {
+  {
+    auto guard = SpinlockGuard::Acquire(state->diagnostic_sink_spinlock);
+    state->diagnostic_sink = diagnostic_sink;
+  }
+  adapter_out->ctx = reinterpret_cast<FuzzTestAdapterCtx*>(new RunnerAdapterCtx{
+      /*need_full_cleanup=*/true, *reinterpret_cast<RunnerCallbacks*>(ctx)});
+
+  adapter_out->SetUpCoverageDomains = SetUpCoverageDomains;
+  adapter_out->GetPresetSeedInputs = GetPresetSeedInputs;
+  adapter_out->GetRandomSeedInput = GetRandomSeedInput;
+  adapter_out->Mutate = Mutate;
+  adapter_out->CrossOver = CrossOver;
+  adapter_out->Execute = Execute;
+  adapter_out->DeserializeInputContent = DeserializeInputContent;
+  adapter_out->UpdateInputMetadata = UpdateInputMetadata;
+  adapter_out->SerializeInputContent = SerializeInputContent;
+  adapter_out->SerializeInputMetadata = SerializeInputMetadata;
+  adapter_out->FreeInput = FreeInput;
+  adapter_out->FreeCtx = FreeCtx;
+}
+
+}  // namespace
+
+// Constructs a FuzzTestAdapterManager and try to run as an engine worker.
 //
-//  Default: Execute ReadOneInputExecuteItAndDumpCoverage() for all inputs.//
-//
-//  Note: argc/argv are used for only ReadOneInputExecuteItAndDumpCoverage().
+// If engine worker is not required, calls
+// ReadOneInputExecuteItAndDumpCoverage() for all input files from argc/argv.
 int RunnerMain(int argc, char** argv, RunnerCallbacks& callbacks) {
   state->centipede_runner_main_executed = true;
 
@@ -1009,32 +933,22 @@
     return EXIT_SUCCESS;
   }
 
-  if (state->flag_helper.HasSwitchFlag("dump_seed_inputs")) {
-    // Seed request.
-    DumpSeedsToDir(callbacks, /*output_dir=*/sancov_state->arg1);
+  FuzzTestAdapterManager manager = {
+      /*ctx=*/reinterpret_cast<FuzzTestAdapterManagerCtx*>(&callbacks),
+      /*GetBinaryId=*/nullptr,
+      /*GetTestName=*/GetTestName,
+      /*ConstructAdapter=*/ConstructAdapter,
+  };
+  const int old_traced = CentipedeSetCurrentThreadTraced(/*traced=*/0);
+  const auto s = FuzzTestWorkerMaybeRun(&manager);
+  CentipedeSetCurrentThreadTraced(old_traced);
+  if (s == kFuzzTestWorkerNotRequired) {
+    for (int i = 1; i < argc; i++) {
+      ReadOneInputExecuteItAndDumpCoverage(argv[i], callbacks);
+    }
     return EXIT_SUCCESS;
   }
-
-  // Inputs / outputs from shmem.
-  if (state->run_time_flags.shmem_size_mb != 0) {
-    if (!sancov_state->arg1 || !sancov_state->arg2) return EXIT_FAILURE;
-    SharedMemoryBlobSequence inputs_blobseq(
-        sancov_state->arg1, state->run_time_flags.shmem_size_mb << 20);
-    SharedMemoryBlobSequence outputs_blobseq(
-        sancov_state->arg2, state->run_time_flags.shmem_size_mb << 20);
-    // Persistent mode loop.
-    if (state->persistent_mode_socket > 0) {
-      return HandlePersistentMode(callbacks, inputs_blobseq, outputs_blobseq);
-    }
-    return HandleSharedMemoryRequest(callbacks, inputs_blobseq,
-                                     outputs_blobseq);
-  }
-
-  // By default, run every input file one-by-one.
-  for (int i = 1; i < argc; i++) {
-    ReadOneInputExecuteItAndDumpCoverage(argv[i], callbacks);
-  }
-  return EXIT_SUCCESS;
+  return s == kFuzzTestWorkerSuccess ? EXIT_SUCCESS : EXIT_FAILURE;
 }
 
 }  // namespace fuzztest::internal
@@ -1157,22 +1071,51 @@
 }
 
 extern "C" void CentipedeSetFailureDescription(const char* description) {
-  using fuzztest::internal::state;
-  if (state->failure_description_path == nullptr) return;
-  if (state->has_failure_description.exchange(true)) return;
-  FILE* f = fopen(state->failure_description_path, "w");
-  if (f == nullptr) {
-    perror("FAILURE: fopen()");
+  if (description == nullptr) {
+    static constexpr std::string_view kMsg =
+        "CentipedeSetFailureDescription called with null description\n";
+    write(STDERR_FILENO, kMsg.data(), kMsg.size());
+    std::_Exit(EXIT_FAILURE);
+  }
+  std::string_view desc_sv = description;
+  static constexpr std::string_view kSetupFailurePrefix = "SETUP FAILURE:";
+  using ::fuzztest::internal::SpinlockGuard;
+  using ::fuzztest::internal::state;
+  auto sink_guard = SpinlockGuard::TryAcquire(state->diagnostic_sink_spinlock);
+  if (!sink_guard) {
+    static constexpr std::string_view kMsg =
+        "Diagnostic sink is busy while setting failure:\n";
+    write(STDERR_FILENO, kMsg.data(), kMsg.size());
+    write(STDERR_FILENO, description, std::strlen(description));
+    std::_Exit(EXIT_FAILURE);
+  }
+  const auto* diagnostic_sink = state->diagnostic_sink;
+  if (diagnostic_sink == nullptr) {
+    static constexpr std::string_view kMsg =
+        "Diagnostic sink is missing while setting failure:\n";
+    write(STDERR_FILENO, kMsg.data(), kMsg.size());
+    write(STDERR_FILENO, description, std::strlen(description));
+    std::_Exit(EXIT_FAILURE);
+  }
+
+  if (desc_sv.substr(0, kSetupFailurePrefix.size()) == kSetupFailurePrefix) {
+    size_t prefix_size = kSetupFailurePrefix.size();
+    while (prefix_size < desc_sv.size() && desc_sv[prefix_size] == ' ') {
+      ++prefix_size;
+    }
+    const FuzzTestBytesView error = {
+        reinterpret_cast<const uint8_t*>(desc_sv.data() + prefix_size),
+        desc_sv.size() - prefix_size,
+    };
+    diagnostic_sink->EmitError(diagnostic_sink->ctx, &error);
     return;
   }
-  const auto len = strlen(description);
-  if (fwrite(description, 1, len, f) != len) {
-    perror("FAILURE: fwrite()");
-  }
-  if (fflush(f) != 0) {
-    perror("FAILURE: fflush()");
-  }
-  if (fclose(f) != 0) {
-    perror("FAILURE: fclose()");
-  }
+
+  const FuzzTestBytesView description_view = {
+      reinterpret_cast<const uint8_t*>(desc_sv.data()),
+      desc_sv.size(),
+  };
+  diagnostic_sink->EmitFinding(diagnostic_sink->ctx, &description_view,
+                               // For now, use the description as the signature.
+                               &description_view);
 }
diff --git a/centipede/runner.h b/centipede/runner.h
index b64f732..2944da3 100644
--- a/centipede/runner.h
+++ b/centipede/runner.h
@@ -23,6 +23,7 @@
 #include <cstdint>
 
 #include "./centipede/byte_array_mutator.h"
+#include "./centipede/engine_abi.h"
 #include "./centipede/knobs.h"
 #include "./centipede/runner_interface.h"
 #include "./centipede/runner_result.h"
@@ -52,7 +53,7 @@
 // TODO(kcc): use a CTOR with absl::kConstInit (will require refactoring).
 struct GlobalRunnerState {
   // Used by LLVMFuzzerMutate and initialized in main().
-  ByteArrayMutator *byte_array_mutator = nullptr;
+  ByteArrayMutator* byte_array_mutator = nullptr;
   Knobs knobs;
 
   GlobalRunnerState();
@@ -92,7 +93,7 @@
   // execution result of the current test input. The object is owned and cleaned
   // up by the state, protected by execution_result_override_mu, and set by
   // `CentipedeSetExecutionResult()`.
-  BatchResult *execution_result_override;
+  BatchResult* execution_result_override;
 
   // Execution stats for the currently executed input.
   ExecutionResult::Stats stats;
@@ -114,6 +115,10 @@
 
   // The Watchdog thread sets this to true.
   std::atomic<bool> watchdog_thread_started;
+
+  // Engine diagnostic sink, protected by a spinlock.
+  std::atomic<bool> diagnostic_sink_spinlock;
+  const FuzzTestDiagnosticSink* diagnostic_sink;
 };
 
 extern ExplicitLifetime<GlobalRunnerState> state;