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}