11import json
22import logging
3+ import time
34import functools as ft
45
56import openai
67from openai import OpenAI , OpenAIError
8+ from openai .types import VectorStore
9+ from openai .types .vector_stores import VectorStoreFileBatch
710from pydantic import BaseModel
11+ from tenacity import (
12+ RetryCallState ,
13+ retry ,
14+ retry_if_exception_type ,
15+ stop_after_attempt ,
16+ wait_exponential ,
17+ )
818
919from app .models import Document , ProviderType
1020
1121logger = logging .getLogger (__name__ )
1222
1323OPENAI_PROVIDER = ProviderType .openai .value
1424
25+ # Under the Celery soft time limit so a hung call can't eat the whole task window.
26+ # SDK-level retries are off (registry.py: max_retries=0); tenacity is the sole
27+ # retry layer, wrapping batch create+index.
28+ OPENAI_TIMEOUT_SECONDS = 30
29+
30+ BATCH_POLL_INTERVAL_SECONDS = 2
31+
32+ # Retry batch create+index on any OpenAI/indexing failure, exponential backoff
33+ # (~2s, 4s, 8s), all inside one Celery soft-time-limit window.
34+ BATCH_INDEX_MAX_ATTEMPTS = 4
35+ BATCH_RETRY_BACKOFF_BASE_SECONDS = 2
36+
37+
38+ def _log_batch_retry (retry_state : RetryCallState ) -> None :
39+ logger .warning (
40+ f"[OpenAIVectorStoreCrud._create_and_index_batch] Batch attempt failed, retrying | "
41+ f"attempt={ retry_state .attempt_number } , "
42+ f"error={ retry_state .outcome .exception () if retry_state .outcome else None } "
43+ )
44+
1545
1646def vs_ls (client : OpenAI , vector_store_id : str ):
1747 kwargs = {}
@@ -85,7 +115,7 @@ def __init__(self, client):
85115
86116
87117class OpenAIVectorStoreCrud (OpenAICrud ):
88- def create (self ):
118+ def create (self ) -> VectorStore :
89119 logger .info (
90120 f"[OpenAIVectorStoreCrud.create] Creating vector store | {{'action': 'create'}}"
91121 )
@@ -101,6 +131,100 @@ def read(self, vector_store_id: str):
101131 )
102132 yield from vs_ls (self .client , vector_store_id )
103133
134+ def _create_file_batch (self , vector_store_id : str , file_ids : list [str ]) -> str :
135+ """Returns the vsfb_ id. poll()'s return deserializes a vector-store body,
136+ so its .id is the vs_ id - take the batch id from create()."""
137+ created = self .client .vector_stores .file_batches .create (
138+ vector_store_id = vector_store_id ,
139+ file_ids = file_ids ,
140+ )
141+ return created .id
142+
143+ def _retrieve_file_batch (
144+ self , batch_id : str , vector_store_id : str
145+ ) -> VectorStoreFileBatch :
146+ return self .client .vector_stores .file_batches .retrieve (
147+ batch_id , vector_store_id = vector_store_id
148+ )
149+
150+ def _poll_file_batch (
151+ self , batch_id : str , vector_store_id : str
152+ ) -> VectorStoreFileBatch :
153+ """Poll until indexing finishes; the Celery soft time limit is the deadline."""
154+ while True :
155+ batch = self ._retrieve_file_batch (batch_id , vector_store_id )
156+ if batch .status != "in_progress" :
157+ return batch
158+ time .sleep (BATCH_POLL_INTERVAL_SECONDS )
159+
160+ def _raise_if_batch_incomplete (
161+ self ,
162+ batch : VectorStoreFileBatch ,
163+ batch_id : str ,
164+ vector_store_id : str ,
165+ docs : list [Document ],
166+ ) -> None :
167+ """Raise on any indexing failure so the batch attempt is retried."""
168+ if batch .file_counts .failed > 0 :
169+ try :
170+ failed_files = self .client .vector_stores .file_batches .list_files (
171+ vector_store_id = vector_store_id ,
172+ batch_id = batch_id ,
173+ filter = "failed" ,
174+ )
175+ doc_by_file_id = {d .file_id [OPENAI_PROVIDER ]: d for d in docs }
176+ parts = []
177+ for f in failed_files :
178+ d = doc_by_file_id .get (f .id )
179+ label = d .fname if d else f .id
180+ msg = f .last_error .message if f .last_error else "no error detail"
181+ parts .append (f"{ label } : { msg } " )
182+ logger .error (
183+ f"[OpenAIVectorStoreCrud._raise_if_batch_incomplete] Files failed to index | "
184+ f"{{'batch_id': '{ batch_id } ', 'failed_files': '{ ', ' .join (parts )} '}}"
185+ )
186+ raise RuntimeError ("; " .join (parts ))
187+ except OpenAIError as err :
188+ logger .warning (
189+ f"[OpenAIVectorStoreCrud._raise_if_batch_incomplete] Could not fetch per-file errors | "
190+ f"{{'batch_id': '{ batch_id } ', 'error': '{ str (err )} '}}"
191+ )
192+ raise
193+
194+ # Only 'completed' is success; a 'cancelled'/'failed' batch with no per-file
195+ # failures slips past the failed-count check above.
196+ if batch .status != "completed" :
197+ error_message = (
198+ f"[OPENAI] Vector store indexing did not complete "
199+ f"(status: { batch .status } ). Retry the collection."
200+ )
201+ logger .error (
202+ f"[OpenAIVectorStoreCrud._raise_if_batch_incomplete] { error_message } | "
203+ f"vector_store_id={ vector_store_id } , batch_id={ batch_id } , "
204+ f"status={ batch .status } "
205+ )
206+ raise RuntimeError (error_message )
207+
208+ @retry (
209+ reraise = True ,
210+ stop = stop_after_attempt (BATCH_INDEX_MAX_ATTEMPTS ),
211+ wait = wait_exponential (multiplier = BATCH_RETRY_BACKOFF_BASE_SECONDS ),
212+ retry = retry_if_exception_type ((OpenAIError , RuntimeError )),
213+ before_sleep = _log_batch_retry ,
214+ )
215+ def _create_and_index_batch (
216+ self , vector_store_id : str , docs : list [Document ]
217+ ) -> tuple [VectorStoreFileBatch , str ]:
218+ """Create the file batch, wait for indexing, verify it completed. Retried as
219+ a unit on any OpenAI/indexing failure; SoftTimeLimitExceeded is not retried
220+ (not an OpenAIError/RuntimeError) so it aborts the task inside the window."""
221+ batch_id = self ._create_file_batch (
222+ vector_store_id , [doc .file_id [OPENAI_PROVIDER ] for doc in docs ]
223+ )
224+ batch = self ._poll_file_batch (batch_id , vector_store_id )
225+ self ._raise_if_batch_incomplete (batch , batch_id , vector_store_id , docs )
226+ return batch , batch_id
227+
104228 def update (
105229 self ,
106230 vector_store_id : str ,
@@ -115,11 +239,7 @@ def update(
115239 )
116240
117241 try :
118- batch = self .client .vector_stores .file_batches .upload_and_poll (
119- vector_store_id = vector_store_id ,
120- files = [],
121- file_ids = [doc .file_id [OPENAI_PROVIDER ] for doc in docs ],
122- )
242+ batch , batch_id = self ._create_and_index_batch (vector_store_id , docs )
123243 except openai .RateLimitError as e :
124244 error_message = (
125245 f"[OPENAI] Rate limit exceeded (code: { e .status_code } ): "
@@ -215,30 +335,9 @@ def update(
215335
216336 logger .info (
217337 f"[OpenAIVectorStoreCrud.update] Batch complete | "
218- f"{{'vector_store_id': '{ vector_store_id } ', "
338+ f"{{'vector_store_id': '{ vector_store_id } ', 'batch_id': ' { batch_id } ', "
219339 f"'completed': { batch .file_counts .completed } , 'failed': { batch .file_counts .failed } }}"
220340 )
221- if batch .file_counts .failed > 0 :
222- try :
223- failed_files = self .client .vector_stores .file_batches .list_files (
224- vector_store_id = vector_store_id ,
225- batch_id = batch .id ,
226- filter = "failed" ,
227- )
228- doc_by_file_id = {d .file_id [OPENAI_PROVIDER ]: d for d in docs }
229- parts = []
230- for f in failed_files :
231- d = doc_by_file_id .get (f .id )
232- label = d .fname if d else f .id
233- msg = f .last_error .message if f .last_error else "no error detail"
234- parts .append (f"{ label } : { msg } " )
235- raise RuntimeError ("; " .join (parts ))
236- except OpenAIError as err :
237- logger .warning (
238- f"[OpenAIVectorStoreCrud.update] Could not fetch per-file errors | "
239- f"{{'batch_id': '{ batch .id } ', 'error': '{ str (err )} '}}"
240- )
241- raise
242341
243342 def delete (self , vector_store_id : str , retries : int = 3 ):
244343 if retries < 1 :
0 commit comments