blob: be928f3510825f03bf21e200221a742965f887c9 [file] [log] [blame]
// 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 com.cloud.network.security;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
import org.apache.logging.log4j.Logger;
import org.apache.logging.log4j.LogManager;
import com.cloud.agent.AgentManager;
import com.cloud.agent.Listener;
import com.cloud.agent.api.Answer;
public class SecurityGroupWorkTracker {
protected Logger logger = LogManager.getLogger(getClass());
protected AtomicLong _discardCount = new AtomicLong(0);
AgentManager _agentMgr;
Listener _answerListener;
int _bufferLength;
Map<Long, Integer> _unackedMessages = new ConcurrentHashMap<Long, Integer>();
public SecurityGroupWorkTracker(AgentManager agentMgr, Listener answerListener, int bufferLength) {
super();
assert (bufferLength >= 1) : "SecurityGroupWorkTracker: Cannot have a zero length buffer";
this._agentMgr = agentMgr;
this._answerListener = answerListener;
this._bufferLength = bufferLength;
}
public boolean canSend(long agentId) {
int currLength = 0;
synchronized (this) {
Integer outstanding = _unackedMessages.get(agentId);
if (outstanding == null) {
outstanding = 0;
_unackedMessages.put(agentId, outstanding);
}
currLength = outstanding.intValue();
if (currLength + 1 > _bufferLength) {
long discarded = _discardCount.incrementAndGet();
//drop it on the floor
logger.debug("SecurityGroupManager: dropping a message because there are more than " + currLength + " outstanding messages, total dropped=" + discarded);
return false;
}
_unackedMessages.put(agentId, ++currLength);
}
return true;
}
public void handleException(long agentId) {
synchronized (this) {
Integer outstanding = _unackedMessages.get(agentId);
if (outstanding != null && outstanding != 0) {
_unackedMessages.put(agentId, --outstanding);
}
}
}
public void processAnswers(long agentId, long seq, Answer[] answers) {
synchronized (this) {
Integer outstanding = _unackedMessages.get(agentId);
if (outstanding != null && outstanding != 0) {
_unackedMessages.put(agentId, --outstanding);
}
}
}
public void processTimeout(long agentId, long seq) {
synchronized (this) {
Integer outstanding = _unackedMessages.get(agentId);
if (outstanding != null && outstanding != 0) {
_unackedMessages.put(agentId, --outstanding);
}
}
}
public void processDisconnect(long agentId) {
synchronized (this) {
_unackedMessages.put(agentId, 0);
}
}
public void processConnect(long agentId) {
synchronized (this) {
_unackedMessages.put(agentId, 0);
}
}
public long getDiscardCount() {
return _discardCount.get();
}
public int getUnackedCount(long agentId) {
Integer outstanding = _unackedMessages.get(agentId);
if (outstanding == null) {
return 0;
}
return outstanding.intValue();
}
}