Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
21 changes: 20 additions & 1 deletion nodejs/consumer.js
Original file line number Diff line number Diff line change
Expand Up @@ -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)) {
Expand Down
40 changes: 39 additions & 1 deletion nodejs/provider.js
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
});
}
}
19 changes: 17 additions & 2 deletions python/src/arete_sdk/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand Down
36 changes: 34 additions & 2 deletions python/src/arete_sdk/provider.py
Original file line number Diff line number Diff line change
@@ -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)
24 changes: 22 additions & 2 deletions rust/src/consumer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,12 +32,32 @@ impl Consumer {
}
}

pub fn watch(&self) -> Result<Receiver<ChangeEvent>, 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<Value>) -> Result<Option<Value>, 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<Receiver<ChangeEvent>, 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+)$")?;
Expand Down
52 changes: 49 additions & 3 deletions rust/src/provider.rs
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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<Value>) -> Result<Option<Value>, 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<Receiver<ChangeEvent>, 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)
}
}