mirror of
https://github.com/status-im/status-python-sdk.git
synced 2026-08-30 21:51:14 +00:00
bot: Initial Status App account
Initial SDK draft 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
This commit is contained in:
+561
@@ -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()
|
||||
+1
-1
@@ -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']}"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user