From b84207e0429ffdd259731dec8ac3a9735e56a364 Mon Sep 17 00:00:00 2001 From: David Rauschenbach Date: Mon, 22 Sep 2025 11:16:49 -0700 Subject: [PATCH 1/4] (Rust) Consumer and Provider classes/structs both require get, put, and watch capabilities --- rust/src/consumer.rs | 24 ++++++++++++++++++-- rust/src/provider.rs | 52 +++++++++++++++++++++++++++++++++++++++++--- 2 files changed, 71 insertions(+), 5 deletions(-) diff --git a/rust/src/consumer.rs b/rust/src/consumer.rs index 0ec6fe4..937a044 100644 --- a/rust/src/consumer.rs +++ b/rust/src/consumer.rs @@ -32,12 +32,32 @@ impl Consumer { } } - pub fn watch(&self) -> Result, Error> { + fn profile_key_prefix(&self) -> String { let system_id = self.context.node.system.id.to_string(); let node_id = &self.context.node.id; let context_id = &self.context.id; let profile = &self.profile; - let key_prefix = format!("cns/{system_id}/nodes/{node_id}/contexts/{context_id}/consumer/{profile}/"); + format!("cns/{system_id}/nodes/{node_id}/contexts/{context_id}/consumer/{profile}/") + } + + fn property_key(&self, property: &str) -> String { + let profile_key_prefix = self.profile_key_prefix(); + format!("{profile_key_prefix}properties/{property}") + } + + pub fn get(&self, property: &str, default_value: Option) -> Result, Error> { + let key = self.property_key(property); + self.client.get(&key, default_value) + } + + pub fn put(&self, property: &str, value: &str) -> Result<(), Error> { + let key = self.property_key(property); + let mut client = self.client.clone(); + client.put(&key, value) + } + + pub fn watch(&self) -> Result, Error> { + let key_prefix = self.profile_key_prefix(); let upstream_rx = self.client.clone().on_update()?; let (tx, rx) = mpsc::channel(); let re = Regex::new(r"connections/(\w+)/properties/(\w+)$")?; diff --git a/rust/src/provider.rs b/rust/src/provider.rs index 7f2e3b5..27a3ad6 100644 --- a/rust/src/provider.rs +++ b/rust/src/provider.rs @@ -1,4 +1,9 @@ use super::{Client, Context, Error}; +use crate::consumer::ChangeEvent; +use regex::Regex; +use serde_json::Value; +use std::sync::mpsc; +use std::sync::mpsc::Receiver; pub struct Provider { client: Client, @@ -15,14 +20,55 @@ impl Provider { } } - pub fn put(&self, property: &str, value: &str) -> Result<(), Error> { + fn profile_key_prefix(&self) -> String { let system_id = self.context.node.system.id.to_string(); let node_id = &self.context.node.id; let context_id = &self.context.id; let profile = &self.profile; - let key = - format!("cns/{system_id}/nodes/{node_id}/contexts/{context_id}/provider/{profile}/properties/{property}"); + format!("cns/{system_id}/nodes/{node_id}/contexts/{context_id}/provider/{profile}/") + } + + fn property_key(&self, property: &str) -> String { + let profile_key_prefix = self.profile_key_prefix(); + format!("{profile_key_prefix}properties/{property}") + } + + pub fn get(&self, property: &str, default_value: Option) -> Result, Error> { + let key = self.property_key(property); + self.client.get(&key, default_value) + } + + pub fn put(&self, property: &str, value: &str) -> Result<(), Error> { + let key = self.property_key(property); let mut client = self.client.clone(); client.put(&key, value) } + + pub fn watch(&self) -> Result, Error> { + let key_prefix = self.profile_key_prefix(); + let upstream_rx = self.client.clone().on_update()?; + let (tx, rx) = mpsc::channel(); + let re = Regex::new(r"connections/(\w+)/properties/(\w+)$")?; + std::thread::spawn(move || { + for event in upstream_rx { + for (k, v) in event.keys.iter() { + if !k.starts_with(&key_prefix) { + continue; + } + if let Some(captures) = re.captures(k) { + let connection = captures[1].to_string(); + let property = captures[2].to_string(); + let value = v.clone(); + let change_event = ChangeEvent { + connection, + property, + value, + }; + tx.send(change_event).unwrap(); + }; + } + } + }); + Ok(rx) + } } From 2c3b780a974e71ddac5a7e0c7004b09d5cf804e7 Mon Sep 17 00:00:00 2001 From: David Rauschenbach Date: Mon, 22 Sep 2025 11:41:07 -0700 Subject: [PATCH 2/4] (NodeJS) Consumer and Provider classes/structs both require get, put, and watch capabilities --- CHANGELOG.md | 4 ++++ nodejs/consumer.js | 21 ++++++++++++++++++++- nodejs/provider.js | 40 +++++++++++++++++++++++++++++++++++++++- 3 files changed, 63 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1d8d451..fee8746 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,9 @@ # Project Arete SDK Changelog +## [Unreleased] +### Added +- [#90](https://github.com/project-arete/sdk/issues/90) Consumer and Provider classes/structs both require get, put, and watch capabilities + ## [0.1.5] - 2025-09-19 ### Added - [#81](https://github.com/project-arete/sdk/issues/81) (Python) Publish tag releases to pypi.org diff --git a/nodejs/consumer.js b/nodejs/consumer.js index 0bf6d3b..2707f27 100644 --- a/nodejs/consumer.js +++ b/nodejs/consumer.js @@ -11,8 +11,27 @@ export class Consumer { this.profile = profile; } + #profileKeyPrefix() { + return `cns/${this.context.node.system.id}/nodes/${this.context.node.id}/contexts/${this.context.id}/consumer/${this.profile}/`; + } + + #propertyKey(property) { + const profileKeyPrefix = this.#profileKeyPrefix(); + return `${profileKeyPrefix}properties/${property}`; + } + + get(property, def = null) { + const key = this.#propertyKey(property); + return this.#client.get(key, def); + } + + put(property, value) { + const key = this.#propertyKey(property); + return this.#client.put(key, value); + } + watch(handler) { - const keyPrefix = `cns/${this.context.node.system.id}/nodes/${this.context.node.id}/contexts/${this.context.id}/consumer/${this.profile}/`; + const keyPrefix = this.#profileKeyPrefix(); const re = new RegExp(`connections/(\\w+)/properties/(\\w+)$`); this.#client.on('update', (event) => { for (let [key, value] of Object.entries(event.keys)) { diff --git a/nodejs/provider.js b/nodejs/provider.js index 1c7be8b..4a5e473 100644 --- a/nodejs/provider.js +++ b/nodejs/provider.js @@ -11,8 +11,46 @@ export class Provider { this.profile = profile; } + #profileKeyPrefix() { + return `cns/${this.context.node.system.id}/nodes/${this.context.node.id}/contexts/${this.context.id}/provider/${this.profile}/`; + } + + #propertyKey(property) { + const profileKeyPrefix = this.#profileKeyPrefix(); + return `${profileKeyPrefix}properties/${property}`; + } + + get(property, def = null) { + const key = this.#propertyKey(property); + return this.#client.get(key, def); + } + put(property, value) { - const key = `cns/${this.context.node.system.id}/nodes/${this.context.node.id}/contexts/${this.context.id}/provider/${this.profile}/properties/${property}`; + const key = this.#propertyKey(property); return this.#client.put(key, value); } + + watch(handler) { + const keyPrefix = this.#profileKeyPrefix(); + const re = new RegExp(`connections/(\\w+)/properties/(\\w+)$`); + this.#client.on('update', (event) => { + for (let [key, value] of Object.entries(event.keys)) { + if (!key.startsWith(keyPrefix)) { + continue; + } + const captures = key.match(re); + if (captures.length < 3) { + continue; + } + const connection = captures[1]; + const property = captures[2]; + const changeEvent = { + connection, + property, + value, + }; + handler(changeEvent); + } + }); + } } From f60095bcdd75cfe1c0e76cbd695960679d682a42 Mon Sep 17 00:00:00 2001 From: David Rauschenbach Date: Mon, 22 Sep 2025 11:43:53 -0700 Subject: [PATCH 3/4] (Python) Consumer and Provider classes/structs both require get, put, and watch capabilities --- python/src/arete_sdk/consumer.py | 19 +++++++++++++++++-- 1 file changed, 17 insertions(+), 2 deletions(-) diff --git a/python/src/arete_sdk/consumer.py b/python/src/arete_sdk/consumer.py index ad4a611..58f21c0 100644 --- a/python/src/arete_sdk/consumer.py +++ b/python/src/arete_sdk/consumer.py @@ -7,12 +7,27 @@ def __init__(self, client, context, profile): self.context = context self.profile = profile - def watch(self, fn): + def profile_key_prefix(self): system_id = self.context.node.system.id node_id = self.context.node.id context_id = self.context.id profile = self.profile - key_prefix = f'cns/{system_id}/nodes/{node_id}/contexts/{context_id}/consumer/{profile}/' + return f'cns/{system_id}/nodes/{node_id}/contexts/{context_id}/consumer/{profile}/' + + def property_key(self, property): + profile_key_prefix = self.profile_key_prefix() + return f'{profile_key_prefix}properties/{property}' + + def get(self, property): + key = self.property_key(property) + self.client.get(key) + + def put(self, property, value): + key = self.property_key(property) + self.client.put(key, value) + + def watch(self, fn): + key_prefix = self.profile_key_prefix() def on_update(event): for key, value in event['keys'].items(): From 0d8ae0a857dc34a5226d86fb3bb3b3490203f1ad Mon Sep 17 00:00:00 2001 From: David Rauschenbach Date: Mon, 22 Sep 2025 11:45:53 -0700 Subject: [PATCH 4/4] (Python) Consumer and Provider classes/structs both require get, put, and watch capabilities --- python/src/arete_sdk/provider.py | 36 ++++++++++++++++++++++++++++++-- 1 file changed, 34 insertions(+), 2 deletions(-) diff --git a/python/src/arete_sdk/provider.py b/python/src/arete_sdk/provider.py index e5995ef..184f8a8 100644 --- a/python/src/arete_sdk/provider.py +++ b/python/src/arete_sdk/provider.py @@ -1,13 +1,45 @@ +import re + + class Provider: def __init__(self, client, context, profile): self.client = client self.context = context self.profile = profile - def put(self, property, value): + def profile_key_prefix(self): system_id = self.context.node.system.id node_id = self.context.node.id context_id = self.context.id profile = self.profile - key = f'cns/{system_id}/nodes/{node_id}/contexts/{context_id}/provider/{profile}/properties/{property}' + return f'cns/{system_id}/nodes/{node_id}/contexts/{context_id}/provider/{profile}/' + + def property_key(self, property): + profile_key_prefix = self.profile_key_prefix() + return f'{profile_key_prefix}properties/{property}' + + def get(self, property): + key = self.property_key(property) + self.client.get(key) + + def put(self, property, value): + key = self.property_key(property) self.client.put(key, value) + + def watch(self, fn): + key_prefix = self.profile_key_prefix() + + def on_update(event): + for key, value in event['keys'].items(): + if not key.startswith(key_prefix): + continue + captures = re.search('connections/(\\w+)/properties/(\\w+)$', key) + connection = captures.group(1) + property = captures.group(2) + change_event = { + 'connection': connection, + 'property': property, + 'value': value, + } + fn(change_event) + self.client.on_update(on_update)