From 0b1a63597a389d1cc7ad7d83e8c3361cd975155f Mon Sep 17 00:00:00 2001 From: Nick Ninov Date: Mon, 9 Mar 2026 12:35:09 +0000 Subject: [PATCH] bot: Status App account Related to: - https://github.com/status-im/data-docs/issues/158 - https://github.com/status-im/data-docs/issues/160 - https://github.com/status-im/data-docs/issues/161 Bug fix: - Incorrect HTTPS URL --- bot/account.py | 561 +++++++++++++++++++++++++++++++++++++++++++++++++ constants.py | 2 +- 2 files changed, 562 insertions(+), 1 deletion(-) create mode 100644 bot/account.py diff --git a/bot/account.py b/bot/account.py new file mode 100644 index 0000000..b80464c --- /dev/null +++ b/bot/account.py @@ -0,0 +1,561 @@ +from typing import Optional, Union, Generator +import requests, datetime, websocket, json, copy, re, queue, threading + +class Signal: + """ + Custom Web Socket app connector that fetches messages / signals from + the running Status Backend Docker container. Currently this class works + only with `/signals` + """ + + # As of now signals should be kept internal until initial + # Python SDK scope is defined. Then the SignalType from the tests + # can be reused. + available_signals = [ + 'messages.new', 'message.delivered', 'node.ready', + 'node.started', 'node.login', 'node.stopped' + ] + + def __init__(self, url: str): + self.__url = url + self.__data = {} + self.__signal_type = None + self.__error_message = None + # For real time data extraction + self.__queue = queue.Queue() + self.__thread = None + + def __on_open(self, ws: websocket.WebSocketApp): + """ + Used to reset class variables before messages / signals + listening starts + """ + self.__data = {} + self.__error_message = None + + def close(self, ws: websocket.WebSocketApp, *args): + """ + Used to reset class variables after messages / signals + have stopped listening or when the object is deleted + """ + self.__signal_type = None + self.__close_thread() + + def __on_error(self, ws: websocket.WebSocketApp, error: str): + """ + Triggered when Status Backend (Docker container) breaks + """ + self.__error_message = f"There was an error with the Status Backend Docker container... Please look at the Docker logs for more information.\nError message: {error}" + + def __get_message(self, ws: websocket.WebSocketApp, signal: str): + """ + Process Status Backend signal and exit. This is used in `get`. + """ + signal: dict = json.loads(signal) + if signal["type"] != self.__signal_type: + return + + event: Optional[dict] = signal.get("event", {}) + if not event: + event = {} + + self.__data = { + "timestamp": datetime.datetime.fromtimestamp(signal["timestamp"]), + "is_error": not isinstance(event.get("error"), type(None)), + "error_message": event.get("error"), + "event": event + } + ws.close() + + def get(self, signal_type: str) -> dict: + """ + Run the the connector to monitor Status Backend for a single message. + NOTE: If not careful, the websocket can end up in an infinite loop + + Parameters: + - `signal_type` - the "type" as it is in Status Backend + + Output: + - the signal data point + """ + self.__signal_type = signal_type + ws = websocket.WebSocketApp( + self.__url, + on_open=self.__on_open, + on_message=self.__get_message, + on_close=self.close, + on_error=self.__on_error + ) + ws.run_forever() + + if self.__error_message: + raise Exception(self.__error_message) + + return copy.deepcopy(self.__data) + + def __listen_message(self, ws: websocket.WebSocketApp, signal: str): + """ + + """ + signal: dict = json.loads(signal) + if signal["type"] != self.__signal_type: + return + + event: Optional[dict] = signal.get("event", {}) or {} + data = { + "timestamp": datetime.datetime.fromtimestamp(signal["timestamp"]), + "is_error": not isinstance(event.get("error"), type(None)), + "error_message": event.get("error"), + "event": event + } + self.__queue.put(data) + + def __close_thread(self): + """ + Called in `listen` / object is deleted and safely close it + """ + if not self.__thread: + return + + if self.__thread.is_alive(): + self.__thread.join(1) + + def listen(self, signal_type: str): + self.__signal_type = signal_type + ws = websocket.WebSocketApp( + self.__url, + on_open=self.__on_open, + on_message=self.__listen_message, + on_close=self.close, + on_error=self.__on_error + ) + self.__thread = threading.Thread(target=ws.run_forever, daemon=True) + self.__thread.start() + while True: + data = self.__queue.get() + if self.__error_message: + raise Exception(self.__error_message) + + yield data + + + +class Account: + + # Enum mappings from original wakuext.py + __mappings = { + "contact_request": { + 0: "none", # No action taken / no association - initial state + 1: "mutual", # Friends + 2: "sent", # Request sent from the bot + 3: "received", # Request sent from another account + 4: "dismissed" # Request cancelled + } + } + + def __init__(self, unix_folder: str = "./data-dir", domain: str = "localhost", port: int = 8080, is_secure: bool = False): + """ + Work with your own Status App account + + Parameters: + - `unix_folder` - where Status Backend files will be initialized. **This folder is required in the Docker container**. Ideally it should be a volume so the account data is persistant if the image is deleted. Folder path is automatically created. + - `domain` - the domain name where Status Backend is running. If running locally it would be `localhost` and if it's running in a container it would be the image's name. + - `port` - the port to connect to Status Backend. Verify the port in the Docker files. + - `is_secure` - if `http` or `https` should be used + """ + self.__timestamp_divisor = 1_000 + self.__kd_iterations = 256000 + self.__unix_folder = unix_folder + self.__is_messenger_launched = False + # Information for the logged in account + self.__info = {} + self.http_base_url = f"http{'s' if is_secure else ''}://{domain}:{port}/statusgo/" + self.ws_base_url = f"ws://{domain}:{port}/" + self.urls = { + "http": { + "initialize": f"{self.http_base_url}InitializeApplication", + "login": f"{self.http_base_url}LoginAccount", + "create": f"{self.http_base_url}CreateAccountAndLogin", + "logout": f"{self.http_base_url}Logout", + "rpc": f"{self.http_base_url}CallRPC" + }, + "socket": { + "signals": f"{self.ws_base_url}signals" + } + } + self.__signal = Signal(self.urls["socket"]["signals"]) + response = requests.post(self.urls["http"]["initialize"], json={ + "dataDir": self.__unix_folder + }) + data: dict = response.json() + accounts: list[dict] = data.get("accounts", []) + if not isinstance(accounts, list): + accounts = [] + + self.__available_accounts = { + account["name"]: { + "key_uid": account["key-uid"], + "created_at": datetime.datetime.fromtimestamp(account["timestamp"]) + } + for account in accounts + } + # In case if there is a hanging logged in session + self.logout() + + def login(self, username: str, password: str): + """ + Login to the given account. If it does not exist, + it will be created and automatically logged in. + + Parameters: + - `username` - your Status display name + - `password` - your Status password + """ + key_uid = self.__available_accounts.get(username, {}).get("key_uid") + is_new_account = isinstance(key_uid, type(None)) + + params = { + "keyUid": key_uid, + "password": password, + 'kdfIterations': self.__kd_iterations + } + if is_new_account: + params = { + "rootDataDir": self.__unix_folder, + "kdfIterations": self.__kd_iterations, + "displayName": username, + "password": password, + "customizationColor": "primary", + "wakuV2LightClient": False, + "thirdpartyServicesEnabled": True, + } + self.logout() + url = self.urls["http"]["login" if not is_new_account else "create"] + response = requests.post(url, json=params) + signal_event = self.__signal.get("node.login") + if signal_event["is_error"]: + raise Exception(f"There was an error with Status Backend...\n{signal_event['error_message']}") + + event: dict = signal_event["event"]["settings"] + self.__info = { + "public_key": event["public-key"], + "emojis": event["emojiHash"], + "key_uid": event["key-uid"], + "mnemonic": event["mnemonic"], + "name": event["display-name"], + "password": password, + "wallet_address": event["address"], + "logged_in_timestamp": datetime.datetime.now() + } + # Messenger can be activated only when logged in + self.__start_messenger() + return self + + def logout(self): + """ + Logout of Status app. In a way this method behaves as a Status cleaner + """ + response = requests.post(self.urls["http"]["logout"]) + self.__info = {} + self.__is_messenger_launched = False + return self + + @property + def info(self) -> dict: + """ + Overall information for currently logged in account. + Can also be used to verify if the user has logged in. + """ + if not self.__info: + raise Exception("Make sure you are logged in to your Status account with login() first...") + return self.__info + + @property + def contacts(self) -> dict[str, dict]: + """ + Get the contacts that the bot has. + This includes contacts that have interacted with us. If a contact has removed us (or the bot has removed us) + + NOTE: We do not use internal state so we can get dynamic values such as: + - Is currently active + - Is currently blocked + - Current display name + - Current bio + + Terminology for Status contact requests: + - approved - when both `contact_state` and `external_contact_state` are `mutual` + - sent request - when `contact_state` is `sent` and `external_contact_state` is `none` + - received request - when `contact_state` is `received` + """ + data = self.__call_rpc("contacts") + raw: list[dict] = data.get("result", []) + if not raw: + return [] + + # dict format can be used in restricting functionality + # such as - `send_message` and `remove_contact` + contacts = { + contact["id"]: { + "public_key": contact["id"], + "chat_id": contact["id"], + "key_uid": contact["compressedKey"], + "emojis": contact["emojiHash"], + "contact_state": self.__mappings["contact_request"][contact["contactRequestState"]], + "external_contact_state": self.__mappings["contact_request"][contact["contactRequestRemoteState"]], + "has_added_us": contact["hasAddedUs"], + "added": contact["added"], + "mutual": contact["mutual"], + "display_name": contact["displayName"], + "bio": contact["bio"], + "wallet_address": contact["address"], + "last_updated": datetime.datetime.fromtimestamp(contact["lastUpdated"] / self.__timestamp_divisor) if contact["lastUpdated"] > 0 else None + } + for contact in raw + } + return contacts + + @property + def communities(self) -> dict[str, dict]: + """ + Get the communities that the bot is in. + NOTE: We do not use internal state so we can get dynamic values such as: + - Current community description + - Current number of community members + - Current channels' names, descriptions and permissions + """ + data = self.__call_rpc("communities") + raw: list[dict] = data.get("result", []) + if not raw: + return [] + + to_datetime = lambda key, mapping: datetime.datetime.fromtimestamp(mapping[key]) if key in mapping else None + communities = [ + { + "id": community["id"], + "name": community["name"], + "verified": community["verified"], + "description": community["description"], + "dialog": community["introMessage"], + "leaving_message": community["outroMessage"], + "tags": community["tags"], + "is_member": community["isMember"], + "joined": community["verified"], + "joined_timestamp": to_datetime("joinedAt", community), + "requested_timestamp": to_datetime("requestedToJoinAt", community), + "encrypted": community["encrypted"], + "members": len(community["members"]), + "channels": [ + { + "id": chat["id"], + "chat_id": community["id"] + chat["id"], + "name": chat["name"], + "description": chat["description"], + "permissions": { + "posting": chat["canPost"], + "viewing": chat["canView"], + "reactions": chat["canPostReactions"], + "token_gated": chat["tokenGated"] + } + } + for chat in community["chats"].values() + ] + } + for community in raw + ] + return communities + + @property + def signal(self) -> Signal: + return self.__signal + + @property + def chats(self) -> list[dict]: + """ + All chats that the bot can send messages to. + This property combines `self.communities` and `self.contacts` chats. + """ + communities = [ + {"type": "channel", "id": chat["chat_id"], "name": f"{community['name']} #{chat['name']}"} + for community in self.communities + for chat in community["channels"] + if chat["permissions"]["posting"] + ] + contacts = [ + {"type": "contact", "id": contact["chat_id"], "name": contact["display_name"]} + for contact in self.contacts.values() + ] + return contacts + communities + + def send_message(self, chat_id: str, message: str): + """ + Send a message to the given chat. + + Parameters: + - `chat_id` - the chat ID can be found in `self.chats` + - `message` - the message that will be sent. Currently only text messages are supported + """ + params = [{ + "chatId": chat_id, + "text": message, + "contentType": 1, # Send message only. Future versions can have different message types (audio, image, etc.) + "responseTo": "" + }] + self.__call_rpc("sendChatMessage", params) + + def listen_messages(self) -> Generator: + """ + Listen for new **RAW** messages continuously. Can be used for real time processing. + """ + self.info + return self.__signal.listen("messages.new") + + def get_messages(self, chat_id: str, start_timestamp: Optional[datetime.datetime] = None, end_timestamp: Optional[datetime.datetime] = None) -> list[dict]: + """ + Get all of the messages in the given start and end timestamps. + Messages are returned in descending order (newest to oldest). + Messages can be fetched for removed contacts as well. + + Parameters: + - `chat_id` - the chat ID can be found in `self.chats` + - `start_timestamp` - the start timestamp for message extraction. If not provided all early messages will be fetched. + - `end_timestamp` - the end timestamp for message extraction. If not provided all latest messages will be fetched. + """ + # NOTE: Order of params matters when making the RCP call + params = { + "chat_id": chat_id, + "cursor": "", + "limit": 500 + } + all_messages = [] + # Keys that need to be converted to datetime.datetime + timestamp_keys = [] + + finished = False + while not finished: + data = self.__call_rpc("chatMessages", list(params.values())) + result: dict[str, Union[str, list[dict]]] = data.get("result", {}) + if result["messages"] and not timestamp_keys: + timestamp_keys = [key for key in result["messages"][0].keys() if "timestamp" in key.lower()] + + for message in result["messages"]: + point = { + self.__camel_to_snake(key): value if key not in timestamp_keys else datetime.datetime.fromtimestamp(value / self.__timestamp_divisor) + for key, value in message.items() + } + + if start_timestamp and point["timestamp"] < start_timestamp: + finished = True + break + + if end_timestamp and point["timestamp"] > end_timestamp: + continue + + all_messages.append(point) + + if len(result["cursor"]) > 0: + params["cursor"] = result["cursor"] + else: + finished = True + + return all_messages + + def add_contact(self, public_key: str, display_name: Optional[str] = None): + """ + Send a contact request / approve a contact. + + Parameters: + - `public_key` - the contact's public key + - `display_name` - this field is required if the `public_key` does not appear in your contacts. This will set their display name (can be different from the one the other user has chosen) + """ + contacts = list(self.contacts.values()) + if not display_name: + for contact in contacts: + if public_key != contact["public_key"]: + continue + + display_name = contact["display_name"] + break + + if not display_name: + raise Exception(f"Cannot add contact {public_key}...\nPlease make sure you add display_name for contacts that you are sending friend requests to!") + + params = [{"id": public_key, "nickname": "", "displayName": display_name, "ensName": ""}] + self.__call_rpc("addContact", params) + + def remove_contact(self, public_key: str) -> bool: + """ + Remove the contact / decline a contact request. + + Parameters: + - `public_key` - the contact's public key + """ + contact_info = self.contacts.get(public_key, {}) + # Cannot remove a contact that is not in your contact + if not contact_info: + return + # Contact has already been removed + if contact_info["contact_state"] == "none": + return + params = [public_key] + self.__call_rpc("removeContact", params) + + def __start_messenger(self): + """ + Start the decentralized messaging service. + This is required for messages to be received / sent. + """ + if self.__is_messenger_launched: + return + self.__call_rpc("startMessenger") + self.__is_messenger_launched = True + + def __del__(self): + """ + Handles automatic logout when calling `del` + and after running `python` + """ + self.logout() + self.__signal.close(None) + + def __call_rpc(self, method_name: str, params: Optional[Union[list, dict]] = None) -> dict: + """ + Make RPC calls to Status Backend + + Parameters: + - `method_name` - the method name as it is in the backend + - `params` - RPC call parameters + + Output: + - the raw output from the RPC method + """ + # Quick initialization check - RPC calls + # can be made only after the user has logged in + self.info + + data = { + 'jsonrpc': '2.0', + # NOTE: Waku may be renamed to Logos Messaging (or something similar) + 'method': f'wakuext_{method_name}', + 'id': None # Original code has an incrementing ID but it does not make a difference + } + if params: + data["params"] = params + + response = requests.get(self.urls["http"]["rpc"], json=data) + return response.json() + + def __camel_to_snake(self, name: str) -> str: + """ + Used to make camel case Status Backend keys + more Pythonic (snake case). Function is used + when the entire raw data point is returned. + + Parameters: + - `name` - camel case dictionary key + + Output: + - snake case `name` + """ + s1 = re.sub(r'(.)([A-Z][a-z]+)', r'\1_\2', name) + s2 = re.sub(r'([a-z0-9])([A-Z])', r'\1_\2', s1) + return s2.lower() diff --git a/constants.py b/constants.py index 595e946..ce88f30 100644 --- a/constants.py +++ b/constants.py @@ -21,5 +21,5 @@ if len(__bot_name) == 0 and os.environ.get("STATUS_USERNAME"): STATUS_BACKEND_PARAMS = { **CONFIG["status_app"]["backend_params"], - "url": f"http:{'s' if CONFIG['status_app']['url']['is_https'] else ''}//{os.environ.get('STATUS_BACKEND_BASE_URL', 'localhost')}:{CONFIG['status_app']['url']['port']}" + "url": f"http{'s' if CONFIG['status_app']['url']['is_https'] else ''}://{os.environ.get('STATUS_BACKEND_BASE_URL', 'localhost')}:{CONFIG['status_app']['url']['port']}" }