12 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or 13 * FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License 14 * version 2 for more details (a copy is included in the LICENSE file that 15 * accompanied this code). 16 * 17 * You should have received a copy of the GNU General Public License version 18 * 2 along with this work; if not, write to the Free Software Foundation, 19 * Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA. 20 * 21 * Please contact Oracle, 500 Oracle Parkway, Redwood Shores, CA 94065 USA 22 * or visit www.oracle.com if you need additional information or have any 23 * questions. 24 */ 25 26 package com.sun.jndi.ldap; 27 28 import java.io.IOException; 29 import java.util.concurrent.BlockingQueue; 30 import java.util.concurrent.LinkedBlockingQueue; 31 import javax.naming.CommunicationException; 32 33 final class LdapRequest { 34 35 LdapRequest next; // Set/read in synchronized Connection methods 36 int msgId; // read-only 37 38 private int gotten = 0; 39 private BlockingQueue<BerDecoder> replies; 40 private int highWatermark = -1; 41 private boolean cancelled = false; 42 private boolean pauseAfterReceipt = false; 43 private boolean completed = false; 44 45 LdapRequest(int msgId, boolean pause) { 46 this(msgId, pause, -1); 47 } 48 49 LdapRequest(int msgId, boolean pause, int replyQueueCapacity) { 50 this.msgId = msgId; 51 this.pauseAfterReceipt = pause; 52 if (replyQueueCapacity == -1) { 53 this.replies = new LinkedBlockingQueue<BerDecoder>(); 54 } else { 55 this.replies = 56 new LinkedBlockingQueue<BerDecoder>(replyQueueCapacity); 57 highWatermark = (replyQueueCapacity * 80) / 100; // 80% capacity 58 } 59 } 60 61 synchronized void cancel() { 62 cancelled = true; 63 64 // Unblock reader of pending request 65 // Should only ever have at most one waiter 66 notify(); 67 } 68 69 synchronized boolean addReplyBer(BerDecoder ber) { 70 if (cancelled) { 71 return false; 72 } 73 74 // Add a new reply to the queue of unprocessed replies. 75 try { 76 replies.put(ber); 77 } catch (InterruptedException e) { 78 // ignore 79 } 80 81 // peek at the BER buffer to check if it is a SearchResultDone PDU 82 try { 83 ber.parseSeq(null); 84 ber.parseInt(); 85 completed = (ber.peekByte() == LdapClient.LDAP_REP_RESULT); 86 } catch (IOException e) { 87 // ignore 88 } 89 ber.reset(); 90 91 notify(); // notify anyone waiting for reply 92 /* 93 * If a queue capacity has been set then trigger a pause when the 94 * queue has filled to 80% capacity. Later, when the queue has drained 95 * then the reader gets unpaused. 96 */ 97 if (highWatermark != -1 && replies.size() >= highWatermark) { 98 return true; // trigger the pause 99 } 100 return pauseAfterReceipt; 101 } 102 103 synchronized BerDecoder getReplyBer() throws CommunicationException { 104 if (cancelled) { 105 throw new CommunicationException("Request: " + msgId + 106 " cancelled"); 107 } 108 109 /* 110 * Remove a reply if the queue is not empty. 111 * poll returns null if queue is empty. 112 */ 113 BerDecoder reply = replies.poll(); 114 return reply; 115 } 116 117 synchronized boolean hasSearchCompleted() { 118 return completed; 119 } 120 } | 12 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or 13 * FITNESS FOR A PARTICULAR PURPOSE. See the GNU General Public License 14 * version 2 for more details (a copy is included in the LICENSE file that 15 * accompanied this code). 16 * 17 * You should have received a copy of the GNU General Public License version 18 * 2 along with this work; if not, write to the Free Software Foundation, 19 * Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA. 20 * 21 * Please contact Oracle, 500 Oracle Parkway, Redwood Shores, CA 94065 USA 22 * or visit www.oracle.com if you need additional information or have any 23 * questions. 24 */ 25 26 package com.sun.jndi.ldap; 27 28 import java.io.IOException; 29 import java.util.concurrent.BlockingQueue; 30 import java.util.concurrent.LinkedBlockingQueue; 31 import javax.naming.CommunicationException; 32 import java.util.concurrent.TimeUnit; 33 34 final class LdapRequest { 35 36 private final static BerDecoder EOF = new BerDecoder(new byte[]{}, -1, 0); 37 38 LdapRequest next; // Set/read in synchronized Connection methods 39 final int msgId; // read-only 40 41 private final BlockingQueue<BerDecoder> replies; 42 private volatile boolean cancelled; 43 private boolean completed; 44 private final boolean pauseAfterReceipt; 45 46 LdapRequest(int msgId, boolean pause, int replyQueueCapacity) { 47 this.msgId = msgId; 48 this.pauseAfterReceipt = pause; 49 if (replyQueueCapacity == -1) { 50 this.replies = new LinkedBlockingQueue<>(); 51 } else { 52 this.replies = new LinkedBlockingQueue<>(8 * replyQueueCapacity / 10); 53 } 54 } 55 56 void cancel() { 57 cancelled = true; 58 replies.offer(EOF); 59 } 60 61 synchronized void close() { 62 try { 63 replies.put(EOF); 64 } catch (InterruptedException e) { 65 Thread.currentThread().interrupt(); 66 } 67 } 68 69 synchronized boolean addReplyBer(BerDecoder ber) { 70 if (cancelled) { 71 return false; 72 } 73 74 // Add a new reply to the queue of unprocessed replies. 75 try { 76 replies.put(ber); 77 } catch (InterruptedException e) { 78 // ignore 79 } 80 81 // peek at the BER buffer to check if it is a SearchResultDone PDU 82 try { 83 ber.parseSeq(null); 84 ber.parseInt(); 85 completed = (ber.peekByte() == LdapClient.LDAP_REP_RESULT); 86 } catch (IOException e) { 87 // ignore 88 } 89 ber.reset(); 90 return pauseAfterReceipt; 91 } 92 93 BerDecoder getReplyBer(long millis) throws CommunicationException, 94 InterruptedException { 95 if (cancelled) { 96 throw new CommunicationException("Request: " + msgId + 97 " cancelled"); 98 } 99 100 BerDecoder result = millis > 0 ? 101 replies.poll(millis, TimeUnit.MILLISECONDS) : replies.take(); 102 103 if (cancelled) { 104 throw new CommunicationException("Request: " + msgId + 105 " cancelled"); 106 } 107 108 if (result == LdapRequest.EOF) { 109 return null; 110 } 111 112 return result; 113 } 114 115 boolean hasSearchCompleted() { 116 return completed; 117 } 118 } |