|
libzypp
9.1.2
|
00001 /*---------------------------------------------------------------------\ 00002 | ____ _ __ __ ___ | 00003 | |__ / \ / / . \ . \ | 00004 | / / \ V /| _/ _/ | 00005 | / /__ | | | | | | | 00006 | /_____||_| |_| |_| | 00007 | | 00008 \---------------------------------------------------------------------*/ 00013 #include <ctype.h> 00014 #include <sys/types.h> 00015 #include <signal.h> 00016 #include <sys/wait.h> 00017 #include <netdb.h> 00018 #include <arpa/inet.h> 00019 00020 #include <vector> 00021 #include <iostream> 00022 #include <algorithm> 00023 00024 00025 #include "zypp/ZConfig.h" 00026 #include "zypp/base/Logger.h" 00027 #include "zypp/media/MediaMultiCurl.h" 00028 #include "zypp/media/MetaLinkParser.h" 00029 00030 using namespace std; 00031 using namespace zypp::base; 00032 00033 namespace zypp { 00034 namespace media { 00035 00036 00038 00039 00040 class multifetchrequest; 00041 00042 // Hack: we derive from MediaCurl just to get the storage space for 00043 // settings, url, curlerrors and the like 00044 00045 class multifetchworker : MediaCurl { 00046 friend class multifetchrequest; 00047 00048 public: 00049 multifetchworker(int no, multifetchrequest &request, const Url &url); 00050 ~multifetchworker(); 00051 void nextjob(); 00052 void run(); 00053 bool checkChecksum(); 00054 bool recheckChecksum(); 00055 void disableCompetition(); 00056 00057 void checkdns(); 00058 void adddnsfd(fd_set &rset, int &maxfd); 00059 void dnsevent(fd_set &rset); 00060 00061 int _workerno; 00062 00063 int _state; 00064 bool _competing; 00065 00066 size_t _blkno; 00067 off_t _blkstart; 00068 size_t _blksize; 00069 bool _noendrange; 00070 00071 double _blkstarttime; 00072 size_t _blkreceived; 00073 off_t _received; 00074 00075 double _avgspeed; 00076 double _maxspeed; 00077 00078 double _sleepuntil; 00079 00080 private: 00081 void stealjob(); 00082 00083 size_t writefunction(void *ptr, size_t size); 00084 static size_t _writefunction(void *ptr, size_t size, size_t nmemb, void *stream); 00085 00086 size_t headerfunction(char *ptr, size_t size); 00087 static size_t _headerfunction(void *ptr, size_t size, size_t nmemb, void *stream); 00088 00089 multifetchrequest *_request; 00090 int _pass; 00091 string _urlbuf; 00092 off_t _off; 00093 size_t _size; 00094 Digest _dig; 00095 00096 pid_t _pid; 00097 int _dnspipe; 00098 }; 00099 00100 #define WORKER_STARTING 0 00101 #define WORKER_LOOKUP 1 00102 #define WORKER_FETCH 2 00103 #define WORKER_DISCARD 3 00104 #define WORKER_DONE 4 00105 #define WORKER_SLEEP 5 00106 #define WORKER_BROKEN 6 00107 00108 00109 00110 class multifetchrequest { 00111 public: 00112 multifetchrequest(const MediaMultiCurl *context, const Pathname &filename, const Url &baseurl, CURLM *multi, FILE *fp, callback::SendReport<DownloadProgressReport> *report, MediaBlockList *blklist, off_t filesize); 00113 ~multifetchrequest(); 00114 00115 void run(std::vector<Url> &urllist); 00116 00117 protected: 00118 friend class multifetchworker; 00119 00120 const MediaMultiCurl *_context; 00121 const Pathname _filename; 00122 Url _baseurl; 00123 00124 FILE *_fp; 00125 callback::SendReport<DownloadProgressReport> *_report; 00126 MediaBlockList *_blklist; 00127 off_t _filesize; 00128 00129 CURLM *_multi; 00130 00131 std::list<multifetchworker *> _workers; 00132 bool _stealing; 00133 bool _havenewjob; 00134 00135 size_t _blkno; 00136 off_t _blkoff; 00137 size_t _activeworkers; 00138 size_t _lookupworkers; 00139 size_t _sleepworkers; 00140 double _minsleepuntil; 00141 bool _finished; 00142 off_t _totalsize; 00143 off_t _fetchedsize; 00144 off_t _fetchedgoodsize; 00145 00146 double _starttime; 00147 double _lastprogress; 00148 00149 double _lastperiodstart; 00150 double _lastperiodfetched; 00151 double _periodavg; 00152 00153 public: 00154 double _timeout; 00155 double _connect_timeout; 00156 double _maxspeed; 00157 int _maxworkers; 00158 }; 00159 00160 #define BLKSIZE 131072 00161 #define MAXURLS 10 00162 00163 00165 00166 static double 00167 currentTime() 00168 { 00169 struct timeval tv; 00170 if (gettimeofday(&tv, NULL)) 00171 return 0; 00172 return tv.tv_sec + tv.tv_usec / 1000000.; 00173 } 00174 00175 size_t 00176 multifetchworker::writefunction(void *ptr, size_t size) 00177 { 00178 size_t len, cnt; 00179 if (_state == WORKER_BROKEN) 00180 return size ? 0 : 1; 00181 00182 double now = currentTime(); 00183 00184 len = size > _size ? _size : size; 00185 if (!len) 00186 { 00187 // kill this job? 00188 return size; 00189 } 00190 00191 if (_blkstart && _off == _blkstart) 00192 { 00193 // make sure that the server replied with "partial content" 00194 // for http requests 00195 char *effurl; 00196 (void)curl_easy_getinfo(_curl, CURLINFO_EFFECTIVE_URL, &effurl); 00197 if (effurl && !strncasecmp(effurl, "http", 4)) 00198 { 00199 long statuscode = 0; 00200 (void)curl_easy_getinfo(_curl, CURLINFO_RESPONSE_CODE, &statuscode); 00201 if (statuscode != 206) 00202 return size ? 0 : 1; 00203 } 00204 } 00205 00206 _blkreceived += len; 00207 _received += len; 00208 00209 _request->_lastprogress = now; 00210 00211 if (_state == WORKER_DISCARD || !_request->_fp) 00212 { 00213 // block is no longer needed 00214 // still calculate the checksum so that we can throw out bad servers 00215 if (_request->_blklist) 00216 _dig.update((const char *)ptr, len); 00217 _off += len; 00218 _size -= len; 00219 return size; 00220 } 00221 if (fseeko(_request->_fp, _off, SEEK_SET)) 00222 return size ? 0 : 1; 00223 cnt = fwrite(ptr, 1, len, _request->_fp); 00224 if (cnt > 0) 00225 { 00226 _request->_fetchedsize += cnt; 00227 if (_request->_blklist) 00228 _dig.update((const char *)ptr, cnt); 00229 _off += cnt; 00230 _size -= cnt; 00231 if (cnt == len) 00232 return size; 00233 } 00234 return cnt; 00235 } 00236 00237 size_t 00238 multifetchworker::_writefunction(void *ptr, size_t size, size_t nmemb, void *stream) 00239 { 00240 multifetchworker *me = reinterpret_cast<multifetchworker *>(stream); 00241 return me->writefunction(ptr, size * nmemb); 00242 } 00243 00244 size_t 00245 multifetchworker::headerfunction(char *p, size_t size) 00246 { 00247 size_t l = size; 00248 if (l > 9 && !strncasecmp(p, "Location:", 9)) 00249 { 00250 string line(p + 9, l - 9); 00251 if (line[l - 10] == '\r') 00252 line.erase(l - 10, 1); 00253 DBG << "#" << _workerno << ": redirecting to" << line << endl; 00254 return size; 00255 } 00256 if (l <= 14 || l >= 128 || strncasecmp(p, "Content-Range:", 14) != 0) 00257 return size; 00258 p += 14; 00259 l -= 14; 00260 while (l && (*p == ' ' || *p == '\t')) 00261 p++, l--; 00262 if (l < 6 || strncasecmp(p, "bytes", 5)) 00263 return size; 00264 p += 5; 00265 l -= 5; 00266 char buf[128]; 00267 memcpy(buf, p, l); 00268 buf[l] = 0; 00269 unsigned long long start, off, filesize; 00270 if (sscanf(buf, "%llu-%llu/%llu", &start, &off, &filesize) != 3) 00271 return size; 00272 if (_request->_filesize == (off_t)-1) 00273 { 00274 WAR << "#" << _workerno << ": setting request filesize to " << filesize << endl; 00275 _request->_filesize = filesize; 00276 if (_request->_totalsize == 0 && !_request->_blklist) 00277 _request->_totalsize = filesize; 00278 } 00279 if (_request->_filesize != (off_t)filesize) 00280 { 00281 DBG << "#" << _workerno << ": filesize mismatch" << endl; 00282 _state = WORKER_BROKEN; 00283 strncpy(_curlError, "filesize mismatch", CURL_ERROR_SIZE); 00284 } 00285 return size; 00286 } 00287 00288 size_t 00289 multifetchworker::_headerfunction(void *ptr, size_t size, size_t nmemb, void *stream) 00290 { 00291 multifetchworker *me = reinterpret_cast<multifetchworker *>(stream); 00292 return me->headerfunction((char *)ptr, size * nmemb); 00293 } 00294 00295 multifetchworker::multifetchworker(int no, multifetchrequest &request, const Url &url) 00296 : MediaCurl(url, Pathname()) 00297 { 00298 _workerno = no; 00299 _request = &request; 00300 _state = WORKER_STARTING; 00301 _competing = false; 00302 _off = _blkstart = 0; 00303 _size = _blksize = 0; 00304 _pass = 0; 00305 _blkno = 0; 00306 _pid = 0; 00307 _dnspipe = -1; 00308 _blkreceived = 0; 00309 _received = 0; 00310 _blkstarttime = 0; 00311 _avgspeed = 0; 00312 _sleepuntil = 0; 00313 _maxspeed = _request->_maxspeed; 00314 _noendrange = false; 00315 00316 Url curlUrl( clearQueryString(url) ); 00317 _urlbuf = curlUrl.asString(); 00318 _curl = _request->_context->fromEasyPool(_url.getHost()); 00319 if (_curl) 00320 DBG << "reused worker from pool" << endl; 00321 if (!_curl && !(_curl = curl_easy_init())) 00322 { 00323 _state = WORKER_BROKEN; 00324 strncpy(_curlError, "curl_easy_init failed", CURL_ERROR_SIZE); 00325 return; 00326 } 00327 try 00328 { 00329 setupEasy(); 00330 } 00331 catch (Exception &ex) 00332 { 00333 curl_easy_cleanup(_curl); 00334 _curl = 0; 00335 _state = WORKER_BROKEN; 00336 strncpy(_curlError, "curl_easy_setopt failed", CURL_ERROR_SIZE); 00337 return; 00338 } 00339 curl_easy_setopt(_curl, CURLOPT_PRIVATE, this); 00340 curl_easy_setopt(_curl, CURLOPT_URL, _urlbuf.c_str()); 00341 curl_easy_setopt(_curl, CURLOPT_WRITEFUNCTION, &_writefunction); 00342 curl_easy_setopt(_curl, CURLOPT_WRITEDATA, this); 00343 if (_request->_filesize == off_t(-1) || !_request->_blklist || !_request->_blklist->haveChecksum(0)) 00344 { 00345 curl_easy_setopt(_curl, CURLOPT_HEADERFUNCTION, &_headerfunction); 00346 curl_easy_setopt(_curl, CURLOPT_HEADERDATA, this); 00347 } 00348 // if this is the same host copy authorization 00349 // (the host check is also what curl does when doing a redirect) 00350 // (note also that unauthorized exceptions are thrown with the request host) 00351 if (url.getHost() == _request->_context->_url.getHost()) 00352 { 00353 _settings.setUsername(_request->_context->_settings.username()); 00354 _settings.setPassword(_request->_context->_settings.password()); 00355 _settings.setAuthType(_request->_context->_settings.authType()); 00356 if ( _settings.userPassword().size() ) 00357 { 00358 curl_easy_setopt(_curl, CURLOPT_USERPWD, _settings.userPassword().c_str()); 00359 string use_auth = _settings.authType(); 00360 if (use_auth.empty()) 00361 use_auth = "digest,basic"; // our default 00362 long auth = CurlAuthData::auth_type_str2long(use_auth); 00363 if( auth != CURLAUTH_NONE) 00364 { 00365 DBG << "#" << _workerno << ": Enabling HTTP authentication methods: " << use_auth 00366 << " (CURLOPT_HTTPAUTH=" << auth << ")" << std::endl; 00367 curl_easy_setopt(_curl, CURLOPT_HTTPAUTH, auth); 00368 } 00369 } 00370 } 00371 checkdns(); 00372 } 00373 00374 multifetchworker::~multifetchworker() 00375 { 00376 if (_curl) 00377 { 00378 if (_state == WORKER_FETCH || _state == WORKER_DISCARD) 00379 curl_multi_remove_handle(_request->_multi, _curl); 00380 if (_state == WORKER_DONE || _state == WORKER_SLEEP) 00381 { 00382 #if LIBCURL_VERSION_NUMBER >= 0x071505 00383 curl_easy_setopt(_curl, CURLOPT_MAX_RECV_SPEED_LARGE, (curl_off_t)0); 00384 #endif 00385 curl_easy_setopt(_curl, CURLOPT_PRIVATE, (void *)0); 00386 curl_easy_setopt(_curl, CURLOPT_WRITEFUNCTION, (void *)0); 00387 curl_easy_setopt(_curl, CURLOPT_WRITEDATA, (void *)0); 00388 curl_easy_setopt(_curl, CURLOPT_HEADERFUNCTION, (void *)0); 00389 curl_easy_setopt(_curl, CURLOPT_HEADERDATA, (void *)0); 00390 _request->_context->toEasyPool(_url.getHost(), _curl); 00391 } 00392 else 00393 curl_easy_cleanup(_curl); 00394 _curl = 0; 00395 } 00396 if (_pid) 00397 { 00398 kill(_pid, SIGKILL); 00399 int status; 00400 while (waitpid(_pid, &status, 0) == -1) 00401 if (errno != EINTR) 00402 break; 00403 _pid = 0; 00404 } 00405 if (_dnspipe != -1) 00406 { 00407 close(_dnspipe); 00408 _dnspipe = -1; 00409 } 00410 // the destructor in MediaCurl doesn't call disconnect() if 00411 // the media is not attached, so we do it here manually 00412 disconnectFrom(); 00413 } 00414 00415 static inline bool env_isset(string name) 00416 { 00417 const char *s = getenv(name.c_str()); 00418 return s && *s ? true : false; 00419 } 00420 00421 void 00422 multifetchworker::checkdns() 00423 { 00424 string host = _url.getHost(); 00425 00426 if (host.empty()) 00427 return; 00428 00429 if (_request->_context->isDNSok(host)) 00430 return; 00431 00432 // no need to do dns checking for numeric hosts 00433 char addrbuf[128]; 00434 if (inet_pton(AF_INET, host.c_str(), addrbuf) == 1) 00435 return; 00436 if (inet_pton(AF_INET6, host.c_str(), addrbuf) == 1) 00437 return; 00438 00439 // no need to do dns checking if we use a proxy 00440 if (!_settings.proxy().empty()) 00441 return; 00442 if (env_isset("all_proxy") || env_isset("ALL_PROXY")) 00443 return; 00444 string schemeproxy = _url.getScheme() + "_proxy"; 00445 if (env_isset(schemeproxy)) 00446 return; 00447 if (schemeproxy != "http_proxy") 00448 { 00449 std::transform(schemeproxy.begin(), schemeproxy.end(), schemeproxy.begin(), ::toupper); 00450 if (env_isset(schemeproxy)) 00451 return; 00452 } 00453 00454 DBG << "checking DNS lookup of " << host << endl; 00455 int pipefds[2]; 00456 if (pipe(pipefds)) 00457 { 00458 _state = WORKER_BROKEN; 00459 strncpy(_curlError, "DNS pipe creation failed", CURL_ERROR_SIZE); 00460 return; 00461 } 00462 _pid = fork(); 00463 if (_pid == pid_t(-1)) 00464 { 00465 close(pipefds[0]); 00466 close(pipefds[1]); 00467 _pid = 0; 00468 _state = WORKER_BROKEN; 00469 strncpy(_curlError, "DNS checker fork failed", CURL_ERROR_SIZE); 00470 return; 00471 } 00472 else if (_pid == 0) 00473 { 00474 close(pipefds[0]); 00475 // XXX: close all other file descriptors 00476 struct addrinfo *ai, aihints; 00477 memset(&aihints, 0, sizeof(aihints)); 00478 aihints.ai_family = PF_UNSPEC; 00479 int tstsock = socket(PF_INET6, SOCK_DGRAM, 0); 00480 if (tstsock == -1) 00481 aihints.ai_family = PF_INET; 00482 else 00483 close(tstsock); 00484 aihints.ai_socktype = SOCK_STREAM; 00485 aihints.ai_flags = AI_CANONNAME; 00486 unsigned int connecttimeout = _request->_connect_timeout; 00487 if (connecttimeout) 00488 alarm(connecttimeout); 00489 signal(SIGALRM, SIG_DFL); 00490 if (getaddrinfo(host.c_str(), NULL, &aihints, &ai)) 00491 _exit(1); 00492 _exit(0); 00493 } 00494 close(pipefds[1]); 00495 _dnspipe = pipefds[0]; 00496 _state = WORKER_LOOKUP; 00497 } 00498 00499 void 00500 multifetchworker::adddnsfd(fd_set &rset, int &maxfd) 00501 { 00502 if (_state != WORKER_LOOKUP) 00503 return; 00504 FD_SET(_dnspipe, &rset); 00505 if (maxfd < _dnspipe) 00506 maxfd = _dnspipe; 00507 } 00508 00509 void 00510 multifetchworker::dnsevent(fd_set &rset) 00511 { 00512 00513 if (_state != WORKER_LOOKUP || !FD_ISSET(_dnspipe, &rset)) 00514 return; 00515 int status; 00516 while (waitpid(_pid, &status, 0) == -1) 00517 { 00518 if (errno != EINTR) 00519 return; 00520 } 00521 _pid = 0; 00522 if (_dnspipe != -1) 00523 { 00524 close(_dnspipe); 00525 _dnspipe = -1; 00526 } 00527 if (!WIFEXITED(status)) 00528 { 00529 _state = WORKER_BROKEN; 00530 strncpy(_curlError, "DNS lookup failed", CURL_ERROR_SIZE); 00531 _request->_activeworkers--; 00532 return; 00533 } 00534 int exitcode = WEXITSTATUS(status); 00535 DBG << "#" << _workerno << ": DNS lookup returned " << exitcode << endl; 00536 if (exitcode != 0) 00537 { 00538 _state = WORKER_BROKEN; 00539 strncpy(_curlError, "DNS lookup failed", CURL_ERROR_SIZE); 00540 _request->_activeworkers--; 00541 return; 00542 } 00543 _request->_context->setDNSok(_url.getHost()); 00544 nextjob(); 00545 } 00546 00547 bool 00548 multifetchworker::checkChecksum() 00549 { 00550 // DBG << "checkChecksum block " << _blkno << endl; 00551 if (!_blksize || !_request->_blklist) 00552 return true; 00553 return _request->_blklist->verifyDigest(_blkno, _dig); 00554 } 00555 00556 bool 00557 multifetchworker::recheckChecksum() 00558 { 00559 // DBG << "recheckChecksum block " << _blkno << endl; 00560 if (!_request->_fp || !_blksize || !_request->_blklist) 00561 return true; 00562 if (fseeko(_request->_fp, _blkstart, SEEK_SET)) 00563 return false; 00564 char buf[4096]; 00565 size_t l = _blksize; 00566 _request->_blklist->createDigest(_dig); // resets digest 00567 while (l) 00568 { 00569 size_t cnt = l > sizeof(buf) ? sizeof(buf) : l; 00570 if (fread(buf, cnt, 1, _request->_fp) != 1) 00571 return false; 00572 _dig.update(buf, cnt); 00573 l -= cnt; 00574 } 00575 return _request->_blklist->verifyDigest(_blkno, _dig); 00576 } 00577 00578 00579 void 00580 multifetchworker::stealjob() 00581 { 00582 if (!_request->_stealing) 00583 { 00584 DBG << "start stealing!" << endl; 00585 _request->_stealing = true; 00586 } 00587 multifetchworker *best = 0; 00588 std::list<multifetchworker *>::iterator workeriter = _request->_workers.begin(); 00589 double now = 0; 00590 for (; workeriter != _request->_workers.end(); ++workeriter) 00591 { 00592 multifetchworker *worker = *workeriter; 00593 if (worker == this) 00594 continue; 00595 if (worker->_pass == -1) 00596 continue; // do not steal! 00597 if (worker->_state == WORKER_DISCARD || worker->_state == WORKER_DONE || worker->_state == WORKER_SLEEP || !worker->_blksize) 00598 continue; // do not steal finished jobs 00599 if (!worker->_avgspeed && worker->_blkreceived) 00600 { 00601 if (!now) 00602 now = currentTime(); 00603 if (now > worker->_blkstarttime) 00604 worker->_avgspeed = worker->_blkreceived / (now - worker->_blkstarttime); 00605 } 00606 if (!best || best->_pass > worker->_pass) 00607 { 00608 best = worker; 00609 continue; 00610 } 00611 if (best->_pass < worker->_pass) 00612 continue; 00613 // if it is the same block, we want to know the best worker, otherwise the worst 00614 if (worker->_blkstart == best->_blkstart) 00615 { 00616 if ((worker->_blksize - worker->_blkreceived) * best->_avgspeed < (best->_blksize - best->_blkreceived) * worker->_avgspeed) 00617 best = worker; 00618 } 00619 else 00620 { 00621 if ((worker->_blksize - worker->_blkreceived) * best->_avgspeed > (best->_blksize - best->_blkreceived) * worker->_avgspeed) 00622 best = worker; 00623 } 00624 } 00625 if (!best) 00626 { 00627 _state = WORKER_DONE; 00628 _request->_activeworkers--; 00629 _request->_finished = true; 00630 return; 00631 } 00632 // do not sleep twice 00633 if (_state != WORKER_SLEEP) 00634 { 00635 if (!_avgspeed && _blkreceived) 00636 { 00637 if (!now) 00638 now = currentTime(); 00639 if (now > _blkstarttime) 00640 _avgspeed = _blkreceived / (now - _blkstarttime); 00641 } 00642 00643 // lets see if we should sleep a bit 00644 DBG << "me #" << _workerno << ": " << _avgspeed << ", size " << best->_blksize << endl; 00645 DBG << "best #" << best->_workerno << ": " << best->_avgspeed << ", size " << (best->_blksize - best->_blkreceived) << endl; 00646 if (_avgspeed && best->_avgspeed && best->_blksize - best->_blkreceived > 0 && 00647 (best->_blksize - best->_blkreceived) * _avgspeed < best->_blksize * best->_avgspeed) 00648 { 00649 if (!now) 00650 now = currentTime(); 00651 double sl = (best->_blksize - best->_blkreceived) / best->_avgspeed * 2; 00652 if (sl > 1) 00653 sl = 1; 00654 DBG << "#" << _workerno << ": going to sleep for " << sl * 1000 << " ms" << endl; 00655 _sleepuntil = now + sl; 00656 _state = WORKER_SLEEP; 00657 _request->_sleepworkers++; 00658 return; 00659 } 00660 } 00661 00662 _competing = true; 00663 best->_competing = true; 00664 _blkstart = best->_blkstart; 00665 _blksize = best->_blksize; 00666 best->_pass++; 00667 _pass = best->_pass; 00668 _blkno = best->_blkno; 00669 run(); 00670 } 00671 00672 void 00673 multifetchworker::disableCompetition() 00674 { 00675 std::list<multifetchworker *>::iterator workeriter = _request->_workers.begin(); 00676 for (; workeriter != _request->_workers.end(); ++workeriter) 00677 { 00678 multifetchworker *worker = *workeriter; 00679 if (worker == this) 00680 continue; 00681 if (worker->_blkstart == _blkstart) 00682 { 00683 if (worker->_state == WORKER_FETCH) 00684 worker->_state = WORKER_DISCARD; 00685 worker->_pass = -1; /* do not steal this one, we already have it */ 00686 } 00687 } 00688 } 00689 00690 00691 void 00692 multifetchworker::nextjob() 00693 { 00694 _noendrange = false; 00695 if (_request->_stealing) 00696 { 00697 stealjob(); 00698 return; 00699 } 00700 00701 MediaBlockList *blklist = _request->_blklist; 00702 if (!blklist) 00703 { 00704 _blksize = BLKSIZE; 00705 if (_request->_filesize != off_t(-1)) 00706 { 00707 if (_request->_blkoff >= _request->_filesize) 00708 { 00709 stealjob(); 00710 return; 00711 } 00712 _blksize = _request->_filesize - _request->_blkoff; 00713 if (_blksize > BLKSIZE) 00714 _blksize = BLKSIZE; 00715 } 00716 } 00717 else 00718 { 00719 MediaBlock blk = blklist->getBlock(_request->_blkno); 00720 while (_request->_blkoff >= blk.off + blk.size) 00721 { 00722 if (++_request->_blkno == blklist->numBlocks()) 00723 { 00724 stealjob(); 00725 return; 00726 } 00727 blk = blklist->getBlock(_request->_blkno); 00728 _request->_blkoff = blk.off; 00729 } 00730 _blksize = blk.off + blk.size - _request->_blkoff; 00731 if (_blksize > BLKSIZE && !blklist->haveChecksum(_request->_blkno)) 00732 _blksize = BLKSIZE; 00733 } 00734 _blkno = _request->_blkno; 00735 _blkstart = _request->_blkoff; 00736 _request->_blkoff += _blksize; 00737 run(); 00738 } 00739 00740 void 00741 multifetchworker::run() 00742 { 00743 char rangebuf[128]; 00744 00745 if (_state == WORKER_BROKEN || _state == WORKER_DONE) 00746 return; // just in case... 00747 if (_noendrange) 00748 sprintf(rangebuf, "%llu-", (unsigned long long)_blkstart); 00749 else 00750 sprintf(rangebuf, "%llu-%llu", (unsigned long long)_blkstart, (unsigned long long)_blkstart + _blksize - 1); 00751 DBG << "#" << _workerno << ": BLK " << _blkno << ":" << rangebuf << " " << _url << endl; 00752 if (curl_easy_setopt(_curl, CURLOPT_RANGE, !_noendrange || _blkstart != 0 ? rangebuf : (char *)0) != CURLE_OK) 00753 { 00754 _request->_activeworkers--; 00755 _state = WORKER_BROKEN; 00756 strncpy(_curlError, "curl_easy_setopt range failed", CURL_ERROR_SIZE); 00757 return; 00758 } 00759 if (curl_multi_add_handle(_request->_multi, _curl) != CURLM_OK) 00760 { 00761 _request->_activeworkers--; 00762 _state = WORKER_BROKEN; 00763 strncpy(_curlError, "curl_multi_add_handle failed", CURL_ERROR_SIZE); 00764 return; 00765 } 00766 _request->_havenewjob = true; 00767 _off = _blkstart; 00768 _size = _blksize; 00769 if (_request->_blklist) 00770 _request->_blklist->createDigest(_dig); // resets digest 00771 _state = WORKER_FETCH; 00772 00773 double now = currentTime(); 00774 _blkstarttime = now; 00775 _blkreceived = 0; 00776 } 00777 00778 00780 00781 00782 multifetchrequest::multifetchrequest(const MediaMultiCurl *context, const Pathname &filename, const Url &baseurl, CURLM *multi, FILE *fp, callback::SendReport<DownloadProgressReport> *report, MediaBlockList *blklist, off_t filesize) : _context(context), _filename(filename), _baseurl(baseurl) 00783 { 00784 _fp = fp; 00785 _report = report; 00786 _blklist = blklist; 00787 _filesize = filesize; 00788 _multi = multi; 00789 _stealing = false; 00790 _havenewjob = false; 00791 _blkno = 0; 00792 if (_blklist) 00793 _blkoff = _blklist->getBlock(0).off; 00794 else 00795 _blkoff = 0; 00796 _activeworkers = 0; 00797 _lookupworkers = 0; 00798 _sleepworkers = 0; 00799 _minsleepuntil = 0; 00800 _finished = false; 00801 _fetchedsize = 0; 00802 _fetchedgoodsize = 0; 00803 _totalsize = 0; 00804 _lastperiodstart = _lastprogress = _starttime = currentTime(); 00805 _lastperiodfetched = 0; 00806 _periodavg = 0; 00807 _timeout = 0; 00808 _connect_timeout = 0; 00809 _maxspeed = 0; 00810 _maxworkers = 0; 00811 if (blklist) 00812 { 00813 for (size_t blkno = 0; blkno < blklist->numBlocks(); blkno++) 00814 { 00815 MediaBlock blk = blklist->getBlock(blkno); 00816 _totalsize += blk.size; 00817 } 00818 } 00819 else if (filesize != off_t(-1)) 00820 _totalsize = filesize; 00821 } 00822 00823 multifetchrequest::~multifetchrequest() 00824 { 00825 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 00826 { 00827 multifetchworker *worker = *workeriter; 00828 *workeriter = NULL; 00829 delete worker; 00830 } 00831 _workers.clear(); 00832 } 00833 00834 void 00835 multifetchrequest::run(std::vector<Url> &urllist) 00836 { 00837 int workerno = 0; 00838 std::vector<Url>::iterator urliter = urllist.begin(); 00839 for (;;) 00840 { 00841 fd_set rset, wset, xset; 00842 int maxfd, nqueue; 00843 00844 if (_finished) 00845 { 00846 DBG << "finished!" << endl; 00847 break; 00848 } 00849 00850 if (_activeworkers < _maxworkers && urliter != urllist.end() && _workers.size() < MAXURLS) 00851 { 00852 // spawn another worker! 00853 multifetchworker *worker = new multifetchworker(workerno++, *this, *urliter); 00854 _workers.push_back(worker); 00855 if (worker->_state != WORKER_BROKEN) 00856 { 00857 _activeworkers++; 00858 if (worker->_state != WORKER_LOOKUP) 00859 { 00860 worker->nextjob(); 00861 } 00862 else 00863 _lookupworkers++; 00864 } 00865 ++urliter; 00866 continue; 00867 } 00868 if (!_activeworkers) 00869 { 00870 WAR << "No more active workers!" << endl; 00871 // show the first worker error we find 00872 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 00873 { 00874 if ((*workeriter)->_state != WORKER_BROKEN) 00875 continue; 00876 ZYPP_THROW(MediaCurlException(_baseurl, "Server error", (*workeriter)->_curlError)); 00877 } 00878 break; 00879 } 00880 00881 FD_ZERO(&rset); 00882 FD_ZERO(&wset); 00883 FD_ZERO(&xset); 00884 00885 curl_multi_fdset(_multi, &rset, &wset, &xset, &maxfd); 00886 00887 if (_lookupworkers) 00888 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 00889 (*workeriter)->adddnsfd(rset, maxfd); 00890 00891 timeval tv; 00892 // if we added a new job we have to call multi_perform once 00893 // to make it show up in the fd set. do not sleep in this case. 00894 tv.tv_sec = 0; 00895 tv.tv_usec = _havenewjob ? 0 : 200000; 00896 if (_sleepworkers && !_havenewjob) 00897 { 00898 if (_minsleepuntil == 0) 00899 { 00900 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 00901 { 00902 multifetchworker *worker = *workeriter; 00903 if (worker->_state != WORKER_SLEEP) 00904 continue; 00905 if (!_minsleepuntil || _minsleepuntil > worker->_sleepuntil) 00906 _minsleepuntil = worker->_sleepuntil; 00907 } 00908 } 00909 double sl = _minsleepuntil - currentTime(); 00910 if (sl < 0) 00911 { 00912 sl = 0; 00913 _minsleepuntil = 0; 00914 } 00915 if (sl < .2) 00916 tv.tv_usec = sl * 1000000; 00917 } 00918 int r = select(maxfd + 1, &rset, &wset, &xset, &tv); 00919 if (r == -1 && errno != EINTR) 00920 ZYPP_THROW(MediaCurlException(_baseurl, "select() failed", "unknown error")); 00921 if (r != 0 && _lookupworkers) 00922 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 00923 { 00924 multifetchworker *worker = *workeriter; 00925 if (worker->_state != WORKER_LOOKUP) 00926 continue; 00927 (*workeriter)->dnsevent(rset); 00928 if (worker->_state != WORKER_LOOKUP) 00929 _lookupworkers--; 00930 } 00931 _havenewjob = false; 00932 00933 // run curl 00934 for (;;) 00935 { 00936 CURLMcode mcode; 00937 int tasks; 00938 mcode = curl_multi_perform(_multi, &tasks); 00939 if (mcode == CURLM_CALL_MULTI_PERFORM) 00940 continue; 00941 if (mcode != CURLM_OK) 00942 ZYPP_THROW(MediaCurlException(_baseurl, "curl_multi_perform", "unknown error")); 00943 break; 00944 } 00945 00946 double now = currentTime(); 00947 00948 // update periodavg 00949 if (now > _lastperiodstart + .5) 00950 { 00951 if (!_periodavg) 00952 _periodavg = (_fetchedsize - _lastperiodfetched) / (now - _lastperiodstart); 00953 else 00954 _periodavg = (_periodavg + (_fetchedsize - _lastperiodfetched) / (now - _lastperiodstart)) / 2; 00955 _lastperiodfetched = _fetchedsize; 00956 _lastperiodstart = now; 00957 } 00958 00959 // wake up sleepers 00960 if (_sleepworkers) 00961 { 00962 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 00963 { 00964 multifetchworker *worker = *workeriter; 00965 if (worker->_state != WORKER_SLEEP) 00966 continue; 00967 if (worker->_sleepuntil > now) 00968 continue; 00969 if (_minsleepuntil == worker->_sleepuntil) 00970 _minsleepuntil = 0; 00971 DBG << "#" << worker->_workerno << ": sleep done, wake up" << endl; 00972 _sleepworkers--; 00973 // nextjob chnages the state 00974 worker->nextjob(); 00975 } 00976 } 00977 00978 // collect all curl results, reschedule new jobs 00979 CURLMsg *msg; 00980 while ((msg = curl_multi_info_read(_multi, &nqueue)) != 0) 00981 { 00982 if (msg->msg != CURLMSG_DONE) 00983 continue; 00984 CURL *easy = msg->easy_handle; 00985 CURLcode cc = msg->data.result; 00986 multifetchworker *worker; 00987 if (curl_easy_getinfo(easy, CURLINFO_PRIVATE, &worker) != CURLE_OK) 00988 ZYPP_THROW(MediaCurlException(_baseurl, "curl_easy_getinfo", "unknown error")); 00989 if (worker->_blkreceived && now > worker->_blkstarttime) 00990 { 00991 if (worker->_avgspeed) 00992 worker->_avgspeed = (worker->_avgspeed + worker->_blkreceived / (now - worker->_blkstarttime)) / 2; 00993 else 00994 worker->_avgspeed = worker->_blkreceived / (now - worker->_blkstarttime); 00995 } 00996 DBG << "#" << worker->_workerno << ": BLK " << worker->_blkno << " done code " << cc << " speed " << worker->_avgspeed << endl; 00997 curl_multi_remove_handle(_multi, easy); 00998 if (cc == CURLE_HTTP_RETURNED_ERROR) 00999 { 01000 long statuscode = 0; 01001 (void)curl_easy_getinfo(easy, CURLINFO_RESPONSE_CODE, &statuscode); 01002 DBG << "HTTP status " << statuscode << endl; 01003 if (statuscode == 416 && !_blklist) /* Range error */ 01004 { 01005 if (_filesize == off_t(-1)) 01006 { 01007 if (!worker->_noendrange) 01008 { 01009 DBG << "#" << worker->_workerno << ": retrying with no end range" << endl; 01010 worker->_noendrange = true; 01011 worker->run(); 01012 continue; 01013 } 01014 worker->_noendrange = false; 01015 worker->stealjob(); 01016 continue; 01017 } 01018 if (worker->_blkstart >= _filesize) 01019 { 01020 worker->nextjob(); 01021 continue; 01022 } 01023 } 01024 } 01025 if (cc == 0) 01026 { 01027 if (!worker->checkChecksum()) 01028 { 01029 WAR << "#" << worker->_workerno << ": checksum error, disable worker" << endl; 01030 worker->_state = WORKER_BROKEN; 01031 strncpy(worker->_curlError, "checksum error", CURL_ERROR_SIZE); 01032 _activeworkers--; 01033 continue; 01034 } 01035 if (worker->_state == WORKER_FETCH) 01036 { 01037 if (worker->_competing) 01038 { 01039 worker->disableCompetition(); 01040 // multiple workers wrote into this block. We already know that our 01041 // data was correct, but maybe some other worker overwrote our data 01042 // with something broken. Thus we have to re-check the block. 01043 if (!worker->recheckChecksum()) 01044 { 01045 DBG << "#" << worker->_workerno << ": recheck checksum error, refetch block" << endl; 01046 // re-fetch! No need to worry about the bad workers, 01047 // they will now be set to DISCARD. At the end of their block 01048 // they will notice that they wrote bad data and go into BROKEN. 01049 worker->run(); 01050 continue; 01051 } 01052 } 01053 _fetchedgoodsize += worker->_blksize; 01054 } 01055 01056 // make bad workers sleep a little 01057 double maxavg = 0; 01058 int maxworkerno = 0; 01059 int numbetter = 0; 01060 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 01061 { 01062 multifetchworker *oworker = *workeriter; 01063 if (oworker->_state == WORKER_BROKEN) 01064 continue; 01065 if (oworker->_avgspeed > maxavg) 01066 { 01067 maxavg = oworker->_avgspeed; 01068 maxworkerno = oworker->_workerno; 01069 } 01070 if (oworker->_avgspeed > worker->_avgspeed) 01071 numbetter++; 01072 } 01073 if (maxavg && !_stealing) 01074 { 01075 double ratio = worker->_avgspeed / maxavg; 01076 ratio = 1 - ratio; 01077 if (numbetter < 3) // don't sleep that much if we're in the top two 01078 ratio = ratio * ratio; 01079 if (ratio > .01) 01080 { 01081 DBG << "#" << worker->_workerno << ": too slow ("<< ratio << ", " << worker->_avgspeed << ", #" << maxworkerno << ": " << maxavg << "), going to sleep for " << ratio * 1000 << " ms" << endl; 01082 worker->_sleepuntil = now + ratio; 01083 worker->_state = WORKER_SLEEP; 01084 _sleepworkers++; 01085 continue; 01086 } 01087 } 01088 01089 // do rate control (if requested) 01090 // should use periodavg, but that's not what libcurl does 01091 if (_maxspeed && now > _starttime) 01092 { 01093 double avg = _fetchedsize / (now - _starttime); 01094 avg = worker->_maxspeed * _maxspeed / avg; 01095 if (avg < _maxspeed / _maxworkers) 01096 avg = _maxspeed / _maxworkers; 01097 if (avg > _maxspeed) 01098 avg = _maxspeed; 01099 if (avg < 1024) 01100 avg = 1024; 01101 worker->_maxspeed = avg; 01102 #if LIBCURL_VERSION_NUMBER >= 0x071505 01103 curl_easy_setopt(worker->_curl, CURLOPT_MAX_RECV_SPEED_LARGE, (curl_off_t)(avg)); 01104 #endif 01105 } 01106 01107 worker->nextjob(); 01108 } 01109 else 01110 { 01111 worker->_state = WORKER_BROKEN; 01112 _activeworkers--; 01113 if (!_activeworkers && !(urliter != urllist.end() && _workers.size() < MAXURLS)) 01114 { 01115 // end of workers reached! goodbye! 01116 worker->evaluateCurlCode(Pathname(), cc, false); 01117 } 01118 } 01119 } 01120 01121 // send report 01122 if (_report) 01123 { 01124 int percent = _totalsize ? (100 * (_fetchedgoodsize + _fetchedsize)) / (_totalsize + _fetchedsize) : 0; 01125 double avg = 0; 01126 if (now > _starttime) 01127 avg = _fetchedsize / (now - _starttime); 01128 if (!(*(_report))->progress(percent, _baseurl, avg, _lastperiodstart == _starttime ? avg : _periodavg)) 01129 ZYPP_THROW(MediaCurlException(_baseurl, "User abort", "cancelled")); 01130 } 01131 01132 if (_timeout && now - _lastprogress > _timeout) 01133 break; 01134 } 01135 01136 if (!_finished) 01137 ZYPP_THROW(MediaTimeoutException(_baseurl)); 01138 01139 // print some download stats 01140 WAR << "overall result" << endl; 01141 for (std::list<multifetchworker *>::iterator workeriter = _workers.begin(); workeriter != _workers.end(); ++workeriter) 01142 { 01143 multifetchworker *worker = *workeriter; 01144 WAR << "#" << worker->_workerno << ": state: " << worker->_state << " received: " << worker->_received << " url: " << worker->_url << endl; 01145 } 01146 } 01147 01148 01150 01151 01152 MediaMultiCurl::MediaMultiCurl(const Url &url_r, const Pathname & attach_point_hint_r) 01153 : MediaCurl(url_r, attach_point_hint_r) 01154 { 01155 MIL << "MediaMultiCurl::MediaMultiCurl(" << url_r << ", " << attach_point_hint_r << ")" << endl; 01156 _multi = 0; 01157 _customHeadersMetalink = 0; 01158 } 01159 01160 MediaMultiCurl::~MediaMultiCurl() 01161 { 01162 if (_customHeadersMetalink) 01163 { 01164 curl_slist_free_all(_customHeadersMetalink); 01165 _customHeadersMetalink = 0; 01166 } 01167 if (_multi) 01168 { 01169 curl_multi_cleanup(_multi); 01170 _multi = 0; 01171 } 01172 std::map<std::string, CURL *>::iterator it; 01173 for (it = _easypool.begin(); it != _easypool.end(); it++) 01174 { 01175 CURL *easy = it->second; 01176 if (easy) 01177 { 01178 curl_easy_cleanup(easy); 01179 it->second = NULL; 01180 } 01181 } 01182 } 01183 01184 void MediaMultiCurl::setupEasy() 01185 { 01186 MediaCurl::setupEasy(); 01187 01188 if (_customHeadersMetalink) 01189 { 01190 curl_slist_free_all(_customHeadersMetalink); 01191 _customHeadersMetalink = 0; 01192 } 01193 struct curl_slist *sl = _customHeaders; 01194 for (; sl; sl = sl->next) 01195 _customHeadersMetalink = curl_slist_append(_customHeadersMetalink, sl->data); 01196 _customHeadersMetalink = curl_slist_append(_customHeadersMetalink, "Accept: */*, application/metalink+xml, application/metalink4+xml"); 01197 } 01198 01199 static bool looks_like_metalink(const Pathname & file) 01200 { 01201 char buf[256], *p; 01202 int fd, l; 01203 if ((fd = open(file.asString().c_str(), O_RDONLY)) == -1) 01204 return false; 01205 while ((l = read(fd, buf, sizeof(buf) - 1)) == -1 && errno == EINTR) 01206 ; 01207 close(fd); 01208 if (l == -1) 01209 return 0; 01210 buf[l] = 0; 01211 p = buf; 01212 while (*p == ' ' || *p == '\t' || *p == '\r' || *p == '\n') 01213 p++; 01214 if (!strncasecmp(p, "<?xml", 5)) 01215 { 01216 while (*p && *p != '>') 01217 p++; 01218 if (*p == '>') 01219 p++; 01220 while (*p == ' ' || *p == '\t' || *p == '\r' || *p == '\n') 01221 p++; 01222 } 01223 bool ret = !strncasecmp(p, "<metalink", 9) ? true : false; 01224 DBG << "looks_like_metalink(" << file << "): " << ret << endl; 01225 return ret; 01226 } 01227 01228 void MediaMultiCurl::doGetFileCopy( const Pathname & filename , const Pathname & target, callback::SendReport<DownloadProgressReport> & report, RequestOptions options ) const 01229 { 01230 Pathname dest = target.absolutename(); 01231 if( assert_dir( dest.dirname() ) ) 01232 { 01233 DBG << "assert_dir " << dest.dirname() << " failed" << endl; 01234 Url url(getFileUrl(filename)); 01235 ZYPP_THROW( MediaSystemException(url, "System error on " + dest.dirname().asString()) ); 01236 } 01237 string destNew = target.asString() + ".new.zypp.XXXXXX"; 01238 char *buf = ::strdup( destNew.c_str()); 01239 if( !buf) 01240 { 01241 ERR << "out of memory for temp file name" << endl; 01242 Url url(getFileUrl(filename)); 01243 ZYPP_THROW(MediaSystemException(url, "out of memory for temp file name")); 01244 } 01245 01246 int tmp_fd = ::mkstemp( buf ); 01247 if( tmp_fd == -1) 01248 { 01249 free( buf); 01250 ERR << "mkstemp failed for file '" << destNew << "'" << endl; 01251 ZYPP_THROW(MediaWriteException(destNew)); 01252 } 01253 destNew = buf; 01254 free( buf); 01255 01256 FILE *file = ::fdopen( tmp_fd, "w" ); 01257 if ( !file ) { 01258 ::close( tmp_fd); 01259 filesystem::unlink( destNew ); 01260 ERR << "fopen failed for file '" << destNew << "'" << endl; 01261 ZYPP_THROW(MediaWriteException(destNew)); 01262 } 01263 DBG << "dest: " << dest << endl; 01264 DBG << "temp: " << destNew << endl; 01265 01266 // set IFMODSINCE time condition (no download if not modified) 01267 if( PathInfo(target).isExist() && !(options & OPTION_NO_IFMODSINCE) ) 01268 { 01269 curl_easy_setopt(_curl, CURLOPT_TIMECONDITION, CURL_TIMECOND_IFMODSINCE); 01270 curl_easy_setopt(_curl, CURLOPT_TIMEVALUE, (long)PathInfo(target).mtime()); 01271 } 01272 else 01273 { 01274 curl_easy_setopt(_curl, CURLOPT_TIMECONDITION, CURL_TIMECOND_NONE); 01275 curl_easy_setopt(_curl, CURLOPT_TIMEVALUE, 0L); 01276 } 01277 // change header to include Accept: metalink 01278 curl_easy_setopt(_curl, CURLOPT_HTTPHEADER, _customHeadersMetalink); 01279 try 01280 { 01281 MediaCurl::doGetFileCopyFile(filename, dest, file, report, options); 01282 } 01283 catch (Exception &ex) 01284 { 01285 ::fclose(file); 01286 filesystem::unlink(destNew); 01287 curl_easy_setopt(_curl, CURLOPT_TIMECONDITION, CURL_TIMECOND_NONE); 01288 curl_easy_setopt(_curl, CURLOPT_TIMEVALUE, 0L); 01289 curl_easy_setopt(_curl, CURLOPT_HTTPHEADER, _customHeaders); 01290 ZYPP_RETHROW(ex); 01291 } 01292 curl_easy_setopt(_curl, CURLOPT_TIMECONDITION, CURL_TIMECOND_NONE); 01293 curl_easy_setopt(_curl, CURLOPT_TIMEVALUE, 0L); 01294 curl_easy_setopt(_curl, CURLOPT_HTTPHEADER, _customHeaders); 01295 long httpReturnCode = 0; 01296 CURLcode infoRet = curl_easy_getinfo(_curl, CURLINFO_RESPONSE_CODE, &httpReturnCode); 01297 if (infoRet == CURLE_OK) 01298 { 01299 DBG << "HTTP response: " + str::numstring(httpReturnCode) << endl; 01300 if ( httpReturnCode == 304 01301 || ( httpReturnCode == 213 && _url.getScheme() == "ftp" ) ) // not modified 01302 { 01303 DBG << "not modified: " << PathInfo(dest) << endl; 01304 return; 01305 } 01306 } 01307 else 01308 { 01309 WAR << "Could not get the reponse code." << endl; 01310 } 01311 01312 bool ismetalink = false; 01313 01314 char *ptr = NULL; 01315 if (curl_easy_getinfo(_curl, CURLINFO_CONTENT_TYPE, &ptr) == CURLE_OK && ptr) 01316 { 01317 string ct = string(ptr); 01318 if (ct.find("application/metalink+xml") == 0 || ct.find("application/metalink4+xml") == 0) 01319 ismetalink = true; 01320 } 01321 01322 if (!ismetalink) 01323 { 01324 // some proxies do not store the content type, so also look at the file to find 01325 // out if we received a metalink (bnc#649925) 01326 fflush(file); 01327 if (looks_like_metalink(Pathname(destNew))) 01328 ismetalink = true; 01329 } 01330 01331 if (ismetalink) 01332 { 01333 bool userabort = false; 01334 fclose(file); 01335 file = NULL; 01336 Pathname failedFile = ZConfig::instance().repoCachePath() / "MultiCurl.failed"; 01337 try 01338 { 01339 MetaLinkParser mlp; 01340 mlp.parse(Pathname(destNew)); 01341 MediaBlockList bl = mlp.getBlockList(); 01342 vector<Url> urls = mlp.getUrls(); 01343 DBG << bl << endl; 01344 file = fopen(destNew.c_str(), "w+"); 01345 if (!file) 01346 ZYPP_THROW(MediaWriteException(destNew)); 01347 if (PathInfo(target).isExist()) 01348 { 01349 DBG << "reusing blocks from file " << target << endl; 01350 bl.reuseBlocks(file, target.asString()); 01351 DBG << bl << endl; 01352 } 01353 if (bl.haveChecksum(1) && PathInfo(failedFile).isExist()) 01354 { 01355 DBG << "reusing blocks from file " << failedFile << endl; 01356 bl.reuseBlocks(file, failedFile.asString()); 01357 DBG << bl << endl; 01358 filesystem::unlink(failedFile); 01359 } 01360 Pathname df = deltafile(); 01361 if (!df.empty()) 01362 { 01363 DBG << "reusing blocks from file " << df << endl; 01364 bl.reuseBlocks(file, df.asString()); 01365 DBG << bl << endl; 01366 } 01367 try 01368 { 01369 multifetch(filename, file, &urls, &report, &bl); 01370 } 01371 catch (MediaCurlException &ex) 01372 { 01373 userabort = ex.errstr() == "User abort"; 01374 ZYPP_RETHROW(ex); 01375 } 01376 } 01377 catch (Exception &ex) 01378 { 01379 // something went wrong. fall back to normal download 01380 if (file) 01381 fclose(file); 01382 file = NULL; 01383 if (PathInfo(destNew).size() >= 63336) 01384 { 01385 ::unlink(failedFile.asString().c_str()); 01386 filesystem::hardlinkCopy(destNew, failedFile); 01387 } 01388 if (userabort) 01389 { 01390 filesystem::unlink(destNew); 01391 ZYPP_RETHROW(ex); 01392 } 01393 file = fopen(destNew.c_str(), "w+"); 01394 if (!file) 01395 ZYPP_THROW(MediaWriteException(destNew)); 01396 MediaCurl::doGetFileCopyFile(filename, dest, file, report, options | OPTION_NO_REPORT_START); 01397 } 01398 } 01399 01400 if (::fchmod( ::fileno(file), filesystem::applyUmaskTo( 0644 ))) 01401 { 01402 ERR << "Failed to chmod file " << destNew << endl; 01403 } 01404 if (::fclose(file)) 01405 { 01406 filesystem::unlink(destNew); 01407 ERR << "Fclose failed for file '" << destNew << "'" << endl; 01408 ZYPP_THROW(MediaWriteException(destNew)); 01409 } 01410 if ( rename( destNew, dest ) != 0 ) 01411 { 01412 ERR << "Rename failed" << endl; 01413 ZYPP_THROW(MediaWriteException(dest)); 01414 } 01415 DBG << "done: " << PathInfo(dest) << endl; 01416 } 01417 01418 void MediaMultiCurl::multifetch(const Pathname & filename, FILE *fp, std::vector<Url> *urllist, callback::SendReport<DownloadProgressReport> *report, MediaBlockList *blklist, off_t filesize) const 01419 { 01420 Url baseurl(getFileUrl(filename)); 01421 if (blklist && filesize == off_t(-1) && blklist->haveFilesize()) 01422 filesize = blklist->getFilesize(); 01423 if (blklist && !blklist->haveBlocks() && filesize != 0) 01424 blklist = 0; 01425 if (blklist && (filesize == 0 || !blklist->numBlocks())) 01426 { 01427 checkFileDigest(baseurl, fp, blklist); 01428 return; 01429 } 01430 if (filesize == 0) 01431 return; 01432 if (!_multi) 01433 { 01434 _multi = curl_multi_init(); 01435 if (!_multi) 01436 ZYPP_THROW(MediaCurlInitException(baseurl)); 01437 } 01438 multifetchrequest req(this, filename, baseurl, _multi, fp, report, blklist, filesize); 01439 req._timeout = _settings.timeout(); 01440 req._connect_timeout = _settings.connectTimeout(); 01441 req._maxspeed = _settings.maxDownloadSpeed(); 01442 req._maxworkers = _settings.maxConcurrentConnections(); 01443 if (req._maxworkers > MAXURLS) 01444 req._maxworkers = MAXURLS; 01445 if (req._maxworkers <= 0) 01446 req._maxworkers = 1; 01447 std::vector<Url> myurllist; 01448 for (std::vector<Url>::iterator urliter = urllist->begin(); urliter != urllist->end(); ++urliter) 01449 { 01450 try 01451 { 01452 string scheme = urliter->getScheme(); 01453 if (scheme == "http" || scheme == "https" || scheme == "ftp") 01454 { 01455 checkProtocol(*urliter); 01456 myurllist.push_back(*urliter); 01457 } 01458 } 01459 catch (...) 01460 { 01461 } 01462 } 01463 if (!myurllist.size()) 01464 myurllist.push_back(baseurl); 01465 req.run(myurllist); 01466 checkFileDigest(baseurl, fp, blklist); 01467 } 01468 01469 void MediaMultiCurl::checkFileDigest(Url &url, FILE *fp, MediaBlockList *blklist) const 01470 { 01471 if (!blklist || !blklist->haveFileChecksum()) 01472 return; 01473 if (fseeko(fp, off_t(0), SEEK_SET)) 01474 ZYPP_THROW(MediaCurlException(url, "fseeko", "seek error")); 01475 Digest dig; 01476 blklist->createFileDigest(dig); 01477 char buf[4096]; 01478 size_t l; 01479 while ((l = fread(buf, 1, sizeof(buf), fp)) > 0) 01480 dig.update(buf, l); 01481 if (!blklist->verifyFileDigest(dig)) 01482 ZYPP_THROW(MediaCurlException(url, "file verification failed", "checksum error")); 01483 } 01484 01485 bool MediaMultiCurl::isDNSok(const string &host) const 01486 { 01487 return _dnsok.find(host) == _dnsok.end() ? false : true; 01488 } 01489 01490 void MediaMultiCurl::setDNSok(const string &host) const 01491 { 01492 _dnsok.insert(host); 01493 } 01494 01495 CURL *MediaMultiCurl::fromEasyPool(const string &host) const 01496 { 01497 if (_easypool.find(host) == _easypool.end()) 01498 return 0; 01499 CURL *ret = _easypool[host]; 01500 _easypool.erase(host); 01501 return ret; 01502 } 01503 01504 void MediaMultiCurl::toEasyPool(const std::string &host, CURL *easy) const 01505 { 01506 CURL *oldeasy = _easypool[host]; 01507 _easypool[host] = easy; 01508 if (oldeasy) 01509 curl_easy_cleanup(oldeasy); 01510 } 01511 01512 } // namespace media 01513 } // namespace zypp 01514
1.7.6.1