This is an automated email from the ASF dual-hosted git repository.

dragos pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-openwhisk.git


The following commit(s) were added to refs/heads/master by this push:
     new 6e45e26  add support for mesos attribute constraints when launching 
containers using MesosContainerFactory (#3559)
6e45e26 is described below

commit 6e45e26cd5c0f177d810e21870b4b9741f67c424
Author: tysonnorris <[email protected]>
AuthorDate: Thu Apr 19 18:07:03 2018 -0700

    add support for mesos attribute constraints when launching containers using 
MesosContainerFactory (#3559)
---
 common/scala/build.gradle                          |  2 +-
 common/scala/src/main/resources/application.conf   |  4 ++
 .../whisk/core/mesos/MesosContainerFactory.scala   | 41 +++++++++++++++--
 .../main/scala/whisk/core/mesos/MesosTask.scala    | 12 +++--
 docs/mesos.md                                      |  4 ++
 .../mesos/test/MesosContainerFactoryTest.scala     | 52 ++++++++++++----------
 6 files changed, 82 insertions(+), 33 deletions(-)

diff --git a/common/scala/build.gradle b/common/scala/build.gradle
index 6c856c2..4a7a82e 100644
--- a/common/scala/build.gradle
+++ b/common/scala/build.gradle
@@ -45,7 +45,7 @@ dependencies {
     compile 'io.kamon:kamon-core_2.11:0.6.7'
     compile 'io.kamon:kamon-statsd_2.11:0.6.7'
     //for mesos
-    compile 'com.adobe.api.platform.runtime:mesos-actor:0.0.4'
+    compile 'com.adobe.api.platform.runtime:mesos-actor:0.0.7'
 }
 
 tasks.withType(ScalaCompile) {
diff --git a/common/scala/src/main/resources/application.conf 
b/common/scala/src/main/resources/application.conf
index 4c36319..e3fadaf 100644
--- a/common/scala/src/main/resources/application.conf
+++ b/common/scala/src/main/resources/application.conf
@@ -162,5 +162,9 @@ whisk {
         role = "*" //see 
http://mesos.apache.org/documentation/latest/roles/#associating-frameworks-with-roles
         failover-timeout = 0 seconds  //Timeout allowed for framework to 
reconnect after disconnection.
         mesos-link-log-message = true //If true, display a link to mesos in 
the static log message, otherwise do not include a link to mesos.
+        constraints = [] //placement constraint strings to use for managed 
containers e.g. ["att1 LIKE v1", "att2 UNLIKE v2"]
+        blackbox-constraints = [] //placement constraints to use for blackbox 
containers
+        constraint-delimiter = " "//used to parse constraint strings
+        teardown-on-exit = true //set to true to disable the mesos framework 
on system exit; set for false for HA deployments
     }
 }
diff --git 
a/common/scala/src/main/scala/whisk/core/mesos/MesosContainerFactory.scala 
b/common/scala/src/main/scala/whisk/core/mesos/MesosContainerFactory.scala
index c9e634e..b0d47a9 100644
--- a/common/scala/src/main/scala/whisk/core/mesos/MesosContainerFactory.scala
+++ b/common/scala/src/main/scala/whisk/core/mesos/MesosContainerFactory.scala
@@ -20,10 +20,14 @@ package whisk.core.mesos
 import akka.actor.ActorRef
 import akka.actor.ActorSystem
 import akka.pattern.ask
+import com.adobe.api.platform.runtime.mesos.Constraint
+import com.adobe.api.platform.runtime.mesos.LIKE
+import com.adobe.api.platform.runtime.mesos.LocalTaskStore
 import com.adobe.api.platform.runtime.mesos.MesosClient
 import com.adobe.api.platform.runtime.mesos.Subscribe
 import com.adobe.api.platform.runtime.mesos.SubscribeComplete
 import com.adobe.api.platform.runtime.mesos.Teardown
+import com.adobe.api.platform.runtime.mesos.UNLIKE
 import java.time.Instant
 import pureconfig.loadConfigOrThrow
 import scala.concurrent.Await
@@ -58,7 +62,11 @@ case class MesosConfig(masterUrl: String,
                        masterPublicUrl: Option[String],
                        role: String,
                        failoverTimeout: FiniteDuration,
-                       mesosLinkLogMessage: Boolean)
+                       mesosLinkLogMessage: Boolean,
+                       constraints: Seq[String],
+                       constraintDelimiter: String,
+                       blackboxConstraints: Seq[String],
+                       teardownOnExit: Boolean) {}
 
 class MesosContainerFactory(config: WhiskConfig,
                             actorSystem: ActorSystem,
@@ -108,6 +116,11 @@ class MesosContainerFactory(config: WhiskConfig,
     } else {
       actionImage.localImageName(config.dockerRegistry, 
config.dockerImagePrefix, Some(config.dockerImageTag))
     }
+    val constraintStrings = if (userProvidedImage) {
+      mesosConfig.blackboxConstraints
+    } else {
+      mesosConfig.constraints
+    }
 
     logging.info(this, s"using Mesos to create a container with image 
$image...")
     MesosTask.create(
@@ -126,9 +139,28 @@ class MesosContainerFactory(config: WhiskConfig,
       //strip any "--" prefixes on parameters (should make this consistent 
everywhere else)
       parameters
         .map({ case (k, v) => if (k.startsWith("--")) (k.replaceFirst("--", 
""), v) else (k, v) })
-        ++ containerArgs.extraArgs)
+        ++ containerArgs.extraArgs,
+      parseConstraints(constraintStrings))
   }
 
+  /**
+   * Validate that constraint strings are well formed, and ignore constraints 
with unknown operators
+   * @param constraintStrings
+   * @param logging
+   * @return
+   */
+  def parseConstraints(constraintStrings: Seq[String])(implicit logging: 
Logging): Seq[Constraint] =
+    constraintStrings.flatMap(cs => {
+      val parts = cs.split(mesosConfig.constraintDelimiter)
+      require(parts.length == 3, "constraint must be in the form 
<attribute><delimiter><operator><delimiter><value>")
+      Seq(LIKE, UNLIKE).find(_.toString == parts(1)) match {
+        case Some(o) => Some(Constraint(parts(0), o, parts(2)))
+        case _ =>
+          logging.warn(this, s"ignoring unsupported constraint operator 
${parts(1)}")
+          None
+      }
+    })
+
   override def init(): Unit = Unit
 
   /** Cleanups any remaining Containers; should block until complete; should 
ONLY be run at shutdown. */
@@ -148,11 +180,12 @@ object MesosContainerFactory {
     actorSystem.actorOf(
       MesosClient
         .props(
-          "whisk-containerfactory-" + UUID(),
+          () => "whisk-containerfactory-" + UUID(),
           "whisk-containerfactory-framework",
           mesosConfig.masterUrl,
           mesosConfig.role,
-          mesosConfig.failoverTimeout))
+          mesosConfig.failoverTimeout,
+          taskStore = new LocalTaskStore))
 
   val counter = new Counter()
   val startTime = Instant.now.getEpochSecond
diff --git a/common/scala/src/main/scala/whisk/core/mesos/MesosTask.scala 
b/common/scala/src/main/scala/whisk/core/mesos/MesosTask.scala
index 1d6ad86..0e50543 100644
--- a/common/scala/src/main/scala/whisk/core/mesos/MesosTask.scala
+++ b/common/scala/src/main/scala/whisk/core/mesos/MesosTask.scala
@@ -24,6 +24,8 @@ import akka.stream.scaladsl.Source
 import akka.util.ByteString
 import akka.util.Timeout
 import com.adobe.api.platform.runtime.mesos.Bridge
+import com.adobe.api.platform.runtime.mesos.CommandDef
+import com.adobe.api.platform.runtime.mesos.Constraint
 import com.adobe.api.platform.runtime.mesos.DeleteTask
 import com.adobe.api.platform.runtime.mesos.Host
 import com.adobe.api.platform.runtime.mesos.Running
@@ -74,9 +76,10 @@ object MesosTask {
              network: String = "bridge",
              dnsServers: Seq[String] = Seq(),
              name: Option[String] = None,
-             parameters: Map[String, Set[String]] = Map())(implicit ec: 
ExecutionContext,
-                                                           log: Logging,
-                                                           as: ActorSystem): 
Future[Container] = {
+             parameters: Map[String, Set[String]] = Map(),
+             constraints: Seq[Constraint] = Seq.empty)(implicit ec: 
ExecutionContext,
+                                                       log: Logging,
+                                                       as: ActorSystem): 
Future[Container] = {
     implicit val tid = transid
 
     log.info(this, s"creating task for image $image...")
@@ -104,7 +107,8 @@ object MesosTask {
       false,
       taskNetwork,
       dnsOrEmpty ++ parameters,
-      environment)
+      Some(CommandDef(environment)),
+      constraints.toSet)
 
     val launched: Future[Running] =
       mesosClientActor.ask(SubmitTask(task))(taskLaunchTimeout).mapTo[Running]
diff --git a/docs/mesos.md b/docs/mesos.md
index bf69508..eefa263 100644
--- a/docs/mesos.md
+++ b/docs/mesos.md
@@ -14,6 +14,10 @@ To enable MesosContainerFactory, use the following TypeSafe 
Config properties
 | `whisk.mesos.role` | optional (default *) | Mesos framework role| any string 
e.g. `openwhisk` |
 | `whisk.mesos.failover-timeout-seconds` | optional (default 0) | how long to 
wait for the framework to reconnect with the same id before tasks are 
terminated  | see 
http://mesos.apache.org/documentation/latest/high-availability-framework-guide/ 
|
 | `whisk.mesos.mesos-link-log-message` | optional (default true) | display a 
log message with a link to Mesos when using the default LogStore (or no log 
message) | Since logs are not available for invoker to collect from Mesos in 
general, you can either use an alternate LogStore or direct users to the Mesos 
ui |   |
+| `whisk.mesos.constraints` | optional (default []) | placement constraint 
strings to use for managed containers | `["att1 LIKE v1", "att2 UNLIKE v2"]` |  
 |
+| `whisk.mesos.blackbox-constraints` | optional (default []) | placement 
constraint strings to use for blackbox containers  | `["att1 LIKE v1", "att2 
UNLIKE v2"]` |   |
+| `whisk.mesos.constraint-delimiter` | optional (default " ") | delimiter used 
to parse constraints |  |   |
+| `whisk.mesos.teardown-on-exit` | optional (default true) | set to true to 
disable the mesos framework on system exit; set for false for HA deployments |  
|   |
 
 To set these properties for your invoker, set the corresponding environment 
variables e.g.,
 ```properties
diff --git 
a/tests/src/test/scala/whisk/core/containerpool/mesos/test/MesosContainerFactoryTest.scala
 
b/tests/src/test/scala/whisk/core/containerpool/mesos/test/MesosContainerFactoryTest.scala
index 9a46509..51c62c5 100644
--- 
a/tests/src/test/scala/whisk/core/containerpool/mesos/test/MesosContainerFactoryTest.scala
+++ 
b/tests/src/test/scala/whisk/core/containerpool/mesos/test/MesosContainerFactoryTest.scala
@@ -23,17 +23,19 @@ import akka.stream.scaladsl.Sink
 import akka.testkit.TestKit
 import akka.testkit.TestProbe
 import com.adobe.api.platform.runtime.mesos.Bridge
+import com.adobe.api.platform.runtime.mesos.CommandDef
+import com.adobe.api.platform.runtime.mesos.Constraint
 import com.adobe.api.platform.runtime.mesos.DeleteTask
+import com.adobe.api.platform.runtime.mesos.LIKE
 import com.adobe.api.platform.runtime.mesos.Running
 import com.adobe.api.platform.runtime.mesos.SubmitTask
 import com.adobe.api.platform.runtime.mesos.Subscribe
 import com.adobe.api.platform.runtime.mesos.SubscribeComplete
 import com.adobe.api.platform.runtime.mesos.TaskDef
+import com.adobe.api.platform.runtime.mesos.UNLIKE
 import com.adobe.api.platform.runtime.mesos.User
 import common.StreamLogging
-import org.apache.mesos.v1.Protos.AgentID
 import org.apache.mesos.v1.Protos.TaskID
-import org.apache.mesos.v1.Protos.TaskInfo
 import org.apache.mesos.v1.Protos.TaskState
 import org.apache.mesos.v1.Protos.TaskStatus
 import org.junit.runner.RunWith
@@ -90,7 +92,7 @@ class MesosContainerFactoryTest
 
   it should "send Subscribe on init" in {
     val wskConfig = new WhiskConfig(Map())
-    val mesosConfig = MesosConfig("http://master:5050";, None, "*", 0.seconds, 
true)
+    val mesosConfig = MesosConfig("http://master:5050";, None, "*", 0.seconds, 
true, Seq.empty, " ", Seq.empty, true)
     new MesosContainerFactory(
       wskConfig,
       system,
@@ -103,8 +105,17 @@ class MesosContainerFactoryTest
     expectMsg(Subscribe)
   }
 
-  it should "send SubmitTask on create" in {
-    val mesosConfig = MesosConfig("http://master:5050";, None, "*", 0.seconds, 
true)
+  it should "send SubmitTask (with constraints) on create" in {
+    val mesosConfig = MesosConfig(
+      "http://master:5050";,
+      None,
+      "*",
+      0.seconds,
+      true,
+      Seq("att1 LIKE v1", "att2 UNLIKE v2"),
+      " ",
+      Seq("bbatt1 LIKE v1", "bbatt2 UNLIKE v2"),
+      true)
 
     val factory =
       new MesosContainerFactory(
@@ -144,11 +155,12 @@ class MesosContainerFactoryTest
           "dns" -> Set("dns1", "dns2"),
           "extra1" -> Set("e1", "e2"),
           "extra2" -> Set("e3", "e4")),
-        Map("__OW_API_HOST" -> wskConfig.wskApiHost))))
+        Some(CommandDef(Map("__OW_API_HOST" -> wskConfig.wskApiHost))),
+        Seq(Constraint("att1", LIKE, "v1"), Constraint("att2", UNLIKE, 
"v2")).toSet)))
   }
 
   it should "send DeleteTask on destroy" in {
-    val mesosConfig = MesosConfig("http://master:5050";, None, "*", 0.seconds, 
true)
+    val mesosConfig = MesosConfig("http://master:5050";, None, "*", 0.seconds, 
true, Seq.empty, " ", Seq.empty, true)
 
     val probe = TestProbe()
     val factory =
@@ -164,7 +176,7 @@ class MesosContainerFactoryTest
 
     probe.expectMsg(Subscribe)
     //emulate successful subscribe
-    probe.reply(new SubscribeComplete)
+    probe.reply(new SubscribeComplete("testid"))
 
     //create the container
     val c = factory.createContainer(
@@ -192,19 +204,15 @@ class MesosContainerFactoryTest
           "dns" -> Set("dns1", "dns2"),
           "extra1" -> Set("e1", "e2"),
           "extra2" -> Set("e3", "e4")),
-        Map("__OW_API_HOST" -> wskConfig.wskApiHost))))
+        Some(CommandDef(Map("__OW_API_HOST" -> wskConfig.wskApiHost))))))
 
     //emulate successful task launch
     val taskId = TaskID.newBuilder().setValue(lastTaskId)
 
     probe.reply(
       Running(
-        TaskInfo
-          .newBuilder()
-          .setName("testTask")
-          .setTaskId(taskId)
-          .setAgentId(AgentID.newBuilder().setValue("testAgentID"))
-          .build(),
+        taskId.getValue,
+        "testAgentID",
         
TaskStatus.newBuilder().setTaskId(taskId).setState(TaskState.TASK_RUNNING).build(),
         "agenthost",
         Seq(30000)))
@@ -223,7 +231,7 @@ class MesosContainerFactoryTest
   }
 
   it should "return static message for logs" in {
-    val mesosConfig = MesosConfig("http://master:5050";, None, "*", 0.seconds, 
true)
+    val mesosConfig = MesosConfig("http://master:5050";, None, "*", 0.seconds, 
true, Seq.empty, " ", Seq.empty, true)
 
     val probe = TestProbe()
     val factory =
@@ -239,7 +247,7 @@ class MesosContainerFactoryTest
 
     probe.expectMsg(Subscribe)
     //emulate successful subscribe
-    probe.reply(new SubscribeComplete)
+    probe.reply(new SubscribeComplete("testid"))
 
     //create the container
     val c = factory.createContainer(
@@ -267,19 +275,15 @@ class MesosContainerFactoryTest
           "other" -> Set("v5", "v6"),
           "extra1" -> Set("e1", "e2"),
           "extra2" -> Set("e3", "e4")),
-        Map("__OW_API_HOST" -> wskConfig.wskApiHost))))
+        Some(CommandDef(Map("__OW_API_HOST" -> wskConfig.wskApiHost))))))
 
     //emulate successful task launch
     val taskId = TaskID.newBuilder().setValue(lastTaskId)
 
     probe.reply(
       Running(
-        TaskInfo
-          .newBuilder()
-          .setName("testTask")
-          .setTaskId(taskId)
-          .setAgentId(AgentID.newBuilder().setValue("testAgentID"))
-          .build(),
+        taskId.getValue,
+        "testAgentID",
         
TaskStatus.newBuilder().setTaskId(taskId).setState(TaskState.TASK_RUNNING).build(),
         "agenthost",
         Seq(30000)))

-- 
To stop receiving notification emails like this one, please contact
[email protected].

Reply via email to