Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Wait for Actor Start message to be processed before returning workflo… #308

Merged
merged 2 commits into from Dec 4, 2015
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Jump to
Jump to file
Failed to load files.
Diff view
Diff view
7 changes: 5 additions & 2 deletions src/main/scala/cromwell/engine/workflow/WorkflowActor.scala
Expand Up @@ -305,7 +305,7 @@ case class WorkflowActor(workflow: WorkflowDescriptor, backend: Backend)
}
}

private def initializeExecutionStore(startMode: StartMode): Unit = {
private def initializeExecutionStore(startMode: StartMode): Future[Unit] = {
val initializationCode = startMode.runInitialization(this)
val futureCaches = for {
_ <- initializationCode
Expand All @@ -320,6 +320,9 @@ case class WorkflowActor(workflow: WorkflowDescriptor, backend: Backend)
case Failure(t) =>
self ! AsyncFailure(t)
}

// "Lose" the actual value on purpose, the caller doesn't care/need to know about it
futureCaches map { _ => () }
}

private def initializeWorkflow: Try[HostInputs] = backend.initializeForWorkflow(workflow)
Expand All @@ -340,7 +343,7 @@ case class WorkflowActor(workflow: WorkflowDescriptor, backend: Backend)
when(WorkflowSubmitted) {
case Event(startMode: StartMode, _) =>
logger.info(s"$startMode message received")
initializeExecutionStore(startMode)
initializeExecutionStore(startMode) pipeTo sender()
stay()
case Event(CachesCreated(startMode), data) =>
logger.info(s"ExecutionStoreCreated($startMode) message received")
Expand Down
Expand Up @@ -3,7 +3,7 @@ package cromwell.engine.workflow
import akka.actor.FSM.SubscribeTransitionCallBack
import akka.actor.{Actor, ActorRef, Props}
import akka.event.{Logging, LoggingReceive}
import akka.pattern.pipe
import akka.pattern.{pipe, ask}
import com.typesafe.config.ConfigFactory
import cromwell.binding
import cromwell.binding._
Expand Down Expand Up @@ -243,15 +243,14 @@ class WorkflowManagerActor(backend: Backend) extends Actor with CromwellActor {
maybeWorkflowId: Option[WorkflowId]): Future[WorkflowId] = {
val workflowId: WorkflowId = maybeWorkflowId.getOrElse(WorkflowId.randomId())
log.info(s"$tag submitWorkflow input id = $maybeWorkflowId, effective id = $workflowId")
val isRestart = maybeWorkflowId.isDefined

val futureId = for {
descriptor <- Future.fromTry(Try(new WorkflowDescriptor(workflowId, source)))
workflowActor = context.actorOf(WorkflowActor.props(descriptor, backend), s"WorkflowActor-$workflowId")
_ <- Future.fromTry(workflowStore.insert(workflowId, workflowActor))
} yield {
val isRestart = maybeWorkflowId.isDefined
workflowActor ! (if (isRestart) Restart else Start)
workflowId
}
_ <- workflowActor ? (if (isRestart) Restart else Start)
} yield workflowId

futureId onFailure {
case e =>
Expand Down