Files
modelTesting/outputs/qwen3-coder-30b-6bit-mlx-rust.rs
adminandClaude 0c79a997f5 Capture kat-coder + qwen3-coder-30b outputs (5 prompts each, API-run, NOT graded yet)
Via tools/grade_run.py against LM Studio (load -> all prompts -> unload -> verify).
kat-coder:      lfu/tts/webhook/automation/rust  (lfu has a real SyntaxError: 'self._ bump_freq')
qwen3-coder-30b: lfu/tts/webhook/automation/rust  (all 4 py parse clean; terse ~2k tok outputs)
Grading deferred to a later session (capture-only mode).

Script hardening:
  - --resident / --unload-all commands (read loaded_instances, POST unload, verify)
  - MEMORY GUARD: aborts if a different model is resident (never stack 30GB models)
  - TTFT self-discovery dump (LM Studio returns empty stats non-streaming;
    fix next session: switch to /api/v1/chat stream:true -> chat.end result.stats)

Co-Authored-By: Claude <noreply@anthropic.com>
2026-07-28 19:22:41 -07:00

284 lines
8.3 KiB
Rust

use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::{mpsc, Mutex, RwLock};
use tokio::time::{sleep, timeout, Duration};
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WatcherConfig {
pub id: u32,
pub interval: Duration,
}
#[derive(Debug, Clone)]
pub struct WatcherState {
pub id: u32,
pub healthy: bool,
}
#[derive(Debug, Clone)]
pub struct AggregatedData {
pub items: Vec<String>,
pub watcher_states: HashMap<u32, WatcherState>,
}
#[derive(thiserror::Error, Debug)]
pub enum FetchError {
#[error("Mock fetch failed")]
MockFetchFailed,
}
#[derive(thiserror::Error, Debug)]
pub enum WatcherManagerError {
#[error("Watcher already exists")]
WatcherExists,
#[error("Watcher does not exist")]
WatcherNotFound,
}
pub struct WatcherManager {
watchers: Arc<RwLock<HashMap<u32, WatcherState>>>,
consumer_tx: mpsc::UnboundedSender<AggregatedData>,
shutdown_token: CancellationToken,
}
impl WatcherManager {
pub fn new() -> Self {
let (consumer_tx, _) = mpsc::unbounded_channel();
Self {
watchers: Arc::new(RwLock::new(HashMap::new())),
consumer_tx,
shutdown_token: CancellationToken::new(),
}
}
pub async fn start(&self) -> Result<(), WatcherManagerError> {
let shutdown_token = self.shutdown_token.clone();
let watchers = self.watchers.clone();
let consumer_tx = self.consumer_tx.clone();
// Start the consumer task
tokio::spawn(async move {
Self::consumer_task(watchers, consumer_tx, shutdown_token).await;
});
Ok(())
}
pub async fn add_watcher(&self, config: WatcherConfig) -> Result<(), WatcherManagerError> {
let mut watchers = self.watchers.write().await;
if watchers.contains_key(&config.id) {
return Err(WatcherManagerError::WatcherExists);
}
let watcher_state = WatcherState {
id: config.id,
healthy: true,
};
watchers.insert(config.id, watcher_state);
drop(watchers);
let shutdown_token = self.shutdown_token.clone();
let watchers = self.watchers.clone();
let consumer_tx = self.consumer_tx.clone();
tokio::spawn(async move {
Self::watcher_task(
config.id,
config.interval,
shutdown_token,
watchers,
consumer_tx,
)
.await;
});
Ok(())
}
pub async fn remove_watcher(&self, id: u32) -> Result<(), WatcherManagerError> {
let mut watchers = self.watchers.write().await;
if !watchers.contains_key(&id) {
return Err(WatcherManagerError::WatcherNotFound);
}
watchers.remove(&id);
Ok(())
}
pub async fn shutdown(&self) {
self.shutdown_token.cancel();
// Give tasks a chance to finish gracefully
sleep(Duration::from_millis(100)).await;
}
pub async fn get_watcher_count(&self) -> usize {
self.watchers.read().await.len()
}
async fn consumer_task(
watchers: Arc<RwLock<HashMap<u32, WatcherState>>>,
consumer_tx: mpsc::UnboundedSender<AggregatedData>,
shutdown_token: CancellationToken,
) {
let mut items = Vec::new();
loop {
tokio::select! {
_ = shutdown_token.cancelled() => {
info!("Consumer task shutting down");
break;
}
_ = sleep(Duration::from_millis(100)) => {
let watcher_states = watchers.read().await.clone();
let data = AggregatedData {
items: items.clone(),
watcher_states,
};
if let Err(e) = consumer_tx.send(data) {
error!("Failed to send aggregated data: {:?}", e);
}
items.clear();
}
}
}
}
async fn watcher_task(
id: u32,
interval: Duration,
shutdown_token: CancellationToken,
watchers: Arc<RwLock<HashMap<u32, WatcherState>>>,
consumer_tx: mpsc::UnboundedSender<AggregatedData>,
) {
let mut consecutive_errors = 0;
loop {
tokio::select! {
_ = shutdown_token.cancelled() => {
info!("Watcher {} shutting down", id);
break;
}
_ = sleep(interval) => {
match Self::mock_fetch(id).await {
Ok(new_items) => {
consecutive_errors = 0;
if let Err(e) = consumer_tx.send(AggregatedData {
items: new_items,
watcher_states: HashMap::new(),
}) {
error!("Failed to send items from watcher {}: {:?}", id, e);
}
}
Err(e) => {
consecutive_errors += 1;
error!("Watcher {} failed to fetch: {:?}", id, e);
if consecutive_errors >= 5 {
info!("Watcher {} marked unhealthy after {} consecutive errors", id, consecutive_errors);
let mut w = watchers.write().await;
if let Some(state) = w.get_mut(&id) {
state.healthy = false;
}
}
}
}
}
}
}
}
async fn mock_fetch(id: u32) -> Result<Vec<String>, FetchError> {
// Simulate network delay
sleep(Duration::from_millis(10)).await;
// ~15% chance of failure
if rand::random::<f64>() < 0.15 {
return Err(FetchError::MockFetchFailed);
}
// Return 0-3 random items
let count = rand::random::<usize>() % 4;
let mut items = Vec::with_capacity(count);
for i in 0..count {
items.push(format!("item_{}_{}", id, i));
}
Ok(items)
}
}
#[cfg(test)]
mod tests {
use super::*;
use tokio::time::timeout;
#[tokio::test]
async fn test_manager_basic_functionality() {
let manager = WatcherManager::new();
manager.start().await.unwrap();
// Add 4 watchers
for i in 0..4 {
manager
.add_watcher(WatcherConfig {
id: i,
interval: Duration::from_millis(50),
})
.await
.unwrap();
}
// Let it run for a bit
sleep(Duration::from_millis(500)).await;
// Check that we have 4 watchers
assert_eq!(manager.get_watcher_count().await, 4);
// Shutdown
manager.shutdown().await;
}
#[tokio::test]
async fn test_unhealthy_watcher() {
let manager = WatcherManager::new();
manager.start().await.unwrap();
// Add a watcher that always fails
manager
.add_watcher(WatcherConfig {
id: 999,
interval: Duration::from_millis(10),
})
.await
.unwrap();
// Wait for it to become unhealthy
sleep(Duration::from_millis(100)).await;
// Check that it's marked as unhealthy
let watchers = manager.watchers.read().await;
assert_eq!(watchers.get(&999).unwrap().healthy, false);
}
#[tokio::test]
async fn test_add_remove_watchers() {
let manager = WatcherManager::new();
manager.start().await.unwrap();
// Add a watcher
manager
.add_watcher(WatcherConfig {
id: 1,
interval: Duration::from_millis(100),
})
.await
.unwrap();
assert_eq!(manager.get_watcher_count().await, 1);
// Remove it
manager.remove_watcher(1).await.unwrap();
assert_eq!(manager.get_watcher_count().await, 0);
}
}