@@ -147,22 +147,22 @@ async def run(self, task: 'Task', session: AsyncSession, resources_provider: Res
147147 for prefix in version_prefixes :
148148 for df , time in self .ingest_prefix (s3 , bucket , f'{ version_path } /{ prefix } ' , version .latest_file_time ,
149149 errors , version .model_id , version .id ):
150+ # For each file, set lock expiry to 240 seconds from now
151+ await lock .extend (240 , replace_ttl = True )
150152 await self .ingestion_backend .log_samples (version , df , session , organization_id , new_scan_time )
151153 version .latest_file_time = max (version .latest_file_time or
152154 pdl .datetime (year = 1970 , month = 1 , day = 1 ), time )
153- # For each file, set lock expiry to 120 seconds from now
154- await lock .extend (120 , replace_ttl = True )
155155
156156 # Ingest labels
157157 for prefix in model_prefixes :
158158 labels_path = f'{ model_path } /labels/{ prefix } '
159159 for df , time in self .ingest_prefix (s3 , bucket , labels_path , model .latest_labels_file_time ,
160160 errors , model_id ):
161+ # For each file, set lock expiry to 240 seconds from now
162+ await lock .extend (240 , replace_ttl = True )
161163 await self .ingestion_backend .log_labels (model , df , session , organization_id )
162164 model .latest_labels_file_time = max (model .latest_labels_file_time
163165 or pdl .datetime (year = 1970 , month = 1 , day = 1 ), time )
164- # For each file, set lock expiry to 120 seconds from now
165- await lock .extend (120 , replace_ttl = True )
166166
167167 model .obj_store_last_scan_time = new_scan_time
168168 except Exception : # pylint: disable=broad-except
0 commit comments