<http://10gen.com>
package com.mongodb.async
package futures
import com.mongodb.async.util._
import org.bson.util.Logging
import org.bson.SerializableBSONObject
sealed trait RequestFuture {
type T
val body: Either[Throwable, T] => Unit
def apply(error: Throwable) = body(Left(error))
def apply[A <% T](result: A) = body(Right(result.asInstanceOf[T]))
protected[futures] var completed = false
}
sealed trait QueryRequestFuture extends RequestFuture {
type DocType
val decoder: SerializableBSONObject[DocType]
}
trait CursorQueryRequestFuture extends QueryRequestFuture {
type T <: Cursor[DocType]
}
trait GetMoreRequestFuture extends QueryRequestFuture {
type T = (Long, Seq[DocType])
}
trait SingleDocQueryRequestFuture extends QueryRequestFuture with Logging {
type T = DocType
val m: Manifest[T]
}
trait FindAndModifyRequestFuture extends QueryRequestFuture with Logging {
type T = DocType
val m: Manifest[T]
}
trait WriteRequestFuture extends RequestFuture {
type T <: (Option[AnyRef] , WriteResult)
}
trait BatchWriteRequestFuture extends RequestFuture {
type T <: (Option[Seq[AnyRef]] , WriteResult)
}
case object NoOpRequestFuture extends RequestFuture with Logging {
type T = Unit
val body = (result: Either[Throwable, Unit]) => result match {
case Right(()) => {}
case Left(t) => log.error(t, "NoOp Command Failed.")
}
override def toString = "{NoopWriteRequestFuture}"
}
object RequestFutures extends Logging {
def getMore[A: SerializableBSONObject](f: Either[Throwable, (Long, Seq[A])] => Unit) =
new GetMoreRequestFuture {
type DocType = A
val body = f
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{GetMoreRequestFuture}"
}
def query[A: SerializableBSONObject](f: Either[Throwable, Cursor[A]] => Unit) =
new CursorQueryRequestFuture {
type DocType = A
type T = Cursor[A]
val body = f
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{CursorQueryRequestFuture}"
}
def find[A: SerializableBSONObject](f: Either[Throwable, Cursor[A]] => Unit) = query(f)
def command[A: SerializableBSONObject: Manifest](f: Either[Throwable, A] => Unit) =
new SingleDocQueryRequestFuture {
type DocType = A
val m = manifest[A]
val body = f
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{SingleDocQueryRequestFuture}"
}
def findOne[A: SerializableBSONObject: Manifest](f: Either[Throwable, A] => Unit) = command(f)
def findAndModify[A: SerializableBSONObject: Manifest](f: Either[Throwable, Option[A]] => Unit) =
new FindAndModifyRequestFuture {
type DocType = A
val m = manifest[A]
val body = (result: Either[Throwable, A]) => {
log.debug("Decoding SingleDocQueryRequestFuture: %s / %s", toString, decoder)
result match {
case Right(doc) => f(Right(Option(doc)))
case Left(t) => t match {
case nme: NoMatchingDocumentError => f(Right(None))
case e => f(Left(e))
}
}
}
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{SingleDocQueryRequestFuture}"
}
def write(f: Either[Throwable, (Option[AnyRef], WriteResult)] => Unit) =
new WriteRequestFuture {
val body = f
override def toString = "{WriteRequestFuture}"
}
def batchWrite(f: Either[Throwable, (Option[Seq[AnyRef]], WriteResult)] => Unit) =
new BatchWriteRequestFuture {
val body = f
override def toString = "{WriteRequestFuture}"
}
}
object SimpleRequestFutures extends Logging {
def findOne[A: SerializableBSONObject: Manifest](f: A => Unit) = command(f)
def command[A: SerializableBSONObject: Manifest](f: A => Unit) =
new SingleDocQueryRequestFuture {
type DocType = A
val m = manifest[A]
val body = (result: Either[Throwable, A]) => {
log.debug("Decoding SingleDocQueryRequestFuture: %s / %s", toString, decoder)
result match {
case Right(doc) => f(doc)
case Left(t) => log.error(t, "Command Failed.")
}
}
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{SimpleSingleDocQueryRequestFuture}"
}
def findAndModify[A: SerializableBSONObject: Manifest](f: Option[A] => Unit) =
new FindAndModifyRequestFuture {
type DocType = A
val m = manifest[A]
val body = (result: Either[Throwable, A]) => {
log.debug("Decoding SingleDocQueryRequestFuture: %s / %s", toString, decoder)
result match {
case Right(doc) => f(Option(doc))
case Left(t) => t match {
case nme: NoMatchingDocumentError => f(None)
case e => log.error(e, "Command Failed.")
}
}
}
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{SimpleSingleDocQueryRequestFuture}"
}
def getMore[A: SerializableBSONObject](f: (Long, Seq[A]) => Unit) =
new GetMoreRequestFuture {
type DocType = A
val body = (result: Either[Throwable, (Long, Seq[A])]) => result match {
case Right((cid, docs)) => f(cid, docs)
case Left(t) => log.error(t, "GetMore Failed."); throw t
}
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{SimpleGetMoreRequestFuture}"
}
def find[T: SerializableBSONObject](f: Cursor[T] => Unit) = query(f)
def query[A: SerializableBSONObject](f: Cursor[A] => Unit) =
new CursorQueryRequestFuture {
type DocType = A
type T = Cursor[A]
val body = (result: Either[Throwable, Cursor[A]]) => result match {
case Right(cursor) => f(cursor)
case Left(t) => log.error(t, "Query Failed."); throw t
}
val decoder = implicitly[SerializableBSONObject[A]]
override def toString = "{SimpleCursorQueryRequestFuture}"
}
def write(f: (Option[AnyRef], WriteResult) => Unit) =
new WriteRequestFuture {
val body = (result: Either[Throwable, (Option[AnyRef], WriteResult)]) => result match {
case Right((oid, wr)) => f(oid, wr)
case Left(t) => log.error(t, "Command Failed.")
}
override def toString = "{SimpleWriteRequestFuture}"
}
def batchWrite(f: (Option[Seq[AnyRef]], WriteResult) => Unit) =
new BatchWriteRequestFuture {
val body = (result: Either[Throwable, (Option[Seq[AnyRef]], WriteResult)]) => result match {
case Right((oids, wr)) => f(oids, wr)
case Left(t) => log.error(t, "Command Failed.")
}
override def toString = "{SimpleWriteRequestFuture}"
}
}