|
|
@@ -1,13 +1,13 @@
|
|
|
package com.cloudera.hue.livy.server
|
|
|
|
|
|
-import com.cloudera.hue.livy.server.sessions.{SessionFailedtoStart, Session}
|
|
|
+import com.cloudera.hue.livy.server.sessions.{SessionFailedToStart, Session}
|
|
|
import com.fasterxml.jackson.core.JsonParseException
|
|
|
import org.json4s.{DefaultFormats, Formats, MappingException}
|
|
|
import org.scalatra._
|
|
|
import org.scalatra.json.JacksonJsonSupport
|
|
|
|
|
|
import scala.concurrent.duration.Duration
|
|
|
-import scala.concurrent.{Await, ExecutionContext, ExecutionContextExecutor}
|
|
|
+import scala.concurrent.{Future, Await, ExecutionContext, ExecutionContextExecutor}
|
|
|
|
|
|
object WebApp {
|
|
|
case class CreateSessionRequest(lang: String)
|
|
|
@@ -43,7 +43,7 @@ class WebApp(sessionManager: SessionManager)
|
|
|
}
|
|
|
|
|
|
val rep = sessionFuture.map {
|
|
|
- case session => Map("id" -> session.id, "state" -> session.state)
|
|
|
+ case session => formatSession(session)
|
|
|
}
|
|
|
|
|
|
// FIXME: this is silently eating exceptions.
|
|
|
@@ -51,35 +51,9 @@ class WebApp(sessionManager: SessionManager)
|
|
|
Await.result(rep, Duration.Inf)
|
|
|
}
|
|
|
|
|
|
- val getStatements = get("/sessions/:sessionId/statements") {
|
|
|
+ get("/sessions/:sessionId") {
|
|
|
sessionManager.get(params("sessionId")) match {
|
|
|
- case Some(session: Session) =>
|
|
|
- val statements = session.statements()
|
|
|
-
|
|
|
- // FIXME: this is silently eating exceptions.
|
|
|
- //new AsyncResult() { val is = statements }
|
|
|
- Await.result(statements, Duration.Inf)
|
|
|
- case None => NotFound("Session not found")
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- val getSession = get("/sessions/:sessionId") {
|
|
|
- sessionManager.get(params("sessionId")) match {
|
|
|
- case Some(session) => Map("id" -> session.id, "state" -> session.state)
|
|
|
- case None => NotFound("Session not found")
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- post("/sessions/:sessionId/statements") {
|
|
|
- val req = parsedBody.extract[ExecuteStatementRequest]
|
|
|
-
|
|
|
- sessionManager.get(params("sessionId")) match {
|
|
|
- case Some(session) =>
|
|
|
- val statement = session.executeStatement(req.statement)
|
|
|
-
|
|
|
- // FIXME: this is silently eating exceptions.
|
|
|
- //new AsyncResult() { val is = statement }
|
|
|
- Await.result(statement, Duration.Inf)
|
|
|
+ case Some(session) => formatSession(session)
|
|
|
case None => NotFound("Session not found")
|
|
|
}
|
|
|
}
|
|
|
@@ -99,7 +73,9 @@ class WebApp(sessionManager: SessionManager)
|
|
|
post("/sessions/:sessionId/interrupt") {
|
|
|
sessionManager.get(params("sessionId")) match {
|
|
|
case Some(session) =>
|
|
|
- val future = session.interrupt()
|
|
|
+ val future = for {
|
|
|
+ _ <- session.interrupt()
|
|
|
+ } yield Accepted()
|
|
|
|
|
|
// FIXME: this is silently eating exceptions.
|
|
|
//new AsyncResult() { val is = for { _ <- future } yield NoContent }
|
|
|
@@ -109,22 +85,47 @@ class WebApp(sessionManager: SessionManager)
|
|
|
}
|
|
|
|
|
|
delete("/sessions/:sessionId") {
|
|
|
- val future = sessionManager.delete(params("sessionId"))
|
|
|
+ val future = for {
|
|
|
+ _ <- sessionManager.delete(params("sessionId"))
|
|
|
+ } yield Accepted()
|
|
|
|
|
|
// FIXME: this is silently eating exceptions.
|
|
|
//new AsyncResult() { val is = for { _ <- future } yield NoContent }
|
|
|
Await.result(future, Duration.Inf)
|
|
|
}
|
|
|
|
|
|
+ get("/sessions/:sessionId/statements") {
|
|
|
+ sessionManager.get(params("sessionId")) match {
|
|
|
+ case Some(session: Session) => session.statements().map(formatStatement)
|
|
|
+ case None => NotFound("Session not found")
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ post("/sessions/:sessionId/statements") {
|
|
|
+ val req = parsedBody.extract[ExecuteStatementRequest]
|
|
|
|
|
|
- val getStatement = get("/sessions/:sessionId/statements/:statementId") {
|
|
|
sessionManager.get(params("sessionId")) match {
|
|
|
case Some(session) =>
|
|
|
- val statement = session.statement(params("statementId").toInt)
|
|
|
+ Future {
|
|
|
+ val statement: Statement = session.executeStatement(req.statement)
|
|
|
|
|
|
- // FIXME: this is silently eating exceptions.
|
|
|
- //new AsyncResult() { val is = statement }
|
|
|
- Await.result(statement, Duration.Inf)
|
|
|
+ // FIXME: this is silently eating exceptions.
|
|
|
+ //new AsyncResult() { val is = statement }
|
|
|
+ Await.result(statement.output, Duration.Inf)
|
|
|
+ }
|
|
|
+
|
|
|
+ Accepted()
|
|
|
+ case None => NotFound("Session not found")
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ get("/sessions/:sessionId/statements/:statementId") {
|
|
|
+ sessionManager.get(params("sessionId")) match {
|
|
|
+ case Some(session) =>
|
|
|
+ session.statement(params("statementId").toInt) match {
|
|
|
+ case Some(statement) => formatStatement(statement)
|
|
|
+ case None => NotFound("Statement not found")
|
|
|
+ }
|
|
|
case None => NotFound("Session not found")
|
|
|
}
|
|
|
}
|
|
|
@@ -132,8 +133,23 @@ class WebApp(sessionManager: SessionManager)
|
|
|
error {
|
|
|
case e: JsonParseException => halt(400, e.getMessage)
|
|
|
case e: MappingException => halt(400, e.getMessage)
|
|
|
- case e: SessionFailedtoStart => halt(500, e.getMessage)
|
|
|
+ case e: SessionFailedToStart => halt(500, e.getMessage)
|
|
|
case e: dispatch.StatusCode => halt(e.code, e.getMessage)
|
|
|
case t => throw t
|
|
|
}
|
|
|
+
|
|
|
+ private def formatSession(session: Session) = {
|
|
|
+ Map(
|
|
|
+ "id" -> session.id,
|
|
|
+ "state" -> session.state.getClass.getSimpleName.toLowerCase
|
|
|
+ )
|
|
|
+ }
|
|
|
+
|
|
|
+ private def formatStatement(statement: Statement) = {
|
|
|
+ Map(
|
|
|
+ "id" -> statement.id,
|
|
|
+ "state" -> statement.state.getClass.getSimpleName.toLowerCase,
|
|
|
+ "output" -> statement.output
|
|
|
+ )
|
|
|
+ }
|
|
|
}
|