Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
191 changes: 127 additions & 64 deletions crates/ctx-history-capture/src/provider/providers/opencode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ use crate::provider::native::{
provider_required_timestamp_millis, provider_role, provider_value_text,
};
use crate::provider::providers::real_content::text_has_real_content;
use crate::provider::sqlite::{sqlite_is_too_big, sqlite_row_ids_with_oversized_value};
use crate::{
CaptureError, ProviderAdapterContext, ProviderNormalizationResult, Result,
KILO_SQLITE_SOURCE_FORMAT, OPENCODE_SQLITE_SOURCE_FORMAT, PROVIDER_MAX_PREVIEW_CHARS,
Expand Down Expand Up @@ -66,6 +67,7 @@ pub(crate) struct OpenCodeMessageSelection {
pub(crate) rows: Vec<OpenCodeMessageRow>,
pub(crate) source_table: Option<&'static str>,
pub(crate) skipped_non_conversational_rows: usize,
pub(crate) skipped_oversized_values: usize,
}

pub(crate) fn normalize_opencode_sqlite(
Expand All @@ -83,17 +85,23 @@ pub(crate) fn normalize_opencode_sqlite(
.iter()
.map(|session| session.id.clone())
.collect::<BTreeSet<_>>();
let message_selection = opencode_session_messages(&conn, dialect, &session_ids)?;
let message_selection = opencode_session_messages(path, &conn, dialect, &session_ids)?;
let mut result = ProviderNormalizationResult::default();
result.summary.skipped += message_selection.skipped_oversized_values;
result.summary.skipped_events += message_selection.skipped_oversized_values;
if message_selection.rows.is_empty() {
push_provider_import_failure(
&mut result.summary,
0,
format!(
"{} SQLite database contained no real conversational message rows",
dialect.display_name
),
);
if message_selection.skipped_oversized_values == 0
|| message_selection.skipped_non_conversational_rows > 0
{
push_provider_import_failure(
&mut result.summary,
0,
format!(
"{} SQLite database contained no real conversational message rows",
dialect.display_name
),
);
}
return Ok(result);
}
let mut session_started = BTreeMap::new();
Expand Down Expand Up @@ -343,61 +351,72 @@ pub(crate) fn opencode_sessions(
}

pub(crate) fn opencode_session_messages(
path: &Path,
conn: &Connection,
dialect: &OpenCodeSqliteDialect,
session_ids: &BTreeSet<String>,
) -> Result<OpenCodeMessageSelection> {
let mut skipped_non_conversational_rows = 0usize;
let mut skipped_oversized_values = 0usize;
if sqlite_table_exists(conn, "session_message")? {
let rows = opencode_session_message_rows(conn, dialect)?;
let (rows, oversized) = opencode_session_message_rows(path, conn, dialect)?;
skipped_oversized_values += oversized;
if opencode_rows_have_import_blocking_errors(&rows, session_ids, dialect) {
return Ok(OpenCodeMessageSelection {
rows,
source_table: Some("session_message"),
skipped_non_conversational_rows,
skipped_oversized_values,
});
}
if opencode_rows_have_real_message_content(&rows) {
return Ok(OpenCodeMessageSelection {
rows,
source_table: Some("session_message"),
skipped_non_conversational_rows,
skipped_oversized_values,
});
}
skipped_non_conversational_rows += rows.len();
}
if sqlite_table_exists(conn, "session_entry")? {
let rows = opencode_session_entry_rows(conn, dialect)?;
let (rows, oversized) = opencode_session_entry_rows(path, conn, dialect)?;
skipped_oversized_values += oversized;
if opencode_rows_have_import_blocking_errors(&rows, session_ids, dialect) {
return Ok(OpenCodeMessageSelection {
rows,
source_table: Some("session_entry"),
skipped_non_conversational_rows,
skipped_oversized_values,
});
}
if opencode_rows_have_real_message_content(&rows) {
return Ok(OpenCodeMessageSelection {
rows,
source_table: Some("session_entry"),
skipped_non_conversational_rows,
skipped_oversized_values,
});
}
skipped_non_conversational_rows += rows.len();
}
if sqlite_table_exists(conn, "message")? {
let rows = opencode_message_rows(conn, dialect)?;
let (rows, oversized) = opencode_message_rows(path, conn, dialect)?;
skipped_oversized_values += oversized;
if opencode_rows_have_import_blocking_errors(&rows, session_ids, dialect) {
return Ok(OpenCodeMessageSelection {
rows,
source_table: Some("message"),
skipped_non_conversational_rows,
skipped_oversized_values,
});
}
if opencode_rows_have_real_message_content(&rows) {
return Ok(OpenCodeMessageSelection {
rows,
source_table: Some("message"),
skipped_non_conversational_rows,
skipped_oversized_values,
});
}
skipped_non_conversational_rows += rows.len();
Expand All @@ -406,13 +425,31 @@ pub(crate) fn opencode_session_messages(
rows: Vec::new(),
source_table: None,
skipped_non_conversational_rows,
skipped_oversized_values,
})
}

fn exclude_ids_clause(ids: &BTreeSet<String>) -> String {
if ids.is_empty() {
String::new()
} else {
let placeholders = std::iter::repeat("?")
.take(ids.len())
.collect::<Vec<_>>()
.join(", ");
format!("where id not in ({placeholders})")
}
}

fn oversized_sqlite_data_row_ids(path: &Path, table: &str) -> Result<BTreeSet<String>> {
sqlite_row_ids_with_oversized_value(path, table, "id", "data")
}

pub(crate) fn opencode_session_message_rows(
path: &Path,
conn: &Connection,
dialect: &OpenCodeSqliteDialect,
) -> Result<Vec<OpenCodeMessageRow>> {
) -> Result<(Vec<OpenCodeMessageRow>, usize)> {
let columns = sqlite_table_columns(conn, "session_message")?;
ensure_sqlite_table_columns(
&columns,
Expand All @@ -429,30 +466,36 @@ pub(crate) fn opencode_session_message_rows(
} else {
("NULL", "id")
};
let oversized_ids = oversized_sqlite_data_row_ids(path, "session_message")?;
let mut skipped_oversized = oversized_ids.len();
let exclude_clause = exclude_ids_clause(&oversized_ids);
let sql = format!(
"select id, session_id, {entry_type}, {seq_expr}, {time_created}, {time_updated}, data \
from session_message order by session_id, {order_expr}"
from session_message \
{exclude_clause} \
order by session_id, {order_expr}",
);
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, Option<i64>>(3)?,
row.get::<_, i64>(4)?,
row.get::<_, i64>(5)?,
row.get::<_, String>(6)?,
))
})?;
let rows = rows
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(CaptureError::from)?;
let mut rows = stmt.query(rusqlite::params_from_iter(oversized_ids.iter()))?;
let mut messages = Vec::new();
let mut next_seq_by_session = BTreeMap::<String, i64>::new();
for (id, session_id, entry_type, seq, time_created, time_updated, data) in rows {
while let Some(row) = rows.next()? {
let id = row.get::<_, String>(0)?;
let data: String = match row.get::<_, String>(6) {
Ok(value) => value,
Err(err) if sqlite_is_too_big(&err) => {
skipped_oversized += 1;
continue;
}
Err(err) => return Err(CaptureError::from(err)),
};
let session_id = row.get::<_, String>(1)?;
let entry_type_raw = row.get::<_, String>(2)?;
let seq = row.get::<_, Option<i64>>(3)?;
let time_created = row.get::<_, i64>(4)?;
let time_updated = row.get::<_, i64>(5)?;
let seq = seq.unwrap_or_else(|| next_opencode_seq(&mut next_seq_by_session, &session_id));
let entry_type = opencode_entry_type_from_data(&entry_type, &data);
let entry_type = opencode_entry_type_from_data(&entry_type_raw, &data);
messages.push(OpenCodeMessageRow {
id,
session_id,
Expand All @@ -463,7 +506,7 @@ pub(crate) fn opencode_session_message_rows(
data,
});
}
Ok(messages)
Ok((messages, skipped_oversized))
}

pub(crate) fn opencode_entry_type_from_data(fallback: &str, data: &str) -> String {
Expand All @@ -477,9 +520,10 @@ pub(crate) fn opencode_entry_type_from_data(fallback: &str, data: &str) -> Strin
}

pub(crate) fn opencode_session_entry_rows(
path: &Path,
conn: &Connection,
dialect: &OpenCodeSqliteDialect,
) -> Result<Vec<OpenCodeMessageRow>> {
) -> Result<(Vec<OpenCodeMessageRow>, usize)> {
let columns = sqlite_table_columns(conn, "session_entry")?;
ensure_sqlite_table_columns(
&columns,
Expand All @@ -493,24 +537,33 @@ pub(crate) fn opencode_session_entry_rows(
"data",
],
)?;
let mut stmt = conn.prepare(
let oversized_ids = oversized_sqlite_data_row_ids(path, "session_entry")?;
let mut skipped_oversized = oversized_ids.len();
let exclude_clause = exclude_ids_clause(&oversized_ids);
let sql = format!(
"select id, session_id, type, time_created, time_updated, data \
from session_entry order by session_id, time_created, id",
)?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, i64>(3)?,
row.get::<_, i64>(4)?,
row.get::<_, String>(5)?,
))
})?;
from session_entry \
{exclude_clause} \
order by session_id, time_created, id",
);
let mut stmt = conn.prepare(&sql)?;
let mut rows = stmt.query(rusqlite::params_from_iter(oversized_ids.iter()))?;
let mut messages = Vec::new();
let mut next_seq_by_session = BTreeMap::<String, i64>::new();
for row in rows {
let (id, session_id, entry_type, time_created, time_updated, data) = row?;
while let Some(row) = rows.next()? {
let id = row.get::<_, String>(0)?;
let data: String = match row.get::<_, String>(5) {
Ok(value) => value,
Err(err) if sqlite_is_too_big(&err) => {
skipped_oversized += 1;
continue;
}
Err(err) => return Err(CaptureError::from(err)),
};
let session_id = row.get::<_, String>(1)?;
let entry_type = row.get::<_, String>(2)?;
let time_created = row.get::<_, i64>(3)?;
let time_updated = row.get::<_, i64>(4)?;
let seq = next_opencode_seq(&mut next_seq_by_session, &session_id);
messages.push(OpenCodeMessageRow {
id,
Expand All @@ -522,36 +575,46 @@ pub(crate) fn opencode_session_entry_rows(
data,
});
}
Ok(messages)
Ok((messages, skipped_oversized))
}

pub(crate) fn opencode_message_rows(
path: &Path,
conn: &Connection,
dialect: &OpenCodeSqliteDialect,
) -> Result<Vec<OpenCodeMessageRow>> {
) -> Result<(Vec<OpenCodeMessageRow>, usize)> {
let columns = sqlite_table_columns(conn, "message")?;
ensure_sqlite_table_columns(
&columns,
&format!("{} SQLite message table", dialect.display_name),
&["id", "session_id", "time_created", "time_updated", "data"],
)?;
let mut stmt = conn.prepare(
let oversized_ids = oversized_sqlite_data_row_ids(path, "message")?;
let mut skipped_oversized = oversized_ids.len();
let exclude_clause = exclude_ids_clause(&oversized_ids);
let sql = format!(
"select id, session_id, time_created, time_updated, data \
from message order by session_id, time_created, id",
)?;
let rows = stmt.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, i64>(3)?,
row.get::<_, String>(4)?,
))
})?;
from message \
{exclude_clause} \
order by session_id, time_created, id",
);
let mut stmt = conn.prepare(&sql)?;
let mut rows = stmt.query(rusqlite::params_from_iter(oversized_ids.iter()))?;
let mut messages = Vec::new();
let mut next_seq_by_session = BTreeMap::<String, i64>::new();
for row in rows {
let (id, session_id, time_created, time_updated, data) = row?;
while let Some(row) = rows.next()? {
let id = row.get::<_, String>(0)?;
let data: String = match row.get::<_, String>(4) {
Ok(value) => value,
Err(err) if sqlite_is_too_big(&err) => {
skipped_oversized += 1;
continue;
}
Err(err) => return Err(CaptureError::from(err)),
};
let session_id = row.get::<_, String>(1)?;
let time_created = row.get::<_, i64>(2)?;
let time_updated = row.get::<_, i64>(3)?;
let seq = next_opencode_seq(&mut next_seq_by_session, &session_id);
let entry_type = serde_json::from_str::<Value>(&data)
.ok()
Expand All @@ -567,7 +630,7 @@ pub(crate) fn opencode_message_rows(
data,
});
}
Ok(messages)
Ok((messages, skipped_oversized))
}

pub(crate) fn next_opencode_seq(
Expand Down
Loading
Loading