Skip to main content

cherenkov/
server.rs

1//! Server startup and shared job ownership. One worker owns the loaded GPU.
2
3use crate::{
4    config::Source,
5    control::{self, State, state::Versioned},
6    prefix_cache::PrefixCache,
7    qwen4_exp,
8    tok::ChatTokenizer,
9    units::BYTES_PER_GB,
10};
11use anyhow::{Context, Result, ensure};
12use serde_json::Value;
13use std::net::{TcpListener, TcpStream};
14use std::sync::{
15    Arc,
16    atomic::{AtomicUsize, Ordering},
17    mpsc,
18};
19use std::time::Duration;
20
21mod failure;
22mod http;
23mod output;
24mod pacer;
25mod registry;
26mod request;
27mod response;
28mod routes;
29mod sessions;
30mod stats;
31mod tool_call;
32mod worker;
33
34pub use stats::UsageStats;
35
36use http::error;
37use routes::connection;
38
39#[derive(Clone, Copy, PartialEq, Eq)]
40enum ApiKind {
41    Chat,
42    Completion,
43}
44
45const MODEL: &str = "cherenkov";
46
47struct Job {
48    stream: TcpStream,
49    body: Value,
50    kind: ApiKind,
51    settings: Versioned,
52    ticket: Arc<registry::Ticket>,
53    session: Option<sessions::Turn>,
54}
55
56pub fn serve(source: Source) -> Result<()> {
57    let config = source.resolve()?;
58    let mut options = config.options();
59    options.repack = source.overrides.repack;
60
61    ensure!(
62        !options.repack || options.build_missing_store,
63        "--repack conflicts with build_missing_store=false"
64    );
65    options.validate()?;
66
67    // Retain the lease until the worker and all GPU resources have been dropped.
68    let model =
69        crate::model::index::resolve_runtime(config.paths()?, &config.model_dir()?, &mut options)?;
70    let model_dir = model.path.clone();
71
72    let cache_bytes = config.cache_bytes();
73    let cache = PrefixCache::new(
74        cache_bytes,
75        config.limits.cache_max_entries,
76        config.limits.cache_idle_seconds,
77    );
78    let state = Arc::new(State::new(source, config.clone()));
79    let requests = Arc::new(registry::Registry::default());
80    let sessions = Arc::new(std::sync::Mutex::new(sessions::Store::new(&config)));
81    let _control = control::Listener::start(&config.server.socket, state.clone())?;
82
83    eprintln!("control socket: {}", config.server.socket.display());
84
85    let listener =
86        TcpListener::bind(("127.0.0.1", config.server.port)).context("binding server")?;
87    let address = listener.local_addr()?;
88    let tok = ChatTokenizer::load(&model_dir)?;
89    let packed = qwen4_exp::packed::Packed::open(&model_dir)?;
90    let gpu = qwen4_exp::gpu::Gpu::load_bounded(
91        &packed,
92        options.max_ctx,
93        &options,
94        config.reserved_bytes(),
95        Some(config.memory_bytes()),
96        config.limits.prefill_quantum,
97    )?;
98
99    ensure!(
100        gpu.allocated_gb() + config.reserved_bytes() as f64 / BYTES_PER_GB as f64
101            <= config.limits.memory_gb,
102        "Metal plus cache/session reservations exceeds configured memory budget"
103    );
104    state.observe(&gpu, &cache, &mut Default::default());
105    state.update(|s| {
106        s.ready = true;
107        s.http_address = Some(address.to_string());
108    });
109
110    if options.cut_weak > 0.0 {
111        eprintln!("WARNING: --cut-weak makes output depend on disk timing and non-reproducible.");
112    }
113
114    eprintln!(
115        "cherenkov serving http://{address}/v1, model {MODEL}, {:.2} GB Metal",
116        gpu.allocated_gb()
117    );
118
119    let (tx, rx) = mpsc::sync_channel::<Job>(config.limits.queued_requests);
120    let connections = state.clone();
121    let network_sessions = sessions.clone();
122
123    std::thread::spawn(move || {
124        let active = Arc::new(AtomicUsize::new(0));
125
126        for stream in listener.incoming() {
127            let Ok(mut stream) = stream else { continue };
128
129            if active.fetch_add(1, Ordering::AcqRel) >= config.limits.http_readers {
130                active.fetch_sub(1, Ordering::AcqRel);
131                stream.set_write_timeout(Some(Duration::from_secs(1))).ok();
132                connections.update(|s| s.rejected_requests += 1);
133
134                let _ = error(&mut stream, 503, "too many connections");
135
136                continue;
137            }
138
139            let tx = tx.clone();
140            let active = active.clone();
141            let state = connections.clone();
142            let requests = requests.clone();
143            let sessions = network_sessions.clone();
144
145            std::thread::spawn(move || {
146                if let Err(e) = connection(stream, &tx, &state, &requests, &sessions) {
147                    eprintln!("HTTP connection: {e:#}");
148                }
149
150                active.fetch_sub(1, Ordering::AcqRel);
151            });
152        }
153    });
154    worker::Worker::new(gpu, &tok, options, cache, state, sessions).run(rx);
155
156    Ok(())
157}
158
159#[cfg(test)]
160#[path = "../tests/unit/server.rs"]
161mod tests;