o
    …>ag:  ã                   @   sª   d Z ddlZddlmZ ddlZddlmZ ddl	m
Z
 ddlmZ ddlmZ ddlmZ ddlZe d¡Ze e ¡ ¡ ed	d
dgƒZedg d¢ƒZG dd„ dƒZdS )z; Helper functions to communicate with replication servers.
é    N)Ú
namedtuple)Úceil)ÚMergeInputReader)Úio)ÚversionÚpyosmiumÚOsmosisStateÚsequenceÚ	timestampÚDownloadResult)ÚidÚreaderÚnewestc                   @   sz   e Zd ZdZddd„Zdd„ Zdd„ Zd d
d„Zd!dd„Z				d"dd„Z	d#dd„Z
d$dd„Zdd„ Zdd„ Zdd„ ZdS )%ÚReplicationServerz· Represents a server that publishes replication data. Replication
        change files allow to keep local OSM data up-to-date without downloading
        the full dataset again.
    úosc.gzc                 C   s   || _ || _t ¡ | _d S )N)ÚbaseurlÚ	diff_typeÚrequestsÚSessionÚsession)ÚselfÚurlr   © r   ú;/usr/lib/python3/dist-packages/osmium/replication/server.pyÚ__init__   s   zReplicationServer.__init__c                 C   s   dd  tj¡i}tj||d�S )Nz
User-Agentzpyosmium/{})Úheaders)Úformatr   Úpyosmium_releaseÚ
urlrequestÚRequest)r   r   r   r   r   r   Úmake_request    s   zReplicationServer.make_requestc                 C   s:   t ƒ }| ¡ D ]
}|d ||d < q| jj| ¡ |dd�S )aÚ   Download a resource from the given URL and return a byte sequence
            of the content.

            This method has no support for cookies or any special authentication
            methods. If you need these, you have to provide your own custom URL
            opener. The method has to return an object which supports the
            `read()` and `readline()` methods to access the content. Example::

                def my_open_url(self, url):
                    opener = urlrequest.build_opener()
                    opener.addheaders = [('X-Fancy-Header', 'important_content')]
                    return opener.open(url)

                svr = ReplicationServer()
                svr.open_url = my_open_url
        é   r   T)r   Ústream)ÚdictÚheader_itemsr   ÚgetÚget_full_url)r   r   r   Úhr   r   r   Úopen_url$   s   zReplicationServer.open_urlé   c                 C   sÒ   |d }|}|   ¡ }|du s||jkrdS tƒ }|dkr`||jkr`z|  |¡}W n   t d¡ d}Y t|ƒdkrA||kr@dS n|| || j¡8 }t d||d ¡ |d7 }|dkr`||jks!t	|d ||jƒS )aÚ   Create a MergeInputReader and download diffs starting with sequence
            id `start_id` into it. `max_size`
            restricts the number of diffs that are downloaded. The download
            stops as soon as either a diff cannot be downloaded or the
            unpacked data in memory exceeds `max_size` kB.

            If some data was downloaded, returns a namedtuple with three fields:
            `id` contains the sequence id of the last downloaded diff, `reader`
            contains the MergeInputReader with the data and `newest` is a
            sequence id of the most recent diff available.

            Returns None if there was an error during download or no new
            data was available.
        r)   Nr   z(Error during diff download. Bailing out.Ú z:Downloaded change %d. (%d kB available in download buffer)r!   )
Úget_state_infor	   r   Úget_diff_blockÚLOGÚdebugÚlenÚ
add_bufferr   r   )r   Ústart_idÚmax_sizeÚ	left_sizeÚ
current_idr   ÚrdÚdiffdatar   r   r   Úcollect_diffs:   s.   
ÿòzReplicationServer.collect_diffsr*   Tc                 C   s0   |   ||¡}|du rdS |jj|||d� |jS )a˜   Download diffs starting with sequence id `start_id`, merge them
            together and then apply them to handler `handler`. `max_size`
            restricts the number of diffs that are downloaded. The download
            stops as soon as either a diff cannot be downloaded or the
            unpacked data in memory exceeds `max_size` kB.

            If `idx` is set, a location cache will be created and applied to
            the way nodes. You should be aware that diff files usually do not
            contain the complete set of nodes when a way is modified. That means
            that you cannot just create a new location cache, apply it to a diff
            and expect to get complete way geometries back. Instead you need to
            do an initial data import using a persistent location cache to
            obtain a full set of node locations and then reuse this location
            cache here when applying diffs.

            Diffs may contain multiple versions of the same object when it was
            changed multiple times during the period covered by the diff. If
            `simplify` is set to False then all versions are returned. If it
            is True (the default) then only the most recent version will be
            sent to the handler.

            The function returns the sequence id of the last diff that was
            downloaded or None if the download failed completely.
        N)ÚidxÚsimplify)r7   r   Úapplyr   )r   Úhandlerr1   r2   r8   r9   Údiffsr   r   r   Úapply_diffsg   s
   zReplicationServer.apply_diffsNc                 C   s  |   ||¡}|du rdS t |¡}	|	 ¡ j}
t ¡ }|
|_|rC| d| j¡ | dt|j	ƒ¡ |  
|j	¡}|durC| d|j d¡¡ |durV| ¡ D ]
\}}| ||¡ qK|du r`t |¡}nt ||¡}|
|_t ||¡}t d¡ |j |	||
¡ |	 ¡  | ¡  |j	|jfS )a×   Download diffs starting with sequence id `start_id`, merge them
            with the data from the OSM file named `infile` and write the result
            into a file with the name `outfile`. The output file must not yet
            exist.

            `max_size` restricts the number of diffs that are downloaded. The
            download stops as soon as either a diff cannot be downloaded or the
            unpacked data in memory exceeds `max_size` kB.

            If `set_replication_header` is true then the URL of the replication
            server and the sequence id and timestamp of the last diff applied
            will be written into the `writer`. Note that this currently works
            only for the PBF format.

            `extra_headers` is a dict with additional header fields to be set.
            Most notably, the 'generator' can be set this way.

            `outformat` sets the format of the output file. If None, the format
            is determined from the file name.

            The function returns a tuple of last downloaded sequence id and
            newest available sequence id if new data has been written or None
            if no data was available or the download failed completely.
        NÚosmosis_replication_base_urlÚ#osmosis_replication_sequence_numberÚosmosis_replication_timestampz%Y-%m-%dT%H:%M:%SZzMerging changes into OSM file.)r7   ÚoioÚReaderÚheaderÚhas_multiple_object_versionsÚHeaderÚsetr   Ústrr   r+   r
   ÚstrftimeÚitemsÚFileÚWriterr-   r.   r   Úapply_to_readerÚcloser   )r   ÚinfileÚoutfiler1   r2   Úset_replication_headerÚextra_headersÚ	outformatr<   r   Úhas_historyr'   ÚinfoÚkÚvÚofÚwriterr   r   r   Úapply_diffs_to_file‰   s4   


z%ReplicationServer.apply_diffs_to_fileFc                 C   s  |   ¡ }|du r
dS ||jks|jdkr|jS d}d}|du rct d|¡ |   |¡}|durI|j|krI|jdks@|jd |jkrC|jS |}d}d}|du r_t||j d ƒ}||kr]|jS |}|du s	 |rqt|j|j d ƒ}n*|j|j  ¡ }|j|j }	||j  ¡ }
|jt|
|	 | ƒ }||jkr›|jd }|   |¡}|du rÃ|d }|du rÃ||jkrÃ|   |¡}|d8 }|du rÃ||jks±|du ræ|d }|du ræ||jk ræ|   |¡}|d7 }|du ræ||jk sÔ|du rí|jS |j|k rõ|}n|}|jd |jk�r|jS qd)aæ   Get the sequence number of the replication file that contains the
            given timestamp. The search algorithm is optimised for replication
            servers that publish updates in regular intervals. For servers
            with irregular change file publication dates 'balanced_search`
            should be set to true so that a standard binary search for the
            sequence will be used. The default is good for all known
            OSM replication services.
        Nr   zTrying with Id %sr!   é   )r+   r
   r	   r-   r.   ÚintÚtotal_secondsr   )r   r
   Úbalanced_searchÚupperÚlowerÚloweridÚnewidÚbase_splitidÚts_intÚseq_intÚgoalÚsplitÚsplitidr   r   r   Útimestamp_to_sequenceÊ   sh   
ï



þ
þ
Ýz'ReplicationServer.timestamp_to_sequencerZ   c           
      C   sl  t |d ƒD ]­}z|  |  |  |¡¡¡}W n ty2 } zt d|t|ƒ¡ W Y d}~ dS d}~ww d}d}t|dƒr@|j	}n|j
}|ƒ D ]\}| d¡}d|v r[|d| d¡… }n| ¡ }|r¢| dd	¡}	t|	ƒd	krq  dS |	d d
kr~t|	d ƒ}qF|	d dkr¢ztj |	d d¡}W n
 ty™   Y  n
w |jtjjd�}qF|dur³|dur³t||d�  S qdS )af   Downloads and returns the state information for the given
            sequence. If the download is successful, a namedtuple with
            `sequence` and `timestamp` is returned, otherwise the function
            returns `None`. `retries` sets the number of times the download
            is retried when pyosmium detects a truncated state file.
        r!   z%Loading state info %s failed with: %sNÚ
iter_lineszutf-8ú#r   ú=rZ   ÚsequenceNumberr
   z%Y-%m-%dT%H\:%M\:%SZ)Útzinfo)r	   r
   )Úranger(   r    Úget_state_urlÚ	Exceptionr-   r.   rG   Úhasattrri   ÚreadlineÚdecodeÚindexÚstriprf   r/   r[   ÚdtÚdatetimeÚstrptimeÚ
ValueErrorÚreplaceÚtimezoneÚutcr   )
r   ÚseqÚretriesÚ_ÚresponseÚerrÚtsÚget_lineÚlineÚkvr   r   r   r+     sH   €þ


ÿ€€z ReplicationServer.get_state_infoc                 C   s.   |   |  |  |¡¡¡}t|dƒr|jS | ¡ S )zö Downloads the diff with the given sequence number and returns
            it as a byte sequence. Throws a :code:`urllib.error.HTTPError`
            (or :code:`urllib2.HTTPError` in python2)
            if the file cannot be downloaded.
        Úcontent)r(   r    Úget_diff_urlrq   r†   Úread)r   r}   Úrespr   r   r   r,   K  s   
z ReplicationServer.get_diff_blockc                 C   s4   |du r	| j d S d| j |d |d d |d f S )zô Returns the URL of the state.txt files for a given sequence id.

            If seq is `None` the URL for the latest state info is returned,
            i.e. the state file in the root directory of the replication
            service.
        Nz
/state.txtz%s/%03i/%03i/%03i.state.txté@B éè  )r   ©r   r}   r   r   r   ro   Z  s
   
ÿzReplicationServer.get_state_urlc                 C   s&   d| j |d |d d |d | jf S )zE Returns the URL to the diff file for the given sequence id.
        z%s/%03i/%03i/%03i.%srŠ   r‹   )r   r   rŒ   r   r   r   r‡   h  s   þÿzReplicationServer.get_diff_url)r   )r)   )r)   r*   T)r)   TNN)F)NrZ   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r    r(   r7   r=   rY   rh   r+   r,   ro   r‡   r   r   r   r   r      s     


-"
þ
A
R/r   )r�   r   Úurllib.requestÚrequestr   rw   rv   Úcollectionsr   Úmathr   Úosmiumr   r   rA   r   ÚloggingÚ	getLoggerr-   Ú
addHandlerÚNullHandlerr   r   r   r   r   r   r   Ú<module>   s    
