Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java
 Thu Mar 24 07:52:35 2011
@@ -22,9 +22,11 @@ import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.TreeMap;
+
 import org.apache.commons.logging.Log;
 import org.apache.commons.logging.LogFactory;
 import org.apache.hadoop.classification.InterfaceAudience.LimitedPrivate;
@@ -32,24 +34,24 @@ import org.apache.hadoop.classification.
 import org.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.net.Node;
 import org.apache.hadoop.security.UserGroupInformation;
+import org.apache.hadoop.yarn.ApplicationID;
+import org.apache.hadoop.yarn.Container;
+import org.apache.hadoop.yarn.ContainerToken;
+import org.apache.hadoop.yarn.NodeID;
+import org.apache.hadoop.yarn.Priority;
+import org.apache.hadoop.yarn.Resource;
+import org.apache.hadoop.yarn.ResourceRequest;
 import org.apache.hadoop.yarn.security.ContainerTokenIdentifier;
+import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ASMEvent;
+import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType;
 import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo;
 import org.apache.hadoop.yarn.server.resourcemanager.scheduler.Application;
-import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker;
-import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTrackerImpl;
 import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeManager;
+import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse;
 import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeType;
 import org.apache.hadoop.yarn.server.resourcemanager.scheduler.Queue;
 import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler;
-import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse;
 import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager;
-import org.apache.hadoop.yarn.ApplicationID;
-import org.apache.hadoop.yarn.Container;
-import org.apache.hadoop.yarn.ContainerToken;
-import org.apache.hadoop.yarn.NodeID;
-import org.apache.hadoop.yarn.Priority;
-import org.apache.hadoop.yarn.Resource;
-import org.apache.hadoop.yarn.ResourceRequest;
 
 @LimitedPrivate("yarn")
 @Evolving
@@ -59,7 +61,6 @@ public class FifoScheduler implements Re
   
   Configuration conf;
   private ContainerTokenSecretManager containerTokenSecretManager;
-  private final ClusterTracker clusterTracker;
   
   // TODO: The memory-block size should be site-configurable?
   public static final int MINIMUM_MEMORY = 1024;
@@ -81,20 +82,14 @@ public class FifoScheduler implements Re
     }
   };
   
-  public FifoScheduler() {
-    this.clusterTracker = createClusterTracker();
-  }
+  public FifoScheduler() {}
   
   public FifoScheduler(Configuration conf,
       ContainerTokenSecretManager containerTokenSecretManager) 
   {
-    this();
     reinitialize(conf, containerTokenSecretManager);
   }
   
-  protected ClusterTracker createClusterTracker() {
-    return new ClusterTrackerImpl(); 
-  }
   
   @Override
   public void reinitialize(Configuration conf,
@@ -141,7 +136,7 @@ public class FifoScheduler implements Re
   private void releaseContainers(Application application, List<Container> 
release) {
     application.releaseContainers(release);
     for (Container container : release) {
-      clusterTracker.releaseContainer(application.getApplicationId(), 
container);
+      releaseContainer(application.getApplicationId(), container);
     }
   }
   
@@ -161,7 +156,6 @@ public class FifoScheduler implements Re
     return applications.get(applicationId);
   }
 
-  @Override
   public synchronized void addApplication(ApplicationID applicationId, 
       String user, String unusedQueue, Priority unusedPriority) 
   throws IOException {
@@ -171,7 +165,6 @@ public class FifoScheduler implements Re
         ", currently active: " + applications.size());
   }
 
-  @Override
   public synchronized void removeApplication(ApplicationID applicationId)
   throws IOException {
     Application application = getApplication(applicationId);
@@ -184,7 +177,7 @@ public class FifoScheduler implements Re
     releaseContainers(application, application.getCurrentContainers());
     
     // Let the cluster know that the applications are done
-    clusterTracker.finishedApplication(applicationId, 
+    finishedApplication(applicationId, 
         application.getAllNodesForApplication());
     
     // Remove the application
@@ -389,9 +382,8 @@ public class FifoScheduler implements Re
         containers.add(container);
       }
       application.allocate(type, node, priority, request, containers);
-      clusterTracker.addAllocatedContainers(node, 
application.getApplicationId(), containers);
+      addAllocatedContainers(node, application.getApplicationId(), containers);
     }
-    
     return assignedContainers;
   }
 
@@ -412,7 +404,7 @@ public class FifoScheduler implements Re
   public synchronized NodeResponse nodeUpdate(NodeInfo node, 
       Map<CharSequence,List<Container>> containers ) {
    
-    NodeResponse nodeResponse = clusterTracker.nodeUpdate(node, containers);
+    NodeResponse nodeResponse = nodeUpdateInternal(node, containers);
     applicationCompletedContainers(nodeResponse.getCompletedContainers());
     LOG.info("Node heartbeat " + node.getNodeID() + " resource = " + 
node.getAvailableResource());
     if (org.apache.hadoop.yarn.server.resourcemanager.resource.Resource.
@@ -426,15 +418,88 @@ public class FifoScheduler implements Re
     // preemption.
     return nodeResponse;
   }  
-  
+
   @Override
-  public NodeInfo addNode(NodeID nodeId,String hostName,
-      Node node, Resource capability) {
-    return clusterTracker.addNode(nodeId, hostName, node, capability);
+  public synchronized void handle(ASMEvent<ApplicationTrackerEventType> event) 
{
+    switch(event.getType()) {
+    case ADD:
+      try {
+        addApplication(event.getAppContext().getApplicationID(), 
event.getAppContext().getUser(),
+            event.getAppContext().getQueue(), 
event.getAppContext().getSubmissionContext().priority);
+      } catch(IOException ie) {
+        LOG.error("Unable to add application " + 
event.getAppContext().getApplicationID(), ie);
+        /** this is fatal we are not able to add applications for scheduling 
**/
+        //TODO handle it later.
+      }
+      break;
+    case REMOVE:
+      try {
+        
+        removeApplication(event.getAppContext().getApplicationID());
+      } catch(IOException ie) {
+        LOG.error("Unable to remove application " + 
event.getAppContext().getApplicationID(), ie);
+      }
+      break;  
+    }
+  }
+  
+  private Map<String, NodeManager> nodes = new HashMap<String, NodeManager>();
+  private Resource clusterResource = new Resource();
+ 
+  public synchronized Resource getClusterResource() {
+    return clusterResource;
   }
 
   @Override
-  public void removeNode(NodeInfo node) {
-    clusterTracker.removeNode(node);
+  public synchronized void removeNode(NodeInfo nodeInfo) {
+    
org.apache.hadoop.yarn.server.resourcemanager.resource.Resource.subtractResource(
+        clusterResource, nodeInfo.getTotalCapability());
+    nodes.remove(nodeInfo.getHostName());
+  }
+  
+  public synchronized boolean isTracked(NodeInfo nodeInfo) {
+    NodeManager node = nodes.get(nodeInfo.getHostName());
+    return (node == null? false: true);
+  }
+ 
+  @Override
+  public synchronized NodeInfo addNode(NodeID nodeId, 
+      String hostName, Node node, Resource capability) {
+    NodeManager nodeManager = new NodeManager(nodeId, hostName, node, 
capability);
+    nodes.put(nodeManager.getHostName(), nodeManager);
+    
org.apache.hadoop.yarn.server.resourcemanager.resource.Resource.addResource(
+        clusterResource, nodeManager.getTotalCapability());
+    return nodeManager;
+  }
+
+  public synchronized boolean releaseContainer(ApplicationID applicationId, 
+      Container container) {
+    // Reap containers
+    LOG.info("Application " + applicationId + " released container " + 
container);
+    NodeManager nodeManager = nodes.get(container.hostName.toString());
+    return nodeManager.releaseContainer(container);
+  }
+  
+  private synchronized NodeResponse nodeUpdateInternal(NodeInfo nodeInfo, 
+      Map<CharSequence,List<Container>> containers) {
+    NodeManager node = nodes.get(nodeInfo.getHostName());
+    LOG.debug("nodeUpdate: node=" + nodeInfo.getHostName() + 
+        " available=" + nodeInfo.getAvailableResource().memory);
+    return node.statusUpdate(containers);
+    
+  }
+
+  public synchronized void addAllocatedContainers(NodeInfo nodeInfo, 
+      ApplicationID applicationId, List<Container> containers) {
+    NodeManager node = nodes.get(nodeInfo.getHostName());
+    node.allocateContainer(applicationId, containers);
+  }
+
+  public synchronized void finishedApplication(ApplicationID applicationId,
+      List<NodeInfo> nodesToNotify) {
+    for (NodeInfo node: nodesToNotify) {
+      NodeManager nodeManager = nodes.get(node.getHostName());
+      nodeManager.notifyFinishedApplication(applicationId);
+    }
   }
 }

Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java
 Thu Mar 24 07:52:35 2011
@@ -1,26 +1,27 @@
 /**
-* 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.
-*/
+ * 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.hadoop.yarn.server.resourcemanager.applicationsmanager;
 
 import java.io.IOException;
 import java.util.Arrays;
 import java.util.List;
+import java.util.concurrent.atomic.AtomicInteger;
 
 import junit.framework.TestCase;
 
@@ -63,19 +64,9 @@ public class TestAMLaunchFailure extends
 
     @Override
     public List<Container> allocate(ApplicationID applicationId,
-    List<ResourceRequest> ask, List<Container> release) throws IOException {
+        List<ResourceRequest> ask, List<Container> release) throws IOException 
{
       return Arrays.asList(container);
     }
-
-    @Override
-    public void addApplication(ApplicationID applicationId, String user,
-    String queue, Priority priority) throws IOException {
-    }
-
-    @Override
-    public void removeApplication(ApplicationID applicationId)
-    throws IOException {
-    }
   }
 
   private class DummyApplicationTracker implements 
EventHandler<ASMEvent<ApplicationTrackerEventType>> {
@@ -90,9 +81,9 @@ public class TestAMLaunchFailure extends
   public class ExtApplicationsManagerImpl extends ApplicationsManagerImpl {
 
     private  class DummyApplicationMasterLauncher implements 
EventHandler<ASMEvent<AMLauncherEventType>> {
-      private Object notify = new Object();
+      private AtomicInteger notify = new AtomicInteger();
       private AppContext app;
-      
+
       public DummyApplicationMasterLauncher(ASMContext context) {
         context.getDispatcher().register(AMLauncherEventType.class, this);
         new TestThread().start();
@@ -101,8 +92,10 @@ public class TestAMLaunchFailure extends
       public void handle(ASMEvent<AMLauncherEventType> appEvent) {
         switch(appEvent.getType()) {
         case LAUNCH:
+          LOG.info("LAUNCH called ");
           app = appEvent.getAppContext();
           synchronized (notify) {
+            notify.addAndGet(1);
             notify.notify();
           }
           break;
@@ -113,27 +106,29 @@ public class TestAMLaunchFailure extends
         public void run() {
           synchronized(notify) {
             try {
-              notify.wait();
+              while (notify.get() == 0) {
+                notify.wait();
+              }
             } catch (InterruptedException e) {
               e.printStackTrace();
             }
             context.getDispatcher().getEventHandler().handle(
-            new ASMEvent<ApplicationEventType>(ApplicationEventType.LAUNCHED,
-            app));
+                new 
ASMEvent<ApplicationEventType>(ApplicationEventType.LAUNCHED,
+                    app));
           }
         }
       }
     }
 
     public ExtApplicationsManagerImpl(
-    ApplicationTokenSecretManager applicationTokenSecretManager,
-    YarnScheduler scheduler) {
+        ApplicationTokenSecretManager applicationTokenSecretManager,
+        YarnScheduler scheduler) {
       super(applicationTokenSecretManager, scheduler, context);
     }
 
     @Override
     protected EventHandler<ASMEvent<AMLauncherEventType>> 
createNewApplicationMasterLauncher(
-    ApplicationTokenSecretManager tokenSecretManager) {
+        ApplicationTokenSecretManager tokenSecretManager) {
       return new DummyApplicationMasterLauncher(context);
     }
   }
@@ -146,6 +141,7 @@ public class TestAMLaunchFailure extends
     Configuration conf = new Configuration();
     new DummyApplicationTracker();
     conf.setLong(YarnConfiguration.AM_EXPIRY_INTERVAL, 3000L);
+    conf.setInt(YarnConfiguration.AM_MAX_RETRIES, 1);
     asmImpl.init(conf);
     asmImpl.start();  
   }

Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java
 Thu Mar 24 07:52:35 2011
@@ -75,14 +75,6 @@ public class TestAMRMRPCResponseId exten
         List<ResourceRequest> ask, List<Container> release) throws IOException 
{
       return null;
     }
-    @Override
-    public void addApplication(ApplicationID applicationId, String user,
-        String queue, Priority priority) throws IOException {
-    }
-    @Override
-    public void removeApplication(ApplicationID applicationId)
-        throws IOException {
-    }
   }
   
   @Before

Added: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java?rev=1084866&view=auto
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java
 (added)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java
 Thu Mar 24 07:52:35 2011
@@ -0,0 +1,300 @@
+package org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import junit.framework.TestCase;
+
+import org.apache.commons.logging.Log;
+import org.apache.commons.logging.LogFactory;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.net.Node;
+import org.apache.hadoop.yarn.ApplicationID;
+import org.apache.hadoop.yarn.ApplicationMaster;
+import org.apache.hadoop.yarn.ApplicationState;
+import org.apache.hadoop.yarn.ApplicationStatus;
+import org.apache.hadoop.yarn.ApplicationSubmissionContext;
+import org.apache.hadoop.yarn.Container;
+import org.apache.hadoop.yarn.ContainerID;
+import org.apache.hadoop.yarn.ContainerToken;
+import org.apache.hadoop.yarn.NodeID;
+import org.apache.hadoop.yarn.Resource;
+import org.apache.hadoop.yarn.ResourceRequest;
+import org.apache.hadoop.yarn.conf.YarnConfiguration;
+import org.apache.hadoop.yarn.event.EventHandler;
+import org.apache.hadoop.yarn.security.ApplicationTokenSecretManager;
+import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager;
+import 
org.apache.hadoop.yarn.server.resourcemanager.ResourceManager.ASMContext;
+import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ASMEvent;
+import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.AMLauncherEventType;
+import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationEventType;
+import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType;
+import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo;
+import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse;
+import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler;
+import org.apache.hadoop.yarn.server.resourcemanager.scheduler.YarnScheduler;
+import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+
+/**
+ * Test to restart the AM on failure.
+ *
+ */
+public class TestAMRestart extends TestCase {
+  private static final Log LOG = LogFactory.getLog(TestAMRestart.class);
+  ApplicationsManagerImpl appImpl;
+  ASMContext asmContext = new ResourceManager.ASMContextImpl();
+  ApplicationTokenSecretManager appTokenSecretManager = 
+    new ApplicationTokenSecretManager();
+  DummyResourceScheduler scheduler;
+  int count = 0;
+  ApplicationID appID;
+  final int maxFailures = 3;
+  AtomicInteger launchNotify = new AtomicInteger();
+  AtomicInteger schedulerNotify = new AtomicInteger();
+  volatile boolean stop = false;
+  int schedulerAddApplication = 0;
+  int schedulerRemoveApplication = 0;
+  int launcherLaunchCalled = 0;
+  int launcherCleanupCalled = 0;
+  ApplicationMasterInfo masterInfo;
+  
+  private class ExtApplicationsManagerImpl extends ApplicationsManagerImpl {
+    public ExtApplicationsManagerImpl(
+        ApplicationTokenSecretManager applicationTokenSecretManager,
+        YarnScheduler scheduler, ASMContext asmContext) {
+      super(applicationTokenSecretManager, scheduler, asmContext);
+    }
+
+    @Override
+    public EventHandler<ASMEvent<AMLauncherEventType>> 
createNewApplicationMasterLauncher(
+        ApplicationTokenSecretManager tokenSecretManager) {
+      return new DummyAMLauncher();
+    }
+  }
+
+  private class DummyAMLauncher implements 
EventHandler<ASMEvent<AMLauncherEventType>> {
+
+    public DummyAMLauncher() {
+      asmContext.getDispatcher().register(AMLauncherEventType.class, this);
+      new Thread() {
+        public void run() {
+          while (!stop) {
+            LOG.info("DEBUG -- waiting for launch");
+            synchronized(launchNotify) {
+              while (launchNotify.get() == 0) {
+                try { 
+                  launchNotify.wait();
+                } catch (InterruptedException e) {
+                }
+              }
+              asmContext.getDispatcher().getEventHandler().handle(new 
+                  ASMEvent<ApplicationEventType>(ApplicationEventType.LAUNCHED,
+                      new TestAppContext(appID)));
+              launchNotify.addAndGet(-1);
+            }
+          }
+        }
+      }.start();
+    }
+
+    @Override
+    public void handle(ASMEvent<AMLauncherEventType> event) {
+      switch (event.getType()) {
+      case CLEANUP:
+        launcherCleanupCalled++;
+        break;
+      case LAUNCH:
+        LOG.info("DEBUG -- launching");
+        launcherLaunchCalled++;
+        synchronized (launchNotify) {
+          launchNotify.addAndGet(1);
+          launchNotify.notify();
+        }
+        break;
+      default:
+        break;
+      }
+    }
+  }
+
+  private class DummyResourceScheduler implements ResourceScheduler {
+    @Override
+    public NodeInfo addNode(NodeID nodeId, String hostName, Node node,
+        Resource capability) {
+      return null;
+    }
+    @Override
+    public void removeNode(NodeInfo node) {
+    }
+    @Override
+    public NodeResponse nodeUpdate(NodeInfo nodeInfo,
+        Map<CharSequence, List<Container>> containers) {
+      return null;
+    }
+
+    @Override
+    public List<Container> allocate(ApplicationID applicationId,
+        List<ResourceRequest> ask, List<Container> release) throws IOException 
{
+      Container container = new Container();
+      container.containerToken = new ContainerToken();
+      container.hostName = "localhost";
+      container.id = new ContainerID();
+      container.id.appID = appID;
+      container.id.id = count;
+      count++;
+      return Arrays.asList(container);
+    }
+
+    @Override
+    public void handle(ASMEvent<ApplicationTrackerEventType> event) {
+      switch (event.getType()) {
+      case ADD:
+        schedulerAddApplication++;
+        break;
+      case REMOVE:
+        schedulerRemoveApplication++;
+        LOG.info("REMOVING app : " + schedulerRemoveApplication);
+        if (schedulerRemoveApplication == maxFailures) {
+          synchronized (schedulerNotify) {
+            schedulerNotify.addAndGet(1);
+            schedulerNotify.notify();
+          }
+        }
+        break;
+      default:
+        break;
+      }
+    }
+
+    @Override
+    public void reinitialize(Configuration conf,
+        ContainerTokenSecretManager secretManager) {
+    }
+  }
+
+  @Before
+  public void setUp() {
+    appID = new ApplicationID();
+    appID.clusterTimeStamp = System.currentTimeMillis();
+    appID.id = 1;
+    scheduler = new DummyResourceScheduler();
+    asmContext.getDispatcher().register(ApplicationTrackerEventType.class, 
scheduler);
+    appImpl = new ExtApplicationsManagerImpl(appTokenSecretManager, scheduler, 
asmContext);
+    Configuration conf = new Configuration();
+    conf.setLong(YarnConfiguration.AM_EXPIRY_INTERVAL, 1000L);
+    conf.setInt(YarnConfiguration.AM_MAX_RETRIES, maxFailures);
+    appImpl.init(conf);
+    appImpl.start();
+  }
+
+  @After
+  public void tearDown() {
+  }
+
+  private void waitForFailed(ApplicationMasterInfo masterInfo, 
ApplicationState 
+      finalState) throws Exception {
+    int count = 0;
+    while(masterInfo.getState() != finalState && count < 10) {
+      Thread.sleep(500);
+      count++;
+    }
+    assertTrue(masterInfo.getState() == finalState);
+  }
+  
+  private class TestAppContext implements AppContext {
+    private ApplicationID appID;
+   
+    public TestAppContext(ApplicationID appID) {
+      this.appID = appID;
+    }
+    @Override
+    public ApplicationSubmissionContext getSubmissionContext() {
+      return null;
+    }
+
+    @Override
+    public Resource getResource() {
+      return null;
+    }
+
+    @Override
+    public ApplicationID getApplicationID() {
+      return appID;
+    }
+
+    @Override
+    public ApplicationStatus getStatus() {
+      return null;
+    }
+
+    @Override
+    public ApplicationMaster getMaster() {
+      return null;
+    }
+
+    @Override
+    public Container getMasterContainer() {
+      return null;
+    }
+
+    @Override
+    public String getUser() {
+      return null;
+    }
+
+    @Override
+    public long getLastSeen() {
+      return 0;
+    }
+
+    @Override
+    public String getName() {
+      return null;
+    }
+
+    @Override
+    public String getQueue() {
+      return null;
+    }
+
+    @Override
+    public int getFailedCount() {
+      return 0;
+    }
+    
+  }
+
+  @Test
+  public void testAMRestart() throws Exception {
+    ApplicationSubmissionContext subContext = new 
ApplicationSubmissionContext();
+    subContext.applicationId = appID;
+    subContext.applicationName = "dummyApp";
+    subContext.command = new ArrayList<CharSequence>();
+    subContext.environment = new HashMap<CharSequence, CharSequence>();
+    subContext.fsTokens = new ArrayList<CharSequence>();
+    subContext.fsTokens_todo = ByteBuffer.wrap(new byte[0]);
+    appImpl.submitApplication(subContext);
+    masterInfo = appImpl.getApplicationMasterInfo(appID);
+    synchronized (schedulerNotify) {
+      while(schedulerNotify.get() == 0) {
+        schedulerNotify.wait();
+      }
+    }
+    assertTrue(launcherCleanupCalled == maxFailures);
+    assertTrue(launcherLaunchCalled == maxFailures);
+    assertTrue(schedulerAddApplication == maxFailures);
+    assertTrue(schedulerRemoveApplication == maxFailures);
+    assertTrue(masterInfo.getFailedCount() == maxFailures);
+    waitForFailed(masterInfo, ApplicationState.FAILED);
+    stop = true;
+  }
+}
\ No newline at end of file

Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java
 Thu Mar 24 07:52:35 2011
@@ -211,6 +211,10 @@ private static class StatusContext imple
   public String getQueue() {
     return null;
   }
+  @Override
+  public int getFailedCount() {
+    return 0;
+  }
 }
 
 private class ApplicationTracker implements 
EventHandler<ASMEvent<ApplicationTrackerEventType>> {

Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java
 Thu Mar 24 07:52:35 2011
@@ -51,9 +51,7 @@ import org.apache.hadoop.yarn.server.res
 import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType;
 import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.SNEventType;
 import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo;
-import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker;
-import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse;
-import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTrackerImpl;
+import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse;
 import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler;
 import org.apache.hadoop.yarn.server.resourcemanager.scheduler.YarnScheduler;
 import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.fifo.FifoScheduler;
@@ -69,25 +67,17 @@ public class TestApplicationCleanup exte
   private static final Log LOG = 
LogFactory.getLog(TestApplicationCleanup.class);
   private AtomicInteger waitForState = new AtomicInteger(0);
   private ResourceScheduler scheduler;
-  private ClusterTracker clusterTracker;
   private final int memoryCapability = 1024;
   private ExtASM asm;
   private static final int memoryNeeded = 100;
 
   private final ASMContext context = new ResourceManager.ASMContextImpl();
 
-  private class ExtFifoScheduler extends FifoScheduler {
-    @Override
-    protected ClusterTracker createClusterTracker()  {
-      clusterTracker = new ClusterTrackerImpl();
-      return clusterTracker;
-    }
-  }
-
   @Before
   public void setUp() {
     new DummyApplicationTracker();
-    scheduler = new ExtFifoScheduler();
+    scheduler = new FifoScheduler();
+    context.getDispatcher().register(ApplicationTrackerEventType.class, 
scheduler);
     asm = new ExtASM(new ApplicationTokenSecretManager(), scheduler);
     asm.init(new Configuration());
   }

Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java
 Thu Mar 24 07:52:35 2011
@@ -39,8 +39,9 @@ import org.apache.hadoop.yarn.server.res
 import 
org.apache.hadoop.yarn.server.resourcemanager.ResourceManager.ASMContext;
 import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ASMEvent;
 import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationEventType;
+import 
org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType;
 import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo;
-import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse;
+import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse;
 import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler;
 import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager;
 import org.junit.After;
@@ -71,15 +72,7 @@ public class TestSchedulerNegotiator ext
       containers.add(container);
       return containers;
     }
-    @Override
-    public void addApplication(ApplicationID applicationId, String user,
-        String unused, Priority priority)
-        throws IOException {
-    }
-    @Override
-    public void removeApplication(ApplicationID applicationId)
-        throws IOException {
-    }
+
     @Override
     public void reinitialize(Configuration conf,
         ContainerTokenSecretManager secretManager) {
@@ -97,6 +90,10 @@ public class TestSchedulerNegotiator ext
     @Override
     public void removeNode(NodeInfo node) {
     }
+
+    @Override
+    public void handle(ASMEvent<ApplicationTrackerEventType> event) {
+    }
   }
   
   @Before

Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java
 Thu Mar 24 07:52:35 2011
@@ -33,8 +33,8 @@ import org.apache.hadoop.yarn.conf.YarnC
 import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo;
 import 
org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.RMResourceTrackerImpl;
 import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeManager;
+import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse;
 import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceListener;
-import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse;
 import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager;
 import org.apache.hadoop.yarn.ApplicationID;
 import org.apache.hadoop.yarn.Container;
@@ -74,7 +74,7 @@ public class TestNMExpiry extends TestCa
   private class TestRMResourceTrackerImpl extends RMResourceTrackerImpl {
     public TestRMResourceTrackerImpl(
         ContainerTokenSecretManager containerTokenSecretManager) {
-      super(containerTokenSecretManager);
+      super(containerTokenSecretManager, new VoidResourceListener());
     }
 
     @Override
@@ -99,7 +99,6 @@ public class TestNMExpiry extends TestCa
   @Before
   public void setUp() {
     resourceTracker = new 
TestRMResourceTrackerImpl(containerTokenSecretManager);
-    resourceTracker.register(new VoidResourceListener());
     Configuration conf = new Configuration();
     conf.setLong(YarnConfiguration.NM_EXPIRY_INTERVAL, 1000);
     resourceTracker.init(conf);

Modified: 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java?rev=1084866&r1=1084865&r2=1084866&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java
 Thu Mar 24 07:52:35 2011
@@ -27,8 +27,8 @@ import org.apache.hadoop.net.Node;
 import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo;
 import 
org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.RMResourceTrackerImpl;
 import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeManager;
+import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse;
 import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceListener;
-import 
org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse;
 import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager;
 import org.apache.hadoop.yarn.ApplicationID;
 import org.apache.hadoop.yarn.Container;
@@ -71,8 +71,7 @@ public class TestRMNMRPCResponseId exten
   
   @Before
   public void setUp() {
-    rmResourceTrackerImpl = new 
RMResourceTrackerImpl(containerTokenSecretManager);
-    rmResourceTrackerImpl.register(listener);
+    rmResourceTrackerImpl = new 
RMResourceTrackerImpl(containerTokenSecretManager, listener);
   }
   
   @After


Reply via email to