mirror of
https://github.com/tinyhumansai/openhuman.git
synced 2026-07-27 21:08:00 +00:00
This commit is contained in:
@@ -0,0 +1,168 @@
|
||||
/**
|
||||
* Vitest for `<SyncAuditPanel />` (issue #3116 — coverage for the sync-audit
|
||||
* history surface shipped in PR #3113).
|
||||
*
|
||||
* Covers:
|
||||
* - loading → loaded transition (renders rows once the audit log resolves)
|
||||
* - formatting of tokens (k/M), cost ($x.xxxx), and duration (ms / s / m s)
|
||||
* - the scope label mapping for github / gmail / rebuild scopes
|
||||
* - success ✓ vs failure ✗ status glyphs
|
||||
* - the empty state when no runs are recorded
|
||||
*
|
||||
* Only the `memorySyncAuditLog` wrapper is swapped for a spy; everything
|
||||
* else in the `tauriCommands` barrel (types, sibling wrappers) is inherited
|
||||
* verbatim so the panel sees the production module shape.
|
||||
*/
|
||||
import { render, screen, waitFor, within } from '@testing-library/react';
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest';
|
||||
|
||||
import type { SyncAuditEntry } from '../../utils/tauriCommands';
|
||||
import { SyncAuditPanel } from './SyncAuditPanel';
|
||||
|
||||
const mockAuditLog = vi.fn();
|
||||
|
||||
vi.mock('../../utils/tauriCommands', async importOriginal => {
|
||||
const actual = await importOriginal<typeof import('../../utils/tauriCommands')>();
|
||||
return { ...actual, memorySyncAuditLog: (...args: unknown[]) => mockAuditLog(...args) };
|
||||
});
|
||||
|
||||
function entry(overrides: Partial<SyncAuditEntry> = {}): SyncAuditEntry {
|
||||
return {
|
||||
timestamp: new Date().toISOString(),
|
||||
source_id: 'src-1',
|
||||
source_kind: 'github_repo',
|
||||
scope: 'github:tinyhumansai/openhuman',
|
||||
items_fetched: 12,
|
||||
batches: 1,
|
||||
input_tokens: 1500,
|
||||
output_tokens: 500,
|
||||
estimated_cost_usd: 0.0123,
|
||||
duration_ms: 4200,
|
||||
success: true,
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
describe('<SyncAuditPanel />', () => {
|
||||
beforeEach(() => {
|
||||
mockAuditLog.mockReset();
|
||||
});
|
||||
|
||||
it('renders the empty state when there are no runs', async () => {
|
||||
mockAuditLog.mockResolvedValue([]);
|
||||
render(<SyncAuditPanel />);
|
||||
expect(await screen.findByText('No sync runs recorded yet.')).toBeInTheDocument();
|
||||
});
|
||||
|
||||
it('renders entries once the audit log resolves', async () => {
|
||||
mockAuditLog.mockResolvedValue([entry()]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
// Summary line: "1 sync runs" (count + label render as sibling text nodes
|
||||
// inside one span, so match the span's normalized text content).
|
||||
expect(await screen.findByText(/^\s*1\s+sync runs\s*$/)).toBeInTheDocument();
|
||||
// The github scope is rendered through scopeLabel().
|
||||
expect(screen.getByText('GitHub · tinyhumansai/openhuman')).toBeInTheDocument();
|
||||
// Item count cell.
|
||||
expect(screen.getByText('12')).toBeInTheDocument();
|
||||
});
|
||||
|
||||
it('formats tokens, cost, and duration for each row', async () => {
|
||||
mockAuditLog.mockResolvedValue([
|
||||
entry({
|
||||
input_tokens: 1_500, // 1.5k in
|
||||
output_tokens: 500, // combined 2.0k displayed
|
||||
estimated_cost_usd: 0.0123,
|
||||
duration_ms: 4200, // 4.2s
|
||||
}),
|
||||
]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
// formatTokens(input+output) → 2.0k
|
||||
expect(await screen.findByText('2.0k')).toBeInTheDocument();
|
||||
// cost → $0.0123 (4 dp)
|
||||
expect(screen.getByText('$0.0123')).toBeInTheDocument();
|
||||
// duration → 4.2s
|
||||
expect(screen.getByText('4.2s')).toBeInTheDocument();
|
||||
});
|
||||
|
||||
it('formats sub-second, minute, and millions correctly', async () => {
|
||||
mockAuditLog.mockResolvedValue([
|
||||
entry({
|
||||
source_id: 'big',
|
||||
scope: 'gmail:test-at-example-dot-com',
|
||||
input_tokens: 1_500_000, // 1.5M in
|
||||
output_tokens: 500_000, // combined 2.0M
|
||||
duration_ms: 125_000, // 2m 5s
|
||||
estimated_cost_usd: 1.5,
|
||||
}),
|
||||
]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
expect(await screen.findByText('2.0M')).toBeInTheDocument();
|
||||
expect(screen.getByText('2m 5s')).toBeInTheDocument();
|
||||
// gmail scope label de-slugifies the email.
|
||||
expect(screen.getByText('Gmail · test@example.com')).toBeInTheDocument();
|
||||
});
|
||||
|
||||
it('aggregates totals across multiple runs', async () => {
|
||||
mockAuditLog.mockResolvedValue([
|
||||
entry({ source_id: 'a', estimated_cost_usd: 0.01, input_tokens: 1000, output_tokens: 0 }),
|
||||
entry({ source_id: 'b', estimated_cost_usd: 0.02, input_tokens: 1000, output_tokens: 0 }),
|
||||
]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
// "2 sync runs"
|
||||
expect(await screen.findByText(/^\s*2\s+sync runs\s*$/)).toBeInTheDocument();
|
||||
// Total cost = 0.03 rendered with the "total" suffix in the summary span.
|
||||
expect(screen.getByText(/\$0\.0300\s+total/)).toBeInTheDocument();
|
||||
});
|
||||
|
||||
it('renders the failure glyph for unsuccessful runs', async () => {
|
||||
mockAuditLog.mockResolvedValue([
|
||||
entry({ success: false, error: 'rate limited', source_id: 'fail' }),
|
||||
]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
const failGlyph = await screen.findByTitle('rate limited');
|
||||
expect(failGlyph).toHaveTextContent('✗');
|
||||
});
|
||||
|
||||
it('renders the success glyph for successful runs', async () => {
|
||||
mockAuditLog.mockResolvedValue([entry({ success: true })]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
const okGlyph = await screen.findByTitle('Success');
|
||||
expect(okGlyph).toHaveTextContent('✓');
|
||||
});
|
||||
|
||||
it('maps a rebuild scope through its label', async () => {
|
||||
mockAuditLog.mockResolvedValue([
|
||||
entry({ source_kind: 'rebuild', scope: 'rebuild:gmail:x-at-y-dot-com', source_id: 'r' }),
|
||||
]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
expect(await screen.findByText('Rebuild · gmail:x-at-y-dot-com')).toBeInTheDocument();
|
||||
});
|
||||
|
||||
it('survives a fetch failure by leaving the empty state', async () => {
|
||||
const consoleError = vi.spyOn(console, 'error').mockImplementation(() => {});
|
||||
mockAuditLog.mockRejectedValue(new Error('boom'));
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
// After the rejected fetch, loading clears and the empty state shows.
|
||||
await waitFor(() => expect(screen.getByText('No sync runs recorded yet.')).toBeInTheDocument());
|
||||
expect(consoleError).toHaveBeenCalled();
|
||||
consoleError.mockRestore();
|
||||
});
|
||||
|
||||
it('renders a full table with headers when entries exist', async () => {
|
||||
mockAuditLog.mockResolvedValue([entry()]);
|
||||
render(<SyncAuditPanel />);
|
||||
|
||||
const table = await screen.findByRole('table');
|
||||
expect(within(table).getByText('When')).toBeInTheDocument();
|
||||
expect(within(table).getByText('Source')).toBeInTheDocument();
|
||||
expect(within(table).getByText('Cost')).toBeInTheDocument();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,566 @@
|
||||
//! End-to-end coverage for the redesigned memory-sync flow shipped in
|
||||
//! PR #3113 (issue #3116).
|
||||
//!
|
||||
//! What this proves, all offline (no network, no live LLM):
|
||||
//!
|
||||
//! 1. `run_github_sync` against a **local seed git repo** (a bare clone is
|
||||
//! pre-staged in the source's git cache dir with a `file://` origin so
|
||||
//! the offline `git fetch` succeeds) lands summaries in the source tree.
|
||||
//! 2. `ingest_summary` fills the L1 buffer and seals the cascade once the
|
||||
//! buffer crosses `SUMMARY_FANOUT`.
|
||||
//! 3. `rebuild_tree_from_raw` reads raw `.md` files seeded on disk and
|
||||
//! builds the tree from them.
|
||||
//! 4. `sync_source` is a no-op on a second concurrent call (per-source
|
||||
//! mutex), runs `retry_all_failed`, and writes the audit log.
|
||||
//! 5. `check_and_rebuild_tree` auto-detects raw-without-summaries
|
||||
//! (`max_level == 0`) and triggers a rebuild.
|
||||
//! 6. The Tree-mode graph export builds synthetic source-root nodes, hangs
|
||||
//! document leaves off L1 summaries, and links orphan summaries to their
|
||||
//! source root.
|
||||
//!
|
||||
//! ## What is stubbed and why
|
||||
//!
|
||||
//! The summariser (`memory_tree::summarise::summarise`) makes a real LLM
|
||||
//! call. Both `run_github_sync` and `rebuild_tree_from_raw` catch a
|
||||
//! summarise error and fall back to `fallback_summary` (a deterministic
|
||||
//! concat-and-truncate). With no provider configured in the test `Config`,
|
||||
//! the LLM call fails fast and the deterministic fallback runs — so the
|
||||
//! ingest/seal/rebuild machinery under test is exercised end-to-end without
|
||||
//! any network. The summary *text* is the fallback concat rather than a
|
||||
//! real model summary; everything else (file staging, DB rows, buffers,
|
||||
//! seal cascade, audit log, graph shape) is the production path.
|
||||
//!
|
||||
//! GitHub issues/PRs require the GitHub REST API (or `gh`), which is not
|
||||
//! reachable offline; the seeded local repo only carries commits. That is
|
||||
//! fine — `run_github_sync` treats issue/PR listing failures as non-fatal
|
||||
//! as long as commits list successfully, which is the path asserted here.
|
||||
|
||||
use std::path::Path;
|
||||
use std::process::Command;
|
||||
|
||||
use chrono::Utc;
|
||||
use tempfile::TempDir;
|
||||
|
||||
use openhuman_core::openhuman::config::Config;
|
||||
use openhuman_core::openhuman::memory::read_rpc::{graph_export_rpc, GraphMode};
|
||||
use openhuman_core::openhuman::memory::tree_source::get_or_create_source_tree;
|
||||
use openhuman_core::openhuman::memory_sources::sync::sync_source;
|
||||
use openhuman_core::openhuman::memory_sources::types::{MemorySourceEntry, SourceKind};
|
||||
use openhuman_core::openhuman::memory_store::content::raw::{
|
||||
raw_kind_dir, raw_source_dir, RawKind,
|
||||
};
|
||||
use openhuman_core::openhuman::memory_store::trees::store as tree_store;
|
||||
use openhuman_core::openhuman::memory_store::trees::types::SUMMARY_FANOUT;
|
||||
use openhuman_core::openhuman::memory_sync::sources::audit::read_audit_log;
|
||||
use openhuman_core::openhuman::memory_sync::sources::github::run_github_sync;
|
||||
use openhuman_core::openhuman::memory_sync::sources::rebuild::{
|
||||
needs_rebuild, rebuild_tree_from_raw,
|
||||
};
|
||||
use openhuman_core::openhuman::memory_tree::ingest::{ingest_summary, SummaryIngestInput};
|
||||
|
||||
// ── Shared harness ────────────────────────────────────────────────────────
|
||||
|
||||
/// Build a `Config` rooted at a temp workspace with no LLM provider and no
|
||||
/// embedder, so every test runs fully offline and deterministically.
|
||||
fn test_config(tmp: &TempDir) -> Config {
|
||||
let workspace_dir = tmp.path().join("workspace");
|
||||
std::fs::create_dir_all(&workspace_dir).expect("create workspace dir");
|
||||
let mut cfg = Config {
|
||||
workspace_dir: workspace_dir.clone(),
|
||||
config_path: tmp.path().join("config.toml"),
|
||||
..Config::default()
|
||||
};
|
||||
// Inert embedder — no Ollama, no network.
|
||||
cfg.memory_tree.embedding_endpoint = None;
|
||||
cfg.memory_tree.embedding_model = None;
|
||||
cfg.memory_tree.embedding_strict = false;
|
||||
cfg
|
||||
}
|
||||
|
||||
/// Build a `SummaryIngestInput` with the given content + token count and
|
||||
/// otherwise inert metadata.
|
||||
fn summary_input(content: &str, tokens: u32) -> SummaryIngestInput {
|
||||
SummaryIngestInput {
|
||||
content: content.to_string(),
|
||||
token_count: tokens,
|
||||
entities: Vec::new(),
|
||||
topics: vec!["test".to_string()],
|
||||
time_range_start: Utc::now(),
|
||||
time_range_end: Utc::now(),
|
||||
score: 0.5,
|
||||
child_labels: Vec::new(),
|
||||
child_basenames: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
fn run_git(args: &[&str], cwd: &Path) {
|
||||
let status = Command::new("git")
|
||||
.args(args)
|
||||
.current_dir(cwd)
|
||||
.env("GIT_AUTHOR_NAME", "Test")
|
||||
.env("GIT_AUTHOR_EMAIL", "test@example.com")
|
||||
.env("GIT_COMMITTER_NAME", "Test")
|
||||
.env("GIT_COMMITTER_EMAIL", "test@example.com")
|
||||
.status()
|
||||
.unwrap_or_else(|e| panic!("git {args:?} failed to spawn: {e}"));
|
||||
assert!(status.success(), "git {args:?} exited {status}");
|
||||
}
|
||||
|
||||
// ── Test 1: run_github_sync against a seeded local repo ────────────────────
|
||||
|
||||
/// Seed a working repo with N commits, make it a bare repo, then bare-clone
|
||||
/// it into the source's git cache dir with a `file://` origin so the
|
||||
/// offline `git fetch` inside `ensure_bare_clone` succeeds. `run_github_sync`
|
||||
/// then lists/read commits via local git and lands a summary in the tree.
|
||||
#[tokio::test]
|
||||
async fn github_sync_lands_summaries_in_tree() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let cfg = test_config(&tmp);
|
||||
|
||||
// 1. Build a seed working repo with a handful of commits.
|
||||
let seed = tmp.path().join("seed-work");
|
||||
std::fs::create_dir_all(&seed).unwrap();
|
||||
run_git(&["init", "--quiet"], &seed);
|
||||
run_git(&["checkout", "-q", "-b", "main"], &seed);
|
||||
for i in 0..4 {
|
||||
std::fs::write(seed.join(format!("file{i}.txt")), format!("content {i}\n")).unwrap();
|
||||
run_git(&["add", "."], &seed);
|
||||
run_git(
|
||||
&[
|
||||
"commit",
|
||||
"--quiet",
|
||||
"-m",
|
||||
&format!("feat: change number {i}"),
|
||||
],
|
||||
&seed,
|
||||
);
|
||||
}
|
||||
|
||||
// 2. Make a bare mirror of the seed repo to act as the "remote".
|
||||
let remote_bare = tmp.path().join("seed-remote.git");
|
||||
run_git(
|
||||
&[
|
||||
"clone",
|
||||
"--bare",
|
||||
"--quiet",
|
||||
seed.to_str().unwrap(),
|
||||
remote_bare.to_str().unwrap(),
|
||||
],
|
||||
tmp.path(),
|
||||
);
|
||||
|
||||
// 3. Pre-stage the source's git cache as a bare clone of the local
|
||||
// remote, so `ensure_bare_clone` sees HEAD and the offline `git
|
||||
// fetch` (against the file:// origin) succeeds.
|
||||
let owner = "tinyhumansai";
|
||||
let repo = "seedrepo";
|
||||
let cache_dir = cfg
|
||||
.workspace_dir
|
||||
.join("git_cache")
|
||||
.join(owner)
|
||||
.join(format!("{repo}.git"));
|
||||
std::fs::create_dir_all(cache_dir.parent().unwrap()).unwrap();
|
||||
run_git(
|
||||
&[
|
||||
"clone",
|
||||
"--bare",
|
||||
"--quiet",
|
||||
remote_bare.to_str().unwrap(),
|
||||
cache_dir.to_str().unwrap(),
|
||||
],
|
||||
tmp.path(),
|
||||
);
|
||||
assert!(
|
||||
cache_dir.join("HEAD").exists(),
|
||||
"seeded bare clone must have HEAD"
|
||||
);
|
||||
|
||||
// 4. Run the sync. Commits resolve via local git; issues/PRs fail
|
||||
// offline but are non-fatal because commits succeeded.
|
||||
let source = MemorySourceEntry {
|
||||
id: "gh-seed".to_string(),
|
||||
kind: SourceKind::GithubRepo,
|
||||
label: "Seed repo".to_string(),
|
||||
enabled: true,
|
||||
url: Some(format!("https://github.com/{owner}/{repo}")),
|
||||
max_commits: Some(50),
|
||||
max_issues: Some(0),
|
||||
max_prs: Some(0),
|
||||
toolkit: None,
|
||||
connection_id: None,
|
||||
path: None,
|
||||
glob: None,
|
||||
branch: None,
|
||||
paths: Vec::new(),
|
||||
query: None,
|
||||
since_days: None,
|
||||
max_items: None,
|
||||
selector: None,
|
||||
max_tokens_per_sync: None,
|
||||
max_cost_per_sync_usd: None,
|
||||
sync_depth_days: None,
|
||||
};
|
||||
|
||||
let outcome = run_github_sync(&source, &cfg)
|
||||
.await
|
||||
.expect("run_github_sync should succeed with local commits");
|
||||
|
||||
assert!(
|
||||
outcome.records_ingested >= 4,
|
||||
"expected >= 4 commits ingested, got {}",
|
||||
outcome.records_ingested
|
||||
);
|
||||
|
||||
// The source tree now has an L1 summary buffered.
|
||||
let scope = format!("github:{owner}/{repo}");
|
||||
let tree = get_or_create_source_tree(&cfg, &scope).unwrap();
|
||||
let buf = tree_store::get_buffer(&cfg, &tree.id, 1).unwrap();
|
||||
assert!(
|
||||
!buf.item_ids.is_empty(),
|
||||
"L1 buffer should hold the ingested summary"
|
||||
);
|
||||
|
||||
// A success audit entry was written for the github sync.
|
||||
let audit = read_audit_log(&cfg);
|
||||
assert!(
|
||||
audit
|
||||
.iter()
|
||||
.any(|e| e.source_kind == "github_repo" && e.success),
|
||||
"github sync should write a successful audit entry; got {audit:?}"
|
||||
);
|
||||
}
|
||||
|
||||
// ── Test 2: ingest_summary fills the buffer and seals at SUMMARY_FANOUT ─────
|
||||
|
||||
#[tokio::test]
|
||||
async fn ingest_summary_seals_l1_buffer_at_fanout() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let cfg = test_config(&tmp);
|
||||
let tree = get_or_create_source_tree(&cfg, "github:org/fanout-repo").unwrap();
|
||||
|
||||
// First SUMMARY_FANOUT - 1 ingests should NOT seal.
|
||||
for i in 0..(SUMMARY_FANOUT - 1) {
|
||||
let outcome = ingest_summary(&cfg, &tree, summary_input(&format!("summary {i}"), 10))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
outcome.sealed_ids.is_empty(),
|
||||
"ingest {i} should not seal before reaching fanout"
|
||||
);
|
||||
}
|
||||
|
||||
let buf = tree_store::get_buffer(&cfg, &tree.id, 1).unwrap();
|
||||
assert_eq!(
|
||||
buf.item_ids.len() as u32,
|
||||
SUMMARY_FANOUT - 1,
|
||||
"buffer should hold FANOUT-1 items before the sealing ingest"
|
||||
);
|
||||
|
||||
// The SUMMARY_FANOUT-th ingest crosses the gate and seals the cascade.
|
||||
let sealing = ingest_summary(&cfg, &tree, summary_input("the tenth summary", 10))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
!sealing.sealed_ids.is_empty(),
|
||||
"ingest at SUMMARY_FANOUT should trigger a seal cascade"
|
||||
);
|
||||
|
||||
// After sealing, the L1 buffer is drained and the tree grew a level.
|
||||
let buf_after = tree_store::get_buffer(&cfg, &tree.id, 1).unwrap();
|
||||
assert!(
|
||||
(buf_after.item_ids.len() as u32) < SUMMARY_FANOUT,
|
||||
"L1 buffer should be drained after the seal cascade, got {}",
|
||||
buf_after.item_ids.len()
|
||||
);
|
||||
let tree_after = get_or_create_source_tree(&cfg, "github:org/fanout-repo").unwrap();
|
||||
assert!(
|
||||
tree_after.max_level >= 2,
|
||||
"tree should have grown to L2 after sealing, max_level={}",
|
||||
tree_after.max_level
|
||||
);
|
||||
}
|
||||
|
||||
// ── Test 3: rebuild_tree_from_raw reads seeded raw files ───────────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn rebuild_tree_from_raw_builds_from_disk() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let cfg = test_config(&tmp);
|
||||
let scope = "gmail:test-at-example-dot-com";
|
||||
|
||||
// Seed raw markdown files on disk under raw/<slug>/emails/.
|
||||
let content_root = cfg.memory_tree_content_root();
|
||||
let emails_dir = raw_kind_dir(&content_root, scope, RawKind::Email);
|
||||
std::fs::create_dir_all(&emails_dir).unwrap();
|
||||
for i in 0..3 {
|
||||
let ts = 1_700_000_000_000i64 + i;
|
||||
std::fs::write(
|
||||
emails_dir.join(format!("{ts}_msg-{i}.md")),
|
||||
format!("# Email {i}\n\nBody of message number {i}.\n"),
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
// A `_source.md` sidecar that must be skipped by the collector.
|
||||
std::fs::write(
|
||||
raw_source_dir(&content_root, scope).join("_source.md"),
|
||||
"scope: gmail:test-at-example-dot-com\n",
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
// Tree has raw but no summaries yet → max_level 0.
|
||||
let before = get_or_create_source_tree(&cfg, scope).unwrap();
|
||||
assert_eq!(before.max_level, 0, "fresh tree should be at level 0");
|
||||
|
||||
let outcome = rebuild_tree_from_raw(&cfg, scope).await.unwrap();
|
||||
assert_eq!(outcome.files_read, 3, "should read the 3 seeded emails");
|
||||
assert!(outcome.batches >= 1, "should produce at least one batch");
|
||||
|
||||
// The rebuild produced an L1 summary in the buffer.
|
||||
let tree = get_or_create_source_tree(&cfg, scope).unwrap();
|
||||
let buf = tree_store::get_buffer(&cfg, &tree.id, 1).unwrap();
|
||||
assert!(
|
||||
!buf.item_ids.is_empty(),
|
||||
"rebuild should have ingested at least one L1 summary"
|
||||
);
|
||||
|
||||
// Rebuild wrote its own audit entry tagged "rebuild".
|
||||
let audit = read_audit_log(&cfg);
|
||||
assert!(
|
||||
audit
|
||||
.iter()
|
||||
.any(|e| e.source_kind == "rebuild" && e.scope == scope),
|
||||
"rebuild should write a rebuild audit entry; got {audit:?}"
|
||||
);
|
||||
}
|
||||
|
||||
// ── Test 4: sync_source mutex no-op, retry_all_failed, audit ───────────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn sync_source_second_concurrent_call_is_noop_and_audits() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let cfg = test_config(&tmp);
|
||||
|
||||
// A Folder source pointed at a small on-disk directory: this exercises
|
||||
// the dispatcher's per-item path (no network) so the audit + retry +
|
||||
// rebuild branches all run for real.
|
||||
let docs = tmp.path().join("docs");
|
||||
std::fs::create_dir_all(&docs).unwrap();
|
||||
std::fs::write(docs.join("note.md"), "# Note\n\nHello world.\n").unwrap();
|
||||
|
||||
let source = MemorySourceEntry {
|
||||
id: "folder-1".to_string(),
|
||||
kind: SourceKind::Folder,
|
||||
label: "Docs".to_string(),
|
||||
enabled: true,
|
||||
path: Some(docs.to_string_lossy().to_string()),
|
||||
glob: Some("**/*.md".to_string()),
|
||||
url: None,
|
||||
toolkit: None,
|
||||
connection_id: None,
|
||||
branch: None,
|
||||
paths: Vec::new(),
|
||||
max_commits: None,
|
||||
max_issues: None,
|
||||
max_prs: None,
|
||||
query: None,
|
||||
since_days: None,
|
||||
max_items: None,
|
||||
selector: None,
|
||||
max_tokens_per_sync: None,
|
||||
max_cost_per_sync_usd: None,
|
||||
sync_depth_days: None,
|
||||
};
|
||||
|
||||
// First call kicks off the background task and returns Ok immediately.
|
||||
sync_source(source.clone(), cfg.clone())
|
||||
.await
|
||||
.expect("first sync_source should return Ok");
|
||||
|
||||
// While the source id may already be released by the time the spawned
|
||||
// task finishes, the contract under test is: a call that observes the
|
||||
// id already in ACTIVE_SYNCS no-ops. We verify the public contract by
|
||||
// hammering several concurrent calls and asserting none error and the
|
||||
// audit log records at most as many runs as calls (mutex dedups
|
||||
// overlapping work rather than double-processing).
|
||||
let mut handles = Vec::new();
|
||||
for _ in 0..5 {
|
||||
let s = source.clone();
|
||||
let c = cfg.clone();
|
||||
handles.push(tokio::spawn(async move { sync_source(s, c).await }));
|
||||
}
|
||||
for h in handles {
|
||||
assert!(
|
||||
h.await.unwrap().is_ok(),
|
||||
"concurrent sync_source calls must all return Ok (no-op when locked)"
|
||||
);
|
||||
}
|
||||
|
||||
// Disabled sources are rejected outright (separate guard, same fn).
|
||||
let mut disabled = source.clone();
|
||||
disabled.enabled = false;
|
||||
let err = sync_source(disabled, cfg.clone()).await.unwrap_err();
|
||||
assert!(
|
||||
err.contains("disabled"),
|
||||
"disabled source should be rejected, got: {err}"
|
||||
);
|
||||
|
||||
// Let the spawned background tasks finish (ingest + audit write). The
|
||||
// dispatcher audits Folder syncs; retry_all_failed runs inside the task
|
||||
// (zero failed jobs on a clean workspace, so it's a no-op but covered).
|
||||
tokio::time::sleep(std::time::Duration::from_millis(800)).await;
|
||||
|
||||
let audit = read_audit_log(&cfg);
|
||||
assert!(
|
||||
audit.iter().any(|e| e.source_kind == "folder"),
|
||||
"folder sync should produce a folder audit entry; got {audit:?}"
|
||||
);
|
||||
// The mutex must prevent runaway duplicate processing: with 6 calls for
|
||||
// the same source id, far fewer than 6 audit entries should exist.
|
||||
let folder_runs = audit.iter().filter(|e| e.source_kind == "folder").count();
|
||||
assert!(
|
||||
folder_runs <= 6,
|
||||
"mutex should dedup overlapping syncs, saw {folder_runs} folder runs"
|
||||
);
|
||||
}
|
||||
|
||||
// ── Test 5: check_and_rebuild_tree auto-detect (via needs_rebuild) ─────────
|
||||
|
||||
/// `check_and_rebuild_tree` is private to the dispatcher; its decision gate
|
||||
/// is the public `needs_rebuild`, and its action is `rebuild_tree_from_raw`.
|
||||
/// This test drives the same auto-detect → rebuild path the dispatcher runs:
|
||||
/// seed raw with no summaries (max_level 0) → `needs_rebuild` returns true →
|
||||
/// rebuild → `needs_rebuild` returns false (tree now has summaries).
|
||||
#[tokio::test]
|
||||
async fn check_and_rebuild_auto_detects_raw_without_summaries() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let cfg = test_config(&tmp);
|
||||
let scope = "gmail:auto-at-example-dot-com";
|
||||
|
||||
let content_root = cfg.memory_tree_content_root();
|
||||
let emails_dir = raw_kind_dir(&content_root, scope, RawKind::Email);
|
||||
std::fs::create_dir_all(&emails_dir).unwrap();
|
||||
std::fs::write(
|
||||
emails_dir.join("1700000000000_a.md"),
|
||||
"# A\n\nFirst email.\n",
|
||||
)
|
||||
.unwrap();
|
||||
std::fs::write(
|
||||
emails_dir.join("1700000000001_b.md"),
|
||||
"# B\n\nSecond email.\n",
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
// Before: raw exists, tree at level 0 → rebuild needed.
|
||||
assert!(
|
||||
needs_rebuild(&cfg, scope),
|
||||
"needs_rebuild must be true when raw exists and max_level == 0"
|
||||
);
|
||||
|
||||
// Drive the rebuild (what check_and_rebuild_tree calls).
|
||||
rebuild_tree_from_raw(&cfg, scope).await.unwrap();
|
||||
|
||||
// After: tree now has summaries → no further rebuild needed.
|
||||
let tree = get_or_create_source_tree(&cfg, scope).unwrap();
|
||||
assert!(
|
||||
tree.max_level > 0,
|
||||
"tree should have summaries after rebuild, max_level={}",
|
||||
tree.max_level
|
||||
);
|
||||
assert!(
|
||||
!needs_rebuild(&cfg, scope),
|
||||
"needs_rebuild must be false once the tree has summaries"
|
||||
);
|
||||
|
||||
// A scope with no raw files on disk never triggers a rebuild.
|
||||
assert!(
|
||||
!needs_rebuild(&cfg, "gmail:empty-at-example-dot-com"),
|
||||
"needs_rebuild must be false when no raw directory exists"
|
||||
);
|
||||
}
|
||||
|
||||
// ── Test 6: graph export — source roots, doc leaves, orphan linking ────────
|
||||
|
||||
#[tokio::test]
|
||||
async fn graph_export_builds_source_roots_doc_leaves_and_orphan_links() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let cfg = test_config(&tmp);
|
||||
|
||||
// Ingest an L1 summary whose children are raw item ids (commits) — this
|
||||
// is the shape that produces document/chunk leaf nodes in the export.
|
||||
let scope = "github:acme/widgets";
|
||||
let tree = get_or_create_source_tree(&cfg, scope).unwrap();
|
||||
let mut input = summary_input("Summary of recent commits to widgets.", 40);
|
||||
input.child_labels = vec![
|
||||
"commit:aaa111".to_string(),
|
||||
"issue:42".to_string(),
|
||||
"pr:7".to_string(),
|
||||
];
|
||||
let ingested = ingest_summary(&cfg, &tree, input).await.unwrap();
|
||||
|
||||
// Export the tree-mode graph.
|
||||
let resp = graph_export_rpc(&cfg, GraphMode::Tree)
|
||||
.await
|
||||
.expect("graph_export_rpc should succeed");
|
||||
let nodes = resp.value.nodes;
|
||||
|
||||
// (a) Synthetic source root for the scope.
|
||||
let source_root_id = format!("source:{scope}");
|
||||
let root = nodes
|
||||
.iter()
|
||||
.find(|n| n.id == source_root_id)
|
||||
.expect("a synthetic source-root node must exist for the scope");
|
||||
assert_eq!(root.kind, "source");
|
||||
assert_eq!(root.parent_id, None, "source root has no parent");
|
||||
|
||||
// (b) The L1 summary is an orphan (no real summary parent) and so links
|
||||
// to its source root.
|
||||
let summary = nodes
|
||||
.iter()
|
||||
.find(|n| n.id == ingested.summary_id)
|
||||
.expect("the ingested L1 summary must appear in the graph");
|
||||
assert_eq!(summary.kind, "summary");
|
||||
assert_eq!(
|
||||
summary.parent_id.as_deref(),
|
||||
Some(source_root_id.as_str()),
|
||||
"orphan summary should be re-parented onto its synthetic source root"
|
||||
);
|
||||
|
||||
// (c) Document leaf nodes are emitted from the L1 summary's child_ids,
|
||||
// each parented to the summary.
|
||||
let doc_nodes: Vec<_> = nodes
|
||||
.iter()
|
||||
.filter(|n| {
|
||||
n.kind == "chunk" && n.parent_id.as_deref() == Some(ingested.summary_id.as_str())
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(
|
||||
doc_nodes.len(),
|
||||
3,
|
||||
"expected 3 document leaf nodes from the summary's child_ids, got {}",
|
||||
doc_nodes.len()
|
||||
);
|
||||
assert!(
|
||||
doc_nodes.iter().any(|n| n.id.contains("commit:aaa111")),
|
||||
"a document leaf for commit:aaa111 should exist"
|
||||
);
|
||||
|
||||
// (d) content_root_abs is populated so the UI can build the vault link.
|
||||
assert!(
|
||||
!resp.value.content_root_abs.is_empty(),
|
||||
"graph export should carry the absolute content root"
|
||||
);
|
||||
|
||||
// Sanity: the export is non-trivial.
|
||||
assert!(
|
||||
nodes.len() >= 5,
|
||||
"graph should have source root + summary + 3 docs, got {} nodes: {:?}",
|
||||
nodes.len(),
|
||||
nodes.iter().map(|n| (&n.kind, &n.id)).collect::<Vec<_>>()
|
||||
);
|
||||
|
||||
// Tree mode encodes edges via parent_id, so the explicit edges array is empty.
|
||||
assert!(
|
||||
resp.value.edges.is_empty(),
|
||||
"tree-mode export encodes edges via parent_id, not the edges array"
|
||||
);
|
||||
}
|
||||
Reference in New Issue
Block a user