Develop · Guide

Write a chain operator.

This guide builds bytecount, an operator that sits between the filesystem and the cache, passes every call through unchanged, and counts the bytes written to each file. It is small, but it has everything a real operator needs.

Time
30 minutes
You need
CLIO Core built from source, on Linux
Result
A module loaded from your own build directory

1. Install CLIO Core into a prefix

An operator builds against an installed CLIO Core, which provides the headers, the libraries and the CMake package. From your CLIO Core source build:

shell
$ cmake --install build --prefix $HOME/clio-inst
$ export PATH=$HOME/clio-inst/bin:$PATH
$ export PYTHONPATH=$HOME/clio-inst/ext

Use the installed clio_run and clio_cte_fuse from now on, so the runtime and your module load the same CLIO Core libraries.

2. Lay out the module

files
bytecount/
  CMakeLists.txt
  clio_mod.yaml
  include/clio_demo/bytecount/
    bytecount_methods.h     method ids
    bytecount_tasks.h       configuration and task types
    bytecount_client.h      client
    bytecount_runtime.h     the operator
  src/
    bytecount_client.cc
    bytecount_runtime.cc    the handlers
    bytecount_exec.cc       dispatch

3. Name the module and its methods

clio_mod.yaml names the module and lists its methods. An operator lists the storage engine methods it intercepts, with the engine's own ids: kPutBlob is 15 and kMultiPutBlob is 48. A write can arrive as either one, so the operator handles both.

yaml · clio_mod.yaml
module_name: bytecount
namespace: clio::demo
version: 1.0.0

kCreate: 0
kDestroy: 1
kMonitor: 9

# Core verbs this operator intercepts. The ids must equal the CTE core's.
kPutBlob: 15
kMultiPutBlob: 48

The namespace clio::demo and the module name bytecount give the library names libclio_demo_bytecount_runtime.so and libclio_demo_bytecount_client.so.

cmake · CMakeLists.txt
cmake_minimum_required(VERSION 3.20)
project(bytecount CXX)
set(CMAKE_CXX_STANDARD 20)

find_package(yaml-cpp REQUIRED)        # clio-core uses it but does not find it itself
find_package(clio-core CONFIG REQUIRED)

add_clio_module_client(SOURCES src/bytecount_client.cc)
add_clio_module_runtime(SOURCES src/bytecount_runtime.cc src/bytecount_exec.cc)

target_link_libraries(clio_demo_bytecount_client PUBLIC clio::cte::core_client)
target_link_libraries(clio_demo_bytecount_runtime PUBLIC clio::cte::core_client)
c++ · bytecount_methods.h
#ifndef CLIO_DEMO_BYTECOUNT_METHODS_H_
#define CLIO_DEMO_BYTECOUNT_METHODS_H_

#include <clio_runtime/clio_runtime.h>
#include <clio_cte/core/autogen/core_methods.h>
#include <string>
#include <vector>

namespace clio::demo::bytecount::Method {

GLOBAL_CROSS_CONST clio::run::u32 kCreate = 0;
GLOBAL_CROSS_CONST clio::run::u32 kDestroy = 1;
GLOBAL_CROSS_CONST clio::run::u32 kMonitor = 9;

// Intercepted core verbs: same ids and task structs as the CTE core.
GLOBAL_CROSS_CONST clio::run::u32 kPutBlob = clio::cte::core::Method::kPutBlob;
GLOBAL_CROSS_CONST clio::run::u32 kMultiPutBlob =
    clio::cte::core::Method::kMultiPutBlob;

// Above the core's id space, so forwarded core ids fit in the model.
GLOBAL_CROSS_CONST clio::run::u32 kMaxMethodId = 101;

/** Names for the dashboard and the scheduler's per-method model. */
inline const std::vector<std::string> &GetMethodNames() {
  static const std::vector<std::string> names = [] {
    std::vector<std::string> v(kMaxMethodId);
    v[kCreate] = "Create";
    v[kDestroy] = "Destroy";
    v[kMonitor] = "Monitor";
    v[kPutBlob] = "PutBlob";
    v[kMultiPutBlob] = "MultiPutBlob";
    return v;
  }();
  return names;
}

}  // namespace clio::demo::bytecount::Method

#endif  // CLIO_DEMO_BYTECOUNT_METHODS_H_

4. Read the configuration

The configuration struct is filled from the module's entry in clio.yaml. An operator needs at least next_pool_id, the pool it forwards to. chimod_lib_name is the name you write as mod_name.

c++ · bytecount_tasks.h
#ifndef CLIO_DEMO_BYTECOUNT_TASKS_H_
#define CLIO_DEMO_BYTECOUNT_TASKS_H_

#include <clio_runtime/clio_runtime.h>
#include <clio_runtime/admin/admin_tasks.h>
#include <clio_cte/core/core_tasks.h>
#include <clio_demo/bytecount/bytecount_methods.h>

#include <string>

namespace clio::demo::bytecount {

/** Settings read from this pool's compose entry in clio.yaml. */
struct BytecountConfig {
  static constexpr const char *chimod_lib_name = "clio_demo_bytecount";

  clio::run::PoolId next_pool_id_;  ///< Pool to forward every call to

  BytecountConfig() : next_pool_id_(clio::run::PoolId::GetNull()) {}
  BytecountConfig(const clio::run::PoolId &, const BytecountConfig &other)
      : next_pool_id_(other.next_pool_id_) {}

  template <class Archive>
  void serialize(Archive &ar) {
    ar(next_pool_id_);
  }

  /** Parse next_pool_id ("major.minor") from the compose YAML. */
  void LoadConfig(const clio::run::PoolConfig &pool_config) {
    if (pool_config.config_.empty()) return;
    YAML::Node node = YAML::Load(pool_config.config_);
    if (node["next_pool_id"]) {
      auto s = node["next_pool_id"].as<std::string>();
      auto dot = s.find('.');
      next_pool_id_ = clio::run::PoolId(std::stoul(s.substr(0, dot)),
                                        std::stoul(s.substr(dot + 1)));
    }
  }
};

using CreateTask = clio::run::admin::GetOrCreatePoolTask<BytecountConfig>;
using MonitorTask = clio::run::admin::MonitorTask;

/** Tear down the container. Carries no fields of its own. */
struct DestroyTask : public clio::run::Task {
  DestroyTask() : clio::run::Task() {}
  DestroyTask(const clio::run::TaskId &task_id,
              const clio::run::PoolId &pool_id,
              const clio::run::PoolQuery &pool_query)
      : clio::run::Task(task_id, pool_id, pool_query, Method::kDestroy) {}

  void AggregateOut(const ctp::ipc::FullPtr<clio::run::Task> &other) {
    Task::AggregateOut(other);
  }
  void Copy(const ctp::ipc::FullPtr<DestroyTask> &other) {
    Task::Copy(other.template Cast<clio::run::Task>());
  }
  template <typename Ar> void SerializeIn(Ar &ar) { Task::SerializeIn(ar); }
  template <typename Ar> void SerializeOut(Ar &ar) { Task::SerializeOut(ar); }
};

}  // namespace clio::demo::bytecount

#endif  // CLIO_DEMO_BYTECOUNT_TASKS_H_
c++ · bytecount_client.h
#ifndef CLIO_DEMO_BYTECOUNT_CLIENT_H_
#define CLIO_DEMO_BYTECOUNT_CLIENT_H_

#include <clio_cte/core/core_client.h>
#include <clio_demo/bytecount/bytecount_tasks.h>

namespace clio::demo::bytecount {

/** The operator speaks the CTE core's interface, so the core client works. */
class Client : public clio::cte::core::Client {
 public:
  using clio::cte::core::Client::Client;
};

}  // namespace clio::demo::bytecount

#endif  // CLIO_DEMO_BYTECOUNT_CLIENT_H_
c++ · bytecount_client.cc
#include <clio_demo/bytecount/bytecount_client.h>

5. Write the operator

The Runtime class derives from CoreInterposer, which supplies ForwardToCore: it sends a task to the next pool and resumes when the task completes there.

c++ · bytecount_runtime.h
#ifndef CLIO_DEMO_BYTECOUNT_RUNTIME_H_
#define CLIO_DEMO_BYTECOUNT_RUNTIME_H_

#include <mutex>
#include <string>
#include <unordered_map>

#include <clio_cte/core/core_interposer.h>
#include <clio_demo/bytecount/bytecount_client.h>
#include <clio_demo/bytecount/bytecount_tasks.h>

namespace clio::demo::bytecount {

/**
 * A data operator that sits in the CTE chain, forwards every call to
 * next_pool_id unchanged, and counts the bytes written to each tag.
 */
class Runtime : public clio::cte::core::CoreInterposer {
 public:
  using CreateParams = BytecountConfig;  // required by CLIO_TASK_CC

  // Handlers for the methods listed in clio_mod.yaml.
  clio::run::TaskResume Create(clio::run::shared_ptr<CreateTask> &task);
  clio::run::TaskResume Destroy(clio::run::shared_ptr<DestroyTask> &task);
  clio::run::TaskResume Monitor(clio::run::shared_ptr<MonitorTask> &task);
  clio::run::TaskResume PutBlob(
      clio::run::shared_ptr<clio::cte::core::PutBlobTask> &task);
  clio::run::TaskResume MultiPutBlob(
      clio::run::shared_ptr<clio::cte::core::MultiPutBlobTask> &task);

  // Container plumbing, implemented in bytecount_exec.cc.
  void Init(const clio::run::PoolId &pool_id, const std::string &pool_name,
            clio::run::u32 container_id = 0) override;
  clio::run::TaskResume Run(clio::run::u32 method,
                            clio::run::shared_ptr<clio::run::Task> task) override;
  clio::run::u64 GetWorkRemaining() const override { return 0; }
  void SaveTask(clio::run::u32 method, clio::run::SaveTaskArchive &ar,
                clio::run::shared_ptr<clio::run::Task> &task) override;
  void LoadTask(clio::run::u32 method, clio::run::LoadTaskArchive &ar,
                clio::run::shared_ptr<clio::run::Task> &task) override;
  clio::run::shared_ptr<clio::run::Task> AllocLoadTask(
      clio::run::u32 method, clio::run::LoadTaskArchive &ar) override;
  void LocalSaveTask(clio::run::u32 method, clio::run::DefaultSaveArchive &ar,
                     clio::run::shared_ptr<clio::run::Task> &task) override;
  void LocalLoadTask(clio::run::u32 method, clio::run::DefaultLoadArchive &ar,
                     clio::run::shared_ptr<clio::run::Task> &task) override;
  clio::run::shared_ptr<clio::run::Task> LocalAllocLoadTask(
      clio::run::u32 method, clio::run::DefaultLoadArchive &ar) override;
  clio::run::shared_ptr<clio::run::Task> NewCopyTask(
      clio::run::u32 method, clio::run::shared_ptr<clio::run::Task> &orig,
      bool deep) override;
  clio::run::shared_ptr<clio::run::Task> NewTask(clio::run::u32 method) override;
  void AggregateOut(clio::run::u32 method,
                    clio::run::shared_ptr<clio::run::Task> &orig,
                    const clio::run::shared_ptr<clio::run::Task> &replica) override;
  void AggregateIn(clio::run::u32 method,
                   clio::run::shared_ptr<clio::run::Task> &agg,
                   const clio::run::shared_ptr<clio::run::Task> &member) override;

 private:
  /** Add n bytes to a tag's total. */
  void Count(const clio::run::UniqueId &tag, clio::run::u64 n);

  BytecountConfig config_;
  std::mutex mu_;  ///< Guards bytes_; never held across a co_await
  std::unordered_map<std::string, clio::run::u64> bytes_;  ///< "major.minor" -> bytes
};

}  // namespace clio::demo::bytecount

#endif  // CLIO_DEMO_BYTECOUNT_RUNTIME_H_

Each handler forwards first and counts afterwards, so it counts only writes that succeeded. context_.replica_ is 0 for the original write, which keeps copies made further down the chain from being counted twice. Monitor answers with a msgpack map, which the dashboard API turns into JSON.

c++ · bytecount_runtime.cc
#include <clio_demo/bytecount/bytecount_runtime.h>
#include <clio_cte/core/blob_batch.h>
#include <clio_ctp/serialize/msgpack_wrapper.h>

namespace clio::demo::bytecount {

clio::run::TaskResume Runtime::Create(clio::run::shared_ptr<CreateTask> &task) {
  CLIO_TASK_BODY_BEGIN
  config_ = task->GetParams();
  interposer_next_pool_ = config_.next_pool_id_;  // where ForwardToCore sends
  task->return_code_ = 0;
  CLIO_CO_RETURN;
  CLIO_TASK_BODY_END
}

clio::run::TaskResume Runtime::Destroy(clio::run::shared_ptr<DestroyTask> &task) {
  CLIO_TASK_BODY_BEGIN
  task->return_code_ = 0;
  CLIO_CO_RETURN;
  CLIO_TASK_BODY_END
}

void Runtime::Count(const clio::run::UniqueId &tag, clio::run::u64 n) {
  std::lock_guard<std::mutex> lock(mu_);
  bytes_[std::to_string(tag.major_) + "." + std::to_string(tag.minor_)] += n;
}

clio::run::TaskResume Runtime::PutBlob(
    clio::run::shared_ptr<clio::cte::core::PutBlobTask> &task) {
  CLIO_TASK_BODY_BEGIN
  // Pass the write down the chain first, then count it if it succeeded.
  CLIO_CO_AWAIT(ForwardToCore(clio::cte::core::Method::kPutBlob,
                              task.template Cast<clio::run::Task>()));
  if (task->GetReturnCode() == 0 && task->context_.replica_ == 0) {
    clio::run::u64 n = 0;
    clio::cte::core::ForEachBlobRegion(
        *task, [&n](const clio::cte::core::BlobRegion &r) {
          n += r.size_;
          return true;
        });
    Count(task->tag_id_, n);
  }
  CLIO_CO_RETURN;
  CLIO_TASK_BODY_END
}

clio::run::TaskResume Runtime::MultiPutBlob(
    clio::run::shared_ptr<clio::cte::core::MultiPutBlobTask> &task) {
  CLIO_TASK_BODY_BEGIN
  CLIO_CO_AWAIT(ForwardToCore(clio::cte::core::Method::kMultiPutBlob,
                              task.template Cast<clio::run::Task>()));
  if (task->GetReturnCode() == 0 && task->context_.replica_ == 0) {
    for (const auto &d : clio::cte::core::DecodeMultiPutDescs(task->descs_)) {
      Count(d.tag_id_, d.size_);
    }
  }
  CLIO_CO_RETURN;
  CLIO_TASK_BODY_END
}

clio::run::TaskResume Runtime::Monitor(clio::run::shared_ptr<MonitorTask> &task) {
  CLIO_TASK_BODY_BEGIN
  // Reply with a msgpack map {"major.minor": bytes}. The dashboard's
  // /api/pools/<id>/monitor route decodes it to JSON.
  msgpack::sbuffer buf;
  msgpack::packer<msgpack::sbuffer> pk(buf);
  {
    std::lock_guard<std::mutex> lock(mu_);
    pk.pack_map(bytes_.size());
    for (const auto &[tag, n] : bytes_) {
      pk.pack(tag);
      pk.pack(n);
    }
  }
  task->results_[container_id_] = std::string(buf.data(), buf.size());
  task->SetReturnCode(0);
  CLIO_CO_RETURN;
  CLIO_TASK_BODY_END
}

}  // namespace clio::demo::bytecount

CLIO_TASK_CC(clio::demo::bytecount::Runtime)
Handlers are coroutines

Handlers run between CLIO_TASK_BODY_BEGIN and CLIO_TASK_BODY_END and wait with CLIO_CO_AWAIT. A worker thread runs other tasks while one waits, so never hold a lock across a CLIO_CO_AWAIT.

6. Add the dispatch

The dispatch file routes each method id to a handler and serializes tasks. The list at the top names the methods the operator handles. Every other id takes the default branch, which hands it to the next pool unchanged. This file is copied from the cache operator in context-transfer-engine/cache/src/autogen/cache_lib_exec.cc, with the method list changed.

src/bytecount_exec.cc
c++ · bytecount_exec.cc
// Container plumbing for the bytecount operator. Methods it handles go to
// its own handlers; every other CTE core method falls through to the
// default branches, which forward it to the next pool unchanged.
#include <clio_demo/bytecount/bytecount_runtime.h>
#include <clio_runtime/clio_runtime.h>
#include <clio_runtime/task.h>

namespace clio::demo::bytecount {

#define BYTECOUNT_FOR_EACH_METHOD(X)              \
  X(kCreate, CreateTask, Create)                   \
  X(kDestroy, DestroyTask, Destroy)                \
  X(kMonitor, MonitorTask, Monitor)                \
  X(kPutBlob, clio::cte::core::PutBlobTask, PutBlob) \
  X(kMultiPutBlob, clio::cte::core::MultiPutBlobTask, MultiPutBlob)

void Runtime::Init(const clio::run::PoolId &pool_id, const std::string &pool_name,
                   clio::run::u32 container_id) {
  clio::run::Container::Init(pool_id, pool_name, container_id);
  DefineModel(Method::kMaxMethodId);
  SetMethodNames(Method::GetMethodNames());
}

clio::run::TaskResume Runtime::Run(clio::run::u32 method,
                             clio::run::shared_ptr<clio::run::Task> task_ptr) {
  CLIO_TASK_BODY_BEGIN
  switch (method) {
#define X(MID, TASK, HANDLER)                                       \
    case Method::MID: {                                             \
      auto &typed = task_ptr.template Cast<TASK>();                 \
      CLIO_CO_AWAIT(HANDLER(typed));                                \
      break;                                                        \
    }
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      CLIO_CO_AWAIT(ForwardToCore(method, task_ptr));
      break;
  }
  CLIO_CO_RETURN;
  CLIO_TASK_BODY_END
}

void Runtime::SaveTask(clio::run::u32 method, clio::run::SaveTaskArchive &archive,
                       clio::run::shared_ptr<clio::run::Task> &task_ptr) {
  switch (method) {
#define X(MID, TASK, HANDLER)                                  \
    case Method::MID:                                          \
      archive << *task_ptr.template Cast<TASK>();              \
      break;
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      ForwardSaveTask(method, archive, task_ptr);
      break;
  }
}

void Runtime::LoadTask(clio::run::u32 method, clio::run::LoadTaskArchive &archive,
                       clio::run::shared_ptr<clio::run::Task> &task_ptr) {
  switch (method) {
#define X(MID, TASK, HANDLER)                                  \
    case Method::MID:                                          \
      archive >> *task_ptr.template Cast<TASK>();              \
      break;
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      ForwardLoadTask(method, archive, task_ptr);
      break;
  }
}

clio::run::shared_ptr<clio::run::Task> Runtime::AllocLoadTask(
    clio::run::u32 method, clio::run::LoadTaskArchive &archive) {
  clio::run::shared_ptr<clio::run::Task> task_ptr = NewTask(method);
  if (!task_ptr.IsNull()) {
    LoadTask(method, archive, task_ptr);
  }
  return task_ptr;
}

void Runtime::LocalLoadTask(clio::run::u32 method, clio::run::DefaultLoadArchive &archive,
                            clio::run::shared_ptr<clio::run::Task> &task_ptr) {
  switch (method) {
#define X(MID, TASK, HANDLER)                                  \
    case Method::MID:                                          \
      archive >> *task_ptr.template Cast<TASK>();              \
      break;
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      ForwardLocalLoadTask(method, archive, task_ptr);
      break;
  }
}

clio::run::shared_ptr<clio::run::Task> Runtime::LocalAllocLoadTask(
    clio::run::u32 method, clio::run::DefaultLoadArchive &archive) {
  clio::run::shared_ptr<clio::run::Task> task_ptr = NewTask(method);
  if (!task_ptr.IsNull()) {
    LocalLoadTask(method, archive, task_ptr);
  }
  return task_ptr;
}

void Runtime::LocalSaveTask(clio::run::u32 method, clio::run::DefaultSaveArchive &archive,
                            clio::run::shared_ptr<clio::run::Task> &task_ptr) {
  switch (method) {
#define X(MID, TASK, HANDLER)                                  \
    case Method::MID:                                          \
      archive << *task_ptr.template Cast<TASK>();              \
      break;
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      ForwardLocalSaveTask(method, archive, task_ptr);
      break;
  }
}

clio::run::shared_ptr<clio::run::Task> Runtime::NewCopyTask(
    clio::run::u32 method, clio::run::shared_ptr<clio::run::Task> &orig_task_ptr, bool deep) {
  auto *ipc_manager = CLIO_IPC;
  if (!ipc_manager) {
    return clio::run::shared_ptr<clio::run::Task>();
  }
  switch (method) {
#define X(MID, TASK, HANDLER)                                            \
    case Method::MID: {                                                  \
      auto new_task = ipc_manager->NewTask<TASK>();                      \
      if (!new_task.IsNull()) {                                          \
        new_task->Copy(ctp::ipc::FullPtr<TASK>(                          \
            orig_task_ptr.template Cast<TASK>().get()));                 \
        return new_task.template Cast<clio::run::Task>();                \
      }                                                                  \
      break;                                                             \
    }
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      return ForwardNewCopyTask(method, orig_task_ptr, deep);
  }
  return clio::run::shared_ptr<clio::run::Task>();
}

clio::run::shared_ptr<clio::run::Task> Runtime::NewTask(clio::run::u32 method) {
  auto *ipc_manager = CLIO_IPC;
  if (!ipc_manager) {
    return clio::run::shared_ptr<clio::run::Task>();
  }
  switch (method) {
#define X(MID, TASK, HANDLER)                                  \
    case Method::MID:                                          \
      return ipc_manager->NewTask<TASK>().template Cast<clio::run::Task>();
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      return ForwardNewTask(method);
  }
}

void Runtime::AggregateOut(clio::run::u32 method, clio::run::shared_ptr<clio::run::Task> &orig_task,
                           const clio::run::shared_ptr<clio::run::Task> &replica_task) {
  switch (method) {
#define X(MID, TASK, HANDLER)                                            \
    case Method::MID:                                                    \
      orig_task.template Cast<TASK>()->AggregateOut(                     \
          ctp::ipc::FullPtr<clio::run::Task>(replica_task.get()));       \
      break;
    BYTECOUNT_FOR_EACH_METHOD(X)
#undef X
    default:
      ForwardAggregateOut(method, orig_task, replica_task);
      break;
  }
}

void Runtime::AggregateIn(clio::run::u32 method, clio::run::shared_ptr<clio::run::Task> &agg_task,
                          const clio::run::shared_ptr<clio::run::Task> &member_task) {
  ForwardAggregateIn(method, agg_task, member_task);
}

#undef BYTECOUNT_FOR_EACH_METHOD

}  // namespace clio::demo::bytecount

7. Build it

shell
$ cmake -S bytecount -B bytecount/build -DCMAKE_PREFIX_PATH=$HOME/clio-inst
$ cmake --build bytecount/build
[ 40%] Built target clio_demo_bytecount_client
[100%] Built target clio_demo_bytecount_runtime

If CMake reports that the target yaml-cpp::yaml-cpp was not found, keep the find_package(yaml-cpp REQUIRED) line above find_package(clio-core).

8. Put it in the chain

Copy your configuration and add the operator right after the clio_cte_cache entry, since it forwards to the cache:

yaml · clio.yaml
  # The bytecount operator: counts bytes written, then forwards to the cache.
  - mod_name: clio_demo_bytecount
    pool_name: bytecount
    pool_query: local
    pool_id: "600.0"
    next_pool_id: "563.0"

Then send the filesystem through it: in the clio_cte_filesystem entry, change next_pool_id from "563.0" to "600.0".

9. Run it

shell
$ export CLIO_REPO_PATH=$HOME/bytecount/build     # where the runtime finds the module
$ export CLIO_SERVER_CONF=$HOME/clio-bytecount.yaml
$ clio_run start &
$ CLIO_WITH_RUNTIME=0 clio_cte_fuse ~/clio-mnt -f &

The runtime log shows Loaded ChiMod: clio_demo_bytecount, and the FUSE daemon prints data pool 600.0 (filesystem chain). Write some files:

shell
$ mkdir -p ~/clio-mnt/runs
$ head -c 3M /dev/urandom > ~/clio-mnt/runs/a.bin
$ head -c 5M /dev/urandom > ~/clio-mnt/runs/b.bin
$ curl -s "http://127.0.0.1:8080/api/pools/600.0/monitor?query=bytes&routing=local"
{"pool_id":"600.0","query":"bytes","results":{"0":{"1610612736.3":5243084,"1610612736.2":3146000, ...}}}

The keys are tag ids. Each file's count is its data plus its small attribute blob. Turn the ids into paths from Python:

python · report.py
import json, os, urllib.request
os.environ["CLIO_WITH_RUNTIME"] = "0"
import clio_cte_core_ext as cte

cte.clio_init(cte.RuntimeMode.kClient, False)
cte.initialize_cte("", cte.PoolQuery.Dynamic())
client = cte.get_cte_client()

# Map "major.minor" tag ids back to paths in the mount.
names = {}
for path in client.TagQuery(".*", 0):
    t = cte.Tag(path).GetTagId()
    names[f"{t.major_}.{t.minor_}"] = path

url = "http://127.0.0.1:8080/api/pools/600.0/monitor?query=bytes&routing=local"
counts = json.load(urllib.request.urlopen(url))["results"]["0"]
for tag, n in sorted(counts.items(), key=lambda kv: -kv[1]):
    print(f"{n:>10}  {names.get(tag, tag)}")
output
   5243084  /runs/b.bin
   3146000  /runs/a.bin
     16384  /runs

Where to go from here

  • Change data on its way down. Read or rewrite the put's payload before you forward it. ForEachBlobRegion visits each region of a put, and MultiPutBatchView resolves a batch's payload; both are in clio_cte/core/blob_batch.h.
  • Intercept reads. Add kGetBlob: 16 to clio_mod.yaml, a handler, and a line in the dispatch list. The cache operator shows a complete read path.
  • Run on several nodes. Calls are routed to the node that owns each blob, so every node keeps its own counts. Query with routing=broadcast to collect them all.
  • Learn from the built-in operators. In context-transfer-engine/, the indexer observes writes like this one does, cache serves reads, and replication writes extra copies.