diff --git a/mempalace/convo_miner.py b/mempalace/convo_miner.py index 0ea784b..7213d96 100644 --- a/mempalace/convo_miner.py +++ b/mempalace/convo_miner.py @@ -11,6 +11,7 @@ Same palace as project mining. Different ingest strategy. import os import sys import json +import hashlib import logging import stat from pathlib import Path @@ -31,6 +32,7 @@ from .palace import ( get_collection, mine_lock, mine_palace_lock, + prefetch_content_hashes, prefetch_mined_set, ) @@ -117,7 +119,14 @@ def _is_regular_source_file(filepath: Path, root: Path) -> bool: pass -def _register_file(collection, source_file: str, wing: str, agent: str, extract_mode: str): +def _register_file( + collection, + source_file: str, + wing: str, + agent: str, + extract_mode: str, + content_hash: Optional[str] = None, +): """Write a sentinel so file_already_mined() returns True for 0-chunk files. Without this, files that normalize to nothing or produce zero chunks are @@ -128,6 +137,11 @@ def _register_file(collection, source_file: str, wing: str, agent: str, extract_ grows past the min-chunk-size floor (e.g. a short session that gets extended) is correctly detected as changed on the next mine instead of being skipped forever by this sentinel. + + Also used to register a file recognized as a content-duplicate of an + already-mined transcript under a different path (see + ``prefetch_content_hashes``) — stamping it here means the next run skips + it via the cheap mtime check instead of re-normalizing and re-hashing it. """ try: source_mtime = os.path.getmtime(source_file) @@ -147,6 +161,8 @@ def _register_file(collection, source_file: str, wing: str, agent: str, extract_ } if source_mtime is not None: meta["source_mtime"] = source_mtime + if content_hash is not None: + meta["content_hash"] = content_hash collection.upsert( documents=[f"[registry] {source_file}"], ids=[sentinel_id], @@ -487,7 +503,15 @@ def _extract_authored_at(filepath): def _file_chunks_locked( - collection, source_file, chunks, wing, room, agent, extract_mode, authored_at=None + collection, + source_file, + chunks, + wing, + room, + agent, + extract_mode, + authored_at=None, + content_hash=None, ): """Lock the source file, purge stale drawers, and upsert fresh chunks. @@ -560,6 +584,8 @@ def _file_chunks_locked( } if source_mtime is not None: meta["source_mtime"] = source_mtime + if content_hash is not None: + meta["content_hash"] = content_hash batch_metas.append(meta) assert_no_collisions(list(zip(batch_ids, batch_metas)), collection) try: @@ -604,6 +630,37 @@ def _is_ai_tool_path(path: Path) -> bool: return False +def _skip_content_duplicate( + dry_run: bool, + collection, + filepath: Path, + content_hash: str, + wing: str, + agent: str, + extract_mode: str, + mined_content_hashes: dict, + progress: str, +) -> bool: + """Register and report a content-duplicate file, if this content_hash is + already filed under a different source_file. Returns True if skipped. + + Repeated exports from Claude/ChatGPT commonly land under a new filename + each run even when the conversation itself is unchanged, so the + source_file-keyed mtime skip never recognizes them. Registering the + duplicate here means the next run skips it via the cheap mtime check + instead of re-normalizing and re-hashing it. + """ + if dry_run: + return False + source_file = str(filepath) + dup_source = mined_content_hashes.get(content_hash) + if dup_source is None or dup_source == source_file: + return False + _register_file(collection, source_file, wing, agent, extract_mode, content_hash) + print(f" = {progress} {filepath.name[:50]:50} duplicate of {Path(dup_source).name}") + return True + + def _is_unchanged_since_last_mine(source_file: str, mined_mtimes: dict) -> bool: """True iff source_file was mined at the current schema AND its on-disk mtime still matches what was stored -- the mtime-aware replacement for @@ -721,6 +778,35 @@ def _compute_hallways_for_wing_safe(wing, collection, drawers_filed, config=None print(f" (hallways skipped: {exc})") +def _normalize_convo_content( + filepath: Path, + source_file: str, + cfg_min_chunk_size: int, + collection, + wing: str, + agent: str, + extract_mode: str, + dry_run: bool, +) -> Optional[str]: + """Normalize a transcript file, registering it as filed when there's + nothing worth mining. Returns None when the caller should skip the file + (normalize failed, or normalized content is too short to chunk). + """ + try: + content = normalize(str(filepath)) + except (OSError, ValueError): + if not dry_run: + _register_file(collection, source_file, wing, agent, extract_mode) + return None + + if not content or len(content.strip()) < cfg_min_chunk_size: + if not dry_run: + _register_file(collection, source_file, wing, agent, extract_mode) + return None + + return content + + def _mine_convos_impl( convo_dir: str, palace_path: str, @@ -773,6 +859,14 @@ def _mine_convos_impl( mined_mtimes: dict = ( prefetch_mined_set(collection, extract_mode=extract_mode) if not dry_run else {} ) + # content_hash -> source_file for transcripts already filed. Repeated + # exports from Claude/ChatGPT commonly land under a new filename each + # run even when the conversation itself is unchanged, so the + # source_file-keyed skip above ("mined_mtimes") never recognizes them — + # this catches the same conversation reappearing at a new path. + mined_content_hashes: dict = ( + prefetch_content_hashes(collection, extract_mode=extract_mode) if not dry_run else {} + ) total_drawers = 0 files_mined = 0 @@ -800,17 +894,37 @@ def _mine_convos_impl( files_skipped += 1 continue - # Normalize format - try: - content = normalize(str(filepath)) - except (OSError, ValueError): - if not dry_run: - _register_file(collection, source_file, wing, agent, extract_mode) + content = _normalize_convo_content( + filepath, + source_file, + cfg_min_chunk_size, + collection, + wing, + agent, + extract_mode, + dry_run, + ) + if content is None: continue - if not content or len(content.strip()) < cfg_min_chunk_size: - if not dry_run: - _register_file(collection, source_file, wing, agent, extract_mode) + content_hash = hashlib.sha256(content.strip().encode("utf-8")).hexdigest() + + # Same conversation, different path: this exact transcript is + # already filed under another source_file (e.g. a fresh Claude/ChatGPT + # export bundle re-exports conversations that haven't changed under + # new filenames). Skip it instead of filing a duplicate set of drawers. + if _skip_content_duplicate( + dry_run, + collection, + filepath, + content_hash, + wing, + agent, + extract_mode, + mined_content_hashes, + progress=f"[{i:4}/{len(files)}]", + ): + files_skipped += 1 continue # Chunk — either exchange pairs or general extraction @@ -872,6 +986,7 @@ def _mine_convos_impl( agent, extract_mode, authored_at=_extract_authored_at(filepath), + content_hash=content_hash, ) if skipped: files_skipped += 1 @@ -879,6 +994,7 @@ def _mine_convos_impl( for r, n in room_delta.items(): room_counts[r] += n + mined_content_hashes[content_hash] = source_file total_drawers += drawers_added files_mined += 1 print(f" + [{i:4}/{len(files)}] {filepath.name[:50]:50} +{drawers_added}") diff --git a/mempalace/palace.py b/mempalace/palace.py index 49ad06c..1230640 100644 --- a/mempalace/palace.py +++ b/mempalace/palace.py @@ -1351,3 +1351,44 @@ def prefetch_mined_set( except Exception: logger.warning("prefetch_mined_set: partial fetch, %d files loaded", len(mined)) return mined + + +def prefetch_content_hashes(collection, extract_mode: Optional[str] = None) -> dict[str, str]: + """Pre-fetch content_hash -> source_file for drawers already filed at the + current NORMALIZE_VERSION, in one bulk pass. + + Repeated exports from Claude/ChatGPT land under a new filename each run + (timestamped bundle, regenerated slug, etc.) even when the conversation + itself hasn't changed. `prefetch_mined_set` only recognizes a file as + already-mined by its exact path, so the same conversation re-exported + under a new path always looked "new" and got re-mined as a duplicate + drawer. This does the same bulk scan but keyed on the SHA-256 of the + normalized transcript text, so the convo miner can recognize "this exact + conversation is already filed under a different path" and skip it. + + Only the first source_file seen for a given hash is kept — good enough + to detect and skip a repeat, the point is not to track every alias. + """ + hashes: dict[str, str] = {} + try: + total = collection.count() + offset = 0 + while offset < total: + batch = collection.get(limit=1000, offset=offset, include=["metadatas"]) + for meta in batch["metadatas"]: + meta = meta or {} + content_hash = meta.get("content_hash") + src = meta.get("source_file") + if not content_hash or not src: + continue + if not _metadata_matches_extract_mode(meta, extract_mode): + continue + version = meta.get("normalize_version", 1) + if version >= NORMALIZE_VERSION and content_hash not in hashes: + hashes[content_hash] = src + if not batch["ids"]: + break + offset += len(batch["ids"]) + except Exception: + logger.warning("prefetch_content_hashes: partial fetch, %d hashes loaded", len(hashes)) + return hashes diff --git a/tests/test_convo_miner.py b/tests/test_convo_miner.py index df35be6..f1ba8a6 100644 --- a/tests/test_convo_miner.py +++ b/tests/test_convo_miner.py @@ -572,6 +572,44 @@ def test_mine_convos_grown_file_purges_stale_drawers_not_additive(capsys): shutil.rmtree(tmpdir, ignore_errors=True) +def test_mine_convos_skips_same_content_under_new_filename(capsys): + """Re-exporting the same conversation from Claude/ChatGPT under a new + filename (fresh export bundle, regenerated slug, etc.) must not create + a duplicate set of drawers -- only the exact-new content should file.""" + tmpdir = tempfile.mkdtemp() + try: + transcript = ( + "> What is the plan?\nStart with the schema, then the API.\n\n" + "> Any risks?\nMigration ordering is the main one.\n" + ) + (Path(tmpdir) / "export_2026-01-01.txt").write_text(transcript) + palace_path = os.path.join(tmpdir, "palace") + mine_convos(tmpdir, palace_path, wing="test") + + client = chromadb.PersistentClient(path=palace_path) + col = client.get_collection("mempalace_drawers") + count_after_first = col.count() + assert count_after_first >= 2 + + # Simulate a later export: the same conversation lands under a new + # filename, alongside one genuinely new conversation. + (Path(tmpdir) / "export_2026-02-01.txt").write_text(transcript) + (Path(tmpdir) / "export_2026-02-01_new.txt").write_text( + "> What's next?\nUNIQUE_SECOND_EXPORT_MARKER covers the new work.\n" + ) + mine_convos(tmpdir, palace_path, wing="test") + out = capsys.readouterr().out + assert "duplicate of export_2026-01-01.txt" in out + + col = client.get_collection("mempalace_drawers") + docs = col.get(include=["documents"])["documents"] + dup_hits = sum(1 for d in docs if "Migration ordering is the main one" in d) + assert dup_hits == 1, f"duplicate transcript re-filed: {dup_hits} copies" + assert any("UNIQUE_SECOND_EXPORT_MARKER" in d for d in docs) + finally: + shutil.rmtree(tmpdir, ignore_errors=True) + + def test_prefetch_mined_set_returns_stored_mtime(): """prefetch_mined_set's dict carries each source_file's stored mtime, not just membership."""