thread/sync_mutex

Code

thread_sync_mutex.odin ¶
168 linesSource

1//
2// Demonstrates how to spawn multiple threads and safely access data 
3// from each by using a mutex.
4//
5package thread_sync_example
6
7import "core:fmt"
8import "core:math/rand"
9import "core:sync"
10import "core:thread"
11import "core:time"
12
13// Defines an arbitrary work item. To simulate the CPU work on
14// these items, each will have a "processing time" that each thread will
15// wait for before continuing onto the next item.
16Work_Item :: struct {
17	item_tag:        i32,
18	processing_time: f32,
19}
20
21create_randomized_queue :: proc(num_items: int) -> (q: [dynamic]Work_Item) {
22	// This initializes the queue with a length of zero, and a capacity of `num_items`.
23	// Pre-allocating space when you know how much you need is good!
24	q = make([dynamic]Work_Item, 0, num_items)
25
26	// Initialize the items in the queue. Each item will have a unique tag,
27	// and a random "processing time".
28	for i in 0 ..< num_items {
29		item: Work_Item
30		item.item_tag = i32(i) + 1
31		// This sets the item's processing time to a value between 0.1 and 0.51 (exclusive).
32		item.processing_time = rand.float32_range(0.1, 0.51)
33		append(&q, item)
34	}
35
36	return
37}
38
39// This is the procedure that we'll be running in the threads that we spawn later.
40process_item :: proc(queue: ^[dynamic]Work_Item, mutex: ^sync.Mutex, thread_identifier: int) {
41	// This proc is essentially an infinite loop that breaks once it no longer has any data to process.
42
43	for {
44		// First we need to get a lock on our mutex. 
45		// That way we know whether we can safely access our queue, or whether 
46		// another thread is using it already.
47		sync.mutex_lock(mutex)
48
49		// This is the critical point where the mutex being locked matters.
50		// Here we attempt to pop the first element off of our queue.
51		item, pop_ok := pop_front_safe(queue)
52
53		// Now that we've got the data we need from the queue, we can unlock our mutex 
54		// to let other threads access the queue to perform their work.
55		sync.mutex_unlock(mutex)
56
57		// If we tried to pop something off but the queue was empty, we have nothing left to
58		// process, so we'll just break out of our ininite loop.
59		// Once the loop ends, our function will return, and the thread will stop.
60		if !pop_ok {
61			break
62		}
63
64		// Now we can do our item processing! Which in this case is just "processing" it for 
65		// the item's `processing_time` in seconds.
66		//
67		// Since `processing_time` is a f32, you need to cast `time.Second` to a f32, 
68		// then back to `time.Duration` to get your fraction of a second.
69		time.sleep(time.Duration(f32(time.Second) * item.processing_time))
70
71		// After we've done our "processing" (sleeping on the job, really), we can print
72		// some info to the console about our item, and the thread that grabbed it.
73		//
74		// `fmt.printfln` and the other `fmt` procs that print to stdout are thread-safe, 
75		// so nothing to worry about here.
76		fmt.printfln(
77			"[THREAD %02d] Item %04d processed in %0.2f seconds.",
78			thread_identifier,
79			item.item_tag,
80			item.processing_time,
81		)
82	}
83}
84
85main :: proc() {
86	// This `RANDOM_SEED` is just a compile-time constant that will
87	// seed the default random generator if specified as a non-zero value.
88	// I added this in to allow for predictable, reproducible outputs.
89	//
90	// When it's 0, Odin's random generator is seeded as it normally would be by default.
91	// Otherwise, this `when` clause kicks in at compile-time 
92	// and will override the default seeding mechanism.
93	//
94	// To specify it yourself, you can just add `--define:RANDOM_SEED=...`
95	// to your `odin build/run` command.
96	RANDOM_SEED: u64 : #config(RANDOM_SEED, 0)
97	when RANDOM_SEED > 0 {
98		state := rand.create(RANDOM_SEED)
99		context.random_generator = rand.default_random_generator(&state)
100	}
101
102	// Initialize a randomized set of data to work off of.
103	// It'll be a dynamic array of `Work_Items`, which essentially just have an ID number and a duration.
104	queue := create_randomized_queue(500)
105
106	// This is a Mutex. (Short for "mutual exclusion lock")
107	// It doesn't actually hold any data, but rather it's used in multi-threaded 
108	// applications as a way to tell other threads when it's safe to access data.
109	//
110	// A Mutex starts in an UNLOCKED state. At any time, you can LOCK a Mutex using `sync.lock`.
111	// If a Mutex is LOCKED, that means when something else tries to LOCK it, it will halt the 
112	// execution of that thread since another thread has already LOCKED it.
113	//
114	// However, once the Mutex is UNLOCKED, any thread can LOCK it for themselves.
115	//
116	// Mutexes can be used to guarantee safe access to data across multiple threads. Once a thread locks it, 
117	// any other threads that also try to lock it will be forced to wait. 
118	// This prevents two threads from reading/writing the same data, which can result in data races.
119	mutex: sync.Mutex
120
121
122	// This constant int is going to define how many threads we actually want to run.
123	MAX_THREADS: int : 8
124	// And here we define an array that's going to hold all of the Threads that we spawn.
125	threads: [MAX_THREADS]^thread.Thread
126
127	// Let's start making some threads.
128	for i in 0 ..< len(threads) {
129		// Let's get which thread number this is so we can pass it to our threaded proc.
130		t_id := i + 1
131		// This is where the magic happens. We're going to create up to our MAX number of Threads and store them in our 
132		// `threads` array.
133		// Since our thread proc takes three arguments, we need a way to pass these in.
134		// Luckily, `create_and_start_with_poly_data` exists! It allows you to pass in function arguments that get 
135		// consumed by the thread proc easily.
136		//
137		// Now to explain exactly what these arguments are:
138		//      &queue  - A pointer to our queue object. We need to pass it by pointer to pop items off of it!
139		//      &mutex  - A pointer to our mutex. This is what our threads will use to signal to each other that they 
140		//				need exclusive access to the queue at the critical point where they access it.
141		//      t_id    - Just the index of our thread + 1, for printing purposes so we can identify who's working.
142		//
143		//      process_item - This is our procedure! The thread is going to make this thing run with all of the 
144		//		previous arguments passed into it.
145		//
146		// With all that out of the way, let's create our thread and store it in `threads` at index `i` for later!
147		threads[i] = thread.create_and_start_with_poly_data3(&queue, &mutex, t_id, process_item)
148	}
149
150	// Now we're going to use `join_multiple` to wait for all of our threads to stop processing.
151	// This is why we're holding onto those threads in our array. You wouldn't want to just let them spin off 
152	// and never check on them again!
153	//
154	// `join_multiple` takes a variable number of Thread pointers (`^Thread`), and BLOCKS the main thread 
155	// until each one of them is finished processing.
156	//
157	// Since we have an array of Thread pointers, we can use the `..` operator to expand all of the array 
158	// items as arguments to `join_multiple`!
159	thread.join_multiple(..threads[:])
160
161	// Once the program ends, we'll clean up after ourselves by destroying each of these threads we created.
162	for t in threads {
163		thread.destroy(t)
164	}
165
166	// Everything's all finished now. Let's print out a "done" message and call it a day!
167	fmt.printfln("Processed all items! Exiting.")
168}

Declarations Used 14