aweisberg commented on code in PR #4708:
URL: https://github.com/apache/cassandra/pull/4708#discussion_r3553765158


##########
src/java/org/apache/cassandra/locator/SatelliteReplicationStrategy.java:
##########
@@ -513,4 +631,1286 @@ public boolean 
hasSameSettings(AbstractReplicationStrategy other)
                primaryDC.equals(otherStrategy.primaryDC) &&
                disabledDCs.equals(otherStrategy.disabledDCs);
     }
+
+    private AbstractReplicationStrategy strategy(ClusterMetadata metadata)
+    {
+        return 
metadata.schema.getKeyspaceMetadata(keyspaceName).replicationStrategy;
+    }
+
+    static abstract class CoordinationPlanner<E extends Endpoints<E>, L 
extends ReplicaLayout<E>, P extends ReplicaPlan<E, P>>
+    {
+        final ClusterMetadata metadata;
+        final Keyspace keyspace;
+        final ConsistencyLevel cl;
+        final SatelliteReplicationStrategy strategy;
+        final String primary;
+        final String satellite;
+        final List<String> dcs;
+        final L fullLayout;
+        final L liveLayout;
+        final Map<String, L> fullLayouts;
+        final Map<String, L> liveLayouts;
+        final Map<String, E> fullEndpoints;
+        final Map<String, E> liveEndpoints;
+
+        public CoordinationPlanner(ClusterMetadata metadata, Keyspace 
keyspace, ConsistencyLevel cl, SatelliteReplicationStrategy strategy, String 
primary, L fullLayout)
+        {
+            this.metadata = metadata;
+            this.keyspace = keyspace;
+            this.cl = cl;
+            this.strategy = strategy;
+            this.primary = primary;
+            this.satellite = strategy.getSatelliteForDC(primary);
+
+            this.dcs = createDcList();
+            Preconditions.checkState(!dcs.isEmpty(), "No DCs available for 
request (primary=%s, all disabled?)", primary);
+
+            Set<String> dcSet = new HashSet<>(dcs);
+            this.fullLayout = filterLayout(fullLayout, rp -> 
dcSet.contains(metadata.locator.location(rp.endpoint()).datacenter));
+            this.liveLayout = filterLayout(this.fullLayout, 
FailureDetector.isReplicaAlive);
+
+            this.fullLayouts = Maps.newHashMapWithExpectedSize(dcs.size());
+            this.liveLayouts = Maps.newHashMapWithExpectedSize(dcs.size());
+
+            this.fullEndpoints = Maps.newHashMapWithExpectedSize(dcs.size());
+            this.liveEndpoints = Maps.newHashMapWithExpectedSize(dcs.size());
+
+            populateDcLayouts();
+        }
+
+        AbstractReplicationStrategy strategy(ClusterMetadata metadata)
+        {
+            return 
metadata.schema.getKeyspaceMetadata(keyspace.getName()).replicationStrategy;
+        }
+
+        List<String> createDcList()
+        {
+            List<String> results = new ArrayList<>();
+            if (!strategy.disabledDCs.contains(primary))
+                results.add(primary);
+
+            if (satellite != null && !strategy.disabledDCs.contains(satellite))
+                results.add(satellite);
+
+            for (String dc : strategy.fullDCs.keySet())
+            {
+                if (dc.equals(primary) || dc.equals(satellite) || 
strategy.disabledDCs.contains(dc))
+                    continue;
+                results.add(dc);
+            }
+
+            Preconditions.checkState(!results.isEmpty());
+            return results;
+        }
+
+        void populateDcLayouts()
+        {
+            for (String dc : dcs)
+            {
+                Predicate<Replica> dcFilter = rp -> 
metadata.locator.location(rp.endpoint()).datacenter.equals(dc);
+                L full = filterLayout(fullLayout, dcFilter);
+                L live = filterLayout(liveLayout, dcFilter);
+
+                fullLayouts.put(dc, full);
+                liveLayouts.put(dc, live);
+
+                fullEndpoints.put(dc, full.natural());
+                liveEndpoints.put(dc, live.natural());
+            }
+
+            if (ADDL_CHECKS_ENABLED)
+                Preconditions.checkState(fullLayouts.size() == dcs.size(), 
"populateDcLayouts: fullLayouts.size()=%s != dcs.size()=%s", 
fullLayouts.size(), dcs.size());
+        }
+
+        ResponseTrackerBuilder<E, L, P> createResponseTrackerBuilder()
+        {
+            switch (cl)
+            {
+                case ONE:
+                case LOCAL_ONE:
+                case NODE_LOCAL:
+                case ANY:
+                case TWO:
+                case THREE:
+                case ALL:
+                    return new ResponseTrackerBuilder.Simple<>(this);
+
+                case QUORUM:
+                case LOCAL_QUORUM:
+                    // LOCAL_QUORUM an QUORUM mean the same thing in SRS
+                    return new ResponseTrackerBuilder.Quorum<>(this);
+
+                case EACH_QUORUM:
+                    return new ResponseTrackerBuilder.EachQuorum<>(this);
+
+                default:
+                    throw new IllegalStateException("Unsupported consistency 
level: " + cl);
+            }
+        }
+
+        ResponseTracker createResponseTracker()
+        {
+            ResponseTrackerBuilder<E, L, P> builder = 
createResponseTrackerBuilder();
+            return builder.build();
+        }
+
+        abstract L filterLayout(L layout, Predicate<Replica> predicate);
+        abstract ResponseTracker createDcResponseTracker(String dc);
+
+        abstract P recomputePlan(ClusterMetadata metadata);
+        abstract P createReplicaPlan();
+
+        interface IndexReadExceptionCheck<E extends Endpoints<E>>
+        {
+            void maybeThrowIndexReadException(int initial, int filtered, 
Map<String, E> candidates, Map<InetAddressAndPort, Index.Status> 
indexStatusMap);
+        }
+
+        static class ForWrite extends CoordinationPlanner<EndpointsForToken, 
ReplicaLayout.ForTokenWrite, ReplicaPlan.ForWrite>
+        {
+            final ReplicaPlans.Selector selector;
+
+            public ForWrite(ClusterMetadata metadata, Keyspace keyspace, 
ConsistencyLevel cl, SatelliteReplicationStrategy strategy, String primary, 
ReplicaLayout.ForTokenWrite fullLayout, ReplicaPlans.Selector selector)
+            {
+                super(metadata, keyspace, cl, strategy, primary, fullLayout);
+                this.selector = selector;
+                if (ADDL_CHECKS_ENABLED)
+                {
+                    
Preconditions.checkState(!strategy.disabledDCs.contains(primary));
+                }
+            }
+
+            @Override
+            ResponseTracker createDcResponseTracker(String dc)
+            {
+                // this basically hardcodes ReplicaPlans.writeAll
+                Predicate<InetAddressAndPort> dcFilter = ep -> 
metadata.locator.location(ep).datacenter.equals(dc);
+
+                ReplicaLayout.ForTokenWrite dcLayout = fullLayouts.get(dc);
+                Preconditions.checkNotNull(dcLayout);
+
+                int blockFor = strategy.calculateQuorum(dc);
+                int totalContacts = dcLayout.all().size();
+
+                // Only count pending replicas that are actually in the 
selected contacts
+                int pendingInDc = dcLayout.pending().size();
+
+                if (pendingInDc == 0)
+                {
+                    return new SimpleResponseTracker(blockFor, totalContacts, 
dcFilter);
+                }
+                else
+                {
+                    int committedContacts = totalContacts - pendingInDc;
+                    int totalBlockFor = blockFor + pendingInDc;
+                    Predicate<InetAddressAndPort> isPending = ep -> 
dcLayout.pending.endpoints().contains(ep);
+                    return new WriteResponseTracker(blockFor, totalBlockFor,
+                                                    committedContacts, 
pendingInDc,
+                                                    isPending, dcFilter);
+                }
+            }
+
+            @Override
+            ReplicaPlan.ForWrite recomputePlan(ClusterMetadata metadata)
+            {
+                Token token = fullLayout.token();
+                return strategy(metadata).planForWrite(metadata, keyspace, cl,
+                                                              cm -> 
ReplicaLayout.forTokenWriteLiveAndDown(cm, keyspace, token),
+                                                       selector).replicas();
+            }
+
+            @Override
+            ReplicaLayout.ForTokenWrite 
filterLayout(ReplicaLayout.ForTokenWrite layout, Predicate<Replica> predicate)
+            {
+                return layout.filter(predicate);
+            }
+
+            private ReplicaLayout.ForTokenWrite 
mergeLayouts(Collection<ReplicaLayout.ForTokenWrite> layouts)
+            {
+                ReplicaCollection.Builder<EndpointsForToken> naturalBuilder = 
null;
+                ReplicaCollection.Builder<EndpointsForToken> pendingBuilder = 
null;
+
+                for (ReplicaLayout.ForTokenWrite layout : layouts)
+                {
+                    EndpointsForToken natural = layout.natural();
+                    EndpointsForToken pending = layout.pending();
+                    if (naturalBuilder == null)
+                    {
+                        naturalBuilder = natural.newBuilder(natural.size());
+                        pendingBuilder = pending.newBuilder(pending.size());
+                    }
+
+                    naturalBuilder.addAll(natural);
+                    pendingBuilder.addAll(pending);
+                }
+
+                return new ReplicaLayout.ForTokenWrite(strategy, 
naturalBuilder.build(), pendingBuilder.build());
+            }
+
+            @Override
+            ReplicaPlan.ForWrite createReplicaPlan()
+            {
+
+                ReplicaLayout.ForTokenWrite fullLayout = 
mergeLayouts(fullLayouts.values());
+                ReplicaLayout.ForTokenWrite liveLayout = 
mergeLayouts(liveLayouts.values());
+                EndpointsForToken contacts = selector.select(cl, fullLayout, 
liveLayout);
+
+                return new ReplicaPlan.ForWrite(keyspace,
+                                                strategy,
+                                                cl,
+                                                fullLayout.pending(),
+                                                fullLayout.all(),
+                                                liveLayout.all(),
+                                                contacts,
+                                                this::recomputePlan,
+                                                metadata.epoch);
+            }
+        }
+
+        static abstract class ForRead<E extends Endpoints<E>, L extends 
ReplicaLayout<E>, P extends ReplicaPlan.ForRead<E, P>> extends 
CoordinationPlanner<E, L, P>
+        {
+            final @Nullable Index.QueryPlan indexQueryPlan;
+            final boolean alwaysSpeculate;
+            final boolean throwOnInsufficientLiveReplicas;
+
+            public ForRead(ClusterMetadata metadata, Keyspace keyspace, 
ConsistencyLevel cl, SatelliteReplicationStrategy strategy, String primary, L 
fullLayout, @Nullable Index.QueryPlan indexQueryPlan, boolean alwaysSpeculate, 
boolean throwOnInsufficientLiveReplicas)
+            {
+                super(metadata, keyspace, cl, strategy, primary, fullLayout);
+                this.indexQueryPlan = indexQueryPlan;
+                this.alwaysSpeculate = alwaysSpeculate;
+                this.throwOnInsufficientLiveReplicas = 
throwOnInsufficientLiveReplicas;
+            }
+
+            @Override
+            ResponseTracker createDcResponseTracker(String dc)
+            {
+                Predicate<InetAddressAndPort> dcFilter = ep -> 
metadata.locator.location(ep).datacenter.equals(dc);
+
+                L dcLayout = fullLayouts.get(dc);
+                int blockFor = strategy.calculateQuorum(dc);
+                int totalContacts = dcLayout.all().size();
+
+                return new SimpleResponseTracker(blockFor, totalContacts, 
dcFilter);
+            }
+
+            Map<String, E> candidates(IndexReadExceptionCheck<E> check)
+            {
+                if (indexQueryPlan == null)
+                    return liveEndpoints;
+
+                Map<String, E> candidates = new HashMap<>();
+                Map<InetAddressAndPort, Index.Status> indexStatusMap = new 
HashMap<>();
+                int initial = 0;
+                int filtered = 0;
+                for (Map.Entry<String, E> entry : liveEndpoints.entrySet())
+                {
+                    E liveEndpoints = entry.getValue();
+                    E queryableEndpoints = 
IndexStatusManager.instance.filterForQuery(liveEndpoints, keyspace.getName(), 
indexQueryPlan, indexStatusMap);
+                    initial += liveEndpoints.size();
+                    filtered += queryableEndpoints.size();
+                    candidates.put(entry.getKey(), queryableEndpoints);
+                }
+
+                if (initial != filtered)
+                    check.maybeThrowIndexReadException(initial, filtered, 
candidates, indexStatusMap);
+
+                return candidates;
+            }
+
+            @Override
+            ResponseTracker createResponseTracker()
+            {
+                ResponseTracker tracker = super.createResponseTracker();
+
+                // read trackers don't wait on pending nodes
+                Preconditions.checkState(tracker.pendingContacts() == 0);
+                return tracker;
+            }
+
+
+            abstract P replicaPlanBuilder(E candidates, E contacts, E 
liveAndDown, int quorumSize);
+
+            ReadReplicaPlanBuilder<E, L, P> replicaPlanBuilder()
+            {
+                switch (cl)
+                {
+                    case ONE:
+                    case LOCAL_ONE:
+                    case NODE_LOCAL:
+                    case ANY:
+                    case TWO:
+                    case THREE:
+                    case ALL:
+                        return new ReadReplicaPlanBuilder.Simple<>(this);
+
+                    case QUORUM:
+                    case LOCAL_QUORUM:
+                    case EACH_QUORUM:
+                        return new ReadReplicaPlanBuilder.DcQuorum<>(this);
+
+                    default:
+                        throw new IllegalStateException("Unsupported 
consistency level: " + cl);
+                }
+            }
+
+            @Override
+            P createReplicaPlan()
+            {
+                ReadReplicaPlanBuilder<E, L, P> builder = replicaPlanBuilder();
+                return builder.build();
+            }
+        }
+
+        static class ForTokenRead extends ForRead<EndpointsForToken, 
ReplicaLayout.ForTokenRead, ReplicaPlan.ForTokenRead>
+        {
+            static final Function<ReplicaPlan<?, ?>, ReplicaPlan.ForWrite> 
READ_REPAIR_PLAN = (ReplicaPlan<?, ?> self) -> {throw new 
IllegalStateException("Cannot create read repair plans for 
SatelliteDatacenterReplicationStrategy");};
+            final Token token;
+            final TableId tableId;
+            final SpeculativeRetryPolicy retry;
+            final ReadCoordinator coordinator;
+
+            public ForTokenRead(ClusterMetadata metadata, Keyspace keyspace, 
Token token, ConsistencyLevel cl, SatelliteReplicationStrategy strategy, String 
primary, ReplicaLayout.ForTokenRead fullLayout, @Nullable Index.QueryPlan 
indexQueryPlan, SpeculativeRetryPolicy retry, boolean 
throwOnInsufficientLiveReplicas, TableId tableId, ReadCoordinator coordinator)
+            {
+                super(metadata, keyspace, cl, strategy, primary, fullLayout, 
indexQueryPlan, retry == AlwaysSpeculativeRetryPolicy.INSTANCE, 
throwOnInsufficientLiveReplicas);
+                this.token = token;
+                this.tableId = tableId;
+                this.retry = retry;
+                this.coordinator = coordinator;
+            }
+
+            @Override
+            ReplicaLayout.ForTokenRead filterLayout(ReplicaLayout.ForTokenRead 
layout, Predicate<Replica> predicate)
+            {
+                return layout.filter(predicate);
+            }
+
+            @Override
+            ReplicaPlan.ForTokenRead recomputePlan(ClusterMetadata metadata)
+            {
+                return strategy(metadata).planForTokenRead(metadata, keyspace, 
tableId, token,
+                                                           indexQueryPlan, cl, 
retry, coordinator).replicas();
+            }
+
+            @Override
+            ReplicaPlan.ForTokenRead replicaPlanBuilder(EndpointsForToken 
candidates, EndpointsForToken contacts, EndpointsForToken liveAndDown, int 
quorumSize)
+            {
+                return new ReplicaPlan.ForTokenRead(keyspace, strategy, cl,
+                                                    candidates, contacts, 
liveAndDown,
+                                                    this::recomputePlan,
+                                                    READ_REPAIR_PLAN,
+                                                    metadata.epoch,
+                                                    quorumSize);
+            }
+        }
+
+        static class ForRangeRead extends ForRead<EndpointsForRange, 
ReplicaLayout.ForRangeRead, ReplicaPlan.ForRangeRead>
+        {
+            static final BiFunction<ReplicaPlan<?, ?>, Token, 
ReplicaPlan.ForWrite> READ_REPAIR_PLAN = (ReplicaPlan<?, ?> self, Token token) 
-> {throw new IllegalStateException("Cannot create read repair plans for 
SatelliteDatacenterReplicationStrategy");};
+            final AbstractBounds<PartitionPosition> range;
+            final int vnodeCount;
+            final TableId tableId;
+
+            public ForRangeRead(ClusterMetadata metadata, Keyspace keyspace, 
AbstractBounds<PartitionPosition> range, int vnodeCount, ConsistencyLevel cl, 
SatelliteReplicationStrategy strategy, String primary, 
ReplicaLayout.ForRangeRead fullLayout, @Nullable Index.QueryPlan 
indexQueryPlan, boolean throwOnInsufficientLiveReplicas, TableId tableId)
+            {
+                super(metadata, keyspace, cl, strategy, primary, fullLayout, 
indexQueryPlan, false, throwOnInsufficientLiveReplicas);
+                this.range = range;
+                this.vnodeCount = vnodeCount;
+                this.tableId = tableId;
+            }
+
+            @Override
+            ReplicaLayout.ForRangeRead filterLayout(ReplicaLayout.ForRangeRead 
layout, Predicate<Replica> predicate)
+            {
+                return layout.filter(predicate);
+            }
+
+            @Override
+            ReplicaPlan.ForRangeRead recomputePlan(ClusterMetadata metadata)
+            {
+                return strategy(metadata).planForRangeRead(metadata, keyspace, 
tableId,
+                                                           indexQueryPlan, cl, 
range, vnodeCount).replicas();
+            }
+
+            @Override
+            ReplicaPlan.ForRangeRead replicaPlanBuilder(EndpointsForRange 
candidates, EndpointsForRange contacts, EndpointsForRange liveAndDown, int 
quorumSize)
+            {
+                return new ReplicaPlan.ForRangeRead(keyspace, strategy, cl, 
range,
+                                                    candidates, contacts, 
liveAndDown, vnodeCount,
+                                                    this::recomputePlan,
+                                                    READ_REPAIR_PLAN,
+                                                    metadata.epoch,
+                                                    quorumSize);
+            }
+        }
+    }
+
+    static abstract class ResponseTrackerBuilder<E extends Endpoints<E>, L 
extends ReplicaLayout<E>, P extends ReplicaPlan<E, P>>
+    {
+        final CoordinationPlanner<E, L, P> planner;
+
+        public ResponseTrackerBuilder(CoordinationPlanner<E, L, P> planner)
+        {
+            this.planner = planner;
+        }
+
+        Map<String, ResponseTracker> responseTrackersForPrimary()
+        {
+            Map<String, ResponseTracker> result = 
Maps.newHashMapWithExpectedSize(planner.dcs.size());
+
+            for (String dc : planner.dcs)
+                result.put(dc, planner.createDcResponseTracker(dc));
+
+            return result;
+        }
+
+        abstract ResponseTracker aggregateTrackers(Map<String, 
ResponseTracker> dcTrackers);
+
+        ResponseTracker build()
+        {
+            Map<String, ResponseTracker> dcTrackers = 
responseTrackersForPrimary();
+            return aggregateTrackers(dcTrackers);
+        }
+
+        static class Simple<E extends Endpoints<E>, L extends 
ReplicaLayout<E>, P extends ReplicaPlan<E, P>> extends 
ResponseTrackerBuilder<E, L, P>
+        {
+            public Simple(CoordinationPlanner<E, L, P> planner)
+            {
+                super(planner);
+            }
+
+            @Override
+            ResponseTracker aggregateTrackers(Map<String, ResponseTracker> 
dcTrackers)
+            {
+                int totalContacts = 0;
+                int totalPending = 0;
+
+
+                List<ResponseTracker> trackers = Lists.newArrayList();
+                Set<String> blockingDatacenters = Sets.newHashSet();
+
+                for (Map.Entry<String, ResponseTracker> entry : 
dcTrackers.entrySet())
+                {
+                    String dc = entry.getKey();
+                    ResponseTracker tracker = entry.getValue();
+                    if (planner.cl == ConsistencyLevel.ALL || 
dc.equals(planner.primary) || dc.equals(planner.satellite))
+                    {
+                        totalContacts += tracker.totalContacts();
+                        totalPending += tracker.pendingContacts();
+                        trackers.add(tracker);
+                        blockingDatacenters.add(dc);
+                    }
+                }
+
+                Predicate<InetAddressAndPort> filter = i -> {
+                    String dc = 
planner.metadata.locator.location(i).datacenter;
+                    return dc != null && blockingDatacenters.contains(dc);
+                };
+
+                Preconditions.checkState(!trackers.isEmpty(), "No blocking DCs 
for CL %s (primary=%s, satellite=%s)", planner.cl, planner.primary, 
planner.satellite);
+
+
+                if (planner.cl == ConsistencyLevel.ALL)
+                    return new SimpleResponseTracker(totalContacts, 
totalContacts, filter);
+
+                int blockFor = planner.cl.blockFor(planner.strategy);
+                int totalBlockFor = blockFor + totalPending;
+
+                Preconditions.checkState(blockFor <= totalContacts, "Simple 
tracker: blockFor=%s > totalContacts=%s for CL %s", blockFor, totalContacts, 
planner.cl);
+
+
+                if (totalPending == 0)
+                    return new SimpleResponseTracker(blockFor, totalContacts, 
filter);
+
+                Predicate<InetAddressAndPort> pending = i -> {
+                    for (ResponseTracker tracker : trackers)
+                    {
+                        if (tracker.isPending(i))
+                            return true;
+                    }
+                    return false;
+                };
+                return new WriteResponseTracker(blockFor, totalBlockFor, 
totalContacts, totalPending, pending, filter);
+            }
+        }
+
+        static class Quorum<E extends Endpoints<E>, L extends 
ReplicaLayout<E>, P extends ReplicaPlan<E, P>> extends 
ResponseTrackerBuilder<E, L, P>
+        {
+            public Quorum(CoordinationPlanner<E, L, P> planner)
+            {
+                super(planner);
+            }
+
+            @Override
+            ResponseTracker aggregateTrackers(Map<String, ResponseTracker> 
dcTrackers)
+            {
+                return new 
CompositeTracker(CompositeTracker.quorum(dcTrackers.size()), new 
ArrayList<>(dcTrackers.values()));
+            }
+        }
+
+        static class EachQuorum<E extends Endpoints<E>, L extends 
ReplicaLayout<E>, P extends ReplicaPlan<E, P>> extends 
ResponseTrackerBuilder<E, L, P>
+        {
+            public EachQuorum(CoordinationPlanner<E, L, P> planner)
+            {
+                super(planner);
+            }
+
+            @Override
+            ResponseTracker aggregateTrackers(Map<String, ResponseTracker> 
dcTrackers)
+            {
+                return new 
CompositeTracker(CompositeTracker.all(dcTrackers.size()), new 
ArrayList<>(dcTrackers.values()));
+            }
+        }
+    }
+
+    static abstract class ReadReplicaPlanBuilder<E extends Endpoints<E>, L 
extends ReplicaLayout<E>, P extends ReplicaPlan.ForRead<E, P>>
+            implements CoordinationPlanner.IndexReadExceptionCheck<E>
+    {
+
+        private static <E extends Endpoints<E>> Iterable<Replica> 
removedSupplier(Map<String, E> initial, Map<String, E> filtered)
+        {
+            return () -> initial.keySet().stream().map(dc -> {
+                Set<InetAddressAndPort> f =  filtered.containsKey(dc) ? 
filtered.get(dc).endpoints() : Collections.emptySet();
+                return initial.get(dc).without(f);
+            }).flatMap(AbstractReplicaCollection::stream).iterator();
+        }
+
+        final CoordinationPlanner.ForRead<E, L, P> planner;
+
+        public ReadReplicaPlanBuilder(CoordinationPlanner.ForRead<E, L, P> 
planner)
+        {
+            this.planner = planner;
+        }
+
+        abstract boolean canSelectReplicasFromDc(String dc, int liveNodes);
+        abstract int minBlockFor();
+        abstract int eligibleCandidates(Map<String, E> candidates);
+
+        abstract boolean haveSufficientContacts(int contactedReplicas, int 
contactedDcs);
+        abstract boolean haveSufficientContactsForDc(String dc, int 
contactedReplicas, int contactedDcs, int dcContacts);
+        abstract boolean haveSufficientLiveNodes(int contactedReplicas, int 
contactedDcs);
+
+        private E allLiveAndDown()
+        {
+            E fullEndpoints = planner.fullLayout.all();
+            ReplicaCollection.Builder<E> fullBuilder = 
fullEndpoints.newBuilder(fullEndpoints.size());
+
+            for (String dc : planner.dcs)
+            {
+                E endpoints = planner.fullEndpoints.get(dc);
+                if (endpoints == null)
+                    continue;
+
+                fullBuilder.addAll(endpoints);
+            }
+
+            return fullBuilder.build();
+        }
+
+        private E allCandidates(Map<String, E> candidates)
+        {
+            E liveEndpoints = planner.liveLayout.all();
+            ReplicaCollection.Builder<E> candidateBuilder = 
liveEndpoints.newBuilder(liveEndpoints.size());
+
+            for (String dc : planner.dcs)
+            {
+                // TODO: make sure the full replica is added if this is the dc 
it's from
+                E dcCandidates = candidates.get(dc);
+                if (dcCandidates == null || !canSelectReplicasFromDc(dc, 
dcCandidates.size()))
+                    continue;
+
+                candidateBuilder.addAll(dcCandidates);
+            }
+
+            return candidateBuilder.build();
+        }
+
+        private E contacts(Map<String, E> candidates)
+        {
+            E liveEndpoints = planner.liveLayout.all();
+
+            ReplicaCollection.Builder<E> contactBuilder = 
liveEndpoints.newBuilder(liveEndpoints.size());
+
+            // find full replica
+            Replica fullReplica = null;
+            String fullDc = null;
+            for (String dc : planner.dcs)
+            {
+                if (fullReplica != null)
+                    break;
+
+                E dcCandidates = candidates.get(dc);
+                if (dcCandidates == null || !canSelectReplicasFromDc(dc, 
dcCandidates.size()))
+                    continue;
+
+                for (Replica replica : dcCandidates)
+                {
+                    if (replica.isFull())
+                    {
+                        fullReplica = replica;
+                        fullDc = dc;
+                        break;
+                    }
+                }
+            }
+
+            if (fullReplica == null)
+            {
+                // No full replicas available - throw error similar to 
assureSufficientLiveReplicas
+                throw UnavailableException.create(ConsistencyLevel.ONE, 
minBlockFor(), 1, eligibleCandidates(candidates), 0);
+            }
+
+            int totalContacts = 0;
+            int totalContactedDcs = 0;
+
+            // add dc with full replicas first
+            {
+                E dcCandidates = candidates.get(fullDc);
+                Preconditions.checkNotNull(dcCandidates);
+
+                contactBuilder.add(fullReplica);
+                totalContacts++;
+                int dcContacts = 1;
+
+                if (!haveSufficientContacts(totalContacts, totalContactedDcs))
+                {
+                    for (Replica replica : dcCandidates)
+                    {
+                        if (replica.equals(fullReplica))
+                            continue;
+
+                        if (haveSufficientContactsForDc(fullDc, totalContacts, 
totalContactedDcs, dcContacts))
+                            break;
+
+                        contactBuilder.add(replica);
+                        totalContacts++;
+                        dcContacts++;
+                    }
+                }
+                totalContactedDcs++;
+            }
+
+            for (String dc : planner.dcs)
+            {
+                if (dc.equals(fullDc))
+                    continue;
+
+                E dcCandidates = candidates.get(dc);
+                if (dcCandidates == null || dcCandidates.isEmpty())
+                    continue;
+
+                if (haveSufficientContacts(totalContacts, totalContactedDcs))
+                    break;
+
+                int dcContacts = 0;
+                for (Replica replica : dcCandidates)
+                {
+                    if (haveSufficientContactsForDc(dc, totalContacts, 
totalContactedDcs, dcContacts))
+                        break;
+
+                    contactBuilder.add(replica);
+                    totalContacts++;
+                    dcContacts++;
+                }
+
+                totalContactedDcs++;
+            }
+
+            if (planner.throwOnInsufficientLiveReplicas && 
!haveSufficientLiveNodes(totalContacts, totalContactedDcs))
+            {
+                int fullReplicas = 0;
+                for (String dc : planner.dcs)
+                {
+                    E dcCandidates = candidates.get(dc);
+                    if (dcCandidates == null || !canSelectReplicasFromDc(dc, 
dcCandidates.size()))
+                        continue;
+                    for (Replica replica : dcCandidates)
+                        if (replica.isFull())
+                            fullReplicas++;
+                }
+
+                throw UnavailableException.create(planner.cl, minBlockFor(), 
1, eligibleCandidates(candidates), fullReplicas);
+            }
+
+            E result = contactBuilder.build();
+
+            if (ADDL_CHECKS_ENABLED)
+            {
+                // check for duplicate contacts
+                Set<InetAddressAndPort> seen = new HashSet<>();
+                for (Replica r : result)
+                    Preconditions.checkState(seen.add(r.endpoint()), 
"Duplicate contact: %s", r.endpoint());
+            }
+
+            return result;
+        }
+
+        public P build()
+        {
+            Map<String, E> candidates = planner.candidates(this);
+
+            E allEndpoints = allLiveAndDown();
+            E allCandidates = allCandidates(candidates);
+            E allContacts = contacts(candidates);
+
+            if (ADDL_CHECKS_ENABLED)
+            {
+                // contacts must be a subset of candidates
+                Set<InetAddressAndPort> candidateSet = 
allCandidates.endpoints();
+                for (Replica contact : allContacts)
+                    
Preconditions.checkState(candidateSet.contains(contact.endpoint()), "Contact %s 
not in candidates", contact.endpoint());
+            }
+
+            return planner.replicaPlanBuilder(allCandidates, allContacts, 
allEndpoints, minBlockFor());
+        }
+
+        static class Simple<E extends Endpoints<E>, L extends 
ReplicaLayout<E>, P extends ReplicaPlan.ForRead<E, P>> extends 
ReadReplicaPlanBuilder<E, L, P>
+        {
+            private final int blockFor;
+            private final int targetContacts;
+
+            public Simple(CoordinationPlanner.ForRead<E, L, P> planner)
+            {
+                super(planner);
+
+                if (planner.cl == ConsistencyLevel.ALL)
+                {
+                    int required = 0;
+                    for (E replicas : planner.fullEndpoints.values())
+                        required += replicas.size();
+
+                    this.blockFor = required;
+                }
+                else
+                {
+                    this.blockFor = planner.cl.blockFor(planner.strategy);
+                }
+
+                this.targetContacts = planner.alwaysSpeculate ? blockFor + 1 : 
blockFor;
+            }
+
+            @Override
+            public void maybeThrowIndexReadException(int initial, int 
filtered, Map<String, E> candidates, Map<InetAddressAndPort, Index.Status> 
indexStatusMap)
+            {
+                if (blockFor <= initial && blockFor > filtered)
+                {
+                    IndexStatusManager.instance.readFailureException(filtered, 
removedSupplier(planner.fullEndpoints, candidates), indexStatusMap, planner.cl, 
blockFor);
+                }
+            }
+
+            @Override
+            boolean canSelectReplicasFromDc(String dc, int liveNodes)
+            {
+                return true;
+            }
+
+            @Override
+            int minBlockFor()
+            {
+                return blockFor;
+            }
+
+            @Override
+            int eligibleCandidates(Map<String, E> candidates)
+            {
+                int eligible = 0;
+                for (Map.Entry<String, E> entry : candidates.entrySet())
+                {
+                    eligible += entry.getValue().size();
+                }
+                return eligible;
+            }
+
+            @Override
+            boolean haveSufficientContacts(int contactedReplicas, int 
contactedDcs)
+            {
+                return contactedReplicas >= targetContacts;
+            }
+
+            @Override
+            boolean haveSufficientContactsForDc(String dc, int 
contactedReplicas, int contactedDcs, int dcContacts)
+            {
+                return haveSufficientContacts(contactedReplicas, contactedDcs);
+            }
+
+            @Override
+            boolean haveSufficientLiveNodes(int contactedReplicas, int 
contactedDcs)
+            {
+                return contactedReplicas >= blockFor;
+            }
+        }
+
+        static class DcQuorum<E extends Endpoints<E>, L extends 
ReplicaLayout<E>, P extends ReplicaPlan.ForRead<E, P>> extends 
ReadReplicaPlanBuilder<E, L, P>
+        {
+            private final int blockForDcs;
+            private final Object2IntHashMap<String> dcQuorums = new 
Object2IntHashMap<>(-1);
+            private boolean approvedSpeculativeReplica = false;
+
+            public DcQuorum(CoordinationPlanner.ForRead<E, L, P> planner)
+            {
+                super(planner);
+                this.blockForDcs = planner.cl != ConsistencyLevel.EACH_QUORUM 
? (planner.dcs.size() / 2) + 1 : planner.dcs.size();
+                for (String dc : planner.dcs)
+                    dcQuorums.put(dc, planner.strategy.calculateQuorum(dc));
+            }
+            private int dcQuorum(String dc)
+            {
+                int v = dcQuorums.get(dc);
+                Preconditions.checkArgument(v >= 0, "dcQuorum must be >= 0");
+                return v;
+            }
+
+            @Override
+            public void maybeThrowIndexReadException(int initial, int 
filtered, Map<String, E> candidates, Map<InetAddressAndPort, Index.Status> 
indexStatusMap)
+            {
+                int required = 0;
+                int available = 0;
+                int liveDcs = 0;
+                for (Map.Entry<String, E> entry : candidates.entrySet())
+                {
+                    int dcQuorum = dcQuorum(entry.getKey());
+                    int dcLive = entry.getValue().size();
+                    if (dcLive >= dcQuorum)
+                        liveDcs++;
+
+                    required += dcQuorum;
+                    available += entry.getValue().size();
+                }
+
+                if (liveDcs <= blockForDcs)
+                {
+                    
IndexStatusManager.instance.readFailureException(available, 
removedSupplier(planner.fullEndpoints, candidates), indexStatusMap, planner.cl, 
required);
+                }
+            }
+
+            @Override
+            boolean canSelectReplicasFromDc(String dc, int liveNodes)
+            {
+                int dcQuorum = dcQuorum(dc);
+                return liveNodes >= dcQuorum;
+            }
+
+            @Override
+            int minBlockFor()
+            {
+                int[] quorums = new int[planner.dcs.size()];
+                for (int i = 0; i < quorums.length; i++)
+                    quorums[i] = dcQuorum(planner.dcs.get(i));
+
+                Arrays.sort(quorums);
+
+                int blockFor = 0;
+                for (int i=0; i<blockForDcs; i++)
+                    blockFor += quorums[i];
+                return blockFor;
+            }
+
+            @Override
+            int eligibleCandidates(Map<String, E> candidates)
+            {
+                int eligible = 0;
+                for (Map.Entry<String, E> entry : candidates.entrySet())
+                {
+                    int dcLive = entry.getValue().size();
+                    int dcQuorum = dcQuorum(entry.getKey());
+                    if (dcLive >= dcQuorum)
+                        eligible += dcLive;
+                }
+
+                return eligible;
+            }
+
+            @Override
+            boolean haveSufficientContacts(int contactedReplicas, int 
contactedDcs)
+            {
+                return contactedDcs >= blockForDcs;
+            }
+
+            @Override
+            boolean haveSufficientContactsForDc(String dc, int 
contactedReplicas, int contactedDcs, int dcContacts)
+            {
+                int dcQuorum = dcQuorum(dc);
+
+                if (dcContacts < dcQuorum)
+                    return false;
+
+                if (dcContacts == dcQuorum && planner.alwaysSpeculate && 
!approvedSpeculativeReplica)
+                {
+                    approvedSpeculativeReplica = true;
+                    return false;
+                }
+
+                return true;
+            }
+
+            @Override
+            boolean haveSufficientLiveNodes(int contactedReplicas, int 
contactedDcs)
+            {
+                return contactedDcs >= blockForDcs;
+            }
+        }
+
+    }
+
+    @Override
+    protected CoordinationPlan.ForWrite planForWriteInternal(ClusterMetadata 
metadata,
+                                                             Keyspace keyspace,
+                                                             ConsistencyLevel 
consistencyLevel,
+                                                             
Function<ClusterMetadata, ReplicaLayout.ForTokenWrite> liveAndDown,
+                                                             
ReplicaPlans.Selector selector)
+    {
+        // writes don't need any special handling during primary failover. We 
just write to the current
+        // primary DC, regardless of any other failover status
+        ReplicaLayout.ForTokenWrite fullLayout = liveAndDown.apply(metadata);
+
+        CoordinationPlanner.ForWrite planner = new 
CoordinationPlanner.ForWrite(metadata, keyspace, consistencyLevel, this, 
primaryDC, fullLayout, selector);
+
+        return new CoordinationPlan.ForWrite(planner.createReplicaPlan(), 
planner.createResponseTracker());
+    }
+
+    /**
+     * Holds the information needed to send satellite commit mutations 
alongside a paxos commit.
+     * The tracker has the primary DC pre-completed, so only 
satellite/secondary DC quorums
+     * need actual responses.
+     */
+    static class SatelliteCommitPlan
+    {
+        final EndpointsForToken liveEndpoints;
+        final EndpointsForToken downEndpoints;
+        final ResponseTracker tracker;
+
+        SatelliteCommitPlan(EndpointsForToken liveEndpoints, EndpointsForToken 
downEndpoints, ResponseTracker tracker)
+        {
+            this.liveEndpoints = liveEndpoints;
+            this.downEndpoints = downEndpoints;
+            this.tracker = tracker;
+        }
+    }
+
+    SatelliteCommitPlan createSatelliteCommitPlan(ClusterMetadata metadata, 
Keyspace keyspace, Token token)
+    {
+        ReplicaLayout.ForTokenWrite fullLayout = 
ReplicaLayout.forTokenWriteLiveAndDown(metadata, keyspace, token);
+        CoordinationPlanner.ForWrite planner = new 
CoordinationPlanner.ForWrite(metadata, keyspace, ConsistencyLevel.QUORUM,
+                                                                             
this, primaryDC, fullLayout, ReplicaPlans.writeAll);
+        ResponseTracker tracker = planner.createResponseTracker();
+
+        // paxos is handling the consensus/commit acks in the primary DC, we 
just need to worry about the witness DCs
+        EndpointsForToken primaryEndpoints = 
planner.fullEndpoints.get(primaryDC);
+        if (primaryEndpoints != null)
+        {
+            for (int i = 0; i < primaryEndpoints.size(); i++)
+                tracker.onResponse(primaryEndpoints.endpoint(i));
+        }
+
+        // Collect satellite/secondary DC endpoints from the filtered layout.
+        // planner.fullLayout is already filtered to only include DCs in the 
coordination plan
+        // (primary DC + its satellite + other full DCs), excluding satellites 
of non-primary DCs.
+        Predicate<Replica> notPrimary = rp -> 
!metadata.locator.location(rp.endpoint()).datacenter.equals(primaryDC);
+        EndpointsForToken allSatellite = 
planner.fullLayout.all().filter(notPrimary);
+        EndpointsForToken liveSatellite = 
allSatellite.filter(FailureDetector.isReplicaAlive);
+        EndpointsForToken downSatellite = allSatellite.filter(rp -> 
!FailureDetector.isReplicaAlive.test(rp));
+
+        return new SatelliteCommitPlan(liveSatellite, downSatellite, tracker);
+    }
+
+    /**
+     * Sends plain mutations to satellite/secondary DCs alongside a paxos 
commit.
+     *
+     * Paxos consensus handles the primary DC writes. This method sends the 
committed mutation to satellite DCs using
+     * the standard SRS write coordination. The primary DC group is 
pre-completed in the tracker so only satellite DC
+     * quorums need actual responses.
+     *
+     * @return a future that completes when satellite DC quorum requirements 
are met (or failed)
+     */
+    @Override
+    public Future<Void> sendPaxosCommitMutations(Agreed commit, boolean 
isUrgent)
+    {
+        ClusterMetadata metadata = ClusterMetadata.current();
+        Keyspace keyspace = Keyspace.open(commit.metadata().keyspace);
+        Token token = commit.partitionKey().getToken();
+
+        // Reject satellite commit writes during TRANSITION_ACK. Paxos 
operations should not be
+        // in progress during this state, but check defensively to avoid 
propagating stale commits.
+        SatelliteFailoverState.FailoverInfo failoverInfo = 
getFailoverInfo(token, metadata);
+        if (failoverInfo.getState() == 
SatelliteFailoverState.State.TRANSITION_ACK)
+            throw new UnavailableException("Paxos commit rejected during 
TRANSITION_ACK failover state",
+                                           ConsistencyLevel.SERIAL, 1, 0);
+
+        SatelliteCommitPlan plan = createSatelliteCommitPlan(metadata, 
keyspace, token);
+        ResponseTracker tracker = plan.tracker;
+        MutationId mutationId = commit.mutation.id();
+        Preconditions.checkState(!mutationId.isNone());
+
+        // Register satellite replicas with mutation tracking service
+        if (!plan.liveEndpoints.isEmpty() || !plan.downEndpoints.isEmpty())
+        {
+            IntHashSet satelliteHostIds = new IntHashSet();
+            for (int i = 0; i < plan.liveEndpoints.size(); i++)
+                
satelliteHostIds.add(metadata.directory.peerId(plan.liveEndpoints.endpoint(i)).id());
+            for (int i = 0; i < plan.downEndpoints.size(); i++)
+                
satelliteHostIds.add(metadata.directory.peerId(plan.downEndpoints.endpoint(i)).id());
+            if (!satelliteHostIds.isEmpty())
+                
MutationTrackingService.instance().sentWriteRequest(commit.makeMutation(), 
satelliteHostIds);
+        }
+
+        // Create promise that resolves when satellite tracker completes
+        AsyncPromise<Void> promise = new AsyncPromise<>();
+
+        // Send plain mutations to live satellite endpoints
+        Message<Mutation> satelliteMessage = 
Message.out(Verb.PAXOS2_COMMIT_REMOTE_REQ, commit.makeMutation(), isUrgent);
+        for (int i = 0, mi = plan.liveEndpoints.size(); i < mi; ++i)
+        {
+            InetAddressAndPort endpoint = plan.liveEndpoints.endpoint(i);
+            logger.trace("Sending satellite commit mutation for {} to {}", 
commit.partitionKey(), endpoint);
+            MessagingService.instance().sendWithCallback(satelliteMessage, 
endpoint, new RequestCallback<NoPayload>()

Review Comment:
   Is there testing added along with the fix?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]


Reply via email to