#pragma once #include "ling3/chat_protocol.h" #include #include #include #include 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 & 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_}}; } }; }