|
2 | 2 |
|
3 | 3 | namespace Stackkit\LaravelGoogleCloudTasksQueue; |
4 | 4 |
|
| 5 | +use Carbon\Carbon; |
| 6 | +use Google\Cloud\Tasks\V2\CloudTasksClient; |
| 7 | +use Google\Cloud\Tasks\V2\HttpMethod; |
| 8 | +use Google\Cloud\Tasks\V2\HttpRequest; |
| 9 | +use Google\Cloud\Tasks\V2\Task; |
| 10 | +use Google\Protobuf\Timestamp; |
5 | 11 | use Illuminate\Contracts\Queue\Queue as QueueContract; |
6 | 12 | use Illuminate\Queue\Queue as LaravelQueue; |
| 13 | +use Illuminate\Support\InteractsWithTime; |
7 | 14 |
|
8 | 15 | class CloudTasksQueue extends LaravelQueue implements QueueContract |
9 | 16 | { |
| 17 | + use InteractsWithTime; |
| 18 | + |
| 19 | + private $client; |
| 20 | + private $default; |
| 21 | + |
| 22 | + public function __construct(array $config, CloudTasksClient $client) |
| 23 | + { |
| 24 | + $this->client = $client; |
| 25 | + $this->default = $config['queue']; |
| 26 | + } |
| 27 | + |
10 | 28 | public function size($queue = null) |
11 | 29 | { |
12 | 30 | // TODO: Implement size() method. |
13 | 31 | } |
14 | 32 |
|
15 | 33 | public function push($job, $data = '', $queue = null) |
16 | 34 | { |
17 | | - // TODO: Implement push() method. |
| 35 | + return $this->pushToCloudTasks($queue, $this->createPayload( |
| 36 | + $job, $this->getQueue($queue), $data |
| 37 | + )); |
18 | 38 | } |
19 | 39 |
|
20 | 40 | public function pushRaw($payload, $queue = null, array $options = []) |
21 | 41 | { |
22 | | - // TODO: Implement pushRaw() method. |
| 42 | + return $this->pushToCloudTasks($queue, $payload); |
23 | 43 | } |
24 | 44 |
|
25 | 45 | public function later($delay, $job, $data = '', $queue = null) |
26 | 46 | { |
27 | | - // TODO: Implement later() method. |
| 47 | + return $this->pushToCloudTasks($queue, $this->createPayload( |
| 48 | + $job, $this->getQueue($queue), $data |
| 49 | + ), $delay); |
| 50 | + } |
| 51 | + |
| 52 | + protected function pushToCloudTasks($queue, $payload, $delay = 0, $attempts = 0) |
| 53 | + { |
| 54 | + $queue = $this->getQueue($queue); |
| 55 | + $queueName = $this->client->queueName(Config::project(), Config::location(), $queue); |
| 56 | + $availableAt = $this->availableAt($delay); |
| 57 | + |
| 58 | + $httpRequest = new HttpRequest(); |
| 59 | + $httpRequest->setUrl(Config::handler()); |
| 60 | + $httpRequest->setHttpMethod(HttpMethod::POST); |
| 61 | + $httpRequest->setBody($payload); |
| 62 | + |
| 63 | + $task = new Task; |
| 64 | + $task->setHttpRequest($httpRequest); |
| 65 | + |
| 66 | + if ($availableAt > time()) { |
| 67 | + $task->setScheduleTime(new Timestamp(['seconds' => $availableAt])); |
| 68 | + } |
| 69 | + |
| 70 | + $this->client->createTask($queueName, $task); |
28 | 71 | } |
29 | 72 |
|
30 | 73 | public function pop($queue = null) |
31 | 74 | { |
32 | 75 | // TODO: Implement pop() method. |
33 | 76 | } |
| 77 | + |
| 78 | + private function getQueue($queue = null) |
| 79 | + { |
| 80 | + return $queue ?: $this->default; |
| 81 | + } |
34 | 82 | } |
0 commit comments