AmpResponse.php 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514
  1. <?php
  2. /*
  3. * This file is part of the Symfony package.
  4. *
  5. * (c) Fabien Potencier <fabien@symfony.com>
  6. *
  7. * For the full copyright and license information, please view the LICENSE
  8. * file that was distributed with this source code.
  9. */
  10. namespace Symfony\Component\HttpClient\Response;
  11. use Amp\ByteStream\StreamException;
  12. use Amp\CancellationTokenSource;
  13. use Amp\Coroutine;
  14. use Amp\Deferred;
  15. use Amp\Http\Client\HttpException;
  16. use Amp\Http\Client\Request;
  17. use Amp\Http\Client\Response;
  18. use Amp\Loop;
  19. use Amp\Promise;
  20. use Amp\Success;
  21. use Psr\Log\LoggerInterface;
  22. use Symfony\Component\HttpClient\Chunk\FirstChunk;
  23. use Symfony\Component\HttpClient\Chunk\InformationalChunk;
  24. use Symfony\Component\HttpClient\Exception\InvalidArgumentException;
  25. use Symfony\Component\HttpClient\Exception\TransportException;
  26. use Symfony\Component\HttpClient\HttpClientTrait;
  27. use Symfony\Component\HttpClient\Internal\AmpBody;
  28. use Symfony\Component\HttpClient\Internal\AmpClientState;
  29. use Symfony\Component\HttpClient\Internal\Canary;
  30. use Symfony\Component\HttpClient\Internal\ClientState;
  31. use Symfony\Contracts\HttpClient\ResponseInterface;
  32. /**
  33. * @author Nicolas Grekas <p@tchwork.com>
  34. *
  35. * @internal
  36. */
  37. final class AmpResponse implements ResponseInterface, StreamableInterface
  38. {
  39. use CommonResponseTrait;
  40. use TransportResponseTrait;
  41. /**
  42. * @var string
  43. */
  44. private static $nextId = 'a';
  45. /**
  46. * @var \Symfony\Component\HttpClient\Internal\AmpClientState
  47. */
  48. private $multi;
  49. /**
  50. * @var mixed[]|null
  51. */
  52. private $options;
  53. /**
  54. * @var \Closure
  55. */
  56. private $onProgress;
  57. /**
  58. * @var string|null
  59. */
  60. private static $delay;
  61. /**
  62. * @internal
  63. * @param \Symfony\Component\HttpClient\Internal\AmpClientState $multi
  64. * @param \Amp\Http\Client\Request $request
  65. * @param mixed[] $options
  66. * @param \Psr\Log\LoggerInterface|null $logger
  67. */
  68. public function __construct($multi, $request, $options, $logger)
  69. {
  70. $this->multi = $multi;
  71. $this->options = &$options;
  72. $this->logger = $logger;
  73. $this->timeout = $options['timeout'];
  74. $this->shouldBuffer = $options['buffer'];
  75. if ($this->inflate = \extension_loaded('zlib') && !$request->hasHeader('accept-encoding')) {
  76. $request->setHeader('Accept-Encoding', 'gzip');
  77. }
  78. $this->initializer = static function (self $response) {
  79. return null !== $response->options;
  80. };
  81. $info = &$this->info;
  82. $headers = &$this->headers;
  83. $canceller = new CancellationTokenSource();
  84. $handle = &$this->handle;
  85. $info['url'] = (string) $request->getUri();
  86. $info['http_method'] = $request->getMethod();
  87. $info['start_time'] = null;
  88. $info['redirect_url'] = null;
  89. $info['original_url'] = $info['url'];
  90. $info['redirect_time'] = 0.0;
  91. $info['redirect_count'] = 0;
  92. $info['size_upload'] = 0.0;
  93. $info['size_download'] = 0.0;
  94. $info['upload_content_length'] = -1.0;
  95. $info['download_content_length'] = -1.0;
  96. $info['user_data'] = $options['user_data'];
  97. $info['max_duration'] = $options['max_duration'];
  98. $info['debug'] = '';
  99. $onProgress = $options['on_progress'] ?? static function () {};
  100. $onProgress = $this->onProgress = static function () use (&$info, $onProgress) {
  101. $info['total_time'] = microtime(true) - $info['start_time'];
  102. $onProgress((int) $info['size_download'], ((int) (1 + $info['download_content_length']) ?: 1) - 1, (array) $info);
  103. };
  104. $pauseDeferred = new Deferred();
  105. $pause = new Success();
  106. $throttleWatcher = null;
  107. $this->id = $id = self::$nextId++;
  108. Loop::defer(static function () use ($request, $multi, $id, &$info, &$headers, $canceller, &$options, $onProgress, &$handle, $logger, &$pause) {
  109. return new Coroutine(self::generateResponse($request, $multi, $id, $info, $headers, $canceller, $options, $onProgress, $handle, $logger, $pause));
  110. });
  111. $info['pause_handler'] = static function (float $duration) use (&$throttleWatcher, &$pauseDeferred, &$pause) {
  112. if (null !== $throttleWatcher) {
  113. Loop::cancel($throttleWatcher);
  114. }
  115. $pause = $pauseDeferred->promise();
  116. if ($duration <= 0) {
  117. $deferred = $pauseDeferred;
  118. $pauseDeferred = new Deferred();
  119. $deferred->resolve();
  120. } else {
  121. $throttleWatcher = Loop::delay(ceil(1000 * $duration), static function () use (&$pauseDeferred) {
  122. $deferred = $pauseDeferred;
  123. $pauseDeferred = new Deferred();
  124. $deferred->resolve();
  125. });
  126. }
  127. };
  128. $multi->lastTimeout = null;
  129. $multi->openHandles[$id] = $id;
  130. ++$multi->responseCount;
  131. $this->canary = new Canary(static function () use ($canceller, $multi, $id) {
  132. $canceller->cancel();
  133. unset($multi->openHandles[$id], $multi->handlesActivity[$id]);
  134. });
  135. }
  136. /**
  137. * @return mixed
  138. * @param string|null $type
  139. */
  140. public function getInfo($type = null)
  141. {
  142. return null !== $type ? $this->info[$type] ?? null : $this->info;
  143. }
  144. public function __sleep()
  145. {
  146. throw new \BadMethodCallException('Cannot serialize '.__CLASS__);
  147. }
  148. public function __wakeup()
  149. {
  150. throw new \BadMethodCallException('Cannot unserialize '.__CLASS__);
  151. }
  152. public function __destruct()
  153. {
  154. try {
  155. $this->doDestruct();
  156. } finally {
  157. // Clear the DNS cache when all requests completed
  158. if (0 >= --$this->multi->responseCount) {
  159. $this->multi->responseCount = 0;
  160. $this->multi->dnsCache = [];
  161. }
  162. }
  163. }
  164. /**
  165. * @param $this $response
  166. * @param mixed[] $runningResponses
  167. */
  168. private static function schedule($response, &$runningResponses)
  169. {
  170. if (isset($runningResponses[0])) {
  171. $runningResponses[0][1][$response->id] = $response;
  172. } else {
  173. $runningResponses[0] = [$response->multi, [$response->id => $response]];
  174. }
  175. if (!isset($response->multi->openHandles[$response->id])) {
  176. $response->multi->handlesActivity[$response->id][] = null;
  177. $response->multi->handlesActivity[$response->id][] = null !== $response->info['error'] ? new TransportException($response->info['error']) : null;
  178. }
  179. }
  180. /**
  181. * @param \Symfony\Component\HttpClient\Internal\ClientState $multi
  182. * @param mixed[]|null $responses
  183. */
  184. private static function perform($multi, &$responses = null)
  185. {
  186. if ($responses) {
  187. foreach ($responses as $response) {
  188. try {
  189. if ($response->info['start_time']) {
  190. $response->info['total_time'] = microtime(true) - $response->info['start_time'];
  191. ($response->onProgress)();
  192. }
  193. } catch (\Throwable $e) {
  194. $multi->handlesActivity[$response->id][] = null;
  195. $multi->handlesActivity[$response->id][] = $e;
  196. }
  197. }
  198. }
  199. }
  200. /**
  201. * @param \Symfony\Component\HttpClient\Internal\ClientState $multi
  202. * @param float $timeout
  203. */
  204. private static function select($multi, $timeout)
  205. {
  206. $timeout += microtime(true);
  207. self::$delay = Loop::defer(static function () use ($timeout) {
  208. if (0 < $timeout -= microtime(true)) {
  209. self::$delay = Loop::delay(ceil(1000 * $timeout), \Closure::fromCallable([Loop::class, 'stop']));
  210. } else {
  211. Loop::stop();
  212. }
  213. });
  214. Loop::run();
  215. return null === self::$delay ? 1 : 0;
  216. }
  217. /**
  218. * @param \Amp\Http\Client\Request $request
  219. * @param \Symfony\Component\HttpClient\Internal\AmpClientState $multi
  220. * @param string $id
  221. * @param mixed[] $info
  222. * @param mixed[] $headers
  223. * @param \Amp\CancellationTokenSource $canceller
  224. * @param mixed[] $options
  225. * @param \Closure $onProgress
  226. * @param \Psr\Log\LoggerInterface|null $logger
  227. * @param \Amp\Promise $pause
  228. */
  229. private static function generateResponse($request, $multi, $id, &$info, &$headers, $canceller, &$options, $onProgress, &$handle, $logger, &$pause)
  230. {
  231. $request->setInformationalResponseHandler(static function (Response $response) use ($multi, $id, &$info, &$headers) {
  232. self::addResponseHeaders($response, $info, $headers);
  233. $multi->handlesActivity[$id][] = new InformationalChunk($response->getStatus(), $response->getHeaders());
  234. self::stopLoop();
  235. });
  236. try {
  237. /* @var Response $response */
  238. if (null === $response = yield from self::getPushedResponse($request, $multi, $info, $headers, $options, $logger)) {
  239. ($logger2 = $logger) ? $logger2->info(sprintf('Request: "%s %s"', $info['http_method'], $info['url'])) : null;
  240. $response = yield from self::followRedirects($request, $multi, $info, $headers, $canceller, $options, $onProgress, $handle, $logger, $pause);
  241. }
  242. $options = null;
  243. $multi->handlesActivity[$id][] = new FirstChunk();
  244. if ('HEAD' === $response->getRequest()->getMethod() || \in_array($info['http_code'], [204, 304], true)) {
  245. $multi->handlesActivity[$id][] = null;
  246. $multi->handlesActivity[$id][] = null;
  247. self::stopLoop();
  248. return;
  249. }
  250. if ($response->hasHeader('content-length')) {
  251. $info['download_content_length'] = (float) $response->getHeader('content-length');
  252. }
  253. $body = $response->getBody();
  254. while (true) {
  255. self::stopLoop();
  256. yield $pause;
  257. if (null === $data = yield $body->read()) {
  258. break;
  259. }
  260. $info['size_download'] += \strlen($data);
  261. $multi->handlesActivity[$id][] = $data;
  262. }
  263. $multi->handlesActivity[$id][] = null;
  264. $multi->handlesActivity[$id][] = null;
  265. } catch (\Throwable $e) {
  266. $multi->handlesActivity[$id][] = null;
  267. $multi->handlesActivity[$id][] = $e;
  268. } finally {
  269. $info['download_content_length'] = $info['size_download'];
  270. }
  271. self::stopLoop();
  272. }
  273. /**
  274. * @param \Amp\Http\Client\Request $originRequest
  275. * @param \Symfony\Component\HttpClient\Internal\AmpClientState $multi
  276. * @param mixed[] $info
  277. * @param mixed[] $headers
  278. * @param \Amp\CancellationTokenSource $canceller
  279. * @param mixed[] $options
  280. * @param \Closure $onProgress
  281. * @param \Psr\Log\LoggerInterface|null $logger
  282. * @param \Amp\Promise $pause
  283. */
  284. private static function followRedirects($originRequest, $multi, &$info, &$headers, $canceller, $options, $onProgress, &$handle, $logger, &$pause)
  285. {
  286. yield $pause;
  287. $originRequest->setBody(new AmpBody($options['body'], $info, $onProgress));
  288. $response = yield $multi->request($options, $originRequest, $canceller->getToken(), $info, $onProgress, $handle);
  289. $previousUrl = null;
  290. while (true) {
  291. self::addResponseHeaders($response, $info, $headers);
  292. $status = $response->getStatus();
  293. if (!\in_array($status, [301, 302, 303, 307, 308], true) || null === $location = $response->getHeader('location')) {
  294. return $response;
  295. }
  296. $urlResolver = new class() {
  297. use HttpClientTrait {
  298. parseUrl as public;
  299. resolveUrl as public;
  300. }
  301. };
  302. try {
  303. $previousUrl = $previousUrl ?? $urlResolver::parseUrl($info['url']);
  304. $location = $urlResolver::parseUrl($location);
  305. $location = $urlResolver::resolveUrl($location, $previousUrl);
  306. $info['redirect_url'] = implode('', $location);
  307. } catch (InvalidArgumentException $exception) {
  308. return $response;
  309. }
  310. if (0 >= $options['max_redirects'] || $info['redirect_count'] >= $options['max_redirects']) {
  311. return $response;
  312. }
  313. ($logger2 = $logger) ? $logger2->info(sprintf('Redirecting: "%s %s"', $status, $info['url'])) : null;
  314. try {
  315. // Discard body of redirects
  316. while (null !== yield $response->getBody()->read()) {
  317. }
  318. } catch (HttpException|StreamException $exception) {
  319. // Ignore streaming errors on previous responses
  320. }
  321. ++$info['redirect_count'];
  322. $info['url'] = $info['redirect_url'];
  323. $info['redirect_url'] = null;
  324. $previousUrl = $location;
  325. $request = new Request($info['url'], $info['http_method']);
  326. $request->setProtocolVersions($originRequest->getProtocolVersions());
  327. $request->setTcpConnectTimeout($originRequest->getTcpConnectTimeout());
  328. $request->setTlsHandshakeTimeout($originRequest->getTlsHandshakeTimeout());
  329. $request->setTransferTimeout($originRequest->getTransferTimeout());
  330. if (\in_array($status, [301, 302, 303], true)) {
  331. $originRequest->removeHeader('transfer-encoding');
  332. $originRequest->removeHeader('content-length');
  333. $originRequest->removeHeader('content-type');
  334. // Do like curl and browsers: turn POST to GET on 301, 302 and 303
  335. if ('POST' === $response->getRequest()->getMethod() || 303 === $status) {
  336. $info['http_method'] = 'HEAD' === $response->getRequest()->getMethod() ? 'HEAD' : 'GET';
  337. $request->setMethod($info['http_method']);
  338. }
  339. } else {
  340. $request->setBody(AmpBody::rewind($response->getRequest()->getBody()));
  341. }
  342. foreach ($originRequest->getRawHeaders() as [$name, $value]) {
  343. $request->addHeader($name, $value);
  344. }
  345. if ($request->getUri()->getAuthority() !== $originRequest->getUri()->getAuthority()) {
  346. $request->removeHeader('authorization');
  347. $request->removeHeader('cookie');
  348. $request->removeHeader('host');
  349. }
  350. yield $pause;
  351. $response = yield $multi->request($options, $request, $canceller->getToken(), $info, $onProgress, $handle);
  352. $info['redirect_time'] = microtime(true) - $info['start_time'];
  353. }
  354. }
  355. /**
  356. * @param \Amp\Http\Client\Response $response
  357. * @param mixed[] $info
  358. * @param mixed[] $headers
  359. */
  360. private static function addResponseHeaders($response, &$info, &$headers)
  361. {
  362. $info['http_code'] = $response->getStatus();
  363. if ($headers) {
  364. $info['debug'] .= "< \r\n";
  365. $headers = [];
  366. }
  367. $h = sprintf('HTTP/%s %s %s', $response->getProtocolVersion(), $response->getStatus(), $response->getReason());
  368. $info['debug'] .= "< {$h}\r\n";
  369. $info['response_headers'][] = $h;
  370. foreach ($response->getRawHeaders() as [$name, $value]) {
  371. $headers[strtolower($name)][] = $value;
  372. $h = $name.': '.$value;
  373. $info['debug'] .= "< {$h}\r\n";
  374. $info['response_headers'][] = $h;
  375. }
  376. $info['debug'] .= "< \r\n";
  377. }
  378. /**
  379. * Accepts pushed responses only if their headers related to authentication match the request.
  380. * @param \Amp\Http\Client\Request $request
  381. * @param \Symfony\Component\HttpClient\Internal\AmpClientState $multi
  382. * @param mixed[] $info
  383. * @param mixed[] $headers
  384. * @param mixed[] $options
  385. * @param \Psr\Log\LoggerInterface|null $logger
  386. */
  387. private static function getPushedResponse($request, $multi, &$info, &$headers, $options, $logger)
  388. {
  389. if ('' !== $options['body']) {
  390. return null;
  391. }
  392. $authority = $request->getUri()->getAuthority();
  393. foreach ($multi->pushedResponses[$authority] ?? [] as $i => [$pushedUrl, $pushDeferred, $pushedRequest, $pushedResponse, $parentOptions]) {
  394. if ($info['url'] !== $pushedUrl || $info['http_method'] !== $pushedRequest->getMethod()) {
  395. continue;
  396. }
  397. foreach ($parentOptions as $k => $v) {
  398. if ($options[$k] !== $v) {
  399. continue 2;
  400. }
  401. }
  402. foreach (['authorization', 'cookie', 'range', 'proxy-authorization'] as $k) {
  403. if ($pushedRequest->getHeaderArray($k) !== $request->getHeaderArray($k)) {
  404. continue 2;
  405. }
  406. }
  407. $response = yield $pushedResponse;
  408. foreach ($response->getHeaderArray('vary') as $vary) {
  409. foreach (preg_split('/\s*+,\s*+/', $vary) as $v) {
  410. if ('*' === $v || ($pushedRequest->getHeaderArray($v) !== $request->getHeaderArray($v) && 'accept-encoding' !== strtolower($v))) {
  411. ($logger2 = $logger) ? $logger2->debug(sprintf('Skipping pushed response: "%s"', $info['url'])) : null;
  412. continue 3;
  413. }
  414. }
  415. }
  416. $pushDeferred->resolve();
  417. ($logger2 = $logger) ? $logger2->debug(sprintf('Accepting pushed response: "%s %s"', $info['http_method'], $info['url'])) : null;
  418. self::addResponseHeaders($response, $info, $headers);
  419. unset($multi->pushedResponses[$authority][$i]);
  420. if (!$multi->pushedResponses[$authority]) {
  421. unset($multi->pushedResponses[$authority]);
  422. }
  423. return $response;
  424. }
  425. }
  426. private static function stopLoop()
  427. {
  428. if (null !== self::$delay) {
  429. Loop::cancel(self::$delay);
  430. self::$delay = null;
  431. }
  432. Loop::defer(\Closure::fromCallable([Loop::class, 'stop']));
  433. }
  434. }