Sariel00's picture
Publish Ling-3.0-tiny RKNN engine and model
3fd1a35 verified
Raw History Blame Contribute Delete
4.95 kB
#pragma once
#include "ling3/chat_protocol.h"
#include <chrono>
#include <condition_variable>
#include <functional>
#include <mutex>
namespace ling3::chat {
// The generation thread parks only AFTER a completed evaluation, BEFORE the
// next one. An acknowledgement belongs to both a server request ID and a
// monotonically increasing sequence; a late control cannot resume another turn.
class FlowGate {
mutable std::mutex mutex_;
std::condition_variable changed_;
std::string id_;
bool active_ = false, enabled_ = false, canceled_ = false, waiting_ = false;
bool pause_requested_ = false, paused_ = false;
bool committing_ = false;
std::uint64_t sequence_ = 0, acknowledged_ = 0;
public:
void Begin(std::string id, bool enabled) {
std::lock_guard lock(mutex_);
id_ = std::move(id); enabled_ = enabled; active_ = true;
canceled_ = waiting_ = false; sequence_ = acknowledged_ = 0;
pause_requested_ = paused_ = false;
committing_ = false;
}
void End() {
std::lock_guard lock(mutex_);
active_ = waiting_ = paused_ = pause_requested_ = committing_ = false; changed_.notify_all();
}
bool Active() const { std::lock_guard lock(mutex_); return active_; }
bool Alive() const { std::lock_guard lock(mutex_); return active_ && !canceled_; }
bool Cancel(const std::string & id = {}) {
std::lock_guard lock(mutex_);
if (!active_ || committing_ || (!id.empty() && id != id_)) return false;
canceled_ = true; changed_.notify_all(); return true;
}
// Linearization point between cancellation and publishing completed state.
// Keep active=true until cleanup, but do not accept cancellation after this.
bool TryCommit() {
std::lock_guard lock(mutex_);
if (!active_ || canceled_ || committing_) return false;
committing_ = true; waiting_ = pause_requested_ = paused_ = false;
changed_.notify_all(); return true;
}
std::uint64_t Prepare() {
std::lock_guard lock(mutex_);
if (!enabled_) return 0;
waiting_ = true; changed_.notify_all(); return ++sequence_;
}
void Pause(const std::string & id) {
std::unique_lock lock(mutex_);
if (!active_ || canceled_ || committing_ || id.empty() || id != id_)
throw Error(409, "stale_flow_control", "request ID is not active");
pause_requested_ = true;
if (!changed_.wait_for(lock, std::chrono::seconds(30), [&] {
return paused_ || waiting_ || !active_ || canceled_ || committing_ || !pause_requested_ || id != id_;
})) {
pause_requested_ = false; changed_.notify_all();
throw Error(409, "pause_timeout", "generation has not reached a safe evaluation boundary");
}
if (!active_ || canceled_ || committing_ || !pause_requested_ || id != id_)
throw Error(409, "pause_interrupted", "generation ended or pause was revoked");
}
void Resume(const std::string & id) {
std::lock_guard lock(mutex_);
if (!active_ || canceled_ || committing_ || id.empty() || id != id_)
throw Error(409, "stale_flow_control", "request ID is not active");
pause_requested_ = false; changed_.notify_all();
}
void Ack(const std::string & id, std::uint64_t sequence) {
std::lock_guard lock(mutex_);
if (!active_ || canceled_ || committing_ || !enabled_ || id != id_ ||
!sequence || sequence != sequence_)
throw Error(409, "stale_flow_control", "request ID or flow sequence is not active");
acknowledged_ = sequence; waiting_ = false; changed_.notify_all();
}
bool Wait(const std::function<bool()> & connected) {
const auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(300);
std::unique_lock lock(mutex_);
while (active_ && !canceled_ && (pause_requested_ || (enabled_ && acknowledged_ < sequence_))) {
paused_ = pause_requested_; changed_.notify_all();
lock.unlock();
const bool alive = connected();
lock.lock();
if (!alive || std::chrono::steady_clock::now() >= deadline) {
canceled_ = true; break;
}
changed_.wait_for(lock, std::chrono::milliseconds(20));
}
waiting_ = paused_ = false; changed_.notify_all();
return active_ && !canceled_;
}
Json Status() const {
std::lock_guard lock(mutex_);
return {{"request_id", active_ ? id_ : ""}, {"active", active_},
{"flow_control", enabled_ ? "ack" : "none"},
{"waiting_for_ack", waiting_}, {"sequence", sequence_},
{"pause_requested", pause_requested_}, {"paused", paused_ || (pause_requested_ && waiting_)},
{"acknowledged", acknowledged_}, {"canceled", canceled_}, {"committing", committing_}};
}
};
}