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
21 changes: 21 additions & 0 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,27 @@ jobs:
${{ matrix.name }}.tar.gz
${{ matrix.name }}.tar.gz.sha256

checksums:
needs: build
runs-on: ubuntu-latest
steps:
- name: Download SHA256 files
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
run: |
gh release download ${{ github.ref_name }} \
--repo ${{ github.repository }} \
--pattern "*.sha256"

- name: Generate checksums.txt
run: |
cat *.sha256 > checksums.txt

- name: Upload checksums.txt
uses: softprops/action-gh-release@v1
with:
files: checksums.txt

publish-crate:
needs: build
runs-on: ubuntu-latest
Expand Down
16 changes: 9 additions & 7 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -113,13 +113,15 @@ CREATE TABLE graph_nodes (

-- Graph edges (RELATED, SUPERSEDES, CO_ACCESSED, CLUSTER_SIBLING)
CREATE TABLE graph_edges (
source_id TEXT, target_id TEXT,
rel_type TEXT, -- 'related', 'supersedes', 'co_accessed', 'cluster_sibling'
strength REAL,
count INTEGER,
label TEXT,
created_at INTEGER,
UNIQUE(source_id, target_id, rel_type)
from_id TEXT NOT NULL,
to_id TEXT NOT NULL,
edge_type TEXT NOT NULL, -- 'related', 'supersedes', 'co_accessed', 'cluster_sibling'
strength REAL NOT NULL DEFAULT 0.0,
rel_type TEXT NOT NULL DEFAULT '',
count INTEGER NOT NULL DEFAULT 0,
last_at INTEGER NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL,
UNIQUE(from_id, to_id, edge_type)
);

-- Access tracking and decay scoring
Expand Down
6 changes: 2 additions & 4 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -38,10 +38,7 @@ serde = { version = "1.0", features = ["derive"] }
serde_json = "1.0"

# Async
tokio = { version = "1", features = ["full"] }

# MCP
jsonrpc-core = "18"
tokio = { version = "1", features = ["rt", "rt-multi-thread", "sync", "time", "macros"] }

# CLI
clap = { version = "4", features = ["derive", "env"] }
Expand All @@ -53,6 +50,7 @@ rusqlite = { version = "0.31", features = ["bundled"] }
# Utils
uuid = { version = "1", features = ["v4", "serde"] }
chrono = { version = "0.4", features = ["serde"] }
sha2 = "0.10"
thiserror = "2"
anyhow = "1"
tracing = "0.1"
Expand Down
27 changes: 27 additions & 0 deletions install.sh
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,33 @@ TMPDIR="$(mktemp -d)"
trap 'rm -rf "$TMPDIR"' EXIT

curl -fsSL "$URL" -o "${TMPDIR}/${TARBALL}"

# Verify SHA256 checksum if checksums.txt is available
CHECKSUMS_URL="https://github.com/${REPO}/releases/download/v${VERSION}/checksums.txt"
if curl -fsSL "$CHECKSUMS_URL" -o "${TMPDIR}/checksums.txt" 2>/dev/null; then
EXPECTED="$(grep "${TARBALL}" "${TMPDIR}/checksums.txt" | awk '{print $1}')"
if [ -n "$EXPECTED" ]; then
if command -v sha256sum &>/dev/null; then
ACTUAL="$(sha256sum "${TMPDIR}/${TARBALL}" | awk '{print $1}')"
elif command -v shasum &>/dev/null; then
ACTUAL="$(shasum -a 256 "${TMPDIR}/${TARBALL}" | awk '{print $1}')"
else
echo "Warning: no sha256sum or shasum found, skipping checksum verification" >&2
ACTUAL=""
fi
if [ -n "$ACTUAL" ] && [ "$ACTUAL" != "$EXPECTED" ]; then
echo "Checksum verification failed!" >&2
echo " Expected: $EXPECTED" >&2
echo " Actual: $ACTUAL" >&2
exit 1
elif [ -n "$ACTUAL" ]; then
echo "Checksum verified."
fi
fi
else
echo "No checksums.txt available, skipping verification."
fi

tar -xzf "${TMPDIR}/${TARBALL}" -C "$TMPDIR"

# Install
Expand Down
2 changes: 2 additions & 0 deletions src/access.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ impl AccessTracker {
SedimentError::Database(format!("Failed to open access database: {}", e))
})?;

conn.execute_batch("PRAGMA journal_mode=WAL;").ok();

conn.execute_batch(
"CREATE TABLE IF NOT EXISTS access_log (
item_id TEXT PRIMARY KEY,
Expand Down
48 changes: 24 additions & 24 deletions src/chunker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -111,40 +111,40 @@ pub fn chunk_content(
/// Split text at sentence boundaries
fn split_at_sentences(text: &str) -> Vec<&str> {
let mut sentences = Vec::new();
let bytes = text.as_bytes();
let mut start = 0;
let mut i = 0;

while i < bytes.len() {
// Look for sentence-ending punctuation followed by space or end
if matches!(bytes[i], b'.' | b'?' | b'!') {
let next_idx = i + 1;
if next_idx >= bytes.len()
|| bytes[next_idx] == b' '
|| bytes[next_idx] == b'\n'
|| bytes[next_idx] == b'\t'
{
// Include the punctuation
let end = i + 1;
if start < end && end <= bytes.len() {
let mut char_indices = text.char_indices().peekable();

while let Some((i, ch)) = char_indices.next() {
// Look for sentence-ending punctuation (ASCII and Unicode)
if matches!(ch, '.' | '?' | '!' | '。' | '?' | '!') {
let end = i + ch.len_utf8();
// Check if followed by whitespace or end of text
let at_end_or_ws = match char_indices.peek() {
None => true,
Some(&(_, next_ch)) => next_ch == ' ' || next_ch == '\n' || next_ch == '\t',
};
if at_end_or_ws {
if start < end {
sentences.push(&text[start..end]);
}
// Skip whitespace after punctuation
i += 1;
while i < bytes.len()
&& (bytes[i] == b' ' || bytes[i] == b'\n' || bytes[i] == b'\t')
{
i += 1;
while let Some(&(_, next_ch)) = char_indices.peek() {
if next_ch == ' ' || next_ch == '\n' || next_ch == '\t' {
char_indices.next();
} else {
break;
}
}
start = i;
continue;
start = match char_indices.peek() {
Some(&(idx, _)) => idx,
None => text.len(),
};
}
}
i += 1;
}

// Add remaining text
if start < bytes.len() {
if start < text.len() {
sentences.push(&text[start..]);
}

Expand Down
42 changes: 32 additions & 10 deletions src/consolidation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ impl ConsolidationQueue {
SedimentError::Database(format!("Failed to open consolidation database: {}", e))
})?;

conn.execute_batch("PRAGMA journal_mode=WAL;").ok();

conn.execute_batch(
"CREATE TABLE IF NOT EXISTS consolidation_queue (
item_id_a TEXT NOT NULL,
Expand Down Expand Up @@ -180,7 +182,11 @@ async fn run_consolidation_batch(
let result = process_candidate(&mut db, access_db_path, candidate).await;
match result {
Ok(status) => {
let _ = queue.mark_processed(&candidate.item_id_a, &candidate.item_id_b, &status);
if let Err(e) =
queue.mark_processed(&candidate.item_id_a, &candidate.item_id_b, &status)
{
tracing::warn!("mark_processed failed: {}", e);
}
info!(
"Consolidated {} <-> {}: {} (similarity: {:.2})",
candidate.item_id_a, candidate.item_id_b, status, candidate.similarity
Expand Down Expand Up @@ -220,26 +226,38 @@ async fn process_candidate(
};

// Transfer edges from old to new
let _ = graph.transfer_edges(&remove.id, &keep.id);
if let Err(e) = graph.transfer_edges(&remove.id, &keep.id) {
tracing::warn!("transfer_edges failed: {}", e);
}

// Create SUPERSEDES edge (preserves lineage for recovery)
let _ = graph.add_supersedes_edge(&keep.id, &remove.id);
if let Err(e) = graph.add_supersedes_edge(&keep.id, &remove.id) {
tracing::warn!("add_supersedes_edge failed: {}", e);
}

// Archive the removed item's content into a RELATED edge label
// so it can be recovered if this was a false positive merge.
// The edge label stores a truncated snapshot; full content is
// preserved in the SUPERSEDES relationship for audit.
let archive_preview = if remove.content.len() > 500 {
format!("{}...", &remove.content[..497])
let archive_preview = if remove.content.chars().count() > 500 {
let cut = remove
.content
.char_indices()
.nth(497)
.map(|(i, _)| i)
.unwrap_or(remove.content.len());
format!("{}...", &remove.content[..cut])
} else {
remove.content.clone()
};
let _ = graph.add_related_edge(
if let Err(e) = graph.add_related_edge(
&keep.id,
&remove.id,
candidate.similarity,
&format!("merged_archive:{}", archive_preview),
);
) {
tracing::warn!("add_related_edge failed: {}", e);
}

// Soft-delete: mark item as expired instead of hard-deleting.
// This allows recovery; expired items are excluded from search
Expand All @@ -249,7 +267,9 @@ async fn process_candidate(
// Fall back to hard delete if expire is not supported
warn!("expire_item failed ({}), falling back to delete", e);
db.delete_item(&remove.id).await?;
let _ = graph.remove_node(&remove.id);
if let Err(e) = graph.remove_node(&remove.id) {
tracing::warn!("remove_node failed: {}", e);
}
}

Ok("merged".to_string())
Expand All @@ -261,12 +281,14 @@ async fn process_candidate(
}
} else {
// Similar but distinct: create RELATED edge
let _ = graph.add_related_edge(
if let Err(e) = graph.add_related_edge(
&candidate.item_id_a,
&candidate.item_id_b,
candidate.similarity,
"similar",
);
) {
tracing::warn!("add_related_edge failed: {}", e);
}
Ok("linked".to_string())
}
}
72 changes: 50 additions & 22 deletions src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ use lancedb::connect;
use lancedb::query::{ExecutableQuery, QueryBase};
use tracing::{debug, info};

use crate::boost_similarity;
use crate::chunker::{ChunkingConfig, chunk_content};
use crate::document::ContentType;
use crate::embedder::{EMBEDDING_DIM, Embedder};
Expand Down Expand Up @@ -288,16 +289,6 @@ impl Database {
item.project_id = self.project_id.clone();
}

// Check for potential conflicts before storing
let potential_conflicts = self
.find_similar_items(
&item.content,
CONFLICT_SIMILARITY_THRESHOLD,
CONFLICT_SEARCH_LIMIT,
)
.await
.unwrap_or_default();

// Determine if we need to chunk
let should_chunk = item.content.len() > CHUNK_THRESHOLD;
item.is_chunked = should_chunk;
Expand Down Expand Up @@ -360,6 +351,19 @@ impl Database {
debug!("Stored item: {} (no chunking)", item.id);
}

// Detect conflicts after storing (informational only, avoids TOCTOU race)
let potential_conflicts = self
.find_similar_items(
&item.content,
CONFLICT_SIMILARITY_THRESHOLD,
CONFLICT_SEARCH_LIMIT,
)
.await
.unwrap_or_default()
.into_iter()
.filter(|c| c.id != item.id)
.collect();

Ok(StoreResult {
id: item.id,
potential_conflicts,
Expand Down Expand Up @@ -751,7 +755,7 @@ impl Database {

// Delete then re-insert with new expiration
table
.delete(&format!("id = '{}'", id))
.delete(&format!("id = '{}'", sanitize_sql_string(id)))
.await
.map_err(|e| SedimentError::Database(format!("Delete for expire failed: {}", e)))?;

Expand Down Expand Up @@ -810,6 +814,31 @@ impl Database {

Ok(stats)
}

/// Delete items whose expires_at timestamp is in the past.
pub async fn cleanup_expired(&self) -> Result<usize> {
let table = match &self.items_table {
Some(t) => t,
None => return Ok(0),
};

let now = Utc::now().timestamp();
let filter = format!("expires_at IS NOT NULL AND expires_at < {}", now);

// Count how many will be deleted
let count = table.count_rows(Some(filter.clone())).await.unwrap_or(0);

if count > 0 {
table
.delete(&filter)
.await
.map_err(|e| SedimentError::Database(format!("Expired cleanup failed: {}", e)))?;

info!("Cleaned up {} expired items", count);
}

Ok(count)
}
}

// ==================== Decay Scoring ====================
Expand Down Expand Up @@ -841,15 +870,6 @@ pub fn score_with_decay(

// ==================== Helper Functions ====================

/// Apply similarity boosting based on project context.
fn boost_similarity(base: f32, item_project: Option<&str>, current_project: Option<&str>) -> f32 {
match (item_project, current_project) {
(Some(m), Some(c)) if m == c => (base * 1.15).min(1.0), // Same project: boost
(Some(_), Some(_)) => base * 0.95, // Different project: slight penalty
_ => base, // Global or no context
}
}

/// Detect content type for smart chunking
fn detect_content_type(content: &str) -> ContentType {
let trimmed = content.trim();
Expand Down Expand Up @@ -1037,12 +1057,20 @@ fn batch_to_items(batch: &RecordBatch) -> Result<Vec<Item>> {
if c.is_null(i) {
None
} else {
Some(Utc.timestamp_opt(c.value(i), 0).unwrap())
Some(
Utc.timestamp_opt(c.value(i), 0)
.single()
.unwrap_or_else(Utc::now),
)
}
});

let created_at = created_at_col
.map(|c| Utc.timestamp_opt(c.value(i), 0).unwrap())
.map(|c| {
Utc.timestamp_opt(c.value(i), 0)
.single()
.unwrap_or_else(Utc::now)
})
.unwrap_or_else(Utc::now);

let item = Item {
Expand Down
Loading