Line data Source code
1 : /******************************************************************************
2 : *
3 : * Project: CPL - Common Portability Library
4 : * Purpose: Implement a write-only file handle using PUT chunked writing
5 : * Author: Even Rouault, even.rouault at spatialys.com
6 : *
7 : ******************************************************************************
8 : * Copyright (c) 2024, Even Rouault <even.rouault at spatialys.com>
9 : *
10 : * SPDX-License-Identifier: MIT
11 : ****************************************************************************/
12 :
13 : #include "cpl_vsil_curl_class.h"
14 :
15 : #ifdef HAVE_CURL
16 :
17 : //! @cond Doxygen_Suppress
18 :
19 : #define unchecked_curl_easy_setopt(handle, opt, param) \
20 : CPL_IGNORE_RET_VAL(curl_easy_setopt(handle, opt, param))
21 :
22 : namespace cpl
23 : {
24 :
25 : /************************************************************************/
26 : /* VSIChunkedWriteHandle() */
27 : /************************************************************************/
28 :
29 3 : VSIChunkedWriteHandle::VSIChunkedWriteHandle(
30 : IVSIS3LikeFSHandler *poFS, const char *pszFilename,
31 3 : IVSIS3LikeHandleHelper *poS3HandleHelper, CSLConstList papszOptions)
32 : : m_poFS(poFS), m_osFilename(pszFilename),
33 : m_poS3HandleHelper(poS3HandleHelper), m_aosOptions(papszOptions),
34 : m_aosHTTPOptions(CPLHTTPGetOptionsFromEnv(pszFilename)),
35 3 : m_oRetryParameters(m_aosHTTPOptions)
36 : {
37 3 : }
38 :
39 : /************************************************************************/
40 : /* ~VSIChunkedWriteHandle() */
41 : /************************************************************************/
42 :
43 6 : VSIChunkedWriteHandle::~VSIChunkedWriteHandle()
44 : {
45 3 : VSIChunkedWriteHandle::Close();
46 3 : delete m_poS3HandleHelper;
47 :
48 3 : if (m_hCurlMulti)
49 : {
50 1 : if (m_hCurl)
51 : {
52 1 : curl_multi_remove_handle(m_hCurlMulti, m_hCurl);
53 1 : curl_easy_cleanup(m_hCurl);
54 : }
55 1 : VSICURLMultiCleanup(m_hCurlMulti);
56 : }
57 3 : CPLFree(m_sWriteFuncHeaderData.pBuffer);
58 6 : }
59 :
60 : /************************************************************************/
61 : /* Close() */
62 : /************************************************************************/
63 :
64 6 : int VSIChunkedWriteHandle::Close()
65 : {
66 6 : int nRet = 0;
67 6 : if (!m_bClosed)
68 : {
69 3 : m_bClosed = true;
70 3 : if (m_hCurlMulti != nullptr)
71 : {
72 1 : nRet = FinishChunkedTransfer();
73 : }
74 : else
75 : {
76 2 : if (!m_bError && !DoEmptyPUT())
77 0 : nRet = -1;
78 : }
79 : }
80 6 : return nRet;
81 : }
82 :
83 : /************************************************************************/
84 : /* InvalidateParentDirectory() */
85 : /************************************************************************/
86 :
87 3 : void VSIChunkedWriteHandle::InvalidateParentDirectory()
88 : {
89 3 : m_poFS->InvalidateCachedData(m_poS3HandleHelper->GetURL().c_str());
90 :
91 3 : std::string osFilenameWithoutSlash(m_osFilename);
92 3 : if (!osFilenameWithoutSlash.empty() && osFilenameWithoutSlash.back() == '/')
93 1 : osFilenameWithoutSlash.pop_back();
94 3 : m_poFS->InvalidateDirContent(
95 6 : CPLGetDirnameSafe(osFilenameWithoutSlash.c_str()));
96 3 : }
97 :
98 : /************************************************************************/
99 : /* Seek() */
100 : /************************************************************************/
101 :
102 0 : int VSIChunkedWriteHandle::Seek(vsi_l_offset nOffset, int nWhence)
103 : {
104 0 : if (!((nWhence == SEEK_SET && nOffset == m_nCurOffset) ||
105 0 : (nWhence == SEEK_CUR && nOffset == 0) ||
106 0 : (nWhence == SEEK_END && nOffset == 0)))
107 : {
108 0 : CPLError(CE_Failure, CPLE_NotSupported,
109 : "Seek not supported on writable %s files",
110 0 : m_poFS->GetFSPrefix().c_str());
111 0 : m_bError = true;
112 0 : return -1;
113 : }
114 0 : return 0;
115 : }
116 :
117 : /************************************************************************/
118 : /* Tell() */
119 : /************************************************************************/
120 :
121 0 : vsi_l_offset VSIChunkedWriteHandle::Tell()
122 : {
123 0 : return m_nCurOffset;
124 : }
125 :
126 : /************************************************************************/
127 : /* Read() */
128 : /************************************************************************/
129 :
130 0 : size_t VSIChunkedWriteHandle::Read(void * /* pBuffer */, size_t /* nBytes */)
131 : {
132 0 : CPLError(CE_Failure, CPLE_NotSupported,
133 : "Read not supported on writable %s files",
134 0 : m_poFS->GetFSPrefix().c_str());
135 0 : m_bError = true;
136 0 : return 0;
137 : }
138 :
139 : /************************************************************************/
140 : /* ReadCallBackBufferChunked() */
141 : /************************************************************************/
142 :
143 3 : size_t VSIChunkedWriteHandle::ReadCallBackBufferChunked(char *buffer,
144 : size_t size,
145 : size_t nitems,
146 : void *instream)
147 : {
148 3 : VSIChunkedWriteHandle *poThis =
149 : static_cast<VSIChunkedWriteHandle *>(instream);
150 3 : if (poThis->m_nChunkedBufferSize == 0)
151 : {
152 : // CPLDebug("VSIChunkedWriteHandle", "Writing 0 byte (finish)");
153 1 : return 0;
154 : }
155 2 : const size_t nSizeMax = size * nitems;
156 2 : size_t nSizeToWrite = nSizeMax;
157 2 : size_t nChunkedBufferRemainingSize =
158 2 : poThis->m_nChunkedBufferSize - poThis->m_nChunkedBufferOff;
159 2 : if (nChunkedBufferRemainingSize < nSizeToWrite)
160 2 : nSizeToWrite = nChunkedBufferRemainingSize;
161 2 : memcpy(buffer,
162 2 : static_cast<const GByte *>(poThis->m_pBuffer) +
163 2 : poThis->m_nChunkedBufferOff,
164 : nSizeToWrite);
165 2 : poThis->m_nChunkedBufferOff += nSizeToWrite;
166 : // CPLDebug("VSIChunkedWriteHandle", "Writing %d bytes", nSizeToWrite);
167 2 : return nSizeToWrite;
168 : }
169 :
170 : /************************************************************************/
171 : /* Write() */
172 : /************************************************************************/
173 :
174 2 : size_t VSIChunkedWriteHandle::Write(const void *pBuffer, size_t nBytes)
175 : {
176 2 : if (m_bError)
177 0 : return 0;
178 :
179 2 : const size_t nBytesToWrite = nBytes;
180 2 : if (nBytesToWrite == 0)
181 0 : return 0;
182 2 : size_t nRet = nBytes;
183 :
184 2 : if (m_hCurlMulti == nullptr)
185 : {
186 1 : m_hCurlMulti = curl_multi_init();
187 : }
188 :
189 2 : WriteFuncStruct sWriteFuncData;
190 4 : CPLHTTPRetryContext oRetryContext(m_oRetryParameters);
191 : // We can only easily retry at the first chunk of a transfer
192 2 : bool bCanRetry = (m_hCurl == nullptr);
193 : bool bRetry;
194 2 : do
195 : {
196 2 : bRetry = false;
197 2 : struct curl_slist *headers = nullptr;
198 2 : if (m_hCurl == nullptr)
199 : {
200 1 : CURL *hCurlHandle = curl_easy_init();
201 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_UPLOAD, 1L);
202 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_READFUNCTION,
203 : ReadCallBackBufferChunked);
204 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_READDATA, this);
205 :
206 1 : VSICURLInitWriteFuncStruct(&sWriteFuncData, nullptr, nullptr,
207 : nullptr);
208 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_WRITEDATA,
209 : &sWriteFuncData);
210 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_WRITEFUNCTION,
211 : VSICurlHandleWriteFunc);
212 :
213 1 : VSICURLInitWriteFuncStruct(&m_sWriteFuncHeaderData, nullptr,
214 : nullptr, nullptr);
215 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_HEADERDATA,
216 : &m_sWriteFuncHeaderData);
217 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_HEADERFUNCTION,
218 : VSICurlHandleWriteFunc);
219 :
220 1 : headers = static_cast<struct curl_slist *>(CPLHTTPSetOptions(
221 1 : hCurlHandle, m_poS3HandleHelper->GetURL().c_str(),
222 1 : m_aosHTTPOptions.List()));
223 2 : headers = VSICurlSetCreationHeadersFromOptions(
224 1 : headers, m_aosOptions.List(), m_osFilename.c_str());
225 1 : headers = m_poS3HandleHelper->GetCurlHeaders("PUT", headers);
226 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_HTTPHEADER,
227 : headers);
228 :
229 1 : m_osCurlErrBuf.resize(CURL_ERROR_SIZE + 1);
230 1 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_ERRORBUFFER,
231 : &m_osCurlErrBuf[0]);
232 :
233 1 : curl_multi_add_handle(m_hCurlMulti, hCurlHandle);
234 1 : m_hCurl = hCurlHandle;
235 : }
236 :
237 2 : m_pBuffer = pBuffer;
238 2 : m_nChunkedBufferOff = 0;
239 2 : m_nChunkedBufferSize = nBytesToWrite;
240 :
241 : // cppcheck-suppress knownConditionTrueFalse
242 4 : while (m_nChunkedBufferOff < m_nChunkedBufferSize && !bRetry)
243 : {
244 : int still_running;
245 :
246 4 : memset(&m_osCurlErrBuf[0], 0, m_osCurlErrBuf.size());
247 :
248 4 : curl_multi_perform(m_hCurlMulti, &still_running);
249 : // cppcheck-suppress knownConditionTrueFalse
250 4 : if (!still_running || m_nChunkedBufferOff == m_nChunkedBufferSize)
251 : break;
252 :
253 : CURLMsg *msg;
254 0 : do
255 : {
256 2 : int msgq = 0;
257 2 : msg = curl_multi_info_read(m_hCurlMulti, &msgq);
258 2 : if (msg && (msg->msg == CURLMSG_DONE))
259 : {
260 0 : CURL *e = msg->easy_handle;
261 0 : if (e == m_hCurl)
262 : {
263 : long response_code;
264 0 : curl_easy_getinfo(m_hCurl, CURLINFO_RESPONSE_CODE,
265 : &response_code);
266 0 : if (response_code != 200 && response_code != 201)
267 : {
268 : // Look if we should attempt a retry
269 0 : if (bCanRetry &&
270 0 : oRetryContext.CanRetry(
271 : static_cast<int>(response_code),
272 0 : m_sWriteFuncHeaderData.pBuffer,
273 : m_osCurlErrBuf.c_str()))
274 : {
275 0 : CPLError(CE_Warning, CPLE_AppDefined,
276 : "HTTP error code: %d - %s. "
277 : "Retrying again in %.1f secs",
278 : static_cast<int>(response_code),
279 0 : m_poS3HandleHelper->GetURL().c_str(),
280 : oRetryContext.GetCurrentDelay());
281 0 : CPLSleep(oRetryContext.GetCurrentDelay());
282 0 : bRetry = true;
283 : }
284 0 : else if (sWriteFuncData.pBuffer != nullptr &&
285 0 : m_poS3HandleHelper->CanRestartOnError(
286 0 : sWriteFuncData.pBuffer,
287 0 : m_sWriteFuncHeaderData.pBuffer, false))
288 : {
289 0 : bRetry = true;
290 : }
291 : else
292 : {
293 0 : CPLError(CE_Failure, CPLE_AppDefined,
294 : "Error %d: %s",
295 : static_cast<int>(response_code),
296 : m_osCurlErrBuf.c_str());
297 :
298 0 : curl_slist_free_all(headers);
299 0 : bRetry = false;
300 : }
301 :
302 0 : curl_multi_remove_handle(m_hCurlMulti, m_hCurl);
303 0 : curl_easy_cleanup(m_hCurl);
304 :
305 0 : CPLFree(sWriteFuncData.pBuffer);
306 0 : CPLFree(m_sWriteFuncHeaderData.pBuffer);
307 :
308 0 : m_hCurl = nullptr;
309 0 : sWriteFuncData.pBuffer = nullptr;
310 0 : m_sWriteFuncHeaderData.pBuffer = nullptr;
311 0 : if (!bRetry)
312 0 : return 0;
313 : }
314 : }
315 : }
316 2 : } while (msg);
317 :
318 2 : CPLMultiPerformWait(m_hCurlMulti);
319 : }
320 :
321 2 : m_nWrittenInPUT += nBytesToWrite;
322 :
323 2 : curl_slist_free_all(headers);
324 :
325 2 : m_pBuffer = nullptr;
326 :
327 2 : if (!bRetry)
328 : {
329 : long response_code;
330 2 : curl_easy_getinfo(m_hCurl, CURLINFO_RESPONSE_CODE, &response_code);
331 2 : if (response_code != 100)
332 : {
333 : // Look if we should attempt a retry
334 0 : if (bCanRetry &&
335 0 : oRetryContext.CanRetry(static_cast<int>(response_code),
336 0 : m_sWriteFuncHeaderData.pBuffer,
337 : m_osCurlErrBuf.c_str()))
338 : {
339 0 : CPLError(CE_Warning, CPLE_AppDefined,
340 : "HTTP error code: %d - %s. "
341 : "Retrying again in %.1f secs",
342 : static_cast<int>(response_code),
343 0 : m_poS3HandleHelper->GetURL().c_str(),
344 : oRetryContext.GetCurrentDelay());
345 0 : CPLSleep(oRetryContext.GetCurrentDelay());
346 0 : bRetry = true;
347 : }
348 0 : else if (sWriteFuncData.pBuffer != nullptr &&
349 0 : m_poS3HandleHelper->CanRestartOnError(
350 0 : sWriteFuncData.pBuffer,
351 0 : m_sWriteFuncHeaderData.pBuffer, false))
352 : {
353 0 : bRetry = true;
354 : }
355 : else
356 : {
357 0 : CPLError(CE_Failure, CPLE_AppDefined, "Error %d: %s",
358 : static_cast<int>(response_code),
359 : m_osCurlErrBuf.c_str());
360 0 : bRetry = false;
361 0 : nRet = 0;
362 : }
363 :
364 0 : curl_multi_remove_handle(m_hCurlMulti, m_hCurl);
365 0 : curl_easy_cleanup(m_hCurl);
366 :
367 0 : CPLFree(sWriteFuncData.pBuffer);
368 0 : CPLFree(m_sWriteFuncHeaderData.pBuffer);
369 :
370 0 : m_hCurl = nullptr;
371 0 : sWriteFuncData.pBuffer = nullptr;
372 0 : m_sWriteFuncHeaderData.pBuffer = nullptr;
373 : }
374 : }
375 : } while (bRetry);
376 :
377 2 : m_nCurOffset += nBytesToWrite;
378 :
379 2 : return nRet;
380 : }
381 :
382 : /************************************************************************/
383 : /* FinishChunkedTransfer() */
384 : /************************************************************************/
385 :
386 1 : int VSIChunkedWriteHandle::FinishChunkedTransfer()
387 : {
388 1 : if (m_hCurl == nullptr)
389 0 : return -1;
390 :
391 2 : NetworkStatisticsFileSystem oContextFS(m_poFS->GetFSPrefix().c_str());
392 2 : NetworkStatisticsFile oContextFile(m_osFilename.c_str());
393 2 : NetworkStatisticsAction oContextAction("Write");
394 :
395 1 : NetworkStatisticsLogger::LogPUT(m_nWrittenInPUT);
396 1 : m_nWrittenInPUT = 0;
397 :
398 1 : m_pBuffer = nullptr;
399 1 : m_nChunkedBufferOff = 0;
400 1 : m_nChunkedBufferSize = 0;
401 :
402 1 : VSICURLMultiPerform(m_hCurlMulti);
403 :
404 : long response_code;
405 1 : curl_easy_getinfo(m_hCurl, CURLINFO_RESPONSE_CODE, &response_code);
406 1 : if (response_code == 200 || response_code == 201)
407 : {
408 1 : InvalidateParentDirectory();
409 : }
410 : else
411 : {
412 0 : CPLError(CE_Failure, CPLE_AppDefined, "Error %d: %s",
413 : static_cast<int>(response_code), m_osCurlErrBuf.c_str());
414 0 : return -1;
415 : }
416 1 : return 0;
417 : }
418 :
419 : /************************************************************************/
420 : /* DoEmptyPUT() */
421 : /************************************************************************/
422 :
423 2 : bool VSIChunkedWriteHandle::DoEmptyPUT()
424 : {
425 2 : bool bSuccess = true;
426 : bool bRetry;
427 4 : CPLHTTPRetryContext oRetryContext(m_oRetryParameters);
428 :
429 4 : NetworkStatisticsFileSystem oContextFS(m_poFS->GetFSPrefix().c_str());
430 4 : NetworkStatisticsFile oContextFile(m_osFilename.c_str());
431 2 : NetworkStatisticsAction oContextAction("Write");
432 :
433 2 : do
434 : {
435 2 : bRetry = false;
436 :
437 2 : PutData putData;
438 2 : putData.pabyData = nullptr;
439 2 : putData.nOff = 0;
440 2 : putData.nTotalSize = 0;
441 :
442 2 : CURL *hCurlHandle = curl_easy_init();
443 2 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_UPLOAD, 1L);
444 2 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_READFUNCTION,
445 : PutData::ReadCallBackBuffer);
446 2 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_READDATA, &putData);
447 2 : unchecked_curl_easy_setopt(hCurlHandle, CURLOPT_INFILESIZE, 0);
448 :
449 : struct curl_slist *headers = static_cast<struct curl_slist *>(
450 2 : CPLHTTPSetOptions(hCurlHandle, m_poS3HandleHelper->GetURL().c_str(),
451 2 : m_aosHTTPOptions.List()));
452 4 : headers = VSICurlSetCreationHeadersFromOptions(
453 2 : headers, m_aosOptions.List(), m_osFilename.c_str());
454 2 : headers = m_poS3HandleHelper->GetCurlHeaders("PUT", headers, "", 0);
455 2 : headers = curl_slist_append(headers, "Expect: 100-continue");
456 :
457 4 : CurlRequestHelper requestHelper;
458 4 : const long response_code = requestHelper.perform(
459 2 : hCurlHandle, headers, m_poFS, m_poS3HandleHelper);
460 :
461 2 : NetworkStatisticsLogger::LogPUT(0);
462 :
463 2 : if (response_code != 200 && response_code != 201)
464 : {
465 : // Look if we should attempt a retry
466 0 : if (oRetryContext.CanRetry(
467 : static_cast<int>(response_code),
468 0 : requestHelper.sWriteFuncHeaderData.pBuffer,
469 : requestHelper.szCurlErrBuf))
470 : {
471 0 : CPLError(CE_Warning, CPLE_AppDefined,
472 : "HTTP error code: %d - %s. "
473 : "Retrying again in %.1f secs",
474 : static_cast<int>(response_code),
475 0 : m_poS3HandleHelper->GetURL().c_str(),
476 : oRetryContext.GetCurrentDelay());
477 0 : CPLSleep(oRetryContext.GetCurrentDelay());
478 0 : bRetry = true;
479 : }
480 0 : else if (requestHelper.sWriteFuncData.pBuffer != nullptr &&
481 0 : m_poS3HandleHelper->CanRestartOnError(
482 0 : requestHelper.sWriteFuncData.pBuffer,
483 0 : requestHelper.sWriteFuncHeaderData.pBuffer, false))
484 : {
485 0 : bRetry = true;
486 : }
487 : else
488 : {
489 0 : CPLDebug("S3", "%s",
490 0 : requestHelper.sWriteFuncData.pBuffer
491 : ? requestHelper.sWriteFuncData.pBuffer
492 : : "(null)");
493 0 : CPLError(CE_Failure, CPLE_AppDefined,
494 : "DoSinglePartPUT of %s failed", m_osFilename.c_str());
495 0 : bSuccess = false;
496 : }
497 : }
498 : else
499 : {
500 2 : InvalidateParentDirectory();
501 : }
502 :
503 2 : if (requestHelper.sWriteFuncHeaderData.pBuffer != nullptr)
504 : {
505 : const char *pzETag =
506 2 : strstr(requestHelper.sWriteFuncHeaderData.pBuffer, "ETag: \"");
507 2 : if (pzETag)
508 : {
509 0 : pzETag += strlen("ETag: \"");
510 0 : const char *pszEndOfETag = strchr(pzETag, '"');
511 0 : if (pszEndOfETag)
512 : {
513 0 : FileProp oFileProp;
514 0 : oFileProp.eExists = EXIST_YES;
515 0 : oFileProp.fileSize = m_nBufferOff;
516 0 : oFileProp.bHasComputedFileSize = true;
517 0 : oFileProp.ETag.assign(pzETag, pszEndOfETag - pzETag);
518 0 : m_poFS->SetCachedFileProp(
519 0 : m_poFS->GetURLFromFilename(m_osFilename.c_str())
520 : .c_str(),
521 : oFileProp);
522 : }
523 : }
524 : }
525 :
526 2 : curl_easy_cleanup(hCurlHandle);
527 : } while (bRetry);
528 4 : return bSuccess;
529 : }
530 :
531 : } // namespace cpl
532 :
533 : //! @endcond
534 :
535 : #endif // HAVE_CURL
|