mirror of
https://github.com/allaunthefox/Research-Stack.git
synced 2026-07-31 03:05:21 +00:00
- source.rs: Add ClawSource adapter that reads .claw/sessions/*.jsonl (skips
LFS stubs), parses both full-JSON and JSONL-stream formats, normalises into
ChatSession/ChatMessage identical to the OpenCode path.
- models.rs: Add title/agent/model fields to ChatSession; remove unused
chrono import.
- normalize.rs: Populate title/agent/model from OpenCodeSession.
- sink.rs: Full rewrite:
- DDL adds title/agent/model columns + ALTER TABLE IF NOT EXISTS migration
path for existing clusters.
- upsert_session and upsert_messages parameter counts match placeholders
($19 and $13 respectively).
- FTS index uses COALESCE to handle NULL text_content.
- sslmode= token stripped before parse so NoTls connect doesn't reject
it; README documents the SSL upgrade path (tokio-postgres-rustls).
- with-serde_json-1 feature added to tokio-postgres dep for JSONB params.
- main.rs:
- Replace DefaultHasher sha256_text with FNV-1a 64-bit (correct hex digest).
- Add ClawSync, List, Get subcommands.
- Remove unused Context import.
- bridge.rs: Rewrite to use bridge_wrapper.py subprocess rather than inline
-c script; resolves wrapper path relative to infra_dir or CARGO_MANIFEST_DIR;
unwraps {"ok":true,"data":{}} envelope; health_check() tests python3 only.
- embed.rs: Remove unused OllamaEmbedRequest import.
- systemd/: Add ene-session-sync.service + .timer (5-min cadence, 2-min
OnBootSec, Persistent=true for catch-up after suspend).
- README.md: Architecture diagram, CLI reference, schema, build/install,
systemd install, Python bridge setup.
Generated with [Devin](https://cli.devin.ai/docs)
Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
108 lines
3.8 KiB
Rust
108 lines
3.8 KiB
Rust
use anyhow::{Context, Result};
|
|
use reqwest::Client;
|
|
use serde_json::json;
|
|
use std::time::Duration;
|
|
use tracing::{debug, info, warn};
|
|
|
|
/// Ollama embedding client.
|
|
pub struct Embedder {
|
|
client: Client,
|
|
base_url: String,
|
|
model: String,
|
|
}
|
|
|
|
impl Embedder {
|
|
pub fn new(base_url: &str, model: &str) -> Self {
|
|
let client = Client::builder()
|
|
.timeout(Duration::from_secs(120))
|
|
.build()
|
|
.unwrap_or_else(|_| Client::new());
|
|
Self {
|
|
client,
|
|
base_url: base_url.trim_end_matches('/').to_string(),
|
|
model: model.to_string(),
|
|
}
|
|
}
|
|
|
|
pub fn from_env() -> Self {
|
|
let base = std::env::var("OLLAMA_HOST").unwrap_or_else(|_| "http://localhost:11434".into());
|
|
let model = std::env::var("OLLAMA_EMBED_MODEL").unwrap_or_else(|_| "nomic-embed-text".into());
|
|
Self::new(&base, &model)
|
|
}
|
|
|
|
/// Embed a single string, returning a 768d (or model-specific) vector.
|
|
pub async fn embed(&self, text: &str) -> Result<Vec<f32>> {
|
|
let url = format!("{}/api/embeddings", self.base_url);
|
|
let payload = json!({
|
|
"model": self.model,
|
|
"prompt": text,
|
|
});
|
|
let resp = self
|
|
.client
|
|
.post(&url)
|
|
.json(&payload)
|
|
.send()
|
|
.await
|
|
.with_context(|| format!("POST {}", url))?;
|
|
if !resp.status().is_success() {
|
|
let status = resp.status();
|
|
let body = resp.text().await.unwrap_or_default();
|
|
anyhow::bail!("Ollama returned {}: {}", status, body);
|
|
}
|
|
let json: serde_json::Value = resp.json().await.context("parse Ollama response")?;
|
|
let embedding = json
|
|
.get("embedding")
|
|
.and_then(|v| v.as_array())
|
|
.context("missing 'embedding' field in Ollama response")?;
|
|
let vec: Vec<f32> = embedding
|
|
.iter()
|
|
.map(|v| v.as_f64().unwrap_or(0.0) as f32)
|
|
.collect();
|
|
debug!("embedded {} chars -> {} dims", text.chars().count(), vec.len());
|
|
Ok(vec)
|
|
}
|
|
|
|
/// Embed multiple strings sequentially (Ollama does not batch natively).
|
|
pub async fn embed_batch(&self, texts: &[String]) -> Result<Vec<Vec<f32>>> {
|
|
let mut out = Vec::with_capacity(texts.len());
|
|
for (i, text) in texts.iter().enumerate() {
|
|
match self.embed(text).await {
|
|
Ok(v) => out.push(v),
|
|
Err(e) => {
|
|
warn!("embedding failed for item {}: {}", i, e);
|
|
out.push(Vec::new());
|
|
}
|
|
}
|
|
}
|
|
info!("embedded {} texts via Ollama {}", texts.len(), self.model);
|
|
Ok(out)
|
|
}
|
|
|
|
/// Check if Ollama is reachable and the model is loaded.
|
|
pub async fn health_check(&self) -> Result<bool> {
|
|
let url = format!("{}/api/tags", self.base_url);
|
|
match self.client.get(&url).send().await {
|
|
Ok(resp) => {
|
|
if !resp.status().is_success() {
|
|
return Ok(false);
|
|
}
|
|
let json: serde_json::Value = resp.json().await.unwrap_or_default();
|
|
let models = json.get("models").and_then(|m| m.as_array());
|
|
if let Some(arr) = models {
|
|
Ok(arr.iter().any(|m| {
|
|
m.get("name")
|
|
.and_then(|n| n.as_str())
|
|
.map(|n| n == self.model || n.starts_with(&format!("{}:", self.model)))
|
|
.unwrap_or(false)
|
|
}))
|
|
} else {
|
|
Ok(false)
|
|
}
|
|
}
|
|
Err(e) => {
|
|
warn!("Ollama health check failed: {}", e);
|
|
Ok(false)
|
|
}
|
|
}
|
|
}
|
|
}
|