Implement FileStorage and MemoryStorage engines
This commit is contained in:
parent
6177abbfa4
commit
6cc9688e49
102
pyrogram/client/storage/file_storage.py
Normal file
102
pyrogram/client/storage/file_storage.py
Normal file
@ -0,0 +1,102 @@
|
|||||||
|
# Pyrogram - Telegram MTProto API Client Library for Python
|
||||||
|
# Copyright (C) 2017-2019 Dan Tès <https://github.com/delivrance>
|
||||||
|
#
|
||||||
|
# This file is part of Pyrogram.
|
||||||
|
#
|
||||||
|
# Pyrogram is free software: you can redistribute it and/or modify
|
||||||
|
# it under the terms of the GNU Lesser General Public License as published
|
||||||
|
# by the Free Software Foundation, either version 3 of the License, or
|
||||||
|
# (at your option) any later version.
|
||||||
|
#
|
||||||
|
# Pyrogram is distributed in the hope that it will be useful,
|
||||||
|
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
# GNU Lesser General Public License for more details.
|
||||||
|
#
|
||||||
|
# You should have received a copy of the GNU Lesser General Public License
|
||||||
|
# along with Pyrogram. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
import base64
|
||||||
|
import json
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import sqlite3
|
||||||
|
from pathlib import Path
|
||||||
|
from sqlite3 import DatabaseError
|
||||||
|
from threading import Lock
|
||||||
|
from typing import Union
|
||||||
|
|
||||||
|
from .memory_storage import MemoryStorage
|
||||||
|
|
||||||
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class FileStorage(MemoryStorage):
|
||||||
|
FILE_EXTENSION = ".session"
|
||||||
|
|
||||||
|
def __init__(self, name: str, workdir: Path):
|
||||||
|
super().__init__(name)
|
||||||
|
|
||||||
|
self.workdir = workdir
|
||||||
|
self.database = workdir / (self.name + self.FILE_EXTENSION)
|
||||||
|
self.conn = None # type: sqlite3.Connection
|
||||||
|
self.lock = Lock()
|
||||||
|
|
||||||
|
# noinspection PyAttributeOutsideInit
|
||||||
|
def migrate_from_json(self, path: Union[str, Path]):
|
||||||
|
log.warning("JSON session storage detected! Pyrogram will now convert it into an SQLite session storage...")
|
||||||
|
|
||||||
|
with open(path, encoding="utf-8") as f:
|
||||||
|
json_session = json.load(f)
|
||||||
|
|
||||||
|
os.remove(path)
|
||||||
|
|
||||||
|
self.open()
|
||||||
|
|
||||||
|
self.dc_id = json_session["dc_id"]
|
||||||
|
self.test_mode = json_session["test_mode"]
|
||||||
|
self.auth_key = base64.b64decode("".join(json_session["auth_key"]))
|
||||||
|
self.user_id = json_session["user_id"]
|
||||||
|
self.date = json_session.get("date", 0)
|
||||||
|
self.is_bot = json_session.get("is_bot", False)
|
||||||
|
|
||||||
|
peers_by_id = json_session.get("peers_by_id", {})
|
||||||
|
peers_by_phone = json_session.get("peers_by_phone", {})
|
||||||
|
|
||||||
|
peers = {}
|
||||||
|
|
||||||
|
for k, v in peers_by_id.items():
|
||||||
|
if v is None:
|
||||||
|
type_ = "group"
|
||||||
|
elif k.startswith("-100"):
|
||||||
|
type_ = "channel"
|
||||||
|
else:
|
||||||
|
type_ = "user"
|
||||||
|
|
||||||
|
peers[int(k)] = [int(k), int(v) if v is not None else None, type_, None, None]
|
||||||
|
|
||||||
|
for k, v in peers_by_phone.items():
|
||||||
|
peers[v][4] = k
|
||||||
|
|
||||||
|
# noinspection PyTypeChecker
|
||||||
|
self.update_peers(peers.values())
|
||||||
|
|
||||||
|
log.warning("Done! The session has been successfully converted from JSON to SQLite storage")
|
||||||
|
|
||||||
|
def open(self):
|
||||||
|
database_exists = os.path.isfile(self.database)
|
||||||
|
|
||||||
|
self.conn = sqlite3.connect(
|
||||||
|
str(self.database),
|
||||||
|
timeout=1,
|
||||||
|
check_same_thread=False
|
||||||
|
)
|
||||||
|
|
||||||
|
try:
|
||||||
|
if not database_exists:
|
||||||
|
self.create()
|
||||||
|
|
||||||
|
with self.conn:
|
||||||
|
self.conn.execute("VACUUM")
|
||||||
|
except DatabaseError:
|
||||||
|
self.migrate_from_json(self.database)
|
241
pyrogram/client/storage/memory_storage.py
Normal file
241
pyrogram/client/storage/memory_storage.py
Normal file
@ -0,0 +1,241 @@
|
|||||||
|
# Pyrogram - Telegram MTProto API Client Library for Python
|
||||||
|
# Copyright (C) 2017-2019 Dan Tès <https://github.com/delivrance>
|
||||||
|
#
|
||||||
|
# This file is part of Pyrogram.
|
||||||
|
#
|
||||||
|
# Pyrogram is free software: you can redistribute it and/or modify
|
||||||
|
# it under the terms of the GNU Lesser General Public License as published
|
||||||
|
# by the Free Software Foundation, either version 3 of the License, or
|
||||||
|
# (at your option) any later version.
|
||||||
|
#
|
||||||
|
# Pyrogram is distributed in the hope that it will be useful,
|
||||||
|
# but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
# GNU Lesser General Public License for more details.
|
||||||
|
#
|
||||||
|
# You should have received a copy of the GNU Lesser General Public License
|
||||||
|
# along with Pyrogram. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
import base64
|
||||||
|
import inspect
|
||||||
|
import logging
|
||||||
|
import sqlite3
|
||||||
|
import struct
|
||||||
|
import time
|
||||||
|
from pathlib import Path
|
||||||
|
from threading import Lock
|
||||||
|
from typing import List, Tuple
|
||||||
|
|
||||||
|
from pyrogram.api import types
|
||||||
|
from pyrogram.client.storage.storage import Storage
|
||||||
|
|
||||||
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class MemoryStorage(Storage):
|
||||||
|
SCHEMA_VERSION = 1
|
||||||
|
USERNAME_TTL = 8 * 60 * 60
|
||||||
|
SESSION_STRING_FMT = ">B?256sI?"
|
||||||
|
SESSION_STRING_SIZE = 351
|
||||||
|
|
||||||
|
def __init__(self, name: str):
|
||||||
|
super().__init__(name)
|
||||||
|
|
||||||
|
self.conn = None # type: sqlite3.Connection
|
||||||
|
self.lock = Lock()
|
||||||
|
|
||||||
|
def create(self):
|
||||||
|
with self.lock, self.conn:
|
||||||
|
with open(Path(__file__).parent / "schema.sql", "r") as schema:
|
||||||
|
self.conn.executescript(schema.read())
|
||||||
|
|
||||||
|
self.conn.execute(
|
||||||
|
"INSERT INTO version VALUES (?)",
|
||||||
|
(self.SCHEMA_VERSION,)
|
||||||
|
)
|
||||||
|
|
||||||
|
self.conn.execute(
|
||||||
|
"INSERT INTO sessions VALUES (?, ?, ?, ?, ?, ?)",
|
||||||
|
(1, None, None, 0, None, None)
|
||||||
|
)
|
||||||
|
|
||||||
|
def _import_session_string(self, string_session: str):
|
||||||
|
decoded = base64.urlsafe_b64decode(string_session + "=" * (-len(string_session) % 4))
|
||||||
|
return struct.unpack(self.SESSION_STRING_FMT, decoded)
|
||||||
|
|
||||||
|
def export_session_string(self):
|
||||||
|
packed = struct.pack(
|
||||||
|
self.SESSION_STRING_FMT,
|
||||||
|
self.dc_id,
|
||||||
|
self.test_mode,
|
||||||
|
self.auth_key,
|
||||||
|
self.user_id,
|
||||||
|
self.is_bot
|
||||||
|
)
|
||||||
|
|
||||||
|
return base64.urlsafe_b64encode(packed).decode().rstrip("=")
|
||||||
|
|
||||||
|
# noinspection PyAttributeOutsideInit
|
||||||
|
def open(self):
|
||||||
|
self.conn = sqlite3.connect(":memory:", check_same_thread=False)
|
||||||
|
self.create()
|
||||||
|
|
||||||
|
if self.name != ":memory:":
|
||||||
|
imported_session_string = self._import_session_string(self.name)
|
||||||
|
|
||||||
|
self.dc_id, self.test_mode, self.auth_key, self.user_id, self.is_bot = imported_session_string
|
||||||
|
self.date = 0
|
||||||
|
|
||||||
|
self.name = ":memory:" + str(self.user_id or "<unknown>")
|
||||||
|
|
||||||
|
# noinspection PyAttributeOutsideInit
|
||||||
|
def save(self):
|
||||||
|
self.date = int(time.time())
|
||||||
|
|
||||||
|
with self.lock:
|
||||||
|
self.conn.commit()
|
||||||
|
|
||||||
|
def close(self):
|
||||||
|
with self.lock:
|
||||||
|
self.conn.close()
|
||||||
|
|
||||||
|
def update_peers(self, peers: List[Tuple[int, int, str, str, str]]):
|
||||||
|
with self.lock:
|
||||||
|
self.conn.executemany(
|
||||||
|
"REPLACE INTO peers (id, access_hash, type, username, phone_number)"
|
||||||
|
"VALUES (?, ?, ?, ?, ?)",
|
||||||
|
peers
|
||||||
|
)
|
||||||
|
|
||||||
|
def clear_peers(self):
|
||||||
|
with self.lock, self.conn:
|
||||||
|
self.conn.execute(
|
||||||
|
"DELETE FROM peers"
|
||||||
|
)
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def _get_input_peer(peer_id: int, access_hash: int, peer_type: str):
|
||||||
|
if peer_type in ["user", "bot"]:
|
||||||
|
return types.InputPeerUser(
|
||||||
|
user_id=peer_id,
|
||||||
|
access_hash=access_hash
|
||||||
|
)
|
||||||
|
|
||||||
|
if peer_type == "group":
|
||||||
|
return types.InputPeerChat(
|
||||||
|
chat_id=-peer_id
|
||||||
|
)
|
||||||
|
|
||||||
|
if peer_type in ["channel", "supergroup"]:
|
||||||
|
return types.InputPeerChannel(
|
||||||
|
channel_id=int(str(peer_id)[4:]),
|
||||||
|
access_hash=access_hash
|
||||||
|
)
|
||||||
|
|
||||||
|
raise ValueError("Invalid peer type")
|
||||||
|
|
||||||
|
def get_peer_by_id(self, peer_id: int):
|
||||||
|
r = self.conn.execute(
|
||||||
|
"SELECT id, access_hash, type FROM peers WHERE id = ?",
|
||||||
|
(peer_id,)
|
||||||
|
).fetchone()
|
||||||
|
|
||||||
|
if r is None:
|
||||||
|
raise KeyError("ID not found")
|
||||||
|
|
||||||
|
return self._get_input_peer(*r)
|
||||||
|
|
||||||
|
def get_peer_by_username(self, username: str):
|
||||||
|
r = self.conn.execute(
|
||||||
|
"SELECT id, access_hash, type, last_update_on FROM peers WHERE username = ?",
|
||||||
|
(username,)
|
||||||
|
).fetchone()
|
||||||
|
|
||||||
|
if r is None:
|
||||||
|
raise KeyError("Username not found")
|
||||||
|
|
||||||
|
if abs(time.time() - r[3]) > self.USERNAME_TTL:
|
||||||
|
raise KeyError("Username expired")
|
||||||
|
|
||||||
|
return self._get_input_peer(*r[:3])
|
||||||
|
|
||||||
|
def get_peer_by_phone_number(self, phone_number: str):
|
||||||
|
r = self.conn.execute(
|
||||||
|
"SELECT id, access_hash, type FROM peers WHERE phone_number = ?",
|
||||||
|
(phone_number,)
|
||||||
|
).fetchone()
|
||||||
|
|
||||||
|
if r is None:
|
||||||
|
raise KeyError("Phone number not found")
|
||||||
|
|
||||||
|
return self._get_input_peer(*r)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def peers_count(self):
|
||||||
|
return self.conn.execute(
|
||||||
|
"SELECT COUNT(*) FROM peers"
|
||||||
|
).fetchone()[0]
|
||||||
|
|
||||||
|
def _get(self):
|
||||||
|
attr = inspect.stack()[1].function
|
||||||
|
|
||||||
|
return self.conn.execute(
|
||||||
|
"SELECT {} FROM sessions".format(attr)
|
||||||
|
).fetchone()[0]
|
||||||
|
|
||||||
|
def _set(self, value):
|
||||||
|
attr = inspect.stack()[1].function
|
||||||
|
|
||||||
|
with self.lock, self.conn:
|
||||||
|
self.conn.execute(
|
||||||
|
"UPDATE sessions SET {} = ?".format(attr),
|
||||||
|
(value,)
|
||||||
|
)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def dc_id(self):
|
||||||
|
return self._get()
|
||||||
|
|
||||||
|
@dc_id.setter
|
||||||
|
def dc_id(self, value):
|
||||||
|
self._set(value)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def test_mode(self):
|
||||||
|
return self._get()
|
||||||
|
|
||||||
|
@test_mode.setter
|
||||||
|
def test_mode(self, value):
|
||||||
|
self._set(value)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def auth_key(self):
|
||||||
|
return self._get()
|
||||||
|
|
||||||
|
@auth_key.setter
|
||||||
|
def auth_key(self, value):
|
||||||
|
self._set(value)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def date(self):
|
||||||
|
return self._get()
|
||||||
|
|
||||||
|
@date.setter
|
||||||
|
def date(self, value):
|
||||||
|
self._set(value)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def user_id(self):
|
||||||
|
return self._get()
|
||||||
|
|
||||||
|
@user_id.setter
|
||||||
|
def user_id(self, value):
|
||||||
|
self._set(value)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def is_bot(self):
|
||||||
|
return self._get()
|
||||||
|
|
||||||
|
@is_bot.setter
|
||||||
|
def is_bot(self, value):
|
||||||
|
self._set(value)
|
34
pyrogram/client/storage/schema.sql
Normal file
34
pyrogram/client/storage/schema.sql
Normal file
@ -0,0 +1,34 @@
|
|||||||
|
CREATE TABLE sessions (
|
||||||
|
dc_id INTEGER PRIMARY KEY,
|
||||||
|
test_mode INTEGER,
|
||||||
|
auth_key BLOB,
|
||||||
|
date INTEGER NOT NULL,
|
||||||
|
user_id INTEGER,
|
||||||
|
is_bot INTEGER
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE TABLE peers (
|
||||||
|
id INTEGER PRIMARY KEY,
|
||||||
|
access_hash INTEGER,
|
||||||
|
type INTEGER NOT NULL,
|
||||||
|
username TEXT,
|
||||||
|
phone_number TEXT,
|
||||||
|
last_update_on INTEGER NOT NULL DEFAULT (CAST(STRFTIME('%s', 'now') AS INTEGER))
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE TABLE version (
|
||||||
|
number INTEGER PRIMARY KEY
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE INDEX idx_peers_id ON peers (id);
|
||||||
|
CREATE INDEX idx_peers_username ON peers (username);
|
||||||
|
CREATE INDEX idx_peers_phone_number ON peers (phone_number);
|
||||||
|
|
||||||
|
CREATE TRIGGER trg_peers_last_update_on
|
||||||
|
AFTER UPDATE
|
||||||
|
ON peers
|
||||||
|
BEGIN
|
||||||
|
UPDATE peers
|
||||||
|
SET last_update_on = CAST(STRFTIME('%s', 'now') AS INTEGER)
|
||||||
|
WHERE id = NEW.id;
|
||||||
|
END;
|
Loading…
Reference in New Issue
Block a user