mcgilman commented on code in PR #11677:
URL: https://github.com/apache/nifi/pull/11677#discussion_r4212096407


##########
nifi-frontend/src/main/frontend/apps/nifi/src/app/pages/connectors/ui/connector-table/connector-table.component.spec.ts:
##########
@@ -17,7 +17,13 @@
 
 import { TestBed } from '@angular/core/testing';
 import { ConnectorTable } from './connector-table.component';
-import { ConnectorAction, ConnectorActionName, ConnectorEntity, 
ConnectorStatus, NiFiCommon } from '@nifi/shared';
+import {
+    ConnectorAction,
+    ConnectorActionName,
+    ConnectorEntity,
+    ConnectorStatus,
+    NiFiCommon
+} from '@nifi/shared';

Review Comment:
   **Lint — restore the single-line import**
   
   This wraps the `@nifi/shared` import that is one line on main. Prettier on 
current main rejects the wrap, and `npx nx run nifi:lint` fails on it. Frontend 
Maven runs that lint during `generate-resources`, so the build fails before 
tests. Put the import back on one line.



##########
nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/manager/ConnectorEntityMerger.java:
##########
@@ -101,6 +101,14 @@ private static void mergeDtos(final ConnectorDTO 
clientDto, final Map<NodeIdenti
             return;
         }
 
+        for (final ConnectorDTO nodeConnector : dtoMap.values()) {
+            if (nodeConnector != null) {
+                if (clientDto.getMultipleVersionsAvailable() == null || 
!Boolean.TRUE.equals(nodeConnector.getMultipleVersionsAvailable())) {
+                    clientDto.setMultipleVersionsAvailable(Boolean.FALSE);

Review Comment:
   **Medium — merge `CHANGE_VERSION` the same way as 
`multipleVersionsAvailable`**
   
   This loop correctly forces `multipleVersionsAvailable` to false when any 
node reports false or null. `availableActions` is still left as the selected 
node's list. `CHANGE_VERSION` is node-local: one node can be stopped while 
another is `UPDATED` with active threads. State merge can then show the 
higher-priority state while the menu still offers Change Version.
   
   Replicated verification calls `verifyCanReload()` on every node before 
commit, so this does not apply a partial bundle change. The menu can still 
offer an action the cluster will reject.
   
   If any node has `CHANGE_VERSION` missing, `allowed` null, or `allowed` 
false, set the client action to not allowed and keep that node's 
`reasonNotAllowed`. Extend `ConnectorEntityMergerTest` with false, null, and 
all-true cases.



##########
nifi-toolkit/nifi-toolkit-cli/src/main/java/org/apache/nifi/toolkit/cli/impl/command/nifi/connectors/ChangeVersionConnector.java:
##########
@@ -0,0 +1,199 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.nifi.toolkit.cli.impl.command.nifi.connectors;
+
+import org.apache.commons.cli.MissingOptionException;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.nifi.toolkit.cli.api.CommandException;
+import org.apache.nifi.toolkit.cli.api.Context;
+import org.apache.nifi.toolkit.cli.impl.command.CommandOption;
+import org.apache.nifi.toolkit.cli.impl.command.nifi.AbstractNiFiCommand;
+import org.apache.nifi.toolkit.cli.impl.result.nifi.ConnectorsResult;
+import org.apache.nifi.toolkit.client.ConnectorClient;
+import org.apache.nifi.toolkit.client.FlowClient;
+import org.apache.nifi.toolkit.client.NiFiClient;
+import org.apache.nifi.toolkit.client.NiFiClientException;
+import org.apache.nifi.web.api.dto.BundleDTO;
+import org.apache.nifi.web.api.dto.ConnectorDTO;
+import org.apache.nifi.web.api.dto.PermissionsDTO;
+import org.apache.nifi.web.api.entity.ConnectorEntity;
+import org.apache.nifi.web.api.entity.ConnectorsEntity;
+
+import java.io.IOException;
+import java.util.HashSet;
+import java.util.Properties;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Command to update the NAR version of Connector instances.
+ */
+public class ChangeVersionConnector extends 
AbstractNiFiCommand<ConnectorsResult> {
+
+    public ChangeVersionConnector() {
+        super("change-version-connector", ConnectorsResult.class);
+    }
+
+    @Override
+    public String getDescription() {
+        return "Changes the version of Connector instances of the specified 
type. If the current version is specified, only instances "
+                + "with that version are updated. Running Connectors are 
stopped before the version is changed and restarted afterward.";
+    }
+
+    @Override
+    protected void doInitialize(final Context context) {
+        addOption(CommandOption.EXT_BUNDLE_GROUP.createOption());
+        addOption(CommandOption.EXT_BUNDLE_ARTIFACT.createOption());
+        addOption(CommandOption.EXT_BUNDLE_VERSION.createOption());
+        addOption(CommandOption.EXT_QUALIFIED_NAME.createOption());
+        addOption(CommandOption.EXT_BUNDLE_CURRENT_VERSION.createOption());
+    }
+
+    @Override
+    public ConnectorsResult doExecute(final NiFiClient client, final 
Properties properties)
+            throws NiFiClientException, IOException, MissingOptionException, 
CommandException {
+
+        final String bundleGroup = getRequiredArg(properties, 
CommandOption.EXT_BUNDLE_GROUP);
+        final String bundleArtifact = getRequiredArg(properties, 
CommandOption.EXT_BUNDLE_ARTIFACT);
+        final String bundleVersion = getRequiredArg(properties, 
CommandOption.EXT_BUNDLE_VERSION);
+        final String qualifiedName = getRequiredArg(properties, 
CommandOption.EXT_QUALIFIED_NAME);
+        final String sourceVersion = getArg(properties, 
CommandOption.EXT_BUNDLE_CURRENT_VERSION);
+
+        final FlowClient flowClient = client.getFlowClient();
+        final ConnectorClient connectorClient = client.getConnectorClient();
+        final ConnectorsEntity connectorsEntity = flowClient.getConnectors();
+        final Set<ConnectorEntity> updatedComponents = new HashSet<>();
+
+        if (connectorsEntity.getConnectors() != null) {
+            for (final ConnectorEntity connector : 
connectorsEntity.getConnectors()) {
+                if (isUnreadable(connector)) {
+                    throw new CommandException("Cannot change version because 
Connector " + connector.getId() + " is unreadable");
+                }
+            }
+
+            for (final ConnectorEntity connector : 
connectorsEntity.getConnectors()) {
+                final BundleDTO bundle = connector.getComponent().getBundle();
+                if (!bundle.getGroup().equals(bundleGroup)
+                        || !bundle.getArtifact().equals(bundleArtifact)
+                        || 
!connector.getComponent().getType().equals(qualifiedName)
+                        || (!StringUtils.isBlank(sourceVersion) && 
!bundle.getVersion().equals(sourceVersion))) {
+                    continue;
+                }
+
+                if (bundleVersion.equals(bundle.getVersion())) {
+                    continue;
+                }
+
+                final String currentState = 
connector.getComponent().getState();
+                if ("TROUBLESHOOTING".equals(currentState)) {
+                    throw new CommandException("Cannot change version of 
Connector " + connector.getId() + " while it is in Troubleshooting");
+                }
+
+                final boolean isRunning = "RUNNING".equals(currentState) || 
"STARTING".equals(currentState);
+                if (isRunning) {
+                    connectorClient.stopConnector(connector);
+                }
+
+                try {
+                    final boolean shouldWaitForReloadable = isRunning || 
!isReloadableState(currentState);
+                    final ConnectorEntity reloadableConnector;
+                    if (shouldWaitForReloadable) {
+                        reloadableConnector = 
waitForConnectorReloadable(connectorClient, connector.getId());
+                    } else {
+                        reloadableConnector = connector;
+                    }
+
+                    final BundleDTO updatedBundle = new BundleDTO(bundleGroup, 
bundleArtifact, bundleVersion);
+                    final ConnectorDTO connectorDto = new ConnectorDTO();
+                    connectorDto.setId(reloadableConnector.getId());
+                    connectorDto.setBundle(updatedBundle);
+
+                    final ConnectorEntity updatedEntity = new 
ConnectorEntity();
+                    
updatedEntity.setRevision(reloadableConnector.getRevision());
+                    updatedEntity.setComponent(connectorDto);
+                    updatedEntity.setId(reloadableConnector.getId());
+
+                    connectorClient.updateConnector(updatedEntity);
+                } catch (final NiFiClientException | IOException | 
CommandException changeVersionFailure) {
+                    if (isRunning) {
+                        try {
+                            final ConnectorEntity currentConnector = 
connectorClient.getConnector(connector.getId());
+                            connectorClient.startConnector(currentConnector);
+                        } catch (final Exception restartFailure) {
+                            changeVersionFailure.addSuppressed(restartFailure);
+                        }
+                    }
+
+                    throw changeVersionFailure;
+                }
+
+                if (isRunning) {
+                    final ConnectorEntity connectorToStart = 
connectorClient.getConnector(connector.getId());
+                    connectorClient.startConnector(connectorToStart);
+                }
+
+                final ConnectorEntity updatedConnector = 
connectorClient.getConnector(connector.getId());
+                updatedComponents.add(updatedConnector);
+            }
+        }
+
+        final ConnectorsEntity resultEntity = new ConnectorsEntity();
+        resultEntity.setConnectors(updatedComponents);
+        return new ConnectorsResult(getResultType(properties), resultEntity);
+    }
+
+    private boolean isUnreadable(final ConnectorEntity connector) {
+        if (connector.getComponent() == null) {
+            return true;
+        }
+
+        final PermissionsDTO permissions = connector.getPermissions();
+        if (permissions == null) {
+            return false;
+        }
+
+        return Boolean.FALSE.equals(permissions.getCanRead());
+    }
+
+    private ConnectorEntity waitForConnectorReloadable(final ConnectorClient 
connectorClient, final String connectorId)
+            throws NiFiClientException, IOException, CommandException {
+        final long deadline = System.currentTimeMillis() + 
TimeUnit.SECONDS.toMillis(60);
+        ConnectorEntity connector = connectorClient.getConnector(connectorId);
+        while (System.currentTimeMillis() < deadline) {
+            final String state = connector.getComponent().getState();
+            if (isReloadableState(state)) {
+                return connector;
+            }
+
+            try {
+                Thread.sleep(200L);
+            } catch (final InterruptedException e) {
+                Thread.currentThread().interrupt();
+                throw new CommandException("Interrupted while waiting for 
Connector " + connectorId + " to stop");
+            }
+
+            connector = connectorClient.getConnector(connectorId);
+        }
+
+        throw new CommandException("Timed out waiting for Connector " + 
connectorId + " to stop; current state is "
+                + connector.getComponent().getState());
+    }
+
+    private boolean isReloadableState(final String state) {

Review Comment:
   **Medium — wait on `CHANGE_VERSION`, not a weaker state list**
   
   `isReloadableState` treats `UPDATED` and `UPDATE_FAILED` as reloadable with 
no active-thread check. `verifyCanReload()` and the `CHANGE_VERSION` action do 
not. The CLI will PUT, get a 409, and, because those states are not `RUNNING` 
or `STARTING`, it will not have stopped the connector first.
   
   After stop-and-wait, require `CHANGE_VERSION` to be allowed on the refreshed 
entity instead of duplicating this state list. Add a CLI test for `UPDATED` 
with active threads.



-- 
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]

Reply via email to