-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathAsynchronousLogic.scala
More file actions
66 lines (51 loc) · 2.38 KB
/
Copy pathAsynchronousLogic.scala
File metadata and controls
66 lines (51 loc) · 2.38 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
package net.officefloor.demo
import cats.effect.IO
import net.officefloor.frame.api.function.ManagedFunctionContext
import reactor.core.publisher.Mono
import zio.ZIO
object AsynchronousLogic {
// START SNIPPET: async
def cats(request: ServerRequest)(implicit repository: AsyncMessageRepository): IO[ServerResponse] =
for {
message <- catsGetMessage(request.getId)
response = new ServerResponse(s"${message.getContent} via Cats")
} yield response
def catsGetMessage(id: Int)(implicit repository: AsyncMessageRepository): IO[Message] =
IO.async(callback => IO {
repository.getMessage(id, result => callback.apply(result))
Some(IO())
})
def zio(request: ServerRequest, repository: AsyncMessageRepository): ZIO[Any, Throwable, ServerResponse] = {
// Server logic
val response = for {
message <- zioGetMessage(request.getId)
response = new ServerResponse(s"${message.getContent} via ZIO")
} yield response
// Provide dependencies
response.provide(new InjectAsyncMessageRepository {
override val asyncMessageRepository = repository
})
}
def zioGetMessage(id: Int): ZIO[InjectAsyncMessageRepository, Throwable, Message] =
ZIO.accessM(env => ZIO.effectAsync(callback => env.asyncMessageRepository.getMessage(id, result => callback.apply(ZIO.fromEither(result)))))
trait InjectAsyncMessageRepository {
val asyncMessageRepository: AsyncMessageRepository
}
def reactor(request: ServerRequest)(implicit repository: AsyncMessageRepository): Mono[ServerResponse] =
reactorGetMessage(request.getId).map(message => new ServerResponse(s"${message.getContent} via Reactor"))
def reactorGetMessage(id: Int)(implicit repository: AsyncMessageRepository): Mono[Message] =
Mono.create { callback =>
repository.getMessage(id, _ match {
case Left(error) => callback.error(error)
case Right(message) => callback.success(message)
})
}
def imperative(request: ServerRequest, repository: AsyncMessageRepository, context: ManagedFunctionContext[_, _]): Unit = {
val async = context.createAsynchronousFlow()
repository.getMessage(request.getId, _ match {
case Left(error) => async.complete(() => throw error)
case Right(message) => async.complete(() => context.setNextFunctionArgument(new ServerResponse(s"${message.getContent} via Imperative")))
})
}
// END SNIPPET: async
}