-
Notifications
You must be signed in to change notification settings - Fork 16
Expand file tree
/
Copy pathProxyImpl.kt
More file actions
108 lines (101 loc) · 3.82 KB
/
Copy pathProxyImpl.kt
File metadata and controls
108 lines (101 loc) · 3.82 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
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
// Copyright (c) 2023 - Restate Software, Inc., Restate GmbH
//
// This file is part of the Restate Java SDK,
// which is released under the MIT license.
//
// You can find a copy of the license in file LICENSE in the root
// directory of this repository or package, or at
// https://github.com/restatedev/sdk-java/blob/main/LICENSE
package dev.restate.sdk.testservices
import dev.restate.common.Request
import dev.restate.common.Target
import dev.restate.sdk.kotlin.*
import dev.restate.sdk.testservices.contracts.Proxy
import dev.restate.serde.Serde
import kotlin.time.Duration
import kotlin.time.Duration.Companion.milliseconds
class ProxyImpl : Proxy {
private fun Proxy.ProxyRequest.toTarget(): Target {
val target =
if (this.virtualObjectKey == null) {
Target.service(this.serviceName, this.handlerName)
} else {
Target.virtualObject(this.serviceName, this.virtualObjectKey, this.handlerName)
}
return if (this.scope != null) target.scoped(this.scope) else target
}
override suspend fun call(request: Proxy.ProxyRequest): ByteArray {
return prepareRequest(
Request.of(request.toTarget(), Serde.RAW, Serde.RAW, request.message).also {
if (request.idempotencyKey != null) {
it.idempotencyKey = request.idempotencyKey
}
if (request.limitKey != null) {
it.limitKey = request.limitKey
}
}
)
.call()
.await()
}
override suspend fun oneWayCall(request: Proxy.ProxyRequest): String =
prepareRequest(
Request.of(request.toTarget(), Serde.RAW, Serde.SLICE, request.message).also {
if (request.idempotencyKey != null) {
it.idempotencyKey = request.idempotencyKey
}
if (request.limitKey != null) {
it.limitKey = request.limitKey
}
}
)
.send(request.delayMillis?.milliseconds ?: Duration.ZERO)
.invocationId()
override suspend fun manyCalls(requests: List<Proxy.ManyCallRequest>) {
val toAwait = mutableListOf<DurableFuture<ByteArray>>()
for (request in requests) {
if (request.oneWayCall) {
prepareRequest(
Request.of(
request.proxyRequest.toTarget(),
Serde.RAW,
Serde.SLICE,
request.proxyRequest.message,
)
.also {
if (request.proxyRequest.idempotencyKey != null) {
it.idempotencyKey = request.proxyRequest.idempotencyKey
}
if (request.proxyRequest.limitKey != null) {
it.limitKey = request.proxyRequest.limitKey
}
}
)
.send(request.proxyRequest.delayMillis?.milliseconds ?: Duration.ZERO)
} else {
val fut =
prepareRequest(
Request.of(
request.proxyRequest.toTarget(),
Serde.RAW,
Serde.RAW,
request.proxyRequest.message,
)
.also {
if (request.proxyRequest.idempotencyKey != null) {
it.idempotencyKey = request.proxyRequest.idempotencyKey
}
if (request.proxyRequest.limitKey != null) {
it.limitKey = request.proxyRequest.limitKey
}
}
)
.call()
if (request.awaitAtTheEnd) {
toAwait.add(fut)
}
}
}
toAwait.toList().awaitAll()
}
}