Compare commits
7 Commits
revolution
...
591cfcb04b
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
591cfcb04b | ||
|
|
3cda57f0bc | ||
|
|
23e7b92d03 | ||
|
|
9f58febe21 | ||
|
|
b1de0d37f7 | ||
|
|
4ff626ab88 | ||
|
|
a45616046d |
@@ -9,14 +9,14 @@ mkdir -p "$OUT"
|
||||
echo "=== Kipinä Node — Binary Build ==="
|
||||
|
||||
# macOS ARM (natiivi)
|
||||
echo "[1/3] macOS ARM64..."
|
||||
echo "[1/4] macOS ARM64..."
|
||||
cd "$SCRIPT_DIR"
|
||||
cargo build --release -p native-node --no-default-features 2>&1 | tail -1
|
||||
cp target/release/native-node "$OUT/kipina-node-macos-arm64"
|
||||
echo " $(ls -lh "$OUT/kipina-node-macos-arm64" | awk '{print $5}')"
|
||||
|
||||
# Linux x86_64 (Docker)
|
||||
echo "[2/3] Linux x86_64..."
|
||||
echo "[2/4] Linux x86_64..."
|
||||
docker run --rm \
|
||||
-v "$SCRIPT_DIR":/app -w /app \
|
||||
--platform linux/amd64 \
|
||||
@@ -25,7 +25,7 @@ docker run --rm \
|
||||
echo " $(ls -lh "$OUT/kipina-node-linux-x86_64" | awk '{print $5}')"
|
||||
|
||||
# Linux ARM64 (Docker)
|
||||
echo "[3/3] Linux ARM64..."
|
||||
echo "[3/4] Linux ARM64..."
|
||||
docker run --rm \
|
||||
-v "$SCRIPT_DIR":/app -w /app \
|
||||
--platform linux/arm64 \
|
||||
@@ -33,6 +33,15 @@ docker run --rm \
|
||||
bash -c "apt-get update -qq && apt-get install -y -qq pkg-config libssl-dev >/dev/null 2>&1 && cargo build --release -p native-node --no-default-features 2>&1 | tail -1 && cp target/release/native-node /app/frontend/public/download/kipina-node-linux-arm64"
|
||||
echo " $(ls -lh "$OUT/kipina-node-linux-arm64" | awk '{print $5}')"
|
||||
|
||||
# Windows x86_64 (Docker + mingw-w64)
|
||||
echo "[4/4] Windows x86_64..."
|
||||
docker run --rm \
|
||||
-v "$SCRIPT_DIR":/app -w /app \
|
||||
--platform linux/amd64 \
|
||||
rust:slim \
|
||||
bash -c "apt-get update -qq && apt-get install -y -qq gcc-mingw-w64-x86-64 pkg-config libssl-dev >/dev/null 2>&1 && rustup target add x86_64-pc-windows-gnu && cargo build --release -p native-node --no-default-features --target x86_64-pc-windows-gnu 2>&1 | tail -1 && cp target/x86_64-pc-windows-gnu/release/native-node.exe /app/frontend/public/download/kipina-node-windows-x86_64.exe"
|
||||
echo " $(ls -lh "$OUT/kipina-node-windows-x86_64.exe" | awk '{print $5}')"
|
||||
|
||||
echo ""
|
||||
echo "=== Binäärit valmiina ==="
|
||||
ls -lh "$OUT"/kipina-node-*
|
||||
|
||||
@@ -35,25 +35,29 @@ if ! git -C "$SCRIPT_DIR" diff --quiet HEAD 2>/dev/null || \
|
||||
echo " Commitoitu: $DEPLOY_MSG"
|
||||
fi
|
||||
|
||||
# 1. Rakennetaan Docker-image lokaalisti
|
||||
echo "[1/4] Rakennetaan image lokaalisti..."
|
||||
# 1. Käännetään native-node-binäärit kaikille alustoille
|
||||
echo "[1/6] Käännetään native-node-binäärit..."
|
||||
./build-binaries.sh
|
||||
|
||||
# 2. Rakennetaan Docker-image lokaalisti
|
||||
echo "[2/6] Rakennetaan image lokaalisti..."
|
||||
docker build --platform linux/amd64 -f Dockerfile.prod -t kipina-agentic:latest .
|
||||
|
||||
# 2. Tallennetaan tiedostoon
|
||||
echo "[2/5] Pakataan image..."
|
||||
# 3. Tallennetaan tiedostoon
|
||||
echo "[3/6] Pakataan image..."
|
||||
docker save kipina-agentic:latest | gzip > /tmp/kipina-agentic.tar.gz
|
||||
echo " Koko: $(du -h /tmp/kipina-agentic.tar.gz | cut -f1)"
|
||||
|
||||
# 3. Siirretään palvelimelle
|
||||
echo "[3/5] Siirretään palvelimelle..."
|
||||
# 4. Siirretään palvelimelle
|
||||
echo "[4/6] Siirretään palvelimelle..."
|
||||
scp $SSH_OPTS /tmp/kipina-agentic.tar.gz $SERVER:/tmp/
|
||||
scp $SSH_OPTS docker-compose.prod.yml Caddyfile.prod $SERVER:$REMOTE_DIR/
|
||||
|
||||
# 4. Ladataan image ja käynnistetään
|
||||
echo "[4/5] Ladataan image palvelimella..."
|
||||
# 5. Ladataan image ja käynnistetään
|
||||
echo "[5/6] Ladataan image palvelimella..."
|
||||
ssh $SSH_OPTS $SERVER "gunzip -c /tmp/kipina-agentic.tar.gz | docker load && rm /tmp/kipina-agentic.tar.gz"
|
||||
|
||||
echo "[5/5] Käynnistetään palvelut uudelleen..."
|
||||
echo "[6/6] Käynnistetään palvelut uudelleen..."
|
||||
ssh $SSH_OPTS $SERVER "cd $REMOTE_DIR && docker compose -f docker-compose.prod.yml down && docker compose -f docker-compose.prod.yml up -d"
|
||||
|
||||
echo "=== Valmis! https://kipina.studio ==="
|
||||
|
||||
BIN
network-poc/frontend/public/download/kipina-node-linux-arm64
Executable file
BIN
network-poc/frontend/public/download/kipina-node-linux-arm64
Executable file
Binary file not shown.
Binary file not shown.
Binary file not shown.
BIN
network-poc/frontend/public/download/kipina-node-windows-x86_64.exe
Executable file
BIN
network-poc/frontend/public/download/kipina-node-windows-x86_64.exe
Executable file
Binary file not shown.
@@ -42,7 +42,9 @@ struct AppState {
|
||||
node_types: Mutex<HashMap<u64, String>>, // node_id → "native" | "browser"
|
||||
node_busy: Mutex<std::collections::HashSet<u64>>, // Solmut joilla on aktiivinen tehtävä
|
||||
pending_task_ids: Mutex<std::collections::HashSet<String>>, // Hubin jakamat task_id:t (gamification-validointi)
|
||||
pending_responses: Mutex<HashMap<String, tokio::sync::oneshot::Sender<serde_json::Value>>>, // task_id → oneshot API-vastaukselle
|
||||
api_rate_limits: Mutex<HashMap<IpAddr, (std::time::Instant, u32)>>, // IP → (ikkuna-alku, pyyntömäärä)
|
||||
node_models: tokio::sync::RwLock<HashMap<u64, serde_json::Value>>, // node_id → ollama tags JSON
|
||||
db: db::NodeDb,
|
||||
}
|
||||
|
||||
@@ -91,6 +93,7 @@ tr:hover td { background:#1c2333; }
|
||||
<div class="tabs">
|
||||
<div class="tab active" onclick="showTab('sessions')">Sessiot</div>
|
||||
<div class="tab" onclick="showTab('pairs')">Tokenisointiparit</div>
|
||||
<div class="tab" onclick="showTab('hardware')">Laitteisto & Mallit</div>
|
||||
</div>
|
||||
|
||||
<div id="sessions" class="panel active">
|
||||
@@ -118,6 +121,19 @@ tr:hover td { background:#1c2333; }
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div id="hardware" class="panel">
|
||||
<div class="stats-grid" id="hardware-stats"></div>
|
||||
<h2 style="margin-top: 10px; margin-bottom: 10px; color: var(--accent); font-size: 16px;">Käytettävissä olevat paikalliset kielimallit</h2>
|
||||
<div class="table-wrap">
|
||||
<table>
|
||||
<thead><tr>
|
||||
<th>Nimi</th><th>Koko</th><th>Parametrit</th>
|
||||
</tr></thead>
|
||||
<tbody id="models-body"></tbody>
|
||||
</table>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<script>
|
||||
function showTab(name) {
|
||||
document.querySelectorAll('.panel').forEach(p => p.classList.remove('active'));
|
||||
@@ -149,12 +165,16 @@ function duration(start, end) {
|
||||
}
|
||||
|
||||
async function load() {
|
||||
const [statsRes, sessionsRes, pairsRes] = await Promise.all([
|
||||
fetch('/api/stats'), fetch('/api/sessions'), fetch('/api/pairs')
|
||||
const [statsRes, sessionsRes, pairsRes, hwRes, modelsRes] = await Promise.all([
|
||||
fetch('/api/stats'), fetch('/api/sessions'), fetch('/api/pairs'),
|
||||
fetch('/api/v1/hardware').catch(() => ({json: async()=>({gpu_name:'', vram_mb:0, ram_mb:0})})),
|
||||
fetch('/api/v1/ollama/tags').catch(() => ({json: async()=>({models:[]})}))
|
||||
]);
|
||||
const stats = await statsRes.json();
|
||||
const sessions = await sessionsRes.json();
|
||||
const pairs = await pairsRes.json();
|
||||
const hw = await hwRes.json().catch(() => ({gpu_name:'', vram_mb:0, ram_mb:0}));
|
||||
const modelsData = await modelsRes.json().catch(() => ({models:[]}));
|
||||
|
||||
// Versio
|
||||
if (stats.version) document.getElementById('admin-version').textContent = 'v' + stats.version;
|
||||
@@ -229,6 +249,24 @@ async function load() {
|
||||
<td>${p.duration_ms||0}ms</td>
|
||||
</tr>`;
|
||||
}).join('');
|
||||
|
||||
// Hardware
|
||||
document.getElementById('hardware-stats').innerHTML = [
|
||||
{v: hw.gpu_name || '-', l: 'Paikallinen GPU tila'},
|
||||
{v: hw.vram_mb ? hw.vram_mb + ' MB' : '-', l: 'GPU Muisti (VRAM)'},
|
||||
{v: hw.ram_mb ? hw.ram_mb + ' MB' : '-', l: 'RAM'},
|
||||
].map(s => `<div class="stat-card"><div class="val">${s.v}</div><div class="label">${s.l}</div></div>`).join('');
|
||||
|
||||
// Models
|
||||
document.getElementById('models-body').innerHTML = (modelsData.models || []).map(m => {
|
||||
const sizeGb = (m.size / (1024*1024*1024)).toFixed(2) + ' GB';
|
||||
const params = m.details?.parameter_size || '-';
|
||||
return `<tr>
|
||||
<td><strong>${m.name}</strong></td>
|
||||
<td>${sizeGb}</td>
|
||||
<td>${params}</td>
|
||||
</tr>`;
|
||||
}).join('');
|
||||
}
|
||||
|
||||
load();
|
||||
@@ -264,7 +302,9 @@ async fn main() {
|
||||
node_types: Mutex::new(HashMap::new()),
|
||||
node_busy: Mutex::new(std::collections::HashSet::new()),
|
||||
pending_task_ids: Mutex::new(std::collections::HashSet::new()),
|
||||
pending_responses: Mutex::new(HashMap::new()),
|
||||
api_rate_limits: Mutex::new(HashMap::new()),
|
||||
node_models: tokio::sync::RwLock::new(HashMap::new()),
|
||||
db: db::NodeDb::new(&std::env::var("DATABASE_PATH").unwrap_or_else(|_| "nodes.db".to_string())),
|
||||
});
|
||||
|
||||
@@ -743,6 +783,12 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, ip: IpAddr) {
|
||||
node_id, ip, hostname, os, cores, ram, allocated
|
||||
);
|
||||
|
||||
// Tallennetaan välitetyt mallit muistiin
|
||||
if let Some(models) = json.get("models") {
|
||||
let mut nm = state.node_models.write().await;
|
||||
nm.insert(node_id, models.clone());
|
||||
}
|
||||
|
||||
if let Some(gpus) = json.get("gpus").and_then(|v| v.as_array()) {
|
||||
for gpu in gpus {
|
||||
tracing::info!(
|
||||
@@ -875,11 +921,18 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, ip: IpAddr) {
|
||||
} else if msg_type == "llm_done" {
|
||||
// Vapautetaan solmu ja tarkistetaan task_id:n aitous
|
||||
state.node_busy.lock().unwrap().remove(&node_id);
|
||||
let valid_task = if let Some(tid) = json.get("task_id").and_then(|v| v.as_str()) {
|
||||
state.pending_task_ids.lock().unwrap().remove(tid)
|
||||
let task_id = json.get("task_id").and_then(|v| v.as_str()).map(|s| s.to_string());
|
||||
let valid_task = if let Some(ref tid) = task_id {
|
||||
state.pending_task_ids.lock().unwrap().remove(tid.as_str())
|
||||
} else {
|
||||
false
|
||||
};
|
||||
|
||||
// Jos API-pyyntö odottaa tätä vastausta, reititetään suoraan oneshot-kanavaan
|
||||
let api_sender = task_id.as_ref().and_then(|tid| {
|
||||
state.pending_responses.lock().unwrap().remove(tid)
|
||||
});
|
||||
|
||||
{
|
||||
let mut json = json;
|
||||
if let Some(obj) = json.as_object_mut() {
|
||||
@@ -899,6 +952,12 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, ip: IpAddr) {
|
||||
state.db.increment_tasks(node_id);
|
||||
obj.insert("node_id".to_string(), serde_json::json!(node_id));
|
||||
}
|
||||
|
||||
if let Some(sender) = api_sender {
|
||||
// API-pyyntö: reititetään vastaus suoraan odottajalle
|
||||
let _ = sender.send(json.clone());
|
||||
}
|
||||
// UI-broadcast jatkuu normaalisti
|
||||
let _ = state.stats_tx.send(json.to_string());
|
||||
|
||||
let active_incentives = state.feature_flags.read().await.get("Insentiivit").copied().unwrap_or(false);
|
||||
@@ -908,7 +967,7 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, ip: IpAddr) {
|
||||
{
|
||||
let mut task_count = state.total_tasks.lock().unwrap();
|
||||
*task_count += 1;
|
||||
|
||||
|
||||
if active_incentives && valid_task {
|
||||
let mut tokens = state.nodes_tokens.lock().unwrap();
|
||||
let balance = tokens.entry(node_id).or_insert(0);
|
||||
@@ -916,7 +975,7 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, ip: IpAddr) {
|
||||
current_balance = *balance;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
if active_incentives && ui_sync {
|
||||
if let Some(tx) = state.node_channels.read().await.get(&node_id) {
|
||||
let msg = serde_json::json!({
|
||||
@@ -926,45 +985,50 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, ip: IpAddr) {
|
||||
let _ = tx.send(msg.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
broadcast_stats(&state).await;
|
||||
}
|
||||
} else if msg_type == "llm_error" {
|
||||
state.node_busy.lock().unwrap().remove(&node_id);
|
||||
if let Some(tid) = json.get("task_id").and_then(|v| v.as_str()) {
|
||||
state.pending_task_ids.lock().unwrap().remove(tid);
|
||||
let task_id = json.get("task_id").and_then(|v| v.as_str()).map(|s| s.to_string());
|
||||
if let Some(ref tid) = task_id {
|
||||
state.pending_task_ids.lock().unwrap().remove(tid.as_str());
|
||||
}
|
||||
// Jos API-pyyntö odottaa, reititetään virhe oneshot-kanavaan
|
||||
let api_sender = task_id.as_ref().and_then(|tid| {
|
||||
state.pending_responses.lock().unwrap().remove(tid)
|
||||
});
|
||||
{
|
||||
let mut json = json;
|
||||
if let Some(obj) = json.as_object_mut() {
|
||||
obj.insert("node_id".to_string(), serde_json::json!(node_id));
|
||||
}
|
||||
if let Some(sender) = api_sender {
|
||||
let _ = sender.send(json.clone());
|
||||
}
|
||||
let _ = state.stats_tx.send(json.to_string());
|
||||
}
|
||||
} else if msg_type == "user_text" {
|
||||
// Käyttäjän lähettämä teksti — broadcastataan pair_taskina ja llm_promptina
|
||||
// Käyttäjän lähettämä teksti — kohdennettu reititys lähettäjäsolmulle
|
||||
let text = json.get("text").and_then(|v| v.as_str()).unwrap_or("").to_string();
|
||||
let task_type = json.get("task_type").and_then(|v| v.as_str()).unwrap_or("tokenize");
|
||||
if !text.is_empty() {
|
||||
let preview: String = text.chars().take(80).collect();
|
||||
tracing::info!("Solmu {} lähetti oman tekstin ({}): \"{}\"", node_id, task_type, preview);
|
||||
match task_type {
|
||||
"tokenize" => {
|
||||
let msg = serde_json::json!({
|
||||
"type": "single_tokenize",
|
||||
"text": text,
|
||||
});
|
||||
let _ = state.stats_tx.send(msg.to_string());
|
||||
}
|
||||
_ => {
|
||||
// LLM-prompti: lähetetään VAIN valitulle mallille, ei kaikille (välttää turhaa ruuhkaa ja busy-tiloja)
|
||||
let prompt = serde_json::json!({
|
||||
"type": "llm_prompt",
|
||||
"prompt": text,
|
||||
"model": task_type,
|
||||
});
|
||||
let _ = state.stats_tx.send(prompt.to_string());
|
||||
}
|
||||
let msg = match task_type {
|
||||
"tokenize" => serde_json::json!({
|
||||
"type": "single_tokenize",
|
||||
"text": text,
|
||||
}),
|
||||
_ => serde_json::json!({
|
||||
"type": "llm_prompt",
|
||||
"prompt": text,
|
||||
"model": task_type,
|
||||
}),
|
||||
};
|
||||
// Lähetetään takaisin lähettäjäsolmulle (käyttäjä haluaa oman tekstinsä tuloksen)
|
||||
if let Some(tx) = state.node_channels.read().await.get(&node_id) {
|
||||
let _ = tx.send(msg.to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -989,6 +1053,7 @@ async fn handle_socket(socket: WebSocket, state: Arc<AppState>, ip: IpAddr) {
|
||||
vram.remove(&node_id);
|
||||
}
|
||||
state.node_types.lock().unwrap().remove(&node_id);
|
||||
state.node_models.write().await.remove(&node_id);
|
||||
tracing::info!("Solmu {} ({}) poistui verkosta.", node_id, ip);
|
||||
broadcast_stats(&state).await;
|
||||
sender_task.abort();
|
||||
@@ -1009,7 +1074,16 @@ struct ChatCompletionResponse {
|
||||
tokens_generated: u64,
|
||||
}
|
||||
|
||||
async fn api_ollama_tags() -> axum::response::Response {
|
||||
async fn api_ollama_tags(
|
||||
axum::extract::State(state): axum::extract::State<Arc<AppState>>,
|
||||
) -> axum::response::Response {
|
||||
// Haetaan natiivisolmun tila muistista — priorisoidaan aito verkko-solmu
|
||||
let node_models = state.node_models.read().await;
|
||||
if let Some((_, models_json)) = node_models.iter().next() {
|
||||
return axum::Json(models_json.clone()).into_response();
|
||||
}
|
||||
|
||||
// Fallback: Haetaan lokaalista infra-Ollamasta ohjaimesta käsin (esim dev ympäristö)
|
||||
let ollama_url = std::env::var("OLLAMA_URL").unwrap_or_else(|_| "http://ollama:11434".to_string());
|
||||
match reqwest::get(format!("{}/api/tags", ollama_url)).await {
|
||||
Ok(resp) => {
|
||||
@@ -1033,11 +1107,10 @@ async fn api_hardware(
|
||||
});
|
||||
|
||||
let (mut vram_mb, mut gpu_name, ram_mb) = if let Some(s) = native {
|
||||
let gpus = s.get("gpus").and_then(|v| v.as_array());
|
||||
let gpu = gpus.and_then(|g| g.first());
|
||||
let vram = gpu.and_then(|g| g.get("vram_total_mb")).and_then(|v| v.as_u64()).unwrap_or(0);
|
||||
let name = gpu.and_then(|g| g.get("name")).and_then(|v| v.as_str()).unwrap_or("").to_string();
|
||||
let ram = s.get("system").and_then(|v| v.get("ram_total_mb")).and_then(|v| v.as_u64()).unwrap_or(0);
|
||||
// Tieto on tietokannassa litteänä
|
||||
let vram = s.get("vram_total_mb").and_then(|v| v.as_u64()).unwrap_or(0);
|
||||
let name = s.get("gpu_name").and_then(|v| v.as_str()).unwrap_or("").to_string();
|
||||
let ram = s.get("ram_mb").and_then(|v| v.as_u64()).unwrap_or(0);
|
||||
(vram, name, ram)
|
||||
} else {
|
||||
(0, String::new(), 0)
|
||||
@@ -1161,8 +1234,9 @@ async fn api_chat_completions(
|
||||
msg.as_object_mut().unwrap().insert("max_tokens".to_string(), serde_json::json!(mt));
|
||||
}
|
||||
|
||||
// Odotuskanava valmiiksi (solmu palauttaa tuloksen stats_tx kautta)
|
||||
let mut rx = state.stats_tx.subscribe();
|
||||
// Oneshot-kanava: solmu palauttaa tuloksen suoraan tälle pyynnölle
|
||||
let (resp_tx, resp_rx) = tokio::sync::oneshot::channel::<serde_json::Value>();
|
||||
state.pending_responses.lock().unwrap().insert(payload.task_id.clone(), resp_tx);
|
||||
|
||||
// Kohdennettu reititys: lähetetään AI-tehtävä suoraan VAIN valitulle solmulle
|
||||
{
|
||||
@@ -1171,48 +1245,34 @@ async fn api_chat_completions(
|
||||
let _ = tx.send(msg.to_string());
|
||||
tracing::info!("Reititettiin API-pyyntö solmulle {} (Malli: {})", target_node_id, payload.model);
|
||||
} else {
|
||||
state.pending_responses.lock().unwrap().remove(&payload.task_id);
|
||||
return (axum::http::StatusCode::SERVICE_UNAVAILABLE, "Verkkovirhe: solmun yhteys katkesi reitityksen aikana").into_response();
|
||||
}
|
||||
}
|
||||
|
||||
let timeout = tokio::time::timeout(std::time::Duration::from_secs(600), async move {
|
||||
loop {
|
||||
let msg_str = match rx.recv().await {
|
||||
Ok(msg) => msg,
|
||||
Err(broadcast::error::RecvError::Lagged(n)) => {
|
||||
tracing::debug!("API-kanava lagged {} viestiä", n);
|
||||
continue;
|
||||
}
|
||||
Err(_) => return Ok(None), // Kanava suljettu
|
||||
};
|
||||
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&msg_str) {
|
||||
if v["type"].as_str() == Some("llm_done") {
|
||||
if let Some(tid) = v["task_id"].as_str() {
|
||||
if tid == payload.task_id {
|
||||
return Ok(Some(ChatCompletionResponse {
|
||||
response: v["response"].as_str().unwrap_or("").to_string(),
|
||||
model: v["model"].as_str().unwrap_or("").to_string(),
|
||||
tokens_generated: v["tokens_generated"].as_u64().unwrap_or(0),
|
||||
}));
|
||||
}
|
||||
}
|
||||
} else if v["type"].as_str() == Some("llm_error") {
|
||||
if let Some(tid) = v["task_id"].as_str() {
|
||||
if tid == payload.task_id {
|
||||
return Err(v["error"].as_str().unwrap_or("Määrittelemätön virhe solmussa").to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
#[allow(unreachable_code)]
|
||||
Ok(None)
|
||||
}).await;
|
||||
let timeout = tokio::time::timeout(std::time::Duration::from_secs(600), resp_rx).await;
|
||||
|
||||
match timeout {
|
||||
Ok(Ok(Some(res))) => axum::Json(res).into_response(),
|
||||
Ok(Ok(None)) => (axum::http::StatusCode::INTERNAL_SERVER_ERROR, "Verkkovirhe: yhteys katkesi").into_response(),
|
||||
Ok(Err(err)) => (axum::http::StatusCode::CONFLICT, err).into_response(),
|
||||
Err(_) => (axum::http::StatusCode::GATEWAY_TIMEOUT, "Aikakatkaisu: solmu ei saanut tehtävää ajoissa valmiiksi").into_response(),
|
||||
Ok(Ok(v)) => {
|
||||
if v["type"].as_str() == Some("llm_error") {
|
||||
let err = v["error"].as_str().unwrap_or("Määrittelemätön virhe solmussa").to_string();
|
||||
(axum::http::StatusCode::CONFLICT, err).into_response()
|
||||
} else {
|
||||
axum::Json(ChatCompletionResponse {
|
||||
response: v["response"].as_str().unwrap_or("").to_string(),
|
||||
model: v["model"].as_str().unwrap_or("").to_string(),
|
||||
tokens_generated: v["tokens_generated"].as_u64().unwrap_or(0),
|
||||
}).into_response()
|
||||
}
|
||||
}
|
||||
Ok(Err(_)) => {
|
||||
// Oneshot-kanava sulkeutui (solmu katosi)
|
||||
state.pending_responses.lock().unwrap().remove(&payload.task_id);
|
||||
(axum::http::StatusCode::INTERNAL_SERVER_ERROR, "Verkkovirhe: yhteys katkesi").into_response()
|
||||
}
|
||||
Err(_) => {
|
||||
state.pending_responses.lock().unwrap().remove(&payload.task_id);
|
||||
(axum::http::StatusCode::GATEWAY_TIMEOUT, "Aikakatkaisu: solmu ei saanut tehtävää ajoissa valmiiksi").into_response()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -78,6 +78,20 @@ impl LlmEngine {
|
||||
}
|
||||
}
|
||||
|
||||
/// Hakee kaikki Ollamaan asennetut mallit
|
||||
pub async fn fetch_models(&self) -> Result<serde_json::Value, String> {
|
||||
let resp = self.client.get(format!("{}/api/tags", self.ollama_url))
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| format!("Ollama tags fetch: {}", e))?;
|
||||
|
||||
if resp.status().is_success() {
|
||||
resp.json().await.map_err(|e| format!("Ollama tags json: {}", e))
|
||||
} else {
|
||||
Err(format!("Ollama tags epäonnistui: {}", resp.status()))
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn generate(&self, prompt: &str, max_tokens: usize) -> Result<GenerateResult, String> {
|
||||
// System prompt tulee agentin konfiguraatiosta (frontend lähettää sen osana promptia).
|
||||
// Tässä ei yliajeta sitä — Ollama saa vain prompt-kentän.
|
||||
|
||||
@@ -222,7 +222,7 @@ fn collect_system_info() -> serde_json::Value {
|
||||
}
|
||||
|
||||
/// Koko auth-viesti hubille
|
||||
fn build_auth_message(allocated_gb: u32) -> String {
|
||||
fn build_auth_message(allocated_gb: u32, model_name: &str, models_data: Option<serde_json::Value>) -> String {
|
||||
let sys = collect_system_info();
|
||||
let gpus = collect_all_gpus();
|
||||
|
||||
@@ -239,7 +239,7 @@ fn build_auth_message(allocated_gb: u32) -> String {
|
||||
"status": "agent_ready",
|
||||
"node_type": "native",
|
||||
"allocated_gb": allocated_gb,
|
||||
"selected_task": "qwen2.5-coder:7b",
|
||||
"selected_task": model_name,
|
||||
"system": sys,
|
||||
});
|
||||
|
||||
@@ -251,6 +251,10 @@ fn build_auth_message(allocated_gb: u32) -> String {
|
||||
msg.as_object_mut().unwrap().insert("gpus".to_string(), json!(gpu_json));
|
||||
}
|
||||
|
||||
if let Some(models) = models_data {
|
||||
msg.as_object_mut().unwrap().insert("models".to_string(), models);
|
||||
}
|
||||
|
||||
msg.to_string()
|
||||
}
|
||||
|
||||
@@ -321,6 +325,22 @@ async fn main() {
|
||||
}
|
||||
};
|
||||
|
||||
let active_model = llm.as_ref().map(|e| e.model_name()).unwrap_or_else(|| "unknown".to_string());
|
||||
tracing::info!("Käytettävä kielimalli konfiguroitu (selected_task): {}", active_model);
|
||||
|
||||
// Haetaan paikalliset mallit hubille lähetettäväksi
|
||||
let mut available_models = None;
|
||||
if let Some(ref engine) = llm {
|
||||
match engine.fetch_models().await {
|
||||
Ok(models) => {
|
||||
available_models = Some(models);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!("Mallilistauksen haku epäonnistui: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Yhdistetään hubiin
|
||||
loop {
|
||||
match connect_async(&hub_url).await {
|
||||
@@ -328,7 +348,7 @@ async fn main() {
|
||||
tracing::info!("Yhdistetty hubiin!");
|
||||
let (mut write, mut read) = ws_stream.split();
|
||||
|
||||
let auth = build_auth_message(allocated_gb);
|
||||
let auth = build_auth_message(allocated_gb, &active_model, available_models.clone());
|
||||
if write.send(Message::Text(auth)).await.is_err() {
|
||||
tracing::error!("Auth-viestin lähetys epäonnistui");
|
||||
continue;
|
||||
|
||||
Reference in New Issue
Block a user