first steps to thread local heaps

This commit is contained in:
Araq 2011-06-02 13:02:40 +02:00
commit 3260702a60
21 changed files with 1072 additions and 633 deletions

View file

@ -1105,8 +1105,8 @@ proc rdSetElemLoc(a: TLoc, setType: PType): PRope =
# before the set operation # before the set operation
result = rdCharLoc(a) result = rdCharLoc(a)
assert(setType.kind == tySet) assert(setType.kind == tySet)
if (firstOrd(setType) != 0): if firstOrd(setType) != 0:
result = ropef("($1-$2)", [result, toRope(firstOrd(setType))]) result = ropef("($1- $2)", [result, toRope(firstOrd(setType))])
proc fewCmps(s: PNode): bool = proc fewCmps(s: PNode): bool =
# this function estimates whether it is better to emit code # this function estimates whether it is better to emit code

View file

@ -431,9 +431,18 @@ proc assignLocalVar(p: BProc, s: PSym) =
proc declareThreadVar(m: BModule, s: PSym) = proc declareThreadVar(m: BModule, s: PSym) =
if optThreads in gGlobalOptions: if optThreads in gGlobalOptions:
app(m.s[cfsVars], "NIM_THREADVAR ") if platform.OS[targetOS].props.contains(ospLacksThreadVars):
app(m.s[cfsVars], getTypeDesc(m, s.loc.t)) # we gather all thread locals var into a struct and put that into
# nim__dat.c; we need to allocate storage for that somehow, can't use
# the thread local storage allocator for it :-(
# XXX we need to adapt expr() too, every reference to a thread local var
# generates quite some code ...
InternalError("no workaround for lack of thread local vars implemented")
else:
app(m.s[cfsVars], "NIM_THREADVAR ")
app(m.s[cfsVars], getTypeDesc(m, s.loc.t))
else:
app(m.s[cfsVars], getTypeDesc(m, s.loc.t))
proc assignGlobalVar(p: BProc, s: PSym) = proc assignGlobalVar(p: BProc, s: PSym) =
if s.loc.k == locNone: if s.loc.k == locNone:
@ -640,8 +649,8 @@ proc genProcAux(m: BModule, prc: PSym) =
assignParam(p, param) assignParam(p, param)
genStmts(p, prc.ast.sons[codePos]) # modifies p.locals, p.init, etc. genStmts(p, prc.ast.sons[codePos]) # modifies p.locals, p.init, etc.
if sfPure in prc.flags: if sfPure in prc.flags:
generatedProc = ropeff("$1 {$n$2$3$4}$n", "define $1 {$n$2$3$4}$n", [header, generatedProc = ropeff("$1 {$n$2$3$4}$n", "define $1 {$n$2$3$4}$n",
p.s[cpsLocals], p.s[cpsInit], p.s[cpsStmts]]) [header, p.s[cpsLocals], p.s[cpsInit], p.s[cpsStmts]])
else: else:
generatedProc = ropeff("$1 {$n", "define $1 {$n", [header]) generatedProc = ropeff("$1 {$n", "define $1 {$n", [header])
if optStackTrace in prc.options: if optStackTrace in prc.options:
@ -756,7 +765,8 @@ proc getFileHeader(cfilenoext: string): PRope =
"; (c) 2011 Andreas Rumpf$n", [toRope(versionAsString)]) "; (c) 2011 Andreas Rumpf$n", [toRope(versionAsString)])
else: else:
result = ropeff("/* Generated by Nimrod Compiler v$1 */$n" & result = ropeff("/* Generated by Nimrod Compiler v$1 */$n" &
"/* (c) 2011 Andreas Rumpf */$n" & "/* Compiled for: $2, $3, $4 */$n" & "/* (c) 2011 Andreas Rumpf */$n" &
"/* Compiled for: $2, $3, $4 */$n" &
"/* Command for C compiler:$n $5 */$n", "/* Command for C compiler:$n $5 */$n",
"; Generated by Nimrod Compiler v$1$n" & "; Generated by Nimrod Compiler v$1$n" &
"; (c) 2011 Andreas Rumpf$n" & "; Compiled for: $2, $3, $4$n" & "; (c) 2011 Andreas Rumpf$n" & "; Compiled for: $2, $3, $4$n" &

View file

@ -27,7 +27,8 @@ type
TInfoOSProp* = enum TInfoOSProp* = enum
ospNeedsPIC, # OS needs PIC for libraries ospNeedsPIC, # OS needs PIC for libraries
ospCaseInsensitive, # OS filesystem is case insensitive ospCaseInsensitive, # OS filesystem is case insensitive
ospPosix # OS is posix-like ospPosix, # OS is posix-like
ospLacksThreadVars # OS lacks proper __threadvar support
TInfoOSProps* = set[TInfoOSProp] TInfoOSProps* = set[TInfoOSProp]
TInfoOS* = tuple[name: string, parDir: string, dllFrmt: string, TInfoOS* = tuple[name: string, parDir: string, dllFrmt: string,
altDirSep: string, objExt: string, newLine: string, altDirSep: string, objExt: string, newLine: string,
@ -129,7 +130,7 @@ const
(name: "MacOSX", parDir: "..", dllFrmt: "lib$1.dylib", altDirSep: ":", (name: "MacOSX", parDir: "..", dllFrmt: "lib$1.dylib", altDirSep: ":",
objExt: ".o", newLine: "\x0A", pathSep: ":", dirSep: "/", objExt: ".o", newLine: "\x0A", pathSep: ":", dirSep: "/",
scriptExt: ".sh", curDir: ".", exeExt: "", extSep: ".", scriptExt: ".sh", curDir: ".", exeExt: "", extSep: ".",
props: {ospNeedsPIC, ospPosix}), props: {ospNeedsPIC, ospPosix, ospLacksThreadVars}),
(name: "EcmaScript", parDir: "..", (name: "EcmaScript", parDir: "..",
dllFrmt: "lib$1.so", altDirSep: "/", dllFrmt: "lib$1.so", altDirSep: "/",
objExt: ".o", newLine: "\x0A", objExt: ".o", newLine: "\x0A",

View file

@ -32,7 +32,8 @@ Core
magic to work. magic to work.
* `threads <threads.html>`_ * `threads <threads.html>`_
Nimrod thread support. Nimrod thread support. **Note**: This is part of the system module. Do not
import it directly.
* `macros <macros.html>`_ * `macros <macros.html>`_
Contains the AST API and documentation of Nimrod for writing macros. Contains the AST API and documentation of Nimrod for writing macros.
@ -228,6 +229,14 @@ Multimedia support
the ``graphics`` module. the ``graphics`` module.
Database support
----------------
* `redis <redis.html>`_
This module implements a redis client. It allows you to connect to a
redis-server instance, send commands and receive replies.
Impure libraries Impure libraries
================ ================
@ -284,6 +293,10 @@ Other
This module contains code for reading from `stdin`:idx:. On UNIX the GNU This module contains code for reading from `stdin`:idx:. On UNIX the GNU
readline library is wrapped and set up. readline library is wrapped and set up.
* `zmq <zmq.html>`_
Nimrod 0mq wrapper. This file contains the low level C wrappers as well as
some higher level constructs.
Wrappers Wrappers
======== ========

14
examples/0mq/client.nim Normal file
View file

@ -0,0 +1,14 @@
import zmq
var connection = zmq.open("tcp://localhost:5555", server=false)
echo("Connecting...")
for i in 0..10:
echo("Sending hello...", i)
send(connection, "Hello")
var reply = receive(connection)
echo("Received ...", reply)
close(connection)

11
examples/0mq/server.nim Normal file
View file

@ -0,0 +1,11 @@
import zmq
var connection = zmq.open("tcp://*:5555", server=true)
while True:
var request = receive(connection)
echo("Received: ", request)
send(connection, "World")
close(connection)

View file

@ -1,348 +0,0 @@
#
#
# Nimrod's Runtime Library
# (c) Copyright 2011 Andreas Rumpf
#
# See the file "copying.txt", included in this
# distribution, for details about the copyright.
#
## Basic thread support for Nimrod. Note that Nimrod's default GC is still
## single-threaded. This means that you MUST turn off the GC while multiple
## threads are executing that allocate GC'ed memory. The alternative is to
## compile with ``--gc:none`` or ``--gc:boehm``.
##
## Example:
##
## .. code-block:: nimrod
##
## var
## thr: array [0..4, TThread[tuple[a,b: int]]]
## L: TLock
##
## proc threadFunc(interval: tuple[a,b: int]) {.procvar.} =
## for i in interval.a..interval.b:
## Aquire(L) # lock stdout
## echo i
## Release(L)
##
## InitLock(L)
##
## GC_disable() # native GC does not support multiple threads yet :-(
## for i in 0..high(thr):
## createThread(thr[i], threadFunc, (i*10, i*10+5))
## for i in 0..high(thr):
## joinThread(thr[i])
## GC_enable()
when not compileOption("threads"):
{.error: "Thread support requires ``--threads:on`` commandline switch".}
when not defined(boehmgc) and not defined(nogc) and false:
{.error: "Thread support requires --gc:boehm or --gc:none".}
include "lib/system/systhread"
# We jump through some hops here to ensure that Nimrod thread procs can have
# the Nimrod calling convention. This is needed because thread procs are
# ``stdcall`` on Windows and ``noconv`` on UNIX. Alternative would be to just
# use ``stdcall`` since it is mapped to ``noconv`` on UNIX anyway. However,
# the current approach will likely result in less problems later when we have
# GC'ed closures in Nimrod.
type
TThreadProcClosure {.pure, final.}[TParam] = object
fn: proc (p: TParam)
threadLocalStorage: pointer
stackBottom: pointer
data: TParam
when defined(windows):
type
THandle = int
TSysThread = THandle
TWinThreadProc = proc (x: pointer): int32 {.stdcall.}
proc CreateThread(lpThreadAttributes: Pointer, dwStackSize: int32,
lpStartAddress: TWinThreadProc,
lpParameter: Pointer,
dwCreationFlags: int32, lpThreadId: var int32): THandle {.
stdcall, dynlib: "kernel32", importc: "CreateThread".}
when false:
proc winSuspendThread(hThread: TSysThread): int32 {.
stdcall, dynlib: "kernel32", importc: "SuspendThread".}
proc winResumeThread(hThread: TSysThread): int32 {.
stdcall, dynlib: "kernel32", importc: "ResumeThread".}
proc WaitForMultipleObjects(nCount: int32,
lpHandles: ptr array[0..10, THandle],
bWaitAll: int32,
dwMilliseconds: int32): int32 {.
stdcall, dynlib: "kernel32", importc: "WaitForMultipleObjects".}
proc WaitForSingleObject(hHandle: THANDLE, dwMilliseconds: int32): int32 {.
stdcall, dynlib: "kernel32", importc: "WaitForSingleObject".}
proc TerminateThread(hThread: THandle, dwExitCode: int32): int32 {.
stdcall, dynlib: "kernel32", importc: "TerminateThread".}
{.push stack_trace:off.}
proc threadProcWrapper[TParam](closure: pointer): int32 {.stdcall.} =
var c = cast[ptr TThreadProcClosure[TParam]](closure)
SetThreadLocalStorage(c.threadLocalStorage)
c.fn(c.data)
# implicitely return 0
{.pop.}
else:
type
TSysThread {.importc: "pthread_t", header: "<sys/types.h>".} = int
Ttimespec {.importc: "struct timespec",
header: "<time.h>", final, pure.} = object
tv_sec: int
tv_nsec: int
proc pthread_create(a1: var TSysThread, a2: ptr int,
a3: proc (x: pointer) {.noconv.},
a4: pointer): cint {.importc: "pthread_create",
header: "<pthread.h>".}
proc pthread_join(a1: TSysThread, a2: ptr pointer): cint {.
importc, header: "<pthread.h>".}
proc pthread_cancel(a1: TSysThread): cint {.
importc: "pthread_cancel", header: "<pthread.h>".}
proc AquireSysTimeoutAux(L: var TSysLock, timeout: var Ttimespec): cint {.
importc: "pthread_mutex_timedlock", header: "<time.h>".}
proc AquireSysTimeout(L: var TSysLock, msTimeout: int) {.inline.} =
var a: Ttimespec
a.tv_sec = msTimeout div 1000
a.tv_nsec = (msTimeout mod 1000) * 1000
var res = AquireSysTimeoutAux(L, a)
if res != 0'i32:
raise newException(EResourceExhausted, $strerror(res))
{.push stack_trace:off.}
proc threadProcWrapper[TParam](closure: pointer) {.noconv.} =
var c = cast[ptr TThreadProcClosure[TParam]](closure)
SetThreadLocalStorage(c.threadLocalStorage)
c.fn(c.data)
{.pop.}
const
noDeadlocks = true # compileOption("deadlockPrevention")
type
TLock* = TSysLock
TThread* {.pure, final.}[TParam] = object ## Nimrod thread.
sys: TSysThread
c: TThreadProcClosure[TParam]
when nodeadlocks:
var
deadlocksPrevented* = 0 ## counts the number of times a
## deadlock has been prevented
proc InitLock*(lock: var TLock) {.inline.} =
## Initializes the lock `lock`.
InitSysLock(lock)
proc OrderedLocks(g: PGlobals): bool =
for i in 0 .. g.locksLen-2:
if g.locks[i] >= g.locks[i+1]: return false
result = true
proc TryAquire*(lock: var TLock): bool {.inline.} =
## Try to aquires the lock `lock`. Returns `true` on success.
when noDeadlocks:
result = TryAquireSys(lock)
if not result: return
# we have to add it to the ordered list. Oh, and we might fail if there#
# there is no space in the array left ...
var g = GetGlobals()
if g.locksLen >= len(g.locks):
ReleaseSys(lock)
raise newException(EResourceExhausted, "cannot aquire additional lock")
# find the position to add:
var p = addr(lock)
var L = g.locksLen-1
var i = 0
while i <= L:
assert g.locks[i] != nil
if g.locks[i] < p: inc(i) # in correct order
elif g.locks[i] == p: return # thread already holds lock
else:
# do the crazy stuff here:
while L >= i:
g.locks[L+1] = g.locks[L]
dec L
g.locks[i] = p
inc(g.locksLen)
assert OrderedLocks(g)
return
# simply add to the end:
g.locks[g.locksLen] = p
inc(g.locksLen)
assert OrderedLocks(g)
else:
result = TryAquireSys(lock)
proc Aquire*(lock: var TLock) =
## Aquires the lock `lock`.
when nodeadlocks:
var g = GetGlobals()
var p = addr(lock)
var L = g.locksLen-1
var i = 0
while i <= L:
assert g.locks[i] != nil
if g.locks[i] < p: inc(i) # in correct order
elif g.locks[i] == p: return # thread already holds lock
else:
# do the crazy stuff here:
if g.locksLen >= len(g.locks):
raise newException(EResourceExhausted, "cannot aquire additional lock")
while L >= i:
ReleaseSys(cast[ptr TSysLock](g.locks[L])[])
g.locks[L+1] = g.locks[L]
dec L
# aquire the current lock:
AquireSys(lock)
g.locks[i] = p
inc(g.locksLen)
# aquire old locks in proper order again:
L = g.locksLen-1
inc i
while i <= L:
AquireSys(cast[ptr TSysLock](g.locks[i])[])
inc(i)
# DANGER: We can only modify this global var if we gained every lock!
# NO! We need an atomic increment. Crap.
discard system.atomicInc(deadlocksPrevented, 1)
assert OrderedLocks(g)
return
# simply add to the end:
if g.locksLen >= len(g.locks):
raise newException(EResourceExhausted, "cannot aquire additional lock")
AquireSys(lock)
g.locks[g.locksLen] = p
inc(g.locksLen)
assert OrderedLocks(g)
else:
AquireSys(lock)
proc Release*(lock: var TLock) =
## Releases the lock `lock`.
when nodeadlocks:
var g = GetGlobals()
var p = addr(lock)
var L = g.locksLen
for i in countdown(L-1, 0):
if g.locks[i] == p:
for j in i..L-2: g.locks[j] = g.locks[j+1]
dec g.locksLen
break
ReleaseSys(lock)
proc joinThread*[TParam](t: TThread[TParam]) {.inline.} =
## waits for the thread `t` until it has terminated.
when hostOS == "windows":
discard WaitForSingleObject(t.sys, -1'i32)
else:
discard pthread_join(t.sys, nil)
proc destroyThread*[TParam](t: var TThread[TParam]) {.inline.} =
## forces the thread `t` to terminate. This is potentially dangerous if
## you don't have full control over `t` and its aquired resources.
when hostOS == "windows":
discard TerminateThread(t.sys, 1'i32)
else:
discard pthread_cancel(t.sys)
proc createThread*[TParam](t: var TThread[TParam],
tp: proc (param: TParam),
param: TParam) =
## creates a new thread `t` and starts its execution. Entry point is the
## proc `tp`. `param` is passed to `tp`.
t.c.threadLocalStorage = AllocThreadLocalStorage()
t.c.data = param
t.c.fn = tp
when hostOS == "windows":
var dummyThreadId: int32
t.sys = CreateThread(nil, 0'i32, threadProcWrapper[TParam],
addr(t.c), 0'i32, dummyThreadId)
else:
if pthread_create(t.sys, nil, threadProcWrapper[TParam], addr(t.c)) != 0:
raise newException(EIO, "cannot create thread")
when isMainModule:
import os
var
thr: array [0..5, TThread[tuple[a, b: int]]]
L, M, N: TLock
proc doNothing() = nil
proc threadFunc(interval: tuple[a, b: int]) {.procvar.} =
doNothing()
for i in interval.a..interval.b:
when nodeadlocks:
case i mod 6
of 0:
Aquire(L) # lock stdout
Aquire(M)
Aquire(N)
of 1:
Aquire(L)
Aquire(N) # lock stdout
Aquire(M)
of 2:
Aquire(M)
Aquire(L)
Aquire(N)
of 3:
Aquire(M)
Aquire(N)
Aquire(L)
of 4:
Aquire(N)
Aquire(M)
Aquire(L)
of 5:
Aquire(N)
Aquire(L)
Aquire(M)
else: assert false
else:
Aquire(L) # lock stdout
Aquire(M)
echo i
os.sleep(10)
when nodeadlocks:
echo "deadlocks prevented: ", deadlocksPrevented
when nodeadlocks:
Release(N)
Release(M)
Release(L)
InitLock(L)
InitLock(M)
InitLock(N)
proc main =
for i in 0..high(thr):
createThread(thr[i], threadFunc, (i*100, i*100+50))
for i in 0..high(thr):
joinThread(thr[i])
GC_disable()
main()
GC_enable()

23
lib/prelude.nim Normal file
View file

@ -0,0 +1,23 @@
#
#
# Nimrod's Runtime Library
# (c) Copyright 2011 Andreas Rumpf
#
# See the file "copying.txt", included in this
# distribution, for details about the copyright.
#
## This is an include file that simply imports common modules for your
## convenience:
##
## .. code-block:: nimrod
## include prelude
##
## Same as:
##
## .. code-block:: nimrod
## import os, strutils, times, parseutils, parseopt
import os, strutils, times, parseutils, parseopt

View file

@ -778,6 +778,12 @@ proc compileOption*(option, arg: string): bool {.
const const
hasThreadSupport = compileOption("threads") hasThreadSupport = compileOption("threads")
hasSharedHeap = false # don't share heaps, so every thread has its own heap
when hasThreadSupport and not hasSharedHeap:
{.pragma: rtlThreadVar, threadvar.}
else:
{.pragma: rtlThreadVar.}
include "system/inclrtl" include "system/inclrtl"
@ -1448,12 +1454,6 @@ proc quit*(errorcode: int = QuitSuccess) {.
when not defined(EcmaScript) and not defined(NimrodVM): when not defined(EcmaScript) and not defined(NimrodVM):
{.push stack_trace: off.} {.push stack_trace: off.}
proc atomicInc*(memLoc: var int, x: int): int {.inline.}
## atomic increment of `memLoc`. Returns the value after the operation.
proc atomicDec*(memLoc: var int, x: int): int {.inline.}
## atomic decrement of `memLoc`. Returns the value after the operation.
proc initGC() proc initGC()
proc initStackBottom() {.inline.} = proc initStackBottom() {.inline.} =
@ -1666,7 +1666,23 @@ when not defined(EcmaScript) and not defined(NimrodVM):
# ---------------------------------------------------------------------------- # ----------------------------------------------------------------------------
include "system/systhread" proc atomicInc*(memLoc: var int, x: int): int {.inline.}
## atomic increment of `memLoc`. Returns the value after the operation.
proc atomicDec*(memLoc: var int, x: int): int {.inline.}
## atomic decrement of `memLoc`. Returns the value after the operation.
include "system/atomics"
type
PSafePoint = ptr TSafePoint
TSafePoint {.compilerproc, final.} = object
prev: PSafePoint # points to next safe point ON THE STACK
status: int
context: C_JmpBuf
when hasThreadSupport:
include "system/threads"
include "system/excpt" include "system/excpt"
# we cannot compile this with stack tracing on # we cannot compile this with stack tracing on
# as it would recurse endlessly! # as it would recurse endlessly!

View file

@ -1,7 +1,7 @@
# #
# #
# Nimrod's Runtime Library # Nimrod's Runtime Library
# (c) Copyright 2009 Andreas Rumpf # (c) Copyright 2011 Andreas Rumpf
# #
# See the file "copying.txt", included in this # See the file "copying.txt", included in this
# distribution, for details about the copyright. # distribution, for details about the copyright.
@ -80,7 +80,7 @@ else:
# system immediately. # system immediately.
const const
ChunkOsReturn = 256 * PageSize ChunkOsReturn = 256 * PageSize # 1 MB
InitialMemoryRequest = ChunkOsReturn div 2 # < ChunkOsReturn! InitialMemoryRequest = ChunkOsReturn div 2 # < ChunkOsReturn!
SmallChunkSize = PageSize SmallChunkSize = PageSize
@ -101,6 +101,7 @@ type
next: ptr TFreeCell # next free cell in chunk (overlaid with refcount) next: ptr TFreeCell # next free cell in chunk (overlaid with refcount)
zeroField: int # 0 means cell is not used (overlaid with typ field) zeroField: int # 0 means cell is not used (overlaid with typ field)
# 1 means cell is manually managed pointer # 1 means cell is manually managed pointer
# otherwise a PNimType is stored in there
PChunk = ptr TBaseChunk PChunk = ptr TBaseChunk
PBigChunk = ptr TBigChunk PBigChunk = ptr TBigChunk
@ -151,6 +152,7 @@ type
TAllocator {.final, pure.} = object TAllocator {.final, pure.} = object
llmem: PLLChunk llmem: PLLChunk
currMem, maxMem, freeMem: int # memory sizes (allocated from OS) currMem, maxMem, freeMem: int # memory sizes (allocated from OS)
lastSize: int # needed for the case that OS gives us pages linearly
freeSmallChunks: array[0..SmallChunkSize div MemAlign-1, PSmallChunk] freeSmallChunks: array[0..SmallChunkSize div MemAlign-1, PSmallChunk]
freeChunksList: PBigChunk # XXX make this a datastructure with O(1) access freeChunksList: PBigChunk # XXX make this a datastructure with O(1) access
chunkStarts: TIntSet chunkStarts: TIntSet
@ -168,9 +170,6 @@ proc getMaxMem(a: var TAllocator): int =
# maximum of these both values here: # maximum of these both values here:
return max(a.currMem, a.maxMem) return max(a.currMem, a.maxMem)
var
allocator: TAllocator
proc llAlloc(a: var TAllocator, size: int): pointer = proc llAlloc(a: var TAllocator, size: int): pointer =
# *low-level* alloc for the memory managers data structures. Deallocation # *low-level* alloc for the memory managers data structures. Deallocation
# is never done. # is never done.
@ -192,10 +191,10 @@ proc IntSetGet(t: TIntSet, key: int): PTrunk =
it = it.next it = it.next
result = nil result = nil
proc IntSetPut(t: var TIntSet, key: int): PTrunk = proc IntSetPut(a: var TAllocator, t: var TIntSet, key: int): PTrunk =
result = IntSetGet(t, key) result = IntSetGet(t, key)
if result == nil: if result == nil:
result = cast[PTrunk](llAlloc(allocator, sizeof(result[]))) result = cast[PTrunk](llAlloc(a, sizeof(result[])))
result.next = t.data[key and high(t.data)] result.next = t.data[key and high(t.data)]
t.data[key and high(t.data)] = result t.data[key and high(t.data)] = result
result.key = key result.key = key
@ -208,8 +207,8 @@ proc Contains(s: TIntSet, key: int): bool =
else: else:
result = false result = false
proc Incl(s: var TIntSet, key: int) = proc Incl(a: var TAllocator, s: var TIntSet, key: int) =
var t = IntSetPut(s, key shr TrunkShift) var t = IntSetPut(a, s, key shr TrunkShift)
var u = key and TrunkMask var u = key and TrunkMask
t.bits[u shr IntShift] = t.bits[u shr IntShift] or (1 shl (u and IntMask)) t.bits[u shr IntShift] = t.bits[u shr IntShift] or (1 shl (u and IntMask))
@ -220,18 +219,6 @@ proc Excl(s: var TIntSet, key: int) =
t.bits[u shr IntShift] = t.bits[u shr IntShift] and not t.bits[u shr IntShift] = t.bits[u shr IntShift] and not
(1 shl (u and IntMask)) (1 shl (u and IntMask))
proc ContainsOrIncl(s: var TIntSet, key: int): bool =
var t = IntSetGet(s, key shr TrunkShift)
if t != nil:
var u = key and TrunkMask
result = (t.bits[u shr IntShift] and (1 shl (u and IntMask))) != 0
if not result:
t.bits[u shr IntShift] = t.bits[u shr IntShift] or
(1 shl (u and IntMask))
else:
Incl(s, key)
result = false
# ------------- chunk management ---------------------------------------------- # ------------- chunk management ----------------------------------------------
proc pageIndex(c: PChunk): int {.inline.} = proc pageIndex(c: PChunk): int {.inline.} =
result = cast[TAddress](c) shr PageShift result = cast[TAddress](c) shr PageShift
@ -241,9 +228,7 @@ proc pageIndex(p: pointer): int {.inline.} =
proc pageAddr(p: pointer): PChunk {.inline.} = proc pageAddr(p: pointer): PChunk {.inline.} =
result = cast[PChunk](cast[TAddress](p) and not PageMask) result = cast[PChunk](cast[TAddress](p) and not PageMask)
assert(Contains(allocator.chunkStarts, pageIndex(result))) #assert(Contains(allocator.chunkStarts, pageIndex(result)))
var lastSize = PageSize
proc requestOsChunks(a: var TAllocator, size: int): PBigChunk = proc requestOsChunks(a: var TAllocator, size: int): PBigChunk =
incCurrMem(a, size) incCurrMem(a, size)
@ -263,6 +248,7 @@ proc requestOsChunks(a: var TAllocator, size: int): PBigChunk =
#echo("Next already allocated!") #echo("Next already allocated!")
next.prevSize = size next.prevSize = size
# set result.prevSize: # set result.prevSize:
var lastSize = if a.lastSize != 0: a.lastSize else: PageSize
var prv = cast[TAddress](result) -% lastSize var prv = cast[TAddress](result) -% lastSize
assert((nxt and PageMask) == 0) assert((nxt and PageMask) == 0)
var prev = cast[PChunk](prv) var prev = cast[PChunk](prv)
@ -271,7 +257,7 @@ proc requestOsChunks(a: var TAllocator, size: int): PBigChunk =
result.prevSize = lastSize result.prevSize = lastSize
else: else:
result.prevSize = 0 # unknown result.prevSize = 0 # unknown
lastSize = size # for next request a.lastSize = size # for next request
proc freeOsChunks(a: var TAllocator, p: pointer, size: int) = proc freeOsChunks(a: var TAllocator, p: pointer, size: int) =
# update next.prevSize: # update next.prevSize:
@ -287,8 +273,8 @@ proc freeOsChunks(a: var TAllocator, p: pointer, size: int) =
dec(a.freeMem, size) dec(a.freeMem, size)
#c_fprintf(c_stdout, "[Alloc] back to OS: %ld\n", size) #c_fprintf(c_stdout, "[Alloc] back to OS: %ld\n", size)
proc isAccessible(p: pointer): bool {.inline.} = proc isAccessible(a: TAllocator, p: pointer): bool {.inline.} =
result = Contains(allocator.chunkStarts, pageIndex(p)) result = Contains(a.chunkStarts, pageIndex(p))
proc contains[T](list, x: T): bool = proc contains[T](list, x: T): bool =
var it = list var it = list
@ -337,7 +323,7 @@ proc updatePrevSize(a: var TAllocator, c: PBigChunk,
prevSize: int) {.inline.} = prevSize: int) {.inline.} =
var ri = cast[PChunk](cast[TAddress](c) +% c.size) var ri = cast[PChunk](cast[TAddress](c) +% c.size)
assert((cast[TAddress](ri) and PageMask) == 0) assert((cast[TAddress](ri) and PageMask) == 0)
if isAccessible(ri): if isAccessible(a, ri):
ri.prevSize = prevSize ri.prevSize = prevSize
proc freeBigChunk(a: var TAllocator, c: PBigChunk) = proc freeBigChunk(a: var TAllocator, c: PBigChunk) =
@ -347,7 +333,7 @@ proc freeBigChunk(a: var TAllocator, c: PBigChunk) =
when coalescRight: when coalescRight:
var ri = cast[PChunk](cast[TAddress](c) +% c.size) var ri = cast[PChunk](cast[TAddress](c) +% c.size)
assert((cast[TAddress](ri) and PageMask) == 0) assert((cast[TAddress](ri) and PageMask) == 0)
if isAccessible(ri) and chunkUnused(ri): if isAccessible(a, ri) and chunkUnused(ri):
assert(not isSmallChunk(ri)) assert(not isSmallChunk(ri))
if not isSmallChunk(ri): if not isSmallChunk(ri):
ListRemove(a.freeChunksList, cast[PBigChunk](ri)) ListRemove(a.freeChunksList, cast[PBigChunk](ri))
@ -357,7 +343,7 @@ proc freeBigChunk(a: var TAllocator, c: PBigChunk) =
if c.prevSize != 0: if c.prevSize != 0:
var le = cast[PChunk](cast[TAddress](c) -% c.prevSize) var le = cast[PChunk](cast[TAddress](c) -% c.prevSize)
assert((cast[TAddress](le) and PageMask) == 0) assert((cast[TAddress](le) and PageMask) == 0)
if isAccessible(le) and chunkUnused(le): if isAccessible(a, le) and chunkUnused(le):
assert(not isSmallChunk(le)) assert(not isSmallChunk(le))
if not isSmallChunk(le): if not isSmallChunk(le):
ListRemove(a.freeChunksList, cast[PBigChunk](le)) ListRemove(a.freeChunksList, cast[PBigChunk](le))
@ -366,7 +352,7 @@ proc freeBigChunk(a: var TAllocator, c: PBigChunk) =
c = cast[PBigChunk](le) c = cast[PBigChunk](le)
if c.size < ChunkOsReturn: if c.size < ChunkOsReturn:
incl(a.chunkStarts, pageIndex(c)) incl(a, a.chunkStarts, pageIndex(c))
updatePrevSize(a, c, c.size) updatePrevSize(a, c, c.size)
ListAdd(a.freeChunksList, c) ListAdd(a.freeChunksList, c)
c.used = false c.used = false
@ -383,7 +369,7 @@ proc splitChunk(a: var TAllocator, c: PBigChunk, size: int) =
rest.prevSize = size rest.prevSize = size
updatePrevSize(a, c, rest.size) updatePrevSize(a, c, rest.size)
c.size = size c.size = size
incl(a.chunkStarts, pageIndex(rest)) incl(a, a.chunkStarts, pageIndex(rest))
ListAdd(a.freeChunksList, rest) ListAdd(a.freeChunksList, rest)
proc getBigChunk(a: var TAllocator, size: int): PBigChunk = proc getBigChunk(a: var TAllocator, size: int): PBigChunk =
@ -410,7 +396,7 @@ proc getBigChunk(a: var TAllocator, size: int): PBigChunk =
result = requestOsChunks(a, size) result = requestOsChunks(a, size)
result.prevSize = 0 # XXX why is this needed? result.prevSize = 0 # XXX why is this needed?
result.used = true result.used = true
incl(a.chunkStarts, pageIndex(result)) incl(a, a.chunkStarts, pageIndex(result))
dec(a.freeMem, size) dec(a.freeMem, size)
proc getSmallChunk(a: var TAllocator): PSmallChunk = proc getSmallChunk(a: var TAllocator): PSmallChunk =
@ -472,7 +458,7 @@ proc rawAlloc(a: var TAllocator, requestedSize: int): pointer =
assert c.size == size assert c.size == size
result = addr(c.data) result = addr(c.data)
assert((cast[TAddress](result) and (MemAlign-1)) == 0) assert((cast[TAddress](result) and (MemAlign-1)) == 0)
assert(isAccessible(result)) assert(isAccessible(a, result))
proc rawDealloc(a: var TAllocator, p: pointer) = proc rawDealloc(a: var TAllocator, p: pointer) =
var c = pageAddr(p) var c = pageAddr(p)
@ -509,7 +495,7 @@ proc rawDealloc(a: var TAllocator, p: pointer) =
freeBigChunk(a, cast[PBigChunk](c)) freeBigChunk(a, cast[PBigChunk](c))
proc isAllocatedPtr(a: TAllocator, p: pointer): bool = proc isAllocatedPtr(a: TAllocator, p: pointer): bool =
if isAccessible(p): if isAccessible(a, p):
var c = pageAddr(p) var c = pageAddr(p)
if not chunkUnused(c): if not chunkUnused(c):
if isSmallChunk(c): if isSmallChunk(c):
@ -522,11 +508,12 @@ proc isAllocatedPtr(a: TAllocator, p: pointer): bool =
var c = cast[PBigChunk](c) var c = cast[PBigChunk](c)
result = p == addr(c.data) and cast[ptr TFreeCell](p).zeroField >% 1 result = p == addr(c.data) and cast[ptr TFreeCell](p).zeroField >% 1
var
allocator {.rtlThreadVar.}: TAllocator
# ---------------------- interface to programs ------------------------------- # ---------------------- interface to programs -------------------------------
when not defined(useNimRtl): when not defined(useNimRtl):
var heapLock: TSysLock
InitSysLock(HeapLock)
proc unlockedAlloc(size: int): pointer {.inline.} = proc unlockedAlloc(size: int): pointer {.inline.} =
result = rawAlloc(allocator, size+sizeof(TFreeCell)) result = rawAlloc(allocator, size+sizeof(TFreeCell))
@ -545,18 +532,18 @@ when not defined(useNimRtl):
assert(not isAllocatedPtr(allocator, x)) assert(not isAllocatedPtr(allocator, x))
proc alloc(size: int): pointer = proc alloc(size: int): pointer =
when hasThreadSupport: AquireSys(HeapLock) when hasThreadSupport and hasSharedHeap: AquireSys(HeapLock)
result = unlockedAlloc(size) result = unlockedAlloc(size)
when hasThreadSupport: ReleaseSys(HeapLock) when hasThreadSupport and hasSharedHeap: ReleaseSys(HeapLock)
proc alloc0(size: int): pointer = proc alloc0(size: int): pointer =
result = alloc(size) result = alloc(size)
zeroMem(result, size) zeroMem(result, size)
proc dealloc(p: pointer) = proc dealloc(p: pointer) =
when hasThreadSupport: AquireSys(HeapLock) when hasThreadSupport and hasSharedHeap: AquireSys(HeapLock)
unlockedDealloc(p) unlockedDealloc(p)
when hasThreadSupport: ReleaseSys(HeapLock) when hasThreadSupport and hasSharedHeap: ReleaseSys(HeapLock)
proc ptrSize(p: pointer): int = proc ptrSize(p: pointer): int =
var x = cast[pointer](cast[TAddress](p) -% sizeof(TFreeCell)) var x = cast[pointer](cast[TAddress](p) -% sizeof(TFreeCell))

41
lib/system/atomics.nim Normal file
View file

@ -0,0 +1,41 @@
#
#
# Nimrod's Runtime Library
# (c) Copyright 2011 Andreas Rumpf
#
# See the file "copying.txt", included in this
# distribution, for details about the copyright.
#
## Atomic operations for Nimrod.
when (defined(gcc) or defined(llvm_gcc)) and hasThreadSupport:
proc sync_add_and_fetch(p: var int, val: int): int {.
importc: "__sync_add_and_fetch", nodecl.}
proc sync_sub_and_fetch(p: var int, val: int): int {.
importc: "__sync_sub_and_fetch", nodecl.}
elif defined(vcc) and hasThreadSupport:
proc sync_add_and_fetch(p: var int, val: int): int {.
importc: "NimXadd", nodecl.}
else:
proc sync_add_and_fetch(p: var int, val: int): int {.inline.} =
inc(p, val)
result = p
proc atomicInc(memLoc: var int, x: int): int =
when hasThreadSupport:
result = sync_add_and_fetch(memLoc, x)
else:
inc(memLoc, x)
result = memLoc
proc atomicDec(memLoc: var int, x: int): int =
when hasThreadSupport:
when defined(sync_sub_and_fetch):
result = sync_sub_and_fetch(memLoc, x)
else:
result = sync_add_and_fetch(memLoc, -x)
else:
dec(memLoc, x)
result = memLoc

View file

@ -10,9 +10,6 @@
# Exception handling code. This is difficult because it has # Exception handling code. This is difficult because it has
# to work if there is no more memory (but it doesn't yet!). # to work if there is no more memory (but it doesn't yet!).
const
MaxLocksPerThread = 10
var var
stackTraceNewLine* = "\n" ## undocumented feature; it is replaced by ``<br>`` stackTraceNewLine* = "\n" ## undocumented feature; it is replaced by ``<br>``
## for CGI applications ## for CGI applications
@ -35,111 +32,10 @@ proc chckRange(i, a, b: int): int {.inline, compilerproc.}
proc chckRangeF(x, a, b: float): float {.inline, compilerproc.} proc chckRangeF(x, a, b: float): float {.inline, compilerproc.}
proc chckNil(p: pointer) {.inline, compilerproc.} proc chckNil(p: pointer) {.inline, compilerproc.}
type
PSafePoint = ptr TSafePoint
TSafePoint {.compilerproc, final.} = object
prev: PSafePoint # points to next safe point ON THE STACK
status: int
context: C_JmpBuf
when hasThreadSupport:
# Support for thread local storage:
when defined(windows):
type
TThreadVarSlot {.compilerproc.} = distinct int32
proc TlsAlloc(): TThreadVarSlot {.
importc: "TlsAlloc", stdcall, dynlib: "kernel32".}
proc TlsSetValue(dwTlsIndex: TThreadVarSlot, lpTlsValue: pointer) {.
importc: "TlsSetValue", stdcall, dynlib: "kernel32".}
proc TlsGetValue(dwTlsIndex: TThreadVarSlot): pointer {.
importc: "TlsGetValue", stdcall, dynlib: "kernel32".}
proc ThreadVarAlloc(): TThreadVarSlot {.compilerproc, inline.} =
result = TlsAlloc()
proc ThreadVarSetValue(s: TThreadVarSlot, value: pointer) {.
compilerproc, inline.} =
TlsSetValue(s, value)
proc ThreadVarGetValue(s: TThreadVarSlot): pointer {.
compilerproc, inline.} =
result = TlsGetValue(s)
else:
{.passL: "-pthread".}
{.passC: "-pthread".}
type
TThreadVarSlot {.importc: "pthread_key_t", pure, final,
header: "<sys/types.h>".} = object
proc pthread_getspecific(a1: TThreadVarSlot): pointer {.
importc: "pthread_getspecific", header: "<pthread.h>".}
proc pthread_key_create(a1: ptr TThreadVarSlot,
destruct: proc (x: pointer) {.noconv.}): int32 {.
importc: "pthread_key_create", header: "<pthread.h>".}
proc pthread_key_delete(a1: TThreadVarSlot): int32 {.
importc: "pthread_key_delete", header: "<pthread.h>".}
proc pthread_setspecific(a1: TThreadVarSlot, a2: pointer): int32 {.
importc: "pthread_setspecific", header: "<pthread.h>".}
proc specificDestroy(mem: pointer) {.noconv.} =
# we really need a thread-safe 'dealloc' here:
dealloc(mem)
proc ThreadVarAlloc(): TThreadVarSlot {.compilerproc, inline.} =
discard pthread_key_create(addr(result), specificDestroy)
proc ThreadVarSetValue(s: TThreadVarSlot, value: pointer) {.
compilerproc, inline.} =
discard pthread_setspecific(s, value)
proc ThreadVarGetValue(s: TThreadVarSlot): pointer {.compilerproc, inline.} =
result = pthread_getspecific(s)
type
TGlobals* {.final, pure.} = object
excHandler: PSafePoint
currException: ref E_Base
framePtr: PFrame
locksLen*: int
locks*: array [0..MaxLocksPerThread-1, pointer]
buf: string # cannot be allocated on the stack!
assertBuf: string # we need a different buffer for
# assert, as it raises an exception and
# exception handler needs the buffer too
gAssertionFailed: ref EAssertionFailed
tempFrames: array [0..127, PFrame] # cannot be allocated on the stack!
data: float # compiler should add thread local variables here!
PGlobals* = ptr TGlobals
# XXX it'd be more efficient to not use a global variable for the
# thread storage slot, but to rely on the implementation to assign slot 0
# for us... ;-)
var globalsSlot = ThreadVarAlloc()
#const globalsSlot = TThreadVarSlot(0)
#assert checkSlot.int == globalsSlot.int
proc NewGlobals(): PGlobals =
result = cast[PGlobals](alloc0(sizeof(TGlobals)))
new(result.gAssertionFailed)
result.buf = newStringOfCap(2000)
result.assertBuf = newStringOfCap(2000)
proc AllocThreadLocalStorage*(): pointer {.inl.} =
isMultiThreaded = true
result = NewGlobals()
proc SetThreadLocalStorage*(p: pointer) {.inl.} =
ThreadVarSetValue(globalsSlot, p)
proc GetGlobals*(): PGlobals {.compilerRtl, inl.} =
result = cast[PGlobals](ThreadVarGetValue(globalsSlot))
# create for the main thread:
ThreadVarSetValue(globalsSlot, NewGlobals())
when hasThreadSupport: when hasThreadSupport:
template ThreadGlobals = template ThreadGlobals =
var globals = GetGlobals() var currentThread = ThisThread()
template `||`(varname: expr): expr = globals.varname template `||`(varname: expr): expr = currentThread.g.varname
else: else:
template ThreadGlobals = nil # nothing template ThreadGlobals = nil # nothing

View file

@ -15,10 +15,6 @@
# stack overflows when traversing deep datastructures. This is comparable to # stack overflows when traversing deep datastructures. This is comparable to
# an incremental and generational GC. It should be well-suited for soft real # an incremental and generational GC. It should be well-suited for soft real
# time applications (like games). # time applications (like games).
#
# Future Improvements:
# * Support for multi-threading. However, locks for the reference counting
# might turn out to be too slow.
const const
CycleIncrease = 2 # is a multiplicative increase CycleIncrease = 2 # is a multiplicative increase
@ -64,10 +60,10 @@ type
stat: TGcStat stat: TGcStat
var var
stackBottom: pointer stackBottom {.rtlThreadVar.}: pointer
gch: TGcHeap gch {.rtlThreadVar.}: TGcHeap
cycleThreshold: int = InitialCycleThreshold cycleThreshold {.rtlThreadVar.}: int = InitialCycleThreshold
recGcLock: int = 0 recGcLock {.rtlThreadVar.}: int = 0
# we use a lock to prevent the garbage collector to be triggered in a # we use a lock to prevent the garbage collector to be triggered in a
# finalizer; the collector should not call itself this way! Thus every # finalizer; the collector should not call itself this way! Thus every
# object allocated by a finalizer will not trigger a garbage collection. # object allocated by a finalizer will not trigger a garbage collection.
@ -186,6 +182,15 @@ proc doOperation(p: pointer, op: TWalkOp)
proc forAllChildrenAux(dest: Pointer, mt: PNimType, op: TWalkOp) proc forAllChildrenAux(dest: Pointer, mt: PNimType, op: TWalkOp)
# we need the prototype here for debugging purposes # we need the prototype here for debugging purposes
when hasThreadSupport and hasSharedHeap:
template `--`(x: expr): expr = atomicDec(x, rcIncrement) <% rcIncrement
template `++`(x: expr): stmt = discard atomicInc(x, rcIncrement)
else:
template `--`(x: expr): expr =
Dec(x, rcIncrement)
x <% rcIncrement
template `++`(x: expr): stmt = Inc(x, rcIncrement)
proc prepareDealloc(cell: PCell) = proc prepareDealloc(cell: PCell) =
if cell.typ.finalizer != nil: if cell.typ.finalizer != nil:
# the finalizer could invoke something that # the finalizer could invoke something that
@ -219,13 +224,13 @@ proc decRef(c: PCell) {.inline.} =
writeCell("broken cell", c) writeCell("broken cell", c)
assert(c.refcount >=% rcIncrement) assert(c.refcount >=% rcIncrement)
#if c.refcount <% rcIncrement: quit("leck mich") #if c.refcount <% rcIncrement: quit("leck mich")
if atomicDec(c.refcount, rcIncrement) <% rcIncrement: if --c.refcount:
rtlAddZCT(c) rtlAddZCT(c)
elif canBeCycleRoot(c): elif canBeCycleRoot(c):
rtlAddCycleRoot(c) rtlAddCycleRoot(c)
proc incRef(c: PCell) {.inline.} = proc incRef(c: PCell) {.inline.} =
discard atomicInc(c.refcount, rcIncrement) ++c.refcount
if canBeCycleRoot(c): if canBeCycleRoot(c):
rtlAddCycleRoot(c) rtlAddCycleRoot(c)
@ -245,10 +250,10 @@ proc asgnRefNoCycle(dest: ppointer, src: pointer) {.compilerProc, inline.} =
# cycle is possible. # cycle is possible.
if src != nil: if src != nil:
var c = usrToCell(src) var c = usrToCell(src)
discard atomicInc(c.refcount, rcIncrement) ++c.refcount
if dest[] != nil: if dest[] != nil:
var c = usrToCell(dest[]) var c = usrToCell(dest[])
if atomicDec(c.refcount, rcIncrement) <% rcIncrement: if --c.refcount:
rtlAddZCT(c) rtlAddZCT(c)
dest[] = src dest[] = src
@ -517,7 +522,17 @@ proc gcMark(p: pointer) {.inline.} =
proc markThreadStacks(gch: var TGcHeap) = proc markThreadStacks(gch: var TGcHeap) =
when hasThreadSupport: when hasThreadSupport:
nil var it = threadList
while it != nil:
# mark registers:
for i in 0 .. high(it.registers): gcMark(it.registers[i])
var sp = cast[TAddress](it.stackBottom)
var max = cast[TAddress](it.stackTop)
# XXX unroll this loop:
while sp <=% max:
gcMark(cast[ppointer](sp)[])
sp = sp +% sizeof(pointer)
it = it.next
# ----------------- stack management -------------------------------------- # ----------------- stack management --------------------------------------
# inspired from Smart Eiffel # inspired from Smart Eiffel
@ -684,7 +699,7 @@ proc unmarkStackAndRegisters(gch: var TGcHeap) =
# decRef(d[i]) inlined: cannot create a cycle and must not aquire lock # decRef(d[i]) inlined: cannot create a cycle and must not aquire lock
var c = d[i] var c = d[i]
# XXX no need for an atomic dec here: # XXX no need for an atomic dec here:
if atomicDec(c.refcount, rcIncrement) <% rcIncrement: if --c.refcount:
addZCT(gch.zct, c) addZCT(gch.zct, c)
assert c.typ != nil assert c.typ != nil
gch.decStack.len = 0 gch.decStack.len = 0

View file

@ -97,7 +97,7 @@ proc reprSetAux(result: var string, p: pointer, typ: PNimType) =
inc(elemCounter) inc(elemCounter)
if typ.size <= 8: if typ.size <= 8:
for i in 0..sizeof(int64)*8-1: for i in 0..sizeof(int64)*8-1:
if (u and (1 shl i)) != 0: if (u and (1'i64 shl int64(i))) != 0'i64:
if elemCounter > 0: add result, ", " if elemCounter > 0: add result, ", "
addSetElem(result, i+typ.node.len, typ.base) addSetElem(result, i+typ.node.len, typ.base)
inc(elemCounter) inc(elemCounter)

View file

@ -1,98 +0,0 @@
#
#
# Nimrod's Runtime Library
# (c) Copyright 2011 Andreas Rumpf
#
# See the file "copying.txt", included in this
# distribution, for details about the copyright.
#
const
maxThreads = 256
SystemInclude = defined(hasThreadSupport)
when not SystemInclude:
# ugly hack: this file is then included from core/threads, so we have
# thread support:
const hasThreadSupport = true
include "lib/system/ansi_c"
when (defined(gcc) or defined(llvm_gcc)) and hasThreadSupport:
proc sync_add_and_fetch(p: var int, val: int): int {.
importc: "__sync_add_and_fetch", nodecl.}
proc sync_sub_and_fetch(p: var int, val: int): int {.
importc: "__sync_sub_and_fetch", nodecl.}
elif defined(vcc) and hasThreadSupport:
proc sync_add_and_fetch(p: var int, val: int): int {.
importc: "NimXadd", nodecl.}
else:
proc sync_add_and_fetch(p: var int, val: int): int {.inline.} =
inc(p, val)
result = p
proc atomicInc(memLoc: var int, x: int): int =
when hasThreadSupport:
result = sync_add_and_fetch(memLoc, x)
else:
inc(memLoc, x)
result = memLoc
proc atomicDec(memLoc: var int, x: int): int =
when hasThreadSupport:
when defined(sync_sub_and_fetch):
result = sync_sub_and_fetch(memLoc, x)
else:
result = sync_add_and_fetch(memLoc, -x)
else:
dec(memLoc, x)
result = memLoc
when defined(Windows):
type
TSysLock {.final, pure.} = object # CRITICAL_SECTION in WinApi
DebugInfo: pointer
LockCount: int32
RecursionCount: int32
OwningThread: int
LockSemaphore: int
Reserved: int32
proc InitSysLock(L: var TSysLock) {.stdcall,
dynlib: "kernel32", importc: "InitializeCriticalSection".}
## Initializes the lock `L`.
proc TryAquireSysAux(L: var TSysLock): int32 {.stdcall,
dynlib: "kernel32", importc: "TryEnterCriticalSection".}
## Tries to aquire the lock `L`.
proc TryAquireSys(L: var TSysLock): bool {.inline.} =
result = TryAquireSysAux(L) != 0'i32
proc AquireSys(L: var TSysLock) {.stdcall,
dynlib: "kernel32", importc: "EnterCriticalSection".}
## Aquires the lock `L`.
proc ReleaseSys(L: var TSysLock) {.stdcall,
dynlib: "kernel32", importc: "LeaveCriticalSection".}
## Releases the lock `L`.
else:
type
TSysLock {.importc: "pthread_mutex_t", pure, final,
header: "<sys/types.h>".} = object
proc InitSysLock(L: var TSysLock, attr: pointer = nil) {.
importc: "pthread_mutex_init", header: "<pthread.h>".}
proc AquireSys(L: var TSysLock) {.
importc: "pthread_mutex_lock", header: "<pthread.h>".}
proc TryAquireSysAux(L: var TSysLock): cint {.
importc: "pthread_mutex_trylock", header: "<pthread.h>".}
proc TryAquireSys(L: var TSysLock): bool {.inline.} =
result = TryAquireSysAux(L) == 0'i32
proc ReleaseSys(L: var TSysLock) {.
importc: "pthread_mutex_unlock", header: "<pthread.h>".}

481
lib/system/threads.nim Executable file
View file

@ -0,0 +1,481 @@
#
#
# Nimrod's Runtime Library
# (c) Copyright 2011 Andreas Rumpf
#
# See the file "copying.txt", included in this
# distribution, for details about the copyright.
#
## Thread support for Nimrod. **Note**: This is part of the system module.
## Do not import it directly. To active thread support you need to compile
## with the ``--threads:on`` command line switch.
##
## Nimrod's memory model for threads is quite different from other common
## programming languages (C, Pascal): Each thread has its own
## (garbage collected) heap and sharing of memory is restricted. This helps
## to prevent race conditions and improves efficiency. See the manual for
## details of this memory model.
##
## Example:
##
## .. code-block:: nimrod
##
## var
## thr: array [0..4, TThread[tuple[a,b: int]]]
## L: TLock
##
## proc threadFunc(interval: tuple[a,b: int]) {.procvar.} =
## for i in interval.a..interval.b:
## Aquire(L) # lock stdout
## echo i
## Release(L)
##
## InitLock(L)
##
## for i in 0..high(thr):
## createThread(thr[i], threadFunc, (i*10, i*10+5))
## joinThreads(thr)
const
maxRegisters = 256 # don't think there is an arch with more registers
maxLocksPerThread* = 10 ## max number of locks a thread can hold
## at the same time
when defined(Windows):
type
TSysLock {.final, pure.} = object # CRITICAL_SECTION in WinApi
DebugInfo: pointer
LockCount: int32
RecursionCount: int32
OwningThread: int
LockSemaphore: int
Reserved: int32
proc InitSysLock(L: var TSysLock) {.stdcall,
dynlib: "kernel32", importc: "InitializeCriticalSection".}
## Initializes the lock `L`.
proc TryAquireSysAux(L: var TSysLock): int32 {.stdcall,
dynlib: "kernel32", importc: "TryEnterCriticalSection".}
## Tries to aquire the lock `L`.
proc TryAquireSys(L: var TSysLock): bool {.inline.} =
result = TryAquireSysAux(L) != 0'i32
proc AquireSys(L: var TSysLock) {.stdcall,
dynlib: "kernel32", importc: "EnterCriticalSection".}
## Aquires the lock `L`.
proc ReleaseSys(L: var TSysLock) {.stdcall,
dynlib: "kernel32", importc: "LeaveCriticalSection".}
## Releases the lock `L`.
type
THandle = int
TSysThread = THandle
TWinThreadProc = proc (x: pointer): int32 {.stdcall.}
proc CreateThread(lpThreadAttributes: Pointer, dwStackSize: int32,
lpStartAddress: TWinThreadProc,
lpParameter: Pointer,
dwCreationFlags: int32,
lpThreadId: var int32): TSysThread {.
stdcall, dynlib: "kernel32", importc: "CreateThread".}
proc winSuspendThread(hThread: TSysThread): int32 {.
stdcall, dynlib: "kernel32", importc: "SuspendThread".}
proc winResumeThread(hThread: TSysThread): int32 {.
stdcall, dynlib: "kernel32", importc: "ResumeThread".}
proc WaitForMultipleObjects(nCount: int32,
lpHandles: ptr TSysThread,
bWaitAll: int32,
dwMilliseconds: int32): int32 {.
stdcall, dynlib: "kernel32", importc: "WaitForMultipleObjects".}
proc WaitForSingleObject(hHandle: TSysThread, dwMilliseconds: int32): int32 {.
stdcall, dynlib: "kernel32", importc: "WaitForSingleObject".}
proc TerminateThread(hThread: TSysThread, dwExitCode: int32): int32 {.
stdcall, dynlib: "kernel32", importc: "TerminateThread".}
type
TThreadVarSlot {.compilerproc.} = distinct int32
proc TlsAlloc(): TThreadVarSlot {.
importc: "TlsAlloc", stdcall, dynlib: "kernel32".}
proc TlsSetValue(dwTlsIndex: TThreadVarSlot, lpTlsValue: pointer) {.
importc: "TlsSetValue", stdcall, dynlib: "kernel32".}
proc TlsGetValue(dwTlsIndex: TThreadVarSlot): pointer {.
importc: "TlsGetValue", stdcall, dynlib: "kernel32".}
proc ThreadVarAlloc(): TThreadVarSlot {.compilerproc, inline.} =
result = TlsAlloc()
proc ThreadVarSetValue(s: TThreadVarSlot, value: pointer) {.
compilerproc, inline.} =
TlsSetValue(s, value)
proc ThreadVarGetValue(s: TThreadVarSlot): pointer {.
compilerproc, inline.} =
result = TlsGetValue(s)
else:
{.passL: "-pthread".}
{.passC: "-pthread".}
type
TSysLock {.importc: "pthread_mutex_t", pure, final,
header: "<sys/types.h>".} = object
proc InitSysLock(L: var TSysLock, attr: pointer = nil) {.
importc: "pthread_mutex_init", header: "<pthread.h>".}
proc AquireSys(L: var TSysLock) {.
importc: "pthread_mutex_lock", header: "<pthread.h>".}
proc TryAquireSysAux(L: var TSysLock): cint {.
importc: "pthread_mutex_trylock", header: "<pthread.h>".}
proc TryAquireSys(L: var TSysLock): bool {.inline.} =
result = TryAquireSysAux(L) == 0'i32
proc ReleaseSys(L: var TSysLock) {.
importc: "pthread_mutex_unlock", header: "<pthread.h>".}
type
TSysThread {.importc: "pthread_t", header: "<sys/types.h>",
final, pure.} = object
Tpthread_attr {.importc: "pthread_attr_t",
header: "<sys/types.h>", final, pure.} = object
Ttimespec {.importc: "struct timespec",
header: "<time.h>", final, pure.} = object
tv_sec: int
tv_nsec: int
proc pthread_attr_init(a1: var TPthread_attr) {.
importc, header: "<pthread.h>".}
proc pthread_attr_setstacksize(a1: var TPthread_attr, a2: int) {.
importc, header: "<pthread.h>".}
proc pthread_create(a1: var TSysThread, a2: var TPthread_attr,
a3: proc (x: pointer) {.noconv.},
a4: pointer): cint {.importc: "pthread_create",
header: "<pthread.h>".}
proc pthread_join(a1: TSysThread, a2: ptr pointer): cint {.
importc, header: "<pthread.h>".}
proc pthread_cancel(a1: TSysThread): cint {.
importc: "pthread_cancel", header: "<pthread.h>".}
proc AquireSysTimeoutAux(L: var TSysLock, timeout: var Ttimespec): cint {.
importc: "pthread_mutex_timedlock", header: "<time.h>".}
proc AquireSysTimeout(L: var TSysLock, msTimeout: int) {.inline.} =
var a: Ttimespec
a.tv_sec = msTimeout div 1000
a.tv_nsec = (msTimeout mod 1000) * 1000
var res = AquireSysTimeoutAux(L, a)
if res != 0'i32: raise newException(EResourceExhausted, $strerror(res))
type
TThreadVarSlot {.importc: "pthread_key_t", pure, final,
header: "<sys/types.h>".} = object
proc pthread_getspecific(a1: TThreadVarSlot): pointer {.
importc: "pthread_getspecific", header: "<pthread.h>".}
proc pthread_key_create(a1: ptr TThreadVarSlot,
destruct: proc (x: pointer) {.noconv.}): int32 {.
importc: "pthread_key_create", header: "<pthread.h>".}
proc pthread_key_delete(a1: TThreadVarSlot): int32 {.
importc: "pthread_key_delete", header: "<pthread.h>".}
proc pthread_setspecific(a1: TThreadVarSlot, a2: pointer): int32 {.
importc: "pthread_setspecific", header: "<pthread.h>".}
proc ThreadVarAlloc(): TThreadVarSlot {.compilerproc, inline.} =
discard pthread_key_create(addr(result), nil)
proc ThreadVarSetValue(s: TThreadVarSlot, value: pointer) {.
compilerproc, inline.} =
discard pthread_setspecific(s, value)
proc ThreadVarGetValue(s: TThreadVarSlot): pointer {.compilerproc, inline.} =
result = pthread_getspecific(s)
type
TGlobals {.final, pure.} = object
excHandler: PSafePoint
currException: ref E_Base
framePtr: PFrame
buf: string # cannot be allocated on the stack!
assertBuf: string # we need a different buffer for
# assert, as it raises an exception and
# exception handler needs the buffer too
gAssertionFailed: ref EAssertionFailed
tempFrames: array [0..127, PFrame] # cannot be allocated on the stack!
data: float # compiler should add thread local variables here!
proc initGlobals(g: var TGlobals) =
new(g.gAssertionFailed)
g.buf = newStringOfCap(2000)
g.assertBuf = newStringOfCap(2000)
type
PGcThread = ptr TGcThread
TGcThread {.pure.} = object
sys: TSysThread
next, prev: PGcThread
stackBottom, stackTop: pointer
stackSize: int
g: TGlobals
locksLen: int
locks: array [0..MaxLocksPerThread-1, pointer]
registers: array[0..maxRegisters-1, pointer] # register contents for GC
# XXX it'd be more efficient to not use a global variable for the
# thread storage slot, but to rely on the implementation to assign slot 0
# for us... ;-)
var globalsSlot = ThreadVarAlloc()
#const globalsSlot = TThreadVarSlot(0)
#assert checkSlot.int == globalsSlot.int
proc ThisThread(): PGcThread {.compilerRtl, inl.} =
result = cast[PGcThread](ThreadVarGetValue(globalsSlot))
# create for the main thread. Note: do not insert this data into the list
# of all threads; it's not to be stopped etc.
when not defined(useNimRtl):
var mainThread: TGcThread
initGlobals(mainThread.g)
ThreadVarSetValue(globalsSlot, addr(mainThread))
var heapLock: TSysLock
InitSysLock(HeapLock)
var
threadList: PGcThread
proc registerThread(t: PGcThread) =
# we need to use the GC global lock here!
AquireSys(HeapLock)
t.prev = nil
t.next = threadList
if threadList != nil:
assert(threadList.prev == nil)
threadList.prev = t
threadList = t
ReleaseSys(HeapLock)
proc unregisterThread(t: PGcThread) =
# we need to use the GC global lock here!
AquireSys(HeapLock)
if t == threadList: threadList = t.next
if t.next != nil: t.next.prev = t.prev
if t.prev != nil: t.prev.next = t.next
# so that a thread can be unregistered twice which might happen if the
# code executes `destroyThread`:
t.next = nil
t.prev = nil
ReleaseSys(HeapLock)
# on UNIX, the GC uses ``SIGFREEZE`` to tell every thread to stop so that
# the GC can examine the stacks?
proc stopTheWord() =
nil
# We jump through some hops here to ensure that Nimrod thread procs can have
# the Nimrod calling convention. This is needed because thread procs are
# ``stdcall`` on Windows and ``noconv`` on UNIX. Alternative would be to just
# use ``stdcall`` since it is mapped to ``noconv`` on UNIX anyway. However,
# the current approach will likely result in less problems later when we have
# GC'ed closures in Nimrod.
type
TThread* {.pure, final.}[TParam] = object of TGcThread ## Nimrod thread.
fn: proc (p: TParam)
data: TParam
template ThreadProcWrapperBody(closure: expr) =
when not hasSharedHeap: initGC() # init the GC for this thread
ThreadVarSetValue(globalsSlot, closure)
var t = cast[ptr TThread[TParam]](closure)
when not hasSharedHeap: stackBottom = addr(t)
t.stackBottom = addr(t)
registerThread(t)
try:
t.fn(t.data)
finally:
unregisterThread(t)
{.push stack_trace:off.}
when defined(windows):
proc threadProcWrapper[TParam](closure: pointer): int32 {.stdcall.} =
ThreadProcWrapperBody(closure)
# implicitely return 0
else:
proc threadProcWrapper[TParam](closure: pointer) {.noconv.} =
ThreadProcWrapperBody(closure)
{.pop.}
proc joinThread*[TParam](t: TThread[TParam]) {.inline.} =
## waits for the thread `t` to finish.
when hostOS == "windows":
discard WaitForSingleObject(t.sys, -1'i32)
else:
discard pthread_join(t.sys, nil)
proc joinThreads*[TParam](t: openArray[TThread[TParam]]) =
## waits for every thread in `t` to finish.
when hostOS == "windows":
var a: array[0..255, TSysThread]
assert a.len >= t.len
for i in 0..t.high: a[i] = t[i].sys
discard WaitForMultipleObjects(t.len, cast[ptr TSysThread](addr(a)), 1, -1)
else:
for i in 0..t.high: joinThread(t[i])
proc destroyThread*[TParam](t: var TThread[TParam]) {.inline.} =
## forces the thread `t` to terminate. This is potentially dangerous if
## you don't have full control over `t` and its aquired resources.
when hostOS == "windows":
discard TerminateThread(t.sys, 1'i32)
else:
discard pthread_cancel(t.sys)
unregisterThread(addr(t.gcInfo))
proc createThread*[TParam](t: var TThread[TParam],
tp: proc (param: TParam),
param: TParam,
stackSize = 1024*256*sizeof(int)) =
## creates a new thread `t` and starts its execution. Entry point is the
## proc `tp`. `param` is passed to `tp`.
t.data = param
t.fn = tp
t.stackSize = stackSize
when hostOS == "windows":
var dummyThreadId: int32
t.sys = CreateThread(nil, stackSize, threadProcWrapper[TParam],
addr(t), 0'i32, dummyThreadId)
else:
var a: Tpthread_attr
pthread_attr_init(a)
pthread_attr_setstacksize(a, stackSize)
if pthread_create(t.sys, a, threadProcWrapper[TParam], addr(t)) != 0:
raise newException(EIO, "cannot create thread")
# --------------------------- lock handling ----------------------------------
type
TLock* = TSysLock ## Nimrod lock
const
noDeadlocks = false # compileOption("deadlockPrevention")
when nodeadlocks:
var
deadlocksPrevented* = 0 ## counts the number of times a
## deadlock has been prevented
proc InitLock*(lock: var TLock) {.inline.} =
## Initializes the lock `lock`.
InitSysLock(lock)
proc OrderedLocks(g: PGcThread): bool =
for i in 0 .. g.locksLen-2:
if g.locks[i] >= g.locks[i+1]: return false
result = true
proc TryAquire*(lock: var TLock): bool {.inline.} =
## Try to aquires the lock `lock`. Returns `true` on success.
when noDeadlocks:
result = TryAquireSys(lock)
if not result: return
# we have to add it to the ordered list. Oh, and we might fail if
# there is no space in the array left ...
var g = ThisThread()
if g.locksLen >= len(g.locks):
ReleaseSys(lock)
raise newException(EResourceExhausted, "cannot aquire additional lock")
# find the position to add:
var p = addr(lock)
var L = g.locksLen-1
var i = 0
while i <= L:
assert g.locks[i] != nil
if g.locks[i] < p: inc(i) # in correct order
elif g.locks[i] == p: return # thread already holds lock
else:
# do the crazy stuff here:
while L >= i:
g.locks[L+1] = g.locks[L]
dec L
g.locks[i] = p
inc(g.locksLen)
assert OrderedLocks(g)
return
# simply add to the end:
g.locks[g.locksLen] = p
inc(g.locksLen)
assert OrderedLocks(g)
else:
result = TryAquireSys(lock)
proc Aquire*(lock: var TLock) =
## Aquires the lock `lock`.
when nodeadlocks:
var g = ThisThread()
var p = addr(lock)
var L = g.locksLen-1
var i = 0
while i <= L:
assert g.locks[i] != nil
if g.locks[i] < p: inc(i) # in correct order
elif g.locks[i] == p: return # thread already holds lock
else:
# do the crazy stuff here:
if g.locksLen >= len(g.locks):
raise newException(EResourceExhausted, "cannot aquire additional lock")
while L >= i:
ReleaseSys(cast[ptr TSysLock](g.locks[L])[])
g.locks[L+1] = g.locks[L]
dec L
# aquire the current lock:
AquireSys(lock)
g.locks[i] = p
inc(g.locksLen)
# aquire old locks in proper order again:
L = g.locksLen-1
inc i
while i <= L:
AquireSys(cast[ptr TSysLock](g.locks[i])[])
inc(i)
# DANGER: We can only modify this global var if we gained every lock!
# NO! We need an atomic increment. Crap.
discard system.atomicInc(deadlocksPrevented, 1)
assert OrderedLocks(g)
return
# simply add to the end:
if g.locksLen >= len(g.locks):
raise newException(EResourceExhausted, "cannot aquire additional lock")
AquireSys(lock)
g.locks[g.locksLen] = p
inc(g.locksLen)
assert OrderedLocks(g)
else:
AquireSys(lock)
proc Release*(lock: var TLock) =
## Releases the lock `lock`.
when nodeadlocks:
var g = ThisThread()
var p = addr(lock)
var L = g.locksLen
for i in countdown(L-1, 0):
if g.locks[i] == p:
for j in i..L-2: g.locks[j] = g.locks[j+1]
dec g.locksLen
break
ReleaseSys(lock)

298
lib/wrappers/zmq.nim Normal file
View file

@ -0,0 +1,298 @@
# Nimrod wrapper of 0mq
# Generated by c2nim with modifications and enhancement from Andreas Rumpf
# Original licence follows:
#
# Copyright (c) 2007-2011 iMatix Corporation
# Copyright (c) 2007-2011 Other contributors as noted in the AUTHORS file
#
# This file is part of 0MQ.
#
# 0MQ is free software; you can redistribute it and/or modify it under
# the terms of the GNU Lesser General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# 0MQ is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Lesser General Public License for more details.
#
# You should have received a copy of the GNU Lesser General Public License
# along with this program. If not, see <http://www.gnu.org/licenses/>.
#
# Generated from zmq version 2.1.5
## Nimrod 0mq wrapper. This file contains the low level C wrappers as well as
## some higher level constructs. The higher level constructs are easily
## recognizable because they are the only ones that have documentation.
{.deadCodeElim: on.}
when defined(windows):
const
zmqdll* = "zmq.dll"
elif defined(macosx):
const
zmqdll* = "libzmq.dylib"
else:
const
zmqdll* = "libzmq.so"
# A number random enough not to collide with different errno ranges on
# different OSes. The assumption is that error_t is at least 32-bit type.
const
HAUSNUMERO* = 156384712
# On Windows platform some of the standard POSIX errnos are not defined.
ENOTSUP* = (HAUSNUMERO + 1)
EPROTONOSUPPORT* = (HAUSNUMERO + 2)
ENOBUFS* = (HAUSNUMERO + 3)
ENETDOWN* = (HAUSNUMERO + 4)
EADDRINUSE* = (HAUSNUMERO + 5)
EADDRNOTAVAIL* = (HAUSNUMERO + 6)
ECONNREFUSED* = (HAUSNUMERO + 7)
EINPROGRESS* = (HAUSNUMERO + 8)
# Native 0MQ error codes.
EFSM* = (HAUSNUMERO + 51)
ENOCOMPATPROTO* = (HAUSNUMERO + 52)
ETERM* = (HAUSNUMERO + 53)
EMTHREAD* = (HAUSNUMERO + 54)
# Maximal size of "Very Small Message". VSMs are passed by value
# to avoid excessive memory allocation/deallocation.
# If VMSs larger than 255 bytes are required, type of 'vsm_size'
# field in msg_t structure should be modified accordingly.
MAX_VSM_SIZE* = 30
POLLIN* = 1
POLLOUT* = 2
POLLERR* = 4
STREAMER* = 1
FORWARDER* = 2
QUEUE* = 3
PAIR* = 0
PUB* = 1
SUB* = 2
REQ* = 3
REP* = 4
DEALER* = 5
ROUTER* = 6
PULL* = 7
PUSH* = 8
XPUB* = 9
XSUB* = 10
XREQ* = DEALER # Old alias, remove in 3.x
XREP* = ROUTER # Old alias, remove in 3.x
UPSTREAM* = PULL # Old alias, remove in 3.x
DOWNSTREAM* = PUSH # Old alias, remove in 3.x
type
# Message types. These integers may be stored in 'content' member of the
# message instead of regular pointer to the data.
TMsgTypes* = enum
DELIMITER = 31,
VSM = 32
# Message flags. MSG_SHARED is strictly speaking not a message flag
# (it has no equivalent in the wire format), however, making it a flag
# allows us to pack the stucture tighter and thus improve performance.
TMsgFlags* = enum
MSG_MORE = 1,
MSG_SHARED = 128,
MSG_MASK = 129 # Merges all the flags
# A message. Note that 'content' is not a pointer to the raw data.
# Rather it is pointer to zmq::msg_content_t structure
# (see src/msg_content.hpp for its definition).
TMsg*{.pure, final.} = object
content*: pointer
flags*: char
vsm_size*: char
vsm_data*: array[0..MAX_VSM_SIZE - 1, char]
TFreeFn = proc (data, hint: pointer) {.noconv.}
TContext {.final, pure.} = object
PContext* = ptr TContext
# Socket Types
TSocket {.final, pure.} = object
PSocket* = ptr TSocket
# Socket options.
TSockOptions* = enum
HWM = 1,
SWAP = 3,
AFFINITY = 4,
IDENTITY = 5,
SUBSCRIBE = 6,
UNSUBSCRIBE = 7,
RATE = 8,
RECOVERY_IVL = 9,
MCAST_LOOP = 10,
SNDBUF = 11,
RCVBUF = 12,
RCVMORE = 13,
FD = 14,
EVENTS = 15,
theTYPE = 16,
LINGER = 17,
RECONNECT_IVL = 18,
BACKLOG = 19,
RECOVERY_IVL_MSEC = 20, # opt. recovery time, reconcile in 3.x
RECONNECT_IVL_MAX = 21
# Send/recv options.
TSendRecvOptions* = enum
NOBLOCK, SNDMORE
TPollItem*{.pure, final.} = object
socket*: PSocket
fd*: cint
events*: cshort
revents*: cshort
# Run-time API version detection
proc version*(major: var cint, minor: var cint, patch: var cint){.cdecl,
importc: "zmq_version", dynlib: zmqdll.}
#****************************************************************************
# 0MQ errors.
#****************************************************************************
# This function retrieves the errno as it is known to 0MQ library. The goal
# of this function is to make the code 100% portable, including where 0MQ
# compiled with certain CRT library (on Windows) is linked to an
# application that uses different CRT library.
proc errno*(): cint{.cdecl, importc: "zmq_errno", dynlib: zmqdll.}
# Resolves system errors and 0MQ errors to human-readable string.
proc strerror*(errnum: cint): cstring {.cdecl, importc: "zmq_strerror",
dynlib: zmqdll.}
#****************************************************************************
# 0MQ message definition.
#****************************************************************************
proc msg_init*(msg: var TMsg): cint{.cdecl, importc: "zmq_msg_init",
dynlib: zmqdll.}
proc msg_init*(msg: var TMsg, size: int): cint{.cdecl,
importc: "zmq_msg_init_size", dynlib: zmqdll.}
proc msg_init*(msg: var TMsg, data: cstring, size: int,
ffn: TFreeFn, hint: pointer): cint{.cdecl,
importc: "zmq_msg_init_data", dynlib: zmqdll.}
proc msg_close*(msg: var TMsg): cint {.cdecl, importc: "zmq_msg_close",
dynlib: zmqdll.}
proc msg_move*(dest, src: var TMsg): cint{.cdecl,
importc: "zmq_msg_move", dynlib: zmqdll.}
proc msg_copy*(dest, src: var TMsg): cint{.cdecl,
importc: "zmq_msg_copy", dynlib: zmqdll.}
proc msg_data*(msg: var TMsg): cstring {.cdecl, importc: "zmq_msg_data",
dynlib: zmqdll.}
proc msg_size*(msg: var TMsg): int {.cdecl, importc: "zmq_msg_size",
dynlib: zmqdll.}
#****************************************************************************
# 0MQ infrastructure (a.k.a. context) initialisation & termination.
#****************************************************************************
proc init*(io_threads: cint): PContext {.cdecl, importc: "zmq_init",
dynlib: zmqdll.}
proc term*(context: PContext): cint {.cdecl, importc: "zmq_term",
dynlib: zmqdll.}
#****************************************************************************
# 0MQ socket definition.
#****************************************************************************
proc socket*(context: PContext, theType: cint): PSocket {.cdecl,
importc: "zmq_socket", dynlib: zmqdll.}
proc close*(s: PSocket): cint{.cdecl, importc: "zmq_close", dynlib: zmqdll.}
proc setsockopt*(s: PSocket, option: cint, optval: pointer,
optvallen: int): cint {.cdecl, importc: "zmq_setsockopt",
dynlib: zmqdll.}
proc getsockopt*(s: PSocket, option: cint, optval: pointer,
optvallen: ptr int): cint{.cdecl,
importc: "zmq_getsockopt", dynlib: zmqdll.}
proc bindAddr*(s: PSocket, address: cstring): cint{.cdecl, importc: "zmq_bind",
dynlib: zmqdll.}
proc connect*(s: PSocket, address: cstring): cint{.cdecl,
importc: "zmq_connect", dynlib: zmqdll.}
proc send*(s: PSocket, msg: var TMsg, flags: cint): cint{.cdecl,
importc: "zmq_send", dynlib: zmqdll.}
proc recv*(s: PSocket, msg: var TMsg, flags: cint): cint{.cdecl,
importc: "zmq_recv", dynlib: zmqdll.}
#****************************************************************************
# I/O multiplexing.
#****************************************************************************
proc poll*(items: ptr TPollItem, nitems: cint, timeout: int): cint{.
cdecl, importc: "zmq_poll", dynlib: zmqdll.}
#****************************************************************************
# Built-in devices
#****************************************************************************
proc device*(device: cint, insocket, outsocket: PSocket): cint{.
cdecl, importc: "zmq_device", dynlib: zmqdll.}
type
EZmq* = object of ESynch ## exception that is raised if something fails
TConnection* {.pure, final.} = object ## a connection
c*: PContext ## the embedded context
s*: PSocket ## the embedded socket
TConnectionMode* = enum ## connection mode
conPAIR = 0,
conPUB = 1,
conSUB = 2,
conREQ = 3,
conREP = 4,
conDEALER = 5,
conROUTER = 6,
conPULL = 7,
conPUSH = 8,
conXPUB = 9,
conXSUB = 10
proc zmqError*() {.noinline, noreturn.} =
## raises EZmq with error message from `zmq.strerror`.
var e: ref EZmq
new(e)
e.msg = $strerror(errno())
raise e
proc open*(address: string, server: bool, mode: TConnectionMode = conDEALER,
numthreads = 4): TConnection =
## opens a new connection. If `server` is true, it uses `bindAddr` for the
## underlying socket, otherwise it opens the socket with `connect`.
result.c = init(cint(numthreads))
if result.c == nil: zmqError()
result.s = socket(result.c, cint(ord(mode)))
if result.s == nil: zmqError()
if server:
if bindAddr(result.s, address) != 0'i32: zmqError()
else:
if connect(result.s, address) != 0'i32: zmqError()
proc close*(c: var TConnection) =
## closes the connection.
if close(c.s) != 0'i32: zmqError()
if term(c.c) != 0'i32: zmqError()
proc send*(c: var TConnection, msg: string) =
## sends a message over the connection.
var m: TMsg
if msg_init(m, msg.len) != 0'i32: zmqError()
copyMem(msg_data(m), cstring(msg), msg.len)
if send(c.s, m, 0'i32) != 0'i32: zmqError()
discard msg_close(m)
proc receive*(c: var TConnection): string =
## receives a message from a connection.
var m: TMsg
if msg_init(m) != 0'i32: zmqError()
if recv(c.s, m, 0'i32) != 0'i32: zmqError()
result = newString(msg_size(m))
copyMem(addr(result[0]), msg_data(m), result.len)
discard msg_close(m)

View file

@ -0,0 +1,69 @@
discard """
outputsub: "101"
"""
import os
const
noDeadlocks = defined(system.deadlocksPrevented)
var
thr: array [0..5, TThread[tuple[a, b: int]]]
L, M, N: TLock
proc doNothing() = nil
proc threadFunc(interval: tuple[a, b: int]) {.procvar.} =
doNothing()
for i in interval.a..interval.b:
when nodeadlocks:
case i mod 6
of 0:
Aquire(L) # lock stdout
Aquire(M)
Aquire(N)
of 1:
Aquire(L)
Aquire(N) # lock stdout
Aquire(M)
of 2:
Aquire(M)
Aquire(L)
Aquire(N)
of 3:
Aquire(M)
Aquire(N)
Aquire(L)
of 4:
Aquire(N)
Aquire(M)
Aquire(L)
of 5:
Aquire(N)
Aquire(L)
Aquire(M)
else: assert false
else:
Aquire(L) # lock stdout
Aquire(M)
echo i
os.sleep(10)
when nodeadlocks:
echo "deadlocks prevented: ", deadlocksPrevented
when nodeadlocks:
Release(N)
Release(M)
Release(L)
InitLock(L)
InitLock(M)
InitLock(N)
proc main =
for i in 0..high(thr):
createThread(thr[i], threadFunc, (i*100, i*100+50))
joinThreads(thr)
main()

View file

@ -1,6 +1,14 @@
* improve ``echo`` code generation for multi-threading
* two issues for thread local heaps:
- must prevent to construct a data structure that contains memory
from different heaps: n.next = otherHeapPtr
- must prevent that GC cleans up memory that other threads can still read...
this can be prevented if the shared heap is simply uncollected (at least
for now)
* add --deadlock_prevention:on|off switch? timeout for locks? * add --deadlock_prevention:on|off switch? timeout for locks?
* make GC fully thread-safe; needs: * make GC fully thread-safe; needs:
- global list of threads
- thread must store its stack boundaries - thread must store its stack boundaries
- GC must traverse these stacks: Even better each thread traverses its - GC must traverse these stacks: Even better each thread traverses its
stack! No need to stop if you can help the GC ;-) stack! No need to stop if you can help the GC ;-)

View file

@ -41,7 +41,8 @@ Changes affecting backwards compatibility
and ``replacef`` operations. and ``replacef`` operations.
- The pointer dereference operation ``p^`` is deprecated and might become - The pointer dereference operation ``p^`` is deprecated and might become
``^p`` in later versions or be dropped entirely since it is rarely used. ``^p`` in later versions or be dropped entirely since it is rarely used.
Use the new notation ``p[]`` to dereference a pointer. Use the new notation ``p[]`` in the rare cases where you need to
dereference a pointer explicitely.
Additions Additions
@ -79,6 +80,7 @@ Additions
- The compiler now might use hashing for string case statements depending - The compiler now might use hashing for string case statements depending
on the number of string literals in the case statement. on the number of string literals in the case statement.
- Added a wrapper for ``redis``. - Added a wrapper for ``redis``.
- Added a wrapper for ``0mq`` via the ``zmq`` module.
- The compiler now supports array, sequence and string slicing. - The compiler now supports array, sequence and string slicing.
- Added ``system.newStringOfCap``. - Added ``system.newStringOfCap``.

View file

@ -24,9 +24,9 @@ file: ticker
doc: "endb;intern;apis;lib;manual;tut1;tut2;nimrodc;overview" doc: "endb;intern;apis;lib;manual;tut1;tut2;nimrodc;overview"
doc: "tools;c2nim;niminst" doc: "tools;c2nim;niminst"
pdf: "manual;lib;tut1;tut2;nimrodc;c2nim;niminst" pdf: "manual;lib;tut1;tut2;nimrodc;c2nim;niminst"
srcdoc: "core/macros;core/threads;core/marshal" srcdoc: "core/macros;core/marshal"
srcdoc: "impure/graphics;pure/sockets" srcdoc: "impure/graphics;pure/sockets"
srcdoc: "system.nim;pure/os;pure/strutils;pure/math" srcdoc: "system.nim;system/threads.nim;pure/os;pure/strutils;pure/math"
srcdoc: "pure/complex;pure/times;pure/osproc;pure/pegs;pure/dynlib" srcdoc: "pure/complex;pure/times;pure/osproc;pure/pegs;pure/dynlib"
srcdoc: "pure/parseopt;pure/hashes;pure/strtabs;pure/lexbase" srcdoc: "pure/parseopt;pure/hashes;pure/strtabs;pure/lexbase"
srcdoc: "pure/parsecfg;pure/parsexml;pure/parsecsv;pure/parsesql" srcdoc: "pure/parsecfg;pure/parsexml;pure/parsecsv;pure/parsesql"
@ -36,8 +36,8 @@ srcdoc: "impure/db_postgres;impure/db_mysql;impure/db_sqlite"
srcdoc: "pure/httpserver;pure/httpclient;pure/smtp;impure/ssl" srcdoc: "pure/httpserver;pure/httpclient;pure/smtp;impure/ssl"
srcdoc: "pure/ropes;pure/unidecode/unidecode;pure/xmldom;pure/xmldomparser" srcdoc: "pure/ropes;pure/unidecode/unidecode;pure/xmldom;pure/xmldomparser"
srcdoc: "pure/xmlparser;pure/htmlparser;pure/xmltree;pure/colors" srcdoc: "pure/xmlparser;pure/htmlparser;pure/xmltree;pure/colors"
srcdoc: "pure/json;pure/base64;pure/scgi;impure/graphics" srcdoc: "pure/json;pure/base64;pure/scgi;pure/redis;impure/graphics"
srcdoc: "impure/rdstdin" srcdoc: "impure/rdstdin;wrappers/zmq"
webdoc: "wrappers/libcurl;pure/md5;wrappers/mysql;wrappers/iup" webdoc: "wrappers/libcurl;pure/md5;wrappers/mysql;wrappers/iup"
webdoc: "wrappers/sqlite3;wrappers/postgres;wrappers/tinyc" webdoc: "wrappers/sqlite3;wrappers/postgres;wrappers/tinyc"