mirror of https://github.com/OpenIdentityPlatform/OpenDJ.git

Valery Kharseko
2 days ago b9b857b0b1f3f9df35ea7ecf9e88e1e5550bd43c
opendj-server-legacy/src/main/java/org/opends/server/backends/cassandra/CASStorage.java
@@ -58,6 +58,7 @@
import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.config.DriverConfigLoader;
import com.datastax.oss.driver.api.core.cql.BoundStatement;
import com.datastax.oss.driver.api.core.cql.PreparedStatement;
import com.datastax.oss.driver.api.core.cql.ResultSet;
import com.datastax.oss.driver.api.core.cql.Row;
@@ -71,7 +72,18 @@
   private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
   private CASBackendCfg config;
   //a cursor starts with a small page so that a point lookup does not transfer thousands of rows,
   //and doubles it while a scan keeps outrunning it (the driver default page is 5000 rows)
   final int initialPageSize=Math.max(1,Integer.getInteger("org.openidentityplatform.opendj.cassandra.fetchsize.initial",32));
   final int maxPageSize=Math.max(initialPageSize,Integer.getInteger("org.openidentityplatform.opendj.cassandra.fetchsize",1000));
   //CursorImpl.forwardInPage outcomes
   final static int POSITIONED=1;
   final static int MISSING=0;
   final static int NEEDS_SEEK=-1;
   final static int PAGE_OUT=-2;
   public CASStorage(CASBackendCfg cfg, ServerContext serverContext) {
      this.config = cfg;
      cfg.addCASChangeListener(this);
@@ -292,42 +304,138 @@
      }
   }
   
   private final class CursorImpl implements Cursor<ByteString, ByteString> {
   // Iterates the (baseDN,indexId) partition in key order.
   // A ResultSet can only be consumed once: rc.iterator() always returns the same iterator,
   // so "restarting" it never rewinds. Every repositioning that cannot be served by moving
   // forward within the already-fetched rows therefore runs a new server-side slice query
   // on the "key" clustering column instead.
   // Pages are sized by the cursor rather than by the driver default of 5000 rows: a point
   // lookup must not transfer thousands of entries, so a cursor starts with a small page and
   // doubles it (up to maxPageSize) every time a scan outruns it. Scanning past the end of a
   // page runs the next slice from the current key instead of letting the driver page with the
   // size the cursor started with.
   final class CursorImpl implements Cursor<ByteString, ByteString> {
      final TreeName treeName;
      final TransactionImpl tx;
      final String tableName=getTableName();
      //visible for tests
      long queryCount;
      int pageSize=initialPageSize;
      ResultSet rc;
      Iterator<Row> iterator;
      //the row the cursor is on, or - when defined is false - the row a failed positionToKey
      //stopped just before, which is where next() resumes (same as pdb)
      Row current=null;
      boolean defined=false;
      //rc holds a DESC page: it must never serve a forward move
      boolean descending=false;
      boolean closed=false;
      public CursorImpl(TransactionImpl tx,TreeName treeName) {
         this.treeName=treeName;
         this.tx=tx;
         rc=full();
         iterator=rc.iterator();
         //lazy: the first navigation decides which query to run, a seek must not pay for a partition scan
      }
      ResultSet full(){
         return execute(
            prepared.get("SELECT key,value FROM "+getTableName()+" WHERE baseDN=:baseDN and indexId=:indexId ORDER BY key").bind()
               .setString("baseDN", treeName.getBaseDN()).setString("indexId", treeName.getIndexId())
         );
      ResultSet select(String condition,ByteSequence key,int rows){
         queryCount++;
         BoundStatement statement=prepared.get("SELECT key,value FROM "+tableName+" WHERE baseDN=:baseDN and indexId=:indexId"+condition).bind()
            .setString("baseDN", treeName.getBaseDN()).setString("indexId", treeName.getIndexId())
            .setPageSize(rows);
         if (key!=null) {
            statement=statement.setByteBuffer("key", ByteBuffer.wrap(key.toByteArray()));
         }
         return execute(statement);
      }
      //runs a new page and positions the cursor on its first row
      boolean slice(String condition,ByteSequence key){
         rc=select(condition,key,pageSize);
         iterator=rc.iterator();
         descending=false;
         if (iterator.hasNext()) {
            current=iterator.next();
            defined=true;
            return true;
         }
         current=null;
         defined=false;
         return false;
      }
      boolean seek(ByteSequence key){
         return slice(" and key>=:key ORDER BY key",key);
      }
      void growPage() {
         pageSize=(int)Math.min(maxPageSize,2L*pageSize);
      }
      //serves a forward repositioning from the rows the driver already fetched: they are the
      //sorted rows following the current one (blob clustering collates in unsigned byte order)
      int forwardInPage(ByteSequence key,boolean exactMatch) {
         if (!defined || descending || rc==null) {
            return NEEDS_SEEK;
         }
         int cmp=key.compareTo(getKey());
         if (cmp==0) {
            return POSITIONED;
         }
         if (cmp<0) { //backward: only the server can rewind
            return NEEDS_SEEK;
         }
         while (rc.getAvailableWithoutFetching()>0) {
            current=iterator.next();
            cmp=key.compareTo(getKey());
            if (cmp==0) {
               return POSITIONED;
            }
            if (cmp<0) { //walked past the key
               if (exactMatch) {
                  defined=false; //the cursor stops just before this row
                  return MISSING;
               }
               return POSITIONED;
            }
         }
         return PAGE_OUT;
      }
      @Override
      public boolean next() {
         try {
            current=iterator.next();
            return true;
         }catch (NoSuchElementException e) {
            current=null;
         if (closed) {
            return false;
         }
         return false;
         if (current!=null && !defined) { //a failed positionToKey stopped just before this row
            defined=true;
            return true;
         }
         if (rc==null) { //lazy cursor: the first navigation runs the scan
            return slice(" ORDER BY key",null);
         }
         if (!descending && rc.getAvailableWithoutFetching()>0) {
            current=iterator.next();
            defined=true;
            return true;
         }
         if (current==null) { //exhausted or explicitly undefined: stay there
            return false;
         }
         if (!descending && rc.isFullyFetched()) { //the server has nothing left either
            current=null;
            defined=false;
            return false;
         }
         //page boundary: continue with a bigger slice starting right after the current key
         growPage();
         return slice(" and key>:key ORDER BY key",getKey());
      }
      @Override
      public boolean isDefined() {
         return current!=null;
         return defined;
      }
      @Override
@@ -354,72 +462,101 @@
         tx.delete(treeName, getKey());
      }
      //a closed cursor is undefined and every navigation on it returns false, like EmptyCursor
      @Override
      public void close() {
         iterator=null;
         closed=true;
         iterator=Collections.emptyIterator();
         current=null;
         defined=false;
         descending=false;
         rc=null;
      }
      @Override
      public boolean positionToKeyOrNext(ByteSequence key) {
         if (!isDefined() || key.compareTo(getKey())<0) { //restart iterator
            iterator=rc.iterator();
         if (closed) {
            return false;
         }
         while (iterator.hasNext()) {
            current=iterator.next();
            if (key.compareTo(getKey())<=0) {
               return true;
            }
         final int served=forwardInPage(key,false);
         if (served>=0) {
            return served==POSITIONED;
         }
         current=null;
         return false;
         if (served==PAGE_OUT) { //the scan outran its page: the next slice should be bigger
            growPage();
         }
         return seek(key);
      }
      @Override
      public boolean positionToKey(ByteSequence key) {
         if (!isDefined() || key.compareTo(getKey())<0) {  //restart iterator
            iterator=rc.iterator();
         if (closed) {
            return false;
         }
         if (isDefined() && key.compareTo(getKey())==0) {
         final int served=forwardInPage(key,true);
         if (served>=0) {
            return served==POSITIONED;
         }
         if (served==PAGE_OUT) {
            growPage();
         }
         if (seek(key) && key.compareTo(getKey())==0) {
            return true;
         }
         while (iterator.hasNext()) {
            current=iterator.next();
            if (key.compareTo(getKey())==0) {
               return true;
            }
         }
         current=null;
         //like jeb/pdb a miss leaves the cursor undefined; the row the seek landed on is the
         //first key after the missing one, so next() resumes there instead of skipping it
         defined=false;
         return false;
      }
      @Override
      public boolean positionToLastKey() {
         while (iterator.hasNext()) {
            current=iterator.next();
         if (closed) {
            return false;
         }
         if (current!=null) {
         rc=select(" ORDER BY key DESC LIMIT 1",null,1);
         iterator=rc.iterator();
         descending=true;
         if (iterator.hasNext()) {
            current=iterator.next(); //nothing follows the last key
            defined=true;
            return true;
         }
         current=null;
         defined=false;
         return false;
      }
      @Override
      public boolean positionToIndex(int index) {
         iterator=rc.iterator();  //restart iterator
         int ct=0;
         while(iterator.hasNext()){
            current=iterator.next();
            if (ct==index) {
               return true;
            }
            ct++;
         if (closed) {
            return false;
         }
         current=null;
         return false;
         if (index<0) {
            rc=null; //an invalid index resets the cursor: next() starts the scan again
            iterator=null;
            current=null;
            defined=false;
            descending=false;
            return false;
         }
         //CQL has no offset clause: restart from the first row and skip. The rows that have to
         //be walked are asked for in one page instead of one round trip per page
         rc=select(" ORDER BY key",null,(int)Math.min(maxPageSize,Math.max(pageSize,index+1L)));
         iterator=rc.iterator();
         descending=false;
         for (int ct=0;ct<=index;ct++) {
            if (!iterator.hasNext()) {
               current=null;
               defined=false;
               return false;
            }
            current=iterator.next();
         }
         defined=true;
         return true;
      }
   }