package com.infomancers.phantomproducer.tests;

import org.junit.Test;

import java.util.*;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ArrayBlockingQueue;
import static java.util.concurrent.TimeUnit.SECONDS;

import com.infomancers.phantomproducer.PhantomProducer;
import com.infomancers.phantomproducer.Work;
import junit.framework.Assert;

/**
 * Copyright (c) 2009, Aviad Ben Dov
 * <p/>
 * All rights reserved.
 * <p/>
 * Redistribution and use in source and binary forms, with or without modification,
 * are permitted provided that the following conditions are met:
 * <p/>
 * 1. Redistributions of source code must retain the above copyright notice, this list
 * of conditions and the following disclaimer.
 * 2. Redistributions in binary form must reproduce the above copyright notice, this
 * list of conditions and the following disclaimer in the documentation and/or other
 * materials provided with the distribution.
 * 3. Neither the name of Infomancers, Ltd. nor the names of its contributors may be
 * used to endorse or promote products derived from this software without specific
 * prior written permission.
 * <p/>
 * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
 * "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
 * LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR
 * A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT OWNER OR
 * CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
 * EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
 * PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
 * PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
 * LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
 * NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
 * SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
 */

public class PhantomProducerTests {

    @Test
    public void sourceIsOneItem_resultOneItem_queueSize10() {
        BlockingQueue<Work<Integer>> resultsQueue = new ArrayBlockingQueue<Work<Integer>>(10);
        PhantomProducer<Integer> phantomProducer = new PhantomProducer<Integer>(resultsQueue, Arrays.asList(1), 10, 1, SECONDS);

        phantomProducer.run();

        Collection<Integer> expected = Arrays.asList(1);

        assertResultsQueueUsingPoll(expected, resultsQueue);
    }

    @Test
    public void sourceIsTwoItem_resultTwoItem_queueSize10() {
        BlockingQueue<Work<Integer>> resultsQueue = new ArrayBlockingQueue<Work<Integer>>(10);
        PhantomProducer<Integer> phantomProducer = new PhantomProducer<Integer>(resultsQueue, Arrays.asList(1, 2), 10, 1, SECONDS);

        phantomProducer.run();

        Collection<Integer> expected = Arrays.asList(1, 2);

        assertResultsQueueUsingPoll(expected, resultsQueue);
    }

    @Test
    public void sourceIsThreeItems_resultThreeItems_queueSize10() {
        BlockingQueue<Work<Integer>> resultsQueue = new ArrayBlockingQueue<Work<Integer>>(10);
        PhantomProducer<Integer> phantomProducer = new PhantomProducer<Integer>(resultsQueue, Arrays.asList(1, 2, 3), 10, 1, SECONDS);

        phantomProducer.run();

        Collection<Integer> expected = Arrays.asList(1, 2, 3);

        assertResultsQueueUsingPoll(expected, resultsQueue);
    }

    @Test
    public void sourceIsFourItems_keepingOne_onlyThreeInResult_queueSize2() throws InterruptedException {
        BlockingQueue<Work<Integer>> resultsQueue = new ArrayBlockingQueue<Work<Integer>>(10);
        //noinspection MismatchedQueryAndUpdateOfCollection
        List<Work<Integer>> keptItems = new ArrayList<Work<Integer>>();
        PhantomProducer<Integer> producer = new PhantomProducer<Integer>(resultsQueue, Arrays.asList(1, 2, 3, 4), 2, 1, SECONDS);

        Thread t = new Thread(producer);
        t.start();


        keptItems.add(resultsQueue.take());
        keptItems.add(resultsQueue.take());

        Thread.sleep(1000);

        Assert.assertEquals(0, resultsQueue.size());

        t.interrupt();

        t.join();
    }

    @Test
    public void sourceIsFourItems_keepingOne_resultFourItems_queueSize3() {

    }

    @Test(timeout = 100000)
    public void sourceIsThreeItems_resultThreeItems_queueSize2_usingThreads() throws InterruptedException {
        BlockingQueue<Work<Integer>> resultsQueue = new ArrayBlockingQueue<Work<Integer>>(10);
        PhantomProducer<Integer> phantomProducer = new PhantomProducer<Integer>(resultsQueue, Arrays.asList(1, 2, 3), 2, 1, SECONDS);

        Thread t = new Thread(phantomProducer);
        t.start();

        Collection<Integer> expected = Arrays.asList(1, 2, 3);

        assertResultsQueueUsingTake(expected, resultsQueue);

        t.join();
    }

    private <T> void assertResultsQueueUsingPoll(Collection<T> expected, BlockingQueue<Work<T>> resultsQueue) {
        int i = 0;
        for (T item : expected) {
            try {
                Assert.assertEquals("Item at " + i, item, resultsQueue.poll().item);
            } catch (Throwable t) {
                t.printStackTrace();
                Assert.fail("Item at " + i);
            }

            i++;
        }
    }

    private <T> void assertResultsQueueUsingTake(Collection<T> expected, BlockingQueue<Work<T>> resultsQueue) {
        int i = 0;
        for (T item : expected) {
            try {
                Work<T> actual;
                while ((actual = resultsQueue.poll(250, TimeUnit.MILLISECONDS)) == null) {
                    System.gc();
                }

                Assert.assertEquals("Item at " + i, item, actual.item);
            } catch (Throwable t) {
                t.printStackTrace();
                Assert.fail("Item at " + i);
            }

            i++;
        }
    }
}
