Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
asyncStream.py120 linesDownload Raw Back to strategy
1"""
2"""
3
4# Created on 2016.07.10
5#
6# Author: Giovanni Cannata
7#
8# Copyright 2016 - 2020 Giovanni Cannata
9#
10# This file is part of ldap3.
11#
12# ldap3 is free software: you can redistribute it and/or modify
13# it under the terms of the GNU Lesser General Public License as published
14# by the Free Software Foundation, either version 3 of the License, or
15# (at your option) any later version.
16#
17# ldap3 is distributed in the hope that it will be useful,
18# but WITHOUT ANY WARRANTY; without even the implied warranty of
19# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
20# GNU Lesser General Public License for more details.
21#
22# You should have received a copy of the GNU Lesser General Public License
23# along with ldap3 in the COPYING and COPYING.LESSER files.
24# If not, see <http://www.gnu.org/licenses/>.
25
26try:
27    from queue import Queue
28except ImportError:  # Python 2
29    # noinspection PyUnresolvedReferences
30    from Queue import Queue
31
32from io import StringIO
33from os import linesep
34
35from ..protocol.rfc2849 import decode_persistent_search_control
36from ..strategy.asynchronous import AsyncStrategy
37from ..core.exceptions import LDAPLDIFError
38from ..utils.conv import prepare_for_stream
39from ..protocol.rfc2849 import persistent_search_response_to_ldif, add_ldif_header
40
41
42# noinspection PyProtectedMember
43class AsyncStreamStrategy(AsyncStrategy):
44    """
45    This strategy is asynchronous. It streams responses in a generator as they appear in the self._responses container
46    """
47    def __init__(self, ldap_connection):
48        AsyncStrategy.__init__(self, ldap_connection)
49        self.can_stream = True
50        self.line_separator = linesep
51        self.all_base64 = False
52        self.stream = None
53        self.order = dict()
54        self._header_added = False
55        self.persistent_search_message_id = None
56        self.streaming = False
57        self.callback = None
58        if ldap_connection.pool_size:
59            self.events = Queue(ldap_connection.pool_size)
60        else:
61            self.events = Queue()
62
63        del self._requests  # remove _requests dict from Async Strategy
64
65    def _start_listen(self):
66        AsyncStrategy._start_listen(self)
67        if self.streaming:
68            if not self.stream or (isinstance(self.stream, StringIO) and self.stream.closed):
69                self.set_stream(StringIO())
70
71    def _stop_listen(self):
72        AsyncStrategy._stop_listen(self)
73        if self.streaming:
74            self.stream.close()
75
76    def accumulate_stream(self, message_id, change):
77        if message_id == self.persistent_search_message_id:
78            with self.async_lock:
79                self._responses[message_id] = []
80            if self.streaming:
81                if not self._header_added and self.stream.tell() == 0:
82                    header = add_ldif_header(['-'])[0]
83                    self.stream.write(prepare_for_stream(header + self.line_separator + self.line_separator))
84                ldif_lines = persistent_search_response_to_ldif(change)
85                if self.stream and ldif_lines and not self.connection.closed:
86                    fragment = self.line_separator.join(ldif_lines)
87                    if not self._header_added and self.stream.tell() == 0:
88                        self._header_added = True
89                        header = add_ldif_header(['-'])[0]
90                        self.stream.write(prepare_for_stream(header + self.line_separator + self.line_separator))
91                    self.stream.write(prepare_for_stream(fragment + self.line_separator + self.line_separator))
92            else:  # strategy is not streaming, events are added to a queue
93                notification = decode_persistent_search_control(change)
94                if notification:
95                    change.update(notification)
96                    del change['controls']['2.16.840.1.113730.3.4.7']
97                if not self.callback:
98                    self.events.put(change)
99                else:
100                    self.callback(change)
101
102    def get_stream(self):
103        if self.streaming:
104            return self.stream
105        return None
106
107    def set_stream(self, value):
108        error = False
109        try:
110            if not value.writable():
111                error = True
112        except (ValueError, AttributeError):
113            error = True
114
115        if error:
116            raise LDAPLDIFError('stream must be writable')
117
118        self.stream = value
119        self.streaming = True
120 
codekingpro/portable-devtools · Team Ai