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