besmarts.core.compute module

besmarts.core.compute

Responsible for setting up and distributing large compute jobs

Architecture-wise, it interfaces the multiprocessing module

besmarts.core.compute.AutoProxy(token, serializer, manager=None, authkey=None, exposed=None, incref=True, manager_owned=False)[source]

Return an auto-proxy for token

class besmarts.core.compute.BaseProxy(*args, **kwds)[source]

Bases: BaseProxy

besmarts.core.compute.Client(address, family=None, authkey=None, timeout=240)[source]

Returns a connection to the address of a Listener

class besmarts.core.compute.Connection(handle, readable=True, writable=True)[source]

Bases: Connection

Connection class based on an arbitrary file descriptor (Unix only), or a socket handle (Windows).

class besmarts.core.compute.Listener(*args, **kwargs)[source]

Bases: Listener

accept()[source]

Accept a connection on the bound socket or named pipe of self.

Returns a Connection object.

besmarts.core.compute.MakeProxyType(name, exposed, _cache={})[source]

Return a proxy type whose methods are given by exposed

besmarts.core.compute.Process(obj, *args, **kwds)[source]
class besmarts.core.compute.Server(*args, **kwargs)[source]

Bases: Server

accept_connection(c, name)[source]

Spawn a new thread to serve this connection

accepter()[source]
decref(c, ident)[source]
handle_request(conn)[source]

Handle a new connection

incref(c, ident)[source]
recv_request(conn, out, err)[source]
run_request(c, request, msg)[source]
send_request(c, msg)[source]
serve_client(conn)[source]

Handle requests from the proxies in a particular process/thread

serve_forever()[source]

Run the server forever

besmarts.core.compute.SocketClient(address, timeout=240)[source]

Return a connection object connected to the socket given by address

besmarts.core.compute.close_pools()[source]
besmarts.core.compute.close_workspaces()[source]
besmarts.core.compute.compute_remote(addr, port, processes=1, queue_size=1)[source]
besmarts.core.compute.dispatch(c, id, methodname, args=(), kwds={})[source]

Send a message to manager using connection c and return response

besmarts.core.compute.dprint(*args, **kwds)[source]
besmarts.core.compute.manager_connect(mgr: BaseManager, success: Event)[source]
besmarts.core.compute.manager_remote_get_iqueue(mgr, timeout=240)[source]
besmarts.core.compute.manager_remote_get_iqueue_thread(mgr, out)[source]
besmarts.core.compute.manager_remote_get_oqueue(mgr, timeout=240)[source]
besmarts.core.compute.manager_remote_get_oqueue_thread(mgr, out)[source]
besmarts.core.compute.manager_remote_get_state(wq, timeout=240)[source]
besmarts.core.compute.manager_remote_get_state_thread(ws, out)[source]
besmarts.core.compute.manager_remote_get_status(wq, timeout=240)[source]
besmarts.core.compute.manager_remote_get_status_thread(ws, out)[source]
besmarts.core.compute.manager_remote_queue_put(oq, obj, timeout=240, n=1)[source]
besmarts.core.compute.manager_remote_queue_put_thread(oq, obj, n, success)[source]
besmarts.core.compute.manager_remote_queue_qsize(q, timeout=240)[source]
besmarts.core.compute.manager_remote_queue_qsize_thread(q, out)[source]
class besmarts.core.compute.myiqueue(maxsize=0)[source]

Bases: Queue

get(block=True, timeout=240)[source]

Remove and return an item from the queue.

If optional args ‘block’ is true and ‘timeout’ is None (the default), block if necessary until an item is available. If ‘timeout’ is a non-negative number, it blocks at most ‘timeout’ seconds and raises the Empty exception if no item was available within that time. Otherwise (‘block’ is false), return an item if one is immediately available, else raise the Empty exception (‘timeout’ is ignored in that case).

class besmarts.core.compute.myqueue(maxsize=0)[source]

Bases: Queue

get(block=True, timeout=240, n=1)[source]

Remove and return an item from the queue.

If optional args ‘block’ is true and ‘timeout’ is None (the default), block if necessary until an item is available. If ‘timeout’ is a non-negative number, it blocks at most ‘timeout’ seconds and raises the Empty exception if no item was available within that time. Otherwise (‘block’ is false), return an item if one is immediately available, else raise the Empty exception (‘timeout’ is ignored in that case).

put(item, block=True, timeout=240, n=1)[source]

Put an item into the queue.

If optional args ‘block’ is true and ‘timeout’ is None (the default), block if necessary until a free slot is available. If ‘timeout’ is a non-negative number, it blocks at most ‘timeout’ seconds and raises the Full exception if no free slot was available within that time. Otherwise (‘block’ is false), put an item on the queue if a free slot is immediately available, else raise the Full exception (‘timeout’ is ignored in that case).

besmarts.core.compute.queue_get_nowait(q, block=False, timeout=240, n=1)[source]
besmarts.core.compute.queue_get_nowait_thread(q, out, n, block, timeout)[source]
besmarts.core.compute.register_pool(p)[source]
besmarts.core.compute.register_workspace(p)[source]
besmarts.core.compute.remote_connect(mgr, timeout=240)[source]
besmarts.core.compute.remote_init_thread(remote_init, shm_proxy, out)[source]
besmarts.core.compute.shm_init(proxy)[source]
class besmarts.core.compute.shm_local(procs_per_task=1, data=None)[source]

Bases: object

get()[source]
remote_init()[source]
besmarts.core.compute.signal_kill_processes(sig, frame)[source]
besmarts.core.compute.thread_name_set(name)[source]
besmarts.core.compute.unregister_pool(p)[source]
besmarts.core.compute.unregister_workspace(p)[source]
class besmarts.core.compute.workqueue(addr, port)[source]

Bases: object

get_state()[source]
get_status()[source]
get_workspaces()[source]
besmarts.core.compute.workqueue_get_workspace(wq: workqueue_remote, addr, port) workspace[source]
besmarts.core.compute.workqueue_is_active(wq)[source]
besmarts.core.compute.workqueue_list_workspaces(wq: workqueue_remote) Dict[source]
class besmarts.core.compute.workqueue_local(addr, port)[source]

Bases: workqueue

close()[source]
get_state()[source]
get_threads()[source]
get_workspaces()[source]
put_workspaces(wss)[source]
remove_workspace(ws)[source]
class besmarts.core.compute.workqueue_manager(address=None, authkey=None, serializer='pickle', ctx=None, *, shutdown_timeout=1.0)[source]

Bases: SyncManager

get_iqueue() Queue[source]
get_oqueue() Queue[source]
get_rqueue() Queue[source]
get_state() Dict[source]
besmarts.core.compute.workqueue_new_workspace(wq: workqueue_local, address=None, shm=None, nproc=-1)[source]
besmarts.core.compute.workqueue_push_workspace(wq: workqueue_local, ws: workspace)[source]
class besmarts.core.compute.workqueue_remote(addr, port)[source]

Bases: workqueue

connect()[source]
get_state()[source]
get_status()[source]
get_workspaces()[source]
put_workspaces(wss)[source]
besmarts.core.compute.workqueue_remote_get_workspaces(wq, timeout=240)[source]
besmarts.core.compute.workqueue_remote_get_workspaces_thread(wss, out, success)[source]
besmarts.core.compute.workqueue_remote_is_active(wq)[source]
besmarts.core.compute.workqueue_remote_put_workspaces(wq, wss, timeout=240)[source]
besmarts.core.compute.workqueue_remote_put_workspaces_thread(mgr, inp, success)[source]
besmarts.core.compute.workqueue_remove_workspace(wq: workqueue_local, ws: workspace_local)[source]
class besmarts.core.compute.workspace(addr, port)[source]

Bases: object

get_iqueue()[source]
get_oqueue()[source]
get_rqueue()[source]
get_state()[source]
get_status()[source]
besmarts.core.compute.workspace_flush(ws: workspace_local, indices, timeout: float = 240, maxwait=None, verbose=True)[source]
besmarts.core.compute.workspace_is_active(ws: workspace_remote)[source]
class besmarts.core.compute.workspace_local(addr, port, shm: shm_local = None, nproc=-1)[source]

Bases: workspace

Assumes that we are process-local to all needed resources and do not need to use the proxy interface i.e. no connection needed

clear_iqueue()[source]
clear_oqueue()[source]
close()[source]
get_iqueue()[source]
get_oqueue()[source]
get_state()[source]
get_status()[source]
manager_start()[source]
pool_close()[source]
pool_start()[source]
reset()[source]
set_status(status)[source]
start()[source]
besmarts.core.compute.workspace_local_remote_gather_thread(ws: workspace_local)[source]

pull results from the remote workers and put them in the local queue

besmarts.core.compute.workspace_local_remote_loadbalance_thread(ws: workspace_local)[source]

pull results from the remote workers and put them in the local queue

besmarts.core.compute.workspace_local_run(ws: workspace_local)[source]

Take jobs from the input queue and distribute to the processing queues, which can be a (low latency) local pool or a (high latency) manager

besmarts.core.compute.workspace_local_run_thread(ws: workspace_local)[source]
besmarts.core.compute.workspace_local_submit(ws, work)[source]
class besmarts.core.compute.workspace_manager(address=None, authkey=None, serializer='pickle', ctx=None, *, shutdown_timeout=1.0)[source]

Bases: SyncManager

create(*args, **kwargs)[source]
get_state() Dict[source]
get_status() workspace_status[source]
get_workspaces() Dict[source]
class besmarts.core.compute.workspace_pool(*args, **kwds)[source]

Bases: Pool

Process(*args, **kwds)
class besmarts.core.compute.workspace_remote(addr, port, nproc=1)[source]

Bases: workspace

Assumes that we need to go through the proxy interface to access resources that are not local i.e. need to go through a connection

close()[source]
connect()[source]
get_iqueue()[source]
get_oqueue()[source]
get_state()[source]
get_status()[source]
iqueue_get(n=1)[source]
iqueue_put(obj, n=1)[source]
oqueue_put(obj, n=1)[source]
start(input_queue_size=2)[source]
besmarts.core.compute.workspace_remote_compute(wq: workqueue_remote, ws: workspace_remote)[source]
besmarts.core.compute.workspace_remote_local_gather_thread(ws: workspace_remote)[source]

pull results from the remote workers and put them in the local queue

besmarts.core.compute.workspace_remote_local_pusher_thread(ws: workspace_remote)[source]

pull results from the remote workers and put them in the local queue

besmarts.core.compute.workspace_remote_shm_init(ws, timeout=240)[source]
besmarts.core.compute.workspace_remote_shm_init_thread(mgr, out)[source]
besmarts.core.compute.workspace_run(distfun: Tuple[Callable, Sequence, Mapping], workspace_address=None)[source]
besmarts.core.compute.workspace_run_init(shm, t0=None)[source]
class besmarts.core.compute.workspace_status[source]

Bases: object

DONE = 5
EMPTY = 0
INACTIVE = 1
INVALID = -1
RUNNING = 4
SUBMITTING = 2
WAITING = 3
besmarts.core.compute.workspace_submit_and_flush(ws, fn, iterable: Dict, chunksize=1, timeout=0.0, batchsize=0, verbose=False, clear=True) Dict[source]