nbio/workers

Code

main.odin ¶
79 linesSource

1/*
2Example shows how you can execute CPU intensive (or blocking IO),
3without blocking the event loop thread.
4
5It creates a thread pool, starts a TCP server, which creates a task on the thread pool
6for each connection (the task currently just sleeps, but could do anything).
7
8Once the work on the worker thread is done, it uses `nbio.next_tick`,
9passing it the main event loop. This will queue that callback onto the main event loop.
10
11When the main thread runs it's next tick, it will see that queued work and execute it's callback.
12That callback finally sends the result to the client, and once it is sent, cleans up.
13
14You can run this with the netcat utility as the client like so: `nc 127.0.0.1 1234`.
15*/
16package main
17
18import "core:nbio"
19import "core:thread"
20import "core:time"
21
22Connection :: struct {
23	loop:   ^nbio.Event_Loop,
24	socket: nbio.TCP_Socket,
25}
26
27main :: proc() {
28	// accept, put work on worker, worker adds back on main thread
29
30	// Set up a thread pool with 2 workers.
31	workers: thread.Pool
32	thread.pool_init(&workers, context.allocator, 2)
33	thread.pool_start(&workers)
34
35	ep, ok := nbio.parse_endpoint("127.0.0.1:1234")
36	assert(ok)
37
38	err := nbio.acquire_thread_event_loop()
39	defer nbio.release_thread_event_loop()
40	assert(err == nil)
41
42	server, listen_err := nbio.listen_tcp(ep)
43	assert(listen_err == nil)
44	nbio.accept_poly(server, &workers, on_accept)
45
46	err = nbio.run()
47	assert(err == nil)
48
49	on_accept :: proc(op: ^nbio.Operation, workers: ^thread.Pool) {
50		assert(op.accept.err == nil)
51
52		// Accept next connection.
53		nbio.accept_poly(op.accept.socket, workers, on_accept)
54
55		// Add work to worker.
56		thread.pool_add_task(workers, context.allocator, do_work, new_clone(Connection{
57			loop   = op.l,
58			socket = op.accept.client,
59		}))
60	}
61
62	do_work :: proc(t: thread.Task) {
63		connection := (^Connection)(t.data)
64
65		// Imagine CPU intensive work.
66		time.sleep(time.Second * 5)
67
68		// Work has been done, we can now tell the client about it.
69		// NOTE: that we pass the event loop of the main IO thread here so it queues it on that.
70		nbio.send_poly(connection.socket, {transmute([]byte)string("Hellope!\n")}, connection, on_sent, l=connection.loop)
71	}
72
73	on_sent :: proc(op: ^nbio.Operation, connection: ^Connection) {
74		assert(op.send.err == nil)
75		// Client got our message, clean up.
76		nbio.close(connection.socket)
77		free(connection)
78	}
79}

Declarations Used 18