Skip to content

Commit c8f166f

Browse files
committed
changed(moonzero): take base64url from moonbase, and declare the derived impls.
Signed-off-by: Leo Cheng (heke1228) <chengkelfan@qq.com>
1 parent ff38366 commit c8f166f

22 files changed

Lines changed: 268 additions & 58 deletions

discov/consul_socket_test.mbt

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ fn consul_test_host() -> String {
1717
///|
1818
/// Register `endpoint` for `service` and pass its TTL check so it is healthy at once.
1919
async fn consul_put_healthy(
20-
sock : ConsulSocket,
20+
sock : @discov.ConsulSocket,
2121
service : String,
2222
id : String,
2323
host : String,
@@ -36,7 +36,7 @@ async fn consul_put_healthy(
3636
///|
3737
async test "integration: consul discovery register/resolve/deregister over a real consul" {
3838
guard @env.get_env_var("MOON_CONSUL_TEST") is Some(_) else { return }
39-
let sock = ConsulSocket::new(consul_test_host())
39+
let sock = @discov.ConsulSocket::new(consul_test_host())
4040
let id1 = "greeter-10.0.0.1-9090"
4141
let id2 = "greeter-10.0.0.2-9090"
4242
// Clean any leftover from a prior run.

discov/discov_test.mbt

Lines changed: 11 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -9,14 +9,14 @@ async test "file registry round-trips a real file: reload picks up a persisted i
99
let dir = @fs.tmpdir(prefix="mz-reg-rt")
1010
let path = dir + "/registry.json"
1111
let writer = @moonzero.PersistentRegistry::new()
12-
persist_registry(path, writer)
13-
let reader = FileRegistry::load(path)
12+
@discov.persist_registry(path, writer)
13+
let reader = @discov.FileRegistry::load(path)
1414
assert_eq(reader.resolve("greet.Greeter").length(), 0)
1515
let _ = writer.register(
1616
"greet.Greeter",
1717
@moonzero.Endpoint::new("10.0.0.7", 9001),
1818
)
19-
persist_registry(path, writer)
19+
@discov.persist_registry(path, writer)
2020
let events = reader.reload()
2121
assert_eq(events.length(), 1)
2222
match events[0] {
@@ -41,12 +41,12 @@ async test "file registry reload emits a Delete when an instance leaves" {
4141
"greet.Greeter",
4242
@moonzero.Endpoint::new("10.0.0.7", 9001),
4343
)
44-
persist_registry(path, writer)
45-
let reader = FileRegistry::load(path)
44+
@discov.persist_registry(path, writer)
45+
let reader = @discov.FileRegistry::load(path)
4646
assert_eq(reader.resolve("greet.Greeter").length(), 1)
4747
let removed = writer.deregister("greet.Greeter", key)
4848
assert_eq(removed, true)
49-
persist_registry(path, writer)
49+
@discov.persist_registry(path, writer)
5050
let events = reader.reload()
5151
assert_eq(events.length(), 1)
5252
match events[0] {
@@ -71,8 +71,8 @@ async test "file registry drives a balancer after a real reload" {
7171
"greet.Greeter",
7272
@moonzero.Endpoint::new("10.0.0.5", 8443),
7373
)
74-
persist_registry(path, writer)
75-
let reader = FileRegistry::load(path)
74+
@discov.persist_registry(path, writer)
75+
let reader = @discov.FileRegistry::load(path)
7676
let _ = reader.reload()
7777
let resolve = reader.resolver()
7878
match @moonzero.pick_first(resolve("greet.Greeter")) {
@@ -92,8 +92,8 @@ async test "file registry watcher fires on a real file change" {
9292
let dir = @fs.tmpdir(prefix="mz-watch")
9393
let path = dir + "/registry.json"
9494
let writer = @moonzero.PersistentRegistry::new()
95-
persist_registry(path, writer)
96-
let reader = FileRegistry::load(path)
95+
@discov.persist_registry(path, writer)
96+
let reader = @discov.FileRegistry::load(path)
9797
let watcher = @fs.Watcher(dir)
9898
let holder : Ref[Array[@moonzero.RegistryEvent]] = Ref([])
9999
let task = g.spawn(() => holder.val = reader.watch_once(watcher))
@@ -102,7 +102,7 @@ async test "file registry watcher fires on a real file change" {
102102
"greet.Greeter",
103103
@moonzero.Endpoint::new("10.0.0.9", 9100),
104104
)
105-
persist_registry(path, writer)
105+
@discov.persist_registry(path, writer)
106106
task.wait()
107107
watcher.close()
108108
let events = holder.val

discov/etcd_socket.mbt

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -144,3 +144,6 @@ pub async fn EtcdSocket::lease_grant(
144144
self.unary_bytes("/etcdserverpb.Lease/LeaseGrant", req.encode()),
145145
)
146146
}
147+
148+
///|
149+
pub extend EtcdSocket with LeaseKeeper::{forget, renew, grant}

discov/etcd_socket_test.mbt

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ fn etcd_test_host() -> String {
1818
///|
1919
/// Grant a lease and Put `endpoint` under `<prefix>/<lease-id>`, returning the lease id.
2020
async fn etcd_put_instance(
21-
sock : EtcdSocket,
21+
sock : @discov.EtcdSocket,
2222
svc_prefix : String,
2323
endpoint : @moonzero.Endpoint,
2424
) -> Int64 {
@@ -34,7 +34,7 @@ async fn etcd_put_instance(
3434
///|
3535
async test "integration: etcd discovery register/resolve/deregister over a real etcd" {
3636
guard @env.get_env_var("MOON_ETCD_TEST") is Some(_) else { return }
37-
let sock = EtcdSocket::connect(etcd_test_host(), 2379)
37+
let sock = @discov.EtcdSocket::connect(etcd_test_host(), 2379)
3838
defer sock.close()
3939
let svc_prefix = @moonzero.etcd_service_prefix("moonzero-ci/", "greeter")
4040
let prefix_bytes = @utf8.encode(svc_prefix)
@@ -71,13 +71,13 @@ async test "integration: etcd watch observes a live PUT over a real watch stream
7171
@async.with_task_group(g => {
7272
let key = @utf8.encode("moonzero-ci/watch/probe")
7373
// Clean any leftover on a setup connection.
74-
let setup = EtcdSocket::connect(etcd_test_host(), 2379)
74+
let setup = @discov.EtcdSocket::connect(etcd_test_host(), 2379)
7575
let _ = setup.delete_range({ key, range_end: b"", prev_kv: false, })
7676
setup.close()
7777
// A concurrent publisher Puts the key once the watch has had time to establish.
7878
let putter = g.spawn(() => {
7979
@async.sleep(500)
80-
let pub_conn = EtcdSocket::connect(etcd_test_host(), 2379)
80+
let pub_conn = @discov.EtcdSocket::connect(etcd_test_host(), 2379)
8181
let _ = pub_conn.put({
8282
key,
8383
value: @utf8.encode("10.0.0.1:9090"),
@@ -86,7 +86,7 @@ async test "integration: etcd watch observes a live PUT over a real watch stream
8686
pub_conn.close()
8787
})
8888
// Watch and read the created ack plus the PUT event the publisher triggers.
89-
let sock = EtcdSocket::connect(etcd_test_host(), 2379)
89+
let sock = @discov.EtcdSocket::connect(etcd_test_host(), 2379)
9090
let responses = sock.watch({ key, range_end: b"", start_revision: 0L, }, 2)
9191
sock.close()
9292
let _ = putter.wait()
@@ -103,7 +103,7 @@ async test "integration: etcd watch observes a live PUT over a real watch stream
103103
}
104104
assert_eq(saw_put, true)
105105
// Cleanup.
106-
let cleanup = EtcdSocket::connect(etcd_test_host(), 2379)
106+
let cleanup = @discov.EtcdSocket::connect(etcd_test_host(), 2379)
107107
let _ = cleanup.delete_range({ key, range_end: b"", prev_kv: false, })
108108
cleanup.close()
109109
})

discov/etcd_watch_test.mbt

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
// The etcd Watch stream, driven against a mock etcd Watch RPC over a real socket: the
2-
// server streams a `created` acknowledgement then a PUT event, and `EtcdSocket::watch` reads
2+
// server streams a `created` acknowledgement then a PUT event, and `@discov.EtcdSocket::watch` reads
33
// and decodes both off the wire. The real-etcd counterpart is gated in the integration test.
44

55
///|
@@ -35,7 +35,7 @@ async test "etcd watch: reads created then a PUT event off a streamed Watch RPC"
3535
})
3636
let ts = g.spawn(() => server.serve(port=18140))
3737
@async.sleep(300)
38-
let sock = EtcdSocket::connect("127.0.0.1", 18140)
38+
let sock = @discov.EtcdSocket::connect("127.0.0.1", 18140)
3939
let responses = sock.watch(
4040
{ key: b"svc/", range_end: b"svc0", start_revision: 0L, },
4141
2,

discov/moon.pkg

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,6 @@ supported_targets = "native"
66

77
import {
88
"moonbitstack/moonzero",
9-
"moonbitstack/moonrpc",
109
"moonbitstack/moonrpc/net",
1110
"moonbitlang/async",
1211
"moonbitlang/async/fs",

discov/redis_limit_test.mbt

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -213,9 +213,9 @@ async test "redis_rate_limit: two assembled apps share one bucket" {
213213
/// arrangement two replicas have. Gated on `MOON_REDIS_TEST`.
214214
async test "integration: two period limiters on a live redis share the window" {
215215
guard @env.get_env_var("MOON_REDIS_TEST") is Some(_) else { return }
216-
let s1 = RedisSocket::connect(redis_test_host(), 6379)
216+
let s1 = @discov.RedisSocket::connect(redis_test_host(), 6379)
217217
defer s1.close()
218-
let s2 = RedisSocket::connect(redis_test_host(), 6379)
218+
let s2 = @discov.RedisSocket::connect(redis_test_host(), 6379)
219219
defer s2.close()
220220
let prefix = "moonzero-ci-limit/"
221221
let key = "period"
@@ -248,9 +248,9 @@ async test "integration: two period limiters on a live redis share the window" {
248248
/// clock so the refill step is exact. Gated on `MOON_REDIS_TEST`.
249249
async test "integration: two token limiters on a live redis share the bucket" {
250250
guard @env.get_env_var("MOON_REDIS_TEST") is Some(_) else { return }
251-
let s1 = RedisSocket::connect(redis_test_host(), 6379)
251+
let s1 = @discov.RedisSocket::connect(redis_test_host(), 6379)
252252
defer s1.close()
253-
let s2 = RedisSocket::connect(redis_test_host(), 6379)
253+
let s2 = @discov.RedisSocket::connect(redis_test_host(), 6379)
254254
defer s2.close()
255255
let key = "moonzero-ci-limit/token"
256256
let _ = s1.command(["DEL", "{" + key + "}.tokens", "{" + key + "}.ts"])

discov/redis_socket_test.mbt

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,10 @@ fn redis_test_host() -> String {
1818
///|
1919
/// SCAN the whole `pattern` keyspace, following the cursor to completion, returning the
2020
/// matched keys as text.
21-
async fn scan_all(sock : RedisSocket, pattern : String) -> Array[String] {
21+
async fn scan_all(
22+
sock : @discov.RedisSocket,
23+
pattern : String,
24+
) -> Array[String] {
2225
let out : Array[String] = []
2326
let mut cursor = "0"
2427
let mut first = true
@@ -50,7 +53,7 @@ async fn scan_all(sock : RedisSocket, pattern : String) -> Array[String] {
5053

5154
///|
5255
/// Delete every key under `pattern` (idempotent cleanup before a run).
53-
async fn scan_clean(sock : RedisSocket, pattern : String) -> Unit {
56+
async fn scan_clean(sock : @discov.RedisSocket, pattern : String) -> Unit {
5457
let keys = scan_all(sock, pattern)
5558
if keys.length() > 0 {
5659
let args = ["DEL"]
@@ -64,7 +67,7 @@ async fn scan_clean(sock : RedisSocket, pattern : String) -> Unit {
6467
///|
6568
async test "integration: RESP discovery register/resolve/deregister over a real redis" {
6669
guard @env.get_env_var("MOON_REDIS_TEST") is Some(_) else { return }
67-
let sock = RedisSocket::connect(redis_test_host(), 6379)
70+
let sock = @discov.RedisSocket::connect(redis_test_host(), 6379)
6871
defer sock.close()
6972
// Liveness over the real socket.
7073
assert_eq(sock.command(["PING"]), @moonzero.SimpleString("PONG"))

discovery.mbt

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -493,3 +493,9 @@ pub fn LoadBalancedChannel::call_bidi_streaming(
493493
Err(status) => Err(status)
494494
}
495495
}
496+
497+
///|
498+
pub extend RegistryEvent with Debug::{to_repr}
499+
500+
///|
501+
pub extend RegistryEvent with Eq::{not_equal, equal}

0 commit comments

Comments
 (0)