001package jmri.jmrix.loconet;
002
003import java.util.concurrent.DelayQueue;
004import java.util.concurrent.Delayed;
005import java.util.concurrent.TimeUnit;
006import javax.annotation.Nonnull;
007import org.slf4j.Logger;
008import org.slf4j.LoggerFactory;
009
010/**
011 * Delay LocoNet messages that need to be throttled.
012 * <p>
013 * A LocoNetThrottledTransmitter object sits in front of a LocoNetInterface
014 * (e.g. TrafficHandler) and meters out specific LocoNet messages.
015 *
016 * <p>
017 * The internal Memo class is used to hold the pending message and the time it's
018 * to be sent. Time computations are in units of milliseconds, as that's all the
019 * accuracy that's needed here.
020 *
021 * @author Bob Jacobsen Copyright (C) 2009
022 */
023public class LocoNetThrottledTransmitter implements LocoNetInterface {
024
025    public LocoNetThrottledTransmitter(@Nonnull LocoNetInterface controller, boolean mTurnoutExtraSpace) {
026        this.controller = controller;
027        this.memo = controller.getSystemConnectionMemo();
028        this.mTurnoutExtraSpace = mTurnoutExtraSpace;
029
030        // calculation is needed time to send on DCC:
031        // msec*nBitsInPacket*packetRepeat/bitRate*safetyFactor
032        minInterval = 1000 * (18 + 3 * 10) * 3 / 16000 * 2;
033
034        if (mTurnoutExtraSpace) {
035            minInterval = minInterval * 4;
036        }
037
038        attachServiceThread();
039    }
040
041    /**
042     * Reference to the system connection memo.
043     */
044    LocoNetSystemConnectionMemo memo = null;
045
046    /**
047     * Set the system connection memo associated with this traffic controller.
048     *
049     * @param m associated systemConnectionMemo object
050     */
051    @Override
052    public void setSystemConnectionMemo(LocoNetSystemConnectionMemo m) {
053        log.debug("LnTrafficController set memo to {}", m.getUserName());
054        memo = m;
055    }
056
057    /**
058     * Get the system connection memo associated with this traffic controller.
059     *
060     * @return the associated systemConnectionMemo object
061     */
062    @Override
063    public LocoNetSystemConnectionMemo getSystemConnectionMemo() {
064        log.debug("getSystemConnectionMemo {} called in LnTC", memo.getUserName());
065        return memo;
066    }
067
068    boolean mTurnoutExtraSpace;
069
070    /**
071     * Request that server thread cease operation, no more messages can be sent.
072     * Note that this returns before the thread is known to be done if it still
073     * has work pending.  If you need to be sure it's done, check and wait on
074     * !running.
075     */
076    public void dispose() {
077        disposed = true;
078
079        // put a shutdown request on the queue after any existing
080        Memo m = new Memo(null, nowMSec(), TimeUnit.MILLISECONDS, false) {
081            @Override
082            boolean requestsShutDown() {
083                return true;
084            }
085        };
086        queue.add(m);
087    }
088
089    volatile boolean disposed = false;
090    volatile boolean running = false;
091
092    // interface being shadowed
093    LocoNetInterface controller;
094
095    // Forward methods to underlying interface
096    @Override
097    public void addLocoNetListener(int mask, LocoNetListener listener) {
098        controller.addLocoNetListener(mask, listener);
099    }
100
101    @Override
102    public void removeLocoNetListener(int mask, LocoNetListener listener) {
103        controller.removeLocoNetListener(mask, listener);
104    }
105
106    @Override
107    public boolean status() {
108        return controller.status();
109    }
110
111    /**
112     * Accept a message to be sent after suitable delay.
113     */
114    @Override
115    public void sendLocoNetMessage(LocoNetMessage msg) {
116        sendLocoNetMessage(msg, false);
117    }
118
119    /**
120     * Accept a message to be sent after suitable delay.
121     *
122     * @param msg  Message to send; will be updated with CRC
123     * @param requestIgnoreEcho  If true: Notify listeners on enqueing message, ignore echo from line.
124     *                           Only in effect if preference "LoconetUpdateSlotOnMessageCreation" is set.
125     */
126    @Override
127    public void sendLocoNetMessage(LocoNetMessage msg, boolean requestIgnoreEcho) {
128        if (disposed) {
129            log.error("Message sent after queue disposed");
130            return;
131        }
132
133        long sendTime = calcSendTimeMSec();
134
135        Memo m = new Memo(msg, sendTime, TimeUnit.MILLISECONDS, requestIgnoreEcho);
136        queue.add(m);
137
138    }
139
140    // minimum time in msec between messages
141    long minInterval;
142
143    long lastSendTimeMSec = 0;
144
145    long calcSendTimeMSec() {
146        // next time is at least now or minInterval after latest so far
147        lastSendTimeMSec = Math.max(nowMSec(), minInterval + lastSendTimeMSec);
148        return lastSendTimeMSec;
149    }
150
151    DelayQueue<Memo> queue = new DelayQueue<Memo>();
152
153    /**
154     * Constant for the name of the Service Thread.
155     * Requires the connection UserName prepending.
156     */
157    public static final String SERVICE_THREAD_NAME = " LocoNetThrottledTransmitter";
158
159    private void attachServiceThread() {
160        theServiceThread = new ServiceThread();
161        theServiceThread.setPriority(Thread.NORM_PRIORITY);
162        theServiceThread.setName( memo.getUserName() + SERVICE_THREAD_NAME);
163        theServiceThread.setDaemon(true);
164        theServiceThread.start();
165    }
166
167    ServiceThread theServiceThread;
168
169    class ServiceThread extends Thread {
170
171        @Override
172        public void run() {
173            running = true;
174            while (true) {
175                try {
176                    Memo m = queue.take();
177
178                    // check for request to shutdown
179                    if (m.requestsShutDown()) {
180                        log.debug("item requests shutdown");
181                        break;
182                    }
183
184                    // normal request
185                    if (log.isDebugEnabled()) {
186                        log.debug("forwarding message: {}", m.getMessage());
187                    }
188                    controller.sendLocoNetMessage(m.getMessage(), m.getRequestIgnoreEcho());
189                    // and go round again
190                } catch (InterruptedException e) {
191                    // request to terminate
192                    this.interrupt();
193                    break;
194                }
195            }
196            running = false;
197        }
198    }
199
200    // a separate method to ease testing by stopping clock
201    static long nowMSec() {
202        return System.currentTimeMillis();
203    }
204
205    static class Memo implements Delayed {
206
207        public Memo(LocoNetMessage msg, long endTime, TimeUnit unit, boolean requestIgnoreEcho) {
208            this.msg = msg;
209            this.endTimeMsec = unit.toMillis(endTime);
210            this.requestIgnoreEcho = requestIgnoreEcho;
211        }
212
213        LocoNetMessage getMessage() {
214            return msg;
215        }
216
217        boolean requestsShutDown() {
218            return false;
219        }
220
221        boolean getRequestIgnoreEcho() {
222            return requestIgnoreEcho;
223        }
224
225        long endTimeMsec;
226        LocoNetMessage msg;
227        boolean requestIgnoreEcho;
228
229        @Override
230        public long getDelay(TimeUnit unit) {
231            long delay = endTimeMsec - nowMSec();
232            return unit.convert(delay, TimeUnit.MILLISECONDS);
233        }
234
235        @Override
236        public int compareTo(Delayed d) {
237            // -1 means this is less than m
238            long delta;
239            if (d instanceof Memo) {
240                delta = this.endTimeMsec - ((Memo)d).endTimeMsec;
241            } else {
242                delta = this.getDelay(TimeUnit.MILLISECONDS)
243                        - d.getDelay(TimeUnit.MILLISECONDS);
244            }
245            if (delta > 0) {
246                return 1;
247            } else if (delta < 0) {
248                return -1;
249            } else {
250                return 0;
251            }
252        }
253
254        // ensure consistent with compareTo
255        @Override
256        public boolean equals(Object o) {
257            if (o == null) {
258                return false;
259            }
260            if (o instanceof Delayed) {
261                return (compareTo((Delayed) o) == 0);
262            } else {
263                return false;
264            }
265        }
266
267        @Override
268        public int hashCode() {
269            return (int) (this.getDelay(TimeUnit.MILLISECONDS) & 0xFFFFFF);
270        }
271    }
272
273    private static final Logger log = LoggerFactory.getLogger(LocoNetThrottledTransmitter.class);
274
275}