Commit 6c428bc9 authored by Jürg Billeter's avatar Jürg Billeter
Browse files

Merge branch 'raoul/cas-refactor' into 'master'

Cas refactor

See merge request BuildStream/buildstream!1071
parents d6587aa0 76f67483
Loading
Loading
Loading
Loading
Loading
+8 −22
Original line number Original line Diff line number Diff line
@@ -19,18 +19,16 @@


import multiprocessing
import multiprocessing
import os
import os
import signal
import string
import string
from collections.abc import Mapping
from collections.abc import Mapping


from ..types import _KeyStrength
from .types import _KeyStrength
from .._exceptions import ArtifactError, CASError, LoadError, LoadErrorReason
from ._exceptions import ArtifactError, CASError, LoadError, LoadErrorReason
from .._message import Message, MessageType
from ._message import Message, MessageType
from .. import _signals
from . import utils
from .. import utils
from . import _yaml
from .. import _yaml


from .cascache import CASRemote, CASRemoteSpec
from ._cas import CASRemote, CASRemoteSpec




CACHE_SIZE_FILE = "cache_size"
CACHE_SIZE_FILE = "cache_size"
@@ -375,20 +373,8 @@ class ArtifactCache():
        remotes = {}
        remotes = {}
        q = multiprocessing.Queue()
        q = multiprocessing.Queue()
        for remote_spec in remote_specs:
        for remote_spec in remote_specs:
            # Use subprocess to avoid creation of gRPC threads in main BuildStream process
            # See https://github.com/grpc/grpc/blob/master/doc/fork_support.md for details
            p = multiprocessing.Process(target=self.cas.initialize_remote, args=(remote_spec, q))


            try:
            error = CASRemote.check_remote(remote_spec, q)
                # Keep SIGINT blocked in the child process
                with _signals.blocked([signal.SIGINT], ignore=False):
                    p.start()

                error = q.get()
                p.join()
            except KeyboardInterrupt:
                utils._kill_process_tree(p.pid)
                raise


            if error and on_failure:
            if error and on_failure:
                on_failure(remote_spec.url, error)
                on_failure(remote_spec.url, error)
@@ -747,7 +733,7 @@ class ArtifactCache():
                                "servers are configured as push remotes.")
                                "servers are configured as push remotes.")


        for remote in push_remotes:
        for remote in push_remotes:
            message_digest = self.cas.push_message(remote, message)
            message_digest = remote.push_message(message)


        return message_digest
        return message_digest


+2 −1
Original line number Original line Diff line number Diff line
@@ -17,4 +17,5 @@
#  Authors:
#  Authors:
#        Tristan Van Berkom <tristan.vanberkom@codethink.co.uk>
#        Tristan Van Berkom <tristan.vanberkom@codethink.co.uk>


from .artifactcache import ArtifactCache, ArtifactCacheSpec, CACHE_SIZE_FILE
from .cascache import CASCache
from .casremote import CASRemote, CASRemoteSpec
+20 −376

File changed and moved.

Preview size limit exceeded, changes collapsed.

+384 −0

File added.

Preview size limit exceeded, changes collapsed.

Loading