File size: 5,539 Bytes
f8b48da
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
"""Remove duplicate uploaded documents from the index.



Duplicates were created by re-uploading the same file before uploads became

content-addressed: each upload landed in a different ``data/uploads/<random>``

folder and was indexed under its own ``document_id``, so the same content could

appear twice in retrieval results and unfairly dominate the ranking.



This script groups uploaded documents by ``(file name, extension)`` and keeps

the newest copy of each, deleting the older ones from PostgreSQL, Qdrant and

keyword postings.



Usage::



    python scripts/dedupe_documents.py --dry-run   # show what would be removed

    python scripts/dedupe_documents.py             # perform the cleanup



Note: the corpus intentionally stores the same topic in several formats

(e.g. ``api-errors.md`` and ``api-errors.html``). Those are NOT duplicates and

are left untouched, because the grouping key includes the file extension.

"""

from __future__ import annotations

import argparse
import os
import re
import sys
from collections import defaultdict
from pathlib import Path

ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(ROOT))

from app.db.models import Chunk, Document, IngestionJob  # noqa: E402
from app.db.session import SessionLocal  # noqa: E402
from app.retrieval.bm25 import BM25Indexer  # noqa: E402
from app.retrieval.vector import VectorStore  # noqa: E402

UPLOADS_MARKER = "uploads"

# Some older uploads were stored under an 8-hex prefixed name to avoid a
# filename collision. Normalise it away so those are recognised as duplicates
# of the original file rather than distinct documents.
_COLLISION_PREFIX_RE = re.compile(r"^[0-9a-f]{8}_")


def _group_key(doc: Document) -> tuple[str, str]:
    """Group by (normalised file name, lower-cased extension)."""
    name = os.path.basename(doc.source or "")
    name = _COLLISION_PREFIX_RE.sub("", name)
    return name.lower(), os.path.splitext(name)[1].lower()


def find_duplicate_uploads(session) -> list[tuple[Document, list[Document]]]:
    """Return ``(keep, [drop, ...])`` pairs for duplicated uploaded files."""
    docs = session.query(Document).filter(Document.source.like(f"%{UPLOADS_MARKER}%")).all()

    groups: dict[tuple[str, str], list[Document]] = defaultdict(list)
    for doc in docs:
        groups[_group_key(doc)].append(doc)

    plan: list[tuple[Document, list[Document]]] = []
    for group in groups.values():
        if len(group) < 2:
            continue
        # Prefer the copy whose stored name has no collision prefix (the
        # original), then the most recently created one.
        group.sort(
            key=lambda d: (
                bool(_COLLISION_PREFIX_RE.match(os.path.basename(d.source or ""))),
                -d.created_at.timestamp(),
            )
        )
        plan.append((group[0], group[1:]))
    return plan


def _delete_document(session, document_id: str) -> int:
    """Delete a document and its chunks from every store. Returns chunks removed."""
    # Order matters: ingestion_jobs and chunks both reference documents, so they
    # must go first or PostgreSQL raises a foreign-key violation.
    session.query(IngestionJob).filter(
        IngestionJob.document_id == document_id
    ).delete(synchronize_session=False)
    removed = (
        session.query(Chunk)
        .filter(Chunk.document_id == document_id)
        .delete(synchronize_session=False)
    )
    session.query(Document).filter(Document.document_id == document_id).delete(
        synchronize_session=False
    )
    session.commit()

    try:
        VectorStore().delete_by_document(document_id)
    except Exception as exc:  # noqa: BLE001
        print(f"    warning: Qdrant delete failed for {document_id[:12]}: {exc}")

    try:
        BM25Indexer().delete_by_document(document_id)
    except Exception as exc:  # noqa: BLE001
        print(f"    warning: keyword delete failed for {document_id[:12]}: {exc}")

    return int(removed)


def main() -> int:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument(
        "--dry-run",
        action="store_true",
        help="List what would be removed without deleting anything",
    )
    args = parser.parse_args()

    session = SessionLocal()
    try:
        plan = find_duplicate_uploads(session)
        if not plan:
            print("No duplicate uploads found.")
            return 0

        total = 0
        for keep, drops in plan:
            print(f"{os.path.basename(keep.source or '')} ({len(drops) + 1} copies)")
            print(f"  KEEP  {keep.document_id[:12]}  {keep.created_at}")
            for doc in drops:
                count = session.query(Chunk).filter_by(document_id=doc.document_id).count()
                if args.dry_run:
                    print(f"  DROP  {doc.document_id[:12]}  chunks={count}  {doc.created_at}")
                else:
                    removed = _delete_document(session, doc.document_id)
                    print(f"  DROPPED {doc.document_id[:12]}  chunks={removed}")
                    total += 1
            print()

        if args.dry_run:
            print(f"Dry run: would remove {sum(len(d) for _, d in plan)} document(s).")
        else:
            print(f"Removed {total} duplicate document(s).")
        return 0
    finally:
        session.close()


if __name__ == "__main__":
    raise SystemExit(main())