Skip to content

Commit dc1dedb

Browse files
committed
feat(moonzero): etcd Watch stream over a real socket — subscribe and observe live key events.
Signed-off-by: 林晨 (Leo Cheng) <chengkelfan@qq.com>
1 parent 52477f2 commit dc1dedb

4 files changed

Lines changed: 122 additions & 1 deletion

File tree

discov/etcd_socket.mbt

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,29 @@ pub async fn EtcdSocket::delete_range(
8383
)
8484
}
8585

86+
///|
87+
/// `Watch.Watch`: subscribe to changes on the key or range in `req` and read the first
88+
/// `count` `WatchResponse`s off the stream — the `created` acknowledgement, then events as
89+
/// the watched keys change. The responses are read inline (no background reader), so the
90+
/// caller collects what it needs and closes the connection; a live subscriber loops calling
91+
/// this. Changes are driven elsewhere (a `Put` on another connection), the way go-zero's
92+
/// discov subscriber watches etcd while publishers register.
93+
pub async fn EtcdSocket::watch(
94+
self : EtcdSocket,
95+
req : @moonzero.EtcdWatchCreateRequest,
96+
count : Int,
97+
) -> Array[@moonzero.EtcdWatchResponse] {
98+
let request = @moonzero.EtcdWatchRequest::Create(req).encode()
99+
let raw = self.channel.stream_take(
100+
"/etcdserverpb.Watch/Watch", request, count,
101+
)
102+
let out : Array[@moonzero.EtcdWatchResponse] = []
103+
for bytes in raw {
104+
out.push(@moonzero.EtcdWatchResponse::decode(bytes))
105+
}
106+
out
107+
}
108+
86109
///|
87110
/// `Lease.LeaseGrant`: obtain a lease with the requested TTL.
88111
pub async fn EtcdSocket::lease_grant(

discov/etcd_socket_test.mbt

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,3 +64,47 @@ async test "integration: etcd discovery register/resolve/deregister over a real
6464
let after = sock.range({ key: prefix_bytes, range_end, limit: 0 })
6565
assert_eq(after.kvs.length(), 0)
6666
}
67+
68+
///|
69+
async test "integration: etcd watch observes a live PUT over a real watch stream" {
70+
guard @env.get_env_var("MOON_ETCD_TEST") is Some(_) else { return }
71+
@async.with_task_group(g => {
72+
let key = @utf8.encode("moonzero-ci/watch/probe")
73+
// Clean any leftover on a setup connection.
74+
let setup = EtcdSocket::connect(etcd_test_host(), 2379)
75+
let _ = setup.delete_range({ key, range_end: b"", prev_kv: false })
76+
setup.close()
77+
// A concurrent publisher Puts the key once the watch has had time to establish.
78+
let putter = g.spawn(() => {
79+
@async.sleep(500)
80+
let pub_conn = EtcdSocket::connect(etcd_test_host(), 2379)
81+
let _ = pub_conn.put({
82+
key,
83+
value: @utf8.encode("10.0.0.1:9090"),
84+
lease: 0L,
85+
})
86+
pub_conn.close()
87+
})
88+
// Watch and read the created ack plus the PUT event the publisher triggers.
89+
let sock = EtcdSocket::connect(etcd_test_host(), 2379)
90+
let responses = sock.watch({ key, range_end: b"", start_revision: 0L }, 2)
91+
sock.close()
92+
let _ = putter.wait()
93+
assert_eq(responses.length() >= 2, true)
94+
assert_eq(responses[0].created, true)
95+
let mut saw_put = false
96+
for r in responses {
97+
for e in r.events {
98+
if e.event_type == @moonzero.EtcdEventType::Put &&
99+
e.kv.value == @utf8.encode("10.0.0.1:9090") {
100+
saw_put = true
101+
}
102+
}
103+
}
104+
assert_eq(saw_put, true)
105+
// Cleanup.
106+
let cleanup = EtcdSocket::connect(etcd_test_host(), 2379)
107+
let _ = cleanup.delete_range({ key, range_end: b"", prev_kv: false })
108+
cleanup.close()
109+
})
110+
}

discov/etcd_watch_test.mbt

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
// 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
3+
// and decodes both off the wire. The real-etcd counterpart is gated in the integration test.
4+
5+
///|
6+
async test "etcd watch: reads created then a PUT event off a streamed Watch RPC" {
7+
@async.with_task_group(g => {
8+
let created : @moonzero.EtcdWatchResponse = {
9+
watch_id: 7L,
10+
created: true,
11+
canceled: false,
12+
events: [],
13+
}
14+
let event : @moonzero.EtcdWatchResponse = {
15+
watch_id: 7L,
16+
created: false,
17+
canceled: false,
18+
events: [
19+
{
20+
event_type: @moonzero.EtcdEventType::Put,
21+
kv: {
22+
key: b"svc/1",
23+
create_revision: 2L,
24+
mod_revision: 3L,
25+
version: 1L,
26+
value: b"10.0.0.1:9090",
27+
lease: 0L,
28+
},
29+
},
30+
],
31+
}
32+
let server = @net.GrpcServer::new()
33+
server.register_server_streaming("/etcdserverpb.Watch/Watch", (_ctx, _req) => {
34+
[created.encode(), event.encode()]
35+
})
36+
let ts = g.spawn(() => server.serve(port=18140))
37+
@async.sleep(300)
38+
let sock = EtcdSocket::connect("127.0.0.1", 18140)
39+
let responses = sock.watch(
40+
{ key: b"svc/", range_end: b"svc0", start_revision: 0L },
41+
2,
42+
)
43+
sock.close()
44+
ts.cancel()
45+
assert_eq(responses.length(), 2)
46+
assert_eq(responses[0].created, true)
47+
assert_eq(responses[1].events.length(), 1)
48+
assert_eq(
49+
responses[1].events[0].event_type == @moonzero.EtcdEventType::Put,
50+
true,
51+
)
52+
assert_eq(responses[1].events[0].kv.value == b"10.0.0.1:9090", true)
53+
})
54+
}

moon.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,6 @@ description = "moonzero — a service framework for MoonBit (← go-zero): confi
2222
import {
2323
"Lfan-ke/moonapi@0.6.2",
2424
"Lfan-ke/moonasgi@0.6.1",
25-
"Lfan-ke/moonrpc@0.8.0",
25+
"Lfan-ke/moonrpc@0.8.1",
2626
"moonbitlang/async@0.20.3",
2727
}

0 commit comments

Comments
 (0)