File size: 5,538 Bytes
0004cda
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
import os
import time
import sys
import json
import dotenv

from pinecone import Pinecone
from app import load_chunks, get_vector_store

# Load environment variables
dotenv.load_dotenv()

CHECKPOINT_FILE = "upload_checkpoint.json"

BATCH_SIZE = 64  # Optimally sized for high-precision paragraph chunks (averaging ~230 tokens each)
SLEEP_BETWEEN_BATCHES = 0.2  # Throttling is now managed by the high-speed Google Gemini API
MAX_RETRIES = 5
BACKOFF_FACTOR = 2

def clear_index():
    api_key = os.getenv("PINECONE_API_KEY")
    if not api_key:
        print("Error: PINECONE_API_KEY is not set.")
        sys.exit(1)
        
    pc = Pinecone(api_key=api_key)
    index_name = "branham-index"
    
    print(f"Connecting to Pinecone index '{index_name}'...")
    idx = pc.Index(index_name)
    
    stats = idx.describe_index_stats()
    total_vectors = stats.get("total_vector_count", 0)
    print(f"Current total vector count before clear: {total_vectors}")
    
    print(f"Deleting all vectors in '{index_name}' to start fresh...")
    try:
        idx.delete(delete_all=True)
    except Exception as e:
        print(f"Index deletion skipped (it might already be completely empty): {e}")
    
    # Wait for deletion to reflect
    time.sleep(5)
    stats = idx.describe_index_stats()
    print(f"Post-clear total vector count: {stats.get('total_vector_count', 0)}")
    print("Pinecone index cleared successfully!\n")

def upload_in_batches(chunks, start_idx=0):
    # Model is loaded inside get_vector_store(), track the exact loading time
    print("Initializing embedding model and Pinecone connection...")
    model_load_start = time.perf_counter()
    vector_store = get_vector_store()
    print(f"Initialization complete in {time.perf_counter() - model_load_start:.1f}s.\n")
    
    total_chunks = len(chunks)
    
    print(f"Starting batch upload of {total_chunks} chunks to Pinecone...")
    print(f"Batch Size: {BATCH_SIZE} | Spacer delay: {SLEEP_BETWEEN_BATCHES}s | Max Retries: {MAX_RETRIES}\n")
    
    # Track elapsed time strictly for encoding + uploading
    start_time = time.perf_counter()
    
    for i in range(start_idx, total_chunks, BATCH_SIZE):
        batch = chunks[i : i + BATCH_SIZE]
        batch_num = (i // BATCH_SIZE) + 1
        total_batches = (total_chunks + BATCH_SIZE - 1) // BATCH_SIZE
        
        # Implement robust retry with exponential backoff
        retries = 0
        while retries <= MAX_RETRIES:
            try:
                # Add to Pinecone
                vector_store.add_documents(batch)
                
                # Progress logging
                elapsed = time.perf_counter() - start_time
                pct = (min(i + BATCH_SIZE, total_chunks) / total_chunks) * 100
                rate = (i + len(batch) - start_idx) / elapsed if elapsed > 0 else 0
                eta = (total_chunks - (i + len(batch))) / rate if rate > 0 else 0
                
                print(
                    f"[{pct:6.2f}%] Batch {batch_num}/{total_batches} uploaded successfully. "
                    f"({min(i + BATCH_SIZE, total_chunks)}/{total_chunks}) | "
                    f"Speed: {rate:.1f} chunks/sec | ETA: {eta/60:.1f} min"
                )
                
                # Write checkpoint to file
                with open(CHECKPOINT_FILE, "w", encoding="utf-8") as f:
                    json.dump({"last_uploaded_index": i + len(batch)}, f)
                break
            except Exception as e:
                retries += 1
                if retries > MAX_RETRIES:
                    print(f"\n[FATAL ERROR] Batch {batch_num} failed completely after {MAX_RETRIES} retries. Error: {e}")
                    raise e
                    
                sleep_time = BACKOFF_FACTOR ** retries
                print(
                    f"\n[WARNING] Error on Batch {batch_num} (Attempt {retries}/{MAX_RETRIES}): {e}. "
                    f"Retrying in {sleep_time}s..."
                )
                time.sleep(sleep_time)
                
        time.sleep(SLEEP_BETWEEN_BATCHES)
        
    total_time = time.perf_counter() - start_time
    print(f"\nSUCCESS! Uploaded {total_chunks} chunks to Pinecone in {total_time/60:.2f} minutes.")
    
    # Delete checkpoint on success
    if os.path.exists(CHECKPOINT_FILE):
        os.remove(CHECKPOINT_FILE)

def main():
    # 1. Load the clean serialized chunks
    print("Loading chunks from 'sermon_chunks.pkl'...")
    chunks = load_chunks()
    if not chunks:
        print("Error: No chunks found in 'sermon_chunks.pkl'. Make sure chunking is complete.")
        return
        
    print(f"Successfully loaded {len(chunks)} sermon chunks.\n")
    
    # 2. Check for checkpoint to see if we can resume
    start_idx = 0
    if os.path.exists(CHECKPOINT_FILE):
        try:
            with open(CHECKPOINT_FILE, "r", encoding="utf-8") as f:
                data = json.load(f)
                start_idx = data.get("last_uploaded_index", 0)
        except Exception as e:
            print(f"[WARNING] Could not read checkpoint file: {e}. Starting fresh.")
            
    if start_idx > 0 and start_idx < len(chunks):
        print(f"[RESUME] Checkpoint found! Resuming upload from chunk index {start_idx}...")
    else:
        # 3. Clear index only when starting fresh
        print("[START FRESH] Starting fresh. Clearing index first...")
        clear_index()
        
    # 4. Upload chunks
    upload_in_batches(chunks, start_idx=start_idx)

if __name__ == "__main__":
    main()