// session_search.mbt — JS-only orchestration.
//
// Cross-target pure helpers (truncate_around_matches, format_conversation,
// _find_all_positions, etc.) live in search_core.mbt. This file keeps
// only the LLM summarization entry points, the ffi.Promise-based
// concurrency plumbing, and the top-level `session_search` async fn.
// ── LLM summarization env ──
///|
fn _get_summarizer_model() -> String {
@ffi.get_env("MNEMO_SUMMARIZER_MODEL")
}
///|
fn _get_summarizer_api_key() -> String {
let key = @ffi.get_env("OPENROUTER_API_KEY")
if key != "" {
key
} else {
@ffi.get_env("MNEMO_SUMMARIZER_API_KEY")
}
}
///|
let _summarize_system_prompt : String = "You summarize conversations concisely."
///|
let _summarize_prompt_prefix : String =
"Summarize this conversation in 3-5 sentences focusing on:\n1. What the user was trying to accomplish\n2. Key decisions or actions taken\n3. Any important outcomes or unresolved issues\n\n\n"
///|
let _summarize_prompt_suffix : String = "\n"
///|
fn _summarize_with_llm(conversation_text : String) -> String {
let model = _get_summarizer_model()
let api_key = _get_summarizer_api_key()
if model == "" || api_key == "" {
return ""
}
let provider = @openai.OpenAIProvider::new(
api_key,
endpoint=@openai.OpenAIEndpoint::OpenRouter,
model=model,
max_tokens=512,
system_prompt=_summarize_system_prompt,
)
let prompt = _summarize_prompt_prefix +
conversation_text +
_summarize_prompt_suffix
let messages = [@llm.Message::user(prompt)]
@llm.collect_text(provider, messages)
}
///| JS-only async variant. Returns a Promise that resolves to the LLM summary
///| (or "" when env is missing / network fails). Enables fan-out via Promise.all.
fn _summarize_with_llm_async(
conversation_text : String
) -> @ffi.Promise[String] {
let model = _get_summarizer_model()
let api_key = _get_summarizer_api_key()
if model == "" || api_key == "" {
return _resolved_string("")
}
let provider = @openai.OpenAIProvider::new(
api_key,
endpoint=@openai.OpenAIEndpoint::OpenRouter,
model=model,
max_tokens=512,
system_prompt=_summarize_system_prompt,
)
let prompt = _summarize_prompt_prefix +
conversation_text +
_summarize_prompt_suffix
let messages = [@llm.Message::user(prompt)]
provider.collect_text_async(messages)
}
///| Promise.resolve(s) — helper for sync-path early returns.
extern "js" fn _resolved_string(s : String) -> @ffi.Promise[String] =
#| (s) => Promise.resolve(s)
///| Promise.all for an array of string-producing promises. Resolves to the
///| array of resolved values in input order. Rejects if any promise rejects
///| — caller should ensure the per-request path already swallows errors into "".
extern "js" fn _promise_all_strings(
promises : Array[@ffi.Promise[String]]
) -> @ffi.Promise[Array[String]] =
#| (promises) => Promise.all(promises)
// ── DB helpers ──
///| Return metadata (id, source, model, started_at, parent_session_id)
///| for a single session as positional strings. Empty array on miss.
fn _sess_get_session_meta(
db : SqliteDb,
session_id : String
) -> Array[String] {
let stmt = prepare(
db,
"SELECT id, source, model, started_at, parent_session_id FROM sessions WHERE id = ?",
)
_stmt_bind_text(stmt, 1, session_id)
let rows = _stmt_all(stmt)
if _rows_length(rows) == 0 {
return []
}
let row = _row_at(rows, 0)
let source = if _row_is_null(row, "source") {
"unknown"
} else {
_row_get_text(row, "source")
}
[
_row_get_text(row, "id"),
source,
_row_get_text(row, "model"),
_row_get_int64(row, "started_at").to_string(),
_row_get_text(row, "parent_session_id"),
]
}
///|
fn srch_resolve_to_root(sdb : SessionDb, session_id : String) -> String {
let mut current = session_id
let visited : Array[String] = []
let mut limit = 0
while limit < 100 {
let mut already = false
for v in visited {
if v == current {
already = true
break
}
}
if already {
break
}
visited.push(current)
let meta = _sess_get_session_meta(sdb.db, current)
if meta.length() == 0 {
break
}
let parent = meta[4]
if parent == "" {
break
}
current = parent
limit = limit + 1
}
current
}
// ── concurrency knob ──
///|
fn _parse_concurrency_js(s : String) -> Int {
(try? @string.parse_int(s.view())).unwrap_or(-1)
}
///| Max concurrent LLM summarizations in flight. Parity with proto's
///| MNEMO_SEARCH_CONCURRENCY default (3, clamped to 5).
fn _get_search_concurrency() -> Int {
let raw = @ffi.get_env("MNEMO_SEARCH_CONCURRENCY")
if raw == "" {
return 3
}
let n = _parse_concurrency_js(raw)
if n < 1 {
1
} else if n > 5 {
5
} else {
n
}
}
// ── session_search orchestration ──
///| Carrier for per-session work that must survive the parallel await step.
priv struct SearchWork {
session_id : String
conversation_text : String
meta : Array[String]
}
///|
pub async fn session_search(
sdb : SessionDb,
query : String,
limit~ : Int = 3,
role_filter~ : String? = None
) -> SessionSearchResponse {
let fts_limit = srch_max(1, srch_min(limit, 5))
let q = query.trim().to_string()
let raw_results = search_messages(sdb, q, limit=50, role_filter~)
if raw_results.length() == 0 {
return { ok: true, query: q, results: [], count: 0, sessions_searched: 0 }
}
// Group by resolved root session, deduplicate, up to fts_limit
let seen_ids : Array[String] = []
let mut ri = 0
while ri < raw_results.length() && seen_ids.length() < fts_limit {
let result = raw_results[ri]
let resolved = srch_resolve_to_root(sdb, result.session_id)
let mut already = false
for s in seen_ids {
if s == resolved {
already = true
break
}
}
if !already {
seen_ids.push(resolved)
}
ri = ri + 1
}
// Gather per-session work (sync DB reads).
let work : Array[SearchWork] = []
for sid in seen_ids {
let messages = get_messages(sdb, sid)
if messages.length() > 0 {
let meta = _sess_get_session_meta(sdb.db, sid)
let fmt_msgs : Array[MsgForFormat] = messages.map(fn(m) {
{ role: m.role, content: m.content, tool_name: m.tool_name, tool_calls: "" }
})
let conversation_text = truncate_around_matches(
format_conversation(fmt_msgs), q,
)
work.push({ session_id: sid, conversation_text, meta })
}
}
// Fan out LLM summarization in chunks (parity with proto Semaphore).
let concurrency = _get_search_concurrency()
let summaries : Array[String] = []
let mut wi = 0
while wi < work.length() {
let end = srch_min(wi + concurrency, work.length())
let chunk : Array[@ffi.Promise[String]] = []
for j in wi.. 0 {
_format_timestamp(w.meta[3])
} else {
"unknown"
}
let source = if w.meta.length() > 0 { w.meta[1] } else { "unknown" }
let model = if w.meta.length() > 0 { w.meta[2] } else { "" }
results.push({ session_id: w.session_id, when, source, model, summary })
}
{ ok: true, query: q, results, count: results.length(), sessions_searched: seen_ids.length() }
}
// ── inline tests ──
///|
async test "session_search: no results when db empty" {
let dir = _mkdtemp_sess()
let sdb = open_session_db(_path_join_2(dir, "s.db"))
let resp = session_search(sdb, "anything")
assert_eq(resp.ok, true)
assert_eq(resp.count, 0)
assert_eq(resp.sessions_searched, 0)
_rmrf_sess(dir)
}
///|
async test "session_search: finds session with matching message" {
let dir = _mkdtemp_sess()
let sdb = open_session_db(_path_join_2(dir, "s.db"))
let sid = create_session(sdb, "cli", "test system")
let _ = append_message(
sdb, sid, "user", 1000L, content=Some("unique_srch_term_xyz"),
)
let resp = session_search(sdb, "unique_srch_term_xyz")
assert_eq(resp.ok, true)
assert_eq(resp.count > 0, true)
_rmrf_sess(dir)
}
///|
test "summarize falls back when env unset" {
let model = @ffi.get_env("MNEMO_SUMMARIZER_MODEL")
let api_key = @ffi.get_env("OPENROUTER_API_KEY")
if model == "" || api_key == "" {
let result = _summarize_with_llm("test transcript")
assert_eq(result, "")
}
}
///|
test "summarize uses LLM when env set" {
let model = @ffi.get_env("MNEMO_SUMMARIZER_MODEL")
let api_key = @ffi.get_env("OPENROUTER_API_KEY")
if model == "" || api_key == "" {
return
}
let transcript = "User: hello.\nAssistant: hi!"
let summary = _summarize_with_llm(transcript)
assert_eq(summary.length() > 0, true)
}