forked from eclipse-zenoh/zenoh-python
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathz_storage.py
More file actions
92 lines (75 loc) · 2.54 KB
/
Copy pathz_storage.py
File metadata and controls
92 lines (75 loc) · 2.54 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
#
# Copyright (c) 2022 ZettaScale Technology
#
# This program and the accompanying materials are made available under the
# terms of the Eclipse Public License 2.0 which is available at
# http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0
# which is available at https://www.apache.org/licenses/LICENSE-2.0.
#
# SPDX-License-Identifier: EPL-2.0 OR Apache-2.0
#
# Contributors:
# ZettaScale Zenoh Team, <zenoh@zettascale.tech>
#
import time
import zenoh
store = {}
def listener(sample: zenoh.Sample):
print(
f">> [Subscriber] Received {sample.kind} ('{sample.key_expr}': '{sample.payload.to_string()}')"
)
if sample.kind == zenoh.SampleKind.DELETE:
store.pop(sample.key_expr, None)
else:
store[sample.key_expr] = sample
def query_handler(query: zenoh.Query):
print(f">> [Queryable ] Received Query '{query.selector}'")
for stored_name, sample in store.items():
if query.key_expr.intersects(stored_name):
query.reply(
sample.key_expr,
sample.payload,
encoding=sample.encoding,
congestion_control=sample.congestion_control,
priority=sample.priority,
express=sample.express,
)
def main(conf: zenoh.Config, key: str, complete: bool):
# initiate logging
zenoh.init_log_from_env_or("error")
print("Opening session...")
with zenoh.open(conf) as session:
print(f"Declaring Subscriber on '{key}'...")
session.declare_subscriber(key, listener)
print(f"Declaring Queryable on '{key}'...")
session.declare_queryable(key, query_handler, complete=complete)
print("Press CTRL-C to quit...")
while True:
time.sleep(1)
# --- Command line argument parsing --- --- --- --- --- ---
if __name__ == "__main__":
import argparse
import json
import common
parser = argparse.ArgumentParser(
prog="z_storage", description="zenoh storage example"
)
common.add_config_arguments(parser)
parser.add_argument(
"--key",
"-k",
dest="key",
default="demo/example/**",
type=str,
help="The key expression matching resources to store.",
)
parser.add_argument(
"--complete",
dest="complete",
default=False,
action="store_true",
help="Declare the storage as complete w.r.t. the key expression.",
)
args = parser.parse_args()
conf = common.get_config_from_args(args)
main(conf, args.key, args.complete)