<?php

class WorkerTask extends TaskBase
{
    public function crawlnewAction($opts = [])
    {
        $crawler = new \ProductDetail();

        if (isset($opts['proxy'])) {
            $crawler->setProxy($opts['proxy']);
        }

        while (true) {

            $sql = 'SELECT asin, store_id FROM product_tmp WHERE status = 0 LIMIT 5000';
            $products = $this->db->fetchAll($sql, null);
            $count = count($products);

            if ($count > 0) {
                $crawler->get($products, function($old, $new) {
                    $this->redis->lPush('insert_product_new', json_encode([
                        'old' => $old,
                        'new' => $new,
                    ]));
                });
            } else {
                echo "\nNo pending products";
                sleep(30);
            }
        }
    }

    public function crawloldAction($opts = [])
    {
        $crawler = new \ProductDetail();

        if (isset($opts['proxy'])) {
            $crawler->setProxy($opts['proxy']);
        }

        if (isset($opts['server'])) {
            $crawler->setServer($opts['server']);
        }

        while (true) {

            $queueLen = (int)$this->redis->lLen('update_product_queue');

            echo "\nupdate_product_queue: $queueLen";

            $arr = $this->redis->brPop('update_product_queue', 0);

            if (isset($arr[1])) {

                $products = json_decode($arr[1]);

                if (empty($products)) {
                    echo "\nProducts empty";
                } else {
                    $crawler->get($products, function($old, $new) {
                        $this->redis->lPush('update_product', json_encode([
                            'old' => $old,
                            'new' => $new,
                        ]));
                    });
                }
            }
        }

        $crawler->complete();
    }

    public function updatequeueAction()
    {
        // get the last elemenet of queue
        $lastUpdate = null;
        $lastElem = $this->redis->lIndex('update_product_queue', -1);
        if ($lastElem) {
            $list = json_decode($lastElem);
            if (is_array($list)) {
                $lastUpdate = $list[count($list) - 1]->updated_at;
            }
        }

        while (true) {

            $queueLen = (int)$this->redis->lLen('update_product_queue');

            echo "\nupdate_product_queue: $queueLen";

            if ($queueLen < 10) {

                $products = getOldestProducts($queueLen > 0 ? $lastUpdate : null);
                $count = count($products);

                echo "\nProducts: $count";
                echo "\nLast Update: $lastUpdate";

                if ($count > 0) {

                    $lastUpdate = $products[$count - 1]->updated_at;
                    $parts = array_chunk($products, 100);

                    foreach ($parts as $part) {
                        $this->redis->lPush('update_product_queue', json_encode($part));
                    }
                }
            }

            echo "\n";
            sleep(10);
        }
    }

    public function updateproductAction()
    {
        while (true) {

            $arr = $this->redis->brPop('update_product', 0);

            if (isset($arr[1])) {

                $data = json_decode($arr[1]);

                if (!isset($data->new) || !isset($data->old)) continue;

                $old = $data->old;
                $new = $data->new;

                if ($new) {

                    printf("\n%s | %s | %s", $old->asin, $new->asin, $new->name);

                    if ($old->asin == $new->asin) {

                        updateProduct($old, $new);

                    } else {

                        // insert new (parent)
                        insertUpdateProduct($old, $new);

                        // remove old (child)
                        updatePartialProduct($old->asin, [
                            'status' => \ProductStatus::DUPLICATED,
                            'parent' => $new->asin,
                        ]);
                    }

                    // insert messages
                    if ($this->hasRankUp($old->rank, $new->rank)) {
                        $message = sprintf(
                            '+%s rank: #%s to #%s',
                            number_format($old->rank > 0 ? $old->rank - $new->rank : $new->rank),
                            number_format($old->rank),
                            number_format($new->rank)
                        );
                        insertMessage($new->asin, $new->mid, 0, $message, [
                            'old' => $old,
                            'new' => $new,
                        ]);
                    }

                    if ($this->hasNewReviews($old->reviews, $new->reviews)) {
                        $message = sprintf(
                            '%d reviews',
                            $new->reviews - $old->reviews,
                            number_format($old->reviews),
                            number_format($new->reviews)
                        );
                        insertMessage($new->asin, $new->mid, 1, $message, [
                            'old' => $old,
                            'new' => $new,
                        ]);
                    }

                } else {

                    printf("\n%s | dog page", $old->asin);

                    updatePartialProduct($old->asin, [
                        'status'     => \ProductStatus::REMOVED,
                        'deleted_at' => getMicrotime(),
                    ]);
                }
            }
        }
    }

    public function insertproducttmpAction()
    {
        while (true) {

            $arr = $this->redis->brPop('insert_product_tmp', 0);

            if (isset($arr[1])) {

                $data = json_decode($arr[1]);

                foreach ($data->asins as $asin) {
                    echo "\n$asin";
                    insertIgnore('product_tmp', [
                        'asin'     => $asin,
                        'store_id' => $data->store_id,
                    ]);
                }
            }
        }
    }

    public function insertproductnewAction()
    {
        while (true) {

            $arr = $this->redis->brPop('insert_product_new', 0);

            if (isset($arr[1])) {

                $data = json_decode($arr[1]);
                $old = $data->old;
                $new = $data->new;

                if ($new) {

                    printf("\n%s | %s | %s", $old->asin, $new->asin, $new->name);

                    if ($new->status !== \ProductStatus::PROCESSING) {

                        // Men's Let That Shit Go Yoga T-Shirt (3XL, Navy) -> child
                        if (preg_match('/^(Men\'s|Women\'s|Kids|Unisex)\s(.*)\s\((.*),(.*)\)$/', $new->name)) {
                            $new->status = \ProductStatus::DUPLICATED;
                        }

                        insertIgnoreProduct($new);

                        updateTmpStatus($old->asin, 1);
                    }

                } else {

                    printf("\n%s | dog page", $old->asin);

                    updateTmpStatus($old->asin, 3);
                }
            }
        }
    }

    public function postelasticAction()
    {
        $elastic = new \GuzzleHttp\Client([
            'base_uri' => 'http://10.136.170.37:9200',
            'proxy' => APPLICATION_ENV === 'production' ? null : 'http://167.99.236.195:51188',
        ]);

        while (true) {

            echo "\n\n-------------";
            echo "\nTime: " . date('Y-m-d H:i:s');
            echo "\nFetch products";
            $sql = 'SELECT
                        asin,
                        name,
                        brand,
                        price,
                        rank,
                        rating,
                        reviews,
                        node,
                        store_id,
                        status,
                        trademark,
                        img_sm,
                        img_lg,
                        first_sold_at,
                        publish_at,
                        created_at,
                        updated_at,
                        deleted_at
                    FROM products WHERE pushed = 0 LIMIT 100';

            $products = $this->db->fetchAll($sql, null);

            $count = count($products);

            echo "\nTotal: $count";

            if ($count > 0) {
                $payload = $this->createPayload($products);
                $response = $elastic->post('/spyamz_products/_bulk', [
                    'headers' => [
                        'Content-Type' => 'application/x-ndjson'
                    ],
                    'body' => $payload
                ]);
                $res = json_decode($response->getBody()->getContents());
                if (isset($res->items)) {
                    $asins = [];
                    foreach ($res->items as $i) {
                        if (isset($i->index->result)) {
                            echo "\n" . $i->index->result . ' | ' . $i->index->_id;
                            if ($i->index->result == 'created' || $i->index->result == 'updated') {
                                $asins[] = $i->index->_id;
                            }
                        } else {
                            echo "\nIndex not found";
                            print_r($i);
                            echo "\n";
                        }
                    }
                    updateProductPushed($asins);
                } else {
                    echo "\nItems not found";
                    sleep(10);
                }
            } else {
                echo "\nNo unpushed products";
                sleep(10);
            }
        }
    }

    public function notifyAction()
    {
        while (true) {

            $sql = 'SELECT id, title, content FROM messages WHERE sent = 0 AND type = 2';
            $messages = $this->db->fetchAll($sql, null);
            $count = count($messages);

            echo "\nMessages: $count | " . date('Y-m-d H:i:s');

            if ($count > 0) {

                // Create the Transport
                $transport = new \Swift_SmtpTransport('smtp.gmail.com', 587, 'tls');
                $transport->setUsername('noreply.spyamz@gmail.com');
                $transport->setPassword('G8Bz8ye5J3EhGXbU');

                // Create the Mailer using your created Transport
                $mailer = new \Swift_Mailer($transport);

                echo "\n-----------------";
                foreach ($messages as $m) {
                    // Create a message
                    $message = new \Swift_Message($m->title, $m->content, 'text/plain', 'UTF-8');
                    $message->setFrom('noreply.spyamz@gmail.com');
                    $message->setTo('phanphuc@gmail.com');

                    echo "\n" . str_pad(substr($m->title, 0, 20), 20);

                    try {
                        // Send the message
                        $mailer->send($message);

                        $this->db->updateAsDict('messages', [
                            'sent' => 1
                        ], [
                            'conditions' => 'id = ?',
                            'bind' => [$m->id],
                        ]);

                        echo ' | success';

                    } catch (\Exception $e) {

                        $this->db->updateAsDict('messages', [
                            'sent' => -1
                        ], [
                            'conditions' => 'id = ?',
                            'bind' => [$m->id],
                        ]);

                        echo ' | ' . $e->getMessage();
                    }

                    echo "\n";
                    sleep(1);
                }
            } else {
                sleep(30);
            }
        }
    }

    private function hasRankUp($old, $new)
    {
        if ($new > 0 && $new < 500000) {
            if ($new < $old) {
                if ($new < 100000) return $old - $new > 1000;
                if ($new < 300000) return $old - $new > 50000;
                if ($new < 400000) return $old - $new > 100000;
                return $old - $new > 150000;
            } else {
                return $old == 0;
            }
        }
        return false;
    }

    private function hasNewReviews($old, $new)
    {
        return $new - $old > 0;
    }

    private function createPayload($products)
    {
        $body = '';
        foreach ($products as $p) {
            $asin = $p->asin;
            unset($p->asin);

            $p->price     = (float)$p->price;
            $p->rating    = (float)$p->rating;
            $p->reviews   = (int)$p->reviews;
            $p->rank      = (int)$p->rank;
            $p->status    = (int)$p->status;
            $p->trademark = (int)$p->trademark;

            if (empty($p->publish_at)) {
                $p->publish_at = substr($p->created_at, 0, 19);
            }

            $body .= '{"index":{"_id":"'.$asin.'"}}' . "\n";
            $body .= json_encode($p) . "\n";
        }
        return $body;
    }
}
