Sync
This commit is contained in:
parent
6f050a029e
commit
776bba907c
3887 changed files with 59894 additions and 7280 deletions
|
|
@ -0,0 +1,70 @@
|
|||
(ns checkpoint.core
|
||||
(:gen-class)
|
||||
(:require [clojure.core.async :as async :refer [go <! >! <!! >!! alts! close!]]
|
||||
[clojure.string :as string]))
|
||||
|
||||
(defn coordinate [ctl-ch resp-ch combine]
|
||||
(go
|
||||
(<! (async/timeout 2000)) ;delay a bit to allow worker setup
|
||||
(loop [members {}, received {}] ;maps by in-channel of out-channels & received data resp.
|
||||
(let [rcvd-count (count received)
|
||||
release #(doseq [outch (vals members)] (go (>! outch %)))
|
||||
received (if (and (pos? rcvd-count) (= rcvd-count (count members)))
|
||||
(do (-> received vals combine release) {})
|
||||
received)
|
||||
[v ch] (alts! (cons ctl-ch (keys members)))]
|
||||
;receive a message on ctrl-ch or any member input channel
|
||||
(if (= ch ctl-ch)
|
||||
(let [[op inch outch] v] ;only a Checkpoint (see below) sends on ctl-ch
|
||||
(condp = op
|
||||
:join (do (>! resp-ch :ok)
|
||||
(recur (assoc members inch outch) received))
|
||||
:part (do (>! resp-ch :ok)
|
||||
(close! inch) (close! outch)
|
||||
(recur (dissoc members inch) (dissoc received inch)))
|
||||
:exit :exit))
|
||||
(if (nil? v) ;is the channel closed?
|
||||
(do
|
||||
(close! (get members ch))
|
||||
(recur (dissoc members ch) (dissoc received ch)))
|
||||
(recur members (assoc received ch v))))))))
|
||||
|
||||
(defprotocol ICheckpoint
|
||||
(join [this])
|
||||
(part [this inch outch]))
|
||||
|
||||
(deftype Checkpoint [ctl-ch resp-ch sync]
|
||||
ICheckpoint
|
||||
(join [this]
|
||||
(let [inch (async/chan), outch (async/chan 1)]
|
||||
(go
|
||||
(>! ctl-ch [:join inch outch])
|
||||
(<! resp-ch)
|
||||
[inch outch])))
|
||||
(part [this inch outch]
|
||||
(go
|
||||
(>! ctl-ch [:part inch outch]))))
|
||||
|
||||
(defn checkpoint [combine]
|
||||
(let [ctl-ch (async/chan), resp-ch (async/chan 1)]
|
||||
(->Checkpoint ctl-ch resp-ch (coordinate ctl-ch resp-ch combine))))
|
||||
|
||||
(defn worker
|
||||
([ckpt repeats] (worker ckpt repeats (fn [& args] nil)))
|
||||
([ckpt repeats mon]
|
||||
(go
|
||||
(let [[send recv] (<! (join ckpt))]
|
||||
(doseq [n (range repeats)]
|
||||
(<! (async/timeout (rand-int 5000)))
|
||||
(>! send n) (mon "sent" n)
|
||||
(<! recv) (mon "recvd"))
|
||||
(part ckpt send recv)))))
|
||||
|
||||
|
||||
(defn -main
|
||||
[& args]
|
||||
(let [ckpt (checkpoint identity)
|
||||
monitor (fn [id]
|
||||
(fn [& args] (println (apply str "worker" id ":" (string/join " " args)))))]
|
||||
(worker ckpt 10 (monitor 1))
|
||||
(worker ckpt 10 (monitor 2))))
|
||||
|
|
@ -0,0 +1,35 @@
|
|||
-module( checkpoint_synchronization ).
|
||||
|
||||
-export( [task/0] ).
|
||||
|
||||
task() ->
|
||||
Pid = erlang:spawn( fun() -> checkpoint_loop([], []) end ),
|
||||
[erlang:spawn(fun() -> random:seed(X, 1, 0), worker_loop(X, 3, Pid) end) || X <- lists:seq(1, 5)],
|
||||
erlang:exit( Pid, normal ).
|
||||
|
||||
|
||||
|
||||
checkpoint_loop( Assemblings, Completes ) ->
|
||||
receive
|
||||
{starting, Worker} -> checkpoint_loop( [Worker | Assemblings], Completes );
|
||||
{done, Worker} ->
|
||||
New_assemblings = lists:delete( Worker, Assemblings ),
|
||||
New_completes = checkpoint_loop_release( New_assemblings, [Worker | Completes] ),
|
||||
checkpoint_loop( New_assemblings, New_completes )
|
||||
end.
|
||||
|
||||
checkpoint_loop_release( [], Completes ) ->
|
||||
[X ! all_complete || X <- Completes],
|
||||
[];
|
||||
checkpoint_loop_release( _Assemblings, Completes ) -> Completes.
|
||||
|
||||
worker_loop( _Worker, 0, _Checkpoint ) -> ok;
|
||||
worker_loop( Worker, N, Checkpoint ) ->
|
||||
Checkpoint ! {starting, erlang:self()},
|
||||
io:fwrite( "Worker ~p ~p~n", [Worker, N] ),
|
||||
timer:sleep( random:uniform(100) ),
|
||||
Checkpoint ! {done, erlang:self()},
|
||||
receive
|
||||
all_complete -> ok
|
||||
end,
|
||||
worker_loop( Worker, N - 1, Checkpoint ).
|
||||
|
|
@ -0,0 +1,28 @@
|
|||
global nWorkers, workers, cv
|
||||
|
||||
procedure main(A)
|
||||
nWorkers := integer(A[1]) | 3
|
||||
cv := condvar()
|
||||
every put(workers := [], worker(!nWorkers))
|
||||
every wait(!workers)
|
||||
end
|
||||
|
||||
procedure worker(n)
|
||||
return thread every !3 do { # Union limits each worker to 3 pieces
|
||||
write(n," is working")
|
||||
delay(?3 * 1000)
|
||||
write(n," is done")
|
||||
countdown()
|
||||
}
|
||||
end
|
||||
|
||||
procedure countdown()
|
||||
critical cv: {
|
||||
if (nWorkers -:= 1) <= 0 then {
|
||||
write("\t\tAll done")
|
||||
nWorkers := *workers
|
||||
return (unlock(cv),signal(cv, 0))
|
||||
}
|
||||
wait(cv)
|
||||
}
|
||||
end
|
||||
|
|
@ -0,0 +1,65 @@
|
|||
:- object(checkpoint).
|
||||
|
||||
:- threaded.
|
||||
|
||||
:- public(run/3).
|
||||
:- mode(run(+integer,+integer,+float), one).
|
||||
:- info(run/3, [
|
||||
comment is 'Assemble items using a team of workers with a maximum time per item assembly.',
|
||||
arguments is ['Workers'-'Number of workers', 'Items'-'Number of items to assemble', 'Time'-'Maximum time in seconds to assemble one item']
|
||||
]).
|
||||
|
||||
:- public(run/0).
|
||||
:- mode(run, one).
|
||||
:- info(run/0, [
|
||||
comment is 'Assemble three items using a team of five workers with a maximum of 0.1 seconds per item assembly.'
|
||||
]).
|
||||
|
||||
:- uses(integer, [between/3]).
|
||||
:- uses(random, [random/3]).
|
||||
|
||||
run(Workers, Items, Time) :-
|
||||
% start the workers
|
||||
forall(
|
||||
between(1, Workers, Worker),
|
||||
threaded_ignore(worker(Worker, Items, Time))
|
||||
),
|
||||
% assemble the items
|
||||
checkpoint_loop(Workers, Items).
|
||||
|
||||
run :-
|
||||
% default values
|
||||
run(5, 3, 0.100).
|
||||
|
||||
checkpoint_loop(_, 0) :-
|
||||
!,
|
||||
write('All assemblies done.'), nl.
|
||||
checkpoint_loop(Workers, Item) :-
|
||||
% wait for all threads to reach the checkpoint
|
||||
forall(
|
||||
between(1, Workers, Worker),
|
||||
threaded_wait(done(Worker, Item))
|
||||
),
|
||||
write('Assembly of item '), write(Item), write(' done.'), nl,
|
||||
% signal the workers to procede to the next assembly
|
||||
NextItem is Item - 1,
|
||||
forall(
|
||||
between(1, Workers, Worker),
|
||||
threaded_notify(next(Worker, NextItem))
|
||||
),
|
||||
checkpoint_loop(Workers, NextItem).
|
||||
|
||||
worker(_, 0, _) :-
|
||||
!.
|
||||
worker(Worker, Item, Time) :-
|
||||
% the time necessary to assemble one item varies between 0.0 and Time seconds
|
||||
random(0.0, Time, AssemblyTime), thread_sleep(AssemblyTime),
|
||||
write('Worker '), write(Worker), write(' item '), write(Item), nl,
|
||||
% notify checkpoint that the worker have done his/her part of this item
|
||||
threaded_notify(done(Worker, Item)),
|
||||
% wait for green light to move to the next item
|
||||
NextItem is Item - 1,
|
||||
threaded_wait(next(Worker, NextItem)),
|
||||
worker(Worker, NextItem, Time).
|
||||
|
||||
:- end_object.
|
||||
|
|
@ -0,0 +1,21 @@
|
|||
| ?- checkpoint::run.
|
||||
Worker 1 item 3
|
||||
Worker 3 item 3
|
||||
Worker 5 item 3
|
||||
Worker 2 item 3
|
||||
Worker 4 item 3
|
||||
Assembly of item 3 done.
|
||||
Worker 4 item 2
|
||||
Worker 1 item 2
|
||||
Worker 5 item 2
|
||||
Worker 3 item 2
|
||||
Worker 2 item 2
|
||||
Assembly of item 2 done.
|
||||
Worker 4 item 1
|
||||
Worker 1 item 1
|
||||
Worker 2 item 1
|
||||
Worker 3 item 1
|
||||
Worker 5 item 1
|
||||
Assembly of item 1 done.
|
||||
All assemblies done.
|
||||
yes
|
||||
Loading…
Add table
Add a link
Reference in a new issue