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
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,16 @@ Notable changes to vla.cpp. Format loosely follows [Keep a Changelog](https://ke

## [Unreleased]

### Added

- **Asynchronous serving.** `vla-server` is now a ZeroMQ ROUTER with prediction
on its own thread, so it keeps receiving while the model runs and a DEALER
client can keep requests in flight; REQ clients are unchanged. `--queue latest`
(default) keeps one pending request per client and answers a superseded one
with `error="superseded"`; `--queue fifo` serves every request in order.
Replies report `latency_ms_queue`. The lerobot fork's `lerobot-vla-cpp
--mode=async` drives it with a timestep-aligned action queue.

## [0.4.0] - 2026-09-30

### Added
Expand Down
10 changes: 7 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -207,9 +207,13 @@ lerobot-vla-cpp --server_address=tcp://127.0.0.1:5555 --arch=smolvla "${ROBOT[@]
relative actions.
- `--task` must match a trained instruction exactly, and the camera keys must
stay `front` and `wrist` in that order.
- `vla-server` answers one request at a time, so the loop is synchronous and
`--n_action_steps` is the feedback rate: 25 at 30 fps leaves ~0.83 s between
observations.
- `--mode=async` runs inference on a background thread and merges each new chunk
into a timestep-aligned action queue, so the arm never waits for the server.
`--chunk_size_threshold` sets how early the next observation goes out, and
`--aggregate_fn_name` how overlapping actions from two chunks are blended. The
default `--mode=sync` executes `--n_action_steps` of each chunk before the
next request, so that number is the feedback rate: 25 at 30 fps leaves ~0.83 s
between observations.

The GR00T paths of the client have not been tested against real checkpoints yet.
Wiring, recording, training and queue sizing are in the
Expand Down
3 changes: 2 additions & 1 deletion docs/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@ source is the detail.
- `src/tokenizer.h` - SentencePiece tokenizers stored in the GGUF, for
`vla-cli --text`.
- `include/vla.h`, `src/vla_c_api.cpp` - the C ABI (`libvla`).
- `src/serving/` - `vla-server` (ZeroMQ + protobuf, action prediction), `vlm-server`
- `src/serving/` - `vla-server` (ZeroMQ ROUTER + protobuf, action prediction on a
worker thread behind a latest-wins request queue), `vlm-server`
(chat), `vla-cli` (one-shot inference) and `vla-bench` (timing).
- `src/kernels/bitvla/` - custom 1.58-bit ternary CUDA kernels for BitVLA.

Expand Down
28 changes: 26 additions & 2 deletions docs/USAGE.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,8 @@ Checkpoints are cached under `$VLA_CACHE` (default `~/.cache/vla`).

## `vla-server`

`vla-server` loads the model once at startup and answers ZeroMQ REQ/REP requests
synchronously.
`vla-server` loads the model once at startup and answers ZeroMQ requests carrying
the protobuf messages in `src/serving/vla.proto`.

```bash
./build/vla-server "$VLA_GGUF"
Expand All @@ -60,6 +60,30 @@ vla-server: bound to tcp://*:5555. ready.
Use `--bind` to change the address and port. Stop the server with `Ctrl-C`.
`vla-server` also takes `-hf user/repo[:file.gguf|:tag]` in place of a checkpoint path.

### Asynchronous clients

The socket is a ZeroMQ ROUTER, and prediction runs on its own thread. A REQ
client sees what it always did: one request, one reply. A DEALER client can
keep several requests in flight, and the server keeps receiving while the model
runs, so a robot's control loop never has to wait for a prediction to send the
next observation. Replies carry the request's `request_id`, which is how a
client matches a chunk to the observation it came from.

`--queue` says what happens to requests that arrive while a prediction is
running:

- `latest` (default): one pending request per client. A newer request from the
same client replaces the pending one, which is answered at once with
`error="superseded"`. The model therefore always works on the freshest
observation each robot has sent, and several robots sharing one server are
served in turn.
- `fifo`: every request is served in arrival order, for benchmarks that want
throughput rather than freshness.

`latency_ms_queue` in the reply is how long the request waited for the predict
thread. A malformed request is rejected on the socket thread, so an error comes
back within milliseconds even mid-prediction.

Clients: the LIBERO and SimplerEnv runners in [EVAL.md](EVAL.md), and the
real-robot client in the README's
[Rollout on a real robot](../README.md#rollout-on-a-real-robot).
Expand Down
137 changes: 137 additions & 0 deletions src/serving/request_queue.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
// Copyright 2026 VinRobotics
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

// The hand-off between vla-server's socket thread and its predict thread.
//
// The socket thread parses and validates a request and pushes a Job; the
// predict thread pops one at a time. A policy server only ever wants the newest
// observation from each robot, so in LATEST mode a push from a client that
// already has a job waiting replaces that job and hands it back to the caller,
// who answers it with a "superseded" error. FIFO keeps every request, bounded
// by a depth cap, for benchmarks that want throughput rather than freshness.
//
// Header-only and free of ZeroMQ, protobuf and the model so tests can drive it.

#pragma once

#include <condition_variable>
#include <cstddef>
#include <deque>
#include <mutex>
#include <optional>
#include <string>
#include <utility>

namespace vla::serving {

enum class QueueMode {
LATEST, ///< One pending job per client; a newer one displaces it.
FIFO, ///< Every job is served in arrival order, up to max_depth.
};

inline bool parse_queue_mode(const std::string & s, QueueMode & out) {
if (s == "latest") { out = QueueMode::LATEST; return true; }
if (s == "fifo") { out = QueueMode::FIFO; return true; }
return false;
}

inline const char * queue_mode_name(QueueMode m) {
return m == QueueMode::LATEST ? "latest" : "fifo";
}

/// What push() did with a job.
enum class PushResult {
QUEUED, ///< Appended; nothing displaced.
REPLACED, ///< Appended, and the same client's older pending job is in `displaced`.
REJECTED, ///< Queue full (FIFO only); the job itself is handed back in `displaced`.
};

template <class Job>
class RequestQueue {
public:
explicit RequestQueue(QueueMode mode, size_t max_depth = 64)
: mode_(mode), max_depth_(max_depth == 0 ? 1 : max_depth) {}

QueueMode mode() const { return mode_; }

/// Jobs are keyed by `key(job)`, the ZeroMQ routing envelope in the server.
template <class KeyFn>
PushResult push(Job job, KeyFn key, std::optional<Job> & displaced) {
std::lock_guard<std::mutex> lk(mu_);
displaced.reset();
if (mode_ == QueueMode::LATEST) {
const auto k = key(job);
for (auto it = q_.begin(); it != q_.end(); ++it) {
if (key(*it) == k) {
displaced = std::move(*it);
q_.erase(it);
q_.push_back(std::move(job));
cv_.notify_one();
return PushResult::REPLACED;
}
}
// No pending job from this client. A bounded queue still protects
// against a flood of distinct identities.
if (q_.size() >= max_depth_) {
displaced = std::move(job);
return PushResult::REJECTED;
}
q_.push_back(std::move(job));
cv_.notify_one();
return PushResult::QUEUED;
}
if (q_.size() >= max_depth_) {
displaced = std::move(job);
return PushResult::REJECTED;
}
q_.push_back(std::move(job));
cv_.notify_one();
return PushResult::QUEUED;
}

/// Blocks until a job is available or stop() was called. Returns false on stop
/// with the queue drained.
bool pop(Job & out) {
std::unique_lock<std::mutex> lk(mu_);
cv_.wait(lk, [&] { return stopped_ || !q_.empty(); });
if (q_.empty())
return false;
out = std::move(q_.front());
q_.pop_front();
return true;
}

/// Wakes pop(). Jobs still queued are returned by subsequent pops until
/// empty, so the worker can answer them before it exits.
void stop() {
std::lock_guard<std::mutex> lk(mu_);
stopped_ = true;
cv_.notify_all();
}

size_t size() const {
std::lock_guard<std::mutex> lk(mu_);
return q_.size();
}

private:
QueueMode mode_;
size_t max_depth_;
mutable std::mutex mu_;
std::condition_variable cv_;
std::deque<Job> q_;
bool stopped_ = false;
};

} // namespace vla::serving
Loading
Loading