new section between containers and questions: lists every name with a state dir under /var/lib/hyperhive/agents/ that doesn't correspond to a live container. shows state size + last-modified age + whether claude creds are kept. two actions per row: - R3V1V3 — queues a spawn approval with the same name (operator approves to recreate; spawn flow reuses prior config + claude creds, no re-login needed) - PURG3 — wipes the agent's state + applied dirs (post /purge-tombstone/ endpoint; refuses if a live container with that name still exists) dashboard also opens agent links in new tabs now (target=_blank + rel=noopener) so the operator's overview tab stays put when they dive into an agent.
167 lines
5.5 KiB
Rust
167 lines
5.5 KiB
Rust
//! Operator question queue. Manager submits via `AskOperator`; the
|
|
//! operator answers via the dashboard. The manager-socket handler long-polls
|
|
//! the store until the answer lands, so claude's `ask_operator` tool call
|
|
//! returns the answer directly as its result.
|
|
|
|
use std::path::Path;
|
|
use std::sync::Mutex;
|
|
use std::time::{SystemTime, UNIX_EPOCH};
|
|
|
|
use anyhow::{Context, Result, bail};
|
|
use rusqlite::{Connection, OptionalExtension, params};
|
|
use serde::Serialize;
|
|
|
|
const SCHEMA: &str = r"
|
|
CREATE TABLE IF NOT EXISTS operator_questions (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
asker TEXT NOT NULL,
|
|
question TEXT NOT NULL,
|
|
options_json TEXT NOT NULL,
|
|
asked_at INTEGER NOT NULL,
|
|
answered_at INTEGER,
|
|
answer TEXT
|
|
);
|
|
CREATE INDEX IF NOT EXISTS idx_operator_questions_pending
|
|
ON operator_questions (id) WHERE answered_at IS NULL;
|
|
";
|
|
|
|
/// Add the `multi` column to pre-existing databases. `ALTER TABLE ADD COLUMN`
|
|
/// has no `IF NOT EXISTS` form in sqlite, so we check `pragma_table_info` first.
|
|
fn ensure_multi_column(conn: &Connection) -> Result<()> {
|
|
let has: bool = conn
|
|
.prepare("SELECT 1 FROM pragma_table_info('operator_questions') WHERE name = 'multi'")?
|
|
.exists([])?;
|
|
if !has {
|
|
conn.execute_batch(
|
|
"ALTER TABLE operator_questions ADD COLUMN multi INTEGER NOT NULL DEFAULT 0;",
|
|
)
|
|
.context("add operator_questions.multi column")?;
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
#[derive(Debug, Clone, Serialize)]
|
|
#[allow(clippy::doc_markdown)]
|
|
pub struct OpQuestion {
|
|
pub id: i64,
|
|
pub asker: String,
|
|
pub question: String,
|
|
pub options: Vec<String>,
|
|
pub multi: bool,
|
|
pub asked_at: i64,
|
|
pub answered_at: Option<i64>,
|
|
pub answer: Option<String>,
|
|
}
|
|
|
|
pub struct OperatorQuestions {
|
|
conn: Mutex<Connection>,
|
|
}
|
|
|
|
impl OperatorQuestions {
|
|
pub fn open(path: &Path) -> Result<Self> {
|
|
if let Some(parent) = path.parent() {
|
|
std::fs::create_dir_all(parent).with_context(|| {
|
|
format!("create operator_questions db parent {}", parent.display())
|
|
})?;
|
|
}
|
|
let conn = Connection::open(path)
|
|
.with_context(|| format!("open operator_questions db {}", path.display()))?;
|
|
conn.execute_batch(SCHEMA)
|
|
.context("apply operator_questions schema")?;
|
|
ensure_multi_column(&conn).context("migrate operator_questions.multi")?;
|
|
Ok(Self {
|
|
conn: Mutex::new(conn),
|
|
})
|
|
}
|
|
|
|
pub fn submit(
|
|
&self,
|
|
asker: &str,
|
|
question: &str,
|
|
options: &[String],
|
|
multi: bool,
|
|
) -> Result<i64> {
|
|
let conn = self.conn.lock().unwrap();
|
|
let options_json = serde_json::to_string(options).unwrap_or_else(|_| "[]".into());
|
|
conn.execute(
|
|
"INSERT INTO operator_questions (asker, question, options_json, multi, asked_at)
|
|
VALUES (?1, ?2, ?3, ?4, ?5)",
|
|
params![asker, question, options_json, i64::from(multi), now_unix()],
|
|
)?;
|
|
Ok(conn.last_insert_rowid())
|
|
}
|
|
|
|
/// Mark the question answered. Returns the original question text so the
|
|
/// caller can include it in any helper event it fires off.
|
|
pub fn answer(&self, id: i64, answer: &str) -> Result<String> {
|
|
let conn = self.conn.lock().unwrap();
|
|
let question: Option<(String, Option<i64>)> = conn
|
|
.query_row(
|
|
"SELECT question, answered_at FROM operator_questions WHERE id = ?1",
|
|
params![id],
|
|
|row| Ok((row.get(0)?, row.get(1)?)),
|
|
)
|
|
.optional()?;
|
|
let Some((question, answered_at)) = question else {
|
|
bail!("question {id} not found");
|
|
};
|
|
if answered_at.is_some() {
|
|
bail!("question {id} already answered");
|
|
}
|
|
conn.execute(
|
|
"UPDATE operator_questions SET answer = ?1, answered_at = ?2 WHERE id = ?3",
|
|
params![answer, now_unix(), id],
|
|
)?;
|
|
Ok(question)
|
|
}
|
|
|
|
#[allow(dead_code)]
|
|
pub fn get(&self, id: i64) -> Result<Option<OpQuestion>> {
|
|
let conn = self.conn.lock().unwrap();
|
|
conn.query_row(
|
|
"SELECT id, asker, question, options_json, multi, asked_at, answered_at, answer
|
|
FROM operator_questions WHERE id = ?1",
|
|
params![id],
|
|
row_to_question,
|
|
)
|
|
.optional()
|
|
.map_err(Into::into)
|
|
}
|
|
|
|
pub fn pending(&self) -> Result<Vec<OpQuestion>> {
|
|
let conn = self.conn.lock().unwrap();
|
|
let mut stmt = conn.prepare(
|
|
"SELECT id, asker, question, options_json, multi, asked_at, answered_at, answer
|
|
FROM operator_questions
|
|
WHERE answered_at IS NULL
|
|
ORDER BY id ASC",
|
|
)?;
|
|
let rows = stmt.query_map([], row_to_question)?;
|
|
rows.collect::<rusqlite::Result<Vec<_>>>()
|
|
.map_err(Into::into)
|
|
}
|
|
}
|
|
|
|
fn row_to_question(row: &rusqlite::Row<'_>) -> rusqlite::Result<OpQuestion> {
|
|
let options_json: String = row.get(3)?;
|
|
let options: Vec<String> = serde_json::from_str(&options_json).unwrap_or_default();
|
|
let multi: i64 = row.get(4)?;
|
|
Ok(OpQuestion {
|
|
id: row.get(0)?,
|
|
asker: row.get(1)?,
|
|
question: row.get(2)?,
|
|
options,
|
|
multi: multi != 0,
|
|
asked_at: row.get(5)?,
|
|
answered_at: row.get(6)?,
|
|
answer: row.get(7)?,
|
|
})
|
|
}
|
|
|
|
fn now_unix() -> i64 {
|
|
SystemTime::now()
|
|
.duration_since(UNIX_EPOCH)
|
|
.ok()
|
|
.and_then(|d| i64::try_from(d.as_secs()).ok())
|
|
.unwrap_or(0)
|
|
}
|