From 80dfbcc0f2d0e2af564f030b005552d02bff18bb Mon Sep 17 00:00:00 2001 From: Fabio Erculiani Date: Thu, 14 Apr 2011 15:17:54 +0200 Subject: [PATCH] [entropy.services] goodbye old and ugly RPC service, R.I.P. --- client/equo.py | 7 +- client/text_repositories.py | 1 - libraries/entropy/client/interfaces/db.py | 1 - .../entropy/client/services/ugc/__init__.py | 0 .../entropy/client/services/ugc/commands.py | 795 ------ .../entropy/client/services/ugc/interfaces.py | 916 ------- libraries/entropy/services/auth_interfaces.py | 786 ------ libraries/entropy/services/authenticators.py | 115 - libraries/entropy/services/commands.py | 73 - libraries/entropy/services/exceptions.py | 31 - libraries/entropy/services/interfaces.py | 2073 -------------- .../entropy/services/repository/__init__.py | 0 .../entropy/services/repository/commands.py | 368 --- .../entropy/services/repository/interfaces.py | 315 --- libraries/entropy/services/skel.py | 419 --- libraries/entropy/services/test/__init__.py | 0 libraries/entropy/services/test/commands.py | 41 - libraries/entropy/services/test/interfaces.py | 41 - libraries/entropy/services/ugc/__init__.py | 0 libraries/entropy/services/ugc/commands.py | 834 ------ libraries/entropy/services/ugc/interfaces.py | 2407 ----------------- .../tests/standalone/sys_big_send_test.py | 27 - .../test_RemoteDatabase_generate_sql.py | 13 - services/repository-services-daemon.example | 169 -- sulfur/src/sulfur/__init__.py | 1 - sulfur/src/sulfur/dialogs.py | 1 - 26 files changed, 1 insertion(+), 9433 deletions(-) delete mode 100644 libraries/entropy/client/services/ugc/__init__.py delete mode 100644 libraries/entropy/client/services/ugc/commands.py delete mode 100644 libraries/entropy/client/services/ugc/interfaces.py delete mode 100644 libraries/entropy/services/auth_interfaces.py delete mode 100644 libraries/entropy/services/authenticators.py delete mode 100644 libraries/entropy/services/commands.py delete mode 100644 libraries/entropy/services/exceptions.py delete mode 100644 libraries/entropy/services/interfaces.py delete mode 100644 libraries/entropy/services/repository/__init__.py delete mode 100644 libraries/entropy/services/repository/commands.py delete mode 100644 libraries/entropy/services/repository/interfaces.py delete mode 100644 libraries/entropy/services/skel.py delete mode 100644 libraries/entropy/services/test/__init__.py delete mode 100644 libraries/entropy/services/test/commands.py delete mode 100644 libraries/entropy/services/test/interfaces.py delete mode 100644 libraries/entropy/services/ugc/__init__.py delete mode 100644 libraries/entropy/services/ugc/commands.py delete mode 100644 libraries/entropy/services/ugc/interfaces.py delete mode 100644 libraries/tests/standalone/sys_big_send_test.py delete mode 100644 libraries/tests/standalone/test_RemoteDatabase_generate_sql.py delete mode 100755 services/repository-services-daemon.example diff --git a/client/equo.py b/client/equo.py index 1c789143a..04891be50 100644 --- a/client/equo.py +++ b/client/equo.py @@ -25,11 +25,6 @@ sys.path.insert(0, '../client') from entropy.exceptions import SystemDatabaseError, OnlineMirrorError, \ RepositoryError, PermissionDenied, FileNotFound, SPMError -try: - from entropy.services.exceptions import ServiceConnectionError -except ImportError: - # backward compatibility - ServiceConnectionError = None try: from entropy.transceivers.exceptions import TransceiverError, \ TransceiverConnectionError @@ -872,7 +867,7 @@ def handle_exception(exc_class, exc_instance, exc_tb): generic_exc_classes = (OnlineMirrorError, RepositoryError, TransceiverError, PermissionDenied, TransceiverConnectionError, - ServiceConnectionError, FileNotFound, SPMError, SystemError) + FileNotFound, SPMError, SystemError) if exc_class in generic_exc_classes: print_error("%s %s. %s." % ( darkred(" * "), exc_instance, _("Cannot continue"),)) diff --git a/client/text_repositories.py b/client/text_repositories.py index e9df63f43..18c45510b 100644 --- a/client/text_repositories.py +++ b/client/text_repositories.py @@ -18,7 +18,6 @@ import os import sys import time -from entropy.services.exceptions import TimeoutError from entropy.const import etpConst, etpUi, const_debug_write from entropy.output import red, darkred, blue, brown, bold, darkgreen, green, \ print_info, print_warning, print_error, purple, teal diff --git a/libraries/entropy/client/interfaces/db.py b/libraries/entropy/client/interfaces/db.py index 4c79ef484..f5f74ca7b 100644 --- a/libraries/entropy/client/interfaces/db.py +++ b/libraries/entropy/client/interfaces/db.py @@ -26,7 +26,6 @@ from entropy.cache import EntropyCacher from entropy.db import EntropyRepository from entropy.exceptions import RepositoryError, SystemDatabaseError, \ PermissionDenied -from entropy.services.exceptions import EntropyServicesError from entropy.security import Repository as RepositorySecurity from entropy.misc import TimeScheduled from entropy.i18n import _ diff --git a/libraries/entropy/client/services/ugc/__init__.py b/libraries/entropy/client/services/ugc/__init__.py deleted file mode 100644 index e69de29bb..000000000 diff --git a/libraries/entropy/client/services/ugc/commands.py b/libraries/entropy/client/services/ugc/commands.py deleted file mode 100644 index dcc0109a6..000000000 --- a/libraries/entropy/client/services/ugc/commands.py +++ /dev/null @@ -1,795 +0,0 @@ -# -*- coding: utf-8 -*- -""" - - @author: Fabio Erculiani - @contact: lxnay@sabayon.org - @copyright: Fabio Erculiani - @license: GPL-2 - - B{Entropy Client Services UGC Base Commands}. - -""" - -import os -from entropy.services.exceptions import TransmissionError, \ - EntropyServicesError, ServiceConnectionError -from entropy.const import etpConst, const_get_stringtype, const_debug_write, \ - const_convert_to_rawstring -from entropy.output import darkblue, bold, blue, darkgreen, darkred, brown -from entropy.i18n import _ -from entropy.core.settings.base import SystemSettings -import entropy.tools -import entropy.dump - -class Base: - - def __init__(self, OutputInterface, Service): - - if not hasattr(OutputInterface, 'output'): - raise AttributeError( - "OutputInterface does not have an output method") - elif not hasattr(OutputInterface.output, '__call__'): - raise AttributeError( - "OutputInterface does not have an output method") - - from entropy.services.ugc.interfaces import Client as Cl - if not isinstance(Service, Cl): - raise AttributeError( - "entropy.services.ugc.interfaces.Client needed") - - import socket, zlib, struct - self.socket, self.zlib, self.struct = socket, zlib, struct - self.Output = OutputInterface - self.Service = Service - self.output_header = '' - self._settings = SystemSettings() - self.standard_answers_map = { - 'all_fine': 0, - 'not_supported_remotely': 1, - 'service_temp_not_avail': 2, - 'command_failed': 3, - 'wrong_answer': 4, - } - - - def handle_standard_answer(self, data, repository = None, arch = None, - product = None): - - do_skip = False - answer_id = self.standard_answers_map['all_fine'] - - # elaborate answer - if data is None: - mytxt = _("feature not supported remotely") - self.Output.output( - "[%s:%s|%s:%s|%s:%s] %s" % ( - darkblue(_("repo")), - bold(str(repository)), - darkred(_("arch")), - bold(str(arch)), - darkgreen(_("product")), - bold(str(product)), - blue(mytxt), - ), - importance = 1, - level = "error", - header = self.output_header - ) - do_skip = True - answer_id = self.standard_answers_map['not_supported_remotely'] - elif not data: - mytxt = _("service temporarily not available") - self.Output.output( - "[%s:%s|%s:%s|%s:%s] %s" % ( - darkblue(_("repo")), - bold(str(repository)), - darkred(_("arch")), - bold(str(arch)), - darkgreen(_("product")), - bold(str(product)), - blue(mytxt), - ), - importance = 1, - level = "error", - header = self.output_header - ) - do_skip = True - answer_id = self.standard_answers_map['service_temp_not_avail'] - elif data == self.Service.answers['no']: - # command failed - mytxt = _("command failed") - self.Output.output( - "[%s:%s|%s:%s|%s:%s] %s" % ( - darkblue(_("repo")), - bold(str(repository)), - darkred(_("arch")), - bold(str(arch)), - darkgreen(_("product")), - bold(str(product)), - blue(mytxt), - ), - importance = 1, - level = "error", - header = self.output_header - ) - do_skip = True - answer_id = self.standard_answers_map['command_failed'] - elif data != self.Service.answers['ok']: - mytxt = _("received wrong answer") - - # do not spam terminal - if isinstance(data, const_get_stringtype()): - if len(data) > 10: - data = data[:10] + "[...]" - - self.Output.output( - "[%s:%s|%s:%s|%s:%s] %s: %s" % ( - darkblue(_("repo")), - bold(str(repository)), - darkred(_("arch")), - bold(str(arch)), - darkgreen(_("product")), - bold(str(product)), - blue(mytxt), - repr(data), - ), - importance = 1, - level = "error", - header = self.output_header - ) - do_skip = True - answer_id = self.standard_answers_map['wrong_answer'] - - return do_skip, answer_id - - def get_result(self, session): - # get the information - cmd = "%s rc" % (session,) - try: - self.Service.transmit(cmd) - except TransmissionError: - entropy.tools.print_traceback() - return None - try: - data = self.Service.receive() - return data - except Exception: - entropy.tools.print_traceback() - return None - - def convert_stream_to_object(self, data, gzipped, repository = None, - arch = None, product = None): - - # unstream object - error = False - try: - data = self.Service.stream_to_object(data, gzipped) - except (EOFError, IOError, self.zlib.error, entropy.dump.pickle.UnpicklingError,): - const_debug_write(__name__, entropy.tools.get_traceback()) - mytxt = _("cannot convert stream into object") - self.Output.output( - "[%s:%s|%s:%s|%s:%s] %s" % ( - darkblue(_("repo")), - bold(str(repository)), - darkred(_("arch")), - bold(str(arch)), - darkgreen(_("product")), - bold(str(product)), - blue(mytxt), - ), - importance = 1, - level = "error", - header = self.output_header - ) - data = None - error = True - return data, error - - def retrieve_command_answer(self, cmd, session_id, repository = None, - arch = None, product = None, compression = False): - - tries = 3 - lasterr = None - while True: - - if tries <= 0: - return lasterr - tries -= 1 - - try: - # send command - self.Service.transmit(cmd) - except TransmissionError: - return None - # receive answer - data = self.Service.receive() - - skip, answer_id = self.handle_standard_answer(data, repository, - arch, product) - if skip: - if tries <= 0: - const_debug_write(__name__, - darkred("skipping command, NOT reconnecting")) - else: - const_debug_write(__name__, - darkred("skipping command, reconnect+retry!")) - const_debug_write(__name__, str(answer_id)) - # reconnect host and retry - self.Service.reconnect_socket() - continue - - data = self.get_result(session_id) - if data is None: - lasterr = None - continue - elif not data: - lasterr = False - continue - - objdata, error = self.convert_stream_to_object(data, compression, - repository, arch, product) - if not error: - return objdata - - def do_generic_handler(self, cmd, session_id, tries = 10, compression = False): - - try: - self.Service.check_socket_connection() - except ServiceConnectionError: - return False, 'connection error' - - while True: - try: - result = self.retrieve_command_answer(cmd, session_id, - compression = compression) - if result is None: - return False, 'command not supported' # untranslated on purpose - return result - except (self.socket.error, self.struct.error,) as err: - try: - self.Service.reconnect_socket() - except self.socket.error as exc: - return False, 'connection error %s' % (exc,) - tries -= 1 - if tries < 1: - return False, 'connection error %s' % (err,) - - def _set_gzip_compression(self, session, do): - self.Service.check_socket_connection() - cmd = "%s %s %s %s zlib" % ( - const_convert_to_rawstring(session), - const_convert_to_rawstring('session_config'), - const_convert_to_rawstring('compression'), - const_convert_to_rawstring(do), - ) - fail_count = 5 - while True: - try: - self.Service.transmit(cmd) - data = self.Service.receive() - except TransmissionError as err: - const_debug_write(__name__, - darkred("_set_gzip_compression: error: " + repr(err))) - fail_count -= 1 - if fail_count == 0: - raise - continue - if data == self.Service.answers['ok']: - return True - return False - - def service_login(self, username, password, session_id): - - cmd = "%s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('login'), - const_convert_to_rawstring(username), - const_convert_to_rawstring(password), - ) - return self.do_generic_handler(cmd, session_id) - - def service_logout(self, username, session_id): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('logout'), - const_convert_to_rawstring(username), - ) - return self.do_generic_handler(cmd, session_id) - - def get_logged_user_data(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('user_data'), - ) - return self.do_generic_handler(cmd, session_id) - - def is_user(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('is_user'), - ) - return self.do_generic_handler(cmd, session_id) - - def is_developer(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('is_developer'), - ) - return self.do_generic_handler(cmd, session_id) - - def is_moderator(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('is_moderator'), - ) - return self.do_generic_handler(cmd, session_id) - - def is_administrator(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('is_administrator'), - ) - return self.do_generic_handler(cmd, session_id) - - def available_commands(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('available_commands'), - ) - return self.do_generic_handler(cmd, session_id) - - -class Client(Base): - - def __init__(self, EntropyInterface, ServiceInterface): - Base.__init__(self, EntropyInterface, ServiceInterface) - - def differential_packages_comparison(self, session_id, idpackages, - repository, arch, product): - - myidlist = const_convert_to_rawstring( - ' '.join([str(x) for x in idpackages])) - cmd = "%s %s %s %s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('repository_server:dbdiff'), - const_convert_to_rawstring(repository), - const_convert_to_rawstring(arch), - const_convert_to_rawstring(product), - const_convert_to_rawstring( - self._settings['repositories']['branch']), - myidlist, - ) - - # enable zlib compression - try: - compression = self._set_gzip_compression(session_id, True) - except EntropyServicesError: - return False, 'connection error' - - data = self.do_generic_handler(cmd, session_id, tries = 5, - compression = compression) - - # disable compression - try: - compression = self._set_gzip_compression(session_id, False) - except EntropyServicesError: - return False, 'connection error' - - return data - - def get_repository_treeupdates(self, session_id, repository, arch, product): - - cmd = "%s %s %s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('repository_server:treeupdates'), - const_convert_to_rawstring(repository), - const_convert_to_rawstring(arch), - const_convert_to_rawstring(product), - const_convert_to_rawstring( - self._settings['repositories']['branch']), - ) - return self.do_generic_handler(cmd, session_id, tries = 5) - - def get_package_sets(self, session_id, repository, arch, product): - - cmd = "%s %s %s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('repository_server:get_package_sets'), - const_convert_to_rawstring(repository), - const_convert_to_rawstring(arch), - const_convert_to_rawstring(product), - const_convert_to_rawstring( - self._settings['repositories']['branch']), - ) - return self.do_generic_handler(cmd, session_id, tries = 5) - - def get_repository_metadata(self, session_id, repository, arch, product): - - cmd = "%s %s %s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring( - 'repository_server:get_repository_metadata'), - const_convert_to_rawstring(repository), - const_convert_to_rawstring(arch), - const_convert_to_rawstring(product), - const_convert_to_rawstring( - self._settings['repositories']['branch']), - ) - return self.do_generic_handler(cmd, session_id, tries = 5) - - def get_strict_package_information(self, session_id, idpackages, - repository, arch, product): - - cmd = "%s %s %s %s %s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('repository_server:pkginfo_strict'), - True, - const_convert_to_rawstring(repository), - const_convert_to_rawstring(arch), - const_convert_to_rawstring(product), - const_convert_to_rawstring( - self._settings['repositories']['branch']), - const_convert_to_rawstring(' '.join([str(x) for x in idpackages])), - ) - - # enable zlib compression - try: - compression = self._set_gzip_compression(session_id, True) - except EntropyServicesError: - return False, 'connection error' - - data = self.do_generic_handler(cmd, session_id, compression = compression) - - # disable compression - try: - compression = self._set_gzip_compression(session_id, False) - except EntropyServicesError: - return False, 'connection error' - - return data - - def ugc_do_download_stats(self, session_id, package_names): - - sub_lists = entropy.tools.split_indexable_into_chunks( - package_names, 100) - - last_srv_rc_data = None - for pkgkeys in sub_lists: - - release_string = '--N/A--' - rel_file = etpConst['systemreleasefile'] - if os.path.isfile(rel_file) and os.access(rel_file, os.R_OK): - with open(rel_file, "r") as f: - release_string = f.read(512) - - hw_hash = self._settings['hw_hash'] - if not hw_hash: - hw_hash = '' - - mydict = { - 'branch': self._settings['repositories']['branch'], - 'release_string': release_string, - 'hw_hash': hw_hash, - 'pkgkeys': ' '.join(pkgkeys), - } - xml_string = entropy.tools.xml_from_dict(mydict) - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:do_download_stats'), - const_convert_to_rawstring(xml_string), - ) - last_srv_rc_data = self.do_generic_handler(cmd, session_id) - if not isinstance(last_srv_rc_data, tuple): - return last_srv_rc_data - elif last_srv_rc_data[0] != True: - return last_srv_rc_data - return last_srv_rc_data - - def ugc_get_downloads(self, session_id, pkgkey): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_downloads'), - pkgkey, - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_get_alldownloads(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_alldownloads'), - ) - - # enable zlib compression - try: - compression = self._set_gzip_compression(session_id, True) - except EntropyServicesError: - return False, 'connection error' - - rc = self.do_generic_handler(cmd, session_id, compression = compression) - - # disable compression - try: - compression = self._set_gzip_compression(session_id, False) - except EntropyServicesError: - return False, 'connection error' - - return rc - - def ugc_do_vote(self, session_id, pkgkey, vote): - - cmd = "%s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:do_vote'), - const_convert_to_rawstring(pkgkey), - const_convert_to_rawstring(vote), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_get_vote(self, session_id, pkgkey): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_vote'), - const_convert_to_rawstring(pkgkey), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_get_allvotes(self, session_id): - - cmd = "%s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_allvotes'), - ) - - # enable zlib compression - try: - compression = self._set_gzip_compression(session_id, True) - except EntropyServicesError: - return False, 'connection error' - - rc = self.do_generic_handler(cmd, session_id, compression = compression) - - # disable compression - try: - compression = self._set_gzip_compression(session_id, False) - except EntropyServicesError: - return False, 'connection error' - - return rc - - def ugc_add_comment(self, session_id, pkgkey, comment, title, keywords): - - mydict = { - 'comment': comment, - 'title': title, - 'keywords': keywords, - } - xml_string = entropy.tools.xml_from_dict(mydict) - - cmd = "%s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:add_comment'), - const_convert_to_rawstring(pkgkey), - const_convert_to_rawstring(xml_string), - ) - - return self.do_generic_handler(cmd, session_id) - - def ugc_edit_comment(self, session_id, iddoc, new_comment, new_title, new_keywords): - - mydict = { - 'comment': new_comment, - 'title': new_title, - 'keywords': new_keywords, - } - xml_string = entropy.tools.xml_from_dict(mydict) - - cmd = "%s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:edit_comment'), - const_convert_to_rawstring(iddoc), - const_convert_to_rawstring(xml_string), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_remove_comment(self, session_id, iddoc): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:remove_comment'), - const_convert_to_rawstring(iddoc), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_remove_image(self, session_id, iddoc): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:remove_image'), - const_convert_to_rawstring(iddoc), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_remove_file(self, session_id, iddoc): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:remove_file'), - const_convert_to_rawstring(iddoc), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_remove_youtube_video(self, session_id, iddoc): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:remove_youtube_video'), - const_convert_to_rawstring(iddoc), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_get_docs(self, session_id, pkgkey): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_alldocs'), - const_convert_to_rawstring(pkgkey), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_get_textdocs(self, session_id, pkgkey): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_textdocs'), - const_convert_to_rawstring(pkgkey), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_get_textdocs_by_identifiers(self, session_id, identifiers): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_textdocs_by_identifiers'), - const_convert_to_rawstring(' '.join([str(x) for x in identifiers])), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_get_documents_by_identifiers(self, session_id, identifiers): - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:get_documents_by_identifiers'), - const_convert_to_rawstring(' '.join([str(x) for x in identifiers])), - ) - return self.do_generic_handler(cmd, session_id) - - def ugc_send_file_stream(self, session_id, file_path): - - if not (os.path.isfile(file_path) and os.access(file_path, os.R_OK)): - return False, False, 'cannot read file_path' - - import zlib - # enable stream - cmd = "%s %s %s on" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('session_config'), - const_convert_to_rawstring('stream'), - ) - status, msg = self.do_generic_handler(cmd, session_id) - if not status: - return False, status, msg - - # enable zlib compression - try: - compression = self._set_gzip_compression(session_id, True) - except EntropyServicesError: - return False, 'connection error' - - # start streamer - stream_status = True - stream_msg = 'ok' - f = open(file_path, "rb") - chunk = f.read(8192) - base_path = os.path.basename(file_path) - transferred = len(chunk) - max_size = entropy.tools.get_file_size(file_path) - while chunk: - - if (not self.Service.quiet) or self.Service.show_progress: - self.Output.output( - "%s, %s: %s" % ( - blue(_("User Generated Content")), - darkgreen(_("sending file")), - darkred(base_path), - ), - importance = 1, - level = "info", - header = brown(" @@ "), - back = True, - count = (transferred, max_size,), - percent = True - ) - - chunk = zlib.compress(chunk, 7) # compression level 1-9 - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - 'stream', - chunk, - ) - status, msg = self.do_generic_handler(cmd, session_id, compression = compression) - if not status: - stream_status = status - stream_msg = msg - break - chunk = f.read(8192) - transferred += len(chunk) - - f.close() - - # disable compression - try: - compression = self._set_gzip_compression(session_id, False) - except EntropyServicesError: - return False, 'connection error' - - # disable config - cmd = "%s %s %s off" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('session_config'), - const_convert_to_rawstring('stream'), - ) - status, msg = self.do_generic_handler(cmd, session_id) - if not status: - return False, status, msg - - return True, stream_status, stream_msg - - def ugc_send_file(self, session_id, pkgkey, file_path, doc_type, title, - description, keywords): - - status, rem_status, err_msg = self.ugc_send_file_stream(session_id, file_path) - if not (status and rem_status): - return False, err_msg - - mydict = { - 'doc_type': str(doc_type), - 'title': title, - 'description': description, - 'keywords': keywords, - 'file_name': os.path.join(pkgkey, os.path.basename(file_path)), - 'real_filename': os.path.basename(file_path), - } - xml_string = entropy.tools.xml_from_dict(mydict) - - cmd = "%s %s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:register_stream'), - pkgkey, - xml_string, - ) - return self.do_generic_handler(cmd, session_id) - - def report_error(self, session_id, error_data): - - import zlib - xml_string = entropy.tools.xml_from_dict_extended(error_data) - xml_comp_string = zlib.compress(xml_string) - - cmd = "%s %s %s" % ( - const_convert_to_rawstring(session_id), - const_convert_to_rawstring('ugc:report_error'), - const_convert_to_rawstring(xml_comp_string), - ) - - return self.do_generic_handler(cmd, session_id) diff --git a/libraries/entropy/client/services/ugc/interfaces.py b/libraries/entropy/client/services/ugc/interfaces.py deleted file mode 100644 index c3b03d0b2..000000000 --- a/libraries/entropy/client/services/ugc/interfaces.py +++ /dev/null @@ -1,916 +0,0 @@ -# -*- coding: utf-8 -*- -""" - - @author: Fabio Erculiani - @contact: lxnay@sabayon.org - @copyright: Fabio Erculiani - @license: GPL-2 - - B{Entropy Client Services UGC Base Interfaces}. - -""" - -import os -import shutil -from entropy.core import Singleton -from entropy.core.settings.base import SystemSettings -from entropy.exceptions import RepositoryError, PermissionDenied -from entropy.services.exceptions import ServiceConnectionError, \ - EntropyServicesError -from entropy.cache import MtimePingus, EntropyCacher -from entropy.const import etpConst, const_setup_file, const_setup_perms -from entropy.i18n import _ - -import entropy.dump - -class Client: - - ssl_connection = True - def __init__(self, entropy_client, quiet = True, show_progress = False): - - from entropy.client.interfaces import Client as Cl - if not isinstance(entropy_client, Cl): - mytxt = "A valid Client based instance is needed" - raise AttributeError(mytxt) - - import socket, threading - self.socket, self.threading = socket, threading - import struct - self.struct = struct - self.Entropy = entropy_client - self.store = AuthStore() - self.quiet = quiet - self.show_progress = show_progress - self.UGCCache = Cache(self) - self.TxLocks = {} - self._cacher = EntropyCacher() - - def connect_to_service(self, repository, timeout = None): - - sys_settings = SystemSettings() - avail_data = sys_settings['repositories']['available'] - # also try excluded repos - if repository not in avail_data: - avail_data = sys_settings['repositories']['excluded'] - if repository not in avail_data: - raise RepositoryError('repository is not available') - - # unsupported by repository? - if 'service_uri' not in avail_data[repository]: - return None - if 'service_port' not in avail_data[repository]: - return None - - url = avail_data[repository]['service_uri'] - port = avail_data[repository]['service_port'] - if self.ssl_connection: - port = avail_data[repository]['ssl_service_port'] - - from entropy.services.ugc.interfaces import Client - from entropy.client.services.ugc.commands import Client as CommandsClient - args = [self.Entropy, CommandsClient] - kwargs = { - 'ssl': self.ssl_connection, - 'quiet': self.quiet, - 'show_progress': self.show_progress - } - if timeout is not None: - kwargs['socket_timeout'] = timeout - - srv = Client(*args, **kwargs) - srv.connect(url, port) - return srv - - def get_service_connection(self, repository, check = True, timeout = None): - if check: - if not self.is_repository_eapi3_aware(repository): - return None - try: - srv = self.connect_to_service(repository, timeout = timeout) - except (RepositoryError, ServiceConnectionError,): - return None - return srv - - def is_repository_eapi3_aware(self, repository, _use_cache = True): - - aware = self.UGCCache._get_live_cache_item(repository, - 'is_repository_eapi3_aware') - if aware is not None: - return aware - - def _get_awareness(): - try: - srv = self.get_service_connection(repository, check = False, - timeout = 6) - if srv is None: - rc = False - else: - session = srv.open_session() - if session is not None: - srv.close_session(session) - srv.disconnect() - rc = True - else: - rc = False - except EntropyServicesError: - rc = False - return rc - - # Life is hard, and socket communication can be very annoying - # over a non-performant connection. So, the only way to circumvent this - # is to cache results somewhere. - cache_id = "entropy.client.interfaces.ugc.is_repository_eapi3_aware_" \ - + repository - pingus = MtimePingus() - - passed = True - if _use_cache: - # are 3 days passed? - passed = pingus.hours_passed(cache_id, 24*3) - - if passed: - aware = _get_awareness() - # update pingus mtime - pingus.ping(cache_id) - else: - # then load data from disk cache - cached = self._cacher.pop(cache_id) - if cached is None: - # no cache on disk - aware = _get_awareness() - self._cacher.push(cache_id, aware, async = False) - pingus.ping(cache_id) - else: - aware = cached - - self.UGCCache._set_live_cache_item(repository, - 'is_repository_eapi3_aware', aware) - return aware - - def read_login(self, repository): - return self.store.read_login(repository) - - def remove_login(self, repository): - return self.store.remove_login(repository) - - def do_login(self, repository, force = False): - - login_data = self.read_login(repository) - if (login_data is not None) and not force: - return True, _('ok') - - aware = self.is_repository_eapi3_aware(repository) - if not aware: - return False, _('repository does not support EAPI3') - - def fake_callback(*args, **kwargs): - return True - - attempts = 3 - while attempts: - - # use input box to read login - input_params = [ - ('username', _('Username'), fake_callback, False), - ('password', _('Password'), fake_callback, True) - ] - login_data = self.Entropy.input_box( - "%s %s %s" % ( - _('Please login against'), repository, _('repository'),), - input_params, - cancel_button = True - ) - if not login_data: - return False, _('login abort') - - # now verify - srv = self.get_service_connection(repository) - if srv is None: - return False, _('connection issues') - session = srv.open_session() - if session is None: - return False, _('cannot open a session') - login_status, login_msg = srv.CmdInterface.service_login( - login_data['username'], login_data['password'], session) - if not login_status: - srv.close_session(session) - srv.disconnect() - self.Entropy.ask_question("%s: %s" % ( - _("Access denied. Login failed"), login_msg,), - responses = [_("Ok")]) - attempts -= 1 - continue - - # login accepted, store it? - srv.close_session(session) - srv.disconnect() - rc = self.Entropy.ask_question( - _("Login successful. Do you want to save these credentials ?")) - save = False - if rc == _("Yes"): - save = True - self.store.store_login(login_data['username'], - login_data['password'], repository, save = save) - return True, _('ok') - - - def login(self, repository, force = False): - - if repository not in self.TxLocks: - self.TxLocks[repository] = self.threading.Lock() - - with self.TxLocks[repository]: - return self.do_login(repository, force = force) - - def logout(self, repository): - return self.store.remove_login(repository) - - def do_cmd(self, repository, login_required, func, args, kwargs): - - if repository not in self.TxLocks: - self.TxLocks[repository] = self.threading.Lock() - - with self.TxLocks[repository]: - - if login_required: - status, err_msg = self.do_login(repository) - if not status: - return False, err_msg - - srv = self.get_service_connection(repository) - if srv is None: - return False, 'no connection' - session = srv.open_session() - if session is None: - return False, 'no session' - args.insert(0, session) - - if login_required: - stored_pass = False - while True: - # login - login_data = self.read_login(repository) - if login_data is None: - status, msg = self.login(repository) - if not status: - return status, msg - username, password = self.read_login(repository) - else: - stored_pass = True - username, password = login_data - logged, error = srv.CmdInterface.service_login(username, - password, session) - if not logged: - if stored_pass: - stored_pass = False - self.remove_login(repository) - continue - srv.close_session(session) - srv.disconnect() - return logged, error - break - - try: - cmd_func = getattr(srv.CmdInterface, func) - except AttributeError: - return False, 'local function not available' - rslt = cmd_func(*args, **kwargs) - try: - srv.close_session(session) - srv.disconnect() - except ServiceConnectionError: - return False, 'no connection' - - return rslt - - def get_comments(self, repository, pkgkey): - return self.do_cmd(repository, False, "ugc_get_textdocs", [pkgkey], {}) - - def get_comments_by_identifiers(self, repository, identifiers): - return self.do_cmd(repository, False, - "ugc_get_textdocs_by_identifiers", [identifiers], {}) - - def get_documents_by_identifiers(self, repository, identifiers): - return self.do_cmd(repository, False, - "ugc_get_documents_by_identifiers", [identifiers], {}) - - def add_comment(self, repository, pkgkey, comment, title, keywords): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_add_comment", - [pkgkey, comment, title, keywords], {}) - - def edit_comment(self, repository, iddoc, new_comment, new_title, new_keywords): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_edit_comment", - [iddoc, new_comment, new_title, new_keywords], {}) - - def remove_comment(self, repository, iddoc): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_remove_comment", [iddoc], {}) - - def add_vote(self, repository, pkgkey, vote): - data = self.do_cmd(repository, True, "ugc_do_vote", [pkgkey, vote], {}) - if isinstance(data, tuple): - voted, add_err_msg = data - else: - return False, 'wrong server answer' - if voted: - self.get_vote(repository, pkgkey) - return voted, add_err_msg - - def get_vote(self, repository, pkgkey): - vote, err_msg = self.do_cmd(repository, False, "ugc_get_vote", [pkgkey], {}) - if isinstance(vote, float): - mydict = {pkgkey: vote} - self.UGCCache.update_vote_cache(repository, mydict) - return vote, err_msg - - def get_all_votes(self, repository): - votes_dict, err_msg = self.do_cmd(repository, False, "ugc_get_allvotes", [], {}) - if isinstance(votes_dict, dict): - self.UGCCache.update_vote_cache(repository, votes_dict) - return votes_dict, err_msg - - def get_downloads(self, repository, pkgkey): - data = self.do_cmd(repository, False, "ugc_get_downloads", [pkgkey], {}) - if isinstance(data, tuple): - downloads, err_msg = data - else: - return False, 'wrong server answer' - if downloads: - mydict = {pkgkey: downloads} - self.UGCCache.update_downloads_cache(repository, mydict) - return downloads, err_msg - - def get_all_downloads(self, repository): - down_dict, err_msg = self.do_cmd(repository, False, "ugc_get_alldownloads", [], {}) - if isinstance(down_dict, dict): - self.UGCCache.update_downloads_cache(repository, down_dict) - return down_dict, err_msg - - def add_download_stats(self, repository, pkgkeys): - return self.do_cmd(repository, False, "ugc_do_download_stats", [pkgkeys], {}) - - def send_file(self, repository, pkgkey, file_path, title, description, keywords): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_send_file", - [pkgkey, file_path, etpConst['ugc_doctypes']['generic_file'], title, description, keywords], {}) - - def remove_file(self, repository, iddoc): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_remove_file", [iddoc], {}) - - def send_image(self, repository, pkgkey, image_path, title, description, keywords): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_send_file", - [pkgkey, image_path, etpConst['ugc_doctypes']['image'], title, description, keywords], {}) - - def remove_image(self, repository, iddoc): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_remove_image", [iddoc], {}) - - def send_youtube_video(self, repository, pkgkey, video_path, title, description, keywords): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_send_file", - [pkgkey, video_path, etpConst['ugc_doctypes']['youtube_video'], title, description, keywords], {}) - - def remove_youtube_video(self, repository, iddoc): - self.UGCCache.clear_alldocs_cache(repository) - return self.do_cmd(repository, True, "ugc_remove_youtube_video", [iddoc], {}) - - def get_docs(self, repository, pkgkey): - data = self.do_cmd(repository, False, "ugc_get_docs", [pkgkey], {}) - if isinstance(data, tuple): - docs_data, err_msg = data - else: - return False, 'wrong server answer' - self.UGCCache.save_alldocs_cache(pkgkey, repository, docs_data) - return docs_data, err_msg - - def get_icon(self, repository, pkgkey): - docs_data, err_msg = self.get_docs(repository, pkgkey) - if not docs_data: - return None - elif not isinstance(docs_data, (tuple, list)): - return None - return self.UGCCache.get_icon_cache(pkgkey, repository) - - def send_document_autosense(self, repository, pkgkey, ugc_type, data, - title, description, keywords): - - if ugc_type == etpConst['ugc_doctypes']['generic_file']: - return self.send_file(repository, pkgkey, data, title, - description, keywords) - elif ugc_type == etpConst['ugc_doctypes']['image']: - return self.send_image(repository, pkgkey, data, title, - description, keywords) - elif ugc_type == etpConst['ugc_doctypes']['youtube_video']: - return self.send_youtube_video(repository, pkgkey, data, title, - description, keywords) - elif ugc_type == etpConst['ugc_doctypes']['comments']: - return self.add_comment(repository, pkgkey, - description, title, keywords) - - return None, 'type not supported locally' - - def remove_document_autosense(self, repository, iddoc, ugc_type): - if ugc_type == etpConst['ugc_doctypes']['generic_file']: - return self.remove_file(repository, iddoc) - elif ugc_type == etpConst['ugc_doctypes']['image']: - return self.remove_image(repository, iddoc) - elif ugc_type == etpConst['ugc_doctypes']['youtube_video']: - return self.remove_youtube_video(repository, iddoc) - elif ugc_type == etpConst['ugc_doctypes']['comments']: - return self.remove_comment(repository, iddoc) - return None, 'type not supported locally' - - def report_error(self, repository, error_data): - return self.do_cmd(repository, False, "report_error", [error_data], {}) - - -class AuthStore(Singleton): - - access_file = etpConst['ugc_accessfile'] - def init_singleton(self): - - from xml.dom import minidom - from xml.parsers import expat - self.expat = expat - self.minidom = minidom - self.setup_store_paths() - try: - self.setup_permissions() - except IOError: - pass - self.store = {} - try: - self.xmldoc = self.minidom.parse(self.access_file) - except (self.expat.ExpatError, IOError, ValueError,): - # ValueError: bad marshal data, wtf!? - self.xmldoc = None - if self.xmldoc is not None: - try: - self.parse_document() - except self.expat.ExpatError: - self.xmldoc = None - self.store = {} - - def setup_store_paths(self): - myhome = os.getenv("HOME") - if myhome is not None: - if os.path.isdir(myhome) and os.access(myhome, os.W_OK): - self.access_file = os.path.join(myhome, ".config/entropy", - os.path.basename(self.access_file)) - self.access_dir = os.path.dirname(self.access_file) - - def setup_permissions(self): - if not os.path.isdir(self.access_dir): - os.makedirs(self.access_dir) - if not os.path.isfile(self.access_file): - f = open(self.access_file, "w") - f.close() - gid = etpConst['entropygid'] - if gid is None: - gid = 0 - - try: - const_setup_file(self.access_dir, gid, 0o700) - except OSError: - pass - try: - const_setup_file(self.access_file, gid, 0o600) - except OSError: - pass - - def parse_document(self): - self.store.clear() - store = self.xmldoc.getElementsByTagName("store")[0] - repositories = store.getElementsByTagName("repository") - for repository in repositories: - repoid = repository.getAttribute("id") - if not repoid: - continue - username = repository.getElementsByTagName("username")[0].firstChild.data.strip() - password = repository.getElementsByTagName("password")[0].firstChild.data.strip() - self.store[repoid] = {'username': username, 'password': password} - - def store_login(self, username, password, repository, save = True): - self.store[repository] = {'username': username, 'password': password} - if save: - self.save_store() - - def save_store(self): - - self.xmldoc = self.minidom.Document() - store = self.xmldoc.createElement("store") - - for repository in self.store: - repo = self.xmldoc.createElement("repository") - repo.setAttribute('id', repository) - # username - username = self.xmldoc.createElement("username") - username_value = self.xmldoc.createTextNode( - self.store[repository]['username']) - username.appendChild(username_value) - repo.appendChild(username) - # password - password = self.xmldoc.createElement("password") - password_value = self.xmldoc.createTextNode( - self.store[repository]['password']) - password.appendChild(password_value) - repo.appendChild(password) - store.appendChild(repo) - - self.xmldoc.appendChild(store) - f = None - try: - f = open(self.access_file, "w") - f.writelines(self.xmldoc.toprettyxml(indent=" ")) - f.flush() - self.setup_permissions() - except IOError: - # no permissions? - pass - finally: - if f is not None: - try: - f.close() - except IOError: - pass - - self.parse_document() - - def remove_login(self, repository, save = True): - if repository in self.store: - del self.store[repository] - if save: - self.save_store() - - def read_login(self, repository): - data = self.store.get(repository) - if data is None: - return None - return data['username'], data['password'] - - -class Cache: - - CACHE_KEYS = { - 'ugc_votes': 'ugc_votes', - 'ugc_downloads': 'ugc_downloads', - 'ugc_docs': 'ugc_docs', - } - CACHE_DIR = os.path.join(etpConst['entropyworkdir'], "ugc_cache") - PKG_ICON_IDENTIFIER = "__icon__" - - def __init__(self, UGCClientInstance): - - if not isinstance(UGCClientInstance, Client): - raise AttributeError( - "entropy.client.services.ugc.Client instance required") - - import threading - self.CacheLock = threading.Lock() - self.Service = UGCClientInstance - self.xcache = {} - - def _get_live_cache_item(self, repository, item): - if repository not in self.xcache: - return None - return self.xcache[repository].get(item) - - def _is_live_cache_item_available(self, repository, item): - if repository not in self.xcache: - return False - return item in self.xcache[repository] - - def _set_live_cache_item(self, repository, item, obj): - if repository not in self.xcache: - self.xcache[repository] = {} - if type(obj) in (list, tuple,): - my_obj = obj[:] - elif type(obj) in (set, frozenset, dict,): - my_obj = obj.copy() - else: - my_obj = obj - self.xcache[repository][item] = my_obj - - def _clear_live_cache_item(self, repository, item): - if repository not in self.xcache: - return - if item not in self.xcache[repository]: - return - del self.xcache[repository][item] - - def _get_store_cache_file(self, iddoc, repository, doc_url, - setup_dir = False): - base_dir = os.path.join(Cache.CACHE_DIR, - Cache.CACHE_KEYS['ugc_docs'], repository, str(iddoc)) - cache_path = os.path.join(base_dir, "_"+doc_url) - cache_dir = os.path.dirname(cache_path) - - if setup_dir and (not os.path.isdir(cache_dir)): - try: - os.makedirs(cache_dir, 0o775) - const_setup_perms(cache_dir, etpConst['entropygid']) - except (OSError, IOError,): - return cache_path - - return cache_path - - def _get_vote_cache_file(self, repository): - return os.path.join(Cache.CACHE_DIR, - self._get_vote_cache_dir(repository)) - - def _get_downloads_cache_file(self, repository): - return os.path.join(Cache.CACHE_DIR, - self._get_downloads_cache_dir(repository)) - - def _get_alldocs_cache_file(self, pkgkey, repository): - return os.path.join(Cache.CACHE_DIR, - self._get_alldocs_cache_dir(repository)+"/"+pkgkey) - - def _get_alldocs_cache_dir(self, repository): - return Cache.CACHE_KEYS['ugc_docs']+"/"+repository - - def _get_downloads_cache_dir(self, repository): - return Cache.CACHE_KEYS['ugc_downloads']+"/"+repository - - def _get_vote_cache_dir(self, repository): - return Cache.CACHE_KEYS['ugc_votes']+"/"+repository - - def _get_vote_cache_key(self, repository): - return 'get_vote_cache_'+repository - - def _get_downloads_cache_key(self, repository): - return 'get_downloads_cache_'+repository - - def _get_alldocs_cache_key(self, repository): - return 'get_package_alldocs_cache_'+repository - - def _clear_cache_dir(self, cache_dir): - try: - shutil.rmtree(cache_dir, True) - except (shutil.Error, OSError, IOError,): - return - - def _complete_cache_path(self, cache_file): - return os.path.join(Cache.CACHE_DIR, cache_file) - - def store_document(self, iddoc, repository, doc_url): - cache_file = self._get_store_cache_file(iddoc, repository, doc_url, - setup_dir = True) - cache_dir = os.path.dirname(cache_file) - - if not os.access(cache_dir, os.W_OK): - raise PermissionDenied("PermissionDenied: %s %s" % ( - _("Cannot write to cache directory"), cache_dir,)) - - if os.path.isfile(cache_file) or os.path.islink(cache_file): - try: - os.remove(cache_file) - except OSError: - raise PermissionDenied("PermissionDenied: %s %s" % ( - _("Cannot remove cache file"), cache_file,)) - - fetcher = self.Service.Entropy._url_fetcher(doc_url, cache_file, - resume = False) - rc = fetcher.download() - if rc in ("-1", "-2", "-3", "-4"): - return None - if not os.path.isfile(cache_file): - return None - - try: - os.chmod(cache_file, 0o664) - if etpConst['entropygid'] is not None: - os.chown(cache_file, -1, etpConst['entropygid']) - except OSError: - raise PermissionDenied("PermissionDenied: %s %s" % ( - _("Cannot write to cache file"), cache_file,)) - - del fetcher - return cache_file - - def get_stored_document(self, iddoc, repository, doc_url): - cache_file = self._get_store_cache_file(iddoc, repository, doc_url) - if os.path.isfile(cache_file) and os.access(cache_file, os.R_OK): - return cache_file - - def update_vote_cache(self, repository, vote_dict): - cached = self.get_vote_cache(repository) - if cached is None: - cached = vote_dict.copy() - else: - cached.update(vote_dict) - self.save_vote_cache(repository, cached) - - def update_downloads_cache(self, repository, down_dict): - cached = self.get_downloads_cache(repository) - if cached is None: - cached = down_dict.copy() - else: - cached.update(down_dict) - self.save_downloads_cache(repository, cached) - - def clear_alldocs_cache(self, repository): - with self.CacheLock: - self._clear_cache_dir( - self._get_alldocs_cache_dir(repository)) - self._clear_live_cache_item(repository, - self._get_alldocs_cache_key(repository)) - - def clear_downloads_cache(self, repository): - with self.CacheLock: - self._clear_cache_dir( - self._get_downloads_cache_dir(repository)) - self._clear_live_cache_item(repository, - self._get_downloads_cache_key(repository)) - - def clear_vote_cache(self, repository): - with self.CacheLock: - self._clear_cache_dir( - self._get_vote_cache_dir(repository)) - self._clear_live_cache_item(repository, - self._get_vote_cache_key(repository)) - - def clear_live_cache(self): - self.xcache.clear() - - def clear_cache(self, repository): - self.clear_alldocs_cache(repository) - self.clear_downloads_cache(repository) - self.clear_vote_cache(repository) - self.xcache.clear() - - def get_vote_cache(self, repository): - cache_key = self._get_vote_cache_key(repository) - cached = self._get_live_cache_item(repository, cache_key) - if cached is not None: - return cached - with self.CacheLock: - cache_file = self._get_vote_cache_file(repository) - try: - data = entropy.dump.loadobj( - cache_file, - complete_path = True) - if data is not None: - self._set_live_cache_item(repository, cache_key, data) - except (IOError, EOFError, OSError): - data = None - return data - - def is_vote_cached(self, repository): - - cache_key = self._get_vote_cache_key(repository) - cached = self._is_live_cache_item_available(repository, cache_key) - if cached: - return True - - cache_file = self._get_vote_cache_file(repository) - avail = os.path.isfile(cache_file) and os.access(cache_file, os.R_OK) - if not avail: - return False - - try: - entropy.dump.loadobj(cache_file, complete_path = True) - except (IOError, EOFError, OSError): - return False - return True - - def get_downloads_cache(self, repository): - cache_key = self._get_downloads_cache_key(repository) - cached = self._get_live_cache_item(repository, cache_key) - if cached is not None: - return cached - with self.CacheLock: - cache_file = self._get_downloads_cache_file(repository) - try: - data = entropy.dump.loadobj(cache_file, complete_path = True) - if data is not None: - self._set_live_cache_item(repository, cache_key, data) - except (IOError, EOFError, OSError): - data = None - return data - - def is_downloads_cached(self, repository): - - cache_key = self._get_downloads_cache_key(repository) - cached = self._is_live_cache_item_available(repository, cache_key) - if cached: - return True - - cache_file = self._get_downloads_cache_file(repository) - avail = os.path.isfile(cache_file) and os.access(cache_file, os.R_OK) - if not avail: - return False - - try: - entropy.dump.loadobj(cache_file, complete_path = True) - except (IOError, EOFError, OSError): - return False - return True - - def get_alldocs_cache(self, pkgkey, repository): - cache_key = self._get_alldocs_cache_key(repository) - cached = self._get_live_cache_item(repository, cache_key) - if isinstance(cached, dict): - if pkgkey in cached: - return cached[pkgkey] - else: - cached = {} - with self.CacheLock: - cache_file = self._get_alldocs_cache_file(pkgkey, repository) - try: - data = entropy.dump.loadobj(cache_file, complete_path = True) - if data is not None: - cached[pkgkey] = data - self._set_live_cache_item(repository, cache_key, cached) - except (IOError, EOFError, OSError): - data = None - return data - - def is_alldocs_cached(self, pkgkey, repository): - - cache_key = self._get_alldocs_cache_key(repository) - cached = self._is_live_cache_item_available(repository, cache_key) - if cached: - return True - - cache_file = self._get_alldocs_cache_file(pkgkey, repository) - avail = os.path.isfile(cache_file) and os.access(cache_file, os.R_OK) - if not avail: - return False - - try: - entropy.dump.loadobj(cache_file, complete_path = True) - except (IOError, EOFError, OSError): - return False - return True - - def get_icon_cache(self, pkgkey, repository): - docs_cache = self.get_alldocs_cache(pkgkey, repository) - if docs_cache is None: - return None - elif not isinstance(docs_cache, (list, tuple)): - return None - - img_doc_type = etpConst['ugc_doctypes']['image'] - icon_id = Cache.PKG_ICON_IDENTIFIER - - try: - icon_elements = [x for x in docs_cache if \ - x['iddoctype'] == img_doc_type and x['title'] == icon_id] - except TypeError: - # cache corruption, invalid data, misc bullshit - return None - if not icon_elements: - return None - # get first (chronologically) icon - icon_elements = sorted(icon_elements, key = lambda x: x['ts']) - return icon_elements[0] - - def save_vote_cache(self, repository, vote_dict): - with self.CacheLock: - self._clear_live_cache_item(repository, - self._get_vote_cache_key(repository)) - - cache_file = self._get_vote_cache_file(repository) - entropy.dump.dumpobj(cache_file, vote_dict, complete_path = True) - - def save_downloads_cache(self, repository, down_dict): - with self.CacheLock: - self._clear_live_cache_item(repository, - self._get_downloads_cache_key(repository)) - cache_file = self._get_downloads_cache_file(repository) - entropy.dump.dumpobj(cache_file, down_dict, complete_path = True) - - def save_alldocs_cache(self, pkgkey, repository, alldocs_dict): - with self.CacheLock: - self._clear_live_cache_item(repository, - self._get_alldocs_cache_key(repository)) - cache_file = self._get_alldocs_cache_file(pkgkey, repository) - entropy.dump.dumpobj(cache_file, alldocs_dict, complete_path = True) - - def get_package_vote(self, repository, pkgkey): - cache = self.get_vote_cache(repository) - if not cache: - return 0.0 - elif not isinstance(cache, dict): - return 0.0 - elif pkgkey not in cache: - return 0.0 - return cache[pkgkey] - - def get_package_downloads(self, repository, pkgkey): - cache = self.get_downloads_cache(repository) - if not cache: - return 0 - elif not isinstance(cache, dict): - return 0 - elif pkgkey not in cache: - return 0 - try: - return int(cache[pkgkey]) - except ValueError: - return 0 diff --git a/libraries/entropy/services/auth_interfaces.py b/libraries/entropy/services/auth_interfaces.py deleted file mode 100644 index 2e87a33b1..000000000 --- a/libraries/entropy/services/auth_interfaces.py +++ /dev/null @@ -1,786 +0,0 @@ -# -*- coding: utf-8 -*- -""" - - @author: Fabio Erculiani - @contact: lxnay@sabayon.org - @copyright: Fabio Erculiani - @license: GPL-2 - - B{Entropy Services Authentication Interfaces}. - -""" - -import os -import time -import random -from entropy.services.skel import Authenticator, RemoteDatabase -from entropy.exceptions import PermissionDenied -from entropy.const import etpConst, const_isstring, const_convert_to_unicode -from entropy.i18n import _ - -import entropy.tools - -class phpBB3Auth(Authenticator, RemoteDatabase): - - def __init__(self): - Authenticator.__init__(self) - RemoteDatabase.__init__(self) - - self.itoa64 = './0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz' - self.USER_NORMAL = 0 - self.USER_INACTIVE = 1 - self.USER_IGNORE = 2 - self.USER_FOUNDER = 3 - self.REGISTERED_USERS_GROUP = 7895 - self.ADMIN_GROUPS = [7893, 7898] - self.MODERATOR_GROUPS = [484] - self.DEVELOPER_GROUPS = [7900] - self.USERNAME_LENGTH_RANGE = list(range(3, 21)) - self.PASSWORD_LENGTH_RANGE = list(range(6, 31)) - self.PRIVMSGS_NO_BOX = -3 - self.NOTIFY_EMAIL = 0 - self.FAKE_USERNAME = 'already_authed' - self.USER_AGENT = "Entropy/%s (compatible; %s; %s: %s %s %s)" % ( - etpConst['entropyversion'], - "Entropy", - "UGC", - os.uname()[0], - os.uname()[4], - os.uname()[2], - ) - self.TABLE_PREFIX = 'phpbb_' - self.do_update_session_table = True - - def validate_username_regex(self, username): - allow_name_chars = self._get_config_value("allow_name_chars") - if allow_name_chars == "USERNAME_CHARS_ANY": - regex = '.+' - elif allow_name_chars == "USERNAME_ALPHA_ONLY": - regex = '[A-Za-z0-9]+' - elif allow_name_chars == "USERNAME_ALPHA_SPACERS": - regex = '[A-Za-z0-9-[\]_+ ]+' - elif allow_name_chars == "USERNAME_LETTER_NUM": - regex = '[a-zA-Z0-9]+' - elif allow_name_chars == "USERNAME_LETTER_NUM_SPACERS": - regex = '[-\]_+ [a-zA-Z0-9]+' - else: # USERNAME_ASCII - regex = '[\x01-\x7F]+' - regex = "^%s$" % (regex,) - import re - myreg = re.compile(regex) - if myreg.match(username): - del myreg - return True - return False - - def does_username_exist(self, username, username_clean): - self.check_connection() - self.cursor.execute('SELECT user_id FROM '+self.TABLE_PREFIX+'users WHERE `username_clean` = %s OR LOWER(`username`) = %s', (username_clean, username.lower(),)) - data = self.cursor.fetchone() - if not data: - return False - if not isinstance(data, dict): - return False - if 'user_id' not in data: - return False - return True - - def does_email_exist(self, email): - self.check_connection() - self.cursor.execute('SELECT user_id FROM '+self.TABLE_PREFIX+'users WHERE `user_email` = %s', (email,)) - data = self.cursor.fetchone() - if not data: - return False - if not isinstance(data, dict): - return False - if 'user_id' not in data: - return False - return True - - def is_username_allowed(self, username): - self.check_connection() - self.cursor.execute('SELECT disallow_id FROM '+self.TABLE_PREFIX+'disallow WHERE `disallow_username` = %s', (username,)) - data = self.cursor.fetchone() - if not data: - return True - if not isinstance(data, dict): - return True - if 'disallow_id' not in data: - return True - return False - - def validate_username_string(self, username, username_clean): - - try: - const_convert_to_unicode(username.encode('utf-8')) - except (UnicodeDecodeError, UnicodeEncodeError,): - return False, 'Invalid username' - if (""" in username) or ("'" in username) or ('"' in username) or \ - (" " in username): - return False, 'Invalid username' - - try: - valid = self.validate_username_regex(username) - except: - return False, 'Username contains bad characters' - if not valid: - return False, 'Invalid username' - - exists = self.does_username_exist(username, username_clean) - if exists: - return False, 'Username already taken' - - allowed = self.is_username_allowed(username) - if not allowed: - return False, 'Username not allowed' - - return True, 'All fine' - - def _generate_email_hash(self, email): - import binascii - return str(binascii.crc32(email.lower())) + str(len(email)) - - def activate_user(self, user_id): - self.check_connection() - self.cursor.execute('UPDATE '+self.TABLE_PREFIX+'users SET user_type = %s WHERE `user_id` = %s', (self.USER_NORMAL, user_id,)) - return True, user_id - - def generate_username_clean(self, username): - import re - username_clean = username.lower() - username_clean = re.sub(r'(?:[\x00-\x1F\x7F]+|(?:\xC2[\x80-\x9F])+)', '', username_clean) - username_clean = re.sub(r' {2,}', ' ', username_clean) - username_clean = username_clean.strip() - return username_clean - - def register_user(self, username, password, email, activate = False): - - if len(username) not in self.USERNAME_LENGTH_RANGE: - return False, 'Username not in range' - if len(password) not in self.PASSWORD_LENGTH_RANGE: - return False, 'Password not in range' - valid = entropy.tools.is_valid_email(email) - if not valid: - return False, 'Invalid email' - - # create the clean one - username_clean = self.generate_username_clean(username) - - # check username validity - status, err_msg = self.validate_username_string(username, username_clean) - if not status: - return False, err_msg - - # check email - exists = self.does_email_exist(email) - if exists: - return False, 'Email already in use' - - # now cross fingers - status, user_id = self.__register(username, username_clean, password, email, activate) - if not status: - return False, 'Invalid username (duplicated)' - - return True, user_id - - - def __register(self, username, username_clean, password, email, activate): - - email_hash = self._generate_email_hash(email) - password_hash = self._get_password_hash(password.encode('utf-8')) - time_now = int(time.time()) - - user_type = self.USER_INACTIVE - if activate: - user_type = self.USER_NORMAL - - registration_data = { - 'username': username, - 'username_clean': username_clean, - 'user_password': password_hash, - 'user_pass_convert': 0, - 'user_email': email.lower(), - 'user_email_hash': email_hash, - 'group_id': self.REGISTERED_USERS_GROUP, - 'user_type': user_type, - 'user_permissions': '', - 'user_timezone': self._get_config_value('board_timezone'), - 'user_dateformat': self._get_config_value('default_dateformat'), - 'user_lang': self._get_config_value('default_lang'), - 'user_style': self._get_config_value('default_style'), - 'user_actkey': '', - 'user_ip': '', - 'user_regdate': time_now, - 'user_passchg': time_now, - 'user_options': 895, # ? don't ask me - 'user_inactive_reason': 0, - 'user_inactive_time': 0, - 'user_lastmark': time_now, - 'user_lastvisit': 0, - 'user_lastpost_time': 0, - 'user_lastpage': '', - 'user_posts': 0, - 'user_dst': self._get_config_value('board_dst'), - 'user_colour': '', - 'user_occ': '', - 'user_interests': '', - 'user_avatar': '', - 'user_avatar_type': 0, - 'user_avatar_width': 0, - 'user_avatar_height': 0, - 'user_new_privmsg': 0, - 'user_unread_privmsg': 0, - 'user_last_privmsg': 0, - 'user_message_rules': 0, - 'user_full_folder': self.PRIVMSGS_NO_BOX, - 'user_emailtime': 0, - 'user_notify': 0, - 'user_notify_pm': 1, - 'user_notify_type': self.NOTIFY_EMAIL, - 'user_allow_pm': 1, - 'user_allow_viewonline': 1, - 'user_allow_viewemail': 1, - 'user_allow_massemail': 1, - 'user_sig': '', - 'user_sig_bbcode_uid': '', - 'user_sig_bbcode_bitfield': '', - 'user_form_salt': self._get_unique_id(), - } - - sql = self._generate_sql('insert', self.TABLE_PREFIX+'users', registration_data) - self.cursor.execute(sql) - user_id = self.cursor.lastrowid - - # now insert into the default group - group_data = { - 'user_id': user_id, - 'group_id': self.REGISTERED_USERS_GROUP, - 'user_pending': 0, - } - sql = self._generate_sql('insert', self.TABLE_PREFIX+'user_group', group_data) - try: - self.cursor.execute(sql) - except self.mysql_exceptions.IntegrityError as e: - # for sure it's about duplicated entry - return False, 1062 - - # set some misc config shit - self._set_config_value('newest_user_id', user_id) - self._set_config_value('newest_username', username) - self._set_config_value('num_users', int(self._get_config_value('num_users'))+1) - self.cursor.execute('SELECT group_colour FROM '+self.TABLE_PREFIX+'groups WHERE group_id = %s', (group_data['group_id'],)) - data = self.cursor.fetchone() - gcolor = None - if isinstance(data, dict): - if 'group_colour' in data: - gcolor = data['group_colour'] - if gcolor: - self._set_config_value('newest_user_colour', gcolor) - - return True, user_id - - - def login(self): - self.check_connection() - self.check_login_data() - - if 'username' not in self.login_data: - raise PermissionDenied('PermissionDenied: %s' % (_('no username specified'),)) - elif 'password' not in self.login_data: - raise PermissionDenied('PermissionDenied: %s' % (_('no password specified'),)) - - if not self.login_data['password']: - raise PermissionDenied('PermissionDenied: %s' % (_('empty password'),)) - elif not self.login_data['username']: - raise PermissionDenied('PermissionDenied: %s' % (_('empty username'),)) - - self.cursor.execute('SELECT * FROM '+self.TABLE_PREFIX+'users WHERE username = %s', (self.login_data['username'],)) - data = self.cursor.fetchone() - if not data: - raise PermissionDenied('PermissionDenied: %s' % (_('user not found'),)) - - if data['user_pass_convert']: - raise PermissionDenied('PermissionDenied: %s' % ( - _('you need to login on the website to update your password format'), - ) - ) - - valid = self._phpbb3_check_hash(self.login_data['password'], data['user_password']) - if not valid: - raise PermissionDenied('PermissionDenied: %s' % (_('wrong password'),)) - - user_type = data['user_type'] - if (user_type == self.USER_INACTIVE) or (user_type == self.USER_IGNORE): - raise PermissionDenied('PermissionDenied: %s' % (_('user inactive'),)) - - banned = self.is_user_banned(data['user_id']) - if banned: - raise PermissionDenied('PermissionDenied: %s' % (_('user banned'),)) - - self.login_data.update(data) - self.logged_in = True - return self.logged_in - - def disconnect(self): - if self.is_logged_in(): - self.logout() - RemoteDatabase.disconnect(self) - - def logout(self): - self.check_connection() - self.check_login_data() - self.logged_in = False - self.login_data.clear() - return True - - def get_user_data(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - self.cursor.execute('SELECT * FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - return self.cursor.fetchone() - - def get_user_birthday(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - self.cursor.execute('SELECT user_birthday FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - bday = self.cursor.fetchone() - if not bday: - return None - elif 'user_birthday' not in bday: - return None - return bday['user_birthday'] - - def get_username(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - self.cursor.execute('SELECT username_clean FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - data = self.cursor.fetchone() - if not data: - return '' - elif 'username_clean' not in data: - return '' - return data['username_clean'] - - def is_developer(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - # search into phpbb_groups - groups = self.get_user_groups() - for group in groups: - if group in self.DEVELOPER_GROUPS: - return True - - return False - - def is_administrator(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - self.cursor.execute('SELECT user_type FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - data = self.cursor.fetchone() - if data: - if data['user_type'] == self.USER_FOUNDER: - return True - - # search into phpbb_groups - groups = self.get_user_groups() - for group in groups: - if group in self.ADMIN_GROUPS: - return True - - return False - - def is_moderator(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - # search into phpbb_groups - groups = self.get_user_groups() - for group in groups: - if group in self.MODERATOR_GROUPS: - return True - - return False - - def is_user(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - if self.is_moderator(): - return False - elif self.is_administrator(): - return False - elif self.is_developer(): - return False - - self.cursor.execute('SELECT user_type,user_id FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - data = self.cursor.fetchone() - if not data: - return False - if self.is_user_banned(data['user_id']): - return False - elif data['user_type'] in [self.USER_NORMAL]: - return True - - return False - - # user == user_id - def is_user_banned(self, user): - self.check_connection() - self.cursor.execute('SELECT ban_userid FROM '+self.TABLE_PREFIX+'banlist WHERE ban_userid = %s', (user,)) - data = self.cursor.fetchone() - if data: - return True - return False - - def is_in_group(self, group): - self.check_connection() - self.check_login_data() - self.check_logged_in() - groups = self.get_user_groups() - if isinstance(group, int): - if group in groups: - return True - elif const_isstring(group): - self.cursor.execute('SELECT group_id FROM '+self.TABLE_PREFIX+'groups WHERE group_name = %s', (group,)) - data = self.cursor.fetchone() - if not data: - return False - elif data['group_id'] in groups: - return True - - return False - - def get_user_groups(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - self.cursor.execute('SELECT '+self.TABLE_PREFIX+'user_group.group_id,'+self.TABLE_PREFIX+'groups.group_name FROM '+self.TABLE_PREFIX+'user_group,'+self.TABLE_PREFIX+'users,'+self.TABLE_PREFIX+'groups WHERE '+self.TABLE_PREFIX+'users.user_id = %s and '+self.TABLE_PREFIX+'users.user_id = '+self.TABLE_PREFIX+'user_group.user_id and '+self.TABLE_PREFIX+'user_group.group_id = '+self.TABLE_PREFIX+'groups.group_id', (self.login_data['user_id'],)) - data = self.cursor.fetchall() - mydata = {} - for mydict in data: - mydata[mydict['group_id']] = mydict['group_name'] - - return mydata - - def get_user_group(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - self.cursor.execute('SELECT group_id FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - data = self.cursor.fetchone() - if data: - if 'group_id' in data: - return data['group_id'] - - return -1 - - def get_user_id(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - return self.login_data['user_id'] - - def get_permission_data(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - self.cursor.execute('SELECT user_permissions FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - return self.cursor.fetchone() - - def update_email(self, email): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - email_hash = self._generate_email_hash(email) - mydata = { - 'user_email_hash': email_hash, - 'user_email': email.lower(), - } - - try: - sql = self._generate_sql("update", self.TABLE_PREFIX+'users', mydata, 'user_id = %s' % (self.login_data['user_id'],)) - self.cursor.execute(sql) - return True - except Exception: - return False - - def update_password_hash(self, password_hash): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - mydata = { - 'user_password': password_hash, - } - - try: - sql = self._generate_sql("update", self.TABLE_PREFIX+'users', mydata, 'user_id = %s' % (self.login_data['user_id'],)) - self.cursor.execute(sql) - return True - except Exception: - return False - - def get_email(self): - self.check_connection() - self.check_login_data() - self.check_logged_in() - self.cursor.execute('SELECT user_email FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (self.login_data['user_id'],)) - data = self.cursor.fetchone() - if not data: - return '' - elif 'user_email' not in data: - return '' - return data['user_email'] - - def update_user_id_profile(self, profile_data): - self.check_connection() - self.check_login_data() - self.check_logged_in() - - # filter valid params - valid_params = [ - "user_icq", "user_yim", "user_msnm", - "user_jabber", "user_website", "user_from", - "user_interests", "user_occ", "user_birthday", - "user_sig" - ] - - my_params = {} - for param in valid_params: - d = profile_data.get(param) - if d is None: - continue - my_params[param] = d - - if not my_params: - return False, 'no parameters' - - - # validate parameters - b_day = my_params.get('user_birthday') - if const_isstring(b_day): - import re - myre = re.compile("(0[1-9]|[12][0-9]|3[01])[-](0[1-9]|1[012])[-](19|20)\d\d") - if not myre.match(b_day): - del my_params['user_birthday'] - - try: - sql = self._generate_sql("update", self.TABLE_PREFIX+'users', my_params, 'user_id = %s' % (self.login_data['user_id'],)) - self.cursor.execute(sql) - return True, None - except Exception as e: - return False, str(e) - - - def _set_config_value(self, config_name, data): - self.cursor.execute('UPDATE '+self.TABLE_PREFIX+'config SET config_value = %s WHERE config_name = %s', (data, config_name,)) - - def _get_config_value(self, config_name): - self.check_connection() - self.cursor.execute('SELECT config_value FROM '+self.TABLE_PREFIX+'config WHERE config_name = %s', (config_name,)) - myconfig = self.cursor.fetchone() - if isinstance(myconfig, dict): - if 'config_value' in myconfig: - return myconfig['config_value'] - return None - - def _update_session_table(self, user_id, ip_address): - self.check_connection() - time_now = int(time.time()) - autologin = self._get_config_value("allow_autologin") - self.cursor.execute('SELECT user_allow_viewonline FROM '+self.TABLE_PREFIX+'users WHERE user_id = %s', (user_id,)) - myuserprefs = self.cursor.fetchone() - session_admin = 0 - session_data = { - 'session_id': None, - 'session_user_id': user_id, - 'session_last_visit': time_now, - 'session_start': time_now, - 'session_time': time_now, - 'session_ip': ip_address, - 'session_browser': self.USER_AGENT, - 'session_forwarded_for': '', - 'session_page': 'index.php', - 'session_viewonline': myuserprefs['user_allow_viewonline'], - 'session_autologin': autologin, - 'session_admin': session_admin, - 'session_forum_id': 0, - } - import hashlib - m = hashlib.md5() - m.update(str(user_id)+str(time_now)+str(self.USER_AGENT)+str(ip_address)+str(autologin)+str(myuserprefs['user_allow_viewonline'])) - session_data['session_id'] = m.hexdigest() - - self.cursor.execute('SELECT * FROM '+self.TABLE_PREFIX+'sessions WHERE session_user_id = %s', (user_id,)) - mydata = self.cursor.fetchone() - do_update = False - if mydata: - do_update = True - # update - session_data['session_id'] = mydata['session_id'] - session_data['session_viewonline'] = mydata['session_viewonline'] - session_data['session_autologin'] = mydata['session_autologin'] - session_data['session_forwarded_for'] = mydata['session_forwarded_for'] - session_data['session_forum_id'] = mydata['session_forum_id'] - session_data['session_page'] = mydata['session_page'] - session_data['session_browser'] = mydata['session_browser'] - session_data['session_admin'] = mydata['session_admin'] - - if do_update: - where = "session_id = '%s'" % (session_data['session_id'],) - del session_data['session_id'] - sql = self._generate_sql('update', self.TABLE_PREFIX+'sessions', session_data, where) - else: - sql = self._generate_sql('insert', self.TABLE_PREFIX+'sessions', session_data) - if sql: - self.cursor.execute(sql) - self.dbconn.commit() - - - def _is_ip_banned(self, ip): - self.check_connection() - self.cursor.execute('SELECT ban_ip FROM '+self.TABLE_PREFIX+'banlist WHERE ban_ip = %s', (ip,)) - data = self.cursor.fetchone() - if data: - return True - return False - - def _get_unique_id(self): - import hashlib - m = hashlib.md5() - rnd = str(abs(hash(os.urandom(1)))) - m.update(rnd) - x = m.hexdigest()[:-16] - del m - return x - - def _get_random_number(self): - myrand = 0 - low_n = 100000 - high_n = 999999 - while (myrand < low_n) or (myrand > high_n): - try: - myrand = hash(os.urandom(1))%high_n - except NotImplementedError: - random.seed() - myrand = random.randint(low_n, high_n) - return myrand - - def _get_password_hash(self, password): - - #random_state = self._get_unique_id() - myrandom = str(self._get_random_number()) - - myhash = self._hash_crypt_private(password, self._hash_gensalt_private(myrandom)) - - if len(myhash) == 34: - return myhash - - import hashlib - m = hashlib.md5() - m.update(myhash) - return m.hexdigest() - - - def _hash_gensalt_private(self, myinput, iteration_count_log2 = 6): - - if (iteration_count_log2 < 4) or (iteration_count_log2 > 31): - iteration_count_log2 = 8 - - myoutput = '$H$' - myoutput += self.itoa64[min(iteration_count_log2 + 5, 30)] - myoutput += self._hash_encode64(myinput, 6) - - return myoutput - - def _hash_crypt_private(self, password, setting): - - myoutput = '*' - # Check for correct hash - if setting[:3] != '$H$': - return myoutput - - count_log2 = self.itoa64.find(setting[3]) - if count_log2 == -1: - count_log2 = 0 - - if (count_log2 < 7) or (count_log2 > 30): - return myoutput - - count = 1 << count_log2 - salt = setting[4:12] - - if len(salt) != 8: - return myoutput - - import hashlib - m = hashlib.md5() - m.update(salt+password) - myhash = m.digest() - while count: - m = hashlib.md5() - m.update(myhash+password) - myhash = m.digest() - count -= 1 - - myoutput = setting[:12] - myoutput += self._hash_encode64(myhash, 16) - - return myoutput - - def _hash_encode64(self, myinput, count): - - output = '' - i = 0 - while i < count: - - value = ord(myinput[i]) - i += 1 - output += self.itoa64[value & 0x3f] - if i < count: - value |= ord(myinput[i]) << 8 - - output += self.itoa64[(value >> 6) & 0x3f] - - if i >= count: - break - i += 1 - - if i < count: - value |= ord(myinput[i]) << 16 - - output += self.itoa64[(value >> 12) & 0x3f] - - if (i >= count): - break - i += 1 - - output += self.itoa64[(value >> 18) & 0x3f] - - return output - - def _phpbb3_check_hash(self, password, myhash): - - if len(myhash) == 34: - return self._hash_crypt_private(password, myhash) == myhash - - import hashlib - m = hashlib.md5() - m.update(password) - rhash = m.hexdigest() - return rhash == myhash diff --git a/libraries/entropy/services/authenticators.py b/libraries/entropy/services/authenticators.py deleted file mode 100644 index d813696e6..000000000 --- a/libraries/entropy/services/authenticators.py +++ /dev/null @@ -1,115 +0,0 @@ -# -*- coding: utf-8 -*- -""" - - @author: Fabio Erculiani - @contact: lxnay@sabayon.org - @copyright: Fabio Erculiani - @license: GPL-2 - - B{Entropy Services Authenticators Interface}. - -""" - -from entropy.const import const_get_stringtype -from entropy.services.auth_interfaces import phpBB3Auth -from entropy.services.skel import SocketAuthenticator -from entropy.exceptions import PermissionDenied - -# Authenticator that can be used by SocketHostInterface based instances -class phpBB3(phpBB3Auth, SocketAuthenticator): - - def __init__(self, HostInterface, *args, **kwargs): - SocketAuthenticator.__init__(self, HostInterface) - phpBB3Auth.__init__(self) - self.set_connection_data(kwargs) - self.connect() - - def set_session(self, session): - self.session = session - session_data = self.HostInterface.sessions.get(self.session) - if not session_data: - return - auth_id = session_data['auth_uid'] - if auth_id: - self.logged_in = True - # fill login_data with fake information - self.login_data = { - 'username': self.FAKE_USERNAME, - 'password': 'look elsewhere, this is not a password', - 'user_id': auth_id - } - ip_address = session_data.get('ip_address') - if ip_address and self.do_update_session_table: - self._update_session_table(auth_id, ip_address) - - def docmd_login(self, arguments): - - # filter n00bs - if not arguments or (len(arguments) != 2): - return False, None, None, 'wrong arguments' - - ip_address = None - session_data = self.HostInterface.sessions.get(self.session) - if session_data: - ip_address = session_data.get('ip_address') - user = arguments[0] - password = arguments[1] - - if ip_address: - if self._is_ip_banned(ip_address): - return False, user, None, "banned IP" - - login_data = {'username': user, 'password': password} - self.set_login_data(login_data) - rc = False - try: - rc = self.login() - except PermissionDenied as e: - return rc, user, None, e.value - - if rc: - uid = self.get_user_id() - is_admin = self.is_administrator() - is_dev = self.is_developer() - is_mod = self.is_moderator() - is_user = self.is_user() - self.HostInterface.sessions[self.session]['admin'] = is_admin - self.HostInterface.sessions[self.session]['developer'] = is_dev - self.HostInterface.sessions[self.session]['moderator'] = is_mod - self.HostInterface.sessions[self.session]['user'] = is_user - if ip_address and uid and self.do_update_session_table: - self._update_session_table(uid, ip_address) - return True, user, uid, "ok" - return rc, user, None, "login failed" - - # if we get here it means we are logged in - def docmd_userdata(self): - data = self.get_user_data() - return True, data, 'ok' - - def docmd_logout(self, myargs): - - # filter n00bs - if (len(myargs) < 1) or (len(myargs) > 1): - return False, None, 'wrong arguments' - - user = myargs[0] - # filter n00bs - if not user or not isinstance(user, const_get_stringtype()): - return False, None, "wrong user" - - if not self.is_logged_in(): - return False, user, "already logged out" - - return True, user, "ok" - - def set_exc_permissions(self, *args, **kwargs): - pass - - def hide_login_data(self, args): - myargs = args[:] - myargs[-1] = 'hidden' - return myargs - - def terminate_instance(self): - self.disconnect() diff --git a/libraries/entropy/services/commands.py b/libraries/entropy/services/commands.py deleted file mode 100644 index d660ae45b..000000000 --- a/libraries/entropy/services/commands.py +++ /dev/null @@ -1,73 +0,0 @@ -# -*- coding: utf-8 -*- -""" - - @author: Fabio Erculiani - @contact: lxnay@sabayon.org - @copyright: Fabio Erculiani - @license: GPL-2 - - B{Entropy Services Command Interfaces}. - -""" -from entropy.services.skel import SocketCommands - -class phpBB3(SocketCommands): - - def __init__(self, HostInterface): - - SocketCommands.__init__(self, HostInterface, inst_name = "phpbb3-commands") - - self.valid_commands = { - 'is_user': { - 'auth': True, - 'built_in': False, - 'cb': self.docmd_is_user, - 'args': ["authenticator"], - 'as_user': False, - 'desc': "returns whether the username linked with the session belongs to a simple user", - 'syntax': " is_user", - 'from': str(self), # from what class - }, - 'is_developer': { - 'auth': True, - 'built_in': False, - 'cb': self.docmd_is_developer, - 'args': ["authenticator"], - 'as_user': False, - 'desc': "returns whether the username linked with the session belongs to a developer", - 'syntax': " is_developer", - 'from': str(self), # from what class - }, - 'is_moderator': { - 'auth': True, - 'built_in': False, - 'cb': self.docmd_is_moderator, - 'args': ["authenticator"], - 'as_user': False, - 'desc': "returns whether the username linked with the session belongs to a moderator", - 'syntax': " is_moderator", - 'from': str(self), # from what class - }, - 'is_administrator': { - 'auth': True, - 'built_in': False, - 'cb': self.docmd_is_administrator, - 'args': ["authenticator"], - 'as_user': False, - 'desc': "returns whether the username linked with the session belongs to an administrator", - 'syntax': " is_administrator", - 'from': str(self), # from what class - }, - } - - def docmd_is_user(self, authenticator): - return authenticator.is_user(), 'ok' - - def docmd_is_developer(self, authenticator): - return authenticator.is_developer(), 'ok' - - def docmd_is_administrator(self, authenticator): - return authenticator.is_administrator(), 'ok' - - def docmd_is_moderator(self, authenticator): - return authenticator.is_moderator(), 'ok' diff --git a/libraries/entropy/services/exceptions.py b/libraries/entropy/services/exceptions.py deleted file mode 100644 index 9cfc9fe47..000000000 --- a/libraries/entropy/services/exceptions.py +++ /dev/null @@ -1,31 +0,0 @@ -# -*- coding: utf-8 -*- -""" - - @author: Fabio Erculiani - @contact: lxnay@sabayon.org - @copyright: Fabio Erculiani - @license: GPL-2 - - B{Entropy Services Exceptions module}. - These are the exceptions raised by Entropy Services RPC system. - -""" -from entropy.exceptions import EntropyException - -class EntropyServicesError(EntropyException): - """ Generic Entropy Services exception. All classes here belong to this. """ - -class TransmissionError(EntropyServicesError): - """ Generic transmission error exception """ - -class BrokenPipe(TransmissionError): - """ Broken pipe transmission error """ - -class SSLTransmissionError(TransmissionError): - """ Error on SSL socket """ - -class ServiceConnectionError(EntropyServicesError): - """Cannot connect to service""" - -class TimeoutError(EntropyServicesError): - """ Timeout error """ diff --git a/libraries/entropy/services/interfaces.py b/libraries/entropy/services/interfaces.py deleted file mode 100644 index f0c60fee3..000000000 --- a/libraries/entropy/services/interfaces.py +++ /dev/null @@ -1,2073 +0,0 @@ -# -*- coding: utf-8 -*- -""" - - @author: Fabio Erculiani - @contact: lxnay@sabayon.org - @copyright: Fabio Erculiani - @license: GPL-2 - - B{Entropy Services Base Interfaces}. - -""" -import sys -import os -import errno -import tempfile -import select -import shutil -import time -import copy - -import entropy.dump -import entropy.tools -from entropy.core.settings.base import SystemSettings -from entropy.const import etpConst, const_setup_perms, const_isstring, \ - const_get_stringtype, const_convert_to_rawstring, etpUi, const_debug_write -from entropy.exceptions import InterruptError, PermissionDenied, DumbException -from entropy.services.exceptions import TimeoutError, ServiceConnectionError -from entropy.services.skel import SocketAuthenticator, SocketCommands -from entropy.i18n import _ -from entropy.output import blue, red, darkgreen, darkred, darkblue, brown, \ - purple - -try: - import SocketServer as socketserver -except ImportError: # Python 3.x - import socketserver - - -class SocketHost: - - import socket - from threading import Thread - - LOG_FILE = os.path.join(etpConst['syslogdir'], "socket.log") - LOG_FILE_SSL = os.path.join(etpConst['syslogdir'], "ssl.socket.log") - - class BasicPamAuthenticator(SocketAuthenticator): - - def __init__(self, HostInterface, *args, **kwargs): - self.valid_auth_types = [ "plain", "shadow", "md5" ] - SocketAuthenticator.__init__(self, HostInterface) - - def docmd_login(self, arguments): - - # filter n00bs - if not arguments or (len(arguments) != 3): - return False, None, None, 'wrong arguments' - - user = arguments[0] - auth_type = arguments[1] - auth_string = arguments[2] - - # check auth type validity - if auth_type not in self.valid_auth_types: - return False, user, None, 'invalid auth type' - - udata = self.__get_user_data(user) - if udata == None: - return False, user, None, 'invalid user' - - uid = udata[2] - # check if user is in the Entropy group - if not entropy.tools.is_user_in_entropy_group(uid): - return False, user, uid, 'user not in %s group' % (etpConst['sysgroup'],) - - # now validate password - valid = self.__validate_auth(user, auth_type, auth_string) - if not valid: - return False, user, uid, 'auth failed' - - if not uid: - self.HostInterface.sessions[self.session]['admin'] = True - else: - self.HostInterface.sessions[self.session]['user'] = True - return True, user, uid, "ok" - - # it we get here is because user is logged in - def docmd_userdata(self): - - auth_uid = self.HostInterface.sessions[self.session]['auth_uid'] - mydata = {} - udata = self.__get_uid_data(auth_uid) - if udata: - mydata['username'] = udata[0] - mydata['uid'] = udata[2] - mydata['gid'] = udata[3] - mydata['references'] = udata[4] - mydata['home'] = udata[5] - mydata['shell'] = udata[6] - return True, mydata, 'ok' - - def __get_uid_data(self, user_id): - import pwd - # check user validty - try: - udata = pwd.getpwuid(user_id) - except KeyError: - return None - return udata - - def __get_user_data(self, user): - import pwd - # check user validty - try: - udata = pwd.getpwnam(user) - except KeyError: - return None - return udata - - def __validate_auth(self, user, auth_type, auth_string): - valid = False - if auth_type == "plain": - valid = self.__do_auth(user, auth_string) - elif auth_type == "shadow": - valid = self.__do_auth(user, auth_string, auth_type = "shadow") - elif auth_type == "md5": - valid = self.__do_auth(user, auth_string, auth_type = "md5") - return valid - - def __do_auth(self, user, password, auth_type = None): - import spwd - - try: - enc_pass = spwd.getspnam(user)[1] - except KeyError: - return False - - if auth_type == None: # plain - import crypt - generated_pass = crypt.crypt(str(password), enc_pass) - elif auth_type == "shadow": - generated_pass = password - elif auth_type == "md5": # md5 - import hashlib - m = hashlib.md5() - m.update(enc_pass) - enc_pass = m.hexdigest() - generated_pass = str(password) - else: # haha, fuck! - generated_pass = None - - if generated_pass == enc_pass: - return True - return False - - def docmd_logout(self, myargs): - - # filter n00bs - if (len(myargs) < 1) or (len(myargs) > 1): - return False, None, 'wrong arguments' - - user = myargs[0] - # filter n00bs - if not user or not const_isstring(user): - return False, None, "wrong user" - - return True, user, "ok" - - def hide_login_data(self, args): - myargs = args[:] - myargs[-1] = 'hidden' - return myargs - - class HostServerMixin(socketserver.ThreadingMixIn, socketserver.TCPServer): - - class ConnWrapper: - ''' - Base class for implementing the rest of the wrappers in this module. - Operates by taking a connection argument which is used when 'self' doesn't - provide the functionality being requested. - ''' - def __init__(self, connection) : - self.connection = connection - - def __getattr__(self, function) : - return getattr(self.connection, function) - - import socket as socket_mod - # This means the main server will not do the equivalent of a - # pthread_join() on the new threads. With this set, Ctrl-C will - # kill the server reliably. - daemon_threads = True - - # By setting this we allow the server to re-bind to the address by - # setting SO_REUSEADDR, meaning you don't have to wait for - # timeouts when you kill the server and the sockets don't get - # closed down correctly. - allow_reuse_address = True - - def __init__(self, server_address, RequestHandlerClass, processor, HostInterface, authorized_clients_only = False): - - self.alive = True - self.socket = self.socket_mod - self.processor = processor - self.server_address = server_address - self.HostInterface = HostInterface - self.SSL = self.HostInterface.SSL - self.real_sock = None - self.ssl_authorized_clients_only = authorized_clients_only - - if self.HostInterface.is_ssl_enabled(): - socketserver.BaseServer.__init__(self, server_address, - RequestHandlerClass) - self.load_ssl_context() - self.make_ssl_connection_alive() - else: - try: - socketserver.TCPServer.__init__(self, server_address, - RequestHandlerClass) - except self.socket_mod.error as e: - if e.errno == errno.EACCESS: - raise ServiceConnectionError("Cannot bind service") - raise - - def load_ssl_context(self): - # setup an SSL context. - self.context = self.SSL['m'].Context(self.SSL['m'].SSLv23_METHOD) - self.context.set_verify(self.SSL['m'].VERIFY_PEER, self.verify_ssl_cb) # ask for a certificate - self.context.set_options(self.SSL['m'].OP_NO_SSLv2) - # load up certificate stuff. - self.context.use_privatekey_file(self.SSL['key']) - self.context.use_certificate_file(self.SSL['cert']) - self.context.load_verify_locations(self.SSL['ca_cert']) - self.context.load_client_ca(self.SSL['ca_cert']) - self.HostInterface.output('SSL context loaded, key: %s - cert: %s, CA cert: %s, CA pkey: %s' % ( - self.SSL['key'], - self.SSL['cert'], - self.SSL['ca_cert'], - self.SSL['ca_pkey'] - ) - ) - - def make_ssl_connection_alive(self): - self.real_sock = self.socket_mod.socket(self.address_family, self.socket_type) - self.socket = self.ConnWrapper(self.SSL['m'].Connection(self.context, self.real_sock)) - self.server_bind() - self.server_activate() - - # this function should do the authentication checking to see that - # the client is who they say they are. - def verify_ssl_cb(self, conn, cert, errnum, depth, ok) : - return ok - - def verify_request(self, request, client_address): - - self.do_ssl = self.HostInterface.SSL - if self.do_ssl: - self.do_ssl = True - else: - self.do_ssl = False - - allowed = self.ip_blacklist_check(client_address[0]) - if allowed: allowed = self.ip_max_connections_check(client_address[0]) - if not allowed: - self.HostInterface.output( - '[from: %s | SSL: %s] connection refused, ip blacklisted or maximum connections per IP reached' % ( - client_address, - self.do_ssl, - ) - ) - return False - - allowed = self.max_connections_check(request) - if not allowed: - self.HostInterface.output( - '[from: %s | SSL: %s] connection refused (max connections reached: %s)' % ( - client_address, - self.do_ssl, - self.HostInterface.max_connections, - ) - ) - return False - - ### let's go! - self.HostInterface.connections += 1 - self.HostInterface.output( - '[from: %s | SSL: %s] connection established (%s of %s max connections)' % ( - client_address, - self.do_ssl, - self.HostInterface.connections, - self.HostInterface.max_connections, - ) - ) - return True - - def ip_blacklist_check(self, client_addr): - if client_addr in self.HostInterface.ip_blacklist: - return False - return True - - def ip_max_connections_check(self, ip_address): - max_conn_per_ip = self.HostInterface.max_connections_per_host - max_conn_per_ip_barrier = self.HostInterface.max_connections_per_host_barrier - per_host_connections = self.HostInterface.per_host_connections - conn_data = per_host_connections.get(ip_address) - if conn_data == None: - per_host_connections[ip_address] = 1 - else: - conn_data += 1 - per_host_connections[ip_address] += 1 - if conn_data > max_conn_per_ip: - self.HostInterface.output( - '[from: %s] ------- :EEK: !! connection closed too many simultaneous connections from host (current: %s | limit: %s) -------' % ( - ip_address, - conn_data, - max_conn_per_ip, - ) - ) - return False - elif conn_data > max_conn_per_ip_barrier: - times = [5, 6, 7, 8] - self.HostInterface.output( - '[from: %s] ------- :EEEK: !! connection warning simultaneous connection barrier reached from host (current: %s | soft limit: %s) -------' % ( - ip_address, - conn_data, - max_conn_per_ip_barrier, - ) - ) - rnd_num = entropy.tools.get_random_number() - time.sleep(times[abs(hash(rnd_num))%len(times)]) - - return True - - def max_connections_check(self, request): - current = self.HostInterface.connections - maximum = self.HostInterface.max_connections - if current >= maximum: - try: - self.HostInterface.transmit( - request, - self.HostInterface.answers['mcr'] - ) - except: - pass - return False - else: - return True - - def serve_forever(self): - while self.alive: - #r,w,e = select.select([self.socket], [], [], 1) - #if r: - self.handle_request() - - # taken from SocketServer.py - def finish_request(self, request, client_address): - """Finish one request by instantiating RequestHandlerClass.""" - self.RequestHandlerClass(request, client_address, self) - - self.HostInterface.output( - '[from: %s] connection closed (%s of %s max connections)' % ( - client_address, - self.HostInterface.connections - 1, - self.HostInterface.max_connections, - ) - ) - per_host_connections = self.HostInterface.per_host_connections - conn_data = per_host_connections.get(client_address[0]) - if conn_data is not None: - if conn_data < 1: - del per_host_connections[client_address[0]] - else: - per_host_connections[client_address[0]] -= 1 - - def close_request(self, request): - if self.HostInterface.connections > 0: - self.HostInterface.connections -= 1 - - class RequestHandler(socketserver.BaseRequestHandler): - - import select - import socket - import gc - timed_out = False - - def __init__(self, request, client_address, server): - - # pre-init attribues - self.__DEBUG = False - self.__buffered_data = None - self.__inst_token = entropy.tools.get_random_number() - self.server = None - self.request = None - self.client_address = None - socketserver.BaseRequestHandler.__init__(self, request, - client_address, server) - self.__data_counter = None - - def _data_receiver(self, ssl_enabled, ssl_exceptions, eos, - max_command_length): - - if self.timed_out: - return True - self.timed_out = True - try: - ready_to_read, ready_to_write, in_error = select.select( - [self.request], [], [], self.default_timeout) - except KeyboardInterrupt: - self.timed_out = True - return True - - buf_len = 16384 - if len(ready_to_read) == 1 and ready_to_read[0] == self.request: - - self.timed_out = False - # for ValueError exception trapping: - data = None - - if self.__DEBUG: - self.server.processor.HostInterface.output( - '[from: %s] request arrived :: counter: %s | buf_data: %s' % ( - self.client_address, - self.__data_counter, - len(self.__buffered_data), - ) - ) - - try: - - if ssl_enabled and hasattr(self.request, 'setblocking'): - # set SSL socket in blocking mode - # this fixes bugs related to data stream flooding - # with SSL - pyOpenSSL, probably because handshake - # and WantRead/WantWrite bullshit is handled - # automatically - self.request.setblocking(True) - - data = self.request.recv(buf_len) - if ssl_enabled: - while self.request.pending(): - data += self.request.recv(buf_len) - - if self.__data_counter is None: - if not data: # client wants to close - return True - elif data == self.server.processor.HostInterface.answers['noop']: - return False - elif len(data) < len(eos): - self.server.processor.HostInterface.output( - 'interrupted: %s, reason: %s - from client: %s - data: "%s" - counter: %s' % ( - self.server.server_address, - "malformed EOS", - self.client_address, - repr(data), - self.__data_counter, - ) - ) - self.__buffered_data = const_convert_to_rawstring('') - return True - mystrlen = data.split(eos)[0] - self.__data_counter = int(mystrlen) - data = data[len(mystrlen)+1:] - self.__data_counter -= len(data) - self.__buffered_data += data - - # command length exceeds our command length limit - if self.__data_counter > max_command_length: - raise InterruptError( - 'InterruptError: command too long: %s, limit: %s' % ( - self.__data_counter, max_command_length,)) - - buf_empty_watchdog_count = 50 # * 0.05 = 2,5 seconds - while self.__data_counter > 0: - data_buf = buf_len - if self.__data_counter < buf_len: - data_buf = self.__data_counter - if ssl_enabled: - x = self.request.recv(data_buf) - else: - x = self.request.recv(data_buf) - xlen = len(x) - self.__data_counter -= xlen - self.__buffered_data += x - # if we did not receive a shit and we still - # need some data, trigger the watchdog - if (xlen == 0) and (self.__data_counter > 0): - buf_empty_watchdog_count -= 1 - time.sleep(0.05) - if buf_empty_watchdog_count < 1: - raise ValueError( - "buffer counter watchdog trigger") - - self.__data_counter = None - except ValueError: - tb = entropy.tools.get_traceback() - print(tb) - self.server.processor.HostInterface.socketLog.write(tb) - self.server.processor.HostInterface.socketLog.write(repr(data)) - self.server.processor.HostInterface.socketLog.write(repr(self)) - self.server.processor.HostInterface.socketLog.write(repr(self.__inst_token)) - self.server.processor.HostInterface.output( - 'interrupted: %s, reason: %s - from client: %s' % ( - self.server.server_address, - "malformed transmission", - self.client_address, - ) - ) - return True - except self.socket.timeout as e: - self.server.processor.HostInterface.output( - 'interrupted: %s, reason: %s - from client: %s' % ( - self.server.server_address, - e, - self.client_address, - ) - ) - return True - except self.socket.sslerror as e: - self.server.processor.HostInterface.output( - 'interrupted: %s, SSL socket error reason: %s - from client: %s' % ( - self.server.server_address, - e, - self.client_address, - ) - ) - return True - except ssl_exceptions['WantX509LookupError']: - return False - except ssl_exceptions['WantReadError']: - self.server.processor.HostInterface._ssl_poll( - self.request, select.POLLIN, 'read') - return False - except ssl_exceptions['WantWriteError']: - self.server.processor.HostInterface._ssl_poll( - self.request, select.POLLOUT, 'read') - return False - except ssl_exceptions['ZeroReturnError']: - return True - except ssl_exceptions['Error'] as e: - self.server.processor.HostInterface.output( - 'interrupted: SSL Error, reason: %s - from client: %s' % ( - e, - self.client_address, - ) - ) - return True - except InterruptError as e: - self.server.processor.HostInterface.output( - 'interrupted: Command Error, reason: %s - from client: %s' % ( - e, - self.client_address, - ) - ) - return True - - if not self.__buffered_data: - return True - - if etpUi['debug']: - const_debug_write(__name__, darkred("=== recv ======== \\")) - const_debug_write(__name__, darkred(repr(self.__buffered_data))) - const_debug_write(__name__, - darkred("=== recv[%s] ======== /" % (len(self.__buffered_data),))) - - cmd = self.server.processor.process(self.__buffered_data, self.request, self.client_address) - if cmd == 'close': - # send KAPUTT signal JA! - self.server.processor.transmit(self.server.processor.HostInterface.answers['cl']) - return True - self.__buffered_data = const_convert_to_rawstring('') - return False - - def fork_lock_acquire(self): - if hasattr(self.server.processor.HostInterface, 'ForkLock'): - x = getattr(self.server.processor.HostInterface, 'ForkLock') - if hasattr(x, 'acquire') and hasattr(x, 'release') and hasattr(x, 'locked'): - x.acquire() - - def fork_lock_release(self): - if hasattr(self.server.processor.HostInterface, 'ForkLock'): - x = getattr(self.server.processor.HostInterface, 'ForkLock') - if hasattr(x, 'acquire') and hasattr(x, 'release') and hasattr(x, 'locked'): - if x.locked(): - x.release() - - def handle(self): - # not using spawnFunction because it causes some mess - # forking this way avoids having memory leaks - if self.server.processor.HostInterface.fork_requests: - self.fork_lock_acquire() - try: - my_timeout = self.server.processor.HostInterface.fork_request_timeout_seconds - pid = os.fork() - seconds = 0 - if pid > 0: # parent here - # pid killer after timeout - passed_away = False - while True: - time.sleep(1) - seconds += 1 - try: - dead = os.waitpid(pid, os.WNOHANG)[0] - except OSError as e: - if e.errno != errno.ECHILD: - raise - dead = True - if passed_away: - break - if dead: - break - if seconds > my_timeout: - self.server.processor.HostInterface.output( - 'interrupted: forked request timeout: %s,%s from client: %s' % ( - seconds, - dead, - self.client_address, - ) - ) - if not dead: - import signal - os.kill(pid, signal.SIGKILL) - passed_away = True # in this way, the process table should be clean - continue - break - else: - self.do_handle() - os._exit(0) - finally: - self.fork_lock_release() - else: - self.do_handle() - #entropy.tools.spawn_function(self.do_handle) - - def do_handle(self): - - self.default_timeout = self.server.processor.HostInterface.timeout - ssl_enabled = self.server.processor.HostInterface.is_ssl_enabled() - ssl_exceptions = self.server.processor.HostInterface.SSL_exceptions - eos = self.server.processor.HostInterface.answers['eos'] - max_command_length = \ - self.server.processor.HostInterface.max_command_length - - while True: - - try: - if self.__DEBUG: - self.server.processor.HostInterface.output( - '[from: %s] calling data_receiver' % ( - self.client_address, - ) - ) - dobreak = self._data_receiver(ssl_enabled, ssl_exceptions, - eos, max_command_length) - if self.__DEBUG: - self.server.processor.HostInterface.output( - '[from: %s] quitting data_receiver :: dobreak: %s' % ( - self.client_address, - dobreak, - ) - ) - if dobreak: - break - except Exception as e: - self.server.processor.HostInterface.output( - 'interrupted: Unhandled exception: %s, error: %s - from client: %s' % ( - Exception, - e, - self.client_address, - ) - ) - # print exception - tb = entropy.tools.get_traceback() - print(tb) - self.server.processor.HostInterface.socketLog.write(tb) - break - - self.request.close() - - def setup(self): - self.__data_counter = None - self.__buffered_data = '' - - - class CommandProcessor: - - import socket - import gc - - def __init__(self, HostInterface): - self.HostInterface = HostInterface - self.channel = None - - def handle_termination_commands(self, data): - if data.strip() in self.HostInterface.termination_commands: - self.HostInterface.output('close: %s' % (self.client_address,)) - self.transmit(self.HostInterface.answers['cl']) - return "close" - - if not data.strip(): - return "ignore" - - def handle_command_string(self, string): - # validate command - args = string.strip().split(" ") - session = args[0] - if (session in self.HostInterface.initialization_commands) or \ - (session in self.HostInterface.no_session_commands) or \ - len(args) < 2: - cmd = args[0] - session = None - else: - cmd = args[1] - args = args[1:] # remove session - - stream_enabled = False - if (session is not None) and session in self.HostInterface.sessions: - stream_enabled = self.HostInterface.sessions[session].get('stream_mode') - - if stream_enabled and (cmd not in self.HostInterface.config_commands): - session_len = 0 - if session: - session_len = len(session)+1 - return cmd, [string[session_len+len(cmd)+1:]], session - else: - myargs = [] - if len(args) > 1: - myargs = args[1:] - - return cmd, myargs, session - - def handle_end_answer(self, cmd, whoops, valid_cmd): - if not valid_cmd: - self.transmit(self.HostInterface.answers['no']) - elif whoops: - self.transmit(self.HostInterface.answers['er']) - elif cmd not in self.HostInterface.no_acked_commands: - self.transmit(self.HostInterface.answers['ok']) - - def validate_command(self, cmd, args, session): - - # answer to invalid commands - if (cmd not in self.HostInterface.valid_commands): - return False, "not a valid command" - - if session == None: - if cmd not in self.HostInterface.no_session_commands: - return False, "need a valid session" - elif session not in self.HostInterface.sessions: - return False, "session is not alive" - - # check if command needs authentication - if session is not None: - auth = self.HostInterface.valid_commands[cmd]['auth'] - if auth: - # are we? - authed = self.HostInterface.sessions[session]['auth_uid'] - if authed == None: - # nope - return False, "not authenticated" - - # keep session alive - if session is not None: - self.HostInterface.set_session_running(session) - self.HostInterface.update_session_time(session) - - return True, "all good" - - def load_authenticator(self): - f, args, kwargs = self.HostInterface.AuthenticatorInst - myinst = f(*args, **kwargs) - return myinst - - def load_service_interface(self, session): - - uid = None - if session is not None: - uid = self.HostInterface.sessions[session]['auth_uid'] - - intf = self.HostInterface.EntropyInstantiation[0] - args = self.HostInterface.EntropyInstantiation[1] - kwds = self.HostInterface.EntropyInstantiation[2] - return intf(*args, **kwds) - - def process(self, data, channel, client_address): - - self.channel = channel - self.client_address = client_address - - term = self.handle_termination_commands(data) - if term: - return term - - cmd, args, session = self.handle_command_string(data) - valid_cmd, reason = self.validate_command(cmd, args, session) - - # decide if we need to load authenticator or Entropy - authenticator = None - cmd_data = self.HostInterface.valid_commands.get(cmd) - if not isinstance(cmd_data, dict): - self.HostInterface.output( - '[from: %s] command error: invalid command: %s' % ( - self.client_address, - cmd, - ) - ) - return "close" - elif (("authenticator" in cmd_data['args']) or (cmd in self.HostInterface.login_pass_commands)): - try: - authenticator = self.load_authenticator() - except ServiceConnectionError as e: - self.HostInterface.output( - '[from: %s] authenticator error: cannot load: %s' % ( - self.client_address, - e, - ) - ) - tb = entropy.tools.get_traceback() - print(tb) - self.HostInterface.socketLog.write(tb) - return "close" - except Exception as e: - self.HostInterface.output( - '[from: %s] authenticator error: cannot load: %s - unknown error' % ( - self.client_address, - e, - ) - ) - tb = entropy.tools.get_traceback() - print(tb) - self.HostInterface.socketLog.write(tb) - return "close" - - p_args = args - if (cmd in self.HostInterface.login_pass_commands) and authenticator is not None: - p_args = authenticator.hide_login_data(p_args) - elif cmd in self.HostInterface.raw_commands: - p_args = ['raw data'] - data_len = len(data) - - # beautify interface if debug mode is on - out_cmd = cmd - out_data_len = data_len - out_session = session - out_valid_cmd = valid_cmd - out_reason = reason - out_args = p_args - - # for god sake, do not spam the whole log file or stdout - if not etpUi['debug'] and len(out_args) > 30: - out_args = out_args[:29] + ["---truncated---"] - if etpUi['debug']: - out_cmd = purple(out_cmd) - out_data_len = brown(str(out_data_len)) - out_session = darkgreen(str(out_session)) - out_valid_cmd = darkgreen(str(out_valid_cmd)) - out_reason = darkblue(str(out_reason)) - out_args = brown(str(out_args)) - - self.HostInterface.output( - '[from: %s] command validation :: called %s: length: %s, args: ' - '%s, session: %s, valid: %s, reason: %s' % ( - self.client_address, - out_cmd, - out_data_len, - out_args, - out_session, - out_valid_cmd, - out_reason, - ) - ) - - whoops = False - if valid_cmd: - - if authenticator is not None: - # now set session - authenticator.set_session(session) - - Entropy = None - if "Entropy" in cmd_data['args']: - Entropy = self.load_service_interface(session) - try: - run_task_out = self.run_task(cmd, args, session, Entropy, authenticator) - if etpUi['debug']: - self.HostInterface.output( - '[from: %s] command executed: result %s' % ( - self.client_address, - run_task_out, - ) - ) - except self.socket.timeout: - self.HostInterface.output( - '[from: %s] command error: timeout, closing connection' % ( - self.client_address, - ) - ) - # close connection - del authenticator - del Entropy - return "close" - except self.socket.error as e: - self.HostInterface.output( - '[from: %s] command error: socket error: %s' % ( - self.client_address, - e, - ) - ) - # close connection - del authenticator - del Entropy - return "close" - except self.HostInterface.SSL_exceptions['SysCallError'] as e: - self.HostInterface.output( - '[from: %s] command error: SSL SysCallError: %s' % ( - self.client_address, - e, - ) - ) - # close connection - del authenticator - del Entropy - return "close" - except Exception as e: - # write to self.HostInterface.socketLog - tb = entropy.tools.get_traceback() - print(tb) - self.HostInterface.socketLog.write(tb) - # store error - self.HostInterface.output( - '[from: %s] command error: %s, type: %s' % ( - self.client_address, - e, - type(e), - ) - ) - if session is not None: - self.HostInterface.store_rc(str(e), session) - whoops = True - - del Entropy - - if session is not None: - self.HostInterface.update_session_time(session) - self.HostInterface.unset_session_running(session) - rcmd = None - try: - self.handle_end_answer(cmd, whoops, valid_cmd) - except (self.socket.error, self.socket.timeout, self.HostInterface.SSL_exceptions['SysCallError'],): - rcmd = "close" - - if authenticator is not None: - authenticator.terminate_instance() - del authenticator - if not self.HostInterface.fork_requests: - self.gc.collect() - return rcmd - - def transmit(self, data): - - if etpUi['debug']: - const_debug_write(__name__, darkblue("=== send ======== \\")) - const_debug_write(__name__, darkblue(repr(data))) - const_debug_write(__name__, - darkblue("=== send[%s] ======== /" % (len(data),))) - - self.HostInterface.transmit(self.channel, data) - - def run_task(self, cmd, args, session, Entropy, authenticator): - - p_args = args - if cmd in self.HostInterface.login_pass_commands: - p_args = authenticator.hide_login_data(p_args) - elif cmd in self.HostInterface.raw_commands: - p_args = ['raw data'] - # for god sake, do not spam the whole log file or stdout - if not etpUi['debug'] and len(p_args) > 30: - p_args = p_args[:29] + ["---truncated---"] - - self.HostInterface.output( - '[from: %s] run_task :: called %s: args: %s, session: %s' % ( - self.client_address, - cmd, - p_args, - session, - ) - ) - - myargs = args - mykwargs = {} - if cmd not in self.HostInterface.raw_commands: - myargs, mykwargs = self._get_args_kwargs(args) - - rc = self.spawn_function(cmd, myargs, mykwargs, session, Entropy, authenticator) - if session is not None and session in self.HostInterface.sessions: - self.HostInterface.store_rc(rc, session) - return rc - - def _get_args_kwargs(self, args): - myargs = [] - mykwargs = {} - - def is_int(x): - try: - int(x) - except ValueError: - return False - return True - - for arg in args: - if (arg.find("=") != -1) and not arg.startswith("="): - x = arg.split("=") - a = x[0] - b = ''.join(x[1:]) - if (b in ("True", "False",)) or is_int(b): - mykwargs[a] = eval(b) - else: - myargs.append(arg) - else: - if (arg in ("True", "False",)) or is_int(arg): - myargs.append(eval(arg)) - else: - myargs.append(arg) - return myargs, mykwargs - - def spawn_function(self, cmd, myargs, mykwargs, session, Entropy, authenticator): - - p_args = myargs - if cmd in self.HostInterface.login_pass_commands: - p_args = authenticator.hide_login_data(p_args) - elif cmd in self.HostInterface.raw_commands: - p_args = ['raw data'] - # for god sake, do not spam the whole log file or stdout - if not etpUi['debug'] and len(p_args) > 30: - p_args = p_args[:29] + ["---truncated---"] - - self.HostInterface.output( - '[from: %s] called %s: args: %s, kwargs: %s' % ( - self.client_address, - cmd, - p_args, - mykwargs, - ) - ) - return self.do_spawn(cmd, myargs, mykwargs, session, Entropy, authenticator) - - def do_spawn(self, cmd, myargs, mykwargs, session, Entropy, authenticator): - - cmd_data = self.HostInterface.valid_commands.get(cmd) - do_fork = cmd_data['as_user'] - f = cmd_data['cb'] - func_args = [] - for arg in cmd_data['args']: - try: - func_args.append(eval(arg)) - except (NameError, SyntaxError): - func_args.append(str(arg)) - - if do_fork: - myfargs = func_args[:] - myfargs.extend(myargs) - return self.fork_task(f, session, authenticator, *myfargs, **mykwargs) - else: - return f(*func_args) - - def fork_task(self, f, session, authenticator, *args, **kwargs): - gid = None - uid = None - if session is not None: - logged_in = self.HostInterface.sessions[session]['auth_uid'] - if logged_in is not None: - uid = logged_in - gid = etpConst['entropygid'] - return entropy.tools.spawn_function(self._do_fork, f, authenticator, uid, gid, *args, **kwargs) - - def _do_fork(self, f, authenticator, uid, gid, *args, **kwargs): - authenticator.set_exc_permissions(uid, gid) - rc = f(*args, **kwargs) - return rc - - class BuiltInCommands(SocketCommands): - - import zlib - - def __init__(self, HostInterface): - - SocketCommands.__init__(self, HostInterface, inst_name = "builtin") - - self.valid_commands = { - 'begin': { - 'auth': False, # does it need authentication ? - 'built_in': True, # is it built-in ? - 'cb': self.docmd_begin, # function to call - 'args': ["self.transmit", "self.client_address"], # arguments to be passed before *args and **kwards, in SocketHostInterface.do_spawn() - 'as_user': False, # do I have to fork the process and run it as logged user? - # needs auth = True - 'desc': "instantiate a session", # description - 'syntax': "begin", # syntax - 'from': str(self), # from what class - }, - 'end': { - 'auth': False, - 'built_in': True, - 'cb': self.docmd_end, - 'args': ["self.transmit", "session"], - 'as_user': False, - 'desc': "end a session", - 'syntax': " end", - 'from': str(self), - }, - 'session_config': { - 'auth': False, - 'built_in': True, - 'cb': self.docmd_session_config, - 'args': ["session", "myargs"], - 'as_user': False, - 'desc': "set session configuration options", - 'syntax': " session_config