summary refs log tree commit diff stats
path: root/tests/stdlib/tchannels_pthread.nim
blob: 3bc0005511f6ca5fd84b22ec09f6d41931ce374a (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
discard """
  targets: "c cpp"
  matrix: "--gc:orc --threads:on; --gc:orc --threads:on -d:blockingTest"
  disabled: "windows"
  disabled: "bsd"
  disabled: "osx"
"""

include std/channels

import std/unittest


type
  ChannelBufKind = enum
    Unbuffered # Unbuffered (blocking) channel
    Buffered   # Buffered (non-blocking channel)


proc capacity(chan: ChannelRaw): int {.inline.} = chan.size
func isBuffered(chan: ChannelRaw): bool =
  chan.size - 1 > 0

when defined(blockingTest):
  const nonBlocking = false
else:
  const nonBlocking = true

type
  Pthread {.importc: "pthread_t", header: "<sys/types.h>".} = distinct culong
  PthreadAttr* {.byref, importc: "pthread_attr_t", header: "<sys/types.h>".} = object
  Errno* = distinct cint

proc pthread_create[T](
      thread: var Pthread,
      attr: ptr PthreadAttr, # In Nim this is a var and how Nim sets a custom stack
      fn: proc (x: ptr T): pointer {.thread, noconv.},
      arg: ptr T
  ): Errno {.header: "<sys/types.h>".}

proc pthread_join(
      thread: Pthread,
      thread_exit_status: ptr pointer
    ): Errno {.header: "<pthread.h>".}

template channel_send_loop(chan: ChannelRaw,
                      data: sink pointer,
                      size: int,
                      body: untyped): untyped =
  while not sendMpmc(chan, data, size, nonBlocking):
    body

template channel_receive_loop(chan: ChannelRaw,
                        data: pointer,
                        size: int,
                        body: untyped): untyped =
  while not recvMpmc(chan, data, size, nonBlocking):
    body


# Without threads:on or release,
# worker threads will crash on popFrame

import std/unittest

type ThreadArgs = object
  ID: int
  chan: ChannelRaw

template Worker(id: int, body: untyped): untyped {.dirty.} =
  if args.ID == id:
    body


const Sender = 1
const Receiver = 0

proc runSuite(
            name: string,
            fn: proc(args: ptr ThreadArgs): pointer {.noconv, gcsafe.}
          ) =
  var chan: ChannelRaw

  for i in Unbuffered .. Buffered:
    if i == Unbuffered:
      chan = allocChannel(size = 32, n = 1)
      check:
        peek(chan) == 0
        capacity(chan) == 1
        isBuffered(chan) == false
        isUnbuffered(chan) == true

    else:
      chan = allocChannel(size = int.sizeof.int, n = 7)
      check:
        peek(chan) == 0
        capacity(chan) == 7
        isBuffered(chan) == true
        isUnbuffered(chan) == false

      var threads: array[2, Pthread]
      var args = [
        ThreadArgs(ID: 0, chan: chan),
        ThreadArgs(ID: 1, chan: chan)
      ]

      discard pthread_create(threads[0], nil, fn, args[0].addr)
      discard pthread_create(threads[1], nil, fn, args[1].addr)

      discard pthread_join(threads[0], nil)
      discard pthread_join(threads[1], nil)

    freeChannel(chan)

# ----------------------------------------------------------------------------------

proc thread_func(args: ptr ThreadArgs): pointer {.noconv.} =

  # Worker RECEIVER:
  # ---------
  # <- chan
  # <- chan
  # <- chan
  #
  # Worker SENDER:
  # ---------
  # chan <- 42
  # chan <- 53
  # chan <- 64
  #

  Worker(Receiver):
    var val: int
    for j in 0 ..< 3:
      channel_receive_loop(args.chan, val.addr, val.sizeof.int):
        # Busy loop, normally it should yield
        discard
      check: val == 42 + j*11

  Worker(Sender):
    var val: int
    check: peek(args.chan) == 0
    for j in 0 ..< 3:
      val = 42 + j*11
      channel_send_loop(args.chan, val.addr, val.sizeof.int):
        # Busy loop, normally it should yield
        discard

  return nil

runSuite("[ChannelRaw] 2 threads can send data", thread_func)

# ----------------------------------------------------------------------------------

iterator pairs(chan: ChannelRaw, T: typedesc): (int, T) =
  var i = 0
  var x: T
  while not isClosed(chan) or peek(chan) > 0:
    let r = recvMpmc(chan, x.addr, x.sizeof.int, true)
    # printf("x: %d, r: %d\n", x, r)
    if r:
      yield (i, x)
      inc i

proc thread_func_2(args: ptr ThreadArgs): pointer {.noconv.} =
  # Worker RECEIVER:
  # ---------
  # <- chan until closed and empty
  #
  # Worker SENDER:
  # ---------
  # chan <- 42, 53, 64, ...

  const N = 100

  Worker(Receiver):
    for j, val in pairs(args.chan, int):
      # TODO: Need special handling that doesn't allocate
      #       in thread with no GC
      #       when check fails
      #
      check: val == 42 + j*11

  Worker(Sender):
    var val: int
    check: peek(args.chan) == 0
    for j in 0 ..< N:
      val = 42 + j*11
      channel_send_loop(args.chan, val.addr, int.sizeof.int):
        discard
    discard channelCloseMpmc(args.chan)

  return nil

runSuite("[ChannelRaw] channel_close, freeChannel, channelCache", thread_func_2)

# ----------------------------------------------------------------------------------

proc isCached(chan: ChannelRaw): bool =
  assert not chan.isNil

  var p = channelCache
  while not p.isNil:
    if chan.itemsize == p.chanSize and
        chan.size == p.chanN:
      for i in 0 ..< p.numCached:
        if chan == p.cache[i]:
          return true
      # No more channel in cache can match
      return false
    p = p.next
  return false

block: # [ChannelRaw] ChannelRaw caching implementation

  # Start from clean cache slate
  freeChannelCache()

  block: # Explicit caches allocation
    check:
      allocChannelCache(int sizeof(char), 4)
      allocChannelCache(int sizeof(int), 8)
      allocChannelCache(int sizeof(ptr float64), 16)

      # Don't create existing channel cache
      not allocChannelCache(int sizeof(char), 4)
      not allocChannelCache(int sizeof(int), 8)
      not allocChannelCache(int sizeof(ptr float64), 16)

    check:
      channelCacheLen == 3

  # ---------------------------------
  var chan, stash: array[10, ChannelRaw]

  block: # Implicit caches allocation

    chan[0] = allocChannel(sizeof(char), 4)
    chan[1] = allocChannel(sizeof(int32), 8)
    chan[2] = allocChannel(sizeof(ptr float64), 16)

    chan[3] = allocChannel(sizeof(char), 5)
    chan[4] = allocChannel(sizeof(int64), 8)
    chan[5] = allocChannel(sizeof(ptr float32), 24)

    # We have caches ready to store specific channel kinds
    check: channelCacheLen == 6 # Cumulated with previous test
    # But they are not in cache while in use
    check:
      not chan[0].isCached
      not chan[1].isCached
      not chan[2].isCached
      not chan[3].isCached
      not chan[4].isCached
      not chan[5].isCached

  block: # Freed channels are returned to cache
    stash[0..5] = chan.toOpenArray(0, 5)
    for i in 0 .. 5:
      # Free the channels
      freeChannel(chan[i])

    check:
      stash[0].isCached
      stash[1].isCached
      stash[2].isCached
      stash[3].isCached
      stash[4].isCached
      stash[5].isCached

  block: # Cached channels are being reused

    chan[6] = allocChannel(sizeof(char), 4)
    chan[7] = allocChannel(sizeof(int32), 8)
    chan[8] = allocChannel(sizeof(ptr float32), 16)
    chan[9] = allocChannel(sizeof(ptr float64), 16)

    # All (itemsize, queue size, implementation) were already allocated
    check: channelCacheLen == 6

    # We reused old channels from cache
    check:
      chan[6] == stash[0]
      chan[7] == stash[1]
      chan[8] == stash[2]
      # chan[9] - required a fresh alloc

  block: # Clearing the cache

    stash[6..9] = chan.toOpenArray(6, 9)

    for i in 6 .. 9:
      freeChannel(chan[i])

    check:
      stash[6].isCached
      stash[7].isCached
      stash[8].isCached
      stash[9].isCached

    freeChannelCache()

    # Check that nothing is cached anymore
    for i in 0 .. 9:
      check: not stash[i].isCached
    # And length is reset to 0
    check: channelCacheLen == 0

    # Cache can grow again
    chan[0] = allocChannel(sizeof((int, float, int32, uint)), 1)
    chan[1] = allocChannel(sizeof(int32), 0)
    chan[2] = allocChannel(sizeof(int32), 0)

    check: channelCacheLen == 2

    # Interleave cache clear and channel free
    freeChannelCache()
    check: channelCacheLen == 0

    freeChannel(chan[0])
    freeChannel(chan[1])
    freeChannel(chan[2])