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
1 change: 1 addition & 0 deletions server/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -554,6 +554,7 @@ add_library(dflash_common STATIC
src/server/tool_hint.cpp
src/server/reasoning.cpp
src/server/tool_memory.cpp
src/server/response_error.cpp
src/server/sse_emitter.cpp
src/server/prefix_cache.cpp
src/server/pin_friendly_prompt.cpp
Expand Down
2 changes: 2 additions & 0 deletions server/src/common/model_backend.h
Original file line number Diff line number Diff line change
Expand Up @@ -207,6 +207,7 @@ struct GenerateRequest {
enum class GenerateErrorCode {
Incomplete,
AdapterUnavailable,
ResourceExhausted,
ContextOverflow,
SamplingUnsupported,
PrefillFailed,
Expand All @@ -221,6 +222,7 @@ constexpr std::string_view generate_error_code(GenerateErrorCode error) {
switch (error) {
case GenerateErrorCode::Incomplete: return "incomplete";
case GenerateErrorCode::AdapterUnavailable: return "adapter_unavailable";
case GenerateErrorCode::ResourceExhausted: return "resource_exhausted";
case GenerateErrorCode::ContextOverflow: return "context_overflow";
case GenerateErrorCode::SamplingUnsupported: return "sampling_unsupported";
case GenerateErrorCode::PrefillFailed: return "prefill_failed";
Expand Down
2 changes: 1 addition & 1 deletion server/src/gemma4/gemma4_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -900,7 +900,7 @@ GenerateResult Gemma4Backend::restore_and_generate_impl(int slot,
if (snap_pos > kvflash_tokens_ - kvflash_pager_.chunk_tokens()) {
std::fprintf(stderr, "[kvflash] restored prefix (%d) exceeds pool %d\n",
snap_pos, kvflash_tokens_);
result.fail(GenerateErrorCode::ContextOverflow);
result.fail(GenerateErrorCode::ResourceExhausted);
return result;
}
kvflash_pager_.reset();
Expand Down
4 changes: 2 additions & 2 deletions server/src/laguna/laguna_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1584,7 +1584,7 @@ GenerateResult LagunaBackend::restore_and_generate_impl(int slot,
N > kvflash_tokens_ - kvflash_pager_.chunk_tokens()) {
std::fprintf(stderr, "[kvflash] restore prompt (%d) exceeds pool %d; "
"raise --kvflash\n", N, kvflash_tokens_);
result.fail(GenerateErrorCode::ContextOverflow);
result.fail(GenerateErrorCode::ResourceExhausted);
return result;
}
if (kvflash_active()) {
Expand Down Expand Up @@ -2743,7 +2743,7 @@ GenerateResult LagunaBackend::generate_hybrid(const GenerateRequest & req,
N > kvflash_tokens_ - kvflash_pager_.chunk_tokens()) {
std::fprintf(stderr, "[kvflash] hybrid prompt (%d) exceeds pool %d; "
"raise --kvflash\n", N, kvflash_tokens_);
result.fail(GenerateErrorCode::ContextOverflow);
result.fail(GenerateErrorCode::ResourceExhausted);
return result;
}

Expand Down
4 changes: 2 additions & 2 deletions server/src/qwen35moe/qwen35moe_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -862,7 +862,7 @@ GenerateResult Qwen35MoeBackend::generate_impl(const GenerateRequest & req,
std::fprintf(stderr,
"[kvflash] hybrid prompt (%d) exceeds pool %d; raise --kvflash "
"or enable pflash compression\n", prompt_len, kvflash_tokens_);
result.fail(GenerateErrorCode::ContextOverflow);
result.fail(GenerateErrorCode::ResourceExhausted);
cleanup_graphs();
return result;
}
Expand Down Expand Up @@ -1526,7 +1526,7 @@ GenerateResult Qwen35MoeBackend::restore_and_generate_impl(int slot,
std::fprintf(stderr,
"[kvflash] hybrid restore prompt (%d) exceeds pool %d; raise "
"--kvflash\n", prompt_len, kvflash_tokens_);
result.fail(GenerateErrorCode::ContextOverflow);
result.fail(GenerateErrorCode::ResourceExhausted);
out_io.emit(-1);
return result;
}
Expand Down
191 changes: 113 additions & 78 deletions server/src/server/http_server.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@

#include "http_server.h"
#include "admission.h"
#include "response_error.h"
#include "sse_emitter.h"
#include "prompt_normalize.h"
#include "tool_hint.h"
Expand Down Expand Up @@ -3566,7 +3567,7 @@ void HttpServer::finalize_generation_cache(
}
}

if (disk_cache_.disabled()) return;
if (disk_cache_.disabled() || !result.ok()) return;

if (!prepared.compressed) {
recent_disk_prompts_.insert(
Expand Down Expand Up @@ -3905,15 +3906,6 @@ void HttpServer::send_nonstream_response(
}
}

std::array<std::string, 2> HttpServer::sse_error_close_chunks(
const std::string & message) {
const json err = {{"error", {
{"message", message},
{"type", "server_error"},
}}};
return {"data: " + err.dump() + "\n\n", "data: [DONE]\n\n"};
}

void HttpServer::worker_loop() {
while (true) {
ServerJob * job = dequeue();
Expand Down Expand Up @@ -3965,19 +3957,6 @@ void HttpServer::process_job(ServerJob * job) {
job->done = true;
job->cv.notify_one();
};
auto fail_request = [&](int status, const std::string & message) {
std::fprintf(stderr, "[server] request failed: %s\n", message.c_str());
if (req.stream) {
stop_job_stream(job);
for (const std::string & chunk : sse_error_close_chunks(message)) {
send_job_bytes(job, chunk.data(), chunk.size());
}
} else {
send_error(fd, status, message);
}
finish_job();
};

std::fprintf(stderr,
"[server] chat START %s format=%s stream=%s prompt_tokens=%zu "
"max_tokens=%d tools=%zu\n",
Expand Down Expand Up @@ -4021,6 +4000,31 @@ void HttpServer::process_job(ServerJob * job) {
}
if (req.stream) start_job_stream(job);

auto fail_request = [&](int status, const std::string & message) {
std::fprintf(stderr, "[server] request failed: %s\n", message.c_str());
ResponseError error;
if (status == 400) {
error = ResponseError::invalid_request(
"invalid_request", message);
} else if (status == 503) {
error = ResponseError::unavailable("unavailable", message);
} else {
error = ResponseError::internal("server_error", message);
}
stop_job_stream(job);
if (req.stream) {
for (const std::string & chunk : emitter.emit_error(error)) {
send_job_bytes(job, chunk.data(), chunk.size());
}
} else {
const json body = build_error_response(
req.format, error, req.response_id);
send_response(fd, response_error_http_status(error),
"application/json", body.dump() + "\n");
}
finish_job();
};

PreparedPrompt prepared = prepare_prompt(req);
if (prepared.error_status != 0) {
fail_request(prepared.error_status, prepared.error);
Expand Down Expand Up @@ -4096,7 +4100,7 @@ void HttpServer::process_job(ServerJob * job) {

// Bandit: update when spec decode actually ran — including 0-accept case,
// which signals the current keep_ratio is too low.
if (!req.session_id.empty() && result.spec_decode_ran) {
if (result.ok() && !req.session_id.empty() && result.spec_decode_ran) {
float old_keep = sessions_.get_keep_ratio(req.session_id);
int old_turn = sessions_.turn_count(req.session_id);
sessions_.update(req.session_id, result.accept_rate);
Expand Down Expand Up @@ -4135,29 +4139,91 @@ void HttpServer::process_job(ServerJob * job) {
agent_turn_cache_hit,
};

// Record performance for /status page.
if (result.ok()) {
PerfRecord perf;
perf.prompt_tokens = (int)req.prompt_tokens.size();
perf.completion_tokens = completion_tokens;
// Use actual prefilled token count: on cache hit the backend only
// prefills the delta beyond the cached prefix, so dividing the full
// prompt size by delta time would be wrong.
const int prefill_tokens =
(std::max)(0, effective_prompt_tokens - cached_prefix_tokens);
perf.prefill_tok_s = (result.prefill_s > 0.0)
? (double)prefill_tokens / result.prefill_s : 0.0;
perf.decode_tok_s = (result.decode_s > 0.0)
? (double)completion_tokens / result.decode_s : 0.0;
perf.accept_rate = result.accept_rate;
perf.cache_hit = cache_hit;
perf.pflash = pflash_compressed;
perf.spec_decode = result.spec_decode_ran;
perf.timestamp = std::chrono::steady_clock::now();
status_.record_perf(perf);
status_.update_completion_tokens(completion_tokens);
broadcast_status();
auto log_done = [&]() {
const auto done_at = std::chrono::steady_clock::now();
const double elapsed_s =
std::chrono::duration<double>(done_at - started_at).count();
const int result_tokens = (int)result.tokens.size();
const int out_tokens = (std::max)(completion_tokens, result_tokens);
const double tok_s = elapsed_s > 0.0 ? out_tokens / elapsed_s : 0.0;
const double decode_tok_s =
result.decode_s > 0.0 ? out_tokens / result.decode_s : 0.0;
const std::string finish = client_disconnected
? "client_disconnect"
: (result.ok() ? emitter.finish_reason() : "error");

std::fprintf(stderr,
"[server] chat DONE %s ok=%s in=%zu effective_in=%zu out=%d "
"%.1fs %.1f tok/s finish=%s restore=%s slot=%d prefix_len=%d "
"prefill=%.1fs decode=%.1fs(%.1ftok/s) error=%s detail=%s\n",
req.response_id.c_str(),
result.ok() ? "true" : "false",
req.prompt_tokens.size(),
effective_prompt.size(),
out_tokens,
elapsed_s,
tok_s,
finish.c_str(),
using_restore ? "true" : "false",
cache_slot,
prefix_len,
result.prefill_s,
result.decode_s,
decode_tok_s,
result.ok() ? "-" : result.error_code().data(),
result.error_detail().empty() ? "-" : result.error_detail().data());
};

// A backend failure terminates the request here. Everything below this
// branch records or frames a successful generation.
if (!result.ok()) {
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
stop_job_stream(job);
if (job->client_disconnected.load(std::memory_order_acquire)) {
client_disconnected = true;
}
if (!client_disconnected) {
const ResponseError error = to_response_error(*result.error);
if (req.stream) {
for (const std::string & chunk : emitter.emit_error(error)) {
if (!send_job_bytes(job, chunk.data(), chunk.size())) {
client_disconnected = true;
break;
}
}
} else {
const json body = build_error_response(
req.format, error, req.response_id);
sock_set_block(fd);
send_response(fd, response_error_http_status(error),
"application/json", body.dump() + "\n");
}
}
log_done();
finish_job();
return;
}

// Record performance for /status page.
PerfRecord perf;
perf.prompt_tokens = (int)req.prompt_tokens.size();
perf.completion_tokens = completion_tokens;
// Use actual prefilled token count: on cache hit the backend only
// prefills the delta beyond the cached prefix, so dividing the full
// prompt size by delta time would be wrong.
const int prefill_tokens =
(std::max)(0, effective_prompt_tokens - cached_prefix_tokens);
perf.prefill_tok_s = (result.prefill_s > 0.0)
? (double)prefill_tokens / result.prefill_s : 0.0;
perf.decode_tok_s = (result.decode_s > 0.0)
? (double)completion_tokens / result.decode_s : 0.0;
perf.accept_rate = result.accept_rate;
perf.cache_hit = cache_hit;
perf.pflash = pflash_compressed;
perf.spec_decode = result.spec_decode_ran;
perf.timestamp = std::chrono::steady_clock::now();
status_.record_perf(perf);
status_.update_completion_tokens(completion_tokens);
broadcast_status();
// Serialize final frames after disabling heartbeat comments so no comment
// can appear after the protocol's [DONE] marker.
stop_job_stream(job);
Expand Down Expand Up @@ -4205,38 +4271,7 @@ void HttpServer::process_job(ServerJob * job) {
req.prompt_tokens.size(), completion_tokens);
}

const auto done_at = std::chrono::steady_clock::now();
const double elapsed_s =
std::chrono::duration<double>(done_at - started_at).count();
const int result_tokens = (int)result.tokens.size();
const int out_tokens = (std::max)(completion_tokens, result_tokens);
const double tok_s = elapsed_s > 0.0 ? out_tokens / elapsed_s : 0.0;
const double decode_tok_s =
result.decode_s > 0.0 ? out_tokens / result.decode_s : 0.0;
const std::string finish = client_disconnected
? "client_disconnect"
: (result.ok() ? emitter.finish_reason() : "error");

std::fprintf(stderr,
"[server] chat DONE %s ok=%s in=%zu effective_in=%zu out=%d "
"%.1fs %.1f tok/s finish=%s restore=%s slot=%d prefix_len=%d "
"prefill=%.1fs decode=%.1fs(%.1ftok/s) error=%s detail=%s\n",
req.response_id.c_str(),
result.ok() ? "true" : "false",
req.prompt_tokens.size(),
effective_prompt.size(),
out_tokens,
elapsed_s,
tok_s,
finish.c_str(),
using_restore ? "true" : "false",
cache_slot,
prefix_len,
result.prefill_s,
result.decode_s,
decode_tok_s,
result.ok() ? "-" : result.error_code().data(),
result.error_detail().empty() ? "-" : result.error_detail().data());
log_done();

// Signal client thread that we're done.
finish_job();
Expand Down
3 changes: 0 additions & 3 deletions server/src/server/http_server.h
Original file line number Diff line number Diff line change
Expand Up @@ -487,9 +487,6 @@ class HttpServer {
std::string format_http_response(
int status, const std::string & content_type,
const std::string & body);
static std::array<std::string, 2> sse_error_close_chunks(
const std::string & message);

// Parse HTTP request from socket.
struct HttpRequest {
std::string method;
Expand Down
Loading
Loading