File size: 4,949 Bytes
3fd1a35
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
#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_}};
    }
};
}