package com.infomancers.phantomproducer;

import java.lang.ref.PhantomReference;
import java.lang.ref.Reference;
import java.lang.ref.ReferenceQueue;
import java.util.HashSet;
import java.util.Set;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;

/**
 * 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 final class PhantomProducer<T> implements Runnable {
    private final BlockingQueue<Work<T>> resultsQueue;
    private final Iterable<T> source;
    private final long timeoutValue;
    private final TimeUnit timeoutUnit;
    private final Semaphore semaphore;
    private final ReferenceQueue<Object> refQueue = new ReferenceQueue<Object>();

    @SuppressWarnings({"MismatchedQueryAndUpdateOfCollection"})
    private final Set<Reference<?>> references = new HashSet<Reference<?>>();

    public PhantomProducer(BlockingQueue<Work<T>> resultsQueue, Iterable<T> source, int limit, long timeoutValue, TimeUnit timeoutUnit) {
        this.resultsQueue = resultsQueue;
        this.source = source;
        this.timeoutValue = timeoutValue;
        this.timeoutUnit = timeoutUnit;
        this.semaphore = new Semaphore(limit);
    }

    public void run() {
        // iterate work from a source of information.
        for (T item : source) {
            // stop gracefully if thread was interrupted
            if (Thread.interrupted()) {
                break;
            }

            // create keep-alive object.
            Object reference = new Object();
            // create phantom reference to the keep-alive object.
            // note that we need to keep this phantom ourselves, otherwise it
            // won't be added to the reference queue.
            references.add(new PhantomReference<Object>(reference, refQueue));

            // create the wrapper with the work
            Work<T> work = new Work<T>(item, reference);
            try {
                // try to get a lock for a while
                while (!semaphore.tryAcquire(timeoutValue, timeoutUnit)) {
                    // if we couldn't get a lock, maybe some items were released in the meantime
                    // and we can release some locks.
                    Reference<?> ref;
                    while ((ref = refQueue.poll()) != null) {
                        semaphore.release();
                        references.remove(ref);
                    }
                }

                // got a lock! add the work to the results queue to be consumed.
                // this could block too, if consumers are busy.
                resultsQueue.offer(work);
            } catch (InterruptedException e) {
                return;
            }
        }
    }
}
